diff --git a/server/src/at_record_server.gleam b/server/src/at_record_server.gleam index 34ebe0e..ed1e5c5 100644 --- a/server/src/at_record_server.gleam +++ b/server/src/at_record_server.gleam @@ -1,22 +1,31 @@ +import at_record/gen/catalog/edit as catalog_edit +import at_record/gen/client as generated_client +import at_record/gen/shelf/entry as shelf_entry import at_record_server/atproto_client import at_record_server/browse import at_record_server/catalog_index import at_record_server/config_env import at_record_server/jetstream_consumer -import at_record_server/known_users.{type Store as KnownUsers} +import at_record_server/known_users.{type KnownUser, type Store as KnownUsers} import at_record_server/oauth/config import at_record_server/oauth/keys import at_record_server/oauth/store import at_record_server/router import at_record_server/wiring import atproto/xrpc.{type Client} +import gleam/dynamic/decode import gleam/erlang/process import gleam/int import gleam/list +import gleam/option.{None, Some} import mist import wisp import wisp/wisp_mist +// The per-user cap on the backfill fan-out; matches `browse`'s own cap so a +// full backfill is never a heavier hit than a single browse load per user. +const per_user_limit = 50 + pub fn main() -> Nil { wisp.configure_logger() let port = config_env.port() @@ -25,13 +34,15 @@ pub fn main() -> Nil { let base_url = config_env.base_url(port) let client = atproto_client.client() let assert Ok(store) = store.start() - let #(sessions, discogs_creds, known_users) = config_env.stores() - let assert Ok(catalog_index) = catalog_index.start() + let #(sessions, discogs_creds, known_users, catalog_index) = + config_env.stores() backfill_catalog_index(client, known_users, catalog_index) jetstream_consumer.start( config_env.jetstream_url(), known_users, catalog_index, + client, + resolver, ) let oauth = config.new( @@ -72,9 +83,11 @@ pub fn main() -> Nil { process.sleep_forever() } -/// One-shot seed for the prototype `catalog_index`: the same per-user -/// `listRecords` fan-out `/browse` runs live on every request, paid once at -/// boot instead. `jetstream_consumer` keeps the index live from here. +/// One-shot seed for `catalog_index`: the same per-user `listRecords` +/// fan-out `/browse` runs live on every request, paid once at boot instead, +/// plus the adoption/edit edges each known user's shelf.entry and +/// catalog.edit records carry. `jetstream_consumer` keeps the index live +/// from here. fn backfill_catalog_index( client: Client, known_users: KnownUsers, @@ -84,6 +97,8 @@ fn backfill_catalog_index( users |> list.each(fn(user) { browse.fetch_user_releases(client, user) |> list.each(index.upsert) + backfill_adoptions(client, user, index) + backfill_edits(client, user, index) }) wisp.log_info( "catalog_index: backfilled from " @@ -91,3 +106,78 @@ fn backfill_catalog_index( <> " known users", ) } + +fn backfill_adoptions( + client: Client, + user: KnownUser, + index: catalog_index.Store, +) -> Nil { + list_records(client, user, shelf_entry.collection) + |> list.each(fn(rec) { + let #(uri, value) = rec + case decode.run(value, shelf_entry.shelf_entry_decoder()) { + Ok(entry) -> + case entry.release { + Some(ref) -> + index.upsert_adoption(catalog_index.Adoption( + entry_uri: uri, + did: user.did, + release_uri: ref.uri, + status: entry.action, + created_at: entry.created_at, + )) + None -> Nil + } + Error(_) -> Nil + } + }) +} + +fn backfill_edits( + client: Client, + user: KnownUser, + index: catalog_index.Store, +) -> Nil { + list_records(client, user, catalog_edit.collection) + |> list.each(fn(rec) { + let #(uri, value) = rec + case decode.run(value, catalog_edit.catalog_edit_decoder()) { + Ok(edit) -> + case edit.subject { + Some(ref) -> + index.upsert_edit(catalog_index.Edit( + edit_uri: uri, + did: user.did, + subject_uri: ref.uri, + entity: edit.entity, + created_at: edit.created_at, + )) + None -> Nil + } + Error(_) -> Nil + } + }) +} + +/// A single page of one user's records in one collection, undecoded: shared +/// by the adoption and edit backfills so there is exactly one place that +/// builds the listRecords query. Any failure yields an empty list; this is +/// a best-effort seed, not a requirement for boot to succeed. +fn list_records( + client: Client, + user: KnownUser, + collection: String, +) -> List(#(String, decode.Dynamic)) { + let params = + generated_client.RepoListRecordsParams( + collection:, + cursor: None, + limit: Some(per_user_limit), + repo: user.did, + reverse: None, + ) + case generated_client.repo_list_records(client, user.pds, params, None) { + Ok(output) -> output.records |> list.map(fn(r) { #(r.uri, r.value) }) + Error(_) -> [] + } +} diff --git a/server/src/at_record_server/config_env.gleam b/server/src/at_record_server/config_env.gleam index ed9f718..9844e0a 100644 --- a/server/src/at_record_server/config_env.gleam +++ b/server/src/at_record_server/config_env.gleam @@ -2,6 +2,8 @@ //// orchestration. Secrets (SECRET_KEY_BASE, STORE_KEY) are required; the rest //// falls back to dev defaults with a log line. +import at_record_server/catalog_index +import at_record_server/catalog_index_postgres import at_record_server/discogs_client import at_record_server/jetstream_consumer import at_record_server/known_users @@ -80,17 +82,18 @@ pub fn cover_cache_directory() -> String { envoy.get("COVER_CACHE_DIRECTORY") |> result.unwrap("./cover-cache") } -/// The session store, the Discogs credential store, and the known-users store: -/// one shared Postgres pool, no sweep ttl on the durable Discogs tokens. The -/// two credential stores are encrypted at rest; known_users is public, so it is -/// not sealed. -pub fn stores() -> #(Store, Store, known_users.Store) { +/// The session store, the Discogs credential store, the known-users store, +/// and the catalog index: one shared Postgres pool, no sweep ttl on the +/// durable Discogs tokens. The two credential stores are encrypted at rest; +/// known_users and the catalog index are public, so neither is sealed. +pub fn stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { let key = store_key() - let #(session_store, discogs_store, known) = raw_stores() + let #(session_store, discogs_store, known, index) = raw_stores() #( sealed_store.wrap(session_store, key), sealed_store.wrap(discogs_store, key), known, + index, ) } @@ -107,11 +110,11 @@ fn store_key() -> gose.Key(String) { } } -fn raw_stores() -> #(Store, Store, known_users.Store) { +fn raw_stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { case envoy.get("DATABASE_URL") { Ok(url) -> case postgres_stores(url) { - Ok(triple) -> triple + Ok(quad) -> quad Error(e) -> { wisp.log_error( "Postgres stores unavailable (" <> e <> "); using in-memory", @@ -130,7 +133,7 @@ fn raw_stores() -> #(Store, Store, known_users.Store) { fn postgres_stores( url: String, -) -> Result(#(Store, Store, known_users.Store), String) { +) -> Result(#(Store, Store, known_users.Store, catalog_index.Store), String) { use conn <- result.try(sessions_postgres.connect_pool(url)) use session_store <- result.try(sessions_postgres.table_store( conn, @@ -142,15 +145,17 @@ fn postgres_stores( "discogs_creds", option.None, )) - use known <- result.map(known_users_postgres.table_store(conn)) - #(session_store, discogs_store, known) + use known <- result.try(known_users_postgres.table_store(conn)) + use index <- result.map(catalog_index_postgres.table_store(conn)) + #(session_store, discogs_store, known, index) } -fn memory_stores() -> #(Store, Store, known_users.Store) { +fn memory_stores() -> #(Store, Store, known_users.Store, catalog_index.Store) { let assert Ok(session_store) = sessions_memory.start() let assert Ok(discogs_store) = sessions_memory.start() let assert Ok(known) = known_users_memory.start() - #(session_store, discogs_store, known) + let assert Ok(index) = catalog_index.start() + #(session_store, discogs_store, known, index) } /// App-level Discogs auth from env. Absent is fine: search still works, just