diff --git a/server/src/at_record_server/browse.gleam b/server/src/at_record_server/browse.gleam index c076496..489e591 100644 --- a/server/src/at_record_server/browse.gleam +++ b/server/src/at_record_server/browse.gleam @@ -119,10 +119,16 @@ pub fn flags(own: Ownership, row: BrowseRow) -> #(Bool, Bool) { } /// The 4-digit year at the head of a `released` string (e.g. "1991-03-27"). +/// Anything shorter than 4 characters, or whose first 4 characters are not +/// all digits, is not a year and yields `None` rather than a truncated +/// number. pub fn year_from_released(released: Option(String)) -> Option(Int) { released |> option.then(fn(r) { - int.parse(string.slice(r, 0, 4)) |> option.from_result + case string.length(r) >= 4 { + True -> int.parse(string.slice(r, 0, 4)) |> option.from_result + False -> None + } }) } diff --git a/server/src/at_record_server/catalog_import.gleam b/server/src/at_record_server/catalog_import.gleam index 2a15c33..d55224a 100644 --- a/server/src/at_record_server/catalog_import.gleam +++ b/server/src/at_record_server/catalog_import.gleam @@ -20,7 +20,6 @@ import at_record_server/oauth/sessions.{type OauthSession} import at_record_server/promotion import at_record_server/provenance import atproto/blob -import atproto/repo import atproto/xrpc.{type Client} import gleam/dict import gleam/int @@ -28,6 +27,8 @@ import gleam/list import gleam/option.{type Option, None, Some} import gleam/result import gleam/set.{type Set} +import gleam/string +import wisp /// The running tally a capped import returns; `done` is false when more pages /// remain (budget spent, rate limited, or a write failed mid-run). @@ -51,9 +52,41 @@ type RunState { genesis_by_id: dict.Dict(String, storage.StoredItem(ShelfEntry)), budget: Int, entity_failures: Int, + details_missing: Int, + skipped_items: Int, ) } +/// The per-page write tally the release loop returns: successful genesis +/// writes, plus the three failure classes the run note surfaces. +type WriteStats { + WriteStats( + written: Int, + entity_failures: Int, + details_missing: Int, + skipped_items: Int, + ) +} + +/// The outcome of one genesis attempt. `Wrote` counts genre-mint failures and +/// whether the release was minted without Discogs details; `SkippedArtistMint` +/// withholds the item (artist mint failed) so it retries; `Stopped` halts the +/// run with a note, carrying any entity failures already banked on the item. +type GenesisResult { + Wrote(entity_failed: Int, details_missing: Bool) + SkippedArtistMint + Stopped(note: String, entity_failed: Int) +} + +/// The three shapes of the shared per-release fetch: details in hand, a rate +/// limit that must stop the run before minting, or a non-429 miss that lets the +/// item mint detail-less (counted so the run note can flag it for amending). +type FetchedDetails { + GotDetails(discogs_client.ReleaseDetails) + DetailsRateLimited + DetailsMissing +} + /// The shape shared by `discogs_client.collection_page` and `.wantlist_page`, /// so a run can be pointed at either source. type PageFetcher = @@ -201,6 +234,8 @@ fn initial_state( genesis_by_id:, budget: cap, entity_failures: 0, + details_missing: 0, + skipped_items: 0, ) } @@ -249,7 +284,7 @@ fn import_pages( ) let planned = list.length(plan.new) + list.length(plan.backfill) let unseen_left = list.length(found.releases) > planned + plan.skipped - let #(written, own, entities, entity_failed, stop_note) = + let #(stats, own, entities, stop_note) = write_releases( ctx, client, @@ -273,22 +308,28 @@ fn import_pages( let acc = ImportRun( ..acc, - imported: acc.imported + written, + imported: acc.imported + stats.written, updated: acc.updated + updated, skipped: acc.skipped + plan.skipped + cover_missing, ) - let budget = state.budget - written - updated + let budget = state.budget - stats.written - updated let last_page = found.page >= found.pages let stop = option.or(stop_note, backfill_note) - let failures = state.entity_failures + entity_failed + let failures = state.entity_failures + stats.entity_failures + let details_missing = state.details_missing + stats.details_missing + let skipped_items = state.skipped_items + stats.skipped_items + let terminal = fn(done: Bool) { + ImportRun( + ..acc, + done:, + note: run_note(stop, failures, details_missing, skipped_items), + ) + } case stop, last_page, unseen_left || budget <= 0 { - Some(note), _, _ -> ImportRun(..acc, done: False, note: Some(note)) - None, True, False -> - ImportRun(..acc, done: True) |> note_entity_failures(failures) - None, True, True -> - ImportRun(..acc, done: False) |> note_entity_failures(failures) - None, False, True -> - ImportRun(..acc, done: False) |> note_entity_failures(failures) + Some(_), _, _ -> terminal(False) + None, True, False -> terminal(True) + None, True, True -> terminal(False) + None, False, True -> terminal(False) None, False, False -> import_pages( ctx, @@ -307,6 +348,8 @@ fn import_pages( entities:, budget:, entity_failures: failures, + details_missing:, + skipped_items:, ), page + 1, acc, @@ -326,89 +369,140 @@ fn write_backfills( ) -> #(Int, Int, Option(String)) { case releases { [] -> #(0, 0, None) - [release, ..rest] -> { - let outcome = { - use genesis <- result.try( - dict.get(genesis_by_id, int.to_string(release.discogs_id)) - |> result.replace_error("no genesis"), - ) - use cover <- result.try( - best_cover(ctx, client, session, release, fetch_details(ctx, release)) - |> option.to_result("no cover"), - ) - ShelfEntry( - subject: Some(RepoStrongRef(cid: genesis.cid, uri: genesis.uri)), - action: "updated", - snapshot: Some(Snapshot( - title: release.title, - artist_display: release.artist, - year: release.year, - format: release.format, - thumb_url: release.thumb_url, - cover: Some(cover), - )), - external_ids: None, - media_grade: None, - sleeve_grade: None, - rating: None, - folder: None, - notes: None, - release: None, - price: None, - counterparty: None, - source: Some( - provenance.source("discogs-import") - |> provenance.with_external(discogs_client.external_id(release)), - ), - created_at: provenance.now_rfc3339(), - ) - |> event_log.append(client, session, _) - |> result.map_error(fn(e) { - case e { - xrpc.BadStatus(status: 429, ..) -> "rate limited" - _ -> "write failed" - } - }) - } - case outcome { - Ok(_) -> { - let #(u, missing, note) = - write_backfills(ctx, client, session, genesis_by_id, rest) - #(u + 1, missing, note) - } - Error("no cover") | Error("no genesis") -> { - let #(u, missing, note) = - write_backfills(ctx, client, session, genesis_by_id, rest) - #(u, missing + 1, note) - } - Error("rate limited") -> #( + [release, ..rest] -> + case fetch_details(ctx, release) { + // A cover-backfill enrichment read that 429s stops the run so the item + // retries, mirroring the genesis path rather than swallowing the 429. + DetailsRateLimited -> #( 0, 0, - Some("PDS rate limited; run again later"), + Some("Discogs rate limited; run again in a minute"), ) - Error(_) -> #(0, 0, Some("PDS write failed; run again to resume")) + fetched -> + backfill_one( + ctx, + client, + session, + genesis_by_id, + release, + rest, + fetched, + ) } + } +} + +// One backfill item, given its already-fetched details. A non-429 details miss +// falls through to the existing missing-cover handling (counted, item skipped). +fn backfill_one( + ctx: Context, + client: Client, + session: OauthSession, + genesis_by_id: dict.Dict(String, storage.StoredItem(ShelfEntry)), + release: discogs_client.Release, + rest: List(discogs_client.Release), + fetched: FetchedDetails, +) -> #(Int, Int, Option(String)) { + let outcome = { + use genesis <- result.try( + dict.get(genesis_by_id, int.to_string(release.discogs_id)) + |> result.replace_error("no genesis"), + ) + use cover <- result.try( + best_cover(ctx, client, session, release, details_option(fetched)) + |> option.to_result("no cover"), + ) + ShelfEntry( + subject: Some(RepoStrongRef(cid: genesis.cid, uri: genesis.uri)), + action: "updated", + snapshot: Some(Snapshot( + title: release.title, + artist_display: release.artist, + year: release.year, + format: release.format, + thumb_url: release.thumb_url, + cover: Some(cover), + )), + external_ids: None, + media_grade: None, + sleeve_grade: None, + rating: None, + folder: None, + notes: None, + release: None, + price: None, + counterparty: None, + source: Some( + provenance.source("discogs-import") + |> provenance.with_external(discogs_client.external_id(release)), + ), + created_at: provenance.now_rfc3339(), + ) + |> event_log.append(client, session, _) + |> result.map_error(fn(e) { + case e { + xrpc.BadStatus(status: 429, ..) -> "rate limited" + _ -> "write failed" + } + }) + } + case outcome { + Ok(_) -> { + let #(u, missing, note) = + write_backfills(ctx, client, session, genesis_by_id, rest) + #(u + 1, missing, note) + } + Error("no cover") | Error("no genesis") -> { + let #(u, missing, note) = + write_backfills(ctx, client, session, genesis_by_id, rest) + #(u, missing + 1, note) } + Error("rate limited") -> #(0, 0, Some("PDS rate limited; run again later")) + Error(_) -> #(0, 0, Some("PDS write failed; run again to resume")) } } -// Surfaces counted entity-write failures via the note when nothing more urgent -// (a stop/rate-limit note) already claims it. -fn note_entity_failures(run: ImportRun, failures: Int) -> ImportRun { - case run.note, failures > 0 { - None, True -> - ImportRun( - ..run, - note: Some( - int.to_string(failures) - <> " artist/genre records failed; run again to fill them", - ), - ) - _, _ -> run +/// The run note, composing the stop reason (if any) with the three banked +/// failure counts. None when nothing needs surfacing. Each clause is phrased to +/// read both standalone and appended after a `; ` to the stop reason. +pub fn run_note( + stop: Option(String), + failures: Int, + details_missing: Int, + skipped_items: Int, +) -> Option(String) { + [ + stop, + clause(failures, " artist/genre records failed; run again to fill them"), + clause( + details_missing, + " releases minted without Discogs details; amend later to enrich", + ), + clause( + skipped_items, + " items skipped: artist mint failed; run again to retry", + ), + ] + |> option.values + |> join_notes +} + +fn clause(count: Int, suffix: String) -> Option(String) { + case count > 0 { + True -> Some(int.to_string(count) <> suffix) + False -> None + } +} + +fn join_notes(parts: List(String)) -> Option(String) { + case parts { + [] -> None + _ -> Some(string.join(parts, "; ")) } } -// Stop at the first failure so a rate limit becomes a partial, resumable run. +// Stop at the first run-halting failure (a 429 or a PDS write error) so it +// becomes a partial, resumable run; artist-mint failures skip just the item. fn write_releases( ctx: Context, client: Client, @@ -418,42 +512,61 @@ fn write_releases( entities: catalog_entities.Entities, releases: List(discogs_client.Release), ) -> #( - Int, + WriteStats, dict.Dict(String, promotion.OwnRelease), catalog_entities.Entities, - Int, Option(String), ) { case releases { - [] -> #(0, own, entities, 0, None) + [] -> #(WriteStats(0, 0, 0, 0), own, entities, None) [release, ..rest] -> { - let #(own, entities, entity_failed, written) = + let #(own, entities, result) = write_genesis(ctx, client, session, action, own, entities, release) - case written { - Ok(_) -> { - let #(n, own, entities, ef, note) = + case result { + Wrote(entity_failed, details_missing) -> { + let #(stats, own, entities, note) = write_releases(ctx, client, session, action, own, entities, rest) - #(n + 1, own, entities, entity_failed + ef, note) + #( + WriteStats( + written: stats.written + 1, + entity_failures: stats.entity_failures + entity_failed, + details_missing: stats.details_missing + + bool_to_int(details_missing), + skipped_items: stats.skipped_items, + ), + own, + entities, + note, + ) } - Error(xrpc.BadStatus(status: 429, ..)) -> #( - 0, - own, - entities, - entity_failed, - Some("PDS rate limited; run again later"), - ) - Error(_) -> #( - 0, + SkippedArtistMint -> { + let #(stats, own, entities, note) = + write_releases(ctx, client, session, action, own, entities, rest) + #( + WriteStats(..stats, skipped_items: stats.skipped_items + 1), + own, + entities, + note, + ) + } + Stopped(note, entity_failed) -> #( + WriteStats(0, entity_failed, 0, 0), own, entities, - entity_failed, - Some("PDS write failed; run again to resume"), + Some(note), ) } } } } +fn bool_to_int(b: Bool) -> Int { + case b { + True -> 1 + False -> 0 + } +} + fn write_genesis( ctx: Context, client: Client, @@ -465,45 +578,62 @@ fn write_genesis( ) -> #( dict.Dict(String, promotion.OwnRelease), catalog_entities.Entities, - Int, - Result(repo.CreatedRecord, xrpc.XrpcError), + GenesisResult, ) { - let details = fetch_details(ctx, release) - let cover = best_cover(ctx, client, session, release, details) - let outcome = - promotion.adopt_or_mint( - ctx.catalog, - client, - session, + case fetch_details(ctx, release) { + // A 429 on the enrichment read stops the run before the item is minted, so + // it retries next run rather than minting a permanently unenriched record. + DetailsRateLimited -> #( own, entities, - release, - cover, - details, - ) - case outcome.report.rate_limited { - // Entity write hit a 429: abort before the shelf event so the item retries. - True -> #( - outcome.own, - outcome.entities, - outcome.report.failed, - Error(xrpc.BadStatus( - status: 429, - error: None, - message: None, - body: "entity write rate limited", - )), + Stopped("Discogs rate limited; run again in a minute", 0), ) - False -> - write_genesis_event( - client, - session, - action, - outcome, - release, - cover, - details, - ) + fetched -> { + let details_missing = fetched == DetailsMissing + case details_missing { + True -> + wisp.log_warning( + "discogs details fetch failed for discogs:" + <> int.to_string(release.discogs_id) + <> "; minting without enrichment", + ) + False -> Nil + } + let details = details_option(fetched) + let cover = best_cover(ctx, client, session, release, details) + let outcome = + promotion.adopt_or_mint( + ctx.catalog, + client, + session, + own, + entities, + release, + cover, + details, + ) + case outcome.report.rate_limited, outcome.item_failed { + // Entity write hit a 429: abort before the shelf event so it retries. + True, _ -> #( + outcome.own, + outcome.entities, + Stopped("PDS rate limited; run again later", outcome.report.failed), + ) + // Artist mint failed: withhold the item rather than mint truncated + // credits. Minted artists are kept so the retry only redoes the failure. + False, True -> #(outcome.own, outcome.entities, SkippedArtistMint) + False, False -> + write_genesis_event( + client, + session, + action, + outcome, + release, + cover, + details_missing, + ) + } + } } } @@ -514,12 +644,11 @@ fn write_genesis_event( outcome: promotion.Outcome, release: discogs_client.Release, cover: Option(blob.Blob), - _details: Option(discogs_client.ReleaseDetails), + details_missing: Bool, ) -> #( dict.Dict(String, promotion.OwnRelease), catalog_entities.Entities, - Int, - Result(repo.CreatedRecord, xrpc.XrpcError), + GenesisResult, ) { let own = outcome.own let promoted = outcome.promoted @@ -559,25 +688,46 @@ fn write_genesis_event( source: Some(source), created_at: provenance.now_rfc3339(), ) - #( - own, - outcome.entities, - outcome.report.failed, - event_log.append(client, session, event), - ) + let result = case event_log.append(client, session, event) { + Ok(_) -> + Wrote( + entity_failed: outcome.report.failed, + details_missing: details_missing, + ) + Error(xrpc.BadStatus(status: 429, ..)) -> + Stopped("PDS rate limited; run again later", outcome.report.failed) + Error(_) -> + Stopped("PDS write failed; run again to resume", outcome.report.failed) + } + #(own, outcome.entities, result) } // The widened per-release fetch, shared by the cover and the mint enrichment. +// A 429 is surfaced distinctly so the run stops instead of swallowing it. fn fetch_details( ctx: Context, release: discogs_client.Release, +) -> FetchedDetails { + case + discogs_client.release_details( + ctx.discogs.send, + ctx.discogs.auth, + release.discogs_id, + ) + { + Ok(details) -> GotDetails(details) + Error("rate limited") -> DetailsRateLimited + Error(_) -> DetailsMissing + } +} + +fn details_option( + fetched: FetchedDetails, ) -> Option(discogs_client.ReleaseDetails) { - discogs_client.release_details( - ctx.discogs.send, - ctx.discogs.auth, - release.discogs_id, - ) - |> option.from_result + case fetched { + GotDetails(details) -> Some(details) + DetailsRateLimited | DetailsMissing -> None + } } // Full-res primary image first (from the shared details), then cover, thumb. diff --git a/server/src/at_record_server/discogs_client.gleam b/server/src/at_record_server/discogs_client.gleam index ec3f148..24061a2 100644 --- a/server/src/at_record_server/discogs_client.gleam +++ b/server/src/at_record_server/discogs_client.gleam @@ -295,6 +295,7 @@ pub fn release_details( json.parse(body, release_details_decoder()) |> map_error(fn(_) { "could not parse discogs release" }) |> result.map(warn_if_no_artists(discogs_id, _)) + Ok(#(429, _)) -> Error("rate limited") Ok(#(status, _)) -> Error("discogs returned status " <> int.to_string(status)) } diff --git a/server/src/at_record_server/handlers/browse.gleam b/server/src/at_record_server/handlers/browse.gleam index ce8aef9..051e354 100644 --- a/server/src/at_record_server/handlers/browse.gleam +++ b/server/src/at_record_server/handlers/browse.gleam @@ -37,13 +37,15 @@ pub fn browse(req: Request, ctx: Context) -> Response { Ok(stored) -> browse_domain.ownership(stored) Error(_) -> browse_domain.empty_ownership() } + // Sort before dedup so "newest wins" is a property of our own code, not an + // assumption about how PDSes happen to order listRecords. let rows = ctx.known_users.list() |> list.flat_map(fn(user) { browse_domain.fetch_user_releases(ctx.atproto.client, user) }) - |> browse_domain.dedup |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) + |> browse_domain.dedup |> list.take(100) json.object([#("releases", json.array(rows, encode_row(own, _)))]) |> json.to_string diff --git a/server/src/at_record_server/handlers/shelf.gleam b/server/src/at_record_server/handlers/shelf.gleam index b25d236..70f3cbb 100644 --- a/server/src/at_record_server/handlers/shelf.gleam +++ b/server/src/at_record_server/handlers/shelf.gleam @@ -364,35 +364,39 @@ fn do_add( [form.cover_url, form.thumb_url], ) } - let #(release, source) = - promote_manual(ctx, client, session, form, cover, details) - ShelfEntry( - subject: None, - action:, - snapshot: Some(Snapshot( - title: form.title, - artist_display: form.artist, - year: form.year, - format: form.format, - thumb_url: form.thumb_url, - cover:, - )), - external_ids:, - media_grade: form.media_grade, - sleeve_grade: None, - rating: None, - folder: None, - notes: form.notes, - release:, - price: None, - counterparty: None, - source: Some(source), - created_at: now_rfc3339(), - ) - |> write_event(client, session) + case promote_manual(ctx, client, session, form, cover, details) { + Error(Nil) -> + error_json(502, "artist records could not be published; try again") + Ok(#(release, source)) -> + ShelfEntry( + subject: None, + action:, + snapshot: Some(Snapshot( + title: form.title, + artist_display: form.artist, + year: form.year, + format: form.format, + thumb_url: form.thumb_url, + cover:, + )), + external_ids:, + media_grade: form.media_grade, + sleeve_grade: None, + rating: None, + folder: None, + notes: form.notes, + release:, + price: None, + counterparty: None, + source: Some(source), + created_at: now_rfc3339(), + ) + |> write_event(client, session) + } } -// A manual add with a discogs id joins the shared catalog like an import does. +// A manual add with a discogs id joins the shared catalog like an import +// does; Error means an artist mint failed and the item must not be written. fn promote_manual( ctx: Context, client: Client, @@ -400,10 +404,10 @@ fn promote_manual( form: AddForm, cover: option.Option(blob.Blob), details: option.Option(discogs_client.ReleaseDetails), -) -> #(option.Option(defs.CatalogRef), defs.Source) { +) -> Result(#(option.Option(defs.CatalogRef), defs.Source), Nil) { let parsed = form.discogs_id |> option.then(parse_id) case parsed { - None -> #(None, provenance.source("manual")) + None -> Ok(#(None, provenance.source("manual"))) Some(discogs_id) -> { let seed = discogs_client.Release( @@ -429,16 +433,18 @@ fn promote_manual( cover, details, ) - case outcome.promoted { - Some(p) -> #( - Some(p.release), - provenance.source("manual") - |> provenance.with_record(RepoStrongRef( - uri: p.release.uri, - cid: p.release.cid, - )), - ) - None -> #(None, provenance.source("manual")) + case outcome.item_failed, outcome.promoted { + True, _ -> Error(Nil) + False, Some(p) -> + Ok(#( + Some(p.release), + provenance.source("manual") + |> provenance.with_record(RepoStrongRef( + uri: p.release.uri, + cid: p.release.cid, + )), + )) + False, None -> Ok(#(None, provenance.source("manual"))) } } } diff --git a/server/src/at_record_server/known_users_memory.gleam b/server/src/at_record_server/known_users_memory.gleam index 8e8ecbd..63cc7f6 100644 --- a/server/src/at_record_server/known_users_memory.gleam +++ b/server/src/at_record_server/known_users_memory.gleam @@ -24,11 +24,34 @@ pub fn start() -> Result(Store, actor.StartError) { Nil }, list: fn() { - actor.call(subject, waiting: 1000, sending: fn(reply) { List(reply) }) + case try_call(subject, 1000, List) { + Ok(users) -> users + Error(Nil) -> [] + } }, ) } +// process.call/actor.call panic on timeout or a dead callee; this mirrors process.call's own monitor-based body but returns Error(Nil) instead, so list() can honor its best-effort contract. +fn try_call( + subject: Subject(msg), + timeout: Int, + make_request: fn(Subject(reply)) -> msg, +) -> Result(reply, Nil) { + use callee <- result.try(process.subject_owner(subject)) + let reply_subject = process.new_subject() + let monitor = process.monitor(callee) + process.send(subject, make_request(reply_subject)) + let reply = + process.new_selector() + |> process.select_map(reply_subject, Ok) + |> process.select_specific_monitor(monitor, fn(_down) { Error(Nil) }) + |> process.selector_receive(timeout) + |> result.unwrap(Error(Nil)) + process.demonitor_process(monitor) + reply +} + fn handle( state: Dict(String, KnownUser), msg: Msg, diff --git a/server/src/at_record_server/musicbrainz_client.gleam b/server/src/at_record_server/musicbrainz_client.gleam index c3af440..4e4eb2a 100644 --- a/server/src/at_record_server/musicbrainz_client.gleam +++ b/server/src/at_record_server/musicbrainz_client.gleam @@ -22,32 +22,73 @@ const host = "musicbrainz.org" const user_agent = "at-record/0.1 (https://tangled.org/@mokkenstorm.dev/at-record)" -/// The MusicBrainz release id for a Discogs release: barcode lookup first -/// when a barcode is known, falling back to the Discogs URL relationship. -/// Entirely best-effort; any HTTP failure returns `None` and logs a warning, -/// never raises, never fails the mint. +/// The MusicBrainz release id for a Discogs release: up to 2 plausible +/// barcode candidates tried in order, falling back to the Discogs URL +/// relationship, bounding the worst case at 3 requests. Every request after +/// the first waits 1000ms first, honoring the 1 req/s etiquette between +/// attempts rather than only once up front. Entirely best-effort; any HTTP +/// failure returns `None` and logs a warning, never raises, never fails the +/// mint. pub fn release_mbid( client: Client, barcodes: List(String), discogs_id: Int, ) -> Option(String) { - process.sleep(1000) - case usable_barcodes(barcodes) { - [barcode, ..] -> + barcodes + |> usable_barcodes + |> list.take(2) + |> attempt_barcodes(client, _, discogs_id, True) +} + +fn attempt_barcodes( + client: Client, + candidates: List(String), + discogs_id: Int, + first: Bool, +) -> Option(String) { + maybe_sleep(first) + case candidates { + [] -> by_discogs_url(client, discogs_id) + [barcode, ..rest] -> case by_barcode(client, barcode, discogs_id) { Some(mbid) -> Some(mbid) - None -> by_discogs_url(client, discogs_id) + None -> attempt_barcodes(client, rest, discogs_id, False) } - [] -> by_discogs_url(client, discogs_id) } } -// Discogs barcodes often carry spaces; MusicBrainz stores them bare, and a -// spaced query degrades to fuzzy token matching instead of an exact hit. -fn usable_barcodes(barcodes: List(String)) -> List(String) { +fn maybe_sleep(first: Bool) -> Nil { + case first { + True -> Nil + False -> process.sleep(1000) + } +} + +/// Discogs barcodes often carry spaces (stripped here, since MusicBrainz +/// stores them bare and a spaced query degrades to fuzzy token matching) or +/// outright junk annotations ("Not printed on packaging", "Text on spine: +/// ..."), which are filtered out entirely: only strings that survive +/// de-spacing as all-digit, 8-14 character candidates are worth querying. +/// Exported because this filtering is real policy, not an implementation +/// detail: it bounds how many candidates `release_mbid` will ever try. +pub fn usable_barcodes(barcodes: List(String)) -> List(String) { barcodes |> list.map(fn(b) { string.replace(b, " ", "") }) - |> list.filter(fn(b) { b != "" }) + |> list.filter(plausible_barcode) +} + +fn plausible_barcode(barcode: String) -> Bool { + let length = string.length(barcode) + length >= 8 + && length <= 14 + && string.to_graphemes(barcode) |> list.all(is_digit) +} + +fn is_digit(grapheme: String) -> Bool { + case grapheme { + "0" | "1" | "2" | "3" | "4" | "5" | "6" | "7" | "8" | "9" -> True + _ -> False + } } fn by_barcode( @@ -112,16 +153,25 @@ fn url_relation_request(discogs_id: Int) -> Request(String) { |> request.set_header("user-agent", user_agent) } -/// The first `releases[].id` from a barcode search response, pure so it can -/// be tested against a fixture without a network call. +/// The id of the first exact-match (`score` 100) release in a barcode search +/// response, pure so it can be tested against a fixture without a network +/// call. A barcode search can return fuzzy token matches scored below 100; +/// only an exact match is trusted here, everything else is treated as a miss +/// so the URL-relation fallback (a direct relationship lookup, inherently +/// exact) gets its chance. pub fn parse_barcode_response(body: String) -> Option(String) { let decoder = { - use releases <- decode.field("releases", decode.list(id_decoder())) + use releases <- decode.field( + "releases", + decode.list(scored_release_decoder()), + ) decode.success(releases) } json.parse(body, decoder) |> result.unwrap([]) + |> list.filter(fn(release) { release.score == 100 }) |> list.first + |> result.map(fn(release) { release.id }) |> option.from_result } @@ -144,6 +194,16 @@ fn id_decoder() -> decode.Decoder(String) { decode.success(id) } +type ScoredRelease { + ScoredRelease(id: String, score: Int) +} + +fn scored_release_decoder() -> decode.Decoder(ScoredRelease) { + use id <- decode.field("id", decode.string) + use score <- decode.optional_field("score", 0, decode.int) + decode.success(ScoredRelease(id:, score:)) +} + fn relation_decoder() -> decode.Decoder(Option(String)) { use release_id <- decode.optional_field( "release", diff --git a/server/src/at_record_server/promotion.gleam b/server/src/at_record_server/promotion.gleam index 301204d..5ab69df 100644 --- a/server/src/at_record_server/promotion.gleam +++ b/server/src/at_record_server/promotion.gleam @@ -1,8 +1,11 @@ //// Adopt-or-mint at write time: resolve the canonical `catalog.release` for a //// Discogs release. A Constellation backlink hit on the natural-key URL //// adopts the existing record; a miss mints our own from the seed data, -//// which seeds the network for the next user. Best-effort throughout: any -//// failure logs and the entry is simply written without a `release` ref. +//// which seeds the network for the next user. Mostly best-effort: a missing +//// backlink or a failed release write just omits the `release` ref. The one +//// hard failure is a failed artist mint: the release would ship with knowingly +//// truncated credits, so `item_failed` is raised and the caller skips the item +//// instead of minting an incomplete immutable record. import at_record/gen/catalog/release as catalog_release import at_record/gen/defs @@ -32,12 +35,15 @@ pub type Promoted { /// The result of an adopt-or-mint: the (possibly grown) own-release map and /// entity dicts, the outcome, and the entity-write report the import surfaces. +/// `item_failed` is set when an artist mint failed and the release was withheld +/// rather than minted with truncated credits; the caller must skip the item. pub type Outcome { Outcome( own: Dict(String, OwnRelease), entities: catalog_entities.Entities, promoted: Option(Promoted), report: catalog_entities.MintReport, + item_failed: Bool, ) } @@ -126,6 +132,7 @@ pub fn adopt_or_mint( entities:, promoted: Some(Promoted(release: existing.ref, origin: "promotion")), report: catalog_entities.empty_report, + item_failed: False, ) Error(Nil) -> case discover(deps, session, id) { @@ -135,9 +142,10 @@ pub fn adopt_or_mint( entities:, promoted: Some(Promoted(release: ref, origin: "adoption")), report: catalog_entities.empty_report, + item_failed: False, ) None -> { - let #(entities, minted, report) = + let #(entities, minted, report, artist_failed) = mint(deps, client, session, entities, release, cover, details, None) case minted { Some(ref) -> @@ -153,8 +161,16 @@ pub fn adopt_or_mint( entities:, promoted: Some(Promoted(release: ref, origin: "promotion")), report:, + item_failed: False, + ) + None -> + Outcome( + own:, + entities:, + promoted: None, + report:, + item_failed: artist_failed, ) - None -> Outcome(own:, entities:, promoted: None, report:) } } } @@ -211,6 +227,8 @@ pub fn parse_at_uri(uri: String) -> Option(#(String, String, String)) { // Artists and genres are minted before the release so a rate-limited entity // write aborts the whole item (nothing seen), letting the next run retry both. +// A non-429 artist-mint failure also withholds the release (fourth tuple slot +// True), since minting it would drop the failed credits from an immutable record. fn mint( deps: Deps, client: Client, @@ -224,18 +242,26 @@ fn mint( catalog_entities.Entities, Option(defs.CatalogRef), catalog_entities.MintReport, + Bool, ) { let id = int.to_string(release.discogs_id) let artists = details |> option.map(fn(d) { d.artists }) |> option.unwrap([]) let #(artist_dict, credits, artist_report) = catalog_entities.ensure_artists(client, session, entities.artists, artists) - case artist_report.rate_limited { - True -> #( + case artist_report.rate_limited, artist_report.failed > 0 { + True, _ -> #( + catalog_entities.Entities(..entities, artists: artist_dict), + None, + artist_report, + False, + ) + False, True -> #( catalog_entities.Entities(..entities, artists: artist_dict), None, artist_report, + True, ) - False -> { + False, False -> { let names = details |> option.map(fn(d) { list.append(d.genres, d.styles) }) @@ -246,7 +272,7 @@ fn mint( catalog_entities.Entities(artists: artist_dict, genres: genre_dict) let report = catalog_entities.combine(artist_report, genre_report) case genre_report.rate_limited { - True -> #(entities, None, report) + True -> #(entities, None, report, False) False -> { let mbid = details @@ -293,6 +319,7 @@ fn mint( external_ids: None, )), report, + False, ) Error(e) -> { wisp.log_warning( @@ -301,7 +328,7 @@ fn mint( <> ": " <> xrpc_error(e), ) - #(entities, None, report) + #(entities, None, report, False) } } } diff --git a/server/test/browse_test.gleam b/server/test/browse_test.gleam index a4d448f..df9532e 100644 --- a/server/test/browse_test.gleam +++ b/server/test/browse_test.gleam @@ -4,6 +4,7 @@ import at_record/storage.{type StoredItem, StoredItem} import at_record_server/browse.{type BrowseRow, BrowseRow} import gleam/list import gleam/option.{type Option, None, Some} +import gleam/string fn row( uri uri: String, @@ -87,6 +88,21 @@ pub fn dedup_keeps_all_rows_without_discogs_id_test() { assert list.length(kept) == 2 } +pub fn dedup_after_sort_keeps_newest_created_at_regardless_of_input_order_test() { + let older = + row(uri: "at://old/1", discogs_id: Some("42"), created_at: "2020-01-01") + let newer = + row(uri: "at://new/1", discogs_id: Some("42"), created_at: "2024-01-01") + let sort_desc_and_dedup = fn(rows: List(BrowseRow)) { + rows + |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) + |> browse.dedup + |> list.map(fn(r) { r.uri }) + } + assert sort_desc_and_dedup([older, newer]) == ["at://new/1"] + assert sort_desc_and_dedup([newer, older]) == ["at://new/1"] +} + pub fn year_from_released_parses_leading_four_digits_test() { assert browse.year_from_released(Some("1991-03-27")) == Some(1991) assert browse.year_from_released(Some("1991")) == Some(1991) @@ -98,6 +114,11 @@ pub fn year_from_released_rejects_absent_or_non_numeric_test() { assert browse.year_from_released(Some("abcd")) == None } +pub fn year_from_released_rejects_short_strings_test() { + assert browse.year_from_released(Some("197")) == None + assert browse.year_from_released(Some("91")) == None +} + pub fn ownership_matches_by_release_uri_test() { let stored = [ genesis( diff --git a/server/test/catalog_import_test.gleam b/server/test/catalog_import_test.gleam new file mode 100644 index 0000000..7c44dd3 --- /dev/null +++ b/server/test/catalog_import_test.gleam @@ -0,0 +1,36 @@ +import at_record_server/catalog_import +import gleam/option.{None, Some} + +pub fn run_note_none_when_clean_test() { + assert catalog_import.run_note(None, 0, 0, 0) == None +} + +// Finding 2: a stop reason must not swallow accumulated entity-failure counts. +pub fn run_note_composes_stop_with_entity_failures_test() { + assert catalog_import.run_note( + Some("Discogs rate limited; run again in a minute"), + 2, + 0, + 0, + ) + == Some( + "Discogs rate limited; run again in a minute; 2 artist/genre records failed; run again to fill them", + ) +} + +pub fn run_note_details_missing_only_test() { + assert catalog_import.run_note(None, 0, 4, 0) + == Some("4 releases minted without Discogs details; amend later to enrich") +} + +pub fn run_note_skipped_items_only_test() { + assert catalog_import.run_note(None, 0, 0, 3) + == Some("3 items skipped: artist mint failed; run again to retry") +} + +pub fn run_note_composes_all_clauses_test() { + assert catalog_import.run_note(Some("stop"), 1, 3, 2) + == Some( + "stop; 1 artist/genre records failed; run again to fill them; 3 releases minted without Discogs details; amend later to enrich; 2 items skipped: artist mint failed; run again to retry", + ) +} diff --git a/server/test/discogs_client_test.gleam b/server/test/discogs_client_test.gleam index 525c6ff..cb35fe7 100644 --- a/server/test/discogs_client_test.gleam +++ b/server/test/discogs_client_test.gleam @@ -142,3 +142,17 @@ pub fn wantlist_page_rate_limit_test() { assert discogs_client.wantlist_page(send, "OAuth ...", "someone", 1) == Error("rate limited") } + +// --- release details --- + +pub fn release_details_rate_limit_test() { + let send = fn(_req) { Ok(#(429, "slow down")) } + assert discogs_client.release_details(send, None, 4577) + == Error("rate limited") +} + +pub fn release_details_other_status_is_error_test() { + let send = fn(_req) { Ok(#(500, "boom")) } + assert discogs_client.release_details(send, None, 4577) + == Error("discogs returned status 500") +} diff --git a/server/test/musicbrainz_client_test.gleam b/server/test/musicbrainz_client_test.gleam index 72e1c42..c8f4101 100644 --- a/server/test/musicbrainz_client_test.gleam +++ b/server/test/musicbrainz_client_test.gleam @@ -9,6 +9,8 @@ const barcode_response = "{\"created\":\"2026-07-09T23:05:03.885Z\",\"count\":2, const empty_barcode_response = "{\"created\":\"2026-07-09T23:05:03.885Z\",\"count\":0,\"offset\":0,\"releases\":[]}" +const low_score_barcode_response = "{\"created\":\"2026-07-09T23:05:03.885Z\",\"count\":1,\"offset\":0,\"releases\":[{\"id\":\"9ec684cf-7529-300f-b6c4-6c56086dfe93\",\"score\":87,\"title\":\"Black Holes and Revelations\",\"artist-credit\":[{\"name\":\"Muse\"}],\"date\":\"2009-08-18\",\"country\":\"US\",\"barcode\":\"825646350919\"}]}" + const url_relation_response = "{\"id\":\"10830f36-6f49-4fa7-a18b-aff1143ac5b3\",\"resource\":\"https://www.discogs.com/release/1890217\",\"relations\":[{\"target-type\":\"release\",\"type\":\"discogs\",\"direction\":\"backward\",\"release\":{\"id\":\"9ec684cf-7529-300f-b6c4-6c56086dfe93\",\"title\":\"Black Holes and Revelations\",\"barcode\":\"825646350919\"}}]}" const empty_url_relation_response = "{\"id\":\"10830f36-6f49-4fa7-a18b-aff1143ac5b3\",\"resource\":\"https://www.discogs.com/release/999999\",\"relations\":[]}" @@ -27,10 +29,26 @@ pub fn parse_barcode_response_none_when_no_releases_test() { == None } +pub fn parse_barcode_response_none_when_score_below_100_test() { + assert musicbrainz_client.parse_barcode_response(low_score_barcode_response) + == None +} + pub fn parse_barcode_response_none_on_malformed_json_test() { assert musicbrainz_client.parse_barcode_response("not json") == None } +pub fn usable_barcodes_filters_junk_and_despaces_test() { + assert musicbrainz_client.usable_barcodes([ + "825 646 350919", + "Not printed on packaging", + "Text on spine: 5099930291623", + "", + "123", + ]) + == ["825646350919"] +} + pub fn parse_url_relation_response_picks_release_id_test() { assert musicbrainz_client.parse_url_relation_response(url_relation_response) == Some("9ec684cf-7529-300f-b6c4-6c56086dfe93") diff --git a/server/test/promotion_test.gleam b/server/test/promotion_test.gleam index f9a9f7b..2d6f6d0 100644 --- a/server/test/promotion_test.gleam +++ b/server/test/promotion_test.gleam @@ -1,3 +1,4 @@ +import at_record/gen/catalog/artist as catalog_artist import at_record/gen/catalog/release as catalog_release import at_record/gen/defs.{CatalogRef} import at_record_server/catalog_deps.{type Deps, Deps} @@ -123,6 +124,26 @@ fn dead_client() -> xrpc.Client { xrpc.Client(send: fn(_) { panic as "no PDS writes expected" }) } +/// A stub PDS that fails every artist createRecord with a non-429 error and +/// answers other collections normally, so an artist mint failure can be forced. +fn artist_failing_client() -> xrpc.Client { + xrpc.Client(send: fn(req) { + let body = bit_array.to_string(req.body) |> result.unwrap("") + let collection = + json.parse(body, decode.at(["collection"], decode.string)) + |> result.unwrap("") + case collection == catalog_artist.collection { + True -> + Ok(response.new(500) |> response.set_body(bit_array.from_string("no"))) + False -> { + let uri = "at://" <> own_did <> "/" <> collection <> "/1" + let out = "{\"uri\": \"" <> uri <> "\", \"cid\": \"bafy\"}" + Ok(response.new(200) |> response.set_body(bit_array.from_string(out))) + } + } + }) +} + /// A stub PDS that answers createRecord with a fixed uri/cid. fn minting_client(uri: String) -> xrpc.Client { xrpc.Client(send: fn(_req) { @@ -424,6 +445,27 @@ pub fn mint_with_unresolved_mbid_keeps_only_discogs_external_id_test() { let assert Some(Promoted(origin: "promotion", ..)) = outcome.promoted } +pub fn failed_artist_mint_withholds_release_test() { + let outcome = + promotion.adopt_or_mint( + no_backlinks_deps(), + artist_failing_client(), + session(), + dict.new(), + no_entities(), + release(42), + None, + Some(details()), + ) + // No release minted, so the item never becomes own or gets a shelf genesis. + assert outcome.promoted == None + assert outcome.item_failed == True + assert dict.size(outcome.own) == 0 + // The caller can act on the report: a counted, non-rate-limited failure. + assert outcome.report.failed >= 1 + assert outcome.report.rate_limited == False +} + pub fn parse_at_uri_extracts_did_collection_rkey_test() { let assert Some(#(did, collection, rkey)) = promotion.parse_at_uri(