diff --git a/server/test/catalog_index_postgres_test.gleam b/server/test/catalog_index_postgres_test.gleam new file mode 100644 --- /dev/null +++ b/server/test/catalog_index_postgres_test.gleam @@ -0,0 +1,108 @@ +//// Round-trip checks for the Postgres backend, gated on `DATABASE_URL`: a +//// no-op (still a "pass") when unset, so `make test` stays green without a +//// real database, but a real regression net for the riskiest bit of +//// `catalog_index_postgres` (the `genres`/`styles` array encode/decode) +//// when run against a live Postgres, e.g. `DATABASE_URL=... gleam test`. + +import at_record_server/browse.{BrowseRow} +import at_record_server/catalog_index.{Adoption, Edit} +import at_record_server/catalog_index_postgres +import at_record_server/oauth/sessions_postgres +import envoy +import gleam/dict +import gleam/option.{None, Some} + +fn with_store(run: fn(catalog_index.Store) -> Nil) -> Nil { + case envoy.get("DATABASE_URL") { + Error(Nil) -> Nil + Ok(url) -> { + let assert Ok(conn) = sessions_postgres.connect_pool(url) + let assert Ok(store) = catalog_index_postgres.table_store(conn) + run(store) + } + } +} + +fn row(uri: String) -> browse.BrowseRow { + BrowseRow( + uri:, + cid: "bafyrow", + title: "Postgres Round Trip", + artist_display: Some("Some Artist"), + genres: ["Electronic", "Ambient"], + styles: ["Drone"], + released: Some("1994"), + country: Some("US"), + cover: None, + thumb_url: None, + discogs_id: Some("42"), + created_at: "2024-01-01T00:00:00Z", + publisher_did: "did:plc:pgtest", + publisher_handle: "pgtest.test", + publisher_pds: "https://pds.test", + supersedes: Some("at://did:plc:pgtest/x/1"), + based_on: None, + ) +} + +pub fn release_round_trip_preserves_arrays_and_nullables_test() { + use store <- with_store() + let uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.release/pg1" + store.upsert(row(uri)) + let assert Ok(found) = store.list() |> find(fn(r) { r.uri == uri }) + assert found.genres == ["Electronic", "Ambient"] + assert found.styles == ["Drone"] + assert found.artist_display == Some("Some Artist") + assert found.supersedes == Some("at://did:plc:pgtest/x/1") + assert found.based_on == None + store.delete(uri) + let assert Error(Nil) = store.list() |> find(fn(r) { r.uri == uri }) + Nil +} + +pub fn adoption_and_edit_round_trip_test() { + use store <- with_store() + let release = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.release/pg2" + let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg-e1" + store.upsert_adoption(Adoption( + entry_uri:, + did: "did:plc:pgtest", + release_uri: release, + status: "acquired", + created_at: "2024-01-01T00:00:00Z", + )) + assert store.adoption_count(release) == 1 + assert dict.get(store.adoption_counts(), release) == Ok(1) + store.delete_adoption(entry_uri) + assert store.adoption_count(release) == 0 + + let edit_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.edit/pg-x1" + store.upsert_edit(Edit( + edit_uri:, + did: "did:plc:pgtest", + subject_uri: release, + entity: "release", + created_at: "2024-01-01T00:00:00Z", + )) + let assert [found] = store.edits_for(release) + assert found.entity == "release" + store.delete_edit(edit_uri) + assert store.edits_for(release) == [] +} + +pub fn cursor_round_trip_test() { + use store <- with_store() + store.save_cursor(1_700_000_000_123_456) + assert store.load_cursor() == Some(1_700_000_000_123_456) +} + +fn find(list: List(a), pred: fn(a) -> Bool) -> Result(a, Nil) { + case list { + [] -> Error(Nil) + [x, ..rest] -> + case pred(x) { + True -> Ok(x) + False -> find(rest, pred) + } + } +} diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam new file mode 100644 --- /dev/null +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -0,0 +1,383 @@ +//// Postgres catalog_index.Store backend (pog), over four create-if-not-exists +//// tables sharing one pool with the other stores: catalog_releases (browse +//// rows), catalog_adoptions (shelf.entry events with a release ref), +//// catalog_edits (catalog.edit events), and jetstream_cursor (a single row +//// tracking Jetstream replay progress). This is derived, rebuildable data, +//// so writes swallow-and-log rather than fail the caller. + +import at_record_server/browse.{type BrowseRow, BrowseRow} +import at_record_server/catalog_index.{ + type Adoption, type Edit, type Store, Edit, Store, +} +import atproto/blob.{Blob} +import gleam/dict.{type Dict} +import gleam/dynamic/decode +import gleam/option.{type Option, None, Some} +import gleam/result +import pog +import wisp + +const cursor_id = "default" + +pub fn table_store(conn: pog.Connection) -> Result(Store, String) { + use _ <- result.try(migrate(conn)) + Ok( + Store( + upsert: fn(row) { upsert_release(conn, row) }, + delete: fn(uri) { delete_release(conn, uri) }, + list: fn() { list_releases(conn) }, + upsert_adoption: fn(adoption) { upsert_adoption(conn, adoption) }, + delete_adoption: fn(entry_uri) { delete_adoption(conn, entry_uri) }, + upsert_edit: fn(edit) { upsert_edit(conn, edit) }, + delete_edit: fn(edit_uri) { delete_edit(conn, edit_uri) }, + adoption_count: fn(release_uri) { adoption_count(conn, release_uri) }, + adoption_counts: fn() { adoption_counts(conn) }, + edits_for: fn(subject_uri) { edits_for(conn, subject_uri) }, + save_cursor: fn(time_us) { save_cursor(conn, time_us) }, + load_cursor: fn() { load_cursor(conn) }, + ), + ) +} + +fn migrate(conn: pog.Connection) -> Result(Nil, String) { + use _ <- result.try( + pog.query( + "create table if not exists catalog_releases ( + uri text primary key, + cid text not null, + title text not null, + artist_display text, + genres text[] not null default '{}', + styles text[] not null default '{}', + released text, + country text, + cover_cid text, + thumb_url text, + discogs_id text, + created_at text not null, + publisher_did text not null, + publisher_handle text not null, + publisher_pds text not null, + supersedes text, + based_on text + )", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_releases migration failed"), + ) + use _ <- result.try( + pog.query( + "create table if not exists catalog_adoptions ( + entry_uri text primary key, + did text not null, + release_uri text not null, + status text not null, + created_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_adoptions migration failed"), + ) + use _ <- result.try( + pog.query( + "create table if not exists catalog_edits ( + edit_uri text primary key, + did text not null, + subject_uri text not null, + entity text not null, + created_at text not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_edits migration failed"), + ) + pog.query( + "create table if not exists jetstream_cursor ( + id text primary key, + time_us bigint not null + )", + ) + |> pog.execute(conn) + |> result.replace_error("jetstream_cursor migration failed") + |> result.map(fn(_) { Nil }) +} + +fn upsert_release(conn: pog.Connection, row: BrowseRow) -> Nil { + let outcome = + pog.query( + "insert into catalog_releases + (uri, cid, title, artist_display, genres, styles, released, country, + cover_cid, thumb_url, discogs_id, created_at, publisher_did, + publisher_handle, publisher_pds, supersedes, based_on) + values ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16, $17) + on conflict (uri) do update set + cid = excluded.cid, + title = excluded.title, + artist_display = excluded.artist_display, + genres = excluded.genres, + styles = excluded.styles, + released = excluded.released, + country = excluded.country, + cover_cid = excluded.cover_cid, + thumb_url = excluded.thumb_url, + discogs_id = excluded.discogs_id, + created_at = excluded.created_at, + publisher_did = excluded.publisher_did, + publisher_handle = excluded.publisher_handle, + publisher_pds = excluded.publisher_pds, + supersedes = excluded.supersedes, + based_on = excluded.based_on", + ) + |> pog.parameter(pog.text(row.uri)) + |> pog.parameter(pog.text(row.cid)) + |> pog.parameter(pog.text(row.title)) + |> pog.parameter(pog.nullable(pog.text, row.artist_display)) + |> pog.parameter(pog.array(pog.text, row.genres)) + |> pog.parameter(pog.array(pog.text, row.styles)) + |> pog.parameter(pog.nullable(pog.text, row.released)) + |> pog.parameter(pog.nullable(pog.text, row.country)) + |> pog.parameter(pog.nullable( + pog.text, + row.cover |> option.map(fn(b) { b.cid }), + )) + |> pog.parameter(pog.nullable(pog.text, row.thumb_url)) + |> pog.parameter(pog.nullable(pog.text, row.discogs_id)) + |> pog.parameter(pog.text(row.created_at)) + |> pog.parameter(pog.text(row.publisher_did)) + |> pog.parameter(pog.text(row.publisher_handle)) + |> pog.parameter(pog.text(row.publisher_pds)) + |> pog.parameter(pog.nullable(pog.text, row.supersedes)) + |> pog.parameter(pog.nullable(pog.text, row.based_on)) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning("catalog_releases upsert failed for " <> row.uri) + } +} + +fn delete_release(conn: pog.Connection, uri: String) -> Nil { + let _ = + pog.query("delete from catalog_releases where uri = $1") + |> pog.parameter(pog.text(uri)) + |> pog.execute(conn) + Nil +} + +fn release_row_decoder() -> decode.Decoder(BrowseRow) { + use uri <- decode.field("uri", decode.string) + use cid <- decode.field("cid", decode.string) + use title <- decode.field("title", decode.string) + use artist_display <- decode.field( + "artist_display", + decode.optional(decode.string), + ) + use genres <- decode.field("genres", decode.list(decode.string)) + use styles <- decode.field("styles", decode.list(decode.string)) + use released <- decode.field("released", decode.optional(decode.string)) + use country <- decode.field("country", decode.optional(decode.string)) + use cover_cid <- decode.field("cover_cid", decode.optional(decode.string)) + use thumb_url <- decode.field("thumb_url", decode.optional(decode.string)) + use discogs_id <- decode.field("discogs_id", decode.optional(decode.string)) + use created_at <- decode.field("created_at", decode.string) + use publisher_did <- decode.field("publisher_did", decode.string) + use publisher_handle <- decode.field("publisher_handle", decode.string) + use publisher_pds <- decode.field("publisher_pds", decode.string) + use supersedes <- decode.field("supersedes", decode.optional(decode.string)) + use based_on <- decode.field("based_on", decode.optional(decode.string)) + decode.success(BrowseRow( + uri:, + cid:, + title:, + artist_display:, + genres:, + styles:, + released:, + country:, + // Only the cid ever gets read back off a browse row's cover (to build + // the same-origin proxy URL), so a placeholder mime/size round-trips + // fine; see handlers/browse.gleam's `encode_row`. + cover: cover_cid + |> option.map(fn(cid) { Blob(cid:, mime_type: "", size: 0) }), + thumb_url:, + discogs_id:, + created_at:, + publisher_did:, + publisher_handle:, + publisher_pds:, + supersedes:, + based_on:, + )) +} + +fn list_releases(conn: pog.Connection) -> List(BrowseRow) { + case + pog.query( + "select uri, cid, title, artist_display, genres, styles, released, + country, cover_cid, thumb_url, discogs_id, created_at, + publisher_did, publisher_handle, publisher_pds, supersedes, based_on + from catalog_releases", + ) + |> pog.returning(release_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + +fn upsert_adoption(conn: pog.Connection, adoption: Adoption) -> Nil { + let outcome = + pog.query( + "insert into catalog_adoptions (entry_uri, did, release_uri, status, created_at) + values ($1, $2, $3, $4, $5) + on conflict (entry_uri) do update set + did = excluded.did, + release_uri = excluded.release_uri, + status = excluded.status, + created_at = excluded.created_at", + ) + |> pog.parameter(pog.text(adoption.entry_uri)) + |> pog.parameter(pog.text(adoption.did)) + |> pog.parameter(pog.text(adoption.release_uri)) + |> pog.parameter(pog.text(adoption.status)) + |> pog.parameter(pog.text(adoption.created_at)) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning( + "catalog_adoptions upsert failed for " <> adoption.entry_uri, + ) + } +} + +fn delete_adoption(conn: pog.Connection, entry_uri: String) -> Nil { + let _ = + pog.query("delete from catalog_adoptions where entry_uri = $1") + |> pog.parameter(pog.text(entry_uri)) + |> pog.execute(conn) + Nil +} + +fn adoption_count(conn: pog.Connection, release_uri: String) -> Int { + let row = { + use count <- decode.field("count", decode.int) + decode.success(count) + } + case + pog.query( + "select count(*)::int as count from catalog_adoptions where release_uri = $1", + ) + |> pog.parameter(pog.text(release_uri)) + |> pog.returning(row) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [count, ..], ..)) -> count + _ -> 0 + } +} + +fn adoption_counts(conn: pog.Connection) -> Dict(String, Int) { + let row = { + use release_uri <- decode.field("release_uri", decode.string) + use count <- decode.field("count", decode.int) + decode.success(#(release_uri, count)) + } + case + pog.query( + "select release_uri, count(*)::int as count from catalog_adoptions group by release_uri", + ) + |> pog.returning(row) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> dict.from_list(rows) + Error(_) -> dict.new() + } +} + +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) + on conflict (edit_uri) do update set + did = excluded.did, + subject_uri = excluded.subject_uri, + entity = excluded.entity, + created_at = excluded.created_at", + ) + |> 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.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> + wisp.log_warning("catalog_edits upsert failed for " <> edit.edit_uri) + } +} + +fn delete_edit(conn: pog.Connection, edit_uri: String) -> Nil { + let _ = + pog.query("delete from catalog_edits where edit_uri = $1") + |> pog.parameter(pog.text(edit_uri)) + |> pog.execute(conn) + Nil +} + +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.parameter(pog.text(subject_uri)) + |> pog.returning(row) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + +fn save_cursor(conn: pog.Connection, time_us: Int) -> Nil { + let outcome = + pog.query( + "insert into jetstream_cursor (id, time_us) values ($1, $2) + on conflict (id) do update set time_us = excluded.time_us", + ) + |> pog.parameter(pog.text(cursor_id)) + |> pog.parameter(pog.int(time_us)) + |> pog.execute(conn) + case outcome { + Ok(_) -> Nil + Error(_) -> wisp.log_warning("jetstream_cursor save failed") + } +} + +fn load_cursor(conn: pog.Connection) -> Option(Int) { + let row = { + use time_us <- decode.field("time_us", decode.int) + decode.success(time_us) + } + case + pog.query("select time_us from jetstream_cursor where id = $1") + |> pog.parameter(pog.text(cursor_id)) + |> pog.returning(row) + |> pog.execute(conn) + { + Ok(pog.Returned(rows: [time_us, ..], ..)) -> Some(time_us) + _ -> None + } +}