diff --git a/e2e/tests/edit-inbox.spec.ts b/e2e/tests/edit-inbox.spec.ts index d56c7e8..25e33a3 100644 --- a/e2e/tests/edit-inbox.spec.ts +++ b/e2e/tests/edit-inbox.spec.ts @@ -61,7 +61,9 @@ test("alice proposes a correction on bob's record; bob reviews and applies it", await alicePage.click("text=PUBLISH AMENDMENT"); await proposed; - // Constellation indexes the new catalog.edit backlink off the same firehose. + // Our own catalog index (not Constellation, since E2 of the appview-first + // roadmap cut the inbox over to it) picks up the new catalog.edit off the + // same firehose — still not instant. await alicePage.waitForTimeout(3000); await bobPage.goto("/inbox"); diff --git a/server/src/at_record_server/catalog_import.gleam b/server/src/at_record_server/catalog_import.gleam index b53ac18..0e22adf 100644 --- a/server/src/at_record_server/catalog_import.gleam +++ b/server/src/at_record_server/catalog_import.gleam @@ -604,6 +604,7 @@ fn write_genesis( let cover = best_cover(ctx, client, session, release, details) let outcome = promotion.adopt_or_mint( + ctx.catalog_index, ctx.catalog, client, session, diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index f09efc0..07ae7fc 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -55,7 +55,12 @@ fn is_active_status(status: String) -> Bool { } /// One catalog.edit event: a proposed change to some catalog entity, keyed by -/// the edit record's own at-uri. +/// the edit record's own at-uri. `fields` is the proposal's `#releaseFields` +/// payload pre-encoded to JSON text (E2 of the appview-first roadmap): storing +/// it as a plain string, not a typed `catalog_edit` ref, keeps this leaf-ish +/// module free of a dependency on the lexicon module (see `catalog/row.gleam`'s +/// docstring for why that separation matters); ingest sites encode it, read +/// sites decode it via `catalog_edit.release_fields_decoder`. pub type Edit { Edit( edit_uri: String, @@ -63,6 +68,9 @@ pub type Edit { subject_uri: String, entity: String, created_at: String, + cid: String, + fields: Option(String), + rationale: Option(String), ) } @@ -105,6 +113,11 @@ pub type EditOps { delete: fn(String) -> Nil, delete_for_did: fn(String) -> Nil, for_subject: fn(String) -> List(Edit), + /// Every edit filed against any of the given subject uris, in one + /// batched read (E2 of the appview-first roadmap): the edit-inbox's + /// `list_proposals` calls this once per request instead of fanning out + /// per release. + list_for_subjects: fn(List(String)) -> List(Edit), ) } @@ -147,6 +160,7 @@ type Msg { AdoptionCounts(Subject(Dict(String, Int))) AdoptionsList(Subject(List(Adoption))) EditsFor(String, Subject(List(Edit))) + EditsForSubjects(List(String), Subject(List(Edit))) SaveCursor(Int) LoadCursor(Subject(Option(Int))) } @@ -249,6 +263,14 @@ pub fn start() -> Result(Store, actor.StartError) { Error(Nil) -> [] } }, + list_for_subjects: fn(subject_uris) { + case + parallel.try_call(subject, 1000, EditsForSubjects(subject_uris, _)) + { + Ok(edits) -> edits + Error(Nil) -> [] + } + }, ), cursor: CursorOps( save: fn(time_us) { @@ -349,6 +371,14 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { ) actor.continue(state) } + EditsForSubjects(subject_uris, reply) -> { + process.send( + reply, + dict.values(state.edits) + |> list.filter(fn(e) { list.contains(subject_uris, e.subject_uri) }), + ) + actor.continue(state) + } SaveCursor(time_us) -> actor.continue(State(..state, cursor: option.Some(time_us))) LoadCursor(reply) -> { diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam index 9129918..c3ca843 100644 --- a/server/src/at_record_server/catalog_index_postgres.gleam +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -49,6 +49,9 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { delete: fn(edit_uri) { delete_edit(conn, edit_uri) }, delete_for_did: fn(did) { delete_edits_for_did(conn, did) }, for_subject: fn(subject_uri) { edits_for(conn, subject_uri) }, + list_for_subjects: fn(subject_uris) { + edits_for_subjects(conn, subject_uris) + }, ), cursor: CursorOps( save: fn(time_us) { save_cursor(conn, time_us) }, @@ -154,6 +157,22 @@ fn migrate(conn: pog.Connection) -> Result(Nil, String) { |> pog.execute(conn) |> result.replace_error("catalog_edits migration failed"), ) + // Idempotent add-column for a table that already existed before + // `cid`/`fields`/`rationale` (E2 of the appview-first roadmap): mirrors the + // catalog_releases format/label/master pattern above. `fields` is JSON text, + // not a typed column, per `catalog_index.Edit`'s docstring. + use _ <- result.try( + pog.query( + "alter table catalog_edits + add column if not exists cid text, + add column if not exists fields text, + add column if not exists rationale text", + ) + |> pog.execute(conn) + |> result.replace_error( + "catalog_edits cid/fields/rationale migration failed", + ), + ) pog.query( "create table if not exists jetstream_cursor ( id text primary key, @@ -466,19 +485,26 @@ fn list_adoptions(conn: pog.Connection) -> List(Adoption) { fn upsert_edit(conn: pog.Connection, edit: Edit) -> Nil { let outcome = pog.query( - "insert into catalog_edits (edit_uri, did, subject_uri, entity, created_at) - values ($1, $2, $3, $4, $5) + "insert into catalog_edits + (edit_uri, did, subject_uri, entity, created_at, cid, fields, rationale) + values ($1, $2, $3, $4, $5, $6, $7, $8) on conflict (edit_uri) do update set did = excluded.did, subject_uri = excluded.subject_uri, entity = excluded.entity, - created_at = excluded.created_at", + created_at = excluded.created_at, + cid = excluded.cid, + fields = excluded.fields, + rationale = excluded.rationale", ) |> pog.parameter(pog.text(edit.edit_uri)) |> pog.parameter(pog.text(edit.did)) |> pog.parameter(pog.text(edit.subject_uri)) |> pog.parameter(pog.text(edit.entity)) |> pog.parameter(pog.text(edit.created_at)) + |> pog.parameter(pog.text(edit.cid)) + |> pog.parameter(pog.nullable(pog.text, edit.fields)) + |> pog.parameter(pog.nullable(pog.text, edit.rationale)) |> pog.execute(conn) case outcome { Ok(_) -> Nil @@ -503,22 +529,52 @@ fn delete_edits_for_did(conn: pog.Connection, did: String) -> Nil { Nil } +const edit_columns = "edit_uri, did, subject_uri, entity, created_at, cid, fields, rationale" + +fn edit_row_decoder() -> decode.Decoder(Edit) { + use edit_uri <- decode.field("edit_uri", decode.string) + use did <- decode.field("did", decode.string) + use subject_uri <- decode.field("subject_uri", decode.string) + use entity <- decode.field("entity", decode.string) + use created_at <- decode.field("created_at", decode.string) + // A pre-E2 row may have been upserted before `cid` existed; degrade to "". + use cid <- decode.optional_field("cid", "", decode.string) + use fields <- decode.field("fields", decode.optional(decode.string)) + use rationale <- decode.field("rationale", decode.optional(decode.string)) + decode.success(Edit( + edit_uri:, + did:, + subject_uri:, + entity:, + created_at:, + cid:, + fields:, + rationale:, + )) +} + fn edits_for(conn: pog.Connection, subject_uri: String) -> List(Edit) { - let row = { - use edit_uri <- decode.field("edit_uri", decode.string) - use did <- decode.field("did", decode.string) - use subject_uri <- decode.field("subject_uri", decode.string) - use entity <- decode.field("entity", decode.string) - use created_at <- decode.field("created_at", decode.string) - decode.success(Edit(edit_uri:, did:, subject_uri:, entity:, created_at:)) - } case - pog.query( - "select edit_uri, did, subject_uri, entity, created_at - from catalog_edits where subject_uri = $1 order by created_at", - ) + pog.query("select " <> edit_columns <> " + from catalog_edits where subject_uri = $1 order by created_at") |> pog.parameter(pog.text(subject_uri)) - |> pog.returning(row) + |> pog.returning(edit_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + +fn edits_for_subjects( + conn: pog.Connection, + subject_uris: List(String), +) -> List(Edit) { + case + pog.query("select " <> edit_columns <> " + from catalog_edits where subject_uri = any($1) order by created_at") + |> pog.parameter(pog.array(pog.text, subject_uris)) + |> pog.returning(edit_row_decoder()) |> pog.execute(conn) { Ok(pog.Returned(rows:, ..)) -> rows diff --git a/server/src/at_record_server/covers.gleam b/server/src/at_record_server/covers.gleam index 94a5f20..84e7064 100644 --- a/server/src/at_record_server/covers.gleam +++ b/server/src/at_record_server/covers.gleam @@ -105,6 +105,10 @@ pub fn copy_foreign_cover( source_did: String, cover: Blob, ) -> Option(Blob) { + // Live at write time (E3 of the appview-first roadmap; filed follow-up, + // not yet converged onto the identity cache): a byte-copy needs the + // source repo's current PDS to fetch from, and this write-time path is + // out of B7's read-path scope. case generated_client.identity_resolve_mini_doc( fetch_client, diff --git a/server/src/at_record_server/edit_inbox.gleam b/server/src/at_record_server/edit_inbox.gleam index 7aef22a..09c7a49 100644 --- a/server/src/at_record_server/edit_inbox.gleam +++ b/server/src/at_record_server/edit_inbox.gleam @@ -1,36 +1,41 @@ -//// The catalog.edit inbox: for each of the caller's own *current* catalog -//// releases (skipping ones already superseded by a newer own release), -//// fan out to Constellation for `catalog.edit` backlinks and fetch each -//// candidate via Slingshot. Best-effort throughout: a Constellation or -//// Slingshot failure for one release just skips that release rather than -//// failing the whole request. The fan-out is capped in both dimensions -//// (releases queried, proposals returned) with a warning logged when either -//// cap trims results, mirroring the browse fan-out's cap pattern. +//// The catalog.edit inbox. `list_proposals` reads entirely off the index +//// (E2 of the appview-first roadmap, hard swap per its decision 6): the +//// caller's own current releases come from `catalog/source.Source` (the +//// same `catalog_index`-backed port browse's resolution layer reads), and +//// pending edits against them come from one batched +//// `catalog_index.EditOps.list_for_subjects` call. The old per-release +//// Constellation-backlink fan-out plus per-proposal Slingshot `fetch_edit` +//// are gone: everything this handler needs is already indexed off the +//// firehose/backfill. `own_current_releases` (a live PDS `listRecords`) and +//// `current_only` survive unchanged: the edit-inbox *apply* flow still +//// needs the full record to mint a superseding release from, which the +//// index's `BrowseRow` doesn't carry. import at_record/gen/catalog/edit as catalog_edit import at_record/gen/catalog/list_edit_proposals.{type ProposalRow, ProposalRow} import at_record/gen/catalog/release as catalog_release import at_record/gen/defs +import at_record_server/catalog/row.{type BrowseRow} import at_record_server/catalog/source as catalog_source -import at_record_server/catalog_deps.{type Deps} +import at_record_server/catalog_index.{type Edit, type EditOps} import at_record_server/identity_resolver import at_record_server/oauth/sessions.{type OauthSession} -import atproto/constellation import atproto/repo import atproto/xrpc.{type Client} import gleam/dict.{type Dict} import gleam/dynamic/decode import gleam/int +import gleam/json import gleam/list import gleam/option.{type Option, None, Some} import gleam/set import gleam/string import wisp -const edit_backlink_source = catalog_edit.collection <> ":subject.uri" - -// How many of the caller's own current releases are queried for backlinks -// per request; releases beyond this (newest-first) are dropped and logged. +// How many of the caller's own current releases are queried per request; +// releases beyond this (newest-first) are dropped and logged. Bounds the +// `list_for_subjects` batch even though it's now one query, not a fan-out, +// so a caller with a huge catalog still gets a predictable response size. pub const release_fanout_cap = 50 // The merged proposal list's final cap, across all queried releases. @@ -42,7 +47,10 @@ pub type OwnRelease { /// The caller's own `catalog.release` records, reduced to the current tip of /// each supersedes chain: a release superseded by another own release is -/// dropped, since a proposal against it is stale. +/// dropped, since a proposal against it is stale. Still a live PDS read: the +/// edit-inbox *apply* flow (`handlers/edit_inbox.apply_to_target`) needs the +/// full decoded record to mint a superseding one from, which the index's +/// `BrowseRow` doesn't carry. pub fn own_current_releases( client: Client, session: OauthSession, @@ -75,48 +83,77 @@ fn own_row_decoder() -> decode.Decoder(OwnRelease) { /// Drop any release another own release's `supersedes` points at: only the /// newest link in each chain is current. pub fn current_only(rows: List(OwnRelease)) -> List(OwnRelease) { + current_only_by(rows, fn(r) { r.ref.uri }, fn(r) { + r.value.supersedes |> option.map(fn(s) { s.uri }) + }) +} + +/// The same supersedes-chain filter as `current_only`, generalized over +/// whichever shape carries a uri and an optional supersedes-uri: shared so +/// the live-fetched `OwnRelease` path (`current_only`, used by apply) and +/// the index-backed `BrowseRow` path (`current_release_rows`, used by +/// `list_proposals`) can never drift apart on what "current" means. +fn current_only_by( + items: List(a), + self_uri: fn(a) -> String, + supersedes_uri: fn(a) -> Option(String), +) -> List(a) { let superseded = - rows - |> list.filter_map(fn(r) { - r.value.supersedes - |> option.map(fn(s) { s.uri }) - |> option.to_result(Nil) - }) + items + |> list.filter_map(fn(r) { supersedes_uri(r) |> option.to_result(Nil) }) |> set.from_list - rows |> list.filter(fn(r) { !set.contains(superseded, r.ref.uri) }) + items |> list.filter(fn(r) { !set.contains(superseded, self_uri(r)) }) +} + +fn current_release_rows(rows: List(BrowseRow)) -> List(BrowseRow) { + current_only_by(rows, fn(r) { r.uri }, fn(r) { r.supersedes }) +} + +/// The caller's own current releases straight off the index: every indexed +/// row published by `did`, reduced to the current supersedes tip, newest +/// first, capped at `release_fanout_cap`. +fn own_current_rows( + variant_source: catalog_source.Source, + did: String, +) -> List(BrowseRow) { + variant_source.releases() + |> list.filter(fn(r) { r.publisher_did == did }) + |> current_release_rows + |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) + |> cap_rows(release_fanout_cap, "own-release fan-out") } /// The proposal rows across every capped own-current release, newest first, -/// capped again as a merged set. Proposer handles are resolved in one -/// batched `resolve_many` call across every backlink from every release, -/// rather than once per proposal, since the same proposer commonly appears -/// against several of the caller's releases. +/// capped again as a merged set (E2 of the appview-first roadmap: one batched +/// `list_for_subjects` read replaces the old per-release Constellation +/// backlink fan-out and per-proposal Slingshot `fetch_edit`). Proposer +/// handles are resolved in one batched `resolve_many` call across every +/// matching edit, rather than once per proposal. pub fn list_proposals( - deps: Deps, - session: OauthSession, - own: List(OwnRelease), variant_source: catalog_source.Source, + edits: EditOps, + session: OauthSession, identity: identity_resolver.Resolver, ) -> List(ProposalRow) { - let adoption_count = subject_adoption_count(variant_source) - let per_release = - cap_releases(own) - |> list.map(fn(owned) { #(owned, backlinks_for(deps, session, owned)) }) + let own = own_current_rows(variant_source, session.did) + let own_uris = own |> list.map(fn(r) { r.uri }) + let by_uri = own |> list.map(fn(r) { #(r.uri, r) }) |> dict.from_list + let subject_edits = + edits.list_for_subjects(own_uris) + |> list.filter(fn(e) { e.entity == "release" && e.did != session.did }) let handles = - per_release - |> list.flat_map(fn(pair) { list.map(pair.1, fn(b) { b.did }) }) + subject_edits + |> list.map(fn(e) { e.did }) |> identity_resolver.resolve_many(identity, _) |> dict.map_values(fn(_did, identity) { identity.handle }) - per_release - |> list.flat_map(fn(pair) { - let #(owned, backlinks) = pair - backlinks - |> list.filter_map(fn(b) { - to_proposal(deps, handles, adoption_count, owned, b) - |> option.to_result(Nil) - }) + let adoption_count = subject_adoption_count(variant_source) + subject_edits + |> list.filter_map(fn(edit) { + to_proposal(handles, adoption_count, by_uri, edit) + |> option.to_result(Nil) }) - // Constellation may return the same linking record more than once. + // A subject could in principle appear more than once if the index ever + // carries a duplicate write; cheap to guard against regardless. |> unique_by_uri |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) |> cap_rows(proposal_cap, "merged proposals") @@ -152,14 +189,6 @@ fn unique_by_uri(rows: List(ProposalRow)) -> List(ProposalRow) { |> list.reverse } -fn cap_releases(own: List(OwnRelease)) -> List(OwnRelease) { - own - |> list.sort(fn(a, b) { - string.compare(b.value.created_at, a.value.created_at) - }) - |> cap_rows(release_fanout_cap, "own-release fan-out") -} - fn cap_rows(rows: List(a), cap: Int, what: String) -> List(a) { case list.length(rows) > cap { False -> rows @@ -177,61 +206,63 @@ fn cap_rows(rows: List(a), cap: Int, what: String) -> List(a) { } } -fn backlinks_for( - deps: Deps, - session: OauthSession, - owned: OwnRelease, -) -> List(constellation.Backlink) { - case deps.backlinks(owned.ref.uri, edit_backlink_source) { - Error(e) -> { - wisp.log_warning( - "edit-inbox: backlinks lookup failed for " - <> owned.ref.uri - <> ": " - <> xrpc.describe(e), - ) - [] - } - Ok(page) -> - page.records - // Self-authored proposals against your own record shouldn't occur - // (amending your own record never files one), but drop them defensively. - |> list.filter(fn(b) { b.did != session.did }) - } -} - fn to_proposal( - deps: Deps, handles: Dict(String, String), adoption_count: fn(String) -> Option(Int), - owned: OwnRelease, - backlink: constellation.Backlink, + by_uri: Dict(String, BrowseRow), + edit: Edit, ) -> Option(ProposalRow) { - let edit_uri = constellation.record_uri(backlink) - use #(cid, edit) <- option.then(deps.fetch_edit(edit_uri)) - use target <- option.then(edit.subject) - use fields <- option.then(release_fields(edit)) - case edit.entity == "release" && target.uri == owned.ref.uri { - False -> None - True -> - Some(ProposalRow( - uri: edit_uri, - cid:, - proposer_did: backlink.did, - proposer_handle: dict.get(handles, backlink.did) |> option.from_result, - target_uri: target.uri, - release_title: owned.value.title, - current: Some(current_fields(owned.value)), - fields: Some(fields), - rationale: edit.rationale, - created_at: edit.created_at, - subject_adoption_count: adoption_count(owned.ref.uri), - )) + use owned <- option.then( + dict.get(by_uri, edit.subject_uri) |> option.from_result, + ) + use fields_json <- option.then(edit.fields) + use fields <- option.then( + json.parse(fields_json, catalog_edit.release_fields_decoder()) + |> option.from_result, + ) + Some(ProposalRow( + uri: edit.edit_uri, + cid: edit.cid, + proposer_did: edit.did, + proposer_handle: dict.get(handles, edit.did) |> option.from_result, + target_uri: owned.uri, + release_title: owned.title, + current: Some(current_fields(owned)), + fields: Some(fields), + rationale: edit.rationale, + created_at: edit.created_at, + subject_adoption_count: adoption_count(owned.uri), + )) +} + +/// The target release's current values for the fields a proposal can touch, +/// shaped like `ReleaseFields` so the client can diff against `fields` +/// without a second fetch. `BrowseRow.genres`/`.styles` don't distinguish +/// "field absent" from "field present but empty" the way the decoded +/// `CatalogRelease` did, so an empty list here reads as `None`. +fn current_fields(row: BrowseRow) -> catalog_edit.ReleaseFields { + catalog_edit.ReleaseFields( + title: Some(row.title), + released: row.released, + country: row.country, + genres: non_empty_list(row.genres), + styles: non_empty_list(row.styles), + thumb_url: None, + cover: None, + ) +} + +fn non_empty_list(items: List(a)) -> Option(List(a)) { + case items { + [] -> None + _ -> Some(items) } } /// The proposal's `#releaseFields` variant, when it carries one (the only /// variant this vertical slice understands; other entities/variants skip). +/// Used by the apply flow against a live-fetched `CatalogEdit`, not the +/// indexed `Edit` (see `handlers/edit_inbox.apply`). pub fn release_fields( edit: catalog_edit.CatalogEdit, ) -> Option(catalog_edit.ReleaseFields) { @@ -240,20 +271,3 @@ pub fn release_fields( _ -> None } } - -/// The target release's current values for the fields a proposal can touch, -/// shaped like `ReleaseFields` so the client can diff against `fields` -/// without a second fetch. -fn current_fields( - value: catalog_release.CatalogRelease, -) -> catalog_edit.ReleaseFields { - catalog_edit.ReleaseFields( - title: Some(value.title), - released: value.released, - country: value.country, - genres: value.genres, - styles: value.styles, - thumb_url: None, - cover: None, - ) -} diff --git a/server/src/at_record_server/handlers/amend.gleam b/server/src/at_record_server/handlers/amend.gleam index 3cdb6e1..f063bca 100644 --- a/server/src/at_record_server/handlers/amend.gleam +++ b/server/src/at_record_server/handlers/amend.gleam @@ -9,6 +9,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 +import at_record_server/browse import at_record_server/context.{ type Context, error_json, require_session, with_pds_client, } @@ -17,6 +18,7 @@ import at_record_server/crate import at_record_server/discogs_client import at_record_server/event_log import at_record_server/external_id +import at_record_server/known_users import at_record_server/oauth/sessions.{type OauthSession} import at_record_server/provenance import at_record_server/shelf_index @@ -356,12 +358,26 @@ pub fn mint_amended( catalog_release.encode_catalog_release(record), ) { - Ok(created) -> + Ok(created) -> { + // Read-your-writes (ADR 0002): without this, the apply-proposal-then- + // view loop (and a manual amend's own next load) would miss the fresh + // release until the firehose catches up. + ctx.catalog_index.releases.upsert(browse.to_browse_row( + known_users.KnownUser( + did: session.did, + handle: session.handle, + pds: session.pds, + ), + created.uri, + created.cid, + record, + )) Some(defs.CatalogRef( uri: created.uri, cid: created.cid, external_ids: None, )) + } Error(_) -> None } } diff --git a/server/src/at_record_server/handlers/crate_overlap.gleam b/server/src/at_record_server/handlers/crate_overlap.gleam index d0e2a5a..540d9db 100644 --- a/server/src/at_record_server/handlers/crate_overlap.gleam +++ b/server/src/at_record_server/handlers/crate_overlap.gleam @@ -45,6 +45,9 @@ fn do_get_crate_overlap( case event_log.load(client, session) { Error(_) -> error_json(502, "could not load your crate from PDS") Ok(own_stored) -> + // Live on every call (E3 of the appview-first roadmap): C6, the + // foreign side reading from the index instead, was left optional + // in lane C and hasn't been picked up. 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") diff --git a/server/src/at_record_server/handlers/edit_inbox.gleam b/server/src/at_record_server/handlers/edit_inbox.gleam index 5b93740..b73e637 100644 --- a/server/src/at_record_server/handlers/edit_inbox.gleam +++ b/server/src/at_record_server/handlers/edit_inbox.gleam @@ -37,15 +37,12 @@ import wisp.{type Request, type Response} const cursor_namespace = "proposals" pub fn list_proposals(req: Request, ctx: Context) -> Response { - use id, session <- require_session(req, ctx) - use client, session <- with_pds_client(ctx, id, session) - let own = edit_inbox.own_current_releases(client, session) + use _id, session <- require_session(req, ctx) let proposals = edit_inbox.list_proposals( - ctx.catalog, - session, - own, ctx.variant_source, + ctx.catalog_index.edits, + session, ctx.atproto.identity, ) let query = wisp.get_query(req) diff --git a/server/src/at_record_server/handlers/shelf.gleam b/server/src/at_record_server/handlers/shelf.gleam index a203bec..5c52776 100644 --- a/server/src/at_record_server/handlers/shelf.gleam +++ b/server/src/at_record_server/handlers/shelf.gleam @@ -662,6 +662,7 @@ fn promote_manual( let entities = catalog_entities.load(client, session) let outcome = promotion.adopt_or_mint( + ctx.catalog_index, ctx.catalog, client, session, diff --git a/server/src/at_record_server/jetstream_consumer.gleam b/server/src/at_record_server/jetstream_consumer.gleam index 0c41ea1..34bc6ba 100644 --- a/server/src/at_record_server/jetstream_consumer.gleam +++ b/server/src/at_record_server/jetstream_consumer.gleam @@ -444,7 +444,10 @@ fn to_browse_row( /// The publisher's #(handle, pds) for a did: a known user is free (no /// network hop); an unrecognized did is resolved via `resolveMiniDoc` and /// cached for the rest of this connection. Any resolution failure falls back -/// to an empty handle/pds rather than blocking indexing on it. +/// to an empty handle/pds rather than blocking indexing on it. Live by design +/// (E3 of the appview-first roadmap): this runs on the background firehose +/// connection, never a user-facing request, so there is no read-path latency +/// to save by routing it through the shared identity cache. fn resolve_publisher(state: State, did: String) -> #(#(String, String), State) { case list.find(state.deps.known_users.list(), fn(u) { u.did == did }) { Ok(user) -> #(#(user.handle, user.pds), state) @@ -579,8 +582,8 @@ fn route_edit(state: State, frame: Frame) -> State { } fn upsert_edit(state: State, frame: Frame) -> State { - case frame.record { - Some(record) -> + case frame.cid, frame.record { + Some(cid), Some(record) -> case decode.run(record, catalog_edit.catalog_edit_decoder()) { Ok(edit) -> case edit.subject { @@ -591,6 +594,9 @@ fn upsert_edit(state: State, frame: Frame) -> State { subject_uri: ref.uri, entity: edit.entity, created_at: edit.created_at, + cid:, + fields: edit_fields_json(edit), + rationale: edit.rationale, )) state } @@ -606,7 +612,30 @@ fn upsert_edit(state: State, frame: Frame) -> State { state } } - None -> state + _, _ -> { + wisp.log_warning( + "jetstream_consumer: " + <> frame.operation + <> " commit for " + <> at_uri(frame) + <> " missing cid/record", + ) + state + } + } +} + +/// The proposal's `#releaseFields` payload, JSON-encoded without a `$type` +/// wrapper (E2 of the appview-first roadmap): the only variant the edit +/// inbox understands today, stored as plain text so `catalog_index` stays +/// free of a `catalog_edit` dependency (see its `Edit` docstring). Any other +/// fields variant (future entity kinds) is dropped -- the read side already +/// only surfaces `entity == "release"` edits. +fn edit_fields_json(edit: catalog_edit.CatalogEdit) -> Option(String) { + case edit.fields { + Some(catalog_edit.CatalogEditFieldsReleaseFields(fields)) -> + Some(json.to_string(catalog_edit.encode_release_fields(fields))) + _ -> None } } diff --git a/server/src/at_record_server/oauth/flow.gleam b/server/src/at_record_server/oauth/flow.gleam index e4c5a3c..91b660e 100644 --- a/server/src/at_record_server/oauth/flow.gleam +++ b/server/src/at_record_server/oauth/flow.gleam @@ -37,6 +37,9 @@ pub fn start_login( cfg: Config, handle: String, ) -> Result(#(String, String), FlowError) { + // Live by design (E3 of the appview-first roadmap): login discovery needs + // the handle's *current* PDS/authorization-server, not a cached one, and + // the index has no reason to carry unauthenticated visitors' identities. use doc <- result.try( generated_client.identity_resolve_mini_doc( cfg.client, diff --git a/server/src/at_record_server/promotion.gleam b/server/src/at_record_server/promotion.gleam index 2d47289..8dad2e4 100644 --- a/server/src/at_record_server/promotion.gleam +++ b/server/src/at_record_server/promotion.gleam @@ -9,10 +9,13 @@ import at_record/gen/catalog/release as catalog_release import at_record/gen/defs +import at_record_server/browse import at_record_server/catalog_deps.{type Deps} import at_record_server/catalog_entities +import at_record_server/catalog_index import at_record_server/discogs_client import at_record_server/external_id +import at_record_server/known_users import at_record_server/oauth/sessions.{type OauthSession} import at_record_server/provenance import atproto/blob.{type Blob} @@ -114,6 +117,7 @@ fn own_row_decoder() -> decode.Decoder( /// Constellation backlink index, else mint (with artist/genre entities and the /// Discogs enrichment when `details` are supplied). Returns the grown maps. pub fn adopt_or_mint( + catalog_index: catalog_index.Store, deps: Deps, client: Client, session: OauthSession, @@ -146,7 +150,17 @@ pub fn adopt_or_mint( ) None -> { let #(entities, minted, report, artist_failed) = - mint(deps, client, session, entities, release, cover, details, None) + mint( + catalog_index, + deps, + client, + session, + entities, + release, + cover, + details, + None, + ) case minted { Some(ref) -> Outcome( @@ -178,6 +192,10 @@ pub fn adopt_or_mint( } // A backlink hit is verified before adoption: index noise isn't trusted, only a confirmed discogs id match is. +// Live by design (E3 of the appview-first roadmap): this is canonical-release +// discovery across every repo on the network, not just known users, so it +// can't be answered from our own index the way browse/scan matching can. +// Constellation is the durable public backlink index for exactly this job. fn discover( deps: Deps, session: OauthSession, @@ -235,6 +253,7 @@ pub fn parse_at_uri(uri: String) -> Option(#(String, String, String)) { // in MB requests, but does mean the call now also happens ahead of an artist // mint that turns out rate-limited or failed. fn mint( + catalog_index: catalog_index.Store, deps: Deps, client: Client, session: OauthSession, @@ -326,16 +345,31 @@ fn mint( catalog_release.encode_catalog_release(record), ) case created { - Ok(created) -> #( - entities, - Some(defs.CatalogRef( - uri: created.uri, - cid: created.cid, - external_ids: None, - )), - report, - False, - ) + Ok(created) -> { + // Read-your-writes (ADR 0002): without this, a freshly minted + // release wouldn't appear in browse/the edit inbox until the + // firehose catches up. + catalog_index.releases.upsert(browse.to_browse_row( + known_users.KnownUser( + did: session.did, + handle: session.handle, + pds: session.pds, + ), + created.uri, + created.cid, + record, + )) + #( + entities, + Some(defs.CatalogRef( + uri: created.uri, + cid: created.cid, + external_ids: None, + )), + report, + False, + ) + } Error(e) -> { wisp.log_warning( "release mint failed for discogs:" diff --git a/server/src/at_record_server/user_backfill.gleam b/server/src/at_record_server/user_backfill.gleam index a88d329..b0b6af8 100644 --- a/server/src/at_record_server/user_backfill.gleam +++ b/server/src/at_record_server/user_backfill.gleam @@ -27,6 +27,7 @@ import atproto/uri import atproto/xrpc.{type Client} import gleam/dynamic/decode import gleam/int +import gleam/json import gleam/list import gleam/option.{type Option, None, Some} import gleam/result @@ -174,12 +175,28 @@ fn backfill_edits( subject_uri: ref.uri, entity: edit.entity, created_at: edit.created_at, + cid: record.cid, + fields: edit_fields_json(edit), + rationale: edit.rationale, )) }) rows |> list.each(index.edits.upsert) list.length(rows) } +/// Mirrors `jetstream_consumer.edit_fields_json`: the proposal's +/// `#releaseFields` payload, JSON-encoded without a `$type` wrapper (E2 of +/// the appview-first roadmap), so both ingest sites store the exact same +/// shape and a read-site `release_fields_decoder` call works regardless of +/// which one wrote the row. +fn edit_fields_json(edit: catalog_edit.CatalogEdit) -> Option(String) { + case edit.fields { + Some(catalog_edit.CatalogEditFieldsReleaseFields(fields)) -> + Some(json.to_string(catalog_edit.encode_release_fields(fields))) + _ -> None + } +} + /// One user's records in one collection, decoded with `decode_row`, paged /// to exhaustion or to `max_pages`. A page's transport failure or the cap /// itself keeps whatever paged so far and logs a warning: this backfill diff --git a/server/test/catalog_index_postgres_test.gleam b/server/test/catalog_index_postgres_test.gleam index ad61178..bb5a616 100644 --- a/server/test/catalog_index_postgres_test.gleam +++ b/server/test/catalog_index_postgres_test.gleam @@ -123,9 +123,16 @@ pub fn adoption_and_edit_round_trip_test() { subject_uri: release, entity: "release", created_at: "2024-01-01T00:00:00Z", + cid: "bafyeditpg", + fields: Some("{\"title\":\"Better Title\"}"), + rationale: Some("typo fix"), )) let assert [found] = store.edits.for_subject(release) assert found.entity == "release" + assert found.cid == "bafyeditpg" + assert found.fields == Some("{\"title\":\"Better Title\"}") + assert found.rationale == Some("typo fix") + assert store.edits.list_for_subjects([release]) == [found] store.edits.delete(edit_uri) assert store.edits.for_subject(release) == [] } diff --git a/server/test/catalog_index_test.gleam b/server/test/catalog_index_test.gleam index c48b4fb..e473e7f 100644 --- a/server/test/catalog_index_test.gleam +++ b/server/test/catalog_index_test.gleam @@ -208,6 +208,9 @@ pub fn edit_upsert_delete_and_edits_for_test() { subject_uri: subject, entity: "release", created_at: "2024-01-01T00:00:00Z", + cid: "bafyedit", + fields: Some("{\"title\":\"Better Title\"}"), + rationale: Some("typo fix"), ) store.edits.upsert(edit) assert store.edits.for_subject(subject) == [edit] @@ -254,6 +257,9 @@ pub fn delete_for_did_drops_only_that_dids_adoptions_and_edits_test() { subject_uri: release, entity: "release", created_at: "2024-01-01T00:00:00Z", + cid: "bafyedit", + fields: None, + rationale: None, )) store.adoptions.delete_for_did("did:plc:a") store.edits.delete_for_did("did:plc:a") @@ -298,3 +304,43 @@ pub fn find_by_barcode_requires_the_already_despaced_column_value_test() { // caller's job to normalize (see `browse.despace`), not this store's. assert store.releases.find_by_barcode("0123 4567 89011") == None } + +pub fn edits_list_for_subjects_batches_across_multiple_releases_test() { + let assert Ok(store) = catalog_index.start() + let r1 = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" + let r2 = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2" + let r3 = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r3" + store.edits.upsert(Edit( + edit_uri: "at://did:plc:a/dev.mokkenstorm.crate.catalog.edit/x1", + did: "did:plc:a", + subject_uri: r1, + entity: "release", + created_at: "2024-01-01T00:00:00Z", + cid: "bafyedit1", + fields: None, + rationale: None, + )) + store.edits.upsert(Edit( + edit_uri: "at://did:plc:b/dev.mokkenstorm.crate.catalog.edit/x2", + did: "did:plc:b", + subject_uri: r2, + entity: "release", + created_at: "2024-01-02T00:00:00Z", + cid: "bafyedit2", + fields: None, + rationale: None, + )) + store.edits.upsert(Edit( + edit_uri: "at://did:plc:c/dev.mokkenstorm.crate.catalog.edit/x3", + did: "did:plc:c", + subject_uri: r3, + entity: "release", + created_at: "2024-01-03T00:00:00Z", + cid: "bafyedit3", + fields: None, + rationale: None, + )) + let found = store.edits.list_for_subjects([r1, r2]) + assert list.length(found) == 2 + assert list.all(found, fn(e) { e.subject_uri == r1 || e.subject_uri == r2 }) +} diff --git a/server/test/edit_inbox_apply_test.gleam b/server/test/edit_inbox_apply_test.gleam index 9b92933..25804e3 100644 --- a/server/test/edit_inbox_apply_test.gleam +++ b/server/test/edit_inbox_apply_test.gleam @@ -7,7 +7,8 @@ import at_record/gen/catalog/release as catalog_release import at_record/gen/defs.{CatalogRef} import at_record/gen/shelf/entry.{type ShelfEntry, ShelfEntry} import at_record_server/catalog_deps.{type Deps, Deps} -import at_record_server/context.{type Context} +import at_record_server/catalog_index +import at_record_server/context.{type Context, Context} import at_record_server/handlers/edit_inbox.{ Applied, CidMismatch, NotYourRecord, ProposalNotFound, ReleaseNotCurrent, apply, @@ -54,6 +55,16 @@ fn stub_context(catalog: Deps) -> Context { ) } +/// Same as `stub_context`, but with a real in-memory `catalog_index.Store` +/// instead of the no-op stub: lets a test read back what the apply flow's +/// write-through (`amend.mint_amended`, via `handlers/amend`) actually wrote. +fn stub_context_with_index( + catalog: Deps, + index: catalog_index.Store, +) -> Context { + Context(..stub_context(catalog), catalog_index: index) +} + fn old_release() -> catalog_release.CatalogRelease { catalog_release.CatalogRelease( ..support.blank_catalog_release(), @@ -255,6 +266,22 @@ pub fn apply_happy_path_mints_supersedes_and_repoints_the_entry_test() { assert release.cid == "bafynew" } +// E2 of the appview-first roadmap: the apply flow's mint goes through +// `amend.mint_amended`, whose write-through means the freshly minted release +// is readable off the index immediately, with no firehose simulation. +pub fn apply_write_through_makes_the_new_release_readable_off_the_index_test() { + let deps = + deps_with_edit(proposal_uri, proposal_cid, a_proposal(old_release_uri)) + let assert Ok(index) = catalog_index.start() + let client = stub_client(fn(_) { Nil }) + let ctx = stub_context_with_index(deps, index) + let result = apply(ctx, client, session(), proposal_uri, proposal_cid) + let assert Ok(Applied(release:)) = result + let assert Some(indexed) = index.releases.get(release.uri) + assert indexed.title == "Spiderland (Remastered)" + assert indexed.supersedes == Some(old_release_uri) +} + pub fn apply_rejects_a_missing_proposal_test() { let result = apply( diff --git a/server/test/edit_inbox_list_test.gleam b/server/test/edit_inbox_list_test.gleam index 982a563..5a8eed3 100644 --- a/server/test/edit_inbox_list_test.gleam +++ b/server/test/edit_inbox_list_test.gleam @@ -1,22 +1,20 @@ //// End-to-end `catalog.listEditProposals` handler tests for cursor -//// pagination. Same `wisp/simulate` + stub PDS client style as -//// `shelf_list_test`; one own release fans out to several proposals. +//// pagination (E2 of the appview-first roadmap: the handler reads entirely +//// off `ctx.variant_source`/`ctx.catalog_index.edits` now, no PDS client or +//// Constellation/Slingshot stubbing needed). One indexed owned release fans +//// out to several indexed proposals. import at_record/gen/catalog/edit as catalog_edit -import at_record/gen/catalog/release as catalog_release -import at_record/gen/defs.{CatalogRef} -import at_record_server/catalog_deps.{type Deps, Deps} -import at_record_server/context.{type Context} +import at_record_server/catalog/row.{type BrowseRow, BrowseRow} +import at_record_server/catalog/source as catalog_source +import at_record_server/catalog_index.{type Edit, Edit} +import at_record_server/context.{type Context, Context} import at_record_server/handlers/edit_inbox as edit_inbox_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/constellation.{Backlink, BacklinksPage} -import atproto/xrpc -import gleam/bit_array import gleam/http -import gleam/http/response import gleam/int import gleam/json import gleam/list @@ -42,34 +40,32 @@ fn a_session() -> sessions.OauthSession { ) } -fn owned_release() -> catalog_release.CatalogRelease { - catalog_release.CatalogRelease( - ..support.blank_catalog_release(), +fn owned_row() -> BrowseRow { + BrowseRow( + uri: owned_release_uri, + cid: "bafyowned", title: "Spiderland", + artist_display: None, + genres: [], + styles: [], + released: None, + country: None, + cover: None, + thumb_url: None, + discogs_id: None, + created_at: "2026-01-01T00:00:00Z", + publisher_did: "did:plc:x", + publisher_handle: "x.test", + publisher_pds: "https://pds.test", + supersedes: None, + based_on: None, + format: None, + label: None, + master: None, + barcode: None, ) } -/// listRecords response for `catalog.release`: one owned release, so -/// `own_current_releases` always resolves the same fan-out target. -fn stub_client() -> xrpc.Client { - let body = - json.object([ - #( - "records", - json.preprocessed_array([ - json.object([ - #("uri", json.string(owned_release_uri)), - #("cid", json.string("bafyowned")), - #("value", catalog_release.encode_catalog_release(owned_release())), - ]), - ]), - ), - ]) - xrpc.Client(send: fn(_req) { - Ok(response.Response(200, [], bit_array.from_string(json.to_string(body)))) - }) -} - fn rkey(n: Int) -> String { "e" <> string.pad_start(int.to_string(n), 3, "0") } @@ -90,73 +86,65 @@ fn release_fields(n: Int) -> catalog_edit.ReleaseFields { ) } -/// `count` proposals against the one owned release, newest (`count`) first -/// once `edit_inbox.list_proposals` sorts by `createdAt` descending. -fn deps_with_proposals(count: Int) -> Deps { - let backlinks = - list.repeat(Nil, count) - |> list.index_map(fn(_, i) { - Backlink( - did: proposer_did, - collection: catalog_edit.collection, - rkey: rkey(i + 1), - ) - }) - Deps( - backlinks: fn(_subject, _source) { - Ok(BacklinksPage(total: count, records: backlinks, cursor: None)) - }, - fetch_release: fn(_) { panic as "fetch_release must not be called" }, - fetch_edit: fn(uri) { - let n = index_of(uri) - Some(#( - "bafyedit" <> int.to_string(n), - catalog_edit.CatalogEdit( - created_at: "2026-02-0" <> int.to_string(n) <> "T00:00:00Z", - entity: "release", - fields: Some( - catalog_edit.CatalogEditFieldsReleaseFields(release_fields(n)), - ), - op: "update", - rationale: None, - source: None, - subject: Some(CatalogRef( - uri: owned_release_uri, - cid: "bafyowned", - external_ids: None, - )), - ), - )) - }, - release_mbid: fn(_, _) { panic as "release_mbid must not be called" }, - ) -} - -/// The backlink index (1-based) whose uri this is, decoded straight back -/// out of `proposal_uri`'s naming scheme. -fn index_of(uri: String) -> Int { - let assert Ok(n) = - uri - |> string.split("/") - |> list.last - |> result_unwrap("") - |> string.drop_start(1) - |> int.parse - n +fn fields_json(n: Int) -> String { + json.to_string(catalog_edit.encode_release_fields(release_fields(n))) +} + +/// `count` indexed edits against the one owned release, newest (`count`) +/// first once `edit_inbox.list_proposals` sorts by `createdAt` descending. +fn edits_for_proposals(count: Int) -> List(Edit) { + list.repeat(Nil, count) + |> list.index_map(fn(_, i) { + let n = i + 1 + Edit( + edit_uri: proposal_uri(n), + did: proposer_did, + subject_uri: owned_release_uri, + entity: "release", + created_at: "2026-02-0" <> int.to_string(n) <> "T00:00:00Z", + cid: "bafyedit" <> int.to_string(n), + fields: Some(fields_json(n)), + rationale: None, + ) + }) } -fn result_unwrap(r: Result(a, b), default: a) -> a { - case r { - Ok(v) -> v - Error(_) -> default - } +fn edits_ops(edits: List(Edit)) -> catalog_index.EditOps { + catalog_index.EditOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + delete_for_did: fn(_) { Nil }, + for_subject: fn(_) { [] }, + list_for_subjects: fn(_) { edits }, + ) } -fn test_context(deps: Deps) -> #(Context, config.Config) { +fn test_context(edits: List(Edit)) -> #(Context, config.Config) { let cfg = - support.stub_config_with(stub_client(), "r", "http://localhost:8080") + support.stub_config_with( + support.unreachable_client(), + "r", + "http://localhost:8080", + ) + let base = + support.stub_context_with( + cfg, + support.unreachable_catalog_deps(), + fn(_req) { Error("unused") }, + [], + ) let ctx = - support.stub_context_with(cfg, deps, fn(_req) { Error("unused") }, []) + Context( + ..base, + variant_source: catalog_source.Source( + releases: fn() { [owned_row()] }, + adoption_count: fn(_) { 0 }, + ), + catalog_index: catalog_index.Store( + ..support.empty_catalog_index(), + edits: edits_ops(edits), + ), + ) #(ctx, cfg) } @@ -181,7 +169,7 @@ fn proposal_uris(body: String) -> List(String) { } pub fn a_full_list_fetch_has_no_next_cursor_test() { - let #(ctx, cfg) = test_context(deps_with_proposals(5)) + let #(ctx, cfg) = test_context(edits_for_proposals(5)) let newest_first = [ proposal_uri(5), proposal_uri(4), @@ -198,7 +186,7 @@ pub fn a_full_list_fetch_has_no_next_cursor_test() { } pub fn a_cursor_past_the_last_proposal_returns_an_empty_page_test() { - let #(ctx, cfg) = test_context(deps_with_proposals(5)) + let #(ctx, cfg) = test_context(edits_for_proposals(5)) let cursor = pagination.encode_cursor("proposals", proposal_uri(1)) let body = list_proposals("?limit=1&cursor=" <> cursor, ctx, cfg) assert proposal_uris(body) == [] @@ -206,7 +194,7 @@ pub fn a_cursor_past_the_last_proposal_returns_an_empty_page_test() { } pub fn a_stale_or_garbage_cursor_restarts_from_the_top_test() { - let #(ctx, cfg) = test_context(deps_with_proposals(5)) + let #(ctx, cfg) = test_context(edits_for_proposals(5)) let unknown = pagination.encode_cursor("proposals", "at://unknown/x/y") ["not-a-real-cursor", unknown] |> list.each(fn(cursor) { @@ -216,7 +204,7 @@ pub fn a_stale_or_garbage_cursor_restarts_from_the_top_test() { } pub fn a_cursor_resumes_after_the_previous_page_test() { - let #(ctx, cfg) = test_context(deps_with_proposals(5)) + let #(ctx, cfg) = test_context(edits_for_proposals(5)) let first = list_proposals("?limit=2", ctx, cfg) let assert Ok(cursor) = support.field_string(first, ["cursor"]) let second = list_proposals("?limit=2&cursor=" <> cursor, ctx, cfg) diff --git a/server/test/edit_inbox_test.gleam b/server/test/edit_inbox_test.gleam index 06f96e7..45917b0 100644 --- a/server/test/edit_inbox_test.gleam +++ b/server/test/edit_inbox_test.gleam @@ -1,19 +1,24 @@ +//// `edit_inbox.list_proposals` tests against a real in-memory `catalog_index` +//// (E2 of the appview-first roadmap): releases and edits are seeded straight +//// into the index rather than stubbed behind Constellation/Slingshot fakes, +//// since the old fan-out is gone. `current_only`/`own_current_releases` +//// (the live-fetched `OwnRelease` path the apply flow still uses) keep their +//// own coverage too. + import at_record/gen/catalog/edit as catalog_edit import at_record/gen/catalog/list_edit_proposals.{ProposalRow} import at_record/gen/catalog/release as catalog_release import at_record/gen/defs.{CatalogRef} import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog/source as catalog_source -import at_record_server/catalog_deps.{type Deps, Deps} -import at_record_server/edit_inbox.{type OwnRelease, OwnRelease} +import at_record_server/catalog_index.{type Edit, Edit} +import at_record_server/edit_inbox.{OwnRelease} import at_record_server/identity_cache.{Identity} import at_record_server/identity_resolver import at_record_server/oauth/sessions.{type OauthSession} -import atproto/constellation.{type Backlink, Backlink, BacklinksPage} -import atproto/uri -import atproto_core/xrpc as core_xrpc import gleam/erlang/process import gleam/int +import gleam/json import gleam/list import gleam/option.{None, Some} import support @@ -40,43 +45,47 @@ fn range(count: Int) -> List(Int) { list.repeat(Nil, count) |> list.index_map(fn(_, i) { i + 1 }) } -fn release( - title: String, +fn own_release_uri(n: Int) -> String { + "at://" + <> own_did + <> "/dev.mokkenstorm.crate.catalog.release/r" + <> int.to_string(n) +} + +/// An indexed `BrowseRow` for the caller's own repo: only `uri`/`supersedes`/ +/// `created_at` vary across the tests below. +fn own_row( + n: Int, + supersedes: option.Option(String), created_at: String, -) -> catalog_release.CatalogRelease { - catalog_release.CatalogRelease( - ..support.blank_catalog_release(), - title:, - created_at:, +) -> BrowseRow { + BrowseRow( + uri: own_release_uri(n), + cid: "bafyrel" <> int.to_string(n), + title: "Spiderland " <> int.to_string(n), + artist_display: None, + genres: ["Rock"], + styles: [], released: Some("1991"), country: Some("US"), - genres: Some(["Rock"]), + cover: None, + thumb_url: None, + discogs_id: None, + created_at:, + publisher_did: own_did, + publisher_handle: "me.test", + publisher_pds: "https://pds.test", + supersedes:, + based_on: None, + format: None, + label: None, + master: None, + barcode: None, ) } -fn own_release( - n: Int, - supersedes: option.Option(defs.CatalogRef), -) -> OwnRelease { - let uri = - "at://" - <> own_did - <> "/dev.mokkenstorm.crate.catalog.release/r" - <> int.to_string(n) - OwnRelease( - ref: CatalogRef( - uri:, - cid: "bafyrel" <> int.to_string(n), - external_ids: None, - ), - value: catalog_release.CatalogRelease( - ..release( - "Spiderland " <> int.to_string(n), - "2026-01-0" <> int.to_string(n % 9 + 1) <> "T00:00:00Z", - ), - supersedes:, - ), - ) +fn own_release(n: Int, created_at: String) -> BrowseRow { + own_row(n, None, created_at) } fn release_fields(title: String) -> catalog_edit.ReleaseFields { @@ -91,152 +100,137 @@ fn release_fields(title: String) -> catalog_edit.ReleaseFields { ) } -fn edit_record( - entity: String, - target_uri: String, - fields: option.Option(catalog_edit.ReleaseFields), -) -> catalog_edit.CatalogEdit { - catalog_edit.CatalogEdit( - created_at: "2026-02-01T00:00:00Z", - entity:, - fields: option.map(fields, catalog_edit.CatalogEditFieldsReleaseFields), - op: "update", - rationale: None, - source: None, - subject: Some(CatalogRef( - uri: target_uri, - cid: "bafyold", - external_ids: None, - )), - ) +fn fields_json(title: String) -> String { + json.to_string(catalog_edit.encode_release_fields(release_fields(title))) } -fn backlink(did: String, n: Int) -> Backlink { - Backlink( +/// One indexed `Edit` targeting `target_uri`, carrying `title` as its only +/// changed field. +fn edit( + edit_uri edit_uri: String, + did did: String, + target_uri target_uri: String, + title title: String, +) -> Edit { + Edit( + edit_uri:, did:, - collection: "dev.mokkenstorm.crate.catalog.edit", - rkey: "e" <> int.to_string(n), - ) -} - -/// A minimal `BrowseRow` fixture: only `uri` varies, everything else is a -/// placeholder since these tests only exercise index membership. -fn indexed_row(uri: String) -> BrowseRow { - BrowseRow( - uri:, - cid: "bafyindexed", - title: "Title", - artist_display: None, - genres: [], - styles: [], - released: None, - country: None, - cover: None, - thumb_url: None, - discogs_id: None, - created_at: "2026-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, - barcode: None, + subject_uri: target_uri, + entity: "release", + created_at: "2026-02-01T00:00:00Z", + cid: "bafyedit", + fields: Some(fields_json(title)), + rationale: None, ) } -/// A `Source` whose `releases` reports exactly `uris` as indexed, so a -/// `subject_adoption_count` lookup against one of them falls through to -/// `count`. -fn indexed_source( - uris: List(String), - count: fn(String) -> Int, -) -> catalog_source.Source { - catalog_source.Source( - releases: fn() { list.map(uris, indexed_row) }, - adoption_count: count, - ) +/// A `Source` reporting `rows` as the whole index (own releases and, in the +/// fan-out cap test, nothing else): `list_proposals` derives "own current" +/// straight off `releases()`, filtered by `publisher_did`. +fn source_over(rows: List(BrowseRow)) -> catalog_source.Source { + catalog_source.Source(releases: fn() { rows }, adoption_count: fn(_) { 0 }) } -/// Deps whose `backlinks` returns one backlink per own release (via `n`) and -/// whose `fetch_edit` looks the matching edit record up from `edits`. -fn deps_with( - backlinks_by_release: fn(String) -> List(Backlink), - edits: fn(String) -> option.Option(#(String, catalog_edit.CatalogEdit)), -) -> Deps { - Deps( - backlinks: fn(subject, _source) { - Ok(BacklinksPage( - total: 0, - records: backlinks_by_release(subject), - cursor: None, - )) - }, - fetch_release: fn(_) { panic as "fetch_release must not be called" }, - fetch_edit: edits, - release_mbid: fn(_, _) { panic as "release_mbid must not be called" }, +/// `edits.list_for_subjects` backed by a fixed list, ignoring the requested +/// subject uris (every test here only ever seeds edits it wants returned). +fn edits_ops(edits: List(Edit)) -> catalog_index.EditOps { + catalog_index.EditOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + delete_for_did: fn(_) { Nil }, + for_subject: fn(_) { [] }, + list_for_subjects: fn(_) { edits }, ) } pub fn current_only_drops_a_release_another_own_release_supersedes_test() { - let old = own_release(1, None) - let new = own_release(2, Some(old.ref)) + let old = + OwnRelease( + ref: CatalogRef( + uri: own_release_uri(1), + cid: "bafyrel1", + external_ids: None, + ), + value: blank_release("2026-01-01T00:00:00Z"), + ) + let new = + OwnRelease( + ref: CatalogRef( + uri: own_release_uri(2), + cid: "bafyrel2", + external_ids: None, + ), + value: catalog_release.CatalogRelease( + ..blank_release("2026-01-02T00:00:00Z"), + supersedes: Some(old.ref), + ), + ) let rows = edit_inbox.current_only([old, new]) assert list.map(rows, fn(r) { r.ref.uri }) == [new.ref.uri] } pub fn current_only_keeps_independent_releases_test() { - let a = own_release(1, None) - let b = own_release(2, None) + let a = + OwnRelease( + ref: CatalogRef( + uri: own_release_uri(1), + cid: "bafyrel1", + external_ids: None, + ), + value: blank_release("2026-01-01T00:00:00Z"), + ) + let b = + OwnRelease( + ref: CatalogRef( + uri: own_release_uri(2), + cid: "bafyrel2", + external_ids: None, + ), + value: blank_release("2026-01-02T00:00:00Z"), + ) let rows = edit_inbox.current_only([a, b]) assert list.length(rows) == 2 } +fn blank_release(created_at: String) -> catalog_release.CatalogRelease { + catalog_release.CatalogRelease(..support.blank_catalog_release(), created_at:) +} + pub fn own_authored_proposal_is_excluded_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(own_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - owned.ref.uri, - Some(release_fields("Better Title")), - ), - )) - }) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> own_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: own_did, + target_uri: owned.uri, + title: "Better Title", + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) assert rows == [] } pub fn foreign_release_entity_proposal_is_included_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - owned.ref.uri, - Some(release_fields("Better Title")), - ), - )) - }) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned.uri, + title: "Better Title", + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) let assert [ @@ -250,90 +244,61 @@ pub fn foreign_release_entity_proposal_is_included_test() { ), ] = rows assert proposer_did == other_did - assert target_uri == owned.ref.uri - assert release_title == owned.value.title + assert target_uri == owned.uri + assert release_title == owned.title assert fields == Some(release_fields("Better Title")) assert current == Some(catalog_edit.ReleaseFields( - title: Some(owned.value.title), - released: owned.value.released, - country: owned.value.country, - genres: owned.value.genres, - styles: owned.value.styles, + title: Some(owned.title), + released: owned.released, + country: owned.country, + genres: Some(owned.genres), + styles: None, thumb_url: None, cover: None, )) } pub fn subject_adoption_count_is_populated_when_the_target_release_is_indexed_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - owned.ref.uri, - Some(release_fields("Better Title")), - ), - )) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned.uri, + title: "Better Title", + ), + ] + let source = + catalog_source.Source(releases: fn() { [owned] }, adoption_count: fn(_) { + 3 }) let rows = edit_inbox.list_proposals( - deps, + source, + edits_ops(edits), session(), - [owned], - indexed_source([owned.ref.uri], fn(_) { 3 }), offline_identity_resolver(), ) let assert [ProposalRow(subject_adoption_count:, ..)] = rows assert subject_adoption_count == Some(3) } -pub fn subject_adoption_count_is_omitted_when_the_target_release_is_not_indexed_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - owned.ref.uri, - Some(release_fields("Better Title")), - ), - )) - }) - let rows = - edit_inbox.list_proposals( - deps, - session(), - [owned], - support.empty_variant_source(), - offline_identity_resolver(), - ) - let assert [ProposalRow(subject_adoption_count:, ..)] = rows - assert subject_adoption_count == None -} - pub fn proposer_handle_resolution_failure_degrades_silently_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - owned.ref.uri, - Some(release_fields("Better Title")), - ), - )) - }) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned.uri, + title: "Better Title", + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) let assert [ProposalRow(proposer_handle:, ..)] = rows @@ -341,66 +306,70 @@ pub fn proposer_handle_resolution_failure_degrades_silently_test() { } pub fn non_release_entity_proposal_is_dropped_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record("artist", owned.ref.uri, Some(release_fields("x"))), - )) - }) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + Edit( + ..edit( + edit_uri: "at://" + <> other_did + <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned.uri, + title: "x", + ), + entity: "artist", + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) assert rows == [] } pub fn proposal_targeting_a_different_release_is_dropped_test() { - let owned = own_release(1, None) - let deps = - deps_with(fn(_subject) { [backlink(other_did, 1)] }, fn(_uri) { - Some(#( - "bafyedit1", - edit_record( - "release", - "at://someone/else/rkey", - Some(release_fields("x")), - ), - )) - }) + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: "at://someone/else/rkey", + title: "x", + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) assert rows == [] } -pub fn a_failed_backlinks_lookup_skips_that_release_rather_than_failing_test() { - let owned = own_release(1, None) - let deps = - Deps( - backlinks: fn(_subject, _source) { - Error(core_xrpc.RequestFailed(core_xrpc.ConnectionFailed("down"))) - }, - fetch_release: fn(_) { panic as "fetch_release must not be called" }, - fetch_edit: fn(_) { panic as "fetch_edit must not be called" }, - release_mbid: fn(_, _) { panic as "release_mbid must not be called" }, - ) +pub fn a_proposal_with_undecodable_fields_is_skipped_rather_than_failing_test() { + let owned = own_release(1, "2026-01-01T00:00:00Z") + let edits = [ + Edit( + ..edit( + edit_uri: "at://" + <> other_did + <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned.uri, + title: "x", + ), + fields: Some("not json"), + ), + ] let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) assert rows == [] @@ -408,99 +377,76 @@ pub fn a_failed_backlinks_lookup_skips_that_release_rather_than_failing_test() { pub fn release_fanout_is_capped_newest_first_test() { let count = edit_inbox.release_fanout_cap + 5 - let owns = range(count) |> list.map(fn(n) { own_release(n, None) }) - // Each release's backlink rkey carries the release's own rkey, so - // `fetch_edit` (keyed only on the resulting edit uri) can rebuild a - // proposal targeting exactly the release that was queried. - let deps = - deps_with( - fn(subject) { - [ - Backlink( - did: other_did, - collection: catalog_edit.collection, - rkey: uri.rkey(subject), - ), - ] - }, - fn(edit_uri) { - let target_uri = - "at://" - <> own_did - <> "/" - <> catalog_release.collection - <> "/" - <> uri.rkey(edit_uri) - Some(#( - "bafyedit", - edit_record("release", target_uri, Some(release_fields("x"))), - )) - }, - ) + let owns = + range(count) + |> list.map(fn(n) { + own_release(n, "2026-01-0" <> int.to_string(n % 9 + 1) <> "T00:00:00Z") + }) + let edits = + owns + |> list.map(fn(owned) { + edit( + edit_uri: owned.uri <> "-edit", + did: other_did, + target_uri: owned.uri, + title: "x", + ) + }) let rows = edit_inbox.list_proposals( - deps, + source_over(owns), + edits_ops(edits), session(), - owns, - support.empty_variant_source(), offline_identity_resolver(), ) assert list.length(rows) == edit_inbox.release_fanout_cap } pub fn merged_proposal_cap_applies_across_a_single_release_test() { - let owned = own_release(1, None) + let owned = own_release(1, "2026-01-01T00:00:00Z") let count = edit_inbox.proposal_cap + 5 - let backlinks = range(count) |> list.map(fn(n) { backlink(other_did, n) }) - let deps = - deps_with(fn(_subject) { backlinks }, fn(_uri) { - Some(#( - "bafyedit", - edit_record("release", owned.ref.uri, Some(release_fields("x"))), - )) + let edits = + range(count) + |> list.map(fn(n) { + edit( + edit_uri: "at://" + <> other_did + <> "/dev.mokkenstorm.crate.catalog.edit/x" + <> int.to_string(n), + did: other_did, + target_uri: owned.uri, + title: "x", + ) }) let rows = edit_inbox.list_proposals( - deps, + source_over([owned]), + edits_ops(edits), session(), - [owned], - support.empty_variant_source(), offline_identity_resolver(), ) assert list.length(rows) == edit_inbox.proposal_cap } // The same proposer DID files against two different own releases; the -// batched pre-resolve (collected across every release before the one -// `resolve_many`) must still call the underlying fetch only once. +// batched pre-resolve must still call the underlying fetch only once. pub fn several_proposals_sharing_one_proposer_did_resolve_with_one_fetch_test() { - let owned1 = own_release(1, None) - let owned2 = own_release(2, None) - let deps = - deps_with( - fn(subject) { - [ - Backlink( - did: other_did, - collection: catalog_edit.collection, - rkey: uri.rkey(subject), - ), - ] - }, - fn(edit_uri) { - let target_uri = - "at://" - <> own_did - <> "/" - <> catalog_release.collection - <> "/" - <> uri.rkey(edit_uri) - Some(#( - "bafyedit" <> uri.rkey(edit_uri), - edit_record("release", target_uri, Some(release_fields("x"))), - )) - }, - ) + let owned1 = own_release(1, "2026-01-01T00:00:00Z") + let owned2 = own_release(2, "2026-01-02T00:00:00Z") + let edits = [ + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x1", + did: other_did, + target_uri: owned1.uri, + title: "x", + ), + edit( + edit_uri: "at://" <> other_did <> "/dev.mokkenstorm.crate.catalog.edit/x2", + did: other_did, + target_uri: owned2.uri, + title: "x", + ), + ] let counter = process.new_subject() let identity = support.counting_identity_resolver(counter, fn(id) { @@ -508,10 +454,9 @@ pub fn several_proposals_sharing_one_proposer_did_resolve_with_one_fetch_test() }) let rows = edit_inbox.list_proposals( - deps, + source_over([owned1, owned2]), + edits_ops(edits), session(), - [owned1, owned2], - support.empty_variant_source(), identity, ) assert list.length(rows) == 2 diff --git a/server/test/jetstream_consumer_test.gleam b/server/test/jetstream_consumer_test.gleam index c3653fa..31c6180 100644 --- a/server/test/jetstream_consumer_test.gleam +++ b/server/test/jetstream_consumer_test.gleam @@ -148,7 +148,7 @@ fn catalog_edit_frame(did: String, time_us: Int, operation: String) -> String { collection: "dev.mokkenstorm.crate.catalog.edit", rkey: "x1", cid: "bafyeditcid", - record: "{\"$type\":\"dev.mokkenstorm.crate.catalog.edit\",\"entity\":\"release\",\"op\":\"propose\",\"createdAt\":\"2024-01-03T00:00:00Z\",\"subject\":{\"cid\":\"bafyrelcid\",\"uri\":\"" + record: "{\"$type\":\"dev.mokkenstorm.crate.catalog.edit\",\"entity\":\"release\",\"op\":\"propose\",\"createdAt\":\"2024-01-03T00:00:00Z\",\"rationale\":\"typo fix\",\"fields\":{\"$type\":\"dev.mokkenstorm.crate.catalog.edit#releaseFields\",\"title\":\"Better Title\"},\"subject\":{\"cid\":\"bafyrelcid\",\"uri\":\"" <> release_uri <> "\"}}", ) @@ -504,6 +504,9 @@ fn seed_catalog_rows_for(index: catalog_index.Store, did: String) -> Nil { subject_uri: release_uri, entity: "release", created_at: "2024-01-01T00:00:00Z", + cid: "bafyeditcid", + fields: None, + rationale: None, )) } @@ -794,6 +797,9 @@ pub fn catalog_edit_indexes_subject_and_entity_then_delete_removes_it_test() { let assert [edit] = index.edits.for_subject(release_uri) assert edit.did == "did:plc:a" assert edit.entity == "release" + assert edit.cid == "bafyeditcid" + assert edit.rationale == Some("typo fix") + assert edit.fields == Some("{\"title\":\"Better Title\"}") let _ = jetstream_consumer.route( state, diff --git a/server/test/promotion_test.gleam b/server/test/promotion_test.gleam index 4826fdd..af824b5 100644 --- a/server/test/promotion_test.gleam +++ b/server/test/promotion_test.gleam @@ -187,6 +187,7 @@ pub fn own_map_hit_reuses_without_network_test() { let own = dict.from_list([#("42", existing)]) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), offline_deps(), dead_client(), session(), @@ -228,6 +229,7 @@ pub fn verified_backlink_is_adopted_test() { ) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, dead_client(), session(), @@ -261,6 +263,7 @@ pub fn unverified_backlink_falls_through_to_mint_test() { let minted_uri = "at://" <> own_did <> "/release/3new" let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, minting_client(minted_uri), session(), @@ -288,6 +291,7 @@ pub fn constellation_outage_still_mints_test() { ) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, minting_client("at://" <> own_did <> "/release/3ddd"), session(), @@ -335,6 +339,7 @@ pub fn mint_with_details_enriches_release_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -364,6 +369,7 @@ pub fn mint_without_details_sets_only_artist_display_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -407,6 +413,7 @@ pub fn mint_with_resolved_mbid_adds_musicbrainz_external_id_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -437,6 +444,7 @@ pub fn mint_with_unresolved_mbid_keeps_only_discogs_external_id_test() { let deps = no_backlinks_deps_with_mbid(fn(_, _) { None }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -469,6 +477,7 @@ pub fn mint_drops_nameless_format_entries_test() { ]) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -526,6 +535,7 @@ pub fn mint_maps_discogs_formats_into_lexicon_shape_test() { ]) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -565,6 +575,7 @@ pub fn mint_with_master_id_adds_discogs_master_external_id_test() { let details_with_master = ReleaseDetails(..details(), master_id: Some(999)) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -608,6 +619,7 @@ pub fn mint_attaches_musicbrainz_id_to_exact_matching_artist_credit_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -648,6 +660,7 @@ pub fn mint_attaches_musicbrainz_id_via_token_match_on_punctuation_drift_test() }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -685,6 +698,7 @@ pub fn mint_leaves_artist_without_musicbrainz_id_on_name_mismatch_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -732,6 +746,7 @@ pub fn dedup_hit_backfills_musicbrainz_id_via_put_record_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -769,6 +784,7 @@ pub fn dedup_hit_with_existing_musicbrainz_id_is_not_touched_test() { ]) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, client, session(), @@ -794,6 +810,7 @@ pub fn failed_backfill_still_mints_release_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), deps, artist_failing_client(), session(), @@ -828,6 +845,7 @@ pub fn mint_leaves_artist_without_musicbrainz_id_on_mb_miss_test() { }) let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), client, session(), @@ -843,6 +861,7 @@ pub fn mint_leaves_artist_without_musicbrainz_id_on_mb_miss_test() { pub fn failed_artist_mint_withholds_release_test() { let outcome = promotion.adopt_or_mint( + support.empty_catalog_index(), no_backlinks_deps(), artist_failing_client(), session(), diff --git a/server/test/support.gleam b/server/test/support.gleam index 2d3fad7..b607737 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -139,6 +139,7 @@ pub fn empty_catalog_index() -> catalog_index.Store { delete: fn(_) { Nil }, delete_for_did: fn(_) { Nil }, for_subject: fn(_) { [] }, + list_for_subjects: fn(_) { [] }, ), cursor: catalog_index.CursorOps(save: fn(_) { Nil }, load: fn() { None }), )