diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index d20ac25..9af2d12 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -74,6 +74,10 @@ pub type ReleaseOps { /// deletion (see `jetstream_consumer`'s account-event handling). delete_for_did: fn(String) -> Nil, list: fn() -> List(BrowseRow), + /// One release by its own at-uri (C5 of the appview-first roadmap): the + /// entry-timeline's release-info read, so it no longer needs a live + /// Slingshot `getRecord`. + get: fn(String) -> Option(BrowseRow), ) } @@ -126,6 +130,7 @@ type Msg { Delete(String) DeleteReleasesForDid(String) List(Subject(List(BrowseRow))) + Get(String, Subject(Option(BrowseRow))) UpsertAdoption(Adoption) DeleteAdoption(String) DeleteAdoptionsForDid(String) @@ -174,6 +179,12 @@ pub fn start() -> Result(Store, actor.StartError) { Error(Nil) -> [] } }, + get: fn(uri) { + case parallel.try_call(subject, 1000, Get(uri, _)) { + Ok(row) -> row + Error(Nil) -> option.None + } + }, ), adoptions: AdoptionOps( upsert: fn(adoption) { @@ -263,6 +274,10 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { process.send(reply, dict.values(state.releases)) actor.continue(state) } + Get(uri, reply) -> { + process.send(reply, dict.get(state.releases, uri) |> option.from_result) + actor.continue(state) + } UpsertAdoption(adoption) -> actor.continue( State( diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam index b3d5e42..484bc97 100644 --- a/server/src/at_record_server/catalog_index_postgres.gleam +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -33,6 +33,7 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { delete: fn(uri) { delete_release(conn, uri) }, delete_for_did: fn(did) { delete_releases_for_did(conn, did) }, list: fn() { list_releases(conn) }, + get: fn(uri) { get_release(conn, uri) }, ), adoptions: AdoptionOps( upsert: fn(adoption) { upsert_adoption(conn, adoption) }, @@ -293,6 +294,24 @@ fn list_releases(conn: pog.Connection) -> List(BrowseRow) { } } +fn get_release(conn: pog.Connection, uri: String) -> Option(BrowseRow) { + case + pog.query( + "select uri, cid, title, artist_display, genres, styles, released, + country, cover_cid, thumb_url, discogs_id, created_at, + publisher_did, publisher_handle, publisher_pds, supersedes, based_on, + format, label, master + from catalog_releases where uri = $1", + ) + |> pog.parameter(pog.text(uri)) + |> pog.returning(release_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [row, ..], ..)) -> Some(row) + _ -> None + } +} + fn upsert_adoption(conn: pog.Connection, adoption: Adoption) -> Nil { let outcome = pog.query( diff --git a/server/src/at_record_server/crate.gleam b/server/src/at_record_server/crate.gleam index 8a2dabb..f329301 100644 --- a/server/src/at_record_server/crate.gleam +++ b/server/src/at_record_server/crate.gleam @@ -10,6 +10,7 @@ import at_record/gen/defs.{ } import at_record/gen/shelf/entry.{type ShelfEntry} import at_record/storage.{type StoredItem} +import at_record_server/catalog/row.{type BrowseRow} import gleam/dict import gleam/json.{type Json} import gleam/list @@ -76,6 +77,17 @@ pub fn status_string(status: Status) -> String { } } +/// The inverse of `status_string`, for folding a stored/flattened status +/// column back into `Status`; an unrecognized value defaults to `Owned` +/// (matching `fold_group`'s seed) rather than failing the read. +pub fn status_from_string(status: String) -> Status { + case status { + "wanted" -> Wanted + "gone" -> Gone + _ -> Owned + } +} + /// The genesis event action implied by a requested status string. pub fn action_from_status(status: String) -> String { case status { @@ -196,3 +208,32 @@ pub fn encode_release_info(release: CatalogRelease) -> Json { ]), ) } + +/// `encode_release_info`'s sibling over the catalog index's `BrowseRow` +/// (C5 of the appview-first roadmap): same field set and omission rules, so +/// swapping the entry-timeline release read from a live fetch to the index +/// changes nothing on the wire. `genres`/`styles` are flat (never absent) on +/// `BrowseRow`, so an empty list maps to the same omitted-field shape a +/// `None` genres/styles produces from a live `CatalogRelease`. +pub fn encode_release_info_row(row: BrowseRow) -> Json { + json.object( + list.flatten([ + opt("artistDisplay", row.artist_display, json.string), + opt("genres", as_option(row.genres), fn(items) { + json.array(items, json.string) + }), + opt("styles", as_option(row.styles), fn(items) { + json.array(items, json.string) + }), + opt("country", row.country, json.string), + opt("released", row.released, json.string), + ]), + ) +} + +fn as_option(items: List(a)) -> Option(List(a)) { + case items { + [] -> None + _ -> Some(items) + } +} diff --git a/server/src/at_record_server/handlers/amend.gleam b/server/src/at_record_server/handlers/amend.gleam index 69ce3bf..3cdb6e1 100644 --- a/server/src/at_record_server/handlers/amend.gleam +++ b/server/src/at_record_server/handlers/amend.gleam @@ -19,8 +19,10 @@ import at_record_server/event_log import at_record_server/external_id import at_record_server/oauth/sessions.{type OauthSession} import at_record_server/provenance +import at_record_server/shelf_index import atproto/blob import atproto/repo +import atproto/uri import atproto/xrpc.{type Client} import gleam/bit_array import gleam/dynamic/decode @@ -208,7 +210,16 @@ fn amend_entry_group( None -> error_json(502, "could not publish the amended release") Some(new_ref) -> { propose_edit(client, session, entry.release, form) - write_repoint(client, session, genesis, entry, form, new_ref, fresh_cover) + write_repoint( + ctx, + client, + session, + genesis, + entry, + form, + new_ref, + fresh_cover, + ) } } } @@ -473,6 +484,7 @@ pub fn repoint_event( } fn write_repoint( + ctx: Context, client: Client, session: OauthSession, genesis: storage.StoredItem(ShelfEntry), @@ -484,12 +496,25 @@ fn write_repoint( let event = repoint_event(genesis, folded, form, new_ref, fresh_cover) case event_log.append(client, session, event) { Error(_) -> error_json(502, "could not write the amend event") - Ok(_) -> + Ok(created) -> { + // Read-your-writes (ADR 0002): a repoint is always an append onto the + // genesis, so its entry_uri is always the genesis's own uri. + shelf_index.record_and_fold( + ctx.shelf_index, + shelf_index.event_row( + event_uri: created.uri, + entry_uri: genesis.uri, + did: session.did, + rkey: uri.rkey(created.uri), + entry: event, + ), + ) json.object([ #("entryId", json.string(folded.entry_id)), #("releaseUri", json.string(new_ref.uri)), ]) |> json.to_string |> wisp.json_response(201) + } } } diff --git a/server/src/at_record_server/handlers/shelf.gleam b/server/src/at_record_server/handlers/shelf.gleam index 79469c4..a203bec 100644 --- a/server/src/at_record_server/handlers/shelf.gleam +++ b/server/src/at_record_server/handlers/shelf.gleam @@ -2,7 +2,6 @@ //// folded crate and per-entry timeline, and writing immutable events (genesis //// add + later actions). Every mutation is a new event, never an edit. -import at_record/gen/catalog/release as catalog_release import at_record/gen/defs.{ type ExternalId, type Price, ExternalId, Snapshot, price_decoder, } @@ -11,6 +10,7 @@ import at_record/gen/shelf/entry.{ type ShelfEntry, ShelfEntry, encode_shelf_entry, } import at_record/storage.{type StoredItem} +import at_record_server/catalog/row.{type BrowseRow} import at_record_server/catalog_entities import at_record_server/context.{ type Context, error_json, require_session, with_pds_client, @@ -25,6 +25,7 @@ import at_record_server/oauth/sessions.{type OauthSession} import at_record_server/pagination import at_record_server/promotion import at_record_server/provenance.{now_rfc3339} +import at_record_server/shelf_index import at_record_server/shelf_owner import atproto/blob import atproto/uri @@ -49,40 +50,99 @@ pub fn list_shelf(req: Request, ctx: Context) -> Response { } fn list_own_shelf(req: Request, ctx: Context) -> Response { - use id, session <- require_session(req, ctx) - use client, session <- with_pds_client(ctx, id, session) - case event_log.load(client, session) { - Error(_) -> error_json(502, "could not load crate from PDS") - Ok(events) -> - shelf_response(req, ctx, session.did, session.handle, events, False) - } + use _id, session <- require_session(req, ctx) + let entries = + ctx.shelf_index.entries.list_for_did(session.did) + |> list.map(shelf_index.to_crate_entry) + shelf_response(req, ctx, session.did, session.handle, entries, False) } fn list_actor_shelf(req: Request, ctx: Context, actor: String) -> Response { case shelf_owner.resolve_actor(ctx, actor) { Error(Nil) -> error_json(404, "could not resolve that user") Ok(#(did, handle, pds)) -> - case shelf_owner.fetch_public_entries(ctx, pds, did) { - Error(shelf_owner.FetchFailed) -> - error_json(502, "could not load that user's crate from their PDS") - Error(shelf_owner.RepoNotFound) -> - error_json(404, "that user's crate could not be found") - Error(shelf_owner.TooManyPages) -> - error_json(413, "that user's crate is too large to page through") - Ok(events) -> shelf_response(req, ctx, did, handle, events, True) + case ensure_actor_seeded(ctx, pds, did) { + Error(err) -> actor_fetch_error_response(err) + Ok(Nil) -> { + let entries = + ctx.shelf_index.entries.list_for_did(did) + |> list.map(shelf_index.to_crate_entry) + shelf_response(req, ctx, did, handle, entries, True) + } } } } -/// The shared read pipeline for own and actor mode: fold to one page, resolve -/// covers, attribute adopted releases, shape the response. `public` is True -/// on the unauthenticated actor path. +/// A never-seen actor's crate is fetched from their PDS exactly once and +/// seeded into the index (the ADR's documented live-read exception), so this +/// request and every later one for the same did come from Postgres. A seen +/// did is a no-op. `fetch_public_entries` is all-or-nothing (a mid-paging +/// failure discards whatever paged so far rather than returning it), so a +/// failed fetch here seeds nothing and propagates the error unchanged -- +/// same 502/404/413 semantics the always-live path had -- and leaves the did +/// unseen so the next visit retries. Marking it seen on failure would cache +/// a stranger's crate as permanently empty: reconcile only re-backfills +/// *known* users, never a visited-but-unreachable stranger. +fn ensure_actor_seeded( + ctx: Context, + pds: String, + did: String, +) -> Result(Nil, shelf_owner.FetchEntriesError) { + case ctx.shelf_index.seen.has_seen(did) { + True -> Ok(Nil) + False -> { + use events <- result.try(shelf_owner.fetch_public_entries(ctx, pds, did)) + events |> list.each(seed_event(ctx, did, _)) + ctx.shelf_index.seen.mark_seen(did) + Ok(Nil) + } + } +} + +fn actor_fetch_error_response(err: shelf_owner.FetchEntriesError) -> Response { + case err { + shelf_owner.FetchFailed -> + error_json(502, "could not load that user's crate from their PDS") + shelf_owner.RepoNotFound -> + error_json(404, "that user's crate could not be found") + shelf_owner.TooManyPages -> + error_json(413, "that user's crate is too large to page through") + } +} + +/// A genesis event's `entry_uri` is its own uri, an append's is its +/// `subject.uri` -- mirrors `jetstream_consumer.upsert_shelf_index_event`'s +/// derivation off a firehose frame. +fn subject_entry_uri(s: StoredItem(ShelfEntry)) -> String { + case s.value.subject { + Some(ref) -> ref.uri + None -> s.uri + } +} + +fn seed_event(ctx: Context, did: String, s: StoredItem(ShelfEntry)) -> Nil { + shelf_index.record_and_fold( + ctx.shelf_index, + shelf_index.event_row( + event_uri: s.uri, + entry_uri: subject_entry_uri(s), + did:, + rkey: s.rkey, + entry: s.value, + ), + ) +} + +/// The shared read pipeline for own and actor mode: filter to one page, +/// resolve covers, attribute adopted releases, shape the response. `entries` +/// are already folded (from the index); `public` is True on the +/// unauthenticated actor path. fn shelf_response( req: Request, ctx: Context, did: String, handle: String, - events: List(StoredItem(ShelfEntry)), + entries: List(crate.CrateEntry), public: Bool, ) -> Response { let query = wisp.get_query(req) @@ -93,7 +153,7 @@ fn shelf_response( True -> "current" False -> view } - let #(page, next_cursor) = folded_page(events, view, query) + let #(page, next_cursor) = folded_page(entries, view, query) let items = list.map(page, resolve_cover(did, _)) json.object( list.flatten([ @@ -113,19 +173,18 @@ fn shelf_response( |> wisp.json_response(200) } -/// One page of the folded view. The fold needs every event to reduce -/// correctly, so pagination slices the folded, filtered, deterministically -/// ordered view rather than the raw log; snapshot-less remnants would fail -/// the frontend decoder, so they never leave the BFF. +/// One page of the folded view: filtered, deterministically ordered, then +/// sliced; snapshot-less remnants would fail the frontend decoder, so they +/// never leave the BFF. fn folded_page( - events: List(StoredItem(ShelfEntry)), + entries: List(crate.CrateEntry), view: String, query: List(#(String, String)), ) -> #(List(crate.CrateEntry), Option(String)) { let cursor = list.key_find(query, "cursor") |> option.from_result let limit = list.key_find(query, "limit") |> result.try(int.parse) |> option.from_result - crate.fold(events) + entries |> list.filter(in_view(view, _)) |> list.filter(fn(e) { e.snapshot != None }) |> list.sort(fn(a, b) { string.compare(a.entry_id, b.entry_id) }) @@ -232,83 +291,80 @@ pub fn get_entry(req: Request, ctx: Context) -> Response { } fn own_entry(req: Request, ctx: Context, entry_id: String) -> Response { - use id, session <- require_session(req, ctx) - use client, session <- with_pds_client(ctx, id, session) - case event_log.load(client, session) { - Error(_) -> error_json(502, "could not load crate from PDS") - Ok(stored) -> - entry_response(ctx, session.did, session.handle, stored, entry_id) - } + use _id, session <- require_session(req, ctx) + entry_response(ctx, session.did, session.handle, entry_id) } fn actor_entry(ctx: Context, actor: String, entry_id: String) -> Response { case shelf_owner.resolve_actor(ctx, actor) { Error(Nil) -> error_json(404, "could not resolve that user") Ok(#(did, handle, pds)) -> - case shelf_owner.fetch_public_entries(ctx, pds, did) { - Error(shelf_owner.FetchFailed) -> - error_json(502, "could not load that user's crate from their PDS") - Error(shelf_owner.RepoNotFound) -> - error_json(404, "that user's crate could not be found") - Error(shelf_owner.TooManyPages) -> - error_json(413, "that user's crate is too large to page through") - Ok(stored) -> entry_response(ctx, did, handle, stored, entry_id) + case ensure_actor_seeded(ctx, pds, did) { + Error(err) -> actor_fetch_error_response(err) + Ok(Nil) -> entry_response(ctx, did, handle, entry_id) } } } +/// One entry's superset from the index: the folded state (`EntryOps.get`) +/// plus its raw event timeline (`EventOps.for_entry`), replacing the old +/// live-load-then-fold. Own and actor mode share this once the actor path +/// has ensured the did is seeded. fn entry_response( ctx: Context, did: String, handle: String, - stored: List(StoredItem(ShelfEntry)), entry_id: String, ) -> Response { - case event_log.entry_events(stored, entry_id) { - Error(Nil) -> error_json(404, "unknown entry") - Ok(events) -> - case crate.fold(events) |> list.first { - Error(Nil) -> error_json(404, "unknown entry") - Ok(folded) -> { - let release = resolve_release(ctx, events) - json.object( - list.flatten([ - [ - #("did", json.string(did)), - #("handle", json.string(handle)), - #("entry", crate.encode_entry(resolve_cover(did, folded))), - #( - "events", - json.array( - list.map(events, fn(s) { s.value }), - encode_shelf_entry, - ), - ), - ], - case release { - Some(r) -> [#("release", crate.encode_release_info(r))] - None -> [] - }, - ]), - ) - |> json.to_string - |> wisp.json_response(200) - } - } + let uri = entry_at_uri(did, entry_id) + case ctx.shelf_index.entries.get(uri) { + None -> error_json(404, "unknown entry") + Some(folded) -> { + let events = + ctx.shelf_index.events.for_entry(uri) |> list.filter_map(decode_event) + let release = resolve_release(ctx, folded) + json.object( + list.flatten([ + [ + #("did", json.string(did)), + #("handle", json.string(handle)), + #( + "entry", + crate.encode_entry(resolve_cover( + did, + shelf_index.to_crate_entry(folded), + )), + ), + #("events", json.array(events, encode_shelf_entry)), + ], + case release { + Some(r) -> [#("release", crate.encode_release_info_row(r))] + None -> [] + }, + ]), + ) + |> json.to_string + |> wisp.json_response(200) + } } } -// Best-effort read: no ref, or a fetch miss, just means the response omits it. +fn entry_at_uri(did: String, entry_id: String) -> String { + "at://" <> did <> "/" <> entry.collection <> "/" <> entry_id +} + +fn decode_event(e: shelf_index.ShelfEntryEvent) -> Result(ShelfEntry, Nil) { + json.parse(e.record_json, entry.shelf_entry_decoder()) + |> result.replace_error(Nil) +} + +// Best-effort read: no ref, or an index miss, just means the response omits +// it -- the same degrade a failed live fetch produced before C5. fn resolve_release( ctx: Context, - events: List(StoredItem(ShelfEntry)), -) -> Option(catalog_release.CatalogRelease) { - crate.fold(events) - |> list.first - |> result.map(fn(e) { e.release }) - |> result.unwrap(None) - |> option.then(fn(ref) { ctx.catalog.fetch_release(ref.uri) }) - |> option.map(fn(pair) { pair.1 }) + folded: shelf_index.FoldedEntry, +) -> Option(BrowseRow) { + folded.release_uri |> option.then(ctx.catalog_index.releases.get) } pub fn purge_entry(req: Request, ctx: Context) -> Response { @@ -338,14 +394,25 @@ fn do_purge( event_log.entry_events(stored, entry_id) |> result.unwrap([]) |> list.each(fn(s) { - let _ = event_log.delete(client, session, s.rkey) - Nil + case event_log.delete(client, session, s.rkey) { + Ok(_) -> write_through_delete(ctx, s) + Error(_) -> Nil + } }) wisp.json_response("{}", 200) } } } +/// Read-your-writes for a purge: without this, a deleted entry would stay +/// visible in the owner's own crate until the firehose catches up. Only +/// called once the PDS delete of this specific event has actually +/// succeeded, mirroring `jetstream_consumer.delete_from_shelf_index`'s +/// `entry_uri` derivation off the event this handler already has in hand. +fn write_through_delete(ctx: Context, s: StoredItem(ShelfEntry)) -> Nil { + shelf_index.delete_and_fold(ctx.shelf_index, s.uri, subject_entry_uri(s)) +} + pub type AddForm { AddForm( title: String, @@ -526,7 +593,7 @@ fn do_add( error_json(502, "artist records could not be published; try again") Ok(#(release, source)) -> genesis_event(form, external_ids, release, source, cover, now_rfc3339()) - |> write_event(client, session) + |> write_event(ctx, client, session) } } @@ -713,13 +780,14 @@ fn do_append( source: Some(provenance.source("manual")), created_at: now_rfc3339(), ) - |> write_event(client, session) + |> write_event(ctx, client, session) } } } fn write_event( event: ShelfEntry, + ctx: Context, client: Client, session: OauthSession, ) -> Response { @@ -727,10 +795,23 @@ fn write_event( Error(_) -> error_json(502, "could not write event to PDS") Ok(created) -> { // The entry id is the genesis TID, regardless of which record this event is. - let entry_id = case event.subject { - Some(ref) -> uri.rkey(ref.uri) - None -> uri.rkey(created.uri) + let entry_uri = case event.subject { + Some(ref) -> ref.uri + None -> created.uri } + let entry_id = uri.rkey(entry_uri) + // Read-your-writes (ADR 0002): without this, the user's own add/append + // would vanish from their crate until the firehose catches up. + shelf_index.record_and_fold( + ctx.shelf_index, + shelf_index.event_row( + event_uri: created.uri, + entry_uri:, + did: session.did, + rkey: uri.rkey(created.uri), + entry: event, + ), + ) json.object([ #("uri", json.string(created.uri)), #("cid", json.string(created.cid)), diff --git a/server/src/at_record_server/shelf_index.gleam b/server/src/at_record_server/shelf_index.gleam index abadcae..fd76e55 100644 --- a/server/src/at_record_server/shelf_index.gleam +++ b/server/src/at_record_server/shelf_index.gleam @@ -6,20 +6,29 @@ //// 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. +//// `record_and_fold`/`delete_and_fold` are the shared write seam the +//// jetstream consumer, the handlers' write-through path, and the +//// never-seen-actor live-fetch seed all call (C2/C3 of the appview-first +//// roadmap); `to_crate_entry` is the read-side mirror, turning a folded row +//// back into the shape the existing response encoders already know. +import at_record/gen/defs.{ + type ExternalId, type Price, type Snapshot, type Source, CatalogRef, + ExternalId, Price, Snapshot, Source, +} +import at_record/gen/repo/strong_ref.{type RepoStrongRef, RepoStrongRef} 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/crate.{type CrateEntry, CrateEntry} import at_record_server/parallel import at_record_server/provenance +import atproto/blob.{Blob} +import atproto/uri 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/option.{type Option, None, Some} import gleam/otp/actor import gleam/result import gleam/set.{type Set} @@ -245,6 +254,110 @@ fn to_folded_entry( ) } +/// The inverse of `to_folded_entry`: reconstitute a `crate.CrateEntry` from +/// one folded row, so the existing `crate.encode_entry`/response-shaping +/// code (own-shelf, actor-shelf, entry timeline) can render an index read +/// exactly like it rendered a freshly-folded PDS read (C3 of the +/// appview-first roadmap). `events` is always `[]`: `encode_entry` never +/// reads it, and the entry-timeline response gets its event list from +/// `EventOps.for_entry` separately, not from the folded row. A reconstructed +/// cover blob's `mime_type`/`size` are placeholders (`""`/`0`): only `cid` +/// ever gets read back off a snapshot's cover (to build the cover-proxy +/// URL), matching `catalog_index_postgres`'s identical `cover_cid` fallback. +pub fn to_crate_entry(folded: FoldedEntry) -> CrateEntry { + CrateEntry( + entry_id: uri.rkey(folded.entry_uri), + status: crate.status_from_string(folded.status), + snapshot: snapshot_from_folded(folded), + media_grade: folded.media_grade, + sleeve_grade: folded.sleeve_grade, + rating: folded.rating, + folder: folded.folder, + notes: folded.notes, + release: release_from_folded(folded), + price: price_from_folded(folded), + counterparty: folded.counterparty, + source: source_from_folded(folded), + created_at: folded.created_at, + updated_at: folded.updated_at, + events: [], + ) +} + +/// A snapshot needs at least `title`+`artist_display` (the lexicon's +/// required fields); an entry that somehow folded without either has no +/// displayable snapshot, matching `folded_page`'s existing +/// `e.snapshot != None` filter that drops such rows before they reach a +/// response anyway. +fn snapshot_from_folded(folded: FoldedEntry) -> Option(Snapshot) { + case folded.title, folded.artist_display { + Some(title), Some(artist_display) -> + Some(Snapshot( + title:, + artist_display:, + year: folded.year, + format: folded.format, + thumb_url: folded.thumb_url, + cover: folded.cover_cid + |> option.map(fn(cid) { Blob(cid:, mime_type: "", size: 0) }), + )) + _, _ -> None + } +} + +fn release_from_folded(folded: FoldedEntry) -> Option(defs.CatalogRef) { + case folded.release_uri, folded.release_cid { + Some(release_uri), Some(release_cid) -> + // external_ids is always written None at genesis/repoint time (see + // promotion.gleam/catalog_entities.gleam/edit_inbox.gleam), so + // omitting it here round-trips every write path unchanged. + Some(CatalogRef(uri: release_uri, cid: release_cid, external_ids: None)) + _, _ -> None + } +} + +fn price_from_folded(folded: FoldedEntry) -> Option(Price) { + case folded.price_amount, folded.price_currency { + Some(amount), Some(currency) -> Some(Price(amount:, currency:)) + _, _ -> None + } +} + +fn external_id_from_folded(folded: FoldedEntry) -> Option(ExternalId) { + case folded.source_external_provider, folded.source_external_id { + Some(provider), Some(id) -> + Some(ExternalId(id:, provider:, url: folded.source_external_url)) + _, _ -> None + } +} + +fn record_ref_from_folded(folded: FoldedEntry) -> Option(RepoStrongRef) { + case folded.source_record_uri, folded.source_record_cid { + Some(record_uri), Some(record_cid) -> + Some(RepoStrongRef(uri: record_uri, cid: record_cid)) + _, _ -> None + } +} + +/// `None` only when every sub-field is absent, so a never-had-a-source entry +/// round-trips to an omitted `source` field rather than an empty `{}` +/// object (`crate.opt` omits on `None`, not on an all-absent record). +fn source_from_folded(folded: FoldedEntry) -> Option(Source) { + let external = external_id_from_folded(folded) + let record = record_ref_from_folded(folded) + case + folded.source_client_agent, + folded.source_origin, + folded.source_origin_url, + external, + record + { + None, None, None, None, None -> None + client_agent, origin, origin_url, external, record -> + Some(Source(client_agent:, external:, origin:, origin_url:, record:)) + } +} + type State { State( events: Dict(String, ShelfEntryEvent), diff --git a/server/test/actor_shelf_test.gleam b/server/test/actor_shelf_test.gleam index 4e1280a..85ee59c 100644 --- a/server/test/actor_shelf_test.gleam +++ b/server/test/actor_shelf_test.gleam @@ -1,23 +1,34 @@ //// Handler tests for the unified shelf endpoints in their unauthenticated -//// actor mode (`?actor=`), faking the resolver and the actor's own PDS via a -//// host-branching `xrpc.Client`. The public-crate view is `view=current`. -//// Also pins the boundary of that mode: with no session and no `actor`, -//// both endpoints mean "my own crate" and demand auth. - -import at_record/gen/catalog/release as catalog_release -import at_record_server/catalog_deps.{type Deps, Deps} -import at_record_server/context.{type Context} +//// actor mode (`?actor=`). Post-C3 (appview-first roadmap), actor reads come +//// from `ctx.shelf_index`, live-fetched from the actor's PDS exactly once +//// per never-seen did. The functional tests here pre-mark the actor seen +//// and seed the index directly (`seeded_context`/`seed`), so they exercise +//// the read/response shaping in isolation from the seeding path; the +//// seeding path itself (one live fetch on success, or an unchanged +//// 502/404/413 propagated on failure with the did left unseen so the next +//// visit retries) gets its own dedicated tests further down. Also pins the +//// boundary of actor mode: with no session and no `actor`, both endpoints +//// mean "my own crate" and demand auth. + +import at_record/gen/defs.{CatalogRef, Snapshot} +import at_record/gen/repo/strong_ref.{RepoStrongRef} +import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} +import at_record_server/catalog/row.{type BrowseRow, BrowseRow} +import at_record_server/catalog_deps.{type Deps} +import at_record_server/catalog_index +import at_record_server/context.{type Context, Context} import at_record_server/handlers/shelf import atproto/xrpc import atproto_core/xrpc as core_xrpc import gleam/bit_array import gleam/dict import gleam/dynamic/decode +import gleam/erlang/process import gleam/http import gleam/http/response import gleam/json import gleam/list -import gleam/option.{Some} +import gleam/option.{None, Some} import support import wisp/simulate @@ -38,79 +49,6 @@ fn entry_uri(rkey: String) -> String { const release_uri = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" -fn record(rkey: String, fields: List(#(String, json.Json))) -> json.Json { - json.object([ - #("uri", json.string(entry_uri(rkey))), - #("cid", json.string("bafy" <> rkey)), - #("value", json.object(fields)), - ]) -} - -fn snapshot(title: String, artist: String) -> #(String, json.Json) { - #( - "snapshot", - json.object([ - #("title", json.string(title)), - #("artistDisplay", json.string(artist)), - ]), - ) -} - -fn owned_genesis(rkey: String) -> json.Json { - record(rkey, [ - #("action", json.string("acquired")), - #("createdAt", json.string("2026-01-01T00:00:00Z")), - snapshot("Spiderland", "Slint"), - #( - "release", - json.object([ - #("cid", json.string("bafyrel")), - #("uri", json.string(release_uri)), - ]), - ), - ]) -} - -fn wanted_genesis(rkey: String) -> json.Json { - record(rkey, [ - #("action", json.string("wanted")), - #("createdAt", json.string("2026-01-02T00:00:00Z")), - snapshot("Loveless", "My Bloody Valentine"), - ]) -} - -/// A genesis plus a later `sold` event referencing it, so folding leaves it -/// `Gone`: the `current` view must exclude it from the response. -fn sold_pair( - genesis_rkey: String, - event_rkey: String, -) -> #(json.Json, json.Json) { - let genesis = - record(genesis_rkey, [ - #("action", json.string("acquired")), - #("createdAt", json.string("2026-01-03T00:00:00Z")), - snapshot("Isn't Anything", "My Bloody Valentine"), - ]) - let sold = - record(event_rkey, [ - #("action", json.string("sold")), - #("createdAt", json.string("2026-01-04T00:00:00Z")), - #( - "subject", - json.object([ - #("cid", json.string("bafy" <> genesis_rkey)), - #("uri", json.string(entry_uri(genesis_rkey))), - ]), - ), - ]) - #(genesis, sold) -} - -fn list_records_body(records: List(json.Json)) -> String { - json.object([#("records", json.preprocessed_array(records))]) - |> json.to_string -} - fn resolve_body() -> String { json.object([ #("did", json.string(actor_did)), @@ -121,14 +59,15 @@ fn resolve_body() -> String { |> json.to_string } -fn network_client(records_body: String) -> xrpc.Client { +/// Answers Slingshot identity resolution only; panics on any other host, so +/// a test using this client proves the live-fetch-and-seed path was never +/// reached (the actor was pre-marked seen and seeded directly instead). +fn resolver_only_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case req.host { "resolver.test" -> Ok(response.Response(200, [], bit_array.from_string(resolve_body()))) - host if host == pds_host -> - Ok(response.Response(200, [], bit_array.from_string(records_body))) - _ -> panic as "unexpected host" + _ -> panic as "actor PDS must not be called once the did is seeded" } }) } @@ -150,8 +89,7 @@ fn pds_unreachable_client() -> xrpc.Client { } // The foreign PDS is reachable but says the repo is missing, distinct from -// `pds_unreachable_client`'s transport failure: this should read as 404, -// not 502. +// `pds_unreachable_client`'s transport failure. fn pds_repo_not_found_client() -> xrpc.Client { xrpc.Client(send: fn(req) { case req.host { @@ -164,22 +102,6 @@ fn pds_repo_not_found_client() -> xrpc.Client { }) } -fn a_release() -> catalog_release.CatalogRelease { - catalog_release.CatalogRelease( - ..support.blank_catalog_release(), - title: "Spiderland", - genres: Some(["Rock"]), - country: Some("US"), - released: Some("1991"), - ) -} - -fn catalog_with_release() -> Deps { - Deps(..support.unreachable_catalog_deps(), fetch_release: fn(_uri) { - Some(#("bafyrel", a_release())) - }) -} - fn test_context(client: xrpc.Client, catalog: Deps) -> Context { support.stub_context_with( support.stub_config(client), @@ -189,17 +111,132 @@ fn test_context(client: xrpc.Client, catalog: Deps) -> Context { ) } +/// The actor pre-marked seen in the index, so `ensure_actor_seeded` never +/// reaches for the live-fetch path -- the standard setup for tests that +/// only care about read/response shaping, not the seeding path itself. +fn seeded_context(client: xrpc.Client, catalog: Deps) -> Context { + let ctx = test_context(client, catalog) + ctx.shelf_index.seen.mark_seen(actor_did) + ctx +} + +fn seed(ctx: Context, rkey: String, e: ShelfEntry) -> Nil { + support.seed_shelf_entry(ctx, actor_did, entry_uri(rkey), rkey, e) +} + +fn owned_genesis() -> ShelfEntry { + ShelfEntry( + ..support.blank_shelf_entry(), + action: "acquired", + created_at: "2026-01-01T00:00:00Z", + snapshot: Some(Snapshot( + artist_display: "Slint", + cover: None, + format: None, + thumb_url: None, + title: "Spiderland", + year: None, + )), + release: Some(CatalogRef( + cid: "bafyrel", + uri: release_uri, + external_ids: None, + )), + ) +} + +fn wanted_genesis() -> ShelfEntry { + ShelfEntry( + ..support.blank_shelf_entry(), + action: "wanted", + created_at: "2026-01-02T00:00:00Z", + snapshot: Some(Snapshot( + artist_display: "My Bloody Valentine", + cover: None, + format: None, + thumb_url: None, + title: "Loveless", + year: None, + )), + ) +} + +/// A genesis plus a later `sold` event referencing it, so folding leaves it +/// `Gone`: the `current` view must exclude it from the response. +fn seed_sold_pair( + ctx: Context, + genesis_rkey: String, + event_rkey: String, +) -> Nil { + seed( + ctx, + genesis_rkey, + ShelfEntry( + ..support.blank_shelf_entry(), + action: "acquired", + created_at: "2026-01-03T00:00:00Z", + snapshot: Some(Snapshot( + artist_display: "My Bloody Valentine", + cover: None, + format: None, + thumb_url: None, + title: "Isn't Anything", + year: None, + )), + ), + ) + seed( + ctx, + event_rkey, + ShelfEntry( + ..support.blank_shelf_entry(), + action: "sold", + created_at: "2026-01-04T00:00:00Z", + subject: Some(RepoStrongRef( + cid: "bafy" <> genesis_rkey, + uri: entry_uri(genesis_rkey), + )), + ), + ) +} + +fn a_release() -> BrowseRow { + BrowseRow( + uri: release_uri, + cid: "bafyrel", + title: "Spiderland", + artist_display: None, + genres: ["Rock"], + styles: [], + released: Some("1991"), + country: Some("US"), + cover: None, + thumb_url: None, + discogs_id: None, + created_at: "2026-01-01T00:00:00Z", + publisher_did: actor_did, + publisher_handle: actor_handle, + publisher_pds: "https://" <> pds_host, + supersedes: None, + based_on: None, + format: None, + label: None, + master: None, + ) +} + +fn catalog_index_with_release() -> catalog_index.Store { + let assert Ok(store) = catalog_index.start() + store.releases.upsert(a_release()) + store +} + pub fn actor_shelf_folds_and_excludes_gone_entries_test() { - let #(sold_genesis, sold_event) = sold_pair("3ccc", "3ccd") - let records = - list_records_body([ - owned_genesis("3aaa"), - wanted_genesis("3bbb"), - sold_genesis, - sold_event, - ]) let ctx = - test_context(network_client(records), support.unreachable_catalog_deps()) + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + seed(ctx, "3aaa", owned_genesis()) + seed(ctx, "3bbb", wanted_genesis()) + seed_sold_pair(ctx, "3ccc", "3ccd") let resp = shelf.list_shelf( simulate.request( @@ -219,32 +256,31 @@ pub fn actor_shelf_folds_and_excludes_gone_entries_test() { assert titles == ["Spiderland", "Loveless"] } -// A foreign-DID release would normally send `via_handles` off to the -// identity resolver; `network_client` panics on any host besides the -// resolver and the actor's own PDS, so this would fail loudly if the actor -// path still fanned out for attribution on someone else's crate. pub fn actor_shelf_skips_via_handles_lookup_for_foreign_release_test() { - let foreign_release = - record("3fff", [ - #("action", json.string("acquired")), - #("createdAt", json.string("2026-01-05T00:00:00Z")), - snapshot("Loveless", "My Bloody Valentine"), - #( - "release", - json.object([ - #("cid", json.string("bafyforeign")), - #( - "uri", - json.string( - "at://did:plc:other/dev.mokkenstorm.crate.catalog.release/r2", - ), - ), - ]), - ), - ]) - let records = list_records_body([foreign_release]) let ctx = - test_context(network_client(records), support.unreachable_catalog_deps()) + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + seed( + ctx, + "3fff", + ShelfEntry( + ..support.blank_shelf_entry(), + action: "acquired", + created_at: "2026-01-05T00:00:00Z", + snapshot: Some(Snapshot( + artist_display: "My Bloody Valentine", + cover: None, + format: None, + thumb_url: None, + title: "Loveless", + year: None, + )), + release: Some(CatalogRef( + cid: "bafyforeign", + uri: "at://did:plc:other/dev.mokkenstorm.crate.catalog.release/r2", + external_ids: None, + )), + ), + ) let resp = shelf.list_shelf( simulate.request( @@ -264,8 +300,10 @@ pub fn actor_shelf_skips_via_handles_lookup_for_foreign_release_test() { } pub fn actor_entry_returns_folded_entry_and_release_test() { - let records = list_records_body([owned_genesis("3aaa")]) - let ctx = test_context(network_client(records), catalog_with_release()) + let ctx = + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + let ctx = Context(..ctx, catalog_index: catalog_index_with_release()) + seed(ctx, "3aaa", owned_genesis()) let resp = shelf.get_entry( simulate.request( @@ -283,9 +321,94 @@ pub fn actor_entry_returns_folded_entry_and_release_test() { assert genres == ["Rock"] } +pub fn actor_entry_for_a_missing_release_omits_it_test() { + let ctx = + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + seed(ctx, "3aaa", owned_genesis()) + let resp = + shelf.get_entry( + simulate.request( + http.Get, + support.xrpc("shelf.getEntry?actor=pub.test&entry=3aaa"), + ), + ctx, + ) + assert resp.status == 200 + assert support.field_present(simulate.read_body(resp), ["release"]) == False +} + +pub fn actor_entry_for_unknown_entry_id_is_404_test() { + let ctx = + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + seed(ctx, "3aaa", owned_genesis()) + let resp = + shelf.get_entry( + simulate.request( + http.Get, + support.xrpc("shelf.getEntry?actor=pub.test&entry=missing"), + ), + ctx, + ) + assert resp.status == 404 +} + +pub fn actor_entry_for_a_repo_the_live_fetch_could_not_find_is_404_test() { + let ctx = + test_context( + pds_repo_not_found_client(), + support.unreachable_catalog_deps(), + ) + let resp = + shelf.get_entry( + simulate.request( + http.Get, + support.xrpc("shelf.getEntry?actor=pub.test&entry=3aaa"), + ), + ctx, + ) + assert resp.status == 404 +} + +pub fn actor_view_history_is_clamped_to_current_test() { + let ctx = + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + seed_sold_pair(ctx, "3eee", "3eef") + let current_resp = + shelf.list_shelf( + simulate.request( + http.Get, + support.xrpc("shelf.listEntries?actor=pub.test&view=current"), + ), + ctx, + ) + let assert Ok(current_ids) = + support.field_nested(simulate.read_body(current_resp), ["items"], [ + "entryId", + ]) + assert current_ids == [] + + let history_resp = + shelf.list_shelf( + simulate.request( + http.Get, + support.xrpc("shelf.listEntries?actor=pub.test&view=history"), + ), + ctx, + ) + let assert Ok(history_ids) = + support.field_nested(simulate.read_body(history_resp), ["items"], [ + "entryId", + ]) + assert history_ids == [] +} + +// A resolver-side failure (the actor doesn't resolve at all) is a genuine +// 404 regardless of the index. A never-seen actor's live-fetch failure +// propagates unchanged (502/404), the same as the always-live path did +// before C3. Missing the `entry` query param is a 400 before actor +// resolution is ever reached, so it needs no seeding either. pub fn error_status_test() { let unreachable = support.unreachable_catalog_deps() - let known_entry = list_records_body([owned_genesis("3aaa")]) [ #( shelf.list_shelf, @@ -305,18 +428,6 @@ pub fn error_status_test() { support.xrpc("shelf.listEntries?actor=pub.test"), 404, ), - #( - shelf.get_entry, - network_client(known_entry), - support.xrpc("shelf.getEntry?actor=pub.test&entry=missing"), - 404, - ), - #( - shelf.get_entry, - pds_repo_not_found_client(), - support.xrpc("shelf.getEntry?actor=pub.test&entry=3aaa"), - 404, - ), #( shelf.get_entry, failing_resolver_client(), @@ -359,12 +470,92 @@ pub fn get_entry_with_no_session_and_no_actor_is_unauthorized_test() { assert resp.status == 401 } -pub fn actor_view_history_is_clamped_to_current_test() { - let #(sold_genesis, sold_event) = sold_pair("3eee", "3eef") - let records = list_records_body([sold_genesis, sold_event]) +// The seeding path itself (C3 of the appview-first roadmap): a never-seen +// actor triggers exactly one live fetch, and on success every later read is +// served from the index without fetching again. A failed fetch is covered +// separately above (`error_status_test`, `actor_live_fetch_*_leaves_the_did_unseen_test`). + +fn record(rkey: String, fields: List(#(String, json.Json))) -> json.Json { + json.object([ + #("uri", json.string(entry_uri(rkey))), + #("cid", json.string("bafy" <> rkey)), + #("value", json.object(fields)), + ]) +} + +fn owned_genesis_record(rkey: String) -> json.Json { + record(rkey, [ + #("action", json.string("acquired")), + #("createdAt", json.string("2026-01-01T00:00:00Z")), + #( + "snapshot", + json.object([ + #("title", json.string("Spiderland")), + #("artistDisplay", json.string("Slint")), + ]), + ), + ]) +} + +fn list_records_body(records: List(json.Json)) -> String { + json.object([#("records", json.preprocessed_array(records))]) + |> json.to_string +} + +/// Answers Slingshot identity resolution unconditionally, and the actor's +/// PDS with `records_body`, counting only the PDS calls on `counter` +/// (identity resolution goes through the same "resolver.test" host on every +/// call regardless of the index's seen state, so it is not the signal here). +fn counting_pds_client( + counter: process.Subject(Nil), + records_body: String, +) -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.host { + "resolver.test" -> + Ok(response.Response(200, [], bit_array.from_string(resolve_body()))) + host if host == pds_host -> { + process.send(counter, Nil) + Ok(response.Response(200, [], bit_array.from_string(records_body))) + } + _ -> panic as "unexpected host" + } + }) +} + +pub fn never_seen_actor_fetches_once_then_serves_from_index_test() { + let calls = process.new_subject() + let client = + counting_pds_client( + calls, + list_records_body([owned_genesis_record("3aaa")]), + ) + let ctx = test_context(client, support.unreachable_catalog_deps()) + let path = support.xrpc("shelf.listEntries?actor=pub.test&view=current") + + let first = shelf.list_shelf(simulate.request(http.Get, path), ctx) + assert first.status == 200 + let assert Ok(first_titles) = + support.field_nested(simulate.read_body(first), ["items"], [ + "snapshot", "title", + ]) + assert first_titles == ["Spiderland"] + assert support.drain_count(calls) == 1 + + let second = shelf.list_shelf(simulate.request(http.Get, path), ctx) + assert second.status == 200 + let assert Ok(second_titles) = + support.field_nested(simulate.read_body(second), ["items"], [ + "snapshot", "title", + ]) + assert second_titles == ["Spiderland"] + assert support.drain_count(calls) == 0 +} + +pub fn seen_but_empty_actor_crate_serves_empty_without_a_live_call_test() { let ctx = - test_context(network_client(records), support.unreachable_catalog_deps()) - let current_resp = + seeded_context(resolver_only_client(), support.unreachable_catalog_deps()) + let resp = shelf.list_shelf( simulate.request( http.Get, @@ -372,23 +563,49 @@ pub fn actor_view_history_is_clamped_to_current_test() { ), ctx, ) - let assert Ok(current_ids) = - support.field_nested(simulate.read_body(current_resp), ["items"], [ - "entryId", - ]) - assert current_ids == [] + assert resp.status == 200 + let assert Ok(items) = + support.field_nested(simulate.read_body(resp), ["items"], ["entryId"]) + assert items == [] +} - let history_resp = +// A live-fetch failure on a never-seen actor must propagate (never degrade +// to an empty 200) and must leave the did unseen, so a stranger whose PDS +// was briefly down isn't cached as permanently empty -- reconcile only +// re-backfills *known* users, never a visited-but-unreachable stranger. The +// status codes themselves are covered by `error_status_test`; these assert +// the `has_seen` side specifically. +pub fn actor_live_fetch_transport_failure_leaves_the_did_unseen_test() { + let ctx = + test_context(pds_unreachable_client(), support.unreachable_catalog_deps()) + let resp = shelf.list_shelf( simulate.request( http.Get, - support.xrpc("shelf.listEntries?actor=pub.test&view=history"), + support.xrpc("shelf.listEntries?actor=pub.test"), ), ctx, ) - let assert Ok(history_ids) = - support.field_nested(simulate.read_body(history_resp), ["items"], [ - "entryId", - ]) - assert history_ids == [] + assert resp.status == 502 + assert ctx.shelf_index.seen.has_seen(actor_did) == False + assert ctx.shelf_index.entries.list_for_did(actor_did) == [] +} + +pub fn actor_live_fetch_repo_not_found_leaves_the_did_unseen_test() { + let ctx = + test_context( + pds_repo_not_found_client(), + support.unreachable_catalog_deps(), + ) + let resp = + shelf.list_shelf( + simulate.request( + http.Get, + support.xrpc("shelf.listEntries?actor=pub.test"), + ), + ctx, + ) + assert resp.status == 404 + assert ctx.shelf_index.seen.has_seen(actor_did) == False + assert ctx.shelf_index.entries.list_for_did(actor_did) == [] } diff --git a/server/test/crate_test.gleam b/server/test/crate_test.gleam index 41d9581..1d9c7d8 100644 --- a/server/test/crate_test.gleam +++ b/server/test/crate_test.gleam @@ -3,6 +3,7 @@ import at_record/gen/defs import at_record/gen/repo/strong_ref.{RepoStrongRef} import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} import at_record/storage.{type StoredItem, StoredItem} +import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/crate import gleam/json import gleam/list @@ -144,6 +145,67 @@ pub fn encode_release_info_omits_absent_fields_test() { assert support.field_present(body, ["released"]) == False } +fn blank_row() -> BrowseRow { + BrowseRow( + uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + cid: "bafyrow", + title: "Spiderland", + artist_display: None, + genres: [], + styles: [], + released: None, + country: None, + cover: None, + thumb_url: None, + discogs_id: None, + created_at: "2024-01-01T00:00:00Z", + publisher_did: "did:plc:pub", + publisher_handle: "pub.test", + publisher_pds: "https://pds.test", + supersedes: None, + based_on: None, + format: None, + label: None, + master: None, + ) +} + +// C5 of the appview-first roadmap: `encode_release_info_row` (the +// catalog-index read) and `encode_release_info` (the old live-fetch read) +// must render byte-identical JSON for equivalent display fields, so the +// entry-timeline's `release` field never changes shape under the SPA. +pub fn encode_release_info_row_matches_the_live_fetch_shape_test() { + let live = + catalog_release.CatalogRelease( + ..blank_release(), + artist_display: Some("Slint"), + genres: Some(["Rock"]), + styles: Some(["Math Rock"]), + country: Some("US"), + released: Some("1991-03-27"), + ) + let indexed = + BrowseRow( + ..blank_row(), + artist_display: Some("Slint"), + genres: ["Rock"], + styles: ["Math Rock"], + country: Some("US"), + released: Some("1991-03-27"), + ) + assert crate.encode_release_info(live) |> json.to_string + == { crate.encode_release_info_row(indexed) |> json.to_string } +} + +pub fn encode_release_info_row_omits_absent_fields_test() { + let body = crate.encode_release_info_row(blank_row()) |> json.to_string + assert support.field_present(body, ["artistDisplay"]) == False + assert support.field_present(body, ["genres"]) == False + assert support.field_present(body, ["styles"]) == False + assert support.field_present(body, ["country"]) == False + assert support.field_present(body, ["released"]) == False +} + fn blank_entry() -> crate.CrateEntry { crate.CrateEntry( entry_id: "3a", diff --git a/server/test/shelf_index_test.gleam b/server/test/shelf_index_test.gleam index dc0f473..8c10123 100644 --- a/server/test/shelf_index_test.gleam +++ b/server/test/shelf_index_test.gleam @@ -1,12 +1,20 @@ //// 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`. +//// arrival, mid-history deletes, purge-to-zero), `seen`/`delete_for_did`, +//// and `to_crate_entry`'s wire-shape parity with the pre-C3 live-fold path +//// (C3 of the appview-first roadmap). +import at_record/gen/defs.{CatalogRef, ExternalId, Price, Snapshot, Source} import at_record/gen/repo/strong_ref.{type RepoStrongRef, RepoStrongRef} import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} +import at_record/storage.{StoredItem} +import at_record_server/crate import at_record_server/shelf_index +import atproto/blob +import gleam/json import gleam/list import gleam/option.{None, Some} +import support /// A genesis event: no `subject`, so `event_uri == entry_uri`. fn genesis( @@ -320,3 +328,184 @@ pub fn delete_for_did_drops_entries_events_and_seen_test() { assert store.seen.has_seen("did:plc:a") == False assert store.entries.get(entry_b) != None } + +/// A genesis carrying every optional field (snapshot with a cover, release, +/// price, counterparty, and a full source with both an external id and a +/// record ref), plus an append that changes the media grade. +fn full_genesis() -> ShelfEntry { + ShelfEntry( + subject: None, + action: "acquired", + snapshot: Some(Snapshot( + title: "Spiderland", + artist_display: "Slint", + year: Some(1991), + format: Some("LP"), + thumb_url: Some("https://example.test/thumb.jpg"), + cover: Some(blob.Blob( + cid: "bafycover", + mime_type: "image/jpeg", + size: 1234, + )), + )), + external_ids: None, + media_grade: Some("VG+"), + sleeve_grade: Some("NM"), + rating: Some(5), + folder: Some("Rock A-M"), + notes: Some("first pressing"), + release: Some(CatalogRef( + uri: "at://did:plc:a/dev.mokkenstorm.crate.catalog.release/r1", + cid: "bafyrel", + external_ids: None, + )), + price: Some(Price(amount: 4000, currency: "EUR")), + counterparty: Some("Record shop"), + source: Some(Source( + client_agent: Some("at-record/1.0"), + external: Some(ExternalId(id: "42", provider: "discogs", url: None)), + origin: Some("manual"), + origin_url: None, + record: Some(RepoStrongRef( + uri: "at://did:plc:a/dev.mokkenstorm.crate.catalog.release/r1", + cid: "bafyrel", + )), + )), + created_at: "2024-01-01T00:00:00Z", + ) +} + +fn full_append(entry_uri: String) -> ShelfEntry { + ShelfEntry( + subject: Some(RepoStrongRef(cid: "bafygenesis", uri: entry_uri)), + action: "regraded", + snapshot: None, + external_ids: None, + media_grade: Some("NM"), + sleeve_grade: None, + rating: None, + folder: None, + notes: None, + release: None, + price: None, + counterparty: None, + source: Some(Source( + client_agent: None, + external: None, + origin: Some("manual"), + origin_url: None, + record: None, + )), + created_at: "2024-02-01T00:00:00Z", + ) +} + +/// `to_crate_entry` (the C3 index-read path) must render the exact same +/// `crate.encode_entry` JSON as `crate.fold` (the pre-C3 live-PDS-read +/// path) did for the same underlying events -- the wire-shape guarantee the +/// SPA depends on. +pub fn to_crate_entry_matches_the_live_fold_wire_shape_test() { + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + let genesis = full_genesis() + let append = full_append(entry_uri) + + let live_folded = + crate.fold([ + StoredItem(uri: entry_uri, cid: "bafygenesis", rkey: "e1", value: genesis), + StoredItem( + uri: entry_uri <> "-e2", + cid: "bafyappend", + rkey: "e2", + value: append, + ), + ]) + let assert [live] = live_folded + + let assert Ok(store) = shelf_index.start() + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did: "did:plc:a", + rkey: "e1", + entry: genesis, + ), + ) + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri <> "-e2", + entry_uri:, + did: "did:plc:a", + rkey: "e2", + entry: append, + ), + ) + let assert Some(folded) = store.entries.get(entry_uri) + let indexed = shelf_index.to_crate_entry(folded) + + // The one accepted deviation: `FoldedEntry` only stores a cover's `cid` + // (matching `catalog_index_postgres`'s identical `cover_cid` fallback), so + // `to_crate_entry` reconstructs a placeholder mime_type/size. Nothing + // downstream ever reads them (only `cid`, to build the cover-proxy URL), + // so the live-fold fixture is normalized the same way before comparing. + let normalize_cover = fn(e: crate.CrateEntry) -> crate.CrateEntry { + crate.CrateEntry( + ..e, + snapshot: option.map(e.snapshot, fn(s) { + defs.Snapshot( + ..s, + cover: option.map(s.cover, fn(b) { + blob.Blob(..b, mime_type: "", size: 0) + }), + ) + }), + ) + } + + assert crate.encode_entry(normalize_cover(live)) |> json.to_string + == { crate.encode_entry(indexed) |> json.to_string } +} + +/// A genesis with no source at all round-trips to an omitted `source` +/// field, not an empty `{}` object (`to_crate_entry`'s `source_from_folded` +/// only synthesizes `None` when every sub-field is absent). +pub fn to_crate_entry_omits_source_when_never_set_test() { + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e2" + let assert Ok(store) = shelf_index.start() + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did: "did:plc:a", + rkey: "e2", + entry: ShelfEntry(..full_genesis(), source: None), + ), + ) + let assert Some(folded) = store.entries.get(entry_uri) + let body = + crate.encode_entry(shelf_index.to_crate_entry(folded)) |> json.to_string + assert support.field_present(body, ["source"]) == False +} + +/// `entry_id` on the reconstructed `CrateEntry` is the rkey, not the full +/// at-uri -- the same shape `crate.fold_group` produces off a genesis +/// `StoredItem`. +pub fn to_crate_entry_entry_id_is_the_rkey_test() { + let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e3" + let assert Ok(store) = shelf_index.start() + shelf_index.record_and_fold( + store, + shelf_index.event_row( + event_uri: entry_uri, + entry_uri:, + did: "did:plc:a", + rkey: "e3", + entry: full_genesis(), + ), + ) + let assert Some(folded) = store.entries.get(entry_uri) + assert shelf_index.to_crate_entry(folded).entry_id == "e3" +} diff --git a/server/test/shelf_list_test.gleam b/server/test/shelf_list_test.gleam index 4c8f0fc..10cc64a 100644 --- a/server/test/shelf_list_test.gleam +++ b/server/test/shelf_list_test.gleam @@ -1,22 +1,19 @@ //// End-to-end `shelf.listEntries` handler tests for cursor pagination. -//// Mirrors `discogs_scan_test`'s `wisp/simulate` + stub PDS client style. +//// Own-mode reads come from the index (C3 of the appview-first roadmap), so +//// entries are seeded directly into `ctx.shelf_index` rather than stubbed +//// off a fake PDS response; `support.unreachable_client()` proves the read +//// path never touches the PDS. import at_record/gen/defs.{Snapshot} -import at_record/gen/shelf/entry.{ - type ShelfEntry, ShelfEntry, encode_shelf_entry, -} +import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} import at_record_server/context.{type Context} import at_record_server/handlers/shelf as shelf_handler import at_record_server/oauth/config import at_record_server/oauth/session_store import at_record_server/oauth/sessions import at_record_server/pagination -import atproto/xrpc -import gleam/bit_array import gleam/http -import gleam/http/response import gleam/int -import gleam/json import gleam/list import gleam/option.{None, Some} import gleam/string @@ -38,8 +35,13 @@ fn a_session() -> sessions.OauthSession { ) } -fn test_context(client: xrpc.Client) -> #(Context, config.Config) { - let cfg = support.stub_config_with(client, "r", "http://localhost:8080") +fn test_context() -> #(Context, config.Config) { + let cfg = + support.stub_config_with( + support.unreachable_client(), + "r", + "http://localhost:8080", + ) let ctx = support.stub_context_with( cfg, @@ -70,31 +72,21 @@ fn stub_entry(action: String, n: Int) -> ShelfEntry { ) } -fn record_json(n: Int, action: String) -> json.Json { - let uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey(n) - json.object([ - #("uri", json.string(uri)), - #("cid", json.string("bafy" <> rkey(n))), - #("value", encode_shelf_entry(stub_entry(action, n))), - ]) -} - /// Entries 1..owned_count are "acquired" (the owned view), the rest up to -/// `total` are "wanted". -fn client_with_entries(owned_count: Int, total: Int) -> xrpc.Client { - let records = - list.repeat(Nil, total) - |> list.index_map(fn(_, i) { - let n = i + 1 - case n <= owned_count { - True -> record_json(n, "acquired") - False -> record_json(n, "wanted") - } - }) - let body = json.object([#("records", json.preprocessed_array(records))]) - xrpc.Client(send: fn(_req) { - Ok(response.Response(200, [], bit_array.from_string(json.to_string(body)))) +/// `total` are "wanted", folded straight into `ctx.shelf_index` as genesis +/// events. +fn seed_entries(ctx: Context, owned_count: Int, total: Int) -> Nil { + list.repeat(Nil, total) + |> list.index_map(fn(_, i) { + let n = i + 1 + let action = case n <= owned_count { + True -> "acquired" + False -> "wanted" + } + let uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey(n) + support.seed_shelf_entry(ctx, did, uri, rkey(n), stub_entry(action, n)) }) + Nil } fn list_shelf(query: String, ctx: Context, cfg: config.Config) -> String { @@ -115,7 +107,8 @@ fn item_ids(body: String) -> List(String) { } pub fn a_full_view_fetch_has_no_next_cursor_test() { - let #(ctx, cfg) = test_context(client_with_entries(5, 5)) + let #(ctx, cfg) = test_context() + seed_entries(ctx, 5, 5) let all = [rkey(1), rkey(2), rkey(3), rkey(4), rkey(5)] ["", "?limit=5"] |> list.each(fn(query) { @@ -126,7 +119,8 @@ pub fn a_full_view_fetch_has_no_next_cursor_test() { } pub fn a_cursor_past_the_last_entry_returns_an_empty_page_test() { - let #(ctx, cfg) = test_context(client_with_entries(5, 5)) + let #(ctx, cfg) = test_context() + seed_entries(ctx, 5, 5) let cursor = pagination.encode_cursor("owned", rkey(5)) let body = list_shelf("?limit=1&cursor=" <> cursor, ctx, cfg) assert item_ids(body) == [] @@ -134,7 +128,8 @@ pub fn a_cursor_past_the_last_entry_returns_an_empty_page_test() { } pub fn a_stale_or_garbage_cursor_restarts_from_the_top_test() { - let #(ctx, cfg) = test_context(client_with_entries(5, 5)) + let #(ctx, cfg) = test_context() + seed_entries(ctx, 5, 5) let unknown_id = pagination.encode_cursor("owned", "g999") ["not-a-real-cursor", unknown_id] |> list.each(fn(cursor) { @@ -146,7 +141,8 @@ pub fn a_stale_or_garbage_cursor_restarts_from_the_top_test() { // A cursor minted under one view must not resume mid-list when replayed // against another view; it restarts from that view's top instead. pub fn a_cursor_from_a_different_view_restarts_within_the_new_view_test() { - let #(ctx, cfg) = test_context(client_with_entries(3, 5)) + let #(ctx, cfg) = test_context() + seed_entries(ctx, 3, 5) let owned_first_page = list_shelf("?limit=2", ctx, cfg) let assert Ok(cursor) = support.field_string(owned_first_page, ["cursor"]) let wanted_body = diff --git a/server/test/shelf_via_handles_test.gleam b/server/test/shelf_via_handles_test.gleam index fd9958d..990e262 100644 --- a/server/test/shelf_via_handles_test.gleam +++ b/server/test/shelf_via_handles_test.gleam @@ -1,24 +1,21 @@ //// `viaHandles` attribution on the own-shelf `shelf.listEntries` endpoint: //// entries whose release was adopted from someone else's repo get that //// repo's handle resolved through the batched identity resolver, deduped -//// per unique foreign DID rather than once per entry. +//// per unique foreign DID rather than once per entry. Own-mode reads come +//// from the index (C3 of the appview-first roadmap), so entries are seeded +//// directly into `ctx.shelf_index` rather than stubbed off a fake PDS. import at_record/gen/defs.{CatalogRef, Snapshot} -import at_record/gen/shelf/entry.{ - type ShelfEntry, ShelfEntry, encode_shelf_entry, -} +import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} import at_record_server/context.{type Context, Atproto, Context} import at_record_server/handlers/shelf as shelf_handler import at_record_server/identity_cache.{Identity} import at_record_server/oauth/session_store import at_record_server/oauth/sessions.{type OauthSession} -import atproto/xrpc -import gleam/bit_array import gleam/dict import gleam/dynamic/decode import gleam/erlang/process import gleam/http -import gleam/http/response import gleam/int import gleam/json import gleam/list @@ -71,27 +68,26 @@ fn foreign_entry(n: Int) -> ShelfEntry { ) } -fn record_json(n: Int) -> json.Json { - let uri = - "at://" <> own_did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey(n) - json.object([ - #("uri", json.string(uri)), - #("cid", json.string("bafy" <> rkey(n))), - #("value", encode_shelf_entry(foreign_entry(n))), - ]) -} - -fn client_with_foreign_entries(count: Int) -> xrpc.Client { - let records = - list.repeat(Nil, count) |> list.index_map(fn(_, i) { record_json(i + 1) }) - let body = json.object([#("records", json.preprocessed_array(records))]) - xrpc.Client(send: fn(_req) { - Ok(response.Response(200, [], bit_array.from_string(json.to_string(body)))) +/// Seeds `count` foreign-release-adopted entries straight into +/// `ctx.shelf_index` as genesis events. +fn seed_foreign_entries(ctx: Context, count: Int) -> Nil { + list.repeat(Nil, count) + |> list.index_map(fn(_, i) { + let n = i + 1 + let uri = + "at://" <> own_did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey(n) + support.seed_shelf_entry(ctx, own_did, uri, rkey(n), foreign_entry(n)) }) + Nil } -fn test_context(client: xrpc.Client) -> Context { - let cfg = support.stub_config_with(client, "r", "http://localhost:8080") +fn test_context() -> Context { + let cfg = + support.stub_config_with( + support.unreachable_client(), + "r", + "http://localhost:8080", + ) support.stub_context_with( cfg, support.unreachable_catalog_deps(), @@ -121,7 +117,8 @@ fn via_handles(body: String) -> dict.Dict(String, String) { pub fn entries_sharing_one_foreign_did_resolve_with_one_fetch_test() { let count = 3 - let ctx = test_context(client_with_foreign_entries(count)) + let ctx = test_context() + seed_foreign_entries(ctx, count) let counter = process.new_subject() let identity = support.counting_identity_resolver(counter, fn(id) { @@ -140,7 +137,8 @@ pub fn entries_sharing_one_foreign_did_resolve_with_one_fetch_test() { } pub fn a_resolution_failure_omits_the_entry_from_via_handles_test() { - let ctx = test_context(client_with_foreign_entries(1)) + let ctx = test_context() + seed_foreign_entries(ctx, 1) let identity = support.counting_identity_resolver(process.new_subject(), fn(_id) { Error(Nil) diff --git a/server/test/shelf_write_through_test.gleam b/server/test/shelf_write_through_test.gleam new file mode 100644 index 0000000..c2569f7 --- /dev/null +++ b/server/test/shelf_write_through_test.gleam @@ -0,0 +1,389 @@ +//// Read-your-writes for the shelf write paths (C3 of the appview-first +//// roadmap): `write_event` (add/append), `do_purge`, and `write_repoint` +//// (amend) synchronously write-through into `ctx.shelf_index` after the PDS +//// write succeeds, so a subsequent index read in the *same* context sees +//// the change immediately, without waiting on the firehose. These are +//// write tests, so the PDS stub stays (matched by `req.path`, mirroring +//// `graph_handler_test.gleam`'s pattern); reads never touch it. + +import at_record/gen/defs +import at_record/gen/shelf/entry +import at_record_server/catalog/row.{BrowseRow} +import at_record_server/catalog_index +import at_record_server/context.{type Context} +import at_record_server/handlers/amend +import at_record_server/handlers/shelf as shelf_handler +import at_record_server/oauth/config +import at_record_server/oauth/session_store +import at_record_server/oauth/sessions +import atproto/xrpc +import gleam/bit_array +import gleam/http +import gleam/http/response +import gleam/json +import gleam/option.{None, Some} +import support +import wisp +import wisp/simulate + +const session_cookie = "ar_oauth_sid" + +const far_future = 9_999_999_999 + +const did = "did:plc:x" + +fn a_session() -> sessions.OauthSession { + support.stub_session_with( + access_token: "at", + refresh_token: "rt", + expires_at: far_future, + ) +} + +fn test_context(client: xrpc.Client) -> #(Context, config.Config) { + let cfg = support.stub_config_with(client, "r", "http://localhost:8080") + let ctx = + support.stub_context_with( + cfg, + support.unreachable_catalog_deps(), + fn(_req) { Error("unused") }, + [], + ) + #(ctx, cfg) +} + +fn authed_post( + path: String, + body: json.Json, + cfg: config.Config, +) -> wisp.Request { + let assert Ok(id) = session_store.create(cfg.sessions, a_session()) + simulate.request(http.Post, path) + |> simulate.json_body(body) + |> simulate.cookie(session_cookie, id, wisp.Signed) +} + +fn authed_get(path: String, cfg: config.Config) -> wisp.Request { + let assert Ok(id) = session_store.create(cfg.sessions, a_session()) + simulate.request(http.Get, path) + |> simulate.cookie(session_cookie, id, wisp.Signed) +} + +fn ok_response( + body: String, +) -> Result(response.Response(BitArray), xrpc.TransportError) { + Ok(response.Response(200, [], bit_array.from_string(body))) +} + +fn create_record_response(collection_rkey: String, cid: String) -> String { + json.object([ + #( + "uri", + json.string( + "at://" + <> did + <> "/dev.mokkenstorm.crate.shelf.entry/" + <> collection_rkey, + ), + ), + #("cid", json.string(cid)), + ]) + |> json.to_string +} + +pub fn add_then_list_shows_the_write_via_the_index_test() { + let client = + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.createRecord" -> + ok_response(create_record_response("e1", "bafygenesis")) + _ -> panic as { "unexpected path: " <> req.path } + } + }) + let #(ctx, cfg) = test_context(client) + + let add_resp = + shelf_handler.add_shelf_item( + authed_post( + support.xrpc("shelf.addEntry"), + json.object([ + #("title", json.string("Spiderland")), + #("artist", json.string("Slint")), + ]), + cfg, + ), + ctx, + ) + assert add_resp.status == 201 + let assert Ok(entry_id) = + support.field_string(simulate.read_body(add_resp), ["entryId"]) + assert entry_id == "e1" + + let list_resp = + shelf_handler.list_shelf( + authed_get(support.xrpc("shelf.listEntries"), cfg), + ctx, + ) + assert list_resp.status == 200 + let assert Ok(titles) = + support.field_nested(simulate.read_body(list_resp), ["items"], [ + "snapshot", "title", + ]) + assert titles == ["Spiderland"] +} + +fn genesis_record_json(rkey: String) -> json.Json { + json.object([ + #( + "uri", + json.string( + "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey, + ), + ), + #("cid", json.string("bafy" <> rkey)), + #( + "value", + json.object([ + #("action", json.string("acquired")), + #("createdAt", json.string("2026-01-01T00:00:00Z")), + #( + "snapshot", + json.object([ + #("title", json.string("Spiderland")), + #("artistDisplay", json.string("Slint")), + ]), + ), + ]), + ), + ]) +} + +pub fn append_then_get_entry_shows_the_event_via_the_index_test() { + let client = + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.getRecord" -> + ok_response(json.to_string(genesis_record_json("e1"))) + "/xrpc/com.atproto.repo.createRecord" -> + ok_response(create_record_response("e2", "bafyappend")) + _ -> panic as { "unexpected path: " <> req.path } + } + }) + let #(ctx, cfg) = test_context(client) + // The genesis is already indexed, as it would be by the time anyone + // appends to it (its own write-through at add time, or the firehose): + // `write_event`'s write-through only re-derives the folded row from every + // stored event, so an append alone can't fold into anything without the + // genesis also being present in the index. + support.seed_shelf_entry( + ctx, + did, + "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/e1", + "e1", + genesis_shelf_entry(), + ) + + let append_resp = + shelf_handler.append_event( + authed_post( + support.xrpc("shelf.appendEvent"), + json.object([ + #("entryId", json.string("e1")), + #("action", json.string("regraded")), + #("mediaGrade", json.string("NM")), + ]), + cfg, + ), + ctx, + ) + assert append_resp.status == 201 + + let entry_resp = + shelf_handler.get_entry( + authed_get(support.xrpc("shelf.getEntry?entry=e1"), cfg), + ctx, + ) + assert entry_resp.status == 200 + let body = simulate.read_body(entry_resp) + let assert Ok(media_grade) = + support.field_string(body, ["entry", "mediaGrade"]) + assert media_grade == "NM" + let assert Ok(actions) = support.field_nested(body, ["events"], ["action"]) + assert actions == ["acquired", "regraded"] +} + +pub fn purge_then_list_no_longer_shows_the_entry_test() { + let client = + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> + ok_response( + json.object([ + #("records", json.preprocessed_array([genesis_record_json("e1")])), + ]) + |> json.to_string, + ) + "/xrpc/com.atproto.repo.deleteRecord" -> ok_response("{}") + _ -> panic as { "unexpected path: " <> req.path } + } + }) + let #(ctx, cfg) = test_context(client) + // The genesis is already indexed (as it would be, either from its own + // write-through at add time or the firehose), so the "before" listEntries + // read reflects it. + support.seed_shelf_entry( + ctx, + did, + "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/e1", + "e1", + genesis_shelf_entry(), + ) + let before = + shelf_handler.list_shelf( + authed_get(support.xrpc("shelf.listEntries"), cfg), + ctx, + ) + let assert Ok(before_ids) = + support.field_nested(simulate.read_body(before), ["items"], ["entryId"]) + assert before_ids == ["e1"] + + let purge_resp = + shelf_handler.purge_entry( + authed_post( + support.xrpc("shelf.purgeEntry"), + json.object([#("entryId", json.string("e1"))]), + cfg, + ), + ctx, + ) + assert purge_resp.status == 200 + + let after = + shelf_handler.list_shelf( + authed_get(support.xrpc("shelf.listEntries"), cfg), + ctx, + ) + let assert Ok(after_ids) = + support.field_nested(simulate.read_body(after), ["items"], ["entryId"]) + assert after_ids == [] +} + +fn genesis_shelf_entry() -> entry.ShelfEntry { + entry.ShelfEntry( + ..support.blank_shelf_entry(), + action: "acquired", + snapshot: Some(defs.Snapshot( + title: "Spiderland", + artist_display: "Slint", + year: None, + format: None, + thumb_url: None, + cover: None, + )), + ) +} + +pub fn amend_then_get_entry_renders_the_new_release_test() { + let entry_uri = "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/e1" + let client = + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> + ok_response( + json.object([ + #("records", json.preprocessed_array([genesis_record_json("e1")])), + ]) + |> json.to_string, + ) + "/xrpc/com.atproto.repo.createRecord" -> { + let body = bit_array.to_string(req.body) |> option.from_result + let collection = + body + |> option.then(fn(b) { + support.field_string(b, ["collection"]) |> option.from_result + }) + case collection { + Some("dev.mokkenstorm.crate.catalog.release") -> + ok_response( + json.object([ + #( + "uri", + json.string( + "at://" + <> did + <> "/dev.mokkenstorm.crate.catalog.release/r2", + ), + ), + #("cid", json.string("bafyrel2")), + ]) + |> json.to_string, + ) + _ -> ok_response(create_record_response("e2", "bafyrepoint")) + } + } + _ -> panic as { "unexpected path: " <> req.path } + } + }) + let #(ctx, cfg) = test_context(client) + // Simulates the entry already having been indexed before this amend (via + // its own earlier write-through, or the firehose): `write_repoint`'s + // write-through only re-derives the folded row from every stored event, + // so a repoint alone can't fold into anything without the genesis also + // being present in the index. + support.seed_shelf_entry(ctx, did, entry_uri, "e1", genesis_shelf_entry()) + let assert Ok(store) = catalog_index.start() + store.releases.upsert(BrowseRow( + uri: "at://" <> did <> "/dev.mokkenstorm.crate.catalog.release/r2", + cid: "bafyrel2", + title: "Spiderland", + artist_display: Some("Slint"), + genres: ["Math Rock"], + styles: [], + released: Some("1991"), + country: None, + cover: None, + thumb_url: None, + discogs_id: None, + created_at: "2026-01-05T00:00:00Z", + publisher_did: did, + publisher_handle: "me.test", + publisher_pds: "https://pds.test", + supersedes: None, + based_on: None, + format: None, + label: None, + master: None, + )) + let ctx = context.Context(..ctx, catalog_index: store) + + let amend_resp = + amend.amend_entry( + authed_post( + support.xrpc("shelf.amendEntry"), + json.object([ + #("entryId", json.string("e1")), + #( + "fields", + json.object([#("genres", json.array(["Math Rock"], json.string))]), + ), + ]), + cfg, + ), + ctx, + ) + assert amend_resp.status == 201 + + let entry_resp = + shelf_handler.get_entry( + authed_get(support.xrpc("shelf.getEntry?entry=e1"), cfg), + ctx, + ) + assert entry_resp.status == 200 + let body = simulate.read_body(entry_resp) + let assert Ok(release_uri) = + support.field_string(body, ["entry", "release", "uri"]) + assert release_uri + == "at://" <> did <> "/dev.mokkenstorm.crate.catalog.release/r2" + let assert Ok(genres) = support.field_strings(body, ["release", "genres"]) + assert genres == ["Math Rock"] +} diff --git a/server/test/support.gleam b/server/test/support.gleam index 98f376f..15c59dc 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -119,6 +119,7 @@ fn empty_catalog_index() -> catalog_index.Store { delete: fn(_) { Nil }, delete_for_did: fn(_) { Nil }, list: fn() { [] }, + get: fn(_) { None }, ), adoptions: catalog_index.AdoptionOps( upsert: fn(_) { Nil }, @@ -146,6 +147,35 @@ fn fresh_shelf_index() -> shelf_index.Store { store } +/// Folds one `ShelfEntry` straight into `ctx.shelf_index`, mirroring what +/// the jetstream consumer (or a write-through handler) would do for the +/// same record -- the standard way a C3-cutover read test seeds the index +/// it now reads from, instead of stubbing a fake PDS response. `uri` is the +/// record's own at-uri, `rkey` its own record key; a genesis's `entry.subject` +/// is `None`, an append's points back at the genesis. +pub fn seed_shelf_entry( + ctx: Context, + did: String, + uri: String, + rkey: String, + shelf_entry: entry.ShelfEntry, +) -> Nil { + let entry_uri = case shelf_entry.subject { + option.Some(ref) -> ref.uri + option.None -> uri + } + shelf_index.record_and_fold( + ctx.shelf_index, + shelf_index.event_row( + event_uri: uri, + entry_uri:, + did:, + rkey:, + entry: shelf_entry, + ), + ) +} + /// 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 {