diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index 712b7ea..01227ba 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -39,6 +39,7 @@ pub fn main() -> Nil { config_env.jetstream_url(), known_users, catalog_index, + shelf_index, client, resolver, identity_cache, diff --git a/server/src/at_record_server/jetstream_consumer.gleam b/server/src/at_record_server/jetstream_consumer.gleam index b4b55bb..6e6f758 100644 --- a/server/src/at_record_server/jetstream_consumer.gleam +++ b/server/src/at_record_server/jetstream_consumer.gleam @@ -18,6 +18,7 @@ import at_record_server/catalog/row import at_record_server/catalog_index import at_record_server/identity_cache.{type Cache, Identity} import at_record_server/known_users +import at_record_server/shelf_index import atproto/xrpc.{type Client} import gleam/dict.{type Dict} import gleam/dynamic/decode.{type Dynamic} @@ -61,6 +62,9 @@ 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, client: Client, resolver: String, /// Threaded from `main` (the same `identity_cache.Cache` passed into @@ -115,11 +119,20 @@ pub fn start( host: String, known_users: known_users.Store, index: catalog_index.Store, + shelf_index: shelf_index.Store, client: Client, resolver: String, identity_cache: Cache, ) -> Nil { - let deps = Deps(known_users:, index:, client:, resolver:, identity_cache:) + let deps = + Deps( + known_users:, + index:, + shelf_index:, + client:, + resolver:, + identity_cache:, + ) process.spawn_unlinked(fn() { supervise(host, deps, base_backoff_ms) }) Nil } @@ -296,6 +309,7 @@ fn apply_account( case active, is_deletion(status) { False, True -> { purge_did(state.deps.index, did) + state.deps.shelf_index.entries.delete_for_did(did) state } False, False -> { @@ -459,36 +473,22 @@ 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_adoption(state, frame) + _ -> upsert_shelf_entry(state, frame) } } -fn upsert_adoption(state: State, frame: Frame) -> State { +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) -> - 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 }), - ), - )) - 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. - None -> state - } + 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 " @@ -503,6 +503,70 @@ fn upsert_adoption(state: State, frame: Frame) -> 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_edit(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { diff --git a/server/test/jetstream_consumer_test.gleam b/server/test/jetstream_consumer_test.gleam index 50ffc58..f8970b7 100644 --- a/server/test/jetstream_consumer_test.gleam +++ b/server/test/jetstream_consumer_test.gleam @@ -9,6 +9,7 @@ import at_record_server/identity_cache.{Identity} import at_record_server/jetstream_consumer.{type Deps, Deps} import at_record_server/known_users import at_record_server/known_users_memory +import at_record_server/shelf_index import atproto/xrpc import gleam/bit_array import gleam/dict @@ -86,6 +87,59 @@ fn shelf_entry_frame( ) } +/// A genesis shelf.entry commit (no `subject`, no release ref): the minimal +/// event that starts a folded entry. +fn shelf_genesis_frame(did: String, time_us: Int, rkey: String) -> String { + frame_json( + did:, + time_us:, + operation: "create", + collection: "dev.mokkenstorm.crate.shelf.entry", + rkey:, + cid: "bafygenesis" <> rkey, + record: "{\"$type\":\"dev.mokkenstorm.crate.shelf.entry\",\"action\":\"acquired\",\"createdAt\":\"2024-01-01T00:00:00Z\"}", + ) +} + +/// An append shelf.entry commit: carries `subject` pointing back at the +/// genesis, no release ref. +fn shelf_append_frame( + did: String, + time_us: Int, + rkey: String, + genesis_uri: String, + action: String, + created_at: String, +) -> String { + frame_json( + did:, + time_us:, + operation: "create", + collection: "dev.mokkenstorm.crate.shelf.entry", + rkey:, + cid: "bafyappend" <> rkey, + record: "{\"$type\":\"dev.mokkenstorm.crate.shelf.entry\",\"action\":\"" + <> action + <> "\",\"createdAt\":\"" + <> created_at + <> "\",\"subject\":{\"cid\":\"bafygenesis\",\"uri\":\"" + <> genesis_uri + <> "\"}}", + ) +} + +fn shelf_entry_delete_frame(did: String, time_us: Int, rkey: String) -> String { + frame_json( + did:, + time_us:, + operation: "delete", + collection: "dev.mokkenstorm.crate.shelf.entry", + rkey:, + cid: "unused", + record: "unused", + ) +} + fn catalog_edit_frame(did: String, time_us: Int, operation: String) -> String { frame_json( did:, @@ -156,6 +210,11 @@ fn a_cache() -> identity_cache.Cache { cache } +fn fresh_shelf_index() -> shelf_index.Store { + let assert Ok(store) = shelf_index.start() + store +} + fn deps( known: known_users.Store, index: catalog_index.Store, @@ -169,10 +228,21 @@ fn deps_with( index: catalog_index.Store, client: xrpc.Client, identity_cache: identity_cache.Cache, +) -> Deps { + deps_full(known, index, fresh_shelf_index(), client, identity_cache) +} + +fn deps_full( + known: known_users.Store, + index: catalog_index.Store, + shelf_index: shelf_index.Store, + client: xrpc.Client, + identity_cache: identity_cache.Cache, ) -> Deps { Deps( known_users: known, index:, + shelf_index:, client:, resolver: "https://resolver.test", identity_cache:, @@ -186,6 +256,24 @@ fn state_over( jetstream_consumer.initial_state(deps(empty_known_users(), index, client)) } +/// `state_over`, but exposing the shelf index the state was built with, so +/// tests can assert on it directly. +fn state_with_shelf_index( + client: xrpc.Client, +) -> #(jetstream_consumer.State, catalog_index.Store, shelf_index.Store) { + let assert Ok(index) = catalog_index.start() + let shelf = fresh_shelf_index() + let state = + jetstream_consumer.initial_state(deps_full( + empty_known_users(), + index, + shelf, + client, + a_cache(), + )) + #(state, index, shelf) +} + fn decode_event(text: String) -> jetstream_consumer.Event { let assert Ok(Some(event)) = json.parse(text, jetstream_consumer.frame_decoder()) @@ -418,11 +506,39 @@ fn seed_catalog_rows_for(index: catalog_index.Store, did: String) -> Nil { )) } +fn shelf_entry_uri_for(did: String) -> String { + "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/e1" +} + +fn seed_shelf_index_for(shelf: shelf_index.Store, did: String) -> Nil { + let entry_uri = shelf_entry_uri_for(did) + shelf_index.record_and_fold( + shelf, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did:, + rkey: "e1", + entry: support.blank_shelf_entry(), + ), + ) +} + pub fn account_deletion_purges_that_dids_catalog_rows_test() { let assert Ok(index) = catalog_index.start() + let shelf = fresh_shelf_index() seed_catalog_rows_for(index, "did:plc:gone") seed_catalog_rows_for(index, "did:plc:stays") - let state = state_over(index, support.unreachable_client()) + seed_shelf_index_for(shelf, "did:plc:gone") + seed_shelf_index_for(shelf, "did:plc:stays") + let state = + jetstream_consumer.initial_state(deps_full( + empty_known_users(), + index, + shelf, + support.unreachable_client(), + a_cache(), + )) let _ = jetstream_consumer.apply_event( state, @@ -435,6 +551,8 @@ pub fn account_deletion_purges_that_dids_catalog_rows_test() { |> list.any(fn(e) { e.did == "did:plc:stays" }) assert index.edits.for_subject(release_uri) |> list.all(fn(e) { e.did != "did:plc:gone" }) + assert shelf.entries.get(shelf_entry_uri_for("did:plc:gone")) == None + assert shelf.entries.get(shelf_entry_uri_for("did:plc:stays")) != None } pub fn account_deactivation_does_not_purge_catalog_rows_test() { @@ -599,6 +717,71 @@ pub fn shelf_entry_without_release_ref_is_ignored_test() { assert index.adoptions.counts() == dict.new() } +/// C2: a shelf.entry commit with no release ref still lands in the shelf +/// index (unlike catalog_adoptions, which only cares about release-ref'd +/// rows). +pub fn shelf_entry_without_release_ref_lands_in_shelf_index_test() { + let #(state, index, shelf) = + state_with_shelf_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() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let assert Some(folded) = shelf.entries.get(entry_uri) + assert folded.status == "owned" + assert folded.did == "did:plc:a" +} + +/// C2: a later append (a sold event, `subject` pointing back at the +/// genesis) re-folds the entry rather than being ignored. +pub fn sold_append_flips_the_folded_status_test() { + let #(state, _index, shelf) = + state_with_shelf_index(support.unreachable_client()) + let genesis_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let state = + jetstream_consumer.route( + state, + decode_frame(shelf_genesis_frame("did:plc:a", 1, "e1")), + ) + let assert Some(before) = shelf.entries.get(genesis_uri) + assert before.status == "owned" + let _ = + jetstream_consumer.route( + state, + decode_frame(shelf_append_frame( + "did:plc:a", + 2, + "e2", + genesis_uri, + "sold", + "2024-02-01T00:00:00Z", + )), + ) + let assert Some(after) = shelf.entries.get(genesis_uri) + assert after.status == "gone" +} + +/// C2: deleting the last remaining event for an entry removes both the raw +/// event mirror and the folded row. +pub fn shelf_entry_delete_removes_the_event_and_folded_row_when_last_test() { + let #(state, _index, shelf) = + state_with_shelf_index(support.unreachable_client()) + let genesis_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let state = + jetstream_consumer.route( + state, + decode_frame(shelf_genesis_frame("did:plc:a", 1, "e1")), + ) + assert shelf.entries.get(genesis_uri) != None + let _ = + jetstream_consumer.route( + state, + decode_frame(shelf_entry_delete_frame("did:plc:a", 2, "e1")), + ) + assert shelf.entries.get(genesis_uri) == None + assert shelf.events.for_entry(genesis_uri) == [] +} + 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())