diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index 95a93cf..945e4fe 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -1,5 +1,6 @@ //// The catalog index port: durable public stats over releases, adoptions -//// (shelf.entry events that carry a release ref), and edits, plus the +//// (shelf.entry events that carry a release ref), edits, and the follow +//// graph, plus the //// Jetstream replay cursor. Backends live in this module (in-memory, dev //// fallback, actor over a Dict) and in catalog_index_postgres (durable); //// wiring picks one at the composition root based on DATABASE_URL. `Store` @@ -14,6 +15,7 @@ import gleam/list import gleam/option.{type Option, Some} import gleam/otp/actor import gleam/result +import gleam/set.{type Set} /// One shelf.entry event that carries a release ref: an edge from a did's /// crate entry to the release it adopted. Keyed by the entry's own at-uri, so @@ -74,6 +76,13 @@ pub type Edit { ) } +/// One graph.follow record, keyed by its own at-uri. The unfollow path's rkey +/// is derived from that uri rather than stored beside it, so they cannot +/// disagree. +pub type Follow { + Follow(follow_uri: String, did: String, subject: String, created_at: String) +} + pub type ReleaseOps { ReleaseOps( upsert: fn(BrowseRow) -> Nil, @@ -126,6 +135,25 @@ pub type EditOps { ) } +pub type FollowOps { + FollowOps( + upsert: fn(Follow) -> Nil, + delete: fn(String) -> Nil, + delete_for_did: fn(String) -> Nil, + /// The dids `did` follows. Empty is ambiguous alone; pair with + /// `has_seen`. + following: fn(String) -> List(String), + /// The unfollow path's lookup, served without a live `listRecords`. + find: fn(String, String) -> Option(Follow), + /// True once `did`'s follows were paged to exhaustion, so an empty + /// `following` reads as empty rather than never-indexed. + has_seen: fn(String) -> Bool, + mark_seen: fn(String) -> Nil, + /// Total indexed follow edges: mirrors `ReleaseOps.count`. + count: fn() -> Int, + ) +} + pub type CursorOps { CursorOps(save: fn(Int) -> Nil, load: fn() -> Option(Int)) } @@ -135,6 +163,7 @@ pub type Store { releases: ReleaseOps, adoptions: AdoptionOps, edits: EditOps, + follows: FollowOps, cursor: CursorOps, ) } @@ -144,6 +173,8 @@ type State { releases: Dict(String, BrowseRow), adoptions: Dict(String, Adoption), edits: Dict(String, Edit), + follows: Dict(String, Follow), + follows_seen: Set(String), cursor: Option(Int), ) } @@ -168,6 +199,14 @@ type Msg { EditsFor(String, Subject(List(Edit))) EditsForSubjects(List(String), Subject(List(Edit))) CountEdits(Subject(Int)) + UpsertFollow(Follow) + DeleteFollow(String) + DeleteFollowsForDid(String) + Following(String, Subject(List(String))) + FindFollow(String, String, Subject(Option(Follow))) + HasSeenFollows(String, Subject(Bool)) + MarkFollowsSeen(String) + CountFollows(Subject(Int)) SaveCursor(Int) LoadCursor(Subject(Option(Int))) } @@ -177,6 +216,8 @@ fn initial_state() -> State { releases: dict.new(), adoptions: dict.new(), edits: dict.new(), + follows: dict.new(), + follows_seen: set.new(), cursor: option.None, ) } @@ -224,6 +265,22 @@ pub fn start() -> Result(Store, actor.StartError) { }, count: fn() { parallel.ask(subject, CountEdits, or: 0) }, ), + follows: FollowOps( + upsert: fn(follow) { process.send(subject, UpsertFollow(follow)) }, + delete: fn(follow_uri) { process.send(subject, DeleteFollow(follow_uri)) }, + delete_for_did: fn(did) { + process.send(subject, DeleteFollowsForDid(did)) + }, + following: fn(did) { parallel.ask(subject, Following(did, _), or: []) }, + find: fn(did, subject_did) { + parallel.ask(subject, FindFollow(did, subject_did, _), or: option.None) + }, + has_seen: fn(did) { + parallel.ask(subject, HasSeenFollows(did, _), or: False) + }, + mark_seen: fn(did) { process.send(subject, MarkFollowsSeen(did)) }, + count: fn() { parallel.ask(subject, CountFollows, or: 0) }, + ), cursor: CursorOps( save: fn(time_us) { process.send(subject, SaveCursor(time_us)) }, load: fn() { parallel.ask(subject, LoadCursor, or: option.None) }, @@ -331,6 +388,54 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { process.send(reply, dict.size(state.edits)) actor.continue(state) } + UpsertFollow(follow) -> + actor.continue( + State( + ..state, + follows: dict.insert(state.follows, follow.follow_uri, follow), + ), + ) + DeleteFollow(follow_uri) -> + actor.continue( + State(..state, follows: dict.delete(state.follows, follow_uri)), + ) + // Only edges authored by `did`: one pointing at them lives in someone + // else's repo. + DeleteFollowsForDid(did) -> + actor.continue( + State( + ..state, + follows: dict.filter(state.follows, fn(_uri, follow) { + follow.did != did + }), + follows_seen: set.delete(state.follows_seen, did), + ), + ) + Following(did, reply) -> { + process.send(reply, following_dids(state, did)) + actor.continue(state) + } + FindFollow(did, subject_did, reply) -> { + process.send( + reply, + dict.values(state.follows) + |> list.find(fn(f) { f.did == did && f.subject == subject_did }) + |> option.from_result, + ) + actor.continue(state) + } + HasSeenFollows(did, reply) -> { + process.send(reply, set.contains(state.follows_seen, did)) + actor.continue(state) + } + MarkFollowsSeen(did) -> + actor.continue( + State(..state, follows_seen: set.insert(state.follows_seen, did)), + ) + CountFollows(reply) -> { + process.send(reply, dict.size(state.follows)) + actor.continue(state) + } SaveCursor(time_us) -> actor.continue(State(..state, cursor: option.Some(time_us))) LoadCursor(reply) -> { @@ -340,6 +445,12 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { } } +fn following_dids(state: State, did: String) -> List(String) { + dict.values(state.follows) + |> list.filter(fn(f) { f.did == did }) + |> list.map(fn(f) { f.subject }) +} + fn find_by_barcode(state: State, barcode: String) -> Option(BrowseRow) { dict.values(state.releases) |> list.find(fn(row) { row.barcode == Some(barcode) }) diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam index fc3c681..6168818 100644 --- a/server/src/at_record_server/catalog_index_postgres.gleam +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -1,16 +1,19 @@ -//// Postgres catalog_index.Store backend (pog), over four create-if-not-exists +//// Postgres catalog_index.Store backend (pog), over six create-if-not-exists //// tables sharing one pool with the other stores: catalog_releases (browse //// rows), catalog_adoptions (shelf.entry events with a release ref), -//// catalog_edits (catalog.edit events), and jetstream_cursor (a single row +//// catalog_edits (catalog.edit events), catalog_follows (graph.follow edges) +//// with catalog_follows_seen (a once-per-did paged-to-exhaustion marker, +//// mirroring shelf_index_seen), and jetstream_cursor (a single row //// tracking Jetstream replay progress). This is derived, rebuildable data, //// so writes swallow-and-log rather than fail the caller. import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog_index.{ - type Adoption, type Edit, type Store, Adoption, AdoptionOps, CursorOps, Edit, - EditOps, ReleaseOps, Store, + type Adoption, type Edit, type Follow, type Store, Adoption, AdoptionOps, + CursorOps, Edit, EditOps, Follow, FollowOps, ReleaseOps, Store, } import at_record_server/db +import at_record_server/provenance import atproto/blob.{Blob} import gleam/dict.{type Dict} import gleam/dynamic/decode @@ -60,6 +63,16 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { }, count: fn() { count_edits(conn) }, ), + follows: FollowOps( + upsert: fn(follow) { upsert_follow(conn, follow) }, + delete: fn(follow_uri) { delete_follow(conn, follow_uri) }, + delete_for_did: fn(did) { delete_follows_for_did(conn, did) }, + following: fn(did) { following(conn, did) }, + find: fn(did, subject) { find_follow(conn, did, subject) }, + has_seen: fn(did) { has_seen_follows(conn, did) }, + mark_seen: fn(did) { mark_follows_seen(conn, did) }, + count: fn() { count_follows(conn) }, + ), cursor: CursorOps( save: fn(time_us) { save_cursor(conn, time_us) }, load: fn() { load_cursor(conn) }, @@ -180,6 +193,37 @@ fn migrate(conn: pog.Connection) -> Result(Nil, String) { "catalog_edits cid/fields/rationale migration failed", ), ) + use _ <- result.try( + pog.query( + "create table if not exists catalog_follows ( + follow_uri text primary key, + did text not null, + subject text not null, + created_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_follows migration failed"), + ) + // Every follow read is scoped to one follower. + use _ <- result.try( + pog.query( + "create index if not exists catalog_follows_did_idx + on catalog_follows (did)", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_follows did index migration failed"), + ) + use _ <- result.try( + pog.query( + "create table if not exists catalog_follows_seen ( + did text primary key, + seen_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_follows_seen migration failed"), + ) pog.query( "create table if not exists jetstream_cursor ( id text primary key, @@ -519,6 +563,103 @@ fn count_edits(conn: pog.Connection) -> Int { |> db.count(conn) } +const follow_columns = "follow_uri, did, subject, created_at" + +fn follow_row_decoder() -> decode.Decoder(Follow) { + use follow_uri <- decode.field("follow_uri", decode.string) + use did <- decode.field("did", decode.string) + use subject <- decode.field("subject", decode.string) + use created_at <- decode.field("created_at", decode.string) + decode.success(Follow(follow_uri:, did:, subject:, created_at:)) +} + +fn upsert_follow(conn: pog.Connection, follow: Follow) -> Nil { + let query = + pog.query( + "insert into catalog_follows (follow_uri, did, subject, created_at) + values ($1, $2, $3, $4) + on conflict (follow_uri) do update set + did = excluded.did, + subject = excluded.subject, + created_at = excluded.created_at", + ) + |> pog.parameter(pog.text(follow.follow_uri)) + |> pog.parameter(pog.text(follow.did)) + |> pog.parameter(pog.text(follow.subject)) + |> pog.parameter(pog.text(follow.created_at)) + db.write( + query, + conn, + "catalog_follows upsert failed for " <> follow.follow_uri, + ) +} + +fn delete_follow(conn: pog.Connection, follow_uri: String) -> Nil { + pog.query("delete from catalog_follows where follow_uri = $1") + |> pog.parameter(pog.text(follow_uri)) + |> db.discard(conn) +} + +fn delete_follows_for_did(conn: pog.Connection, did: String) -> Nil { + pog.query("delete from catalog_follows where did = $1") + |> pog.parameter(pog.text(did)) + |> db.discard(conn) + pog.query("delete from catalog_follows_seen where did = $1") + |> pog.parameter(pog.text(did)) + |> db.discard(conn) +} + +fn following(conn: pog.Connection, did: String) -> List(String) { + let row = { + use subject <- decode.field("subject", decode.string) + decode.success(subject) + } + pog.query("select subject from catalog_follows where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.returning(row) + |> db.all(conn) +} + +fn find_follow( + conn: pog.Connection, + did: String, + subject: String, +) -> Option(Follow) { + pog.query("select " <> follow_columns <> " + from catalog_follows where did = $1 and subject = $2 limit 1") + |> pog.parameter(pog.text(did)) + |> pog.parameter(pog.text(subject)) + |> pog.returning(follow_row_decoder()) + |> db.first(conn) +} + +fn has_seen_follows(conn: pog.Connection, did: String) -> Bool { + let row = { + use did <- decode.field("did", decode.string) + decode.success(did) + } + pog.query("select did from catalog_follows_seen where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.returning(row) + |> db.exists(conn) +} + +fn mark_follows_seen(conn: pog.Connection, did: String) -> Nil { + let query = + pog.query( + "insert into catalog_follows_seen (did, seen_at) values ($1, $2) + on conflict (did) do nothing", + ) + |> pog.parameter(pog.text(did)) + |> pog.parameter(pog.text(provenance.now_rfc3339())) + db.write(query, conn, "catalog_follows_seen mark failed for " <> did) +} + +fn count_follows(conn: pog.Connection) -> Int { + pog.query("select count(*)::int as count from catalog_follows") + |> db.count(conn) +} + fn save_cursor(conn: pog.Connection, time_us: Int) -> Nil { let query = pog.query( diff --git a/server/src/at_record_server/graph_follows.gleam b/server/src/at_record_server/graph_follows.gleam index a104699..a466fc0 100644 --- a/server/src/at_record_server/graph_follows.gleam +++ b/server/src/at_record_server/graph_follows.gleam @@ -10,6 +10,7 @@ import atproto/xrpc.{type Client, type XrpcError} import gleam/dynamic/decode import gleam/list import gleam/option.{type Option} +import gleam/result pub fn load( client: Client, @@ -26,23 +27,25 @@ pub fn load( } /// Create a follow record for `subject_did`; duplicates are the caller's -/// problem to guard against, per the lexicon's doc comment. +/// problem to guard against, per the lexicon's doc comment. Returns the record +/// written, so a caller mirroring it into the index stores the timestamp that +/// actually reached the PDS. pub fn create( client: Client, session: OauthSession, subject_did: String, -) -> Result(repo.CreatedRecord, XrpcError) { +) -> Result(#(repo.CreatedRecord, GraphFollow), XrpcError) { + let record = + GraphFollow(created_at: provenance.now_rfc3339(), subject: subject_did) repo.create_record( client, session.pds, session.access_token, session.did, follow.collection, - follow.encode_graph_follow(GraphFollow( - created_at: provenance.now_rfc3339(), - subject: subject_did, - )), + follow.encode_graph_follow(record), ) + |> result.map(fn(created) { #(created, record) }) } /// The stored follow record for `subject_did`, if any. diff --git a/server/src/at_record_server/handlers/debug_index.gleam b/server/src/at_record_server/handlers/debug_index.gleam index b97ebac..2793855 100644 --- a/server/src/at_record_server/handlers/debug_index.gleam +++ b/server/src/at_record_server/handlers/debug_index.gleam @@ -22,6 +22,7 @@ fn json_status(ctx: Context) -> Response { #("shelf_index_seen", json.int(ctx.shelf_index.seen.count())), #("catalog_releases", json.int(ctx.catalog_index.releases.count())), #("catalog_edits", json.int(ctx.catalog_index.edits.count())), + #("catalog_follows", json.int(ctx.catalog_index.follows.count())), #("jetstream_cursor", cursor_json(ctx.catalog_index.cursor.load())), ]) |> json.to_string diff --git a/server/src/at_record_server/handlers/feed.gleam b/server/src/at_record_server/handlers/feed.gleam index fdc8d18..4fbf042 100644 --- a/server/src/at_record_server/handlers/feed.gleam +++ b/server/src/at_record_server/handlers/feed.gleam @@ -1,6 +1,9 @@ //// HTTP glue for the feed skeleton: collapses the caller's followed dids' //// recent adoptions into feed items, falling back to a network-wide slice //// (everyone, viewer included) when the follow graph produces nothing. +//// +//// Both halves come from the index, so a warm request touches no PDS +//// (ADR 0002); only a never-seen viewer is seeded from their own repo once. import at_record/gen/feed/get_feed_skeleton as feed_gen import at_record_server/catalog_index.{type Adoption} @@ -28,7 +31,27 @@ const network_namespace = "feed-network" pub fn get_feed_skeleton(req: Request, ctx: Context) -> Response { use id, session <- require_session(req, ctx) - use client, session <- with_pds_client(ctx, id, session) + case ctx.catalog_index.follows.has_seen(session.did) { + True -> + respond( + req, + ctx, + session.did, + ctx.catalog_index.follows.following(session.did), + ) + False -> { + use client, session <- with_pds_client(ctx, id, session) + respond(req, ctx, session.did, seed_follows(ctx, client, session)) + } + } +} + +fn respond( + req: Request, + ctx: Context, + viewer_did: String, + followed: List(String), +) -> Response { let query = wisp.get_query(req) let cursor = list.key_find(query, "cursor") |> option.from_result let limit = @@ -39,9 +62,8 @@ pub fn get_feed_skeleton(req: Request, ctx: Context) -> Response { |> option.unwrap(pagination.default_limit), ) let adoptions = ctx.catalog_index.adoptions.list() - let followed = followed_dids(client, session) let #(items, next_cursor, fallback) = - resolve_page(adoptions, followed, session.did, cursor, limit) + resolve_page(adoptions, followed, viewer_did, cursor, limit) json.object([ #("items", json.array(items, feed_gen.encode_feed_item)), #("fallback", json.bool(fallback)), @@ -54,13 +76,29 @@ pub fn get_feed_skeleton(req: Request, ctx: Context) -> Response { |> wisp.json_response(200) } -/// The caller's followed dids; a load failure (PDS hiccup) reads the same as -/// an empty follow graph rather than failing the whole feed -- that just -/// falls back to the network-wide slice below. -fn followed_dids(client: Client, session: OauthSession) -> List(String) { +// ADR 0002 allows this live read for "an actor the index has never seen". A +// failure stays unmarked so the next request retries rather than freezing the +// viewer into an empty follow list. +fn seed_follows( + ctx: Context, + client: Client, + session: OauthSession, +) -> List(String) { case graph_follows.load(client, session) { - Ok(stored) -> list.map(stored, fn(s) { s.value.subject }) Error(_) -> [] + Ok(stored) -> { + stored + |> list.each(fn(s) { + ctx.catalog_index.follows.upsert(catalog_index.Follow( + follow_uri: s.uri, + did: session.did, + subject: s.value.subject, + created_at: s.value.created_at, + )) + }) + ctx.catalog_index.follows.mark_seen(session.did) + list.map(stored, fn(s) { s.value.subject }) + } } } diff --git a/server/src/at_record_server/handlers/graph.gleam b/server/src/at_record_server/handlers/graph.gleam index 97d21f3..2e64125 100644 --- a/server/src/at_record_server/handlers/graph.gleam +++ b/server/src/at_record_server/handlers/graph.gleam @@ -1,15 +1,19 @@ //// HTTP glue for the caller's follow graph: create/delete a graph.follow -//// record in their own repo. +//// record in their own repo, mirroring each write into the follow index so +//// the feed sees it without waiting for the firehose to echo it back. +import at_record/gen/graph/follow.{type GraphFollow} +import at_record_server/catalog_index import at_record_server/context.{ type Context, error_json, require_session, with_pds_client, } import at_record_server/graph_follows import at_record_server/oauth/sessions.{type OauthSession} +import atproto/uri import atproto/xrpc.{type Client} import gleam/dynamic/decode import gleam/json -import gleam/option.{None, Some} +import gleam/option.{type Option, None, Some} import wisp.{type Request, type Response} fn subject_decoder() -> decode.Decoder(String) { @@ -35,13 +39,29 @@ fn do_follow( use client, session <- with_pds_client(ctx, id, session) case graph_follows.create(client, session, subject_did) { Error(_) -> error_json(502, "could not write follow record to PDS") - Ok(created) -> + Ok(#(created, record)) -> { + index_follow(ctx, session, created.uri, record) json.object([#("uri", json.string(created.uri))]) |> json.to_string |> wisp.json_response(200) + } } } +fn index_follow( + ctx: Context, + session: OauthSession, + follow_uri: String, + record: GraphFollow, +) -> Nil { + ctx.catalog_index.follows.upsert(catalog_index.Follow( + follow_uri:, + did: session.did, + subject: record.subject, + created_at: record.created_at, + )) +} + pub fn unfollow(req: Request, ctx: Context) -> Response { use id, session <- require_session(req, ctx) use body <- wisp.require_json(req) @@ -58,23 +78,50 @@ fn do_unfollow( subject_did: String, ) -> Response { use client, session <- with_pds_client(ctx, id, session) - case graph_follows.load(client, session) { - Error(_) -> error_json(502, "could not load your follows from PDS") - Ok(stored) -> - case graph_follows.find(stored, subject_did) { - None -> error_json(404, "not following that subject") - Some(found) -> delete_and_confirm(client, session, found.rkey) + case stored_follow(ctx, client, session, subject_did) { + Error(Nil) -> error_json(502, "could not load your follows from PDS") + Ok(None) -> error_json(404, "not following that subject") + Ok(Some(follow_uri)) -> delete_and_confirm(ctx, client, session, follow_uri) + } +} + +// Index-first; a caller the index has never seen still falls back to their +// own repo, same exception the feed makes. +fn stored_follow( + ctx: Context, + client: Client, + session: OauthSession, + subject_did: String, +) -> Result(Option(String), Nil) { + case ctx.catalog_index.follows.has_seen(session.did) { + True -> + Ok( + ctx.catalog_index.follows.find(session.did, subject_did) + |> option.map(fn(f) { f.follow_uri }), + ) + False -> + case graph_follows.load(client, session) { + Error(_) -> Error(Nil) + Ok(stored) -> + Ok( + graph_follows.find(stored, subject_did) + |> option.map(fn(s) { s.uri }), + ) } } } fn delete_and_confirm( + ctx: Context, client: Client, session: OauthSession, - rkey: String, + follow_uri: String, ) -> Response { - case graph_follows.delete(client, session, rkey) { + case graph_follows.delete(client, session, uri.rkey(follow_uri)) { Error(_) -> error_json(502, "could not delete follow record on PDS") - Ok(Nil) -> wisp.json_response("{}", 200) + Ok(Nil) -> { + ctx.catalog_index.follows.delete(follow_uri) + wisp.json_response("{}", 200) + } } } diff --git a/server/src/at_record_server/jetstream_consumer.gleam b/server/src/at_record_server/jetstream_consumer.gleam index 34bc6ba..51380e7 100644 --- a/server/src/at_record_server/jetstream_consumer.gleam +++ b/server/src/at_record_server/jetstream_consumer.gleam @@ -1,6 +1,6 @@ //// The live tail for `catalog_index`: a `stratus` websocket client -//// subscribed to Jetstream across all three of our collections -//// (catalog.release, shelf.entry, catalog.edit), routing commits into the +//// subscribed to Jetstream across all four of our collections +//// (catalog.release, shelf.entry, catalog.edit, graph.follow), routing commits into the //// store as they arrive. A supervising process restarts the socket with //// exponential backoff on any drop, resubscribing with a `cursor` a small //// buffer behind the last persisted position so a reconnect never misses an @@ -12,6 +12,7 @@ import at_record/gen/catalog/edit as catalog_edit import at_record/gen/catalog/release as catalog_release import at_record/gen/client as generated_client +import at_record/gen/graph/follow as graph_follow import at_record/gen/shelf/entry as shelf_entry import at_record_server/browse import at_record_server/catalog/row @@ -196,6 +197,7 @@ fn build_url(host: String, cursor: Option(Int)) -> String { #("wantedCollections", catalog_release.collection), #("wantedCollections", shelf_entry.collection), #("wantedCollections", catalog_edit.collection), + #("wantedCollections", graph_follow.collection), ] let params = case cursor { Some(time_us) -> [ @@ -338,6 +340,7 @@ fn purge_did(index: catalog_index.Store, did: String) -> Nil { index.releases.delete_for_did(did) index.adoptions.delete_for_did(did) index.edits.delete_for_did(did) + index.follows.delete_for_did(did) wisp.log_info( "jetstream_consumer: purged catalog rows for deleted account " <> did, ) @@ -352,6 +355,7 @@ pub fn route(state: State, frame: Frame) -> State { c if c == catalog_release.collection -> route_release(state, frame) c if c == shelf_entry.collection -> route_shelf_entry(state, frame) c if c == catalog_edit.collection -> route_edit(state, frame) + c if c == graph_follow.collection -> route_follow(state, frame) _ -> state } } @@ -571,6 +575,54 @@ fn delete_from_shelf_index(state: State, frame: Frame) -> Nil { } } +fn route_follow(state: State, frame: Frame) -> State { + case frame.operation { + "delete" -> { + state.deps.index.follows.delete(at_uri(frame)) + state + } + _ -> upsert_follow(state, frame) + } +} + +// Deliberately never marks the author seen: one firehose edge is no evidence +// the index holds the rest of their graph. +fn upsert_follow(state: State, frame: Frame) -> State { + case frame.record { + Some(record) -> + case decode.run(record, graph_follow.graph_follow_decoder()) { + Ok(follow) -> { + state.deps.index.follows.upsert(catalog_index.Follow( + follow_uri: at_uri(frame), + did: frame.did, + subject: follow.subject, + created_at: follow.created_at, + )) + state + } + Error(e) -> { + wisp.log_warning( + "jetstream_consumer: graph.follow decode failed for " + <> at_uri(frame) + <> ": " + <> string.inspect(e), + ) + state + } + } + None -> { + wisp.log_warning( + "jetstream_consumer: " + <> frame.operation + <> " commit for " + <> at_uri(frame) + <> " missing record", + ) + state + } + } +} + fn route_edit(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { diff --git a/server/src/at_record_server/user_backfill.gleam b/server/src/at_record_server/user_backfill.gleam index 35bec66..f34bcd8 100644 --- a/server/src/at_record_server/user_backfill.gleam +++ b/server/src/at_record_server/user_backfill.gleam @@ -3,7 +3,8 @@ //// decoder, mirroring `shelf_owner.fetch_pages`'s cursor-following. Used by //// both the boot backfill (`at_record_server`, every known user) and the //// first-login backfill (`handlers/oauth`, one new user) so the two seed -//// catalog.release, shelf.entry, and catalog.edit records identically -- +//// catalog.release, shelf.entry, catalog.edit, and graph.follow records +//// identically -- //// closing the gap where login only ever seeded releases, and the one //// where either backfill silently truncated at a single page. //// @@ -17,6 +18,7 @@ import at_record/gen/catalog/edit as catalog_edit import at_record/gen/catalog/release as catalog_release import at_record/gen/client as generated_client +import at_record/gen/graph/follow as graph_follow import at_record/gen/repo/list_records.{type RecordEntry} import at_record/gen/shelf/entry as shelf_entry import at_record_server/browse @@ -49,7 +51,13 @@ const max_pages = 50 /// reconcile pass's drift signal, compared against the shelf index's own /// per-did event count. pub type BackfillSummary { - BackfillSummary(releases: Int, adoptions: Int, shelf_events: Int, edits: Int) + BackfillSummary( + releases: Int, + adoptions: Int, + shelf_events: Int, + edits: Int, + follows: Int, + ) } /// Page-to-exhaustion, best-effort seed of one user's `catalog.release`, @@ -71,6 +79,7 @@ pub fn backfill_user( let shelf_events = backfill_shelf_index(user, shelf_index, shelf_records) shelf_index.seen.mark_seen(user.did) let edits = backfill_edits(client, user, index) + let follows = backfill_follows(client, user, index) wisp.log_info( "user_backfill: seeded " <> int.to_string(releases) @@ -80,10 +89,41 @@ pub fn backfill_user( <> int.to_string(shelf_events) <> " shelf event(s), " <> int.to_string(edits) - <> " edit(s) for " + <> " edit(s), " + <> int.to_string(follows) + <> " follow(s) for " <> user.did, ) - BackfillSummary(releases:, adoptions:, shelf_events:, edits:) + BackfillSummary(releases:, adoptions:, shelf_events:, edits:, follows:) +} + +// Marks seen only when every page landed: a truncated run would make +// `has_seen`'s "fully indexed" promise a lie. An empty graph still marks. +fn backfill_follows( + client: Client, + user: KnownUser, + index: catalog_index.Store, +) -> Int { + let #(rows, complete) = + fetch_all_checked(client, user, graph_follow.collection, fn(record) { + use follow <- result.try(decode_or_warn( + record, + graph_follow.collection, + graph_follow.graph_follow_decoder(), + )) + Ok(catalog_index.Follow( + follow_uri: record.uri, + did: user.did, + subject: follow.subject, + created_at: follow.created_at, + )) + }) + rows |> list.each(index.follows.upsert) + case complete { + True -> index.follows.mark_seen(user.did) + False -> Nil + } + list.length(rows) } fn backfill_releases( @@ -247,6 +287,18 @@ fn fetch_all( fetch_all_raw(client, user, collection) |> list.filter_map(decode_row) } +/// `fetch_all` plus whether every page landed, for callers that go on to +/// claim a repo is fully indexed. +fn fetch_all_checked( + client: Client, + user: KnownUser, + collection: String, + decode_row: fn(RecordEntry) -> Result(a, Nil), +) -> #(List(a), Bool) { + let #(records, complete) = fetch_pages(client, user, collection, None, [], 0) + #(list.filter_map(records, decode_row), complete) +} + /// `fetch_all`'s undecoded form: the raw paged rows, so a caller needing two /// different projections off the same collection (`backfill_adoptions` and /// `backfill_shelf_index`, both over `shelf.entry`) pages it once and @@ -257,7 +309,7 @@ fn fetch_all_raw( user: KnownUser, collection: String, ) -> List(RecordEntry) { - fetch_pages(client, user, collection, None, [], 0) + fetch_pages(client, user, collection, None, [], 0).0 } fn fetch_pages( @@ -267,7 +319,7 @@ fn fetch_pages( cursor: Option(String), acc: List(RecordEntry), page: Int, -) -> List(RecordEntry) { +) -> #(List(RecordEntry), Bool) { case page >= max_pages { True -> { wisp.log_warning( @@ -278,7 +330,7 @@ fn fetch_pages( <> " for " <> user.did, ) - acc + #(acc, False) } False -> { let params = @@ -299,12 +351,12 @@ fn fetch_pages( <> " for " <> user.did, ) - acc + #(acc, False) } Ok(output) -> { let all = list.append(acc, output.records) case output.cursor { - Some("") | None -> all + Some("") | None -> #(all, True) Some(next) -> fetch_pages(client, user, collection, Some(next), all, page + 1) } diff --git a/server/test/catalog_index_postgres_test.gleam b/server/test/catalog_index_postgres_test.gleam index a1f44eb..f78dc24 100644 --- a/server/test/catalog_index_postgres_test.gleam +++ b/server/test/catalog_index_postgres_test.gleam @@ -9,7 +9,9 @@ import envoy import exception import gleam/dict import gleam/erlang/process +import gleam/list import gleam/option.{None, Some} +import gleam/string // A single-connection pool, killed as soon as this call returns (even if // `run` itself crashes): each test opens its own fresh pool, and a normal @@ -243,3 +245,41 @@ fn find(list: List(a), pred: fn(a) -> Bool) -> Result(a, Nil) { } } } + +fn follow(did: String, subject: String, rkey: String) -> catalog_index.Follow { + catalog_index.Follow( + follow_uri: "at://" <> did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey, + did:, + subject:, + created_at: "2024-01-01T00:00:00Z", + ) +} + +pub fn follow_round_trip_test() { + use store <- with_store() + let did = "did:plc:follow-rt" + store.follows.upsert(follow(did, "did:plc:subj-a", "fr1")) + store.follows.upsert(follow(did, "did:plc:subj-b", "fr2")) + assert list.sort(store.follows.following(did), string.compare) + == ["did:plc:subj-a", "did:plc:subj-b"] + + let assert Some(found) = store.follows.find(did, "did:plc:subj-a") + assert found.created_at == "2024-01-01T00:00:00Z" + assert store.follows.find(did, "did:plc:nobody") == None + + store.follows.delete(follow(did, "did:plc:subj-a", "fr1").follow_uri) + assert store.follows.following(did) == ["did:plc:subj-b"] + store.follows.delete_for_did(did) + assert store.follows.following(did) == [] +} + +pub fn follow_seen_round_trip_test() { + use store <- with_store() + let did = "did:plc:follow-seen" + assert store.follows.has_seen(did) == False + store.follows.mark_seen(did) + assert store.follows.has_seen(did) == True + // A purge clears the mark too, so a re-backfill is not skipped. + store.follows.delete_for_did(did) + assert store.follows.has_seen(did) == False +} diff --git a/server/test/catalog_index_test.gleam b/server/test/catalog_index_test.gleam index e473e7f..eabc75b 100644 --- a/server/test/catalog_index_test.gleam +++ b/server/test/catalog_index_test.gleam @@ -3,6 +3,7 @@ import at_record_server/catalog_index.{type Adoption, Adoption, Edit} import gleam/dict import gleam/list import gleam/option.{type Option, None, Some} +import gleam/string fn row(uri: String) -> BrowseRow { row_with_barcode(uri, None) @@ -344,3 +345,60 @@ pub fn edits_list_for_subjects_batches_across_multiple_releases_test() { assert list.length(found) == 2 assert list.all(found, fn(e) { e.subject_uri == r1 || e.subject_uri == r2 }) } + +fn follow(did: String, subject: String, rkey: String) -> catalog_index.Follow { + catalog_index.Follow( + follow_uri: "at://" <> did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey, + did:, + subject:, + created_at: "2024-01-01T00:00:00Z", + ) +} + +pub fn following_lists_only_that_dids_subjects_test() { + let assert Ok(store) = catalog_index.start() + store.follows.upsert(follow("did:a", "did:b", "f1")) + store.follows.upsert(follow("did:a", "did:c", "f2")) + store.follows.upsert(follow("did:z", "did:b", "f3")) + assert list.sort(store.follows.following("did:a"), string.compare) + == ["did:b", "did:c"] + assert store.follows.following("did:z") == ["did:b"] + assert store.follows.count() == 3 +} + +pub fn following_is_empty_until_a_did_is_seen_test() { + let assert Ok(store) = catalog_index.start() + assert store.follows.has_seen("did:a") == False + store.follows.mark_seen("did:a") + assert store.follows.has_seen("did:a") == True + assert store.follows.following("did:a") == [] +} + +pub fn find_returns_the_edge_for_one_subject_test() { + let assert Ok(store) = catalog_index.start() + store.follows.upsert(follow("did:a", "did:b", "f1")) + let assert Some(found) = store.follows.find("did:a", "did:b") + assert found.follow_uri == "at://did:a/dev.mokkenstorm.crate.graph.follow/f1" + assert store.follows.find("did:a", "did:nope") == None + assert store.follows.find("did:other", "did:b") == None +} + +pub fn deleting_a_follow_drops_only_that_edge_test() { + let assert Ok(store) = catalog_index.start() + store.follows.upsert(follow("did:a", "did:b", "f1")) + store.follows.upsert(follow("did:a", "did:c", "f2")) + store.follows.delete("at://did:a/dev.mokkenstorm.crate.graph.follow/f1") + assert store.follows.following("did:a") == ["did:c"] +} + +pub fn delete_for_did_drops_that_dids_follows_and_seen_mark_test() { + let assert Ok(store) = catalog_index.start() + store.follows.upsert(follow("did:a", "did:b", "f1")) + store.follows.upsert(follow("did:z", "did:b", "f2")) + store.follows.mark_seen("did:a") + store.follows.delete_for_did("did:a") + assert store.follows.following("did:a") == [] + assert store.follows.has_seen("did:a") == False + // The edge pointing *at* did:a lives in someone else's repo and stays. + assert store.follows.following("did:z") == ["did:b"] +} diff --git a/server/test/feed_handler_test.gleam b/server/test/feed_handler_test.gleam index 2cc031d..be82cee 100644 --- a/server/test/feed_handler_test.gleam +++ b/server/test/feed_handler_test.gleam @@ -294,3 +294,72 @@ pub fn limit_is_capped_at_the_lexicon_default_test() { assert item_count(body) == 100 assert support.field_present(body, ["cursor"]) } + +/// Any PDS call at all blows up: the assertion for a warm feed request. +fn panicking_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { panic as { "unexpected PDS call: " <> req.path } }) +} + +fn indexed_follow(rkey: String, subject: String) -> catalog_index.Follow { + catalog_index.Follow( + follow_uri: "at://" + <> viewer_did + <> "/dev.mokkenstorm.crate.graph.follow/" + <> rkey, + did: viewer_did, + subject:, + created_at: "2026-01-01T00:00:00Z", + ) +} + +pub fn an_indexed_follow_graph_is_served_without_touching_the_pds_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: "did:plc:f1", + rkey: "e1", + release_uri: release_uri("r1"), + created_at: nth_day(3), + )) + store.adoptions.upsert(adoption( + did: "did:plc:n1", + rkey: "e2", + release_uri: release_uri("r2"), + created_at: nth_day(2), + )) + store.follows.upsert(indexed_follow("3aaa", "did:plc:f1")) + store.follows.mark_seen(viewer_did) + let #(ctx, cfg) = test_context(panicking_client(), store) + let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) + assert res.status == 200 + assert item_actors(simulate.read_body(res)) == ["did:plc:f1"] +} + +pub fn a_never_seen_viewer_is_seeded_from_their_pds_once_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: "did:plc:f1", + rkey: "e1", + release_uri: release_uri("r1"), + created_at: nth_day(3), + )) + let client = + network_client( + list_records_body([follow_record_json("3aaa", "did:plc:f1")]), + ) + let #(ctx, cfg) = test_context(client, store) + assert store.follows.has_seen(viewer_did) == False + let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) + assert res.status == 200 + assert item_actors(simulate.read_body(res)) == ["did:plc:f1"] + // Seeded and marked, so the next request never asks the PDS again. + assert store.follows.following(viewer_did) == ["did:plc:f1"] + assert store.follows.has_seen(viewer_did) == True +} + +pub fn a_failed_seed_leaves_the_viewer_unmarked_for_a_retry_test() { + let assert Ok(store) = catalog_index.start() + let #(ctx, cfg) = test_context(support.unreachable_client(), store) + let res = feed.get_feed_skeleton(authed_get(feed_path(), cfg), ctx) + assert res.status == 200 + assert store.follows.has_seen(viewer_did) == False +} diff --git a/server/test/graph_handler_test.gleam b/server/test/graph_handler_test.gleam index 18590b7..a6d0698 100644 --- a/server/test/graph_handler_test.gleam +++ b/server/test/graph_handler_test.gleam @@ -3,7 +3,8 @@ import at_record/gen/graph/follow.{type GraphFollow, GraphFollow} import at_record/storage.{type StoredItem, StoredItem} -import at_record_server/context.{type Context} +import at_record_server/catalog_index +import at_record_server/context.{type Context, Context} import at_record_server/graph_follows import at_record_server/handlers/graph import at_record_server/oauth/config @@ -175,3 +176,87 @@ pub fn unfollow_returns_404_when_no_follow_exists_test() { let resp = graph.unfollow(req, ctx) assert resp.status == 404 } + +const viewer_did = "did:plc:x" + +fn follow_uri(rkey: String) -> String { + "at://" <> viewer_did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey +} + +fn test_context_with( + client: xrpc.Client, + index: catalog_index.Store, +) -> #(Context, config.Config) { + let cfg = support.stub_config(client) + let ctx = + Context( + ..support.stub_context_with( + cfg, + support.unreachable_catalog_deps(), + fn(_req) { Error("unused") }, + [], + ), + catalog_index: index, + ) + #(ctx, cfg) +} + +/// Fails any listRecords: the assertion that a lookup came from the index. +fn no_list_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") + "/xrpc/com.atproto.repo.createRecord" -> + ok_response(create_record_response("3ccc")) + _ -> panic as { "unexpected PDS call: " <> req.path } + } + }) +} + +pub fn follow_mirrors_the_written_record_into_the_index_test() { + let assert Ok(index) = catalog_index.start() + let #(ctx, cfg) = test_context_with(no_list_client(), index) + let req = + authed_post( + support.xrpc("graph.followUser"), + json.object([#("subject", json.string(target_did))]), + cfg, + ) + assert graph.follow(req, ctx).status == 200 + assert index.follows.following(viewer_did) == [target_did] + let assert Some(edge) = index.follows.find(viewer_did, target_did) + assert edge.follow_uri == follow_uri("3ccc") +} + +pub fn unfollow_finds_the_rkey_in_the_index_without_a_pds_list_test() { + let assert Ok(index) = catalog_index.start() + index.follows.upsert(catalog_index.Follow( + follow_uri: follow_uri("3ddd"), + did: viewer_did, + subject: target_did, + created_at: "2026-01-01T00:00:00Z", + )) + index.follows.mark_seen(viewer_did) + let #(ctx, cfg) = test_context_with(no_list_client(), index) + let req = + authed_post( + support.xrpc("graph.unfollowUser"), + json.object([#("subject", json.string(target_did))]), + cfg, + ) + assert graph.unfollow(req, ctx).status == 200 + assert index.follows.following(viewer_did) == [] +} + +pub fn unfollow_404s_from_the_index_when_the_edge_is_absent_test() { + let assert Ok(index) = catalog_index.start() + index.follows.mark_seen(viewer_did) + let #(ctx, cfg) = test_context_with(no_list_client(), index) + let req = + authed_post( + support.xrpc("graph.unfollowUser"), + json.object([#("subject", json.string(target_did))]), + cfg, + ) + assert graph.unfollow(req, ctx).status == 404 +} diff --git a/server/test/jetstream_consumer_test.gleam b/server/test/jetstream_consumer_test.gleam index e8cb54e..6c6f909 100644 --- a/server/test/jetstream_consumer_test.gleam +++ b/server/test/jetstream_consumer_test.gleam @@ -854,3 +854,91 @@ pub fn replay_from_rewinds_by_the_buffer_test() { pub fn replay_from_never_goes_negative_test() { assert jetstream_consumer.replay_from(500) == 0 } + +fn follow_frame( + did: String, + time_us: Int, + operation: String, + rkey: String, + subject: String, +) -> String { + frame_json( + did:, + time_us:, + operation:, + collection: "dev.mokkenstorm.crate.graph.follow", + rkey:, + cid: "bafyfollowcid", + record: "{\"$type\":\"dev.mokkenstorm.crate.graph.follow\",\"createdAt\":\"2024-02-01T00:00:00Z\",\"subject\":\"" + <> subject + <> "\"}", + ) +} + +pub fn follow_commit_indexes_the_edge_test() { + let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) + let _ = + jetstream_consumer.route( + state, + decode_frame(follow_frame("did:plc:a", 1, "create", "f1", "did:plc:b")), + ) + assert index.follows.following("did:plc:a") == ["did:plc:b"] + let assert Some(edge) = index.follows.find("did:plc:a", "did:plc:b") + assert edge.created_at == "2024-02-01T00:00:00Z" +} + +pub fn follow_commit_does_not_claim_the_graph_is_complete_test() { + let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) + let _ = + jetstream_consumer.route( + state, + decode_frame(follow_frame("did:plc:a", 1, "create", "f1", "did:plc:b")), + ) + assert index.follows.has_seen("did:plc:a") == False +} + +pub fn follow_delete_frame_drops_the_edge_test() { + let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) + let state = + jetstream_consumer.route( + state, + decode_frame(follow_frame("did:plc:a", 1, "create", "f1", "did:plc:b")), + ) + let _ = + jetstream_consumer.route( + state, + decode_frame(follow_frame("did:plc:a", 2, "delete", "f1", "did:plc:b")), + ) + assert index.follows.following("did:plc:a") == [] +} + +pub fn account_deletion_purges_that_dids_follows_test() { + let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) + let state = + jetstream_consumer.route( + state, + decode_frame(follow_frame("did:plc:gone", 1, "create", "f1", "did:plc:b")), + ) + let state = + jetstream_consumer.route( + state, + decode_frame(follow_frame( + "did:plc:stays", + 2, + "create", + "f2", + "did:plc:gone", + )), + ) + let _ = + jetstream_consumer.apply_event( + state, + decode_event(account_frame("did:plc:gone", 3, False, Some("deleted"))), + ) + assert index.follows.following("did:plc:gone") == [] + assert index.follows.following("did:plc:stays") == ["did:plc:gone"] +} diff --git a/server/test/support.gleam b/server/test/support.gleam index c88830e..d0a8c03 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -146,6 +146,16 @@ pub fn empty_catalog_index() -> catalog_index.Store { list_for_subjects: fn(_) { [] }, count: fn() { 0 }, ), + follows: catalog_index.FollowOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + delete_for_did: fn(_) { Nil }, + following: fn(_) { [] }, + find: fn(_, _) { None }, + has_seen: fn(_) { False }, + mark_seen: fn(_) { Nil }, + count: fn() { 0 }, + ), cursor: catalog_index.CursorOps(save: fn(_) { Nil }, load: fn() { None }), ) } diff --git a/server/test/user_backfill_test.gleam b/server/test/user_backfill_test.gleam index e4bec42..70ae801 100644 --- a/server/test/user_backfill_test.gleam +++ b/server/test/user_backfill_test.gleam @@ -6,12 +6,14 @@ import at_record/gen/catalog/edit as catalog_edit import at_record/gen/catalog/release as catalog_release +import at_record/gen/graph/follow as graph_follow import at_record/gen/shelf/entry as shelf_entry import at_record_server/catalog_index import at_record_server/known_users.{KnownUser} import at_record_server/shelf_index import at_record_server/user_backfill import atproto/xrpc +import atproto_core/xrpc as core_xrpc import gleam/bit_array import gleam/http/request.{type Request} import gleam/http/response @@ -115,17 +117,18 @@ fn empty_page_body() -> String { } fn collection_of(req: Request(BitArray)) -> String { + let known = [ + catalog_release.collection, + shelf_entry.collection, + catalog_edit.collection, + graph_follow.collection, + ] case req.query { Some(q) -> - case string.contains(q, catalog_release.collection) { - True -> catalog_release.collection - False -> - case string.contains(q, shelf_entry.collection) { - True -> shelf_entry.collection - False -> catalog_edit.collection - } - } - None -> "" + known + |> list.find(fn(c) { string.contains(q, c) }) + |> result.unwrap(catalog_edit.collection) + None -> catalog_edit.collection } } @@ -352,3 +355,75 @@ pub fn backfill_user_marks_seen_even_with_no_shelf_entries_test() { assert shelf.entries.list_for_did("did:plc:me") == [] assert shelf.seen.has_seen("did:plc:me") == True } + +fn follow_record(rkey: String, subject: String) -> json.Json { + record( + "at://did:plc:me/dev.mokkenstorm.crate.graph.follow/" <> rkey, + "bafy" <> rkey, + json.object([ + #("createdAt", json.string("2026-02-01T00:00:00Z")), + #("subject", json.string(subject)), + ]), + ) +} + +fn follows_client(pages: fn(String) -> String) -> xrpc.Client { + xrpc.Client(send: fn(req) { + case collection_of(req) { + col if col == graph_follow.collection -> + ok_body(pages(option.unwrap(req.query, ""))) + _ -> ok_body(empty_page_body()) + } + }) +} + +pub fn backfill_seeds_follows_and_marks_the_did_seen_test() { + let assert Ok(store) = catalog_index.start() + let assert Ok(shelf) = shelf_index.start() + let client = + follows_client(fn(_q) { + page_body([follow_record("f1", "did:plc:b")], None) + }) + let summary = user_backfill.backfill_user(client, a_user(), store, shelf) + assert summary.follows == 1 + assert store.follows.following("did:plc:me") == ["did:plc:b"] + assert store.follows.has_seen("did:plc:me") == True +} + +pub fn backfill_marks_seen_even_with_an_empty_follow_graph_test() { + let assert Ok(store) = catalog_index.start() + let assert Ok(shelf) = shelf_index.start() + let summary = + user_backfill.backfill_user( + follows_client(fn(_q) { empty_page_body() }), + a_user(), + store, + shelf, + ) + assert summary.follows == 0 + assert store.follows.has_seen("did:plc:me") == True +} + +pub fn a_truncated_follow_run_keeps_rows_but_does_not_mark_seen_test() { + let assert Ok(store) = catalog_index.start() + let assert Ok(shelf) = shelf_index.start() + // First page has a cursor; the follow-up request fails, so the run is + // partial and must not claim the graph is fully indexed. + let client = + xrpc.Client(send: fn(req) { + case + collection_of(req), + string.contains(option.unwrap(req.query, ""), "cursor=next") + { + col, False if col == graph_follow.collection -> + ok_body(page_body([follow_record("f1", "did:plc:b")], Some("next"))) + col, True if col == graph_follow.collection -> + Error(core_xrpc.ConnectionFailed("pds gone")) + _, _ -> ok_body(empty_page_body()) + } + }) + let summary = user_backfill.backfill_user(client, a_user(), store, shelf) + assert summary.follows == 1 + assert store.follows.following("did:plc:me") == ["did:plc:b"] + assert store.follows.has_seen("did:plc:me") == False +}