//// 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 atproto/xrpc import crate/gen/defs import crate/gen/shelf/entry import crate_server/catalog/row.{BrowseRow} import crate_server/catalog_index import crate_server/context.{type Context} import crate_server/handlers/amend import crate_server/handlers/shelf as shelf_handler import crate_server/oauth/config import crate_server/oauth/session_store import crate_server/oauth/sessions 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, barcode: 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"] }