From b6b42db564acde45b92175dcbed3f01ecbe3e934 Mon Sep 17 00:00:00 2001 From: Niels Mokkenstorm Date: Thu, 16 Jul 2026 22:57:18 +0200 Subject: [PATCH] refactor(server): split catalog_index.Store into grouped sub-records Store's 12 flat fn fields (upsert/delete/list, upsert_adoption/..., etc.) grouped into ReleaseOps/AdoptionOps/EditOps/CursorOps by the entity each operates on, so call sites read as store.adoptions.count(uri) rather than a wall of same-shaped top-level fields. Both backends, jetstream_consumer, the boot/login backfills, wiring, and every test fake updated. --- server/src/at_record_server.gleam | 7 +- .../src/at_record_server/catalog_index.gleam | 173 ++++++++++-------- .../catalog_index_postgres.gleam | 33 ++-- .../src/at_record_server/handlers/oauth.gleam | 2 +- .../at_record_server/jetstream_consumer.gleam | 16 +- server/src/at_record_server/wiring.gleam | 4 +- server/test/catalog_index_postgres_test.gleam | 46 ++--- server/test/catalog_index_test.gleam | 80 ++++---- server/test/jetstream_consumer_test.gleam | 16 +- server/test/oauth_test.gleam | 2 +- server/test/support.gleam | 29 +-- 11 files changed, 227 insertions(+), 181 deletions(-) diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index ed1e5c5..100e0a1 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -96,7 +96,8 @@ fn backfill_catalog_index( let users = known_users.list() users |> list.each(fn(user) { - browse.fetch_user_releases(client, user) |> list.each(index.upsert) + browse.fetch_user_releases(client, user) + |> list.each(index.releases.upsert) backfill_adoptions(client, user, index) backfill_edits(client, user, index) }) @@ -119,7 +120,7 @@ fn backfill_adoptions( Ok(entry) -> case entry.release { Some(ref) -> - index.upsert_adoption(catalog_index.Adoption( + index.adoptions.upsert(catalog_index.Adoption( entry_uri: uri, did: user.did, release_uri: ref.uri, @@ -145,7 +146,7 @@ fn backfill_edits( Ok(edit) -> case edit.subject { Some(ref) -> - index.upsert_edit(catalog_index.Edit( + index.edits.upsert(catalog_index.Edit( edit_uri: uri, did: user.did, subject_uri: ref.uri, diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index e0de4c1..d128259 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -2,9 +2,9 @@ //// (shelf.entry events that carry a release ref), and edits, plus the //// Jetstream replay cursor. Backends live in this module (in-memory, dev //// fallback, actor over a Dict) and in catalog_index_postgres (durable); -//// wiring picks one at the composition root based on DATABASE_URL. Kept -//// additive: `upsert`/`delete`/`list` signatures never change so existing -//// callers compile untouched. +//// wiring picks one at the composition root based on DATABASE_URL. `Store` +//// groups its operations by the entity they act on, so a call site reads +//// naturally: `store.adoptions.count(uri)`. import at_record_server/catalog/row.{type BrowseRow} import at_record_server/parallel @@ -55,20 +55,41 @@ pub type Edit { ) } -pub type Store { - Store( +pub type ReleaseOps { + ReleaseOps( upsert: fn(BrowseRow) -> Nil, delete: fn(String) -> Nil, list: fn() -> List(BrowseRow), - upsert_adoption: fn(Adoption) -> Nil, - delete_adoption: fn(String) -> Nil, - upsert_edit: fn(Edit) -> Nil, - delete_edit: fn(String) -> Nil, - adoption_count: fn(String) -> Int, - adoption_counts: fn() -> Dict(String, Int), - edits_for: fn(String) -> List(Edit), - save_cursor: fn(Int) -> Nil, - load_cursor: fn() -> Option(Int), + ) +} + +pub type AdoptionOps { + AdoptionOps( + upsert: fn(Adoption) -> Nil, + delete: fn(String) -> Nil, + count: fn(String) -> Int, + counts: fn() -> Dict(String, Int), + ) +} + +pub type EditOps { + EditOps( + upsert: fn(Edit) -> Nil, + delete: fn(String) -> Nil, + for_subject: fn(String) -> List(Edit), + ) +} + +pub type CursorOps { + CursorOps(save: fn(Int) -> Nil, load: fn() -> Option(Int)) +} + +pub type Store { + Store( + releases: ReleaseOps, + adoptions: AdoptionOps, + edits: EditOps, + cursor: CursorOps, ) } @@ -111,64 +132,72 @@ pub fn start() -> Result(Store, actor.StartError) { ) let subject = started.data Store( - upsert: fn(row) { - process.send(subject, Upsert(row)) - Nil - }, - delete: fn(uri) { - process.send(subject, Delete(uri)) - Nil - }, - list: fn() { - case parallel.try_call(subject, 1000, List) { - Ok(rows) -> rows - Error(Nil) -> [] - } - }, - upsert_adoption: fn(adoption) { - process.send(subject, UpsertAdoption(adoption)) - Nil - }, - delete_adoption: fn(entry_uri) { - process.send(subject, DeleteAdoption(entry_uri)) - Nil - }, - upsert_edit: fn(edit) { - process.send(subject, UpsertEdit(edit)) - Nil - }, - delete_edit: fn(edit_uri) { - process.send(subject, DeleteEdit(edit_uri)) - Nil - }, - adoption_count: fn(release_uri) { - case parallel.try_call(subject, 1000, AdoptionCount(release_uri, _)) { - Ok(count) -> count - Error(Nil) -> 0 - } - }, - adoption_counts: fn() { - case parallel.try_call(subject, 1000, AdoptionCounts) { - Ok(counts) -> counts - Error(Nil) -> dict.new() - } - }, - edits_for: fn(subject_uri) { - case parallel.try_call(subject, 1000, EditsFor(subject_uri, _)) { - Ok(edits) -> edits - Error(Nil) -> [] - } - }, - save_cursor: fn(time_us) { - process.send(subject, SaveCursor(time_us)) - Nil - }, - load_cursor: fn() { - case parallel.try_call(subject, 1000, LoadCursor) { - Ok(cursor) -> cursor - Error(Nil) -> option.None - } - }, + releases: ReleaseOps( + upsert: fn(row) { + process.send(subject, Upsert(row)) + Nil + }, + delete: fn(uri) { + process.send(subject, Delete(uri)) + Nil + }, + list: fn() { + case parallel.try_call(subject, 1000, List) { + Ok(rows) -> rows + Error(Nil) -> [] + } + }, + ), + adoptions: AdoptionOps( + upsert: fn(adoption) { + process.send(subject, UpsertAdoption(adoption)) + Nil + }, + delete: fn(entry_uri) { + process.send(subject, DeleteAdoption(entry_uri)) + Nil + }, + count: fn(release_uri) { + case parallel.try_call(subject, 1000, AdoptionCount(release_uri, _)) { + Ok(count) -> count + Error(Nil) -> 0 + } + }, + counts: fn() { + case parallel.try_call(subject, 1000, AdoptionCounts) { + Ok(counts) -> counts + Error(Nil) -> dict.new() + } + }, + ), + edits: EditOps( + upsert: fn(edit) { + process.send(subject, UpsertEdit(edit)) + Nil + }, + delete: fn(edit_uri) { + process.send(subject, DeleteEdit(edit_uri)) + Nil + }, + for_subject: fn(subject_uri) { + case parallel.try_call(subject, 1000, EditsFor(subject_uri, _)) { + Ok(edits) -> edits + Error(Nil) -> [] + } + }, + ), + cursor: CursorOps( + save: fn(time_us) { + process.send(subject, SaveCursor(time_us)) + Nil + }, + load: fn() { + case parallel.try_call(subject, 1000, LoadCursor) { + Ok(cursor) -> cursor + Error(Nil) -> option.None + } + }, + ), ) } diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam index ceca66a..ac2680c 100644 --- a/server/src/at_record_server/catalog_index_postgres.gleam +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -7,7 +7,8 @@ import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog_index.{ - type Adoption, type Edit, type Store, Edit, Store, + type Adoption, type Edit, type Store, AdoptionOps, CursorOps, Edit, EditOps, + ReleaseOps, Store, } import atproto/blob.{Blob} import gleam/dict.{type Dict} @@ -26,22 +27,28 @@ const active_status_filter = "status not in ('sold', 'dropped')" pub fn table_store(conn: pog.Connection) -> Result(Store, String) { use _ <- result.try(migrate(conn)) - Ok( - Store( + Ok(Store( + releases: ReleaseOps( 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) }, ), - ) + adoptions: AdoptionOps( + upsert: fn(adoption) { upsert_adoption(conn, adoption) }, + delete: fn(entry_uri) { delete_adoption(conn, entry_uri) }, + count: fn(release_uri) { adoption_count(conn, release_uri) }, + counts: fn() { adoption_counts(conn) }, + ), + edits: EditOps( + upsert: fn(edit) { upsert_edit(conn, edit) }, + delete: fn(edit_uri) { delete_edit(conn, edit_uri) }, + for_subject: fn(subject_uri) { edits_for(conn, subject_uri) }, + ), + cursor: CursorOps( + save: fn(time_us) { save_cursor(conn, time_us) }, + load: fn() { load_cursor(conn) }, + ), + )) } fn migrate(conn: pog.Connection) -> Result(Nil, String) { diff --git a/server/src/at_record_server/handlers/oauth.gleam b/server/src/at_record_server/handlers/oauth.gleam index 649dd93..2f19e32 100644 --- a/server/src/at_record_server/handlers/oauth.gleam +++ b/server/src/at_record_server/handlers/oauth.gleam @@ -195,7 +195,7 @@ fn backfill_user_catalog(ctx: Context, session: sessions.OauthSession) -> Nil { pds: session.pds, ) let rows = browse.fetch_user_releases(ctx.atproto.client, user) - rows |> list.each(ctx.catalog_index.upsert) + rows |> list.each(ctx.catalog_index.releases.upsert) wisp.log_info( "catalog_index: backfilled " <> int.to_string(list.length(rows)) diff --git a/server/src/at_record_server/jetstream_consumer.gleam b/server/src/at_record_server/jetstream_consumer.gleam index dac2922..b1912f1 100644 --- a/server/src/at_record_server/jetstream_consumer.gleam +++ b/server/src/at_record_server/jetstream_consumer.gleam @@ -108,7 +108,7 @@ pub fn start( } fn supervise(host: String, deps: Deps, backoff_ms: Int) -> Nil { - let cursor = deps.index.load_cursor() + let cursor = deps.index.cursor.load() case connect(host, deps, cursor) { Ok(started) -> { wisp.log_info("jetstream_consumer: connected to " <> host) @@ -237,7 +237,7 @@ fn at_uri(frame: Frame) -> String { fn route_release(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { - state.deps.index.delete(at_uri(frame)) + state.deps.index.releases.delete(at_uri(frame)) state } _ -> upsert_release(state, frame) @@ -251,7 +251,7 @@ fn upsert_release(state: State, frame: Frame) -> State { Ok(value) -> { let #(publisher, next_state) = resolve_publisher(state, frame.did) let #(handle, pds) = publisher - next_state.deps.index.upsert(to_browse_row( + next_state.deps.index.releases.upsert(to_browse_row( frame, cid, value, @@ -343,7 +343,7 @@ fn fetch_identity(deps: Deps, did: String) -> Result(#(String, String), Nil) { fn route_shelf_entry(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { - state.deps.index.delete_adoption(at_uri(frame)) + state.deps.index.adoptions.delete(at_uri(frame)) state } _ -> upsert_adoption(state, frame) @@ -357,7 +357,7 @@ fn upsert_adoption(state: State, frame: Frame) -> State { Ok(entry) -> case entry.release { Some(ref) -> { - state.deps.index.upsert_adoption(catalog_index.Adoption( + state.deps.index.adoptions.upsert(catalog_index.Adoption( entry_uri: at_uri(frame), did: frame.did, release_uri: ref.uri, @@ -388,7 +388,7 @@ fn upsert_adoption(state: State, frame: Frame) -> State { fn route_edit(state: State, frame: Frame) -> State { case frame.operation { "delete" -> { - state.deps.index.delete_edit(at_uri(frame)) + state.deps.index.edits.delete(at_uri(frame)) state } _ -> upsert_edit(state, frame) @@ -402,7 +402,7 @@ fn upsert_edit(state: State, frame: Frame) -> State { Ok(edit) -> case edit.subject { Some(ref) -> { - state.deps.index.upsert_edit(catalog_index.Edit( + state.deps.index.edits.upsert(catalog_index.Edit( edit_uri: at_uri(frame), did: frame.did, subject_uri: ref.uri, @@ -430,7 +430,7 @@ fn upsert_edit(state: State, frame: Frame) -> State { fn maybe_persist_cursor(state: State, time_us: Int) -> State { case state.pending + 1 >= persist_every { True -> { - state.deps.index.save_cursor(time_us) + state.deps.index.cursor.save(time_us) State(..state, pending: 0) } False -> State(..state, pending: state.pending + 1) diff --git a/server/src/at_record_server/wiring.gleam b/server/src/at_record_server/wiring.gleam index 3701d7d..6111489 100644 --- a/server/src/at_record_server/wiring.gleam +++ b/server/src/at_record_server/wiring.gleam @@ -53,8 +53,8 @@ pub fn context( /// network round trip. fn variant_source(catalog_index: catalog_index.Store) -> catalog_source.Source { catalog_source.Source( - releases: catalog_index.list, - adoption_count: catalog_index.adoption_count, + releases: catalog_index.releases.list, + adoption_count: catalog_index.adoptions.count, ) } diff --git a/server/test/catalog_index_postgres_test.gleam b/server/test/catalog_index_postgres_test.gleam index a1752b5..9f2167e 100644 --- a/server/test/catalog_index_postgres_test.gleam +++ b/server/test/catalog_index_postgres_test.gleam @@ -48,15 +48,15 @@ fn row(uri: String) -> BrowseRow { 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 }) + store.releases.upsert(row(uri)) + let assert Ok(found) = store.releases.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 }) + store.releases.delete(uri) + let assert Error(Nil) = store.releases.list() |> find(fn(r) { r.uri == uri }) Nil } @@ -64,45 +64,45 @@ 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( + store.adoptions.upsert(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 + assert store.adoptions.count(release) == 1 + assert dict.get(store.adoptions.counts(), release) == Ok(1) + store.adoptions.delete(entry_uri) + assert store.adoptions.count(release) == 0 let edit_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.edit/pg-x1" - store.upsert_edit(Edit( + store.edits.upsert(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) + let assert [found] = store.edits.for_subject(release) assert found.entity == "release" - store.delete_edit(edit_uri) - assert store.edits_for(release) == [] + store.edits.delete(edit_uri) + assert store.edits.for_subject(release) == [] } pub fn gone_edges_are_excluded_but_reacquiring_counts_again_test() { use store <- with_store() let release = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.release/pg3" let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg-e2" - store.upsert_adoption(Adoption( + store.adoptions.upsert(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 - store.upsert_adoption(Adoption( + assert store.adoptions.count(release) == 1 + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:pgtest", release_uri: release, @@ -110,23 +110,23 @@ pub fn gone_edges_are_excluded_but_reacquiring_counts_again_test() { created_at: "2024-02-01T00:00:00Z", )) // Gone, but still stored: the count excludes it without a delete. - assert store.adoption_count(release) == 0 - assert dict.get(store.adoption_counts(), release) == Error(Nil) - store.upsert_adoption(Adoption( + assert store.adoptions.count(release) == 0 + assert dict.get(store.adoptions.counts(), release) == Error(Nil) + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:pgtest", release_uri: release, status: "acquired", created_at: "2024-03-01T00:00:00Z", )) - assert store.adoption_count(release) == 1 - store.delete_adoption(entry_uri) + assert store.adoptions.count(release) == 1 + store.adoptions.delete(entry_uri) } 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) + store.cursor.save(1_700_000_000_123_456) + assert store.cursor.load() == Some(1_700_000_000_123_456) } fn find(list: List(a), pred: fn(a) -> Bool) -> Result(a, Nil) { diff --git a/server/test/catalog_index_test.gleam b/server/test/catalog_index_test.gleam index 705e4cf..6f15a09 100644 --- a/server/test/catalog_index_test.gleam +++ b/server/test/catalog_index_test.gleam @@ -28,70 +28,74 @@ fn row(uri: String) -> BrowseRow { pub fn upsert_and_list_test() { let assert Ok(store) = catalog_index.start() - store.upsert(row("at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1")) - store.upsert(row("at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2")) - assert list.length(store.list()) == 2 + store.releases.upsert(row( + "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + )) + store.releases.upsert(row( + "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2", + )) + assert list.length(store.releases.list()) == 2 } pub fn delete_removes_row_test() { let assert Ok(store) = catalog_index.start() let uri = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - store.upsert(row(uri)) - store.delete(uri) - assert store.list() == [] + store.releases.upsert(row(uri)) + store.releases.delete(uri) + assert store.releases.list() == [] } pub fn upsert_overwrites_by_uri_test() { let assert Ok(store) = catalog_index.start() let uri = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - store.upsert(row(uri)) - store.upsert(BrowseRow(..row(uri), title: "Retitled")) - let assert [only] = store.list() + store.releases.upsert(row(uri)) + store.releases.upsert(BrowseRow(..row(uri), title: "Retitled")) + let assert [only] = store.releases.list() assert only.title == "Retitled" } pub fn adoption_upsert_delete_and_count_test() { let assert Ok(store) = catalog_index.start() let release = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri: "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1", did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", )) - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri: "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e2", did: "did:plc:b", release_uri: release, status: "wanted", created_at: "2024-01-02T00:00:00Z", )) - assert store.adoption_count(release) == 2 - assert store.adoption_counts() == dict.from_list([#(release, 2)]) - store.delete_adoption("at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1") - assert store.adoption_count(release) == 1 + assert store.adoptions.count(release) == 2 + assert store.adoptions.counts() == dict.from_list([#(release, 2)]) + store.adoptions.delete("at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1") + assert store.adoptions.count(release) == 1 } pub fn adoption_upsert_by_entry_uri_overwrites_test() { let assert Ok(store) = catalog_index.start() let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" let release = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", )) - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:a", release_uri: release, status: "wanted", created_at: "2024-02-01T00:00:00Z", )) - assert store.adoption_count(release) == 1 + assert store.adoptions.count(release) == 1 } pub fn gone_edges_do_not_count_but_stay_stored_test() { @@ -100,21 +104,21 @@ pub fn gone_edges_do_not_count_but_stay_stored_test() { let sold_entry = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" let dropped_entry = "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e2" let owned_entry = "at://did:plc:c/dev.mokkenstorm.crate.shelf.entry/e3" - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri: sold_entry, did: "did:plc:a", release_uri: release, status: "sold", created_at: "2024-01-01T00:00:00Z", )) - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri: dropped_entry, did: "did:plc:b", release_uri: release, status: "dropped", created_at: "2024-01-01T00:00:00Z", )) - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri: owned_entry, did: "did:plc:c", release_uri: release, @@ -124,52 +128,52 @@ pub fn gone_edges_do_not_count_but_stay_stored_test() { // Only the owned edge counts; the gone edges stay stored (no delete // required to make them stop counting), so `adoption_counts` never even // mentions a release whose only edges are all gone. - assert store.adoption_count(release) == 1 - assert store.adoption_counts() == dict.from_list([#(release, 1)]) + assert store.adoptions.count(release) == 1 + assert store.adoptions.counts() == dict.from_list([#(release, 1)]) } pub fn gone_then_reacquired_entry_counts_again_test() { let assert Ok(store) = catalog_index.start() let entry_uri = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" let release = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - store.upsert_adoption(Adoption( + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", )) - assert store.adoption_count(release) == 1 - store.upsert_adoption(Adoption( + assert store.adoptions.count(release) == 1 + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:a", release_uri: release, status: "sold", created_at: "2024-02-01T00:00:00Z", )) - assert store.adoption_count(release) == 0 - store.upsert_adoption(Adoption( + assert store.adoptions.count(release) == 0 + store.adoptions.upsert(Adoption( entry_uri:, did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-03-01T00:00:00Z", )) - assert store.adoption_count(release) == 1 + assert store.adoptions.count(release) == 1 } pub fn edit_upsert_delete_and_edits_for_test() { let assert Ok(store) = catalog_index.start() let subject = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" let edit_uri = "at://did:plc:a/dev.mokkenstorm.crate.catalog.edit/x1" - store.upsert_edit(Edit( + store.edits.upsert(Edit( edit_uri:, did: "did:plc:a", subject_uri: subject, entity: "release", created_at: "2024-01-01T00:00:00Z", )) - assert store.edits_for(subject) + assert store.edits.for_subject(subject) == [ Edit( edit_uri:, @@ -179,15 +183,15 @@ pub fn edit_upsert_delete_and_edits_for_test() { created_at: "2024-01-01T00:00:00Z", ), ] - store.delete_edit(edit_uri) - assert store.edits_for(subject) == [] + store.edits.delete(edit_uri) + assert store.edits.for_subject(subject) == [] } pub fn cursor_save_and_load_test() { let assert Ok(store) = catalog_index.start() - assert store.load_cursor() == None - store.save_cursor(1_700_000_000_000_000) - assert store.load_cursor() == Some(1_700_000_000_000_000) - store.save_cursor(1_700_000_001_000_000) - assert store.load_cursor() == Some(1_700_000_001_000_000) + assert store.cursor.load() == None + store.cursor.save(1_700_000_000_000_000) + assert store.cursor.load() == Some(1_700_000_000_000_000) + store.cursor.save(1_700_000_001_000_000) + assert store.cursor.load() == Some(1_700_000_001_000_000) } diff --git a/server/test/jetstream_consumer_test.gleam b/server/test/jetstream_consumer_test.gleam index bc9f286..e72de51 100644 --- a/server/test/jetstream_consumer_test.gleam +++ b/server/test/jetstream_consumer_test.gleam @@ -198,7 +198,7 @@ pub fn release_commit_upserts_into_the_index_test() { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) let _ = jetstream_consumer.route(state, frame) - let assert [row] = index.list() + let assert [row] = index.releases.list() assert row.uri == release_uri assert row.title == "Test Album" assert row.publisher_handle == "pub.test" @@ -225,7 +225,7 @@ pub fn release_delete_removes_from_the_index_test() { jetstream_consumer.frame_decoder(), ) let _ = jetstream_consumer.route(state, delete) - assert index.list() == [] + assert index.releases.list() == [] } pub fn shelf_entry_with_release_ref_creates_an_adoption_edge_test() { @@ -240,7 +240,7 @@ pub fn shelf_entry_with_release_ref_creates_an_adoption_edge_test() { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) let _ = jetstream_consumer.route(state, frame) - assert index.adoption_count(release_uri) == 1 + assert index.adoptions.count(release_uri) == 1 } pub fn shelf_entry_without_release_ref_is_ignored_test() { @@ -255,7 +255,7 @@ pub fn shelf_entry_without_release_ref_is_ignored_test() { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) let _ = jetstream_consumer.route(state, frame) - assert index.adoption_counts() == dict.new() + assert index.adoptions.counts() == dict.new() } pub fn shelf_entry_delete_removes_the_adoption_edge_test() { @@ -278,7 +278,7 @@ pub fn shelf_entry_delete_removes_the_adoption_edge_test() { jetstream_consumer.frame_decoder(), ) let _ = jetstream_consumer.route(state, delete) - assert index.adoption_count(release_uri) == 0 + assert index.adoptions.count(release_uri) == 0 } pub fn catalog_edit_indexes_subject_and_entity_test() { @@ -293,7 +293,7 @@ pub fn catalog_edit_indexes_subject_and_entity_test() { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) let _ = jetstream_consumer.route(state, frame) - let assert [edit] = index.edits_for(release_uri) + let assert [edit] = index.edits.for_subject(release_uri) assert edit.did == "did:plc:a" assert edit.entity == "release" } @@ -318,7 +318,7 @@ pub fn catalog_edit_delete_removes_the_edit_test() { jetstream_consumer.frame_decoder(), ) let _ = jetstream_consumer.route(state, delete) - assert index.edits_for(release_uri) == [] + assert index.edits.for_subject(release_uri) == [] } pub fn unknown_did_resolves_via_resolve_mini_doc_test() { @@ -333,7 +333,7 @@ pub fn unknown_did_resolves_via_resolve_mini_doc_test() { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) let _ = jetstream_consumer.route(state, frame) - let assert [row] = index.list() + let assert [row] = index.releases.list() assert row.publisher_handle == "stranger.test" assert row.publisher_pds == "https://stranger-pds.test" } diff --git a/server/test/oauth_test.gleam b/server/test/oauth_test.gleam index 2673457..0c43dc0 100644 --- a/server/test/oauth_test.gleam +++ b/server/test/oauth_test.gleam @@ -408,7 +408,7 @@ pub fn callback_success_backfills_new_users_releases_into_catalog_index_test() { let resp = oauth_handler.callback(req, ctx) assert resp.status == 303 - let assert [row] = index.list() + let assert [row] = index.releases.list() assert row.title == "Loveless" assert row.publisher_did == "did:plc:abc" } diff --git a/server/test/support.gleam b/server/test/support.gleam index 30c0ed1..40e357f 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -45,18 +45,23 @@ fn known_users_of(users: List(known_users.KnownUser)) -> known_users.Store { fn empty_catalog_index() -> catalog_index.Store { catalog_index.Store( - upsert: fn(_) { Nil }, - delete: fn(_) { Nil }, - list: fn() { [] }, - upsert_adoption: fn(_) { Nil }, - delete_adoption: fn(_) { Nil }, - upsert_edit: fn(_) { Nil }, - delete_edit: fn(_) { Nil }, - adoption_count: fn(_) { 0 }, - adoption_counts: fn() { dict.new() }, - edits_for: fn(_) { [] }, - save_cursor: fn(_) { Nil }, - load_cursor: fn() { None }, + releases: catalog_index.ReleaseOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + list: fn() { [] }, + ), + adoptions: catalog_index.AdoptionOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + count: fn(_) { 0 }, + counts: fn() { dict.new() }, + ), + edits: catalog_index.EditOps( + upsert: fn(_) { Nil }, + delete: fn(_) { Nil }, + for_subject: fn(_) { [] }, + ), + cursor: catalog_index.CursorOps(save: fn(_) { Nil }, load: fn() { None }), ) } -- 2.51.2