diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index 100e0a1..c04328b 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -126,6 +126,9 @@ fn backfill_adoptions( release_uri: ref.uri, status: entry.action, created_at: entry.created_at, + source: catalog_index.source_label( + option.then(entry.source, fn(s) { s.origin }), + ), )) None -> Nil } diff --git a/server/src/at_record_server/catalog_index.gleam b/server/src/at_record_server/catalog_index.gleam index d128259..1aa11cd 100644 --- a/server/src/at_record_server/catalog_index.gleam +++ b/server/src/at_record_server/catalog_index.gleam @@ -11,7 +11,7 @@ import at_record_server/parallel import gleam/dict.{type Dict} import gleam/erlang/process.{type Subject} import gleam/list -import gleam/option.{type Option} +import gleam/option.{type Option, Some} import gleam/otp/actor import gleam/result @@ -26,9 +26,20 @@ pub type Adoption { release_uri: String, status: String, created_at: String, + source: Option(String), ) } +/// The `Adoption.source` label derived from a shelf.entry's provenance +/// origin: only the discogs-import case is folded to the shorter, +/// feed-policy-facing label ("discogs"); anything else passes through as-is. +pub fn source_label(origin: Option(String)) -> Option(String) { + case origin { + Some("discogs-import") -> Some("discogs") + other -> other + } +} + /// `status` carries the raw shelf.entry `action` (see `jetstream_consumer` /// and `at_record_server`'s backfill, both of which store `entry.action` /// verbatim), not the folded `crate.Status` vocabulary. These are the two @@ -69,6 +80,9 @@ pub type AdoptionOps { delete: fn(String) -> Nil, count: fn(String) -> Int, counts: fn() -> Dict(String, Int), + /// Every active (not gone) adoption; callers filter/sort/paginate + /// further, so this is deliberately unfiltered beyond that. + list: fn() -> List(Adoption), ) } @@ -112,6 +126,7 @@ type Msg { DeleteEdit(String) AdoptionCount(String, Subject(Int)) AdoptionCounts(Subject(Dict(String, Int))) + AdoptionsList(Subject(List(Adoption))) EditsFor(String, Subject(List(Edit))) SaveCursor(Int) LoadCursor(Subject(Option(Int))) @@ -169,6 +184,12 @@ pub fn start() -> Result(Store, actor.StartError) { Error(Nil) -> dict.new() } }, + list: fn() { + case parallel.try_call(subject, 1000, AdoptionsList) { + Ok(rows) -> rows + Error(Nil) -> [] + } + }, ), edits: EditOps( upsert: fn(edit) { @@ -238,6 +259,10 @@ fn handle(state: State, msg: Msg) -> actor.Next(State, Msg) { process.send(reply, all_counts(state)) actor.continue(state) } + AdoptionsList(reply) -> { + process.send(reply, active_adoptions(state)) + actor.continue(state) + } EditsFor(subject_uri, reply) -> { process.send( reply, @@ -263,6 +288,11 @@ fn count_for(state: State, release_uri: String) -> Int { |> list.length } +fn active_adoptions(state: State) -> List(Adoption) { + dict.values(state.adoptions) + |> list.filter(fn(a) { is_active_status(a.status) }) +} + fn all_counts(state: State) -> Dict(String, Int) { dict.values(state.adoptions) |> list.filter(fn(a) { is_active_status(a.status) }) diff --git a/server/src/at_record_server/catalog_index_postgres.gleam b/server/src/at_record_server/catalog_index_postgres.gleam index 371ce73..e9a5f7d 100644 --- a/server/src/at_record_server/catalog_index_postgres.gleam +++ b/server/src/at_record_server/catalog_index_postgres.gleam @@ -7,8 +7,8 @@ import at_record_server/catalog/row.{type BrowseRow, BrowseRow} import at_record_server/catalog_index.{ - type Adoption, type Edit, type Store, AdoptionOps, CursorOps, Edit, EditOps, - ReleaseOps, Store, + type Adoption, type Edit, type Store, Adoption, AdoptionOps, CursorOps, Edit, + EditOps, ReleaseOps, Store, } import atproto/blob.{Blob} import gleam/dict.{type Dict} @@ -38,6 +38,7 @@ pub fn table_store(conn: pog.Connection) -> Result(Store, String) { delete: fn(entry_uri) { delete_adoption(conn, entry_uri) }, count: fn(release_uri) { adoption_count(conn, release_uri) }, counts: fn() { adoption_counts(conn) }, + list: fn() { list_adoptions(conn) }, ), edits: EditOps( upsert: fn(edit) { upsert_edit(conn, edit) }, @@ -108,6 +109,15 @@ fn migrate(conn: pog.Connection) -> Result(Nil, String) { |> pog.execute(conn) |> result.replace_error("catalog_adoptions migration failed"), ) + // Idempotent add-column for a table that already existed before `source`: + // mirrors the catalog_releases format/label/master pattern above. + use _ <- result.try( + pog.query( + "alter table catalog_adoptions add column if not exists source text", + ) + |> pog.execute(conn) + |> result.replace_error("catalog_adoptions source migration failed"), + ) use _ <- result.try( pog.query( "create table if not exists catalog_edits ( @@ -273,19 +283,21 @@ fn list_releases(conn: pog.Connection) -> List(BrowseRow) { 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) + "insert into catalog_adoptions (entry_uri, did, release_uri, status, created_at, source) + values ($1, $2, $3, $4, $5, $6) on conflict (entry_uri) do update set did = excluded.did, release_uri = excluded.release_uri, status = excluded.status, - created_at = excluded.created_at", + created_at = excluded.created_at, + source = excluded.source", ) |> 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.parameter(pog.nullable(pog.text, adoption.source)) |> pog.execute(conn) case outcome { Ok(_) -> Nil @@ -338,6 +350,35 @@ fn adoption_counts(conn: pog.Connection) -> Dict(String, Int) { } } +fn adoption_row_decoder() -> decode.Decoder(Adoption) { + use entry_uri <- decode.field("entry_uri", decode.string) + use did <- decode.field("did", decode.string) + use release_uri <- decode.field("release_uri", decode.string) + use status <- decode.field("status", decode.string) + use created_at <- decode.field("created_at", decode.string) + use source <- decode.field("source", decode.optional(decode.string)) + decode.success(Adoption( + entry_uri:, + did:, + release_uri:, + status:, + created_at:, + source:, + )) +} + +fn list_adoptions(conn: pog.Connection) -> List(Adoption) { + case + pog.query("select entry_uri, did, release_uri, status, created_at, source + from catalog_adoptions where " <> active_status_filter) + |> pog.returning(adoption_row_decoder()) + |> pog.execute(conn) + { + Ok(pog.Returned(rows:, ..)) -> rows + Error(_) -> [] + } +} + fn upsert_edit(conn: pog.Connection, edit: Edit) -> Nil { let outcome = pog.query( diff --git a/server/src/at_record_server/feed_skeleton.gleam b/server/src/at_record_server/feed_skeleton.gleam new file mode 100644 index 0000000..16e9604 --- /dev/null +++ b/server/src/at_record_server/feed_skeleton.gleam @@ -0,0 +1,210 @@ +//// Pure feed-skeleton collapse policy: folds a page of raw adoption rows +//// (already filtered to the intended audience, sorted `created_at` desc) +//// into the reason-tagged `FeedItem`s the lexicon defines. Precedence: +//// subject convergence beats same-actor batching, which beats a lone +//// passthrough. No IO, no clock -- every window comparison works off the +//// rows' own `created_at` strings. + +import at_record/gen/defs +import at_record/gen/feed/get_feed_skeleton.{ + type FeedItem, EntryRef, FeedItem, FeedItemReasonReasonActorBatch, + FeedItemReasonReasonImport, FeedItemReasonReasonSingle, + FeedItemReasonReasonSubjectConverge, ReasonActorBatch, ReasonImport, + ReasonSingle, ReasonSubjectConverge, +} +import at_record_server/catalog_index.{type Adoption} +import atproto/uri +import gleam/dict +import gleam/list +import gleam/option.{None, Some} +import gleam/order +import gleam/time/duration +import gleam/time/timestamp.{type Timestamp} + +/// The span within which same-actor or same-subject rows collapse together. +const batch_window_hours = 6 + +/// >= this many same-did rows in one window collapse to an actor batch. +const actor_batch_threshold = 2 + +/// >= this many same-did rows in one window is always an import, batch or not. +const import_count_threshold = 5 + +/// >= this many distinct dids on the same release in one window converge. +const converge_threshold = 2 + +const discogs_source = "discogs" + +/// One adoption row paired with its parsed timestamp, so window comparisons +/// never re-parse `created_at`. +type Timed = + #(Adoption, Timestamp) + +pub fn collapse(adoptions: List(Adoption)) -> List(FeedItem) { + let #(converge_groups, leftover_groups) = + cluster_by_window(adoptions, fn(a) { a.release_uri }) + |> list.partition(fn(cluster) { + distinct_dids(cluster) >= converge_threshold + }) + let converge_items = list.map(converge_groups, to_converge_item) + let actor_items = + leftover_groups + |> list.flatten + |> list.map(fn(pair) { pair.0 }) + |> cluster_by_window(fn(a) { a.did }) + |> list.map(to_actor_item) + list.append(converge_items, actor_items) + |> list.sort(fn(a, b) { timestamp.compare(b.1, a.1) }) + |> list.map(fn(pair) { pair.0 }) +} + +fn distinct_dids(cluster: List(Timed)) -> Int { + cluster |> list.map(fn(pair) { pair.0.did }) |> list.unique |> list.length +} + +fn to_converge_item(cluster: List(Timed)) -> #(FeedItem, Timestamp) { + let #(latest_row, latest_ts) = latest_of(cluster) + let entries = + cluster + |> list.map(fn(pair) { pair.0.did }) + |> list.unique + |> list.map(fn(did) { + let #(row, _) = + latest_of(list.filter(cluster, fn(pair) { pair.0.did == did })) + EntryRef(actor: did, entry_id: entry_id_of(row)) + }) + // `cid` is only known once catalog.getRelease(uri) is hydrated downstream; + // this placeholder mirrors the same pattern as BrowseRow's cover_cid + // round-trip in catalog_index_postgres. + let subject = + defs.CatalogRef(cid: "", external_ids: None, uri: latest_row.release_uri) + let item = + FeedItem( + entries:, + reason: FeedItemReasonReasonSubjectConverge(ReasonSubjectConverge( + action: reason_action(latest_row.status), + subject:, + )), + ) + #(item, latest_ts) +} + +fn to_actor_item(cluster: List(Timed)) -> #(FeedItem, Timestamp) { + let #(latest_row, latest_ts) = latest_of(cluster) + let #(earliest_row, _) = earliest_of(cluster) + let count = list.length(cluster) + let has_discogs = + list.any(cluster, fn(pair) { pair.0.source == Some(discogs_source) }) + let entries = + list.map(cluster, fn(pair) { + EntryRef(actor: pair.0.did, entry_id: entry_id_of(pair.0)) + }) + let reason = case has_discogs, count >= import_count_threshold { + True, _ -> + FeedItemReasonReasonImport(ReasonImport(source: Some(discogs_source))) + False, True -> FeedItemReasonReasonImport(ReasonImport(source: None)) + False, False -> + case count >= actor_batch_threshold { + True -> + FeedItemReasonReasonActorBatch(ReasonActorBatch( + action: reason_action(latest_row.status), + window_end: latest_row.created_at, + window_start: earliest_row.created_at, + )) + False -> + FeedItemReasonReasonSingle( + ReasonSingle(action: reason_action(latest_row.status)), + ) + } + } + #(FeedItem(entries:, reason:), latest_ts) +} + +/// The feed reason vocabulary is only "owned"/"wanted" -- everything but an +/// explicit `wanted` shelf.entry action (acquisitions, regrades, ...) reads +/// as owned, since gone rows never reach `collapse` (the store already +/// filters them out). +fn reason_action(status: String) -> String { + case status { + "wanted" -> "wanted" + _ -> "owned" + } +} + +fn entry_id_of(a: Adoption) -> String { + uri.rkey(a.entry_uri) +} + +fn latest_of(cluster: List(Timed)) -> Timed { + let assert Ok(row) = list.last(cluster) + row +} + +fn earliest_of(cluster: List(Timed)) -> Timed { + let assert Ok(row) = list.first(cluster) + row +} + +fn parse_ts(a: Adoption) -> Timestamp { + case timestamp.parse_rfc3339(a.created_at) { + Ok(ts) -> ts + Error(_) -> timestamp.from_unix_seconds(0) + } +} + +/// Groups `rows` by `key_of`, then splits each group into window-bounded +/// clusters (ascending by time, each cluster's span from its earliest row +/// capped at `batch_window_hours`). +fn cluster_by_window( + rows: List(Adoption), + key_of: fn(Adoption) -> String, +) -> List(List(Timed)) { + rows + |> list.map(fn(a) { #(a, parse_ts(a)) }) + |> list.group(fn(pair) { key_of(pair.0) }) + |> dict.values + |> list.flat_map(fn(group) { + group + |> list.sort(fn(a, b) { timestamp.compare(a.1, b.1) }) + |> split_into_windows + }) +} + +fn split_into_windows(rows: List(Timed)) -> List(List(Timed)) { + do_split(rows, []) +} + +fn do_split(rows: List(Timed), acc: List(List(Timed))) -> List(List(Timed)) { + case rows { + [] -> list.reverse(acc) + [first, ..rest] -> { + let #(cluster, remaining) = take_within_window(first.1, rest, [first]) + do_split(remaining, [list.reverse(cluster), ..acc]) + } + } +} + +fn take_within_window( + start: Timestamp, + rows: List(Timed), + acc: List(Timed), +) -> #(List(Timed), List(Timed)) { + case rows { + [] -> #(acc, []) + [next, ..rest] -> + case within_window(start, next.1) { + True -> take_within_window(start, rest, [next, ..acc]) + False -> #(acc, rows) + } + } +} + +fn within_window(start: Timestamp, candidate: Timestamp) -> Bool { + // `difference(left, right)` is `right - left`, so this is `candidate - start` + // (non-negative: candidate is always >= start, rows are sorted ascending). + duration.compare( + timestamp.difference(start, candidate), + duration.hours(batch_window_hours), + ) + != order.Gt +} diff --git a/server/src/at_record_server/handlers/feed.gleam b/server/src/at_record_server/handlers/feed.gleam new file mode 100644 index 0000000..f39d544 --- /dev/null +++ b/server/src/at_record_server/handlers/feed.gleam @@ -0,0 +1,114 @@ +//// HTTP glue for the feed skeleton: collapses the caller's followed dids' +//// recent adoptions into feed items, falling back to a network-wide slice +//// (everyone, viewer included) when the follow graph produces nothing. + +import at_record/gen/feed/get_feed_skeleton as feed_gen +import at_record_server/catalog_index.{type Adoption} +import at_record_server/context.{type Context, require_session, with_pds_client} +import at_record_server/feed_skeleton +import at_record_server/graph_follows +import at_record_server/oauth/sessions.{type OauthSession} +import at_record_server/pagination +import atproto/xrpc.{type Client} +import gleam/int +import gleam/json +import gleam/list +import gleam/option.{type Option, None, Some} +import gleam/result +import gleam/string +import wisp.{type Request, type Response} + +// No colons: `pagination`'s cursor format joins `namespace <> ":" <> last_id` +// and splits on the *first* colon to decode, so a namespace containing one +// (a bare did always does, e.g. "did:plc:...") would never round-trip. +const following_namespace_prefix = "feed-following-" + +const network_namespace = "feed-network" + +pub fn get_feed_skeleton(req: Request, ctx: Context) -> Response { + use id, session <- require_session(req, ctx) + use client, session <- with_pds_client(ctx, id, session) + let query = wisp.get_query(req) + let cursor = list.key_find(query, "cursor") |> option.from_result + let limit = + Some( + list.key_find(query, "limit") + |> result.try(int.parse) + |> option.from_result + |> option.unwrap(pagination.default_limit), + ) + let adoptions = ctx.catalog_index.adoptions.list() + let followed = followed_dids(client, session) + let #(items, next_cursor, fallback) = + resolve_page(adoptions, followed, session.did, cursor, limit) + json.object([ + #("items", json.array(items, feed_gen.encode_feed_item)), + #("fallback", json.bool(fallback)), + ..case next_cursor { + Some(c) -> [#("cursor", json.string(c))] + None -> [] + } + ]) + |> json.to_string + |> wisp.json_response(200) +} + +/// The caller's followed dids; a load failure (PDS hiccup) reads the same as +/// an empty follow graph rather than failing the whole feed -- that just +/// falls back to the network-wide slice below. +fn followed_dids(client: Client, session: OauthSession) -> List(String) { + case graph_follows.load(client, session) { + Ok(stored) -> list.map(stored, fn(s) { s.value.subject }) + Error(_) -> [] + } +} + +/// Followed-mode first; falls back to network-wide (unfiltered, viewer +/// included -- simpler than re-excluding them from their own discovery feed) +/// when the follow graph is empty, or when a first request (no cursor yet) +/// finds nothing from people the caller follows. A cursor already scoped to +/// one namespace never resumes the other: `pagination.page` restarts a +/// foreign-namespace cursor from the top of whichever list it's handed. +fn resolve_page( + adoptions: List(Adoption), + followed: List(String), + viewer_did: String, + cursor: Option(String), + limit: Option(Int), +) -> #(List(feed_gen.FeedItem), Option(String), Bool) { + let followed_rows = + adoptions + |> list.filter(fn(a) { + a.did != viewer_did && list.contains(followed, a.did) + }) + |> sort_desc + let #(followed_page, followed_cursor) = + pagination.page( + followed_rows, + fn(a: Adoption) { a.entry_uri }, + following_namespace_prefix <> string.replace(viewer_did, ":", "-"), + cursor, + limit, + ) + let needs_fallback = + list.is_empty(followed) + || { cursor == None && list.is_empty(followed_page) } + case needs_fallback { + True -> { + let #(network_page, network_cursor) = + pagination.page( + sort_desc(adoptions), + fn(a: Adoption) { a.entry_uri }, + network_namespace, + cursor, + limit, + ) + #(feed_skeleton.collapse(network_page), network_cursor, True) + } + False -> #(feed_skeleton.collapse(followed_page), followed_cursor, False) + } +} + +fn sort_desc(rows: List(Adoption)) -> List(Adoption) { + list.sort(rows, fn(a, b) { string.compare(b.created_at, a.created_at) }) +} diff --git a/server/src/at_record_server/jetstream_consumer.gleam b/server/src/at_record_server/jetstream_consumer.gleam index 48bc67c..55bbd19 100644 --- a/server/src/at_record_server/jetstream_consumer.gleam +++ b/server/src/at_record_server/jetstream_consumer.gleam @@ -366,6 +366,9 @@ fn upsert_adoption(state: State, frame: Frame) -> State { release_uri: ref.uri, status: entry.action, created_at: entry.created_at, + source: catalog_index.source_label( + option.then(entry.source, fn(s) { s.origin }), + ), )) state } diff --git a/server/src/at_record_server/router.gleam b/server/src/at_record_server/router.gleam index d851ddc..456b93d 100644 --- a/server/src/at_record_server/router.gleam +++ b/server/src/at_record_server/router.gleam @@ -8,6 +8,7 @@ import at_record_server/handlers/cover_proxy import at_record_server/handlers/crate_overlap import at_record_server/handlers/discogs import at_record_server/handlers/edit_inbox +import at_record_server/handlers/feed import at_record_server/handlers/graph import at_record_server/handlers/oauth import at_record_server/handlers/public_shelf @@ -84,6 +85,8 @@ fn dispatch_xrpc( crate_overlap.get_crate_overlap(req, ctx) "dev.mokkenstorm.crate.graph.followUser", Post -> graph.follow(req, ctx) "dev.mokkenstorm.crate.graph.unfollowUser", Post -> graph.unfollow(req, ctx) + "dev.mokkenstorm.crate.feed.getFeedSkeleton", Get -> + feed.get_feed_skeleton(req, ctx) "dev.mokkenstorm.crate.discogs.searchReleases", Get -> discogs.search(req, ctx) "dev.mokkenstorm.crate.discogs.searchArtists", Get -> diff --git a/server/test/catalog_index_postgres_test.gleam b/server/test/catalog_index_postgres_test.gleam index c6e72d9..8c70137 100644 --- a/server/test/catalog_index_postgres_test.gleam +++ b/server/test/catalog_index_postgres_test.gleam @@ -59,6 +59,7 @@ fn adoption( release_uri:, status:, created_at:, + source: None, ) } @@ -140,6 +141,24 @@ pub fn gone_edges_are_excluded_but_reacquiring_counts_again_test() { store.adoptions.delete(entry_uri) } +pub fn adoption_list_round_trips_source_test() { + use store <- with_store() + let release = "at://did:plc:pgtest/dev.mokkenstorm.crate.catalog.release/pg4" + let entry_uri = "at://did:plc:pgtest/dev.mokkenstorm.crate.shelf.entry/pg-e3" + store.adoptions.upsert(Adoption( + entry_uri:, + did: "did:plc:pgtest", + release_uri: release, + status: "acquired", + created_at: "2024-01-01T00:00:00Z", + source: Some("discogs"), + )) + let assert Ok(found) = + store.adoptions.list() |> find(fn(a) { a.entry_uri == entry_uri }) + assert found.source == Some("discogs") + store.adoptions.delete(entry_uri) +} + pub fn cursor_round_trip_test() { use store <- with_store() store.cursor.save(1_700_000_000_123_456) diff --git a/server/test/catalog_index_test.gleam b/server/test/catalog_index_test.gleam index 7514574..7c103c2 100644 --- a/server/test/catalog_index_test.gleam +++ b/server/test/catalog_index_test.gleam @@ -36,7 +36,14 @@ fn adoption( status status: String, created_at created_at: String, ) -> Adoption { - Adoption(entry_uri:, did: "did:plc:a", release_uri:, status:, created_at:) + Adoption( + entry_uri:, + did: "did:plc:a", + release_uri:, + status:, + created_at:, + source: None, + ) } pub fn upsert_and_list_test() { @@ -161,6 +168,30 @@ pub fn gone_then_reacquired_entry_counts_again_test() { assert store.adoptions.count(release) == 1 } +pub fn adoptions_list_returns_only_active_rows_with_source_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", + release_uri: release, + status: "acquired", + created_at: "2024-01-01T00:00:00Z", + source: Some("discogs"), + )) + store.adoptions.upsert(Adoption( + entry_uri: "at://did:plc:b/dev.mokkenstorm.crate.shelf.entry/e2", + did: "did:plc:b", + release_uri: release, + status: "sold", + created_at: "2024-01-02T00:00:00Z", + source: None, + )) + let assert [only] = store.adoptions.list() + assert only.did == "did:plc:a" + assert only.source == Some("discogs") +} + 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" diff --git a/server/test/feed_handler_test.gleam b/server/test/feed_handler_test.gleam new file mode 100644 index 0000000..c5e0c22 --- /dev/null +++ b/server/test/feed_handler_test.gleam @@ -0,0 +1,301 @@ +//// End-to-end wiring tests for the `feed.getFeedSkeleton` handler: followed +//// filtering, viewer exclusion, the fallback flip to network-wide when the +//// follow graph produces nothing, cursor round-tripping within a namespace, +//// and the limit cap. + +import at_record_server/catalog_index.{type Adoption, Adoption} +import at_record_server/context.{type Context, Context} +import at_record_server/handlers/feed +import at_record_server/oauth/config +import at_record_server/oauth/session_store +import at_record_server/oauth/sessions +import atproto/xrpc +import gleam/bit_array +import gleam/dynamic/decode +import gleam/http +import gleam/http/response +import gleam/int +import gleam/json +import gleam/list +import gleam/option.{None} +import gleam/time/calendar +import gleam/time/duration +import gleam/time/timestamp +import support +import wisp +import wisp/simulate + +const viewer_did = "did:plc:x" + +const session_cookie = "ar_oauth_sid" + +const far_future = 9_999_999_999 + +const feed_path = "/xrpc/dev.mokkenstorm.crate.feed.getFeedSkeleton" + +fn a_session() -> sessions.OauthSession { + support.stub_session_with( + access_token: "at", + refresh_token: "rt", + expires_at: far_future, + ) +} + +fn adoption( + did did: String, + rkey rkey: String, + release_uri release_uri: String, + status status: String, + created_at created_at: String, +) -> Adoption { + Adoption( + entry_uri: "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey, + did:, + release_uri:, + status:, + created_at:, + source: None, + ) +} + +/// The nth day since the epoch, RFC3339 -- a cheap way to generate rows +/// spaced far enough apart (24h) that none of them ever collapse together. +fn nth_day(n: Int) -> String { + timestamp.from_unix_seconds(0) + |> timestamp.add(duration.hours(24 * n)) + |> timestamp.to_rfc3339(calendar.utc_offset) +} + +fn follow_record_json(rkey: String, subject: String) -> json.Json { + json.object([ + #( + "uri", + json.string( + "at://" <> viewer_did <> "/dev.mokkenstorm.crate.graph.follow/" <> rkey, + ), + ), + #("cid", json.string("bafyfollow" <> rkey)), + #( + "value", + json.object([ + #("createdAt", json.string("2026-01-01T00:00:00Z")), + #("subject", json.string(subject)), + ]), + ), + ]) +} + +fn list_records_body(records: List(json.Json)) -> String { + json.object([#("records", json.preprocessed_array(records))]) + |> json.to_string +} + +fn network_client(follows_body: String) -> xrpc.Client { + xrpc.Client(send: fn(req) { + case req.path { + "/xrpc/com.atproto.repo.listRecords" -> + Ok(response.Response(200, [], bit_array.from_string(follows_body))) + _ -> panic as { "unexpected path: " <> req.path } + } + }) +} + +fn test_context( + client: xrpc.Client, + index: catalog_index.Store, +) -> #(Context, config.Config) { + let cfg = support.stub_config(client) + let ctx = + Context( + ..support.stub_context_with( + cfg, + support.unreachable_catalog_deps(), + fn(_req) { Error("unused") }, + [], + ), + catalog_index: index, + ) + #(ctx, cfg) +} + +fn authed_get(path: String, cfg: config.Config) -> wisp.Request { + let assert Ok(id) = session_store.create(cfg.sessions, a_session()) + simulate.request(http.Get, path) + |> simulate.cookie(session_cookie, id, wisp.Signed) +} + +fn item_actors(body: String) -> List(String) { + let decoder = + decode.at( + ["items"], + decode.list(decode.at( + ["entries"], + decode.list(decode.at(["actor"], decode.string)), + )), + ) + let assert Ok(nested) = json.parse(body, decoder) + list.flatten(nested) +} + +fn item_count(body: String) -> Int { + let assert Ok(items) = + json.parse(body, decode.at(["items"], decode.list(decode.dynamic))) + list.length(items) +} + +pub fn followed_rows_are_included_and_viewer_and_non_followed_are_excluded_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: "did:plc:f1", + rkey: "e1", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + status: "acquired", + created_at: nth_day(3), + )) + store.adoptions.upsert(adoption( + did: "did:plc:n1", + rkey: "e2", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2", + status: "acquired", + created_at: nth_day(2), + )) + store.adoptions.upsert(adoption( + did: viewer_did, + rkey: "e3", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r3", + status: "acquired", + created_at: nth_day(1), + )) + let client = + network_client( + list_records_body([follow_record_json("3aaa", "did:plc:f1")]), + ) + let #(ctx, cfg) = test_context(client, store) + let resp = feed.get_feed_skeleton(authed_get(feed_path, cfg), ctx) + assert resp.status == 200 + let body = simulate.read_body(resp) + assert support.field_bool(body, ["fallback"]) == Ok(False) + assert item_actors(body) == ["did:plc:f1"] +} + +pub fn empty_follow_graph_falls_back_to_network_wide_including_the_viewer_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: viewer_did, + rkey: "e1", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + status: "acquired", + created_at: nth_day(2), + )) + store.adoptions.upsert(adoption( + did: "did:plc:n1", + rkey: "e2", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2", + status: "acquired", + created_at: nth_day(1), + )) + let client = network_client(list_records_body([])) + let #(ctx, cfg) = test_context(client, store) + let resp = feed.get_feed_skeleton(authed_get(feed_path, cfg), ctx) + assert resp.status == 200 + let body = simulate.read_body(resp) + assert support.field_bool(body, ["fallback"]) == Ok(True) + let actors = item_actors(body) + assert list.contains(actors, viewer_did) + assert list.contains(actors, "did:plc:n1") +} + +pub fn nonempty_follows_with_no_followed_adoptions_falls_back_on_the_first_page_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: "did:plc:n1", + rkey: "e1", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + status: "acquired", + created_at: nth_day(1), + )) + // Followed, but they have no adoptions of their own. + let client = + network_client( + list_records_body([follow_record_json("3bbb", "did:plc:f1")]), + ) + let #(ctx, cfg) = test_context(client, store) + let resp = feed.get_feed_skeleton(authed_get(feed_path, cfg), ctx) + assert resp.status == 200 + let body = simulate.read_body(resp) + assert support.field_bool(body, ["fallback"]) == Ok(True) + assert item_actors(body) == ["did:plc:n1"] +} + +pub fn cursor_stays_in_the_followed_namespace_across_pages_test() { + let assert Ok(store) = catalog_index.start() + store.adoptions.upsert(adoption( + did: "did:plc:f1", + rkey: "e1", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r1", + status: "acquired", + created_at: nth_day(3), + )) + store.adoptions.upsert(adoption( + did: "did:plc:f2", + rkey: "e2", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r2", + status: "acquired", + created_at: nth_day(2), + )) + store.adoptions.upsert(adoption( + did: "did:plc:f3", + rkey: "e3", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r3", + status: "acquired", + created_at: nth_day(1), + )) + let client = + network_client( + list_records_body([ + follow_record_json("3aaa", "did:plc:f1"), + follow_record_json("3bbb", "did:plc:f2"), + follow_record_json("3ccc", "did:plc:f3"), + ]), + ) + let #(ctx, cfg) = test_context(client, store) + let first = + feed.get_feed_skeleton(authed_get(feed_path <> "?limit=1", cfg), ctx) + let first_body = simulate.read_body(first) + assert support.field_bool(first_body, ["fallback"]) == Ok(False) + assert item_actors(first_body) == ["did:plc:f1"] + let assert Ok(cursor) = support.field_string(first_body, ["cursor"]) + let second = + feed.get_feed_skeleton( + authed_get(feed_path <> "?limit=1&cursor=" <> cursor, cfg), + ctx, + ) + let second_body = simulate.read_body(second) + assert support.field_bool(second_body, ["fallback"]) == Ok(False) + assert item_actors(second_body) == ["did:plc:f2"] +} + +pub fn limit_is_capped_at_the_lexicon_default_test() { + let assert Ok(store) = catalog_index.start() + list.repeat(0, 105) + |> list.index_map(fn(_, i) { i + 1 }) + |> list.each(fn(n) { + store.adoptions.upsert(adoption( + did: "did:plc:f1", + rkey: "e" <> int.to_string(n), + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/r" + <> int.to_string(n), + status: "acquired", + created_at: nth_day(n), + )) + }) + let client = + network_client( + list_records_body([follow_record_json("3aaa", "did:plc:f1")]), + ) + let #(ctx, cfg) = test_context(client, store) + let resp = feed.get_feed_skeleton(authed_get(feed_path, cfg), ctx) + let body = simulate.read_body(resp) + assert item_count(body) == 100 + assert support.field_present(body, ["cursor"]) +} diff --git a/server/test/feed_skeleton_test.gleam b/server/test/feed_skeleton_test.gleam new file mode 100644 index 0000000..b8bd2ea --- /dev/null +++ b/server/test/feed_skeleton_test.gleam @@ -0,0 +1,371 @@ +//// Pure collapse-policy tests for `feed_skeleton.collapse`: the reason +//// precedence (converge > import > actorBatch > single), the 6h window +//// boundary, and that ordering by recency survives collapsing. + +import at_record/gen/feed/get_feed_skeleton.{ + FeedItem, FeedItemReasonReasonActorBatch, FeedItemReasonReasonImport, + FeedItemReasonReasonSingle, FeedItemReasonReasonSubjectConverge, + ReasonActorBatch, ReasonImport, ReasonSingle, +} +import at_record_server/catalog_index.{type Adoption, Adoption} +import at_record_server/feed_skeleton +import gleam/list +import gleam/option.{None, Some} + +const release_a = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/ra" + +const release_b = "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/rb" + +fn entry_uri(did: String, rkey: String) -> String { + "at://" <> did <> "/dev.mokkenstorm.crate.shelf.entry/" <> rkey +} + +fn row( + did did: String, + rkey rkey: String, + release_uri release_uri: String, + status status: String, + created_at created_at: String, + source source: option.Option(String), +) -> Adoption { + Adoption( + entry_uri: entry_uri(did, rkey), + did:, + release_uri:, + status:, + created_at:, + source:, + ) +} + +pub fn single_passthrough_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + ] + let assert [FeedItem(entries:, reason: FeedItemReasonReasonSingle(inner))] = + feed_skeleton.collapse(rows) + assert entries + == [get_feed_skeleton.EntryRef(actor: "did:plc:a", entry_id: "e1")] + assert inner == ReasonSingle(action: "owned") +} + +pub fn single_passthrough_reports_wanted_action_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "wanted", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + ] + let assert [FeedItem(reason: FeedItemReasonReasonSingle(inner), ..)] = + feed_skeleton.collapse(rows) + assert inner.action == "wanted" +} + +pub fn actor_batch_at_the_window_boundary_is_included_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e2", + release_uri: release_b, + status: "acquired", + created_at: "2026-01-01T06:00:00Z", + source: None, + ), + ] + let assert [FeedItem(entries:, reason: FeedItemReasonReasonActorBatch(inner))] = + feed_skeleton.collapse(rows) + assert list.length(entries) == 2 + assert inner + == ReasonActorBatch( + action: "owned", + window_end: "2026-01-01T06:00:00Z", + window_start: "2026-01-01T00:00:00Z", + ) +} + +pub fn actor_batch_just_past_the_window_boundary_splits_into_singles_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e2", + release_uri: release_b, + status: "acquired", + created_at: "2026-01-01T06:00:01Z", + source: None, + ), + ] + let items = feed_skeleton.collapse(rows) + assert list.length(items) == 2 + assert list.all(items, fn(item) { + case item.reason { + FeedItemReasonReasonSingle(_) -> True + _ -> False + } + }) +} + +pub fn import_via_count_collapses_five_same_did_rows_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T01:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e2", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T02:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e3", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T03:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e4", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T04:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e5", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T05:00:00Z", + source: None, + ), + ] + let assert [FeedItem(entries:, reason: FeedItemReasonReasonImport(inner))] = + feed_skeleton.collapse(rows) + assert list.length(entries) == 5 + assert inner == ReasonImport(source: None) +} + +pub fn import_via_source_collapses_even_a_single_row_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: Some("discogs"), + ), + ] + let assert [FeedItem(reason: FeedItemReasonReasonImport(inner), ..)] = + feed_skeleton.collapse(rows) + assert inner == ReasonImport(source: Some("discogs")) +} + +pub fn import_via_source_wins_over_actor_batch_for_a_small_cluster_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: Some("discogs"), + ), + row( + did: "did:plc:a", + rkey: "e2", + release_uri: release_b, + status: "acquired", + created_at: "2026-01-01T01:00:00Z", + source: None, + ), + ] + let assert [FeedItem(reason: FeedItemReasonReasonImport(inner), ..)] = + feed_skeleton.collapse(rows) + assert inner == ReasonImport(source: Some("discogs")) +} + +pub fn subject_converge_with_two_actors_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + row( + did: "did:plc:b", + rkey: "e2", + release_uri: release_a, + status: "wanted", + created_at: "2026-01-01T01:00:00Z", + source: None, + ), + ] + let assert [ + FeedItem(entries:, reason: FeedItemReasonReasonSubjectConverge(inner)), + ] = feed_skeleton.collapse(rows) + assert list.length(entries) == 2 + assert inner.subject.uri == release_a + // The freshest row's action ("wanted" at 01:00) drives the single action label. + assert inner.action == "wanted" +} + +pub fn subject_converge_with_three_actors_test() { + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + row( + did: "did:plc:b", + rkey: "e2", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T01:00:00Z", + source: None, + ), + row( + did: "did:plc:c", + rkey: "e3", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T02:00:00Z", + source: None, + ), + ] + let assert [FeedItem(entries:, ..)] = feed_skeleton.collapse(rows) + assert list.length(entries) == 3 +} + +pub fn converge_takes_precedence_over_the_same_did_also_batching_test() { + // did:plc:a adopts release_a, then did:plc:b adopts the same release soon + // after (converge), while did:plc:a also has a second, unrelated row close + // enough in time to have otherwise formed an actor batch with the first. + let rows = [ + row( + did: "did:plc:a", + rkey: "e1", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + row( + did: "did:plc:b", + rkey: "e2", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-01T01:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e3", + release_uri: release_b, + status: "acquired", + created_at: "2026-01-01T02:00:00Z", + source: None, + ), + ] + let items = feed_skeleton.collapse(rows) + // e1+e2 converge into one item; e3 (did:plc:a's only remaining row) is a + // lone single, not merged with the consumed e1. + assert list.length(items) == 2 + let assert Ok(converge) = + list.find(items, fn(item) { + case item.reason { + FeedItemReasonReasonSubjectConverge(_) -> True + _ -> False + } + }) + assert list.length(converge.entries) == 2 + let assert Ok(single) = + list.find(items, fn(item) { + case item.reason { + FeedItemReasonReasonSingle(_) -> True + _ -> False + } + }) + assert single.entries + == [ + get_feed_skeleton.EntryRef(actor: "did:plc:a", entry_id: "e3"), + ] +} + +pub fn ordering_by_recency_is_preserved_for_unrelated_singles_test() { + let rows = [ + row( + did: "did:plc:c", + rkey: "e3", + release_uri: release_a, + status: "acquired", + created_at: "2026-01-03T00:00:00Z", + source: None, + ), + row( + did: "did:plc:b", + rkey: "e2", + release_uri: release_b, + status: "acquired", + created_at: "2026-01-02T00:00:00Z", + source: None, + ), + row( + did: "did:plc:a", + rkey: "e1", + release_uri: "at://did:plc:pub/dev.mokkenstorm.crate.catalog.release/rc", + status: "acquired", + created_at: "2026-01-01T00:00:00Z", + source: None, + ), + ] + let items = feed_skeleton.collapse(rows) + let dids = + list.map(items, fn(item) { + let assert [only] = item.entries + only.actor + }) + assert dids == ["did:plc:c", "did:plc:b", "did:plc:a"] +} diff --git a/server/test/support.gleam b/server/test/support.gleam index 66039cf..f02ec6c 100644 --- a/server/test/support.gleam +++ b/server/test/support.gleam @@ -55,6 +55,7 @@ fn empty_catalog_index() -> catalog_index.Store { delete: fn(_) { Nil }, count: fn(_) { 0 }, counts: fn() { dict.new() }, + list: fn() { [] }, ), edits: catalog_index.EditOps( upsert: fn(_) { Nil }, @@ -240,6 +241,10 @@ pub fn field_int(body: String, path: List(String)) -> Result(Int, Nil) { json.parse(body, decode.at(path, decode.int)) |> result.replace_error(Nil) } +pub fn field_bool(body: String, path: List(String)) -> Result(Bool, Nil) { + json.parse(body, decode.at(path, decode.bool)) |> result.replace_error(Nil) +} + pub fn field_strings( body: String, path: List(String),