//// Frame decoding, routing into a real in-memory catalog_index (assertions //// read the store rather than a hand-rolled spy), identity resolution, and //// the pure backoff/replay-cursor helpers. import at_record_server/catalog_index import at_record_server/jetstream_consumer.{type Deps, Deps} import at_record_server/known_users import atproto/xrpc import gleam/bit_array import gleam/dict import gleam/erlang/process import gleam/http/response import gleam/int import gleam/json import gleam/list import gleam/option.{None, Some} import support const release_uri = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" /// One raw Jetstream frame; `commit.record`/`cid` are omitted for a delete. fn frame_json( did did: String, time_us time_us: Int, operation operation: String, collection collection: String, rkey rkey: String, cid cid: String, record record: String, ) -> String { "{\"did\":\"" <> did <> "\",\"time_us\":" <> int.to_string(time_us) <> ",\"kind\":\"commit\",\"commit\":{\"rev\":\"1\",\"operation\":\"" <> operation <> "\",\"collection\":\"" <> collection <> "\",\"rkey\":\"" <> rkey <> "\"" <> case operation { "delete" -> "" _ -> ",\"cid\":\"" <> cid <> "\",\"record\":" <> record } <> "}}" } fn release_frame(did: String, time_us: Int, operation: String) -> String { frame_json( did:, time_us:, operation:, collection: "dev.mokkenstorm.crate.catalog.release", rkey: "r1", cid: "bafyrelcid", record: "{\"$type\":\"dev.mokkenstorm.crate.catalog.release\",\"createdAt\":\"2024-01-01T00:00:00Z\",\"title\":\"Test Album\",\"artistDisplay\":\"Test Artist\",\"genres\":[\"Electronic\"],\"styles\":[\"Techno\"]}", ) } fn shelf_entry_frame( did: String, time_us: Int, operation: String, with_release: Bool, ) -> String { let release_field = case with_release { True -> ",\"release\":{\"cid\":\"bafyrelcid\",\"uri\":\"" <> release_uri <> "\"}" False -> "" } frame_json( did:, time_us:, operation:, collection: "dev.mokkenstorm.crate.shelf.entry", rkey: "e1", cid: "bafyentrycid", record: "{\"$type\":\"dev.mokkenstorm.crate.shelf.entry\",\"action\":\"acquired\",\"createdAt\":\"2024-01-02T00:00:00Z\"" <> release_field <> "}", ) } fn catalog_edit_frame(did: String, time_us: Int, operation: String) -> String { frame_json( did:, time_us:, operation:, collection: "dev.mokkenstorm.crate.catalog.edit", rkey: "x1", cid: "bafyeditcid", record: "{\"$type\":\"dev.mokkenstorm.crate.catalog.edit\",\"entity\":\"release\",\"op\":\"propose\",\"createdAt\":\"2024-01-03T00:00:00Z\",\"subject\":{\"cid\":\"bafyrelcid\",\"uri\":\"" <> release_uri <> "\"}}", ) } fn known_user_store() -> known_users.Store { known_users.Store(upsert: fn(_) { Nil }, list: fn() { [ known_users.KnownUser( did: "did:plc:pub", handle: "pub.test", pds: "https://pds.test", ), ] }) } fn empty_known_users() -> known_users.Store { known_users.Store(upsert: fn(_) { Nil }, list: fn() { [] }) } fn resolve_body() -> String { "{\"did\":\"did:plc:stranger\",\"handle\":\"stranger.test\",\"pds\":\"https://stranger-pds.test\",\"signing_key\":\"zTest\"}" } fn resolving_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case req.host { "resolver.test" -> Ok(response.Response(200, [], bit_array.from_string(resolve_body()))) _ -> panic as "unexpected host" } }) } /// Counts calls via a subject since the client closure must stay a pure fn. fn counting_resolver_client(calls: process.Subject(Nil)) -> xrpc.Client { xrpc.Client(send: fn(req) { case req.host { "resolver.test" -> { process.send(calls, Nil) Ok(response.Response(200, [], bit_array.from_string(resolve_body()))) } _ -> panic as "unexpected host" } }) } fn count_messages(subject: process.Subject(Nil)) -> Int { case process.receive(subject, 10) { Ok(Nil) -> 1 + count_messages(subject) Error(Nil) -> 0 } } fn deps( known: known_users.Store, index: catalog_index.Store, client: xrpc.Client, ) -> Deps { Deps(known_users: known, index:, client:, resolver: "https://resolver.test") } fn state_over( index: catalog_index.Store, client: xrpc.Client, ) -> jetstream_consumer.State { jetstream_consumer.initial_state(deps(empty_known_users(), index, client)) } fn decode_frame(text: String) -> jetstream_consumer.Frame { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) frame } pub fn frame_decoder_reads_did_and_time_us_test() { let frame = decode_frame(release_frame("did:plc:pub", 1_700_000_000_000_000, "create")) assert frame.did == "did:plc:pub" assert frame.time_us == 1_700_000_000_000_000 } pub fn non_commit_kind_decodes_to_none_test() { let text = "{\"did\":\"did:plc:a\",\"time_us\":1,\"kind\":\"identity\"}" assert json.parse(text, jetstream_consumer.frame_decoder()) == Ok(None) } /// `cid`/`record` are present on create, absent on delete, per collection. pub fn frame_decoder_decodes_every_collection_and_operation_test() { let release = "dev.mokkenstorm.crate.catalog.release" let shelf_entry = "dev.mokkenstorm.crate.shelf.entry" let catalog_edit = "dev.mokkenstorm.crate.catalog.edit" let cases = [ #(release_frame("did:plc:pub", 1, "create"), release, "r1", True), #(release_frame("did:plc:pub", 2, "delete"), release, "r1", False), #( shelf_entry_frame("did:plc:a", 3, "create", True), shelf_entry, "e1", True, ), #( shelf_entry_frame("did:plc:a", 4, "delete", True), shelf_entry, "e1", False, ), #(catalog_edit_frame("did:plc:a", 5, "create"), catalog_edit, "x1", True), #(catalog_edit_frame("did:plc:a", 6, "delete"), catalog_edit, "x1", False), ] cases |> list.each(fn(c) { let #(text, collection, rkey, has_record) = c let frame = decode_frame(text) assert frame.collection == collection assert frame.rkey == rkey case has_record { True -> { assert frame.cid != None assert frame.record != None } False -> { assert frame.cid == None assert frame.record == None } } }) } pub fn release_commit_upserts_then_delete_removes_from_the_index_test() { let assert Ok(index) = catalog_index.start() let state = jetstream_consumer.initial_state(deps( known_user_store(), index, support.unreachable_client(), )) let state = jetstream_consumer.route( state, decode_frame(release_frame("did:plc:pub", 1, "create")), ) let assert [row] = index.releases.list() assert row.uri == release_uri assert row.title == "Test Album" assert row.publisher_handle == "pub.test" assert row.publisher_pds == "https://pds.test" let _ = jetstream_consumer.route( state, decode_frame(release_frame("did:plc:pub", 2, "delete")), ) assert index.releases.list() == [] } pub fn shelf_entry_with_release_ref_creates_then_delete_removes_the_adoption_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(shelf_entry_frame("did:plc:a", 1, "create", True)), ) assert index.adoptions.count(release_uri) == 1 let _ = jetstream_consumer.route( state, decode_frame(shelf_entry_frame("did:plc:a", 2, "delete", True)), ) assert index.adoptions.count(release_uri) == 0 } pub fn shelf_entry_without_release_ref_is_ignored_test() { let assert Ok(index) = catalog_index.start() let state = state_over(index, support.unreachable_client()) let frame = decode_frame(shelf_entry_frame("did:plc:a", 1, "create", False)) let _ = jetstream_consumer.route(state, frame) assert index.adoptions.counts() == dict.new() } pub fn catalog_edit_indexes_subject_and_entity_then_delete_removes_it_test() { let assert Ok(index) = catalog_index.start() let state = state_over(index, support.unreachable_client()) let state = jetstream_consumer.route( state, decode_frame(catalog_edit_frame("did:plc:a", 1, "create")), ) let assert [edit] = index.edits.for_subject(release_uri) assert edit.did == "did:plc:a" assert edit.entity == "release" let _ = jetstream_consumer.route( state, decode_frame(catalog_edit_frame("did:plc:a", 2, "delete")), ) assert index.edits.for_subject(release_uri) == [] } pub fn unknown_did_resolves_via_resolve_mini_doc_test() { let assert Ok(index) = catalog_index.start() let state = state_over(index, resolving_client()) let frame = decode_frame(release_frame("did:plc:stranger", 1, "create")) let _ = jetstream_consumer.route(state, frame) let assert [row] = index.releases.list() assert row.publisher_handle == "stranger.test" assert row.publisher_pds == "https://stranger-pds.test" } pub fn unknown_did_resolution_is_cached_within_a_connection_test() { let assert Ok(index) = catalog_index.start() let calls = process.new_subject() let state = state_over(index, counting_resolver_client(calls)) let state = jetstream_consumer.route( state, decode_frame(release_frame("did:plc:stranger", 1, "create")), ) let _ = jetstream_consumer.route( state, decode_frame(release_frame("did:plc:stranger", 2, "create")), ) // Two events, same unknown did: the second resolution hits the cache. assert count_messages(calls) == 1 } pub fn next_backoff_doubles_and_caps_test() { assert jetstream_consumer.next_backoff(1000) == 2000 assert jetstream_consumer.next_backoff(2000) == 4000 assert jetstream_consumer.next_backoff(40_000) == 60_000 assert jetstream_consumer.next_backoff(60_000) == 60_000 } pub fn replay_from_rewinds_by_the_buffer_test() { assert jetstream_consumer.replay_from(10_000_000) == 8_000_000 } pub fn replay_from_never_goes_negative_test() { assert jetstream_consumer.replay_from(500) == 0 }