diff --git a/server/test/catalog_index_test.gleam b/server/test/catalog_index_test.gleam index 6f15a09..662f2a9 100644 --- a/server/test/catalog_index_test.gleam +++ b/server/test/catalog_index_test.gleam @@ -1,5 +1,5 @@ import at_record_server/catalog/row.{type BrowseRow, BrowseRow} -import at_record_server/catalog_index.{Adoption, Edit} +import at_record_server/catalog_index.{type Adoption, Adoption, Edit} import gleam/dict import gleam/list import gleam/option.{None, Some} @@ -26,6 +26,17 @@ fn row(uri: String) -> BrowseRow { ) } +/// `did` is never asserted on in this file (adoptions are keyed and counted +/// by `entry_uri`/`release_uri`/`status`), so it's a fixed placeholder here. +fn adoption( + entry_uri entry_uri: String, + release_uri release_uri: String, + status status: String, + created_at created_at: String, +) -> Adoption { + Adoption(entry_uri:, did: "did:plc:a", release_uri:, status:, created_at:) +} + pub fn upsert_and_list_test() { let assert Ok(store) = catalog_index.start() store.releases.upsert(row( @@ -57,23 +68,22 @@ pub fn upsert_overwrites_by_uri_test() { 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.adoptions.upsert(Adoption( - entry_uri: "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1", - did: "did:plc:a", + let entry_a = "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1" + store.adoptions.upsert(adoption( + entry_uri: entry_a, release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", )) - store.adoptions.upsert(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.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") + store.adoptions.delete(entry_a) assert store.adoptions.count(release) == 1 } @@ -81,16 +91,14 @@ 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.adoptions.upsert(Adoption( + store.adoptions.upsert(adoption( entry_uri:, - did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", )) - store.adoptions.upsert(Adoption( + store.adoptions.upsert(adoption( entry_uri:, - did: "did:plc:a", release_uri: release, status: "wanted", created_at: "2024-02-01T00:00:00Z", @@ -101,26 +109,20 @@ pub fn adoption_upsert_by_entry_uri_overwrites_test() { pub fn gone_edges_do_not_count_but_stay_stored_test() { let assert Ok(store) = catalog_index.start() let release = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" - 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.adoptions.upsert(Adoption( - entry_uri: sold_entry, - did: "did:plc:a", + store.adoptions.upsert(adoption( + entry_uri: "at://did:plc:a/dev.mokkenstorm.crate.shelf.entry/e1", release_uri: release, status: "sold", created_at: "2024-01-01T00:00:00Z", )) - store.adoptions.upsert(Adoption( - entry_uri: dropped_entry, - did: "did:plc:b", + store.adoptions.upsert(adoption( + entry_uri: "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e2", release_uri: release, status: "dropped", created_at: "2024-01-01T00:00:00Z", )) - store.adoptions.upsert(Adoption( - entry_uri: owned_entry, - did: "did:plc:c", + store.adoptions.upsert(adoption( + entry_uri: "at://did:plc:c/dev.mokkenstorm.crate.shelf.entry/e3", release_uri: release, status: "acquired", created_at: "2024-01-01T00:00:00Z", @@ -136,25 +138,22 @@ 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.adoptions.upsert(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.adoptions.count(release) == 1 - store.adoptions.upsert(Adoption( + store.adoptions.upsert(adoption( entry_uri:, - did: "did:plc:a", release_uri: release, status: "sold", created_at: "2024-02-01T00:00:00Z", )) assert store.adoptions.count(release) == 0 - store.adoptions.upsert(Adoption( + store.adoptions.upsert(adoption( entry_uri:, - did: "did:plc:a", release_uri: release, status: "acquired", created_at: "2024-03-01T00:00:00Z", @@ -166,23 +165,16 @@ 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.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(subject) - == [ - Edit( - edit_uri:, - did: "did:plc:a", - subject_uri: subject, - entity: "release", - created_at: "2024-01-01T00:00:00Z", - ), - ] + let edit = + Edit( + edit_uri:, + did: "did:plc:a", + subject_uri: subject, + entity: "release", + created_at: "2024-01-01T00:00:00Z", + ) + store.edits.upsert(edit) + assert store.edits.for_subject(subject) == [edit] store.edits.delete(edit_uri) assert store.edits.for_subject(subject) == [] } diff --git a/server/test/catalog_stats_test.gleam b/server/test/catalog_stats_test.gleam index f498f00..bbfb6b6 100644 --- a/server/test/catalog_stats_test.gleam +++ b/server/test/catalog_stats_test.gleam @@ -1,11 +1,12 @@ -import at_record_server/catalog/source.{type VariantRow, VariantRow} +import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog/stats import atproto/blob.{Blob} import gleam/dict +import gleam/list import gleam/option.{None, Some} -fn bare_row(uri: String) -> VariantRow { - VariantRow( +fn bare_row(uri: String) -> BrowseRow { + BrowseRow( uri:, cid: "bafyrow", title: "Title", @@ -26,43 +27,47 @@ fn bare_row(uri: String) -> VariantRow { ) } -pub fn completeness_is_zero_when_every_optional_field_is_absent_test() { - assert stats.completeness(bare_row("at://a/1")) == 0.0 -} - -pub fn completeness_is_one_when_every_optional_field_is_present_test() { - let full = - VariantRow( - ..bare_row("at://a/1"), - artist_display: Some("Slint"), - released: Some("1991"), - country: Some("US"), - cover: Some(Blob(cid: "bafycover", mime_type: "image/jpeg", size: 100)), - thumb_url: Some("https://example.test/thumb.jpg"), - discogs_id: Some("42"), - ) - assert stats.completeness(full) == 1.0 -} - -pub fn completeness_is_a_fraction_for_partial_rows_test() { - let half = - VariantRow( - ..bare_row("at://a/1"), - artist_display: Some("Slint"), - country: Some("US"), - released: Some("1991"), - ) - assert stats.completeness(half) == 3.0 /. 6.0 -} - -pub fn completeness_ignores_chain_edges_test() { - let with_chain = - VariantRow( - ..bare_row("at://a/1"), - supersedes: Some("at://old/1"), - based_on: Some("at://origin/1"), - ) - assert stats.completeness(with_chain) == 0.0 +pub fn completeness_reflects_how_many_optional_fields_are_populated_test() { + let cases = [ + #("every optional field absent", bare_row("at://a/1"), 0.0), + #( + "every optional field present", + BrowseRow( + ..bare_row("at://a/1"), + artist_display: Some("Slint"), + released: Some("1991"), + country: Some("US"), + cover: Some(Blob(cid: "bafycover", mime_type: "image/jpeg", size: 100)), + thumb_url: Some("https://example.test/thumb.jpg"), + discogs_id: Some("42"), + ), + 1.0, + ), + #( + "half the optional fields present", + BrowseRow( + ..bare_row("at://a/1"), + artist_display: Some("Slint"), + country: Some("US"), + released: Some("1991"), + ), + 3.0 /. 6.0, + ), + #( + "chain edges (supersedes/based_on) don't count toward completeness", + BrowseRow( + ..bare_row("at://a/1"), + supersedes: Some("at://old/1"), + based_on: Some("at://origin/1"), + ), + 0.0, + ), + ] + cases + |> list.each(fn(c) { + let #(_label, row, expected) = c + assert stats.completeness(row) == expected + }) } pub fn compute_reads_adoption_off_the_injected_counter_test() { diff --git a/server/test/catalog_strategy_test.gleam b/server/test/catalog_strategy_test.gleam index 5e1de4d..2c4170e 100644 --- a/server/test/catalog_strategy_test.gleam +++ b/server/test/catalog_strategy_test.gleam @@ -1,12 +1,12 @@ -import at_record_server/catalog/source.{type VariantRow, VariantRow} +import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog/stats.{VariantStats} import at_record_server/catalog/strategy import gleam/dict import gleam/list import gleam/option.{None, Some} -fn variant_row(uri: String, created_at: String) -> VariantRow { - VariantRow( +fn variant_row(uri: String, created_at: String) -> BrowseRow { + BrowseRow( uri:, cid: "bafyrow", title: "Title", @@ -28,8 +28,8 @@ fn variant_row(uri: String, created_at: String) -> VariantRow { } fn stats_of( - entries: List(#(VariantRow, Int, String, Float)), -) -> #(List(VariantRow), dict.Dict(String, stats.VariantStats)) { + entries: List(#(BrowseRow, Int, String, Float)), +) -> #(List(BrowseRow), dict.Dict(String, stats.VariantStats)) { let rows = entries |> list.map(fn(e) { e.0 }) let table = entries @@ -57,52 +57,66 @@ pub fn default_strategy_is_most_adopted_test() { assert strategy.default_strategy() == strategy.MostAdopted } -pub fn resolve_most_adopted_picks_the_higher_adoption_count_test() { - let a = variant_row("at://a/1", "2020-01-01") - let b = variant_row("at://b/1", "2020-01-01") - let #(rows, table) = - stats_of([#(a, 3, "2020-01-01", 0.5), #(b, 9, "2020-01-01", 0.5)]) - assert strategy.resolve(rows, table, strategy.MostAdopted) == Some(b) +/// One case per strategy: `b`'s stats beat `a`'s on exactly the metric that +/// strategy names first, so the winner proves `resolve` reads the right +/// field. +pub fn resolve_picks_the_winner_by_the_strategys_primary_metric_test() { + let cases = [ + #(strategy.MostAdopted, #(3, "2020-01-01", 0.5), #(9, "2020-01-01", 0.5)), + #(strategy.MostRecent, #(0, "2020-01-01", 0.5), #(0, "2024-01-01", 0.5)), + #(strategy.MostComplete, #(0, "2020-01-01", 0.2), #(0, "2020-01-01", 0.9)), + ] + cases + |> list.each(fn(c) { + let #(chosen, a_stats, b_stats) = c + let #(a_adoption, a_created_at, a_completeness) = a_stats + let #(b_adoption, b_created_at, b_completeness) = b_stats + let a = variant_row("at://a/1", a_created_at) + let b = variant_row("at://b/1", b_created_at) + let #(rows, table) = + stats_of([ + #(a, a_adoption, a_created_at, a_completeness), + #(b, b_adoption, b_created_at, b_completeness), + ]) + assert strategy.resolve(rows, table, chosen) == Some(b) + }) } -pub fn resolve_most_recent_picks_the_newer_created_at_test() { - let a = variant_row("at://a/1", "2020-01-01") - let b = variant_row("at://b/1", "2024-01-01") - let #(rows, table) = - stats_of([#(a, 0, "2020-01-01", 0.5), #(b, 0, "2024-01-01", 0.5)]) - assert strategy.resolve(rows, table, strategy.MostRecent) == Some(b) -} - -pub fn resolve_most_complete_picks_the_higher_completeness_test() { - let a = variant_row("at://a/1", "2020-01-01") - let b = variant_row("at://b/1", "2020-01-01") - let #(rows, table) = - stats_of([#(a, 0, "2020-01-01", 0.2), #(b, 0, "2020-01-01", 0.9)]) - assert strategy.resolve(rows, table, strategy.MostComplete) == Some(b) -} - -pub fn resolve_breaks_a_full_tie_on_created_at_then_uri_test() { - let a = variant_row("at://z/1", "2020-01-01") - let b = variant_row("at://a/1", "2020-01-01") - let #(rows, table) = - stats_of([#(a, 5, "2020-01-01", 0.5), #(b, 5, "2020-01-01", 0.5)]) - // Equal adoption -> falls through to created_at (also equal) -> uri asc: - // "at://a/1" sorts before "at://z/1". - assert strategy.resolve(rows, table, strategy.MostAdopted) == Some(b) -} - -pub fn resolve_breaks_a_metric_tie_on_created_at_before_uri_test() { - let older_but_alphabetically_later = variant_row("at://z/1", "2019-01-01") - let newer_but_alphabetically_earlier = variant_row("at://a/1", "2021-01-01") - let #(rows, table) = - stats_of([ - #(older_but_alphabetically_later, 5, "2019-01-01", 0.5), - #(newer_but_alphabetically_earlier, 5, "2021-01-01", 0.5), - ]) - // Equal adoption -> tie-break is created_at ascending, which picks the - // 2019 row even though "at://z/1" sorts after "at://a/1". - assert strategy.resolve(rows, table, strategy.MostAdopted) - == Some(older_but_alphabetically_later) +/// The tie-break cascade: primary metric, then created_at ascending, then +/// uri ascending. +pub fn resolve_breaks_ties_on_created_at_then_uri_test() { + let cases = [ + // Full tie (adoption and created_at both equal) -> falls through to uri + // ascending: "at://a/1" beats "at://z/1". + #( + #("at://z/1", 5, "2020-01-01"), + #("at://a/1", 5, "2020-01-01"), + "at://a/1", + ), + // Adoption tie only -> decided by created_at ascending, even though that + // picks the uri that sorts later. + #( + #("at://z/1", 5, "2019-01-01"), + #("at://a/1", 5, "2021-01-01"), + "at://z/1", + ), + ] + cases + |> list.each(fn(c) { + let #(a_spec, b_spec, expected_uri) = c + let #(a_uri, a_adoption, a_created_at) = a_spec + let #(b_uri, b_adoption, b_created_at) = b_spec + let a = variant_row(a_uri, a_created_at) + let b = variant_row(b_uri, b_created_at) + let #(rows, table) = + stats_of([ + #(a, a_adoption, a_created_at, 0.5), + #(b, b_adoption, b_created_at, 0.5), + ]) + let assert Some(winner) = + strategy.resolve(rows, table, strategy.MostAdopted) + assert winner.uri == expected_uri + }) } pub fn resolve_of_an_empty_set_is_none_test() { diff --git a/server/test/jetstream_consumer_test.gleam b/server/test/jetstream_consumer_test.gleam index e72de51..bb78671 100644 --- a/server/test/jetstream_consumer_test.gleam +++ b/server/test/jetstream_consumer_test.gleam @@ -13,72 +13,90 @@ import gleam/erlang/process import gleam/http/response import gleam/int import gleam/json +import gleam/list import gleam/option.{None, Some} import support const release_uri = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1" -fn release_frame(did: String, time_us: Int, operation: String) -> String { +/// One raw Jetstream frame, `commit.record` included unless `operation` is +/// `"delete"` (matches the real firehose: a delete commit carries no cid or +/// record). +fn frame_json( + did did: String, + time_us time_us: Int, + operation operation: String, + collection collection: String, + rkey rkey: String, + cid cid: String, + record record: String, +) -> String { "{\"did\":\"" <> did <> "\",\"time_us\":" <> int.to_string(time_us) <> ",\"kind\":\"commit\",\"commit\":{\"rev\":\"1\",\"operation\":\"" <> operation - <> "\",\"collection\":\"dev.mokkenstorm.crate.catalog.release\",\"rkey\":\"r1\"" + <> "\",\"collection\":\"" + <> collection + <> "\",\"rkey\":\"" + <> rkey + <> "\"" <> case operation { "delete" -> "" - _ -> - ",\"cid\":\"bafyrelcid\",\"record\":{\"$type\":\"dev.mokkenstorm.crate.catalog.release\",\"createdAt\":\"2024-01-01T00:00:00Z\",\"title\":\"Test Album\",\"artistDisplay\":\"Test Artist\",\"genres\":[\"Electronic\"],\"styles\":[\"Techno\"]}" + _ -> ",\"cid\":\"" <> cid <> "\",\"record\":" <> record } <> "}}" } +fn release_frame(did: String, time_us: Int, operation: String) -> String { + frame_json( + did:, + time_us:, + operation:, + collection: "dev.mokkenstorm.crate.catalog.release", + rkey: "r1", + cid: "bafyrelcid", + record: "{\"$type\":\"dev.mokkenstorm.crate.catalog.release\",\"createdAt\":\"2024-01-01T00:00:00Z\",\"title\":\"Test Album\",\"artistDisplay\":\"Test Artist\",\"genres\":[\"Electronic\"],\"styles\":[\"Techno\"]}", + ) +} + fn shelf_entry_frame( did: String, time_us: Int, operation: String, with_release: Bool, ) -> String { - "{\"did\":\"" - <> did - <> "\",\"time_us\":" - <> int.to_string(time_us) - <> ",\"kind\":\"commit\",\"commit\":{\"rev\":\"1\",\"operation\":\"" - <> operation - <> "\",\"collection\":\"dev.mokkenstorm.crate.shelf.entry\",\"rkey\":\"e1\"" - <> case operation { - "delete" -> "" - _ -> - ",\"cid\":\"bafyentrycid\",\"record\":{\"$type\":\"dev.mokkenstorm.crate.shelf.entry\",\"action\":\"acquired\",\"createdAt\":\"2024-01-02T00:00:00Z\"" - <> case with_release { - True -> - ",\"release\":{\"cid\":\"bafyrelcid\",\"uri\":\"" - <> release_uri - <> "\"}" - False -> "" - } - <> "}" + let release_field = case with_release { + True -> + ",\"release\":{\"cid\":\"bafyrelcid\",\"uri\":\"" <> release_uri <> "\"}" + False -> "" } - <> "}}" + frame_json( + did:, + time_us:, + operation:, + collection: "dev.mokkenstorm.crate.shelf.entry", + rkey: "e1", + cid: "bafyentrycid", + record: "{\"$type\":\"dev.mokkenstorm.crate.shelf.entry\",\"action\":\"acquired\",\"createdAt\":\"2024-01-02T00:00:00Z\"" + <> release_field + <> "}", + ) } fn catalog_edit_frame(did: String, time_us: Int, operation: String) -> String { - "{\"did\":\"" - <> did - <> "\",\"time_us\":" - <> int.to_string(time_us) - <> ",\"kind\":\"commit\",\"commit\":{\"rev\":\"1\",\"operation\":\"" - <> operation - <> "\",\"collection\":\"dev.mokkenstorm.crate.catalog.edit\",\"rkey\":\"x1\"" - <> case operation { - "delete" -> "" - _ -> - ",\"cid\":\"bafyeditcid\",\"record\":{\"$type\":\"dev.mokkenstorm.crate.catalog.edit\",\"entity\":\"release\",\"op\":\"propose\",\"createdAt\":\"2024-01-03T00:00:00Z\",\"subject\":{\"cid\":\"bafyrelcid\",\"uri\":\"" + frame_json( + did:, + time_us:, + operation:, + 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\":\"" <> release_uri - <> "\"}}" - } - <> "}}" + <> "\"}}", + ) } fn known_user_store() -> known_users.Store { @@ -141,36 +159,24 @@ fn deps( Deps(known_users: known, index:, client:, resolver: "https://resolver.test") } -pub fn release_frame_decodes_all_top_level_fields_test() { - let text = release_frame("did:plc:pub", 1_700_000_000_000_000, "create") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) - assert frame.did == "did:plc:pub" - assert frame.time_us == 1_700_000_000_000_000 - assert frame.operation == "create" - assert frame.collection == "dev.mokkenstorm.crate.catalog.release" - assert frame.rkey == "r1" - assert frame.cid == Some("bafyrelcid") - assert frame.record != None +fn state_over( + index: catalog_index.Store, + client: xrpc.Client, +) -> jetstream_consumer.State { + jetstream_consumer.initial_state(deps(empty_known_users(), index, client)) } -pub fn shelf_entry_frame_decodes_test() { - let text = - shelf_entry_frame("did:plc:a", 1_700_000_000_000_001, "create", True) +fn decode_frame(text: String) -> jetstream_consumer.Frame { let assert Ok(Some(frame)) = json.parse(text, jetstream_consumer.frame_decoder()) - assert frame.collection == "dev.mokkenstorm.crate.shelf.entry" - assert frame.rkey == "e1" - assert frame.record != None + frame } -pub fn catalog_edit_frame_decodes_test() { - let text = catalog_edit_frame("did:plc:a", 1_700_000_000_000_002, "create") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) - assert frame.collection == "dev.mokkenstorm.crate.catalog.edit" - assert frame.rkey == "x1" - assert frame.record != None +pub fn frame_decoder_reads_did_and_time_us_test() { + let frame = + decode_frame(release_frame("did:plc:pub", 1_700_000_000_000_000, "create")) + assert frame.did == "did:plc:pub" + assert frame.time_us == 1_700_000_000_000_000 } pub fn non_commit_kind_decodes_to_none_test() { @@ -178,12 +184,47 @@ pub fn non_commit_kind_decodes_to_none_test() { assert json.parse(text, jetstream_consumer.frame_decoder()) == Ok(None) } -pub fn delete_frame_has_no_cid_or_record_test() { - let text = release_frame("did:plc:pub", 2, "delete") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) - assert frame.cid == None - assert frame.record == None +/// `commit.record`/`cid` are only present on a create/update commit, never a +/// delete, across every collection this consumer subscribes to. +pub fn frame_decoder_decodes_every_collection_and_operation_test() { + let release = "dev.mokkenstorm.crate.catalog.release" + let shelf_entry = "dev.mokkenstorm.crate.shelf.entry" + let catalog_edit = "dev.mokkenstorm.crate.catalog.edit" + let cases = [ + #(release_frame("did:plc:pub", 1, "create"), release, "r1", True), + #(release_frame("did:plc:pub", 2, "delete"), release, "r1", False), + #( + shelf_entry_frame("did:plc:a", 3, "create", True), + shelf_entry, + "e1", + True, + ), + #( + shelf_entry_frame("did:plc:a", 4, "delete", True), + shelf_entry, + "e1", + False, + ), + #(catalog_edit_frame("did:plc:a", 5, "create"), catalog_edit, "x1", True), + #(catalog_edit_frame("did:plc:a", 6, "delete"), catalog_edit, "x1", False), + ] + cases + |> list.each(fn(c) { + let #(text, collection, rkey, has_record) = c + let frame = decode_frame(text) + assert frame.collection == collection + assert frame.rkey == rkey + case has_record { + True -> { + assert frame.cid != None + assert frame.record != None + } + False -> { + assert frame.cid == None + assert frame.record == None + } + } + }) } pub fn release_commit_upserts_into_the_index_test() { @@ -194,9 +235,7 @@ pub fn release_commit_upserts_into_the_index_test() { index, support.unreachable_client(), )) - let text = release_frame("did:plc:pub", 1, "create") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) + let frame = decode_frame(release_frame("did:plc:pub", 1, "create")) let _ = jetstream_consumer.route(state, frame) let assert [row] = index.releases.list() assert row.uri == release_uri @@ -213,85 +252,55 @@ pub fn release_delete_removes_from_the_index_test() { index, support.unreachable_client(), )) - let assert Ok(Some(create)) = - json.parse( - release_frame("did:plc:pub", 1, "create"), - jetstream_consumer.frame_decoder(), + let state = + jetstream_consumer.route( + state, + decode_frame(release_frame("did:plc:pub", 1, "create")), ) - let state = jetstream_consumer.route(state, create) - let assert Ok(Some(delete)) = - json.parse( - release_frame("did:plc:pub", 2, "delete"), - jetstream_consumer.frame_decoder(), + let _ = + jetstream_consumer.route( + state, + decode_frame(release_frame("did:plc:pub", 2, "delete")), ) - let _ = jetstream_consumer.route(state, delete) assert index.releases.list() == [] } pub fn shelf_entry_with_release_ref_creates_an_adoption_edge_test() { let assert Ok(index) = catalog_index.start() - let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - support.unreachable_client(), - )) - let text = shelf_entry_frame("did:plc:a", 1, "create", True) - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) + let state = state_over(index, support.unreachable_client()) + let frame = decode_frame(shelf_entry_frame("did:plc:a", 1, "create", True)) let _ = jetstream_consumer.route(state, frame) assert index.adoptions.count(release_uri) == 1 } pub fn shelf_entry_without_release_ref_is_ignored_test() { let assert Ok(index) = catalog_index.start() - let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - support.unreachable_client(), - )) - let text = shelf_entry_frame("did:plc:a", 1, "create", False) - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) + let state = state_over(index, support.unreachable_client()) + let frame = decode_frame(shelf_entry_frame("did:plc:a", 1, "create", False)) let _ = jetstream_consumer.route(state, frame) assert index.adoptions.counts() == dict.new() } pub fn shelf_entry_delete_removes_the_adoption_edge_test() { let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - support.unreachable_client(), - )) - let assert Ok(Some(create)) = - json.parse( - shelf_entry_frame("did:plc:a", 1, "create", True), - jetstream_consumer.frame_decoder(), + jetstream_consumer.route( + state, + decode_frame(shelf_entry_frame("did:plc:a", 1, "create", True)), ) - let state = jetstream_consumer.route(state, create) - let assert Ok(Some(delete)) = - json.parse( - shelf_entry_frame("did:plc:a", 2, "delete", True), - jetstream_consumer.frame_decoder(), + let _ = + jetstream_consumer.route( + state, + decode_frame(shelf_entry_frame("did:plc:a", 2, "delete", True)), ) - let _ = jetstream_consumer.route(state, delete) assert index.adoptions.count(release_uri) == 0 } pub fn catalog_edit_indexes_subject_and_entity_test() { let assert Ok(index) = catalog_index.start() - let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - support.unreachable_client(), - )) - let text = catalog_edit_frame("did:plc:a", 1, "create") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) + let state = state_over(index, support.unreachable_client()) + let frame = decode_frame(catalog_edit_frame("did:plc:a", 1, "create")) let _ = jetstream_consumer.route(state, frame) let assert [edit] = index.edits.for_subject(release_uri) assert edit.did == "did:plc:a" @@ -300,38 +309,24 @@ pub fn catalog_edit_indexes_subject_and_entity_test() { pub fn catalog_edit_delete_removes_the_edit_test() { let assert Ok(index) = catalog_index.start() + let state = state_over(index, support.unreachable_client()) let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - support.unreachable_client(), - )) - let assert Ok(Some(create)) = - json.parse( - catalog_edit_frame("did:plc:a", 1, "create"), - jetstream_consumer.frame_decoder(), + jetstream_consumer.route( + state, + decode_frame(catalog_edit_frame("did:plc:a", 1, "create")), ) - let state = jetstream_consumer.route(state, create) - let assert Ok(Some(delete)) = - json.parse( - catalog_edit_frame("did:plc:a", 2, "delete"), - jetstream_consumer.frame_decoder(), + let _ = + jetstream_consumer.route( + state, + decode_frame(catalog_edit_frame("did:plc:a", 2, "delete")), ) - let _ = jetstream_consumer.route(state, delete) assert index.edits.for_subject(release_uri) == [] } pub fn unknown_did_resolves_via_resolve_mini_doc_test() { let assert Ok(index) = catalog_index.start() - let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - resolving_client(), - )) - let text = release_frame("did:plc:stranger", 1, "create") - let assert Ok(Some(frame)) = - json.parse(text, jetstream_consumer.frame_decoder()) + let state = state_over(index, resolving_client()) + let frame = decode_frame(release_frame("did:plc:stranger", 1, "create")) let _ = jetstream_consumer.route(state, frame) let assert [row] = index.releases.list() assert row.publisher_handle == "stranger.test" @@ -341,24 +336,17 @@ pub fn unknown_did_resolves_via_resolve_mini_doc_test() { pub fn unknown_did_resolution_is_cached_within_a_connection_test() { let assert Ok(index) = catalog_index.start() let calls = process.new_subject() + let state = state_over(index, counting_resolver_client(calls)) let state = - jetstream_consumer.initial_state(deps( - empty_known_users(), - index, - counting_resolver_client(calls), - )) - let assert Ok(Some(first)) = - json.parse( - release_frame("did:plc:stranger", 1, "create"), - jetstream_consumer.frame_decoder(), + jetstream_consumer.route( + state, + decode_frame(release_frame("did:plc:stranger", 1, "create")), ) - let state = jetstream_consumer.route(state, first) - let assert Ok(Some(second)) = - json.parse( - release_frame("did:plc:stranger", 2, "create"), - jetstream_consumer.frame_decoder(), + let _ = + jetstream_consumer.route( + state, + decode_frame(release_frame("did:plc:stranger", 2, "create")), ) - let _ = jetstream_consumer.route(state, second) // Two events from the same unknown did, but the network is only asked once // -- the second resolution comes from the in-process cache. assert count_messages(calls) == 1