diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index 722662d..5535716 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -7,6 +7,7 @@ import at_record_server/known_users.{type Store as KnownUsers} import at_record_server/oauth/config import at_record_server/oauth/keys import at_record_server/oauth/store +import at_record_server/reconcile import at_record_server/router import at_record_server/user_backfill import at_record_server/wiring @@ -41,6 +42,15 @@ pub fn main() -> Nil { client, resolver, ) + // Unconditional, same as jetstream_consumer.start above: neither is gated + // on any deploy flag, so the cheap reconcile pass runs wherever the + // firehose consumer does. + reconcile.start( + config_env.reconcile_interval_ms(), + client, + known_users, + catalog_index, + ) let oauth = config.new( client:, diff --git a/server/src/at_record_server/config_env.gleam b/server/src/at_record_server/config_env.gleam index 0c08496..f31410d 100644 --- a/server/src/at_record_server/config_env.gleam +++ b/server/src/at_record_server/config_env.gleam @@ -31,6 +31,12 @@ const default_identity_ttl_seconds = 3600 /// Default `identity_cache` negative TTL: 60 seconds. const default_identity_negative_ttl_seconds = 60 +/// Default reconcile-pass interval: 6 hours, in milliseconds. Cheap version +/// of ADR 0002's consistency model (decision 2, 2026-07-20): a periodic full +/// re-backfill rather than a `rev`-watermark drift check, deferred until +/// user count makes full re-pages wasteful. +const default_reconcile_interval_ms = 21_600_000 + pub fn port() -> Int { env_int("PORT", default_port) } @@ -48,6 +54,11 @@ pub fn identity_negative_ttl_seconds() -> Int { ) } +/// How often the reconcile pass re-backfills every known user. +pub fn reconcile_interval_ms() -> Int { + env_int("RECONCILE_INTERVAL_MS", default_reconcile_interval_ms) +} + /// The cookie signing secret. Required; a random fallback would hide /// misconfiguration and drop sessions on every restart. pub fn secret_key_base() -> String { diff --git a/server/src/at_record_server/reconcile.gleam b/server/src/at_record_server/reconcile.gleam new file mode 100644 index 0000000..63fd91f --- /dev/null +++ b/server/src/at_record_server/reconcile.gleam @@ -0,0 +1,71 @@ +//// The cheap reconcile pass (ADR 0002, decision 2): a supervised loop that +//// every `RECONCILE_INTERVAL_MS` re-runs `user_backfill.backfill_user` for +//// every known user, self-healing any drift between the index and a user's +//// repo (missed events, restarts, a partial earlier backfill). Safe to +//// re-run because every write it triggers is an upsert. Mirrors +//// `jetstream_consumer`'s spawn/sleep/supervise shape, but on a fixed +//// interval rather than reconnect-with-backoff; the `rev`-watermark drift +//// check from the ADR is deferred until full re-pages get too expensive. + +import at_record_server/catalog_index +import at_record_server/clock +import at_record_server/known_users.{type KnownUser} +import at_record_server/user_backfill +import atproto/xrpc.{type Client} +import gleam/erlang/process +import gleam/int +import gleam/list +import gleam/string +import wisp + +/// Fire-and-forget: spawns a supervising process that re-backfills every +/// known user every `interval_ms`, for as long as the server runs; never +/// blocks or crashes the caller. +pub fn start( + interval_ms: Int, + client: Client, + known_users: known_users.Store, + index: catalog_index.Store, +) -> Nil { + process.spawn_unlinked(fn() { loop(interval_ms, client, known_users, index) }) + Nil +} + +fn loop( + interval_ms: Int, + client: Client, + known_users: known_users.Store, + index: catalog_index.Store, +) -> Nil { + process.sleep(interval_ms) + run_sweep(client, known_users, index) + loop(interval_ms, client, known_users, index) +} + +/// One reconcile sweep: re-backfills every known user in `users_for_sweep` +/// order, then logs the user count and wall-clock duration. Exposed for +/// tests: takes its dependencies as plain arguments, so a fake client/store +/// drives it without a real timer. +pub fn run_sweep( + client: Client, + known_users: known_users.Store, + index: catalog_index.Store, +) -> Nil { + let started_at = clock.now_seconds() + let users = users_for_sweep(known_users.list()) + users |> list.each(user_backfill.backfill_user(client, _, index)) + wisp.log_info( + "reconcile: swept " + <> int.to_string(list.length(users)) + <> " user(s) in " + <> int.to_string(clock.now_seconds() - started_at) + <> "s", + ) +} + +/// The sweep order: known users sorted by did, so repeated sweeps process +/// users in a stable, log-diffable order. Pure so it is unit-testable +/// without a real store or timer. +pub fn users_for_sweep(users: List(KnownUser)) -> List(KnownUser) { + users |> list.sort(fn(a, b) { string.compare(a.did, b.did) }) +} diff --git a/server/test/reconcile_test.gleam b/server/test/reconcile_test.gleam new file mode 100644 index 0000000..901f228 --- /dev/null +++ b/server/test/reconcile_test.gleam @@ -0,0 +1,110 @@ +//// `reconcile.users_for_sweep`'s pure ordering, and `run_sweep`'s behavior +//// against a real in-memory `catalog_index` and a host-branching stub xrpc +//// client (see `user_backfill_test.gleam`/`jetstream_consumer_test.gleam`'s +//// `resolving_client` for the pattern): every known user gets backfilled, +//// and re-running the sweep is idempotent (no duplicate rows). + +import at_record/gen/catalog/release as catalog_release +import at_record_server/catalog_index +import at_record_server/known_users.{type KnownUser, KnownUser} +import at_record_server/reconcile +import atproto/xrpc +import gleam/bit_array +import gleam/http/request.{type Request} +import gleam/http/response +import gleam/json +import gleam/list +import gleam/option.{None} +import gleam/string + +fn user_a() -> KnownUser { + KnownUser(did: "did:plc:a", handle: "a.test", pds: "https://pds-a.test") +} + +fn user_b() -> KnownUser { + KnownUser(did: "did:plc:b", handle: "b.test", pds: "https://pds-b.test") +} + +fn known_user_store(users: List(KnownUser)) -> known_users.Store { + known_users.Store(upsert: fn(_) { Nil }, list: fn() { users }) +} + +fn release_uri(did: String) -> String { + "at://" <> did <> "/dev.mokkenstorm.crate.catalog.release/r1" +} + +fn release_record(did: String) -> json.Json { + json.object([ + #("uri", json.string(release_uri(did))), + #("cid", json.string("bafy" <> did)), + #( + "value", + json.object([ + #("title", json.string("Title for " <> did)), + #("createdAt", json.string("2026-01-01T00:00:00Z")), + ]), + ), + ]) +} + +fn page_body(records: List(json.Json)) -> String { + json.object([#("records", json.preprocessed_array(records))]) + |> json.to_string +} + +fn empty_page_body() -> String { + page_body([]) +} + +fn collection_is_release(req: Request(BitArray)) -> Bool { + case req.query { + option.Some(q) -> string.contains(q, catalog_release.collection) + None -> False + } +} + +/// One release record per known user's own `pds` host (mirroring how +/// `user_backfill` calls `repo_list_records(client, user.pds, ...)`, so the +/// request's host is the per-user signal), empty pages for everything else. +fn per_user_release_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + let body = case collection_is_release(req), req.host { + True, "pds-a.test" -> page_body([release_record("did:plc:a")]) + True, "pds-b.test" -> page_body([release_record("did:plc:b")]) + _, _ -> empty_page_body() + } + Ok(response.Response(200, [], bit_array.from_string(body))) + }) +} + +pub fn users_for_sweep_sorts_by_did_test() { + assert reconcile.users_for_sweep([user_b(), user_a()]) == [user_a(), user_b()] +} + +pub fn users_for_sweep_on_no_users_is_empty_test() { + assert reconcile.users_for_sweep([]) == [] +} + +pub fn run_sweep_backfills_every_known_user_test() { + let assert Ok(index) = catalog_index.start() + reconcile.run_sweep( + per_user_release_client(), + known_user_store([user_a(), user_b()]), + index, + ) + let uris = index.releases.list() |> list.map(fn(r) { r.uri }) + assert list.contains(uris, release_uri("did:plc:a")) + assert list.contains(uris, release_uri("did:plc:b")) + assert list.length(uris) == 2 +} + +/// Re-running the sweep must never duplicate rows: every write it triggers +/// is an upsert keyed by uri, the same idempotence `catalog_index` already +/// guarantees for the firehose and the one-shot boot backfill. +pub fn run_sweep_is_idempotent_across_repeated_runs_test() { + let assert Ok(index) = catalog_index.start() + let store = known_user_store([user_a(), user_b()]) + reconcile.run_sweep(per_user_release_client(), store, index) + reconcile.run_sweep(per_user_release_client(), store, index) + assert list.length(index.releases.list()) == 2 +}