//// The catalog.edit inbox. `list_proposals` reads entirely off the index //// (E2 of the appview-first roadmap, hard swap per its decision 6): the //// caller's own current releases come from `catalog/source.Source` (the //// same `catalog_index`-backed port browse's resolution layer reads), and //// pending edits against them come from one batched //// `catalog_index.EditOps.list_for_subjects` call. The old per-release //// Constellation-backlink fan-out plus per-proposal Slingshot `fetch_edit` //// are gone: everything this handler needs is already indexed off the //// firehose/backfill. `own_current_releases` (a live PDS `listRecords`) and //// `current_only` survive unchanged: the edit-inbox *apply* flow still //// needs the full record to mint a superseding release from, which the //// index's `BrowseRow` doesn't carry. import atproto/repo import atproto/xrpc.{type Client} import crate/gen/catalog/edit as catalog_edit import crate/gen/catalog/list_edit_proposals.{type ProposalRow, ProposalRow} import crate/gen/catalog/release as catalog_release import crate/gen/defs import crate_server/catalog/row.{type BrowseRow} import crate_server/catalog/source as catalog_source import crate_server/catalog_index.{type Edit, type EditOps} import crate_server/identity_resolver import crate_server/oauth/sessions.{type OauthSession} import gleam/dict.{type Dict} import gleam/dynamic/decode import gleam/int import gleam/json import gleam/list import gleam/option.{type Option, None, Some} import gleam/set import gleam/string import wisp // How many of the caller's own current releases are queried per request; // releases beyond this (newest-first) are dropped and logged. Bounds the // `list_for_subjects` batch even though it's now one query, not a fan-out, // so a caller with a huge catalog still gets a predictable response size. pub const release_fanout_cap = 50 // The merged proposal list's final cap, across all queried releases. pub const proposal_cap = 50 pub type OwnRelease { OwnRelease(ref: defs.CatalogRef, value: catalog_release.CatalogRelease) } /// The caller's own `catalog.release` records, reduced to the current tip of /// each supersedes chain: a release superseded by another own release is /// dropped, since a proposal against it is stale. Still a live PDS read: the /// edit-inbox *apply* flow (`handlers/edit_inbox.apply_to_target`) needs the /// full decoded record to mint a superseding one from, which the index's /// `BrowseRow` doesn't carry. pub fn own_current_releases( client: Client, session: OauthSession, ) -> List(OwnRelease) { case repo.list_records( client, session.pds, session.access_token, session.did, catalog_release.collection, own_row_decoder(), ) { Error(_) -> [] Ok(rows) -> current_only(rows) } } fn own_row_decoder() -> decode.Decoder(OwnRelease) { use uri <- decode.field("uri", decode.string) use cid <- decode.field("cid", decode.string) use value <- decode.field("value", catalog_release.catalog_release_decoder()) decode.success(OwnRelease( ref: defs.CatalogRef(uri:, cid:, external_ids: None), value:, )) } /// Drop any release another own release's `supersedes` points at: only the /// newest link in each chain is current. pub fn current_only(rows: List(OwnRelease)) -> List(OwnRelease) { current_only_by(rows, fn(r) { r.ref.uri }, fn(r) { r.value.supersedes |> option.map(fn(s) { s.uri }) }) } /// The same supersedes-chain filter as `current_only`, generalized over /// whichever shape carries a uri and an optional supersedes-uri: shared so /// the live-fetched `OwnRelease` path (`current_only`, used by apply) and /// the index-backed `BrowseRow` path (`current_release_rows`, used by /// `list_proposals`) can never drift apart on what "current" means. fn current_only_by( items: List(a), self_uri: fn(a) -> String, supersedes_uri: fn(a) -> Option(String), ) -> List(a) { let superseded = items |> list.filter_map(fn(r) { supersedes_uri(r) |> option.to_result(Nil) }) |> set.from_list items |> list.filter(fn(r) { !set.contains(superseded, self_uri(r)) }) } fn current_release_rows(rows: List(BrowseRow)) -> List(BrowseRow) { current_only_by(rows, fn(r) { r.uri }, fn(r) { r.supersedes }) } /// The caller's own current releases straight off the index: every indexed /// row published by `did`, reduced to the current supersedes tip, newest /// first, capped at `release_fanout_cap`. fn own_current_rows( variant_source: catalog_source.Source, did: String, ) -> List(BrowseRow) { variant_source.releases() |> list.filter(fn(r) { r.publisher_did == did }) |> current_release_rows |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) |> cap_rows(release_fanout_cap, "own-release fan-out") } /// The proposal rows across every capped own-current release, newest first, /// capped again as a merged set (E2 of the appview-first roadmap: one batched /// `list_for_subjects` read replaces the old per-release Constellation /// backlink fan-out and per-proposal Slingshot `fetch_edit`). Proposer /// handles are resolved in one batched `resolve_many` call across every /// matching edit, rather than once per proposal. pub fn list_proposals( variant_source: catalog_source.Source, edits: EditOps, session: OauthSession, identity: identity_resolver.Resolver, ) -> List(ProposalRow) { let own = own_current_rows(variant_source, session.did) let own_uris = own |> list.map(fn(r) { r.uri }) let by_uri = own |> list.map(fn(r) { #(r.uri, r) }) |> dict.from_list let subject_edits = edits.list_for_subjects(own_uris) |> list.filter(fn(e) { e.entity == "release" && e.did != session.did }) let handles = subject_edits |> list.map(fn(e) { e.did }) |> identity_resolver.resolve_many(identity, _) |> dict.map_values(fn(_did, identity) { identity.handle }) let adoption_count = subject_adoption_count(variant_source) subject_edits |> list.filter_map(fn(edit) { to_proposal(handles, adoption_count, by_uri, edit) |> option.to_result(Nil) }) // A subject could in principle appear more than once if the index ever // carries a duplicate write; cheap to guard against regardless. |> unique_by_uri |> list.sort(fn(a, b) { string.compare(b.created_at, a.created_at) }) |> cap_rows(proposal_cap, "merged proposals") } /// A release uri's adoption count from the same catalog index the browse /// pipeline reads (owned + wanted, excluding gone; see /// `catalog_index.adoption_count`), or `None` when the release hasn't been /// indexed yet (best-effort: ingestion lag, not an error). fn subject_adoption_count( variant_source: catalog_source.Source, ) -> fn(String) -> Option(Int) { let indexed = variant_source.releases() |> list.map(fn(r) { r.uri }) |> set.from_list fn(uri) { case set.contains(indexed, uri) { True -> Some(variant_source.adoption_count(uri)) False -> None } } } fn unique_by_uri(rows: List(ProposalRow)) -> List(ProposalRow) { rows |> list.fold(#(set.new(), []), fn(acc, row) { let #(seen, kept) = acc case set.contains(seen, row.uri) { True -> acc False -> #(set.insert(seen, row.uri), [row, ..kept]) } }) |> fn(acc) { acc.1 } |> list.reverse } fn cap_rows(rows: List(a), cap: Int, what: String) -> List(a) { case list.length(rows) > cap { False -> rows True -> { wisp.log_warning( "edit-inbox: capped " <> what <> " at " <> int.to_string(cap) <> ", dropped " <> int.to_string(list.length(rows) - cap), ) list.take(rows, cap) } } } fn to_proposal( handles: Dict(String, String), adoption_count: fn(String) -> Option(Int), by_uri: Dict(String, BrowseRow), edit: Edit, ) -> Option(ProposalRow) { use owned <- option.then( dict.get(by_uri, edit.subject_uri) |> option.from_result, ) use fields_json <- option.then(edit.fields) use fields <- option.then( json.parse(fields_json, catalog_edit.release_fields_decoder()) |> option.from_result, ) Some(ProposalRow( uri: edit.edit_uri, cid: edit.cid, proposer_did: edit.did, proposer_handle: dict.get(handles, edit.did) |> option.from_result, target_uri: owned.uri, release_title: owned.title, current: Some(current_fields(owned)), fields: Some(fields), rationale: edit.rationale, created_at: edit.created_at, subject_adoption_count: adoption_count(owned.uri), )) } /// The target release's current values for the fields a proposal can touch, /// shaped like `ReleaseFields` so the client can diff against `fields` /// without a second fetch. `BrowseRow.genres`/`.styles` don't distinguish /// "field absent" from "field present but empty" the way the decoded /// `CatalogRelease` did, so an empty list here reads as `None`. fn current_fields(row: BrowseRow) -> catalog_edit.ReleaseFields { catalog_edit.ReleaseFields( title: Some(row.title), released: row.released, country: row.country, genres: non_empty_list(row.genres), styles: non_empty_list(row.styles), thumb_url: None, cover: None, ) } fn non_empty_list(items: List(a)) -> Option(List(a)) { case items { [] -> None _ -> Some(items) } } /// The proposal's `#releaseFields` variant, when it carries one (the only /// variant this vertical slice understands; other entities/variants skip). /// Used by the apply flow against a live-fetched `CatalogEdit`, not the /// indexed `Edit` (see `handlers/edit_inbox.apply`). pub fn release_fields( edit: catalog_edit.CatalogEdit, ) -> Option(catalog_edit.ReleaseFields) { case edit.fields { Some(catalog_edit.CatalogEditFieldsReleaseFields(f)) -> Some(f) _ -> None } }