//// `user_backfill.backfill_user` tests against a host-branching stub xrpc //// client (see `actor_shelf_test.gleam`/`crate_overlap_test.gleam` for the //// pattern): cursor-chained multi-page release paging, the page-cap //// degrade-and-keep-partial behaviour, and the login-parity regression //// (releases, adoptions, and edits all seeded, not just releases). import atproto/xrpc import atproto_core/xrpc as core_xrpc import crate/gen/catalog/edit as catalog_edit import crate/gen/catalog/release as catalog_release import crate/gen/graph/follow as graph_follow import crate/gen/shelf/entry as shelf_entry import crate_server/catalog_index import crate_server/known_users.{KnownUser} import crate_server/shelf_index import crate_server/user_backfill import gleam/bit_array import gleam/http/request.{type Request} import gleam/http/response import gleam/int import gleam/json import gleam/list import gleam/option.{None, Some} import gleam/result import gleam/string import support const pds_host = "pds.test" /// `[1, 2, .., count]`: `gleam/list` has no `range`, so build it the way /// `pagination_test.gleam` does. fn ints_up_to(count: Int) -> List(Int) { list.repeat(Nil, count) |> list.index_map(fn(_, i) { i + 1 }) } fn a_user() -> known_users.KnownUser { KnownUser(did: "did:plc:me", handle: "me.test", pds: "https://" <> pds_host) } fn record(uri: String, cid: String, value: json.Json) -> json.Json { json.object([ #("uri", json.string(uri)), #("cid", json.string(cid)), #("value", value), ]) } fn release_uri(rkey: String) -> String { "at://did:plc:me/dev.mokkenstorm.crate.catalog.release/" <> rkey } fn release_record(rkey: String) -> json.Json { record( release_uri(rkey), "bafy" <> rkey, json.object([ #("title", json.string("Title " <> rkey)), #("createdAt", json.string("2026-01-01T00:00:00Z")), ]), ) } fn shelf_entry_record(rkey: String, release_rkey: String) -> json.Json { record( "at://did:plc:me/dev.mokkenstorm.crate.shelf.entry/" <> rkey, "bafy" <> rkey, json.object([ #("action", json.string("acquired")), #("createdAt", json.string("2026-01-01T00:00:00Z")), #( "release", json.object([ #("cid", json.string("bafyrel")), #("uri", json.string(release_uri(release_rkey))), ]), ), ]), ) } fn catalog_edit_record(rkey: String, subject_rkey: String) -> json.Json { record( "at://did:plc:me/dev.mokkenstorm.crate.catalog.edit/" <> rkey, "bafy" <> rkey, json.object([ #("createdAt", json.string("2026-01-01T00:00:00Z")), #("entity", json.string("release")), #("op", json.string("propose")), #( "subject", json.object([ #("cid", json.string("bafyrel")), #("uri", json.string(release_uri(subject_rkey))), ]), ), ]), ) } fn page_body( records: List(json.Json), cursor: option.Option(String), ) -> String { json.object( list.flatten([ [#("records", json.preprocessed_array(records))], case cursor { Some(c) -> [#("cursor", json.string(c))] None -> [] }, ]), ) |> json.to_string } fn empty_page_body() -> String { page_body([], None) } 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) -> known |> list.find(fn(c) { string.contains(q, c) }) |> result.unwrap(catalog_edit.collection) None -> catalog_edit.collection } } fn cursor_of(req: Request(BitArray)) -> option.Option(String) { use q <- option.then(req.query) case string.split(q, "cursor=") { [_, rest] -> rest |> string.split("&") |> list.first |> option.from_result _ -> None } } fn ok_body( body: String, ) -> Result(response.Response(BitArray), xrpc.TransportError) { Ok(response.Response(200, [], bit_array.from_string(body))) } /// Three cursor-chained pages of `catalog.release` (40 + 40 + 15 = 95 /// records, past the pre-fix 50-record single-page ceiling), and a single /// empty page for the other two collections. fn multi_page_release_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req), cursor_of(req) { col, None if col == catalog_release.collection -> ok_body(page_body( ints_up_to(40) |> list.map(fn(n) { release_record("p1-" <> int.to_string(n)) }), Some("p2"), )) col, Some("p2") if col == catalog_release.collection -> ok_body(page_body( ints_up_to(40) |> list.map(fn(n) { release_record("p2-" <> int.to_string(n)) }), Some("p3"), )) col, Some("p3") if col == catalog_release.collection -> ok_body(page_body( ints_up_to(15) |> list.map(fn(n) { release_record("p3-" <> int.to_string(n)) }), None, )) _, _ -> ok_body(empty_page_body()) } }) } pub fn backfill_user_pages_releases_to_exhaustion_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( multi_page_release_client(), a_user(), store, shelf, support.fresh_follow_index(), ) assert list.length(store.releases.list()) == 95 } /// A client that never runs out of cursor: every response points at another /// page. Backfill must still terminate (at `max_pages`), landing exactly one /// distinct release per page paged before the cap, and no more. fn looping_release_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req) { col if col == catalog_release.collection -> { let n = case cursor_of(req) { Some(c) -> int.parse(c) |> result.unwrap(0) None -> 0 } ok_body(page_body( [release_record("loop-" <> int.to_string(n))], Some(int.to_string(n + 1)), )) } _ -> ok_body(empty_page_body()) } }) } pub fn backfill_user_caps_pages_and_keeps_partial_results_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( looping_release_client(), a_user(), store, shelf, support.fresh_follow_index(), ) // 50 pages, one distinct release apiece, then the cap stops further paging. assert list.length(store.releases.list()) == 50 } /// One record per collection: the regression case for the login-backfill /// parity bug, where only releases ever landed in the index. fn one_of_each_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req) { col if col == catalog_release.collection -> ok_body(page_body([release_record("r1")], None)) col if col == shelf_entry.collection -> ok_body(page_body([shelf_entry_record("e1", "r1")], None)) col if col == catalog_edit.collection -> ok_body(page_body([catalog_edit_record("d1", "r1")], None)) _ -> ok_body(empty_page_body()) } }) } /// A `catalog.release` record missing the required `title` field, so /// `catalog_release_decoder` fails on it; paired with a well-formed record to /// prove the malformed one is dropped rather than crashing the rest of the /// page. fn undecodable_release_record(rkey: String) -> json.Json { record( release_uri(rkey), "bafy" <> rkey, json.object([#("createdAt", json.string("2026-01-01T00:00:00Z"))]), ) } fn one_valid_one_undecodable_release_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req) { col if col == catalog_release.collection -> ok_body(page_body( [release_record("r1"), undecodable_release_record("bad")], None, )) _ -> ok_body(empty_page_body()) } }) } /// The regression case for `decode_or_warn`: a page mixing one well-formed /// and one undecodable `catalog.release` record must seed the well-formed /// one and silently (bar the warning log) drop the other, not crash the /// whole page. pub fn backfill_user_drops_undecodable_release_records_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( one_valid_one_undecodable_release_client(), a_user(), store, shelf, support.fresh_follow_index(), ) assert list.length(store.releases.list()) == 1 } /// A `catalog.edit` record missing the required `entity` field, so /// `catalog_edit_decoder` fails on it. fn undecodable_edit_record(rkey: String) -> json.Json { record( "at://did:plc:me/dev.mokkenstorm.crate.catalog.edit/" <> rkey, "bafy" <> rkey, json.object([ #("createdAt", json.string("2026-01-01T00:00:00Z")), #("op", json.string("propose")), ]), ) } fn one_valid_one_undecodable_edit_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case collection_of(req) { col if col == catalog_release.collection -> ok_body(page_body([release_record("r1")], None)) col if col == catalog_edit.collection -> ok_body(page_body( [catalog_edit_record("d1", "r1"), undecodable_edit_record("bad")], None, )) _ -> ok_body(empty_page_body()) } }) } pub fn backfill_user_drops_undecodable_edit_records_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( one_valid_one_undecodable_edit_client(), a_user(), store, shelf, support.fresh_follow_index(), ) assert list.length(store.edits.for_subject(release_uri("r1"))) == 1 } pub fn backfill_user_seeds_adoptions_and_edits_not_just_releases_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( one_of_each_client(), a_user(), store, shelf, support.fresh_follow_index(), ) assert list.length(store.releases.list()) == 1 assert list.length(store.adoptions.list()) == 1 assert list.length(store.edits.for_subject(release_uri("r1"))) == 1 } /// C4 of the appview-first roadmap: the shared `shelf.entry` page also feeds /// `shelf_index`, and the user is marked seen so a later actor-mode read /// never re-triggers a live fetch for this did. pub fn backfill_user_seeds_the_shelf_index_and_marks_seen_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() user_backfill.backfill_user( one_of_each_client(), a_user(), store, shelf, support.fresh_follow_index(), ) let entry_uri = "at://did:plc:me/dev.mokkenstorm.crate.shelf.entry/e1" let assert Some(folded) = shelf.entries.get(entry_uri) assert folded.did == "did:plc:me" assert folded.release_uri == Some(release_uri("r1")) assert shelf.seen.has_seen("did:plc:me") == True } /// A user with no shelf.entry records at all is still marked seen (an empty /// crate, not a never-backfilled one), matching the never-seen-actor /// live-fetch seed's behaviour on the read path. pub fn backfill_user_marks_seen_even_with_no_shelf_entries_test() { let assert Ok(store) = catalog_index.start() let assert Ok(shelf) = shelf_index.start() let no_shelf_client = fn() { xrpc.Client(send: fn(req) { case collection_of(req) { col if col == catalog_release.collection -> ok_body(page_body([release_record("r1")], None)) _ -> ok_body(empty_page_body()) } }) } user_backfill.backfill_user( no_shelf_client(), a_user(), store, shelf, support.fresh_follow_index(), ) 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 follows = support.fresh_follow_index() 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, follows) assert summary.follows == 1 assert follows.edges.following("did:plc:me") == ["did:plc:b"] assert follows.seen.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 follows = support.fresh_follow_index() let summary = user_backfill.backfill_user( follows_client(fn(_q) { empty_page_body() }), a_user(), store, shelf, follows, ) assert summary.follows == 0 assert follows.seen.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() let follows = support.fresh_follow_index() // 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, follows) assert summary.follows == 1 assert follows.edges.following("did:plc:me") == ["did:plc:b"] assert follows.seen.has_seen("did:plc:me") == False }