//// The live tail for `catalog_index`: a `stratus` websocket client //// 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 //// event. Unknown publishers are resolved via `resolveMiniDoc` and cached //// in-process for the life of the connection. Jetstream also emits //// `identity` (handle changes) and `account` (activation/deletion) events //// alongside commits; see `apply_identity`/`apply_account` below. import atproto/xrpc.{type Client} import crate/gen/catalog/edit as catalog_edit import crate/gen/catalog/release as catalog_release import crate/gen/client as generated_client import crate/gen/graph/follow as graph_follow import crate/gen/shelf/entry as shelf_entry import crate_server/browse import crate_server/catalog/row import crate_server/catalog_index import crate_server/follow_index import crate_server/identity_cache.{type Cache, Identity} import crate_server/known_users import crate_server/shelf_index import gleam/dict.{type Dict} import gleam/dynamic/decode.{type Dynamic} import gleam/erlang/process.{type Subject} import gleam/http/request import gleam/int import gleam/json import gleam/list import gleam/option.{type Option, None, Some} import gleam/otp/actor import gleam/result import gleam/string import gleam/uri as gleam_uri import stratus import wisp /// A public Bluesky-operated instance; override with `JETSTREAM_URL` for a /// local devnet or a different region. pub const default_host = "https://jetstream1.us-east.bsky.network/subscribe" /// The first retry after a drop; doubles on every consecutive failure, reset /// to this floor as soon as a connection succeeds. pub const base_backoff_ms = 1000 /// The retry ceiling: Jetstream itself may be down for a while, and there is /// no point hammering it more often than once a minute. pub const max_backoff_ms = 60_000 /// How far behind the last persisted cursor a fresh connection replays from, /// in Jetstream's microsecond timestamps: enough to cover the handful of /// events that might have landed between the last persisted cursor and the /// disconnect. const replay_buffer_us = 2_000_000 /// Persist the cursor every this-many processed events rather than on every /// message, so a fast-moving firehose doesn't turn into a write per event. const persist_every = 20 /// Exposed (not opaque) so tests can build one directly to drive `route`. pub type Deps { Deps( known_users: known_users.Store, index: catalog_index.Store, /// Wired dark (C1+C2 of the appview-first roadmap): every shelf.entry /// commit lands here now, but no read path consumes it yet. shelf_index: shelf_index.Store, follow_index: follow_index.Store, client: Client, resolver: String, /// Threaded from `main` (the same `identity_cache.Cache` passed into /// `Context`), rather than `Context` itself: the consumer has no session /// or request, so only the narrow capability it needs crosses in. identity_cache: Cache, ) } /// The stratus actor's threaded state: `deps` for the store/identity /// capabilities, `cache` for in-process did -> #(handle, pds) resolutions /// (reset on reconnect; resolution is cheap enough that this is fine), and /// `pending` counting events since the cursor was last persisted. Exposed /// (not opaque) so tests can build one directly and drive `route`. pub type State { State(deps: Deps, cache: Dict(String, #(String, String)), pending: Int) } pub fn initial_state(deps: Deps) -> State { State(deps:, cache: dict.new(), pending: 0) } /// One decoded Jetstream commit event, collection-agnostic: `route` decodes /// `record` per-collection and dispatches to the right store operation. /// Exposed (not opaque) so tests can inspect a decoded frame's fields /// directly. pub type Frame { Frame( did: String, time_us: Int, operation: String, collection: String, rkey: String, cid: Option(String), record: Option(Dynamic), ) } /// A decoded Jetstream event of any kind: a repo commit (already handled by /// `route`), a handle change, or an account state change. Exposed (not /// opaque) so tests can pattern-match a decoded event directly. pub type Event { CommitEvent(Frame) IdentityEvent(did: String, handle: Option(String), time_us: Int) AccountEvent(did: String, active: Bool, status: Option(String), time_us: Int) } /// Fire-and-forget: spawns a supervising process that reconnects with /// backoff for as long as the server runs; never blocks or crashes the /// caller. pub fn start( host: String, known_users: known_users.Store, index: catalog_index.Store, shelf_index: shelf_index.Store, follow_index: follow_index.Store, client: Client, resolver: String, identity_cache: Cache, ) -> Nil { let deps = Deps( known_users:, index:, shelf_index:, follow_index:, client:, resolver:, identity_cache:, ) process.spawn_unlinked(fn() { supervise(host, deps, base_backoff_ms) }) Nil } fn supervise(host: String, deps: Deps, backoff_ms: Int) -> Nil { let cursor = deps.index.cursor.load() case connect(host, deps, cursor) { Ok(started) -> { wisp.log_info("jetstream_consumer: connected to " <> host) let monitor = process.monitor(started.pid) let _ = process.new_selector() |> process.select_specific_monitor(monitor, fn(_down) { Nil }) |> process.selector_receive_forever wisp.log_warning("jetstream_consumer: connection lost, reconnecting") supervise(host, deps, base_backoff_ms) } Error(e) -> { wisp.log_warning( "jetstream_consumer: " <> e <> "; retrying in " <> int.to_string(backoff_ms) <> "ms", ) process.sleep(backoff_ms) supervise(host, deps, next_backoff(backoff_ms)) } } } /// The next backoff after a failed attempt: doubles, capped at /// `max_backoff_ms`. Pure so the schedule is unit-testable without a clock /// or a socket. pub fn next_backoff(current_ms: Int) -> Int { case current_ms * 2 > max_backoff_ms { True -> max_backoff_ms False -> current_ms * 2 } } fn connect( host: String, deps: Deps, cursor: Option(Int), ) -> Result(actor.Started(Subject(stratus.InternalMessage(Nil))), String) { let url = build_url(host, cursor) use req <- result.try( request.to(url) |> result.replace_error("invalid Jetstream URL: " <> url), ) stratus.new(req, initial_state(deps)) |> stratus.on_message(handle_message) |> stratus.start |> result.map_error(fn(e) { "Jetstream connect failed: " <> string.inspect(e) }) } fn build_url(host: String, cursor: Option(Int)) -> String { let collections = [ #("wantedCollections", catalog_release.collection), #("wantedCollections", shelf_entry.collection), #("wantedCollections", catalog_edit.collection), #("wantedCollections", graph_follow.collection), ] let params = case cursor { Some(time_us) -> [ #("cursor", int.to_string(replay_from(time_us))), ..collections ] None -> collections } host <> "?" <> gleam_uri.query_to_string(params) } /// The cursor a fresh connection should replay from: a small buffer behind /// the last persisted position, never negative. pub fn replay_from(time_us: Int) -> Int { case time_us > replay_buffer_us { True -> time_us - replay_buffer_us False -> 0 } } fn handle_message( state: State, msg: stratus.Message(Nil), _conn: stratus.Connection, ) -> stratus.Next(State, Nil) { case msg { stratus.Text(text) -> stratus.continue(handle_frame(state, text)) stratus.Binary(_) -> stratus.continue(state) stratus.User(_) -> stratus.continue(state) } } fn handle_frame(state: State, text: String) -> State { case json.parse(text, frame_decoder()) { Ok(Some(event)) -> apply_event(state, event) Ok(None) -> state Error(e) -> { // Not fatal: a batch of malformed frames on a stream with millions of // events shouldn't take the whole consumer down. wisp.log_warning( "jetstream_consumer: frame decode failed: " <> string.inspect(e), ) state } } } /// Dispatches a decoded event to the right handler and advances the replay /// cursor by its own `time_us` -- every kind does this, not just commits, so /// skipping identity/account events here would let the cursor lag behind /// them and replay more than the intended buffer on reconnect. Exposed for /// tests, same rationale as `route`: a real in-memory `catalog_index`/ /// `identity_cache`/`known_users` store makes a fine spy. pub fn apply_event(state: State, event: Event) -> State { case event { CommitEvent(frame) -> route(state, frame) |> maybe_persist_cursor(frame.time_us) IdentityEvent(did:, handle:, time_us:) -> apply_identity(state, did, handle) |> maybe_persist_cursor(time_us) AccountEvent(did:, active:, status:, time_us:) -> apply_account(state, did, active, status) |> maybe_persist_cursor(time_us) } } /// A handle change: refreshes the identity cache entry for `did` if one /// exists (never creates one -- a cold cache stays cold until something asks /// for it) and updates the stored `known_users` handle, so browse rows and /// attribution stay accurate after a rename. `handle` is optional per /// Jetstream's own schema (identity events can arrive before the handle is /// resolved); absent handle is a no-op. fn apply_identity(state: State, did: String, handle: Option(String)) -> State { case handle { None -> state Some(handle) -> { refresh_cached_identity(state.deps.identity_cache, did, handle) refresh_known_user_handle(state.deps.known_users, did, handle) state } } } fn refresh_cached_identity(cache: Cache, did: String, handle: String) -> Nil { case dict.get(cache.get_many([did]).hits, did) { Ok(identity) -> cache.put(did, Identity(..identity, handle:)) Error(Nil) -> Nil } } fn refresh_known_user_handle( known_users: known_users.Store, did: String, handle: String, ) -> Nil { case list.find(known_users.list(), fn(u) { u.did == did }) { Ok(user) -> known_users.upsert(known_users.KnownUser(..user, handle:)) Error(Nil) -> Nil } } /// An account state change. Only a genuine deletion purges the did's catalog /// rows; deactivation (reversible) and moderation states (takendown, /// suspended) are left in the index untouched, matching the ADR's /// "honour deletions, not accumulate them" posture without over-reacting to /// a temporary state. fn apply_account( state: State, did: String, active: Bool, status: Option(String), ) -> State { case active, is_deletion(status) { False, True -> { purge_did(state.deps.index, did) state.deps.shelf_index.entries.delete_for_did(did) state.deps.follow_index.edges.delete_for_did(did) state } False, False -> { wisp.log_info( "jetstream_consumer: account " <> did <> " inactive (" <> option.unwrap(status, "no status") <> "); not a deletion, leaving the index untouched", ) state } True, _ -> state } } /// Jetstream's account `status` values include "deactivated" (reversible), /// "takendown"/"suspended" (moderation, also reversible), and "deleted" /// (permanent). Only "deleted" should ever purge catalog rows. Pure so the /// distinction is unit-testable without a real store or timer. pub fn is_deletion(status: Option(String)) -> Bool { status == Some("deleted") } 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) wisp.log_info( "jetstream_consumer: purged catalog rows for deleted account " <> did, ) } /// Dispatches a decoded frame into the right store operation for its /// collection. Exposed for tests: constructs no side effects beyond the /// `Deps` capabilities it is handed, so a real in-memory `catalog_index` /// store makes a fine spy. pub fn route(state: State, frame: Frame) -> State { case frame.collection { 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 } } fn at_uri(frame: Frame) -> String { "at://" <> frame.did <> "/" <> frame.collection <> "/" <> frame.rkey } fn route_release(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { state.deps.index.releases.delete(at_uri(frame)) state } _ -> upsert_release(state, frame) } } fn upsert_release(state: State, frame: Frame) -> State { case frame.cid, frame.record { Some(cid), Some(record) -> case decode.run(record, catalog_release.catalog_release_decoder()) { Ok(value) -> { let #(publisher, next_state) = resolve_publisher(state, frame.did) let #(handle, pds) = publisher next_state.deps.index.releases.upsert(to_browse_row( frame, cid, value, handle, pds, )) next_state } Error(e) -> { wisp.log_warning( "jetstream_consumer: release decode failed for " <> at_uri(frame) <> ": " <> string.inspect(e), ) state } } _, _ -> { wisp.log_warning( "jetstream_consumer: " <> frame.operation <> " commit for " <> at_uri(frame) <> " missing cid/record", ) state } } } fn to_browse_row( frame: Frame, cid: String, value: catalog_release.CatalogRelease, handle: String, pds: String, ) -> row.BrowseRow { row.BrowseRow( uri: at_uri(frame), cid:, title: value.title, artist_display: value.artist_display, genres: value.genres |> option.unwrap([]), styles: value.styles |> option.unwrap([]), released: value.released, country: value.country, cover: value.cover, thumb_url: value.thumb_url, discogs_id: browse.discogs_id(value), created_at: value.created_at, publisher_did: frame.did, publisher_handle: handle, publisher_pds: pds, supersedes: value.supersedes |> option.map(fn(ref) { ref.uri }), based_on: value.based_on |> option.map(fn(ref) { ref.uri }), format: browse.format_display(value), label: browse.label_display(value), master: value.master |> option.map(fn(ref) { ref.uri }), barcode: browse.barcode_of(value.identifiers), ) } /// The publisher's #(handle, pds) for a did: a known user is free (no /// network hop); an unrecognized did is resolved via `resolveMiniDoc` and /// cached for the rest of this connection. Any resolution failure falls back /// to an empty handle/pds rather than blocking indexing on it. Live by design /// (E3 of the appview-first roadmap): this runs on the background firehose /// connection, never a user-facing request, so there is no read-path latency /// to save by routing it through the shared identity cache. fn resolve_publisher(state: State, did: String) -> #(#(String, String), State) { case list.find(state.deps.known_users.list(), fn(u) { u.did == did }) { Ok(user) -> #(#(user.handle, user.pds), state) Error(Nil) -> case dict.get(state.cache, did) { Ok(pair) -> #(pair, state) Error(Nil) -> { let pair = fetch_identity(state.deps, did) |> result.unwrap(#("", "")) #(pair, State(..state, cache: dict.insert(state.cache, did, pair))) } } } } fn fetch_identity(deps: Deps, did: String) -> Result(#(String, String), Nil) { generated_client.identity_resolve_mini_doc( deps.client, deps.resolver, generated_client.IdentityResolveMiniDocParams(identifier: did), None, ) |> result.map(fn(doc) { #(doc.handle, doc.pds) }) |> result.replace_error(Nil) } fn route_shelf_entry(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { state.deps.index.adoptions.delete(at_uri(frame)) delete_from_shelf_index(state, frame) state } _ -> upsert_shelf_entry(state, frame) } } fn upsert_shelf_entry(state: State, frame: Frame) -> State { case frame.record { Some(record) -> case decode.run(record, shelf_entry.shelf_entry_decoder()) { Ok(entry) -> { upsert_adoption_edge(state, frame, entry) upsert_shelf_index_event(state, frame, entry) state } Error(e) -> { wisp.log_warning( "jetstream_consumer: shelf.entry decode failed for " <> at_uri(frame) <> ": " <> string.inspect(e), ) state } } None -> state } } /// Not every shelf.entry event carries a release ref (a plain /// acquire-by-barcode append, a note, ...); only adoption edges matter to /// the catalog index. fn upsert_adoption_edge( state: State, frame: Frame, entry: shelf_entry.ShelfEntry, ) -> Nil { case entry.release { Some(ref) -> state.deps.index.adoptions.upsert(catalog_index.Adoption( entry_uri: at_uri(frame), did: frame.did, release_uri: ref.uri, status: entry.action, created_at: entry.created_at, source: catalog_index.source_label( option.then(entry.source, fn(s) { s.origin }), ), )) None -> Nil } } /// Unlike `upsert_adoption_edge`, every shelf.entry commit lands here /// (C2 of the appview-first roadmap): the shelf index needs the whole event /// log to fold, not just release-ref'd rows. fn upsert_shelf_index_event( state: State, frame: Frame, entry: shelf_entry.ShelfEntry, ) -> Nil { let entry_uri = case entry.subject { Some(ref) -> ref.uri None -> at_uri(frame) } shelf_index.record_and_fold( state.deps.shelf_index, shelf_index.event_row( event_uri: at_uri(frame), entry_uri:, did: frame.did, rkey: frame.rkey, entry:, ), ) } /// A delete frame carries no record body, so the deleted event's `entry_uri` /// (genesis vs. append) has to come from our own stored mirror rather than /// the firehose payload. fn delete_from_shelf_index(state: State, frame: Frame) -> Nil { let event_uri = at_uri(frame) case state.deps.shelf_index.events.get(event_uri) { Some(existing) -> shelf_index.delete_and_fold( state.deps.shelf_index, event_uri, existing.entry_uri, ) None -> Nil } } fn route_follow(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { state.deps.follow_index.edges.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.follow_index.edges.upsert(follow_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" -> { state.deps.index.edits.delete(at_uri(frame)) state } _ -> upsert_edit(state, frame) } } fn upsert_edit(state: State, frame: Frame) -> State { case frame.cid, frame.record { Some(cid), Some(record) -> case decode.run(record, catalog_edit.catalog_edit_decoder()) { Ok(edit) -> case edit.subject { Some(ref) -> { state.deps.index.edits.upsert(catalog_index.Edit( edit_uri: at_uri(frame), did: frame.did, subject_uri: ref.uri, entity: edit.entity, created_at: edit.created_at, cid:, fields: edit_fields_json(edit), rationale: edit.rationale, )) state } None -> state } Error(e) -> { wisp.log_warning( "jetstream_consumer: catalog.edit decode failed for " <> at_uri(frame) <> ": " <> string.inspect(e), ) state } } _, _ -> { wisp.log_warning( "jetstream_consumer: " <> frame.operation <> " commit for " <> at_uri(frame) <> " missing cid/record", ) state } } } /// The proposal's `#releaseFields` payload, JSON-encoded without a `$type` /// wrapper (E2 of the appview-first roadmap): the only variant the edit /// inbox understands today, stored as plain text so `catalog_index` stays /// free of a `catalog_edit` dependency (see its `Edit` docstring). Any other /// fields variant (future entity kinds) is dropped -- the read side already /// only surfaces `entity == "release"` edits. fn edit_fields_json(edit: catalog_edit.CatalogEdit) -> Option(String) { case edit.fields { Some(catalog_edit.CatalogEditFieldsReleaseFields(fields)) -> Some(json.to_string(catalog_edit.encode_release_fields(fields))) _ -> None } } fn maybe_persist_cursor(state: State, time_us: Int) -> State { case state.pending + 1 >= persist_every { True -> { state.deps.index.cursor.save(time_us) State(..state, pending: 0) } False -> State(..state, pending: state.pending + 1) } } /// Exposed for tests: decodes one raw Jetstream event's top-level fields /// (`did`/`time_us`/`kind` plus the kind-named payload object: `commit`, /// `identity`, or `account`), returning `None` for any other/future kind /// Jetstream might send. pub fn frame_decoder() -> decode.Decoder(Option(Event)) { use did <- decode.field("did", decode.string) use time_us <- decode.field("time_us", decode.int) use kind <- decode.field("kind", decode.string) case kind { "commit" -> decode.field( "commit", commit_decoder(did, time_us) |> decode.map(fn(frame) { Some(CommitEvent(frame)) }), decode.success, ) "identity" -> decode.field( "identity", identity_decoder(time_us) |> decode.map(Some), decode.success, ) "account" -> decode.field( "account", account_decoder(time_us) |> decode.map(Some), decode.success, ) _ -> decode.success(None) } } fn commit_decoder(did: String, time_us: Int) -> decode.Decoder(Frame) { use operation <- decode.field("operation", decode.string) use collection <- decode.field("collection", decode.string) use rkey <- decode.field("rkey", decode.string) use cid <- decode.optional_field( "cid", None, decode.string |> decode.map(Some), ) use record <- decode.optional_field( "record", None, decode.dynamic |> decode.map(Some), ) decode.success(Frame( did:, time_us:, operation:, collection:, rkey:, cid:, record:, )) } /// `{did, handle, seq, time}`; `handle` is treated as optional even though /// Jetstream normally sends it, since a missing handle is a harmless no-op /// downstream (`apply_identity`) rather than a decode failure. fn identity_decoder(time_us: Int) -> decode.Decoder(Event) { use did <- decode.field("did", decode.string) use handle <- decode.optional_field( "handle", None, decode.string |> decode.map(Some), ) decode.success(IdentityEvent(did:, handle:, time_us:)) } /// `{active, did, seq, time, status?}`; `status` is present only when /// `active` is false (deactivated/suspended/takendown/deleted). fn account_decoder(time_us: Int) -> decode.Decoder(Event) { use did <- decode.field("did", decode.string) use active <- decode.field("active", decode.bool) use status <- decode.optional_field( "status", None, decode.string |> decode.map(Some), ) decode.success(AccountEvent(did:, active:, status:, time_us:)) }