diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index f7ea2e5..712b7ea 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -27,7 +27,7 @@ pub fn main() -> Nil { let base_url = config_env.base_url(port) let client = atproto_client.client() let assert Ok(store) = store.start() - let #(sessions, discogs_creds, known_users, catalog_index) = + let #(sessions, discogs_creds, known_users, catalog_index, shelf_index) = config_env.stores() let assert Ok(identity_cache) = identity_cache.start( @@ -72,6 +72,7 @@ pub fn main() -> Nil { discogs_creds:, known_users:, catalog_index:, + shelf_index:, identity_cache:, oauth:, ) diff --git a/server/src/at_record_server/config_env.gleam b/server/src/at_record_server/config_env.gleam index f31410d..0fd13dc 100644 --- a/server/src/at_record_server/config_env.gleam +++ b/server/src/at_record_server/config_env.gleam @@ -13,6 +13,8 @@ import at_record_server/oauth/sessions.{type Store} import at_record_server/oauth/sessions_memory import at_record_server/oauth/sessions_postgres import at_record_server/sealed_store +import at_record_server/shelf_index +import at_record_server/shelf_index_postgres import atproto/constellation import atproto/identity import envoy @@ -113,17 +115,25 @@ pub fn cover_cache_directory() -> String { } /// The session store, the Discogs credential store, the known-users store, -/// and the catalog index: one shared Postgres pool, no sweep ttl on the -/// durable Discogs tokens. The two credential stores are encrypted at rest; -/// known_users and the catalog index are public, so neither is sealed. -pub fn stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { +/// the catalog index, and the shelf index: one shared Postgres pool, no +/// sweep ttl on the durable Discogs tokens. The two credential stores are +/// encrypted at rest; known_users and both indexes are public, so none of +/// them are sealed. +pub fn stores() -> #( + Store, + Store, + known_users.Store, + catalog_index.Store, + shelf_index.Store, +) { let key = store_key() - let #(session_store, discogs_store, known, index) = raw_stores() + let #(session_store, discogs_store, known, index, shelf) = raw_stores() #( sealed_store.wrap(session_store, key), sealed_store.wrap(discogs_store, key), known, index, + shelf, ) } @@ -140,11 +150,17 @@ fn store_key() -> gose.Key(String) { } } -fn raw_stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { +fn raw_stores() -> #( + Store, + Store, + known_users.Store, + catalog_index.Store, + shelf_index.Store, +) { case envoy.get("DATABASE_URL") { Ok(url) -> case postgres_stores(url) { - Ok(quad) -> quad + Ok(quintet) -> quintet Error(e) -> { wisp.log_error( "Postgres stores unavailable (" <> e <> "); using in-memory", @@ -163,7 +179,10 @@ fn raw_stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { fn postgres_stores( url: String, -) -> Result(#(Store, Store, known_users.Store, catalog_index.Store), String) { +) -> Result( + #(Store, Store, known_users.Store, catalog_index.Store, shelf_index.Store), + String, +) { use conn <- result.try(sessions_postgres.connect_pool(url)) use session_store <- result.try(sessions_postgres.table_store( conn, @@ -176,16 +195,24 @@ fn postgres_stores( option.None, )) use known <- result.try(known_users_postgres.table_store(conn)) - use index <- result.map(catalog_index_postgres.table_store(conn)) - #(session_store, discogs_store, known, index) -} - -fn memory_stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { + use index <- result.try(catalog_index_postgres.table_store(conn)) + use shelf <- result.map(shelf_index_postgres.table_store(conn)) + #(session_store, discogs_store, known, index, shelf) +} + +fn memory_stores() -> #( + Store, + Store, + known_users.Store, + catalog_index.Store, + shelf_index.Store, +) { let assert Ok(session_store) = sessions_memory.start() let assert Ok(discogs_store) = sessions_memory.start() let assert Ok(known) = known_users_memory.start() let assert Ok(index) = catalog_index.start() - #(session_store, discogs_store, known, index) + let assert Ok(shelf) = shelf_index.start() + #(session_store, discogs_store, known, index, shelf) } /// App-level Discogs auth from env. Absent is fine: search still works, just diff --git a/server/src/at_record_server/context.gleam b/server/src/at_record_server/context.gleam index 4f99035..5af5d1b 100644 --- a/server/src/at_record_server/context.gleam +++ b/server/src/at_record_server/context.gleam @@ -13,6 +13,7 @@ import at_record_server/oauth/authed import at_record_server/oauth/config.{type Config} import at_record_server/oauth/session_store import at_record_server/oauth/sessions.{type OauthSession, type Store} +import at_record_server/shelf_index import atproto/xrpc.{type Client} import gleam/json import gleam/option.{type Option, None, Some} @@ -51,6 +52,9 @@ pub type Context { catalog: catalog_deps.Deps, known_users: known_users.Store, catalog_index: catalog_index.Store, + /// Wired dark (C1+C2 of the appview-first roadmap): no read path + /// consumes this yet. See `shelf_index`. + shelf_index: shelf_index.Store, /// The resolution layer's read port: candidate release rows (chain edges /// included) plus a public adoption tally. See `catalog/source`. variant_source: catalog_source.Source, diff --git a/server/src/at_record_server/shelf_index.gleam b/server/src/at_record_server/shelf_index.gleam new file mode 100644 index 0000000..abadcae --- /dev/null +++ b/server/src/at_record_server/shelf_index.gleam @@ -0,0 +1,411 @@ +//// The shelf index port: durable per-entry event mirrors plus the folded +//// crate view they derive, so a crate and an entry timeline can be read +//// from Postgres instead of the owner's live PDS. Backends live in this +//// module (in-memory, dev fallback, actor over a Dict) and in +//// shelf_index_postgres (durable); wiring picks one at the composition +//// root. `Store` groups its operations by the entity they act on, mirroring +//// `catalog_index`'s shape: `store.entries.get(uri)`, `store.events.upsert(row)`. +//// +//// Wired dark (C1+C2 of the appview-first roadmap): nothing reads this store +//// yet. `record_and_fold`/`delete_and_fold` are the shared seam the +//// jetstream consumer (C2) and the later write-through path (C3) both call. + +import at_record/gen/shelf/entry as shelf_entry +import at_record/storage.{type StoredItem, StoredItem} +import at_record_server/crate.{type CrateEntry} +import at_record_server/parallel +import at_record_server/provenance +import gleam/dict.{type Dict} +import gleam/erlang/process.{type Subject} +import gleam/json +import gleam/list +import gleam/option.{type Option, None} +import gleam/otp/actor +import gleam/result +import gleam/set.{type Set} +import gleam/string + +/// One raw shelf.entry commit, mirrored verbatim so re-folding never touches +/// the network. `entry_uri` is the owning entry's identity: the genesis +/// record's own at-uri (equal to `event_uri` for the genesis itself, or the +/// append's `subject.uri` otherwise). `record_json` is the full encoded +/// `ShelfEntry`, needed (not just the folded fields) so a re-fold is lossless +/// even for actions this index doesn't otherwise project. +pub type ShelfEntryEvent { + ShelfEntryEvent( + event_uri: String, + entry_uri: String, + did: String, + rkey: String, + action: String, + created_at: String, + record_json: String, + indexed_at: String, + ) +} + +/// The folded projection of one crate entry: `crate.CrateEntry` flattened to +/// scalar columns (no JSON blobs), matching `catalog_releases`' style. Every +/// field beyond `entry_uri`/`did`/`status`/`created_at`/`updated_at` is +/// last-write-wins per `crate.fold`, so it mirrors `CrateEntry` field for +/// field: snapshot flattens to `title`/`artist_display`/`year`/`format`/ +/// `thumb_url`/`cover_cid`, `release` to `release_uri`/`release_cid`, `price` +/// to `price_amount`/`price_currency`, and `source` to its five components. +pub type FoldedEntry { + FoldedEntry( + entry_uri: String, + did: String, + status: String, + title: Option(String), + artist_display: Option(String), + year: Option(Int), + format: Option(String), + thumb_url: Option(String), + cover_cid: Option(String), + media_grade: Option(String), + sleeve_grade: Option(String), + rating: Option(Int), + folder: Option(String), + notes: Option(String), + release_uri: Option(String), + release_cid: Option(String), + price_amount: Option(Int), + price_currency: Option(String), + counterparty: Option(String), + source_client_agent: Option(String), + source_origin: Option(String), + source_origin_url: Option(String), + source_external_provider: Option(String), + source_external_id: Option(String), + source_external_url: Option(String), + source_record_uri: Option(String), + source_record_cid: Option(String), + created_at: String, + updated_at: String, + indexed_at: String, + ) +} + +pub type EventOps { + EventOps( + upsert: fn(ShelfEntryEvent) -> Nil, + delete: fn(String) -> Nil, + /// One event by its own at-uri; used to recover `entry_uri` on a + /// firehose delete frame, which carries no record body to read it from. + get: fn(String) -> Option(ShelfEntryEvent), + /// Every event belonging to one logical entry (genesis + appends), + /// ordered by `rkey`. + for_entry: fn(String) -> List(ShelfEntryEvent), + ) +} + +pub type EntryOps { + EntryOps( + get: fn(String) -> Option(FoldedEntry), + list_for_did: fn(String) -> List(FoldedEntry), + upsert: fn(FoldedEntry) -> Nil, + delete: fn(String) -> Nil, + /// Drops every folded entry (and its events) owned by `did`, honoring an + /// account deletion; parity with `catalog_index.AdoptionOps.delete_for_did`. + delete_for_did: fn(String) -> Nil, + ) +} + +pub type SeenOps { + SeenOps( + /// Whether this did's crate has ever been indexed, so an empty crate + /// reads as genuinely empty rather than never-seen. + has_seen: fn(String) -> Bool, + mark_seen: fn(String) -> Nil, + ) +} + +pub type Store { + Store(events: EventOps, entries: EntryOps, seen: SeenOps) +} + +/// Build one `ShelfEntryEvent` row from a decoded `ShelfEntry`, stamping +/// `indexed_at` at construction time so every caller (the jetstream +/// consumer, and later the write-through handlers) gets it for free instead +/// of each re-deriving "now". +pub fn event_row( + event_uri event_uri: String, + entry_uri entry_uri: String, + did did: String, + rkey rkey: String, + entry entry: shelf_entry.ShelfEntry, +) -> ShelfEntryEvent { + ShelfEntryEvent( + event_uri:, + entry_uri:, + did:, + rkey:, + action: entry.action, + created_at: entry.created_at, + record_json: json.to_string(shelf_entry.encode_shelf_entry(entry)), + indexed_at: provenance.now_rfc3339(), + ) +} + +/// Store one event, then re-derive the folded row for its entry from every +/// stored event -- the single write path the jetstream consumer (C2) and the +/// write-through handlers (C3) both call. +pub fn record_and_fold(store: Store, event: ShelfEntryEvent) -> Nil { + store.events.upsert(event) + refold(store, event.entry_uri) +} + +/// Remove one event, then re-derive `entry_uri`'s folded row. A genesis +/// deletion leaves no events with an absent `subject`, so `refold` deletes +/// the folded row rather than upserting one; a mid-history append deletion +/// self-heals the fold from the remaining events. +pub fn delete_and_fold( + store: Store, + event_uri: String, + entry_uri: String, +) -> Nil { + store.events.delete(event_uri) + refold(store, entry_uri) +} + +/// Re-derive `entry_uri`'s folded row from every stored event: decode each +/// `record_json` back into a `ShelfEntry`, run the existing `crate.fold` over +/// the group, and upsert the result -- or delete the folded row when no +/// genesis remains (out-of-order arrival, or the genesis itself was purged). +/// Idempotent by construction, so replays and reconnects are always safe. +pub fn refold(store: Store, entry_uri: String) -> Nil { + let events = store.events.for_entry(entry_uri) + let stored = events |> list.filter_map(decode_event) + case crate.fold(stored) { + [] -> store.entries.delete(entry_uri) + [folded, ..] -> { + let did = + events |> list.first |> result.map(fn(e) { e.did }) |> result.unwrap("") + store.entries.upsert(to_folded_entry(entry_uri, did, folded)) + } + } +} + +fn decode_event( + event: ShelfEntryEvent, +) -> Result(StoredItem(shelf_entry.ShelfEntry), Nil) { + json.parse(event.record_json, shelf_entry.shelf_entry_decoder()) + |> result.map(fn(value) { + StoredItem(uri: event.event_uri, cid: "", rkey: event.rkey, value:) + }) + |> result.replace_error(Nil) +} + +fn to_folded_entry( + entry_uri: String, + did: String, + entry: CrateEntry, +) -> FoldedEntry { + let snapshot = entry.snapshot + let source = entry.source + let external = option.then(source, fn(s) { s.external }) + FoldedEntry( + entry_uri:, + did:, + status: crate.status_string(entry.status), + title: option.map(snapshot, fn(s) { s.title }), + artist_display: option.map(snapshot, fn(s) { s.artist_display }), + year: option.then(snapshot, fn(s) { s.year }), + format: option.then(snapshot, fn(s) { s.format }), + thumb_url: option.then(snapshot, fn(s) { s.thumb_url }), + cover_cid: option.then(snapshot, fn(s) { s.cover }) + |> option.map(fn(b) { b.cid }), + media_grade: entry.media_grade, + sleeve_grade: entry.sleeve_grade, + rating: entry.rating, + folder: entry.folder, + notes: entry.notes, + release_uri: option.map(entry.release, fn(r) { r.uri }), + release_cid: option.map(entry.release, fn(r) { r.cid }), + price_amount: option.map(entry.price, fn(p) { p.amount }), + price_currency: option.map(entry.price, fn(p) { p.currency }), + counterparty: entry.counterparty, + source_client_agent: option.then(source, fn(s) { s.client_agent }), + source_origin: option.then(source, fn(s) { s.origin }), + source_origin_url: option.then(source, fn(s) { s.origin_url }), + source_external_provider: option.map(external, fn(e) { e.provider }), + source_external_id: option.map(external, fn(e) { e.id }), + source_external_url: option.then(external, fn(e) { e.url }), + source_record_uri: option.map( + option.then(source, fn(s) { s.record }), + fn(r) { r.uri }, + ), + source_record_cid: option.map( + option.then(source, fn(s) { s.record }), + fn(r) { r.cid }, + ), + created_at: entry.created_at, + updated_at: entry.updated_at, + indexed_at: provenance.now_rfc3339(), + ) +} + +type State { + State( + events: Dict(String, ShelfEntryEvent), + entries: Dict(String, FoldedEntry), + seen: Set(String), + ) +} + +type Msg { + UpsertEvent(ShelfEntryEvent) + DeleteEvent(String) + GetEvent(String, Subject(Option(ShelfEntryEvent))) + EventsForEntry(String, Subject(List(ShelfEntryEvent))) + GetEntry(String, Subject(Option(FoldedEntry))) + ListForDid(String, Subject(List(FoldedEntry))) + UpsertEntry(FoldedEntry) + DeleteEntry(String) + DeleteForDid(String) + HasSeen(String, Subject(Bool)) + MarkSeen(String) +} + +fn initial_state() -> State { + State(events: dict.new(), entries: dict.new(), seen: set.new()) +} + +pub fn start() -> Result(Store, actor.StartError) { + use started <- result.map( + actor.new(initial_state()) |> actor.on_message(handle) |> actor.start, + ) + let subject = started.data + Store( + events: EventOps( + upsert: fn(event) { + process.send(subject, UpsertEvent(event)) + Nil + }, + delete: fn(event_uri) { + process.send(subject, DeleteEvent(event_uri)) + Nil + }, + get: fn(event_uri) { + case parallel.try_call(subject, 1000, GetEvent(event_uri, _)) { + Ok(event) -> event + Error(Nil) -> None + } + }, + for_entry: fn(entry_uri) { + case parallel.try_call(subject, 1000, EventsForEntry(entry_uri, _)) { + Ok(events) -> events + Error(Nil) -> [] + } + }, + ), + entries: EntryOps( + get: fn(entry_uri) { + case parallel.try_call(subject, 1000, GetEntry(entry_uri, _)) { + Ok(entry) -> entry + Error(Nil) -> None + } + }, + list_for_did: fn(did) { + case parallel.try_call(subject, 1000, ListForDid(did, _)) { + Ok(entries) -> entries + Error(Nil) -> [] + } + }, + upsert: fn(entry) { + process.send(subject, UpsertEntry(entry)) + Nil + }, + delete: fn(entry_uri) { + process.send(subject, DeleteEntry(entry_uri)) + Nil + }, + delete_for_did: fn(did) { + process.send(subject, DeleteForDid(did)) + Nil + }, + ), + seen: SeenOps( + has_seen: fn(did) { + case parallel.try_call(subject, 1000, HasSeen(did, _)) { + Ok(seen) -> seen + Error(Nil) -> False + } + }, + mark_seen: fn(did) { + process.send(subject, MarkSeen(did)) + Nil + }, + ), + ) +} + +fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { + case msg { + UpsertEvent(event) -> + actor.continue( + State( + ..state, + events: dict.insert(state.events, event.event_uri, event), + ), + ) + DeleteEvent(event_uri) -> + actor.continue( + State(..state, events: dict.delete(state.events, event_uri)), + ) + GetEvent(event_uri, reply) -> { + process.send( + reply, + dict.get(state.events, event_uri) |> option.from_result, + ) + actor.continue(state) + } + EventsForEntry(entry_uri, reply) -> { + process.send(reply, events_for(state, entry_uri)) + actor.continue(state) + } + GetEntry(entry_uri, reply) -> { + process.send( + reply, + dict.get(state.entries, entry_uri) |> option.from_result, + ) + actor.continue(state) + } + ListForDid(did, reply) -> { + process.send( + reply, + dict.values(state.entries) |> list.filter(fn(e) { e.did == did }), + ) + actor.continue(state) + } + UpsertEntry(entry) -> + actor.continue( + State( + ..state, + entries: dict.insert(state.entries, entry.entry_uri, entry), + ), + ) + DeleteEntry(entry_uri) -> + actor.continue( + State(..state, entries: dict.delete(state.entries, entry_uri)), + ) + DeleteForDid(did) -> + actor.continue(State( + events: dict.filter(state.events, fn(_uri, e) { e.did != did }), + entries: dict.filter(state.entries, fn(_uri, e) { e.did != did }), + seen: set.delete(state.seen, did), + )) + HasSeen(did, reply) -> { + process.send(reply, set.contains(state.seen, did)) + actor.continue(state) + } + MarkSeen(did) -> + actor.continue(State(..state, seen: set.insert(state.seen, did))) + } +} + +fn events_for(state: State, entry_uri: String) -> List(ShelfEntryEvent) { + dict.values(state.events) + |> list.filter(fn(e) { e.entry_uri == entry_uri }) + |> list.sort(fn(a, b) { string.compare(a.rkey, b.rkey) }) +} diff --git a/server/src/at_record_server/shelf_index_postgres.gleam b/server/src/at_record_server/shelf_index_postgres.gleam new file mode 100644 index 0000000..5e0a6c0 --- /dev/null +++ b/server/src/at_record_server/shelf_index_postgres.gleam @@ -0,0 +1,495 @@ +//// Postgres shelf_index.Store backend (pog), over three create-if-not-exists +//// tables sharing one pool with the other stores: shelf_entry_events (the +//// raw event mirror, so re-folding never touches the network), shelf_entries +//// (the folded crate view, derived by re-running crate.fold over one +//// entry's rows per commit), and shelf_index_seen (a once-per-did marker so +//// an empty crate is distinguishable from a never-indexed one). This is +//// derived, rebuildable data, so writes swallow-and-log rather than fail +//// the caller. + +import at_record_server/provenance +import at_record_server/shelf_index.{ + type FoldedEntry, type ShelfEntryEvent, EntryOps, EventOps, FoldedEntry, + SeenOps, ShelfEntryEvent, Store, +} +import gleam/dynamic/decode +import gleam/option.{type Option, None, Some} +import gleam/result +import pog +import wisp + +pub fn table_store(conn: pog.Connection) -> Result(shelf_index.Store, String) { + use _ <- result.try(migrate(conn)) + Ok(Store( + events: EventOps( + upsert: fn(event) { upsert_event(conn, event) }, + delete: fn(event_uri) { delete_event(conn, event_uri) }, + get: fn(event_uri) { get_event(conn, event_uri) }, + for_entry: fn(entry_uri) { events_for(conn, entry_uri) }, + ), + entries: EntryOps( + get: fn(entry_uri) { get_entry(conn, entry_uri) }, + list_for_did: fn(did) { list_for_did(conn, did) }, + upsert: fn(entry) { upsert_entry(conn, entry) }, + delete: fn(entry_uri) { delete_entry(conn, entry_uri) }, + delete_for_did: fn(did) { delete_for_did(conn, did) }, + ), + seen: SeenOps(has_seen: fn(did) { has_seen(conn, did) }, mark_seen: fn(did) { + mark_seen(conn, did) + }), + )) +} + +fn migrate(conn: pog.Connection) -> Result(Nil, String) { + use _ <- result.try( + pog.query( + "create table if not exists shelf_entry_events ( + event_uri text primary key, + entry_uri text not null, + did text not null, + rkey text not null, + action text not null, + created_at text not null, + record_json text not null, + indexed_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_entry_events migration failed"), + ) + use _ <- result.try( + pog.query( + "create index if not exists shelf_entry_events_entry_rkey_idx + on shelf_entry_events (entry_uri, rkey)", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_entry_events index migration failed"), + ) + use _ <- result.try( + pog.query( + "create table if not exists shelf_entries ( + entry_uri text primary key, + did text not null, + status text not null, + title text, + artist_display text, + year int, + format text, + thumb_url text, + cover_cid text, + media_grade text, + sleeve_grade text, + rating int, + folder text, + notes text, + release_uri text, + release_cid text, + price_amount int, + price_currency text, + counterparty text, + source_client_agent text, + source_origin text, + source_origin_url text, + source_external_provider text, + source_external_id text, + source_external_url text, + source_record_uri text, + source_record_cid text, + created_at text not null, + updated_at text not null, + indexed_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_entries migration failed"), + ) + use _ <- result.try( + pog.query( + "create index if not exists shelf_entries_did_status_idx + on shelf_entries (did, status)", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_entries did/status index migration failed"), + ) + use _ <- result.try( + pog.query( + "create index if not exists shelf_entries_release_uri_idx + on shelf_entries (release_uri)", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_entries release_uri index migration failed"), + ) + pog.query( + "create table if not exists shelf_index_seen ( + did text primary key, + seen_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("shelf_index_seen migration failed") + |> result.map(fn(_) { Nil }) +} + +fn upsert_event(conn: pog.Connection, event: ShelfEntryEvent) -> Nil { + let outcome = + pog.query( + "insert into shelf_entry_events + (event_uri, entry_uri, did, rkey, action, created_at, record_json, indexed_at) + values ($1, $2, $3, $4, $5, $6, $7, $8) + on conflict (event_uri) do update set + entry_uri = excluded.entry_uri, + did = excluded.did, + rkey = excluded.rkey, + action = excluded.action, + created_at = excluded.created_at, + record_json = excluded.record_json, + indexed_at = excluded.indexed_at", + ) + |> pog.parameter(pog.text(event.event_uri)) + |> pog.parameter(pog.text(event.entry_uri)) + |> pog.parameter(pog.text(event.did)) + |> pog.parameter(pog.text(event.rkey)) + |> pog.parameter(pog.text(event.action)) + |> pog.parameter(pog.text(event.created_at)) + |> pog.parameter(pog.text(event.record_json)) + |> pog.parameter(pog.text(event.indexed_at)) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning( + "shelf_entry_events upsert failed for " <> event.event_uri, + ) + } +} + +fn delete_event(conn: pog.Connection, event_uri: String) -> Nil { + let _ = + pog.query("delete from shelf_entry_events where event_uri = $1") + |> pog.parameter(pog.text(event_uri)) + |> pog.execute(conn) + Nil +} + +fn event_row_decoder() -> decode.Decoder(ShelfEntryEvent) { + use event_uri <- decode.field("event_uri", decode.string) + use entry_uri <- decode.field("entry_uri", decode.string) + use did <- decode.field("did", decode.string) + use rkey <- decode.field("rkey", decode.string) + use action <- decode.field("action", decode.string) + use created_at <- decode.field("created_at", decode.string) + use record_json <- decode.field("record_json", decode.string) + use indexed_at <- decode.field("indexed_at", decode.string) + decode.success(ShelfEntryEvent( + event_uri:, + entry_uri:, + did:, + rkey:, + action:, + created_at:, + record_json:, + indexed_at:, + )) +} + +fn get_event( + conn: pog.Connection, + event_uri: String, +) -> Option(ShelfEntryEvent) { + case + pog.query( + "select event_uri, entry_uri, did, rkey, action, created_at, record_json, indexed_at + from shelf_entry_events where event_uri = $1", + ) + |> pog.parameter(pog.text(event_uri)) + |> pog.returning(event_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [event, ..], ..)) -> Some(event) + _ -> None + } +} + +fn events_for( + conn: pog.Connection, + entry_uri: String, +) -> List(ShelfEntryEvent) { + case + pog.query( + "select event_uri, entry_uri, did, rkey, action, created_at, record_json, indexed_at + from shelf_entry_events where entry_uri = $1 order by rkey", + ) + |> pog.parameter(pog.text(entry_uri)) + |> pog.returning(event_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + +const entry_columns = "entry_uri, did, status, title, artist_display, year, + format, thumb_url, cover_cid, media_grade, sleeve_grade, rating, folder, + notes, release_uri, release_cid, price_amount, price_currency, counterparty, + source_client_agent, source_origin, source_origin_url, + source_external_provider, source_external_id, source_external_url, + source_record_uri, source_record_cid, created_at, updated_at, indexed_at" + +fn upsert_entry(conn: pog.Connection, entry: FoldedEntry) -> Nil { + let outcome = + pog.query("insert into shelf_entries (" <> entry_columns <> ") + values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, + $16, $17, $18, $19, $20, $21, $22, $23, $24, $25, $26, $27, $28, $29, $30) + on conflict (entry_uri) do update set + did = excluded.did, + status = excluded.status, + title = excluded.title, + artist_display = excluded.artist_display, + year = excluded.year, + format = excluded.format, + thumb_url = excluded.thumb_url, + cover_cid = excluded.cover_cid, + media_grade = excluded.media_grade, + sleeve_grade = excluded.sleeve_grade, + rating = excluded.rating, + folder = excluded.folder, + notes = excluded.notes, + release_uri = excluded.release_uri, + release_cid = excluded.release_cid, + price_amount = excluded.price_amount, + price_currency = excluded.price_currency, + counterparty = excluded.counterparty, + source_client_agent = excluded.source_client_agent, + source_origin = excluded.source_origin, + source_origin_url = excluded.source_origin_url, + source_external_provider = excluded.source_external_provider, + source_external_id = excluded.source_external_id, + source_external_url = excluded.source_external_url, + source_record_uri = excluded.source_record_uri, + source_record_cid = excluded.source_record_cid, + created_at = excluded.created_at, + updated_at = excluded.updated_at, + indexed_at = excluded.indexed_at") + |> pog.parameter(pog.text(entry.entry_uri)) + |> pog.parameter(pog.text(entry.did)) + |> pog.parameter(pog.text(entry.status)) + |> pog.parameter(pog.nullable(pog.text, entry.title)) + |> pog.parameter(pog.nullable(pog.text, entry.artist_display)) + |> pog.parameter(pog.nullable(pog.int, entry.year)) + |> pog.parameter(pog.nullable(pog.text, entry.format)) + |> pog.parameter(pog.nullable(pog.text, entry.thumb_url)) + |> pog.parameter(pog.nullable(pog.text, entry.cover_cid)) + |> pog.parameter(pog.nullable(pog.text, entry.media_grade)) + |> pog.parameter(pog.nullable(pog.text, entry.sleeve_grade)) + |> pog.parameter(pog.nullable(pog.int, entry.rating)) + |> pog.parameter(pog.nullable(pog.text, entry.folder)) + |> pog.parameter(pog.nullable(pog.text, entry.notes)) + |> pog.parameter(pog.nullable(pog.text, entry.release_uri)) + |> pog.parameter(pog.nullable(pog.text, entry.release_cid)) + |> pog.parameter(pog.nullable(pog.int, entry.price_amount)) + |> pog.parameter(pog.nullable(pog.text, entry.price_currency)) + |> pog.parameter(pog.nullable(pog.text, entry.counterparty)) + |> pog.parameter(pog.nullable(pog.text, entry.source_client_agent)) + |> pog.parameter(pog.nullable(pog.text, entry.source_origin)) + |> pog.parameter(pog.nullable(pog.text, entry.source_origin_url)) + |> pog.parameter(pog.nullable(pog.text, entry.source_external_provider)) + |> pog.parameter(pog.nullable(pog.text, entry.source_external_id)) + |> pog.parameter(pog.nullable(pog.text, entry.source_external_url)) + |> pog.parameter(pog.nullable(pog.text, entry.source_record_uri)) + |> pog.parameter(pog.nullable(pog.text, entry.source_record_cid)) + |> pog.parameter(pog.text(entry.created_at)) + |> pog.parameter(pog.text(entry.updated_at)) + |> pog.parameter(pog.text(entry.indexed_at)) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning("shelf_entries upsert failed for " <> entry.entry_uri) + } +} + +fn delete_entry(conn: pog.Connection, entry_uri: String) -> Nil { + let _ = + pog.query("delete from shelf_entries where entry_uri = $1") + |> pog.parameter(pog.text(entry_uri)) + |> pog.execute(conn) + Nil +} + +/// Honors an account deletion: drops every folded entry and raw event owned +/// by `did`. Deactivation is not a deletion and must never reach this. +fn delete_for_did(conn: pog.Connection, did: String) -> Nil { + let _ = + pog.query("delete from shelf_entries where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.execute(conn) + let _ = + pog.query("delete from shelf_entry_events where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.execute(conn) + let _ = + pog.query("delete from shelf_index_seen where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.execute(conn) + Nil +} + +fn entry_row_decoder() -> decode.Decoder(FoldedEntry) { + use entry_uri <- decode.field("entry_uri", decode.string) + use did <- decode.field("did", decode.string) + use status <- decode.field("status", decode.string) + use title <- decode.field("title", decode.optional(decode.string)) + use artist_display <- decode.field( + "artist_display", + decode.optional(decode.string), + ) + use year <- decode.field("year", decode.optional(decode.int)) + use format <- decode.field("format", decode.optional(decode.string)) + use thumb_url <- decode.field("thumb_url", decode.optional(decode.string)) + use cover_cid <- decode.field("cover_cid", decode.optional(decode.string)) + use media_grade <- decode.field("media_grade", decode.optional(decode.string)) + use sleeve_grade <- decode.field( + "sleeve_grade", + decode.optional(decode.string), + ) + use rating <- decode.field("rating", decode.optional(decode.int)) + use folder <- decode.field("folder", decode.optional(decode.string)) + use notes <- decode.field("notes", decode.optional(decode.string)) + use release_uri <- decode.field("release_uri", decode.optional(decode.string)) + use release_cid <- decode.field("release_cid", decode.optional(decode.string)) + use price_amount <- decode.field("price_amount", decode.optional(decode.int)) + use price_currency <- decode.field( + "price_currency", + decode.optional(decode.string), + ) + use counterparty <- decode.field( + "counterparty", + decode.optional(decode.string), + ) + use source_client_agent <- decode.field( + "source_client_agent", + decode.optional(decode.string), + ) + use source_origin <- decode.field( + "source_origin", + decode.optional(decode.string), + ) + use source_origin_url <- decode.field( + "source_origin_url", + decode.optional(decode.string), + ) + use source_external_provider <- decode.field( + "source_external_provider", + decode.optional(decode.string), + ) + use source_external_id <- decode.field( + "source_external_id", + decode.optional(decode.string), + ) + use source_external_url <- decode.field( + "source_external_url", + decode.optional(decode.string), + ) + use source_record_uri <- decode.field( + "source_record_uri", + decode.optional(decode.string), + ) + use source_record_cid <- decode.field( + "source_record_cid", + decode.optional(decode.string), + ) + use created_at <- decode.field("created_at", decode.string) + use updated_at <- decode.field("updated_at", decode.string) + use indexed_at <- decode.field("indexed_at", decode.string) + decode.success(FoldedEntry( + entry_uri:, + did:, + status:, + title:, + artist_display:, + year:, + format:, + thumb_url:, + cover_cid:, + media_grade:, + sleeve_grade:, + rating:, + folder:, + notes:, + release_uri:, + release_cid:, + price_amount:, + price_currency:, + counterparty:, + source_client_agent:, + source_origin:, + source_origin_url:, + source_external_provider:, + source_external_id:, + source_external_url:, + source_record_uri:, + source_record_cid:, + created_at:, + updated_at:, + indexed_at:, + )) +} + +fn get_entry(conn: pog.Connection, entry_uri: String) -> Option(FoldedEntry) { + case + pog.query( + "select " <> entry_columns <> " from shelf_entries where entry_uri = $1", + ) + |> pog.parameter(pog.text(entry_uri)) + |> pog.returning(entry_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [entry, ..], ..)) -> Some(entry) + _ -> None + } +} + +fn list_for_did(conn: pog.Connection, did: String) -> List(FoldedEntry) { + case + pog.query( + "select " <> entry_columns <> " from shelf_entries where did = $1", + ) + |> pog.parameter(pog.text(did)) + |> pog.returning(entry_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + +fn has_seen(conn: pog.Connection, did: String) -> Bool { + let row = { + use did <- decode.field("did", decode.string) + decode.success(did) + } + case + pog.query("select did from shelf_index_seen where did = $1") + |> pog.parameter(pog.text(did)) + |> pog.returning(row) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [_, ..], ..)) -> True + _ -> False + } +} + +fn mark_seen(conn: pog.Connection, did: String) -> Nil { + let outcome = + pog.query( + "insert into shelf_index_seen (did, seen_at) values ($1, $2) + on conflict (did) do nothing", + ) + |> pog.parameter(pog.text(did)) + |> pog.parameter(pog.text(provenance.now_rfc3339())) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> wisp.log_warning("shelf_index_seen mark failed for " <> did) + } +} diff --git a/server/src/at_record_server/wiring.gleam b/server/src/at_record_server/wiring.gleam index 87c36c2..89e79ef 100644 --- a/server/src/at_record_server/wiring.gleam +++ b/server/src/at_record_server/wiring.gleam @@ -13,6 +13,7 @@ import at_record_server/identity_cache.{type Cache, type Identity, Identity} import at_record_server/identity_resolver.{Resolver} import at_record_server/musicbrainz_client import at_record_server/promotion +import at_record_server/shelf_index import atproto/constellation import atproto/xrpc.{type Client} import gleam/dynamic/decode @@ -29,6 +30,7 @@ pub fn context( discogs_creds discogs_creds, known_users known_users, catalog_index catalog_index, + shelf_index shelf_index: shelf_index.Store, identity_cache identity_cache: Cache, oauth oauth, ) -> Context { @@ -46,6 +48,7 @@ pub fn context( catalog:, known_users:, catalog_index:, + shelf_index:, variant_source: variant_source(catalog_index), oauth:, ) diff --git a/server/test/shelf_index_postgres_test.gleam b/server/test/shelf_index_postgres_test.gleam new file mode 100644 index 0000000..0640d86 --- /dev/null +++ b/server/test/shelf_index_postgres_test.gleam @@ -0,0 +1,136 @@ +//// Round-trip checks for the Postgres backend, gated on `DATABASE_URL` (a +//// no-op when unset); run with a live Postgres to regression-test it. + +import at_record/gen/repo/strong_ref.{type RepoStrongRef} +import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} +import at_record_server/oauth/sessions_postgres +import at_record_server/shelf_index +import at_record_server/shelf_index_postgres +import envoy +import gleam/option.{None, Some} + +fn with_store(run: fn(shelf_index.Store) -> Nil) -> Nil { + case envoy.get("DATABASE_URL") { + Error(Nil) -> Nil + Ok(url) -> { + let assert Ok(conn) = sessions_postgres.connect_pool(url) + let assert Ok(store) = shelf_index_postgres.table_store(conn) + run(store) + } + } +} + +fn blank_entry( + action action: String, + created_at created_at: String, + subject subject: option.Option(RepoStrongRef), +) -> ShelfEntry { + ShelfEntry( + action:, + counterparty: Some("discogs:seller123"), + created_at:, + external_ids: None, + folder: None, + media_grade: Some("VG+"), + notes: Some("first pressing"), + price: None, + rating: None, + release: None, + sleeve_grade: None, + snapshot: None, + source: None, + subject:, + ) +} + +pub fn event_round_trip_test() { + use store <- with_store() + let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg1" + let event = + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did: "did:plc:pgtest", + rkey: "pg1", + entry: blank_entry( + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + subject: None, + ), + ) + store.events.upsert(event) + let assert Some(found) = store.events.get(entry_uri) + assert found.entry_uri == entry_uri + assert found.did == "did:plc:pgtest" + assert found.action == "acquired" + let assert [only] = store.events.for_entry(entry_uri) + assert only.event_uri == entry_uri + store.events.delete(entry_uri) + assert store.events.get(entry_uri) == None +} + +pub fn record_and_fold_round_trips_flattened_columns_test() { + use store <- with_store() + let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg2" + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did: "did:plc:pgtest", + rkey: "pg2", + entry: blank_entry( + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + subject: None, + ), + ), + ) + let assert Some(found) = store.entries.get(entry_uri) + assert found.did == "did:plc:pgtest" + assert found.status == "owned" + assert found.media_grade == Some("VG+") + assert found.notes == Some("first pressing") + assert found.counterparty == Some("discogs:seller123") + assert found.created_at == "2024-01-01T00:00:00Z" + + let assert [in_list] = store.entries.list_for_did("did:plc:pgtest") + assert in_list.entry_uri == entry_uri + + shelf_index.delete_and_fold(store, entry_uri, entry_uri) + assert store.entries.get(entry_uri) == None + assert store.events.for_entry(entry_uri) == [] +} + +pub fn seen_round_trip_test() { + use store <- with_store() + let did = "did:plc:pgtest-seen" + assert store.seen.has_seen(did) == False + store.seen.mark_seen(did) + assert store.seen.has_seen(did) == True +} + +pub fn delete_for_did_round_trip_test() { + use store <- with_store() + let did = "did:plc:pgtest-purge" + let entry_uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/pg3" + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did:, + rkey: "pg3", + entry: blank_entry( + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + subject: None, + ), + ), + ) + store.seen.mark_seen(did) + store.entries.delete_for_did(did) + assert store.entries.get(entry_uri) == None + assert store.events.for_entry(entry_uri) == [] + assert store.seen.has_seen(did) == False +} diff --git a/server/test/shelf_index_test.gleam b/server/test/shelf_index_test.gleam new file mode 100644 index 0000000..dc0f473 --- /dev/null +++ b/server/test/shelf_index_test.gleam @@ -0,0 +1,322 @@ +//// In-memory `shelf_index.Store` coverage: fold semantics via +//// `record_and_fold`/`delete_and_fold` (genesis-only, appends, out-of-order +//// arrival, mid-history deletes, purge-to-zero), plus `seen`/`delete_for_did`. + +import at_record/gen/repo/strong_ref.{type RepoStrongRef, RepoStrongRef} +import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} +import at_record_server/shelf_index +import gleam/list +import gleam/option.{None, Some} + +/// A genesis event: no `subject`, so `event_uri == entry_uri`. +fn genesis( + entry_uri entry_uri: String, + did did: String, + rkey rkey: String, + action action: String, + created_at created_at: String, +) -> shelf_index.ShelfEntryEvent { + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did:, + rkey:, + entry: blank_entry(action:, created_at:, subject: None, media_grade: None), + ) +} + +/// An append event: carries `subject` pointing back at the genesis, its own +/// `event_uri`/`rkey`. +fn append( + entry_uri entry_uri: String, + event_uri event_uri: String, + did did: String, + rkey rkey: String, + action action: String, + created_at created_at: String, +) -> shelf_index.ShelfEntryEvent { + shelf_index.event_row( + event_uri:, + entry_uri:, + did:, + rkey:, + entry: blank_entry( + action:, + created_at:, + subject: Some(RepoStrongRef(cid: "bafygenesis", uri: entry_uri)), + media_grade: None, + ), + ) +} + +/// An append event that also sets `mediaGrade`, for the mid-history delete +/// test. +fn regrade_append( + entry_uri entry_uri: String, + event_uri event_uri: String, + did did: String, + rkey rkey: String, + created_at created_at: String, + grade grade: String, +) -> shelf_index.ShelfEntryEvent { + shelf_index.event_row( + event_uri:, + entry_uri:, + did:, + rkey:, + entry: blank_entry( + action: "regraded", + created_at:, + subject: Some(RepoStrongRef(cid: "bafygenesis", uri: entry_uri)), + media_grade: Some(grade), + ), + ) +} + +fn blank_entry( + action action: String, + created_at created_at: String, + subject subject: option.Option(RepoStrongRef), + media_grade media_grade: option.Option(String), +) -> ShelfEntry { + ShelfEntry( + action:, + counterparty: None, + created_at:, + external_ids: None, + folder: None, + media_grade:, + notes: None, + price: None, + rating: None, + release: None, + sleeve_grade: None, + snapshot: None, + source: None, + subject:, + ) +} + +pub fn genesis_only_fold_test() { + let assert Ok(store) = shelf_index.start() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + genesis( + entry_uri:, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + let assert Some(folded) = store.entries.get(entry_uri) + assert folded.did == "did:plc:a" + assert folded.status == "owned" + assert folded.created_at == "2024-01-01T00:00:00Z" + assert folded.updated_at == "2024-01-01T00:00:00Z" +} + +pub fn genesis_and_appends_status_transitions_test() { + let assert Ok(store) = shelf_index.start() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + genesis( + entry_uri:, + did: "did:plc:a", + rkey: "e1", + action: "wanted", + created_at: "2024-01-01T00:00:00Z", + ), + ) + let assert Some(wanted) = store.entries.get(entry_uri) + assert wanted.status == "wanted" + + shelf_index.record_and_fold( + store, + append( + entry_uri:, + event_uri: entry_uri <> "-e2", + did: "did:plc:a", + rkey: "e2", + action: "acquired", + created_at: "2024-02-01T00:00:00Z", + ), + ) + let assert Some(owned) = store.entries.get(entry_uri) + assert owned.status == "owned" + assert owned.updated_at == "2024-02-01T00:00:00Z" + + shelf_index.record_and_fold( + store, + append( + entry_uri:, + event_uri: entry_uri <> "-e3", + did: "did:plc:a", + rkey: "e3", + action: "sold", + created_at: "2024-03-01T00:00:00Z", + ), + ) + let assert Some(gone) = store.entries.get(entry_uri) + assert gone.status == "gone" + assert gone.updated_at == "2024-03-01T00:00:00Z" +} + +pub fn out_of_order_append_before_genesis_is_a_noop_that_self_heals_test() { + let assert Ok(store) = shelf_index.start() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + append( + entry_uri:, + event_uri: entry_uri <> "-e2", + did: "did:plc:a", + rkey: "e2", + action: "sold", + created_at: "2024-02-01T00:00:00Z", + ), + ) + // No genesis yet, so folding finds nothing to project: a no-op, not a + // partial or wrong row. + assert store.entries.get(entry_uri) == None + // The event is still mirrored even though it didn't fold into anything. + assert list.length(store.events.for_entry(entry_uri)) == 1 + + shelf_index.record_and_fold( + store, + genesis( + entry_uri:, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + let assert Some(healed) = store.entries.get(entry_uri) + assert healed.status == "gone" +} + +pub fn delete_mid_history_refolds_correctly_test() { + let assert Ok(store) = shelf_index.start() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let regrade_uri = entry_uri <> "-e2" + shelf_index.record_and_fold( + store, + genesis( + entry_uri:, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + shelf_index.record_and_fold( + store, + regrade_append( + entry_uri:, + event_uri: regrade_uri, + did: "did:plc:a", + rkey: "e2", + created_at: "2024-02-01T00:00:00Z", + grade: "NM", + ), + ) + let assert Some(before) = store.entries.get(entry_uri) + assert before.media_grade == Some("NM") + + shelf_index.delete_and_fold(store, regrade_uri, entry_uri) + let assert Some(after) = store.entries.get(entry_uri) + assert after.media_grade == None + assert after.status == "owned" + assert list.length(store.events.for_entry(entry_uri)) == 1 +} + +pub fn purge_to_zero_deletes_the_folded_row_test() { + let assert Ok(store) = shelf_index.start() + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + genesis( + entry_uri:, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + assert store.entries.get(entry_uri) != None + shelf_index.delete_and_fold(store, entry_uri, entry_uri) + assert store.entries.get(entry_uri) == None + assert store.events.for_entry(entry_uri) == [] +} + +pub fn has_seen_and_mark_seen_test() { + let assert Ok(store) = shelf_index.start() + assert store.seen.has_seen("did:plc:a") == False + store.seen.mark_seen("did:plc:a") + assert store.seen.has_seen("did:plc:a") == True +} + +pub fn list_for_did_returns_only_that_dids_entries_test() { + let assert Ok(store) = shelf_index.start() + let entry_a = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let entry_b = "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + genesis( + entry_uri: entry_a, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + shelf_index.record_and_fold( + store, + genesis( + entry_uri: entry_b, + did: "did:plc:b", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + let assert [only] = store.entries.list_for_did("did:plc:a") + assert only.entry_uri == entry_a +} + +pub fn delete_for_did_drops_entries_events_and_seen_test() { + let assert Ok(store) = shelf_index.start() + let entry_a = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let entry_b = "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e1" + shelf_index.record_and_fold( + store, + genesis( + entry_uri: entry_a, + did: "did:plc:a", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + shelf_index.record_and_fold( + store, + genesis( + entry_uri: entry_b, + did: "did:plc:b", + rkey: "e1", + action: "acquired", + created_at: "2024-01-01T00:00:00Z", + ), + ) + store.seen.mark_seen("did:plc:a") + + store.entries.delete_for_did("did:plc:a") + + assert store.entries.get(entry_a) == None + assert store.events.for_entry(entry_a) == [] + assert store.seen.has_seen("did:plc:a") == False + assert store.entries.get(entry_b) != None +} diff --git a/server/test/support.gleam b/server/test/support.gleam index 04bd7a9..98f376f 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -19,6 +19,7 @@ import at_record_server/oauth/keys import at_record_server/oauth/sessions import at_record_server/oauth/sessions_memory import at_record_server/oauth/store +import at_record_server/shelf_index import at_record_server/wiring import atproto/xrpc import atproto_core/xrpc as core_xrpc @@ -137,6 +138,14 @@ fn empty_catalog_index() -> catalog_index.Store { ) } +/// A real in-memory `shelf_index.Store`: cheap to start, and tests that +/// actually exercise it can read it back directly instead of hand-rolling a +/// spy, mirroring `empty_catalog_index`'s dev-fallback backend. +fn fresh_shelf_index() -> shelf_index.Store { + let assert Ok(store) = shelf_index.start() + store +} + /// A `variant_source` with no candidate rows and a zero adoption count for /// every uri; use `stub_context_with_variant_source` to override it. pub fn empty_variant_source() -> catalog_source.Source { @@ -198,6 +207,7 @@ pub fn stub_context_with( catalog:, known_users: known_users_of(known_users), catalog_index: empty_catalog_index(), + shelf_index: fresh_shelf_index(), variant_source: empty_variant_source(), oauth: cfg, )