//! Walk `com.atproto.sync.listRepos` and feed newly discovered repos into the //! resync queue for the dispatcher to process. //! //! Each page of results is processed and the cursor persisted before moving //! on, so the walk can be safely resumed after a restart. Already-known repos //! (any state) are skipped — the dispatcher's retry mechanism handles repos //! that need re-syncing. //! //! Accounts listed as non-active (takendown, suspended, deactivated, deleted) //! have their status recorded without queuing a resync — there's nothing to //! fetch. use std::sync::Arc; use jacquard_api::com_atproto::sync::list_repos::{ListRepos, RepoStatus}; use jacquard_common::{ error::ClientErrorKind, types::{crypto::PublicKey, string::Did}, url::{Host, Url}, {IntoStatic, xrpc::XrpcExt}, }; use tokio::time::Duration; use tokio_util::sync::CancellationToken; use tracing::{debug, error, info, trace, warn}; use crate::{ error::Result, http::ThrottledClient, identity::Resolver, storage::{ self, DbRef, backfill_progress::{BackfillProgress, get, set}, repo::{AccountStatus, RepoInfo, RepoState}, }, sync::discovery_queue::{DiscoveryItem, DiscoveryQueue}, util::TokenExt, }; use std::sync::atomic::Ordering; use std::time::SystemTime; const PAGE_LIMIT: i64 = 500; // --------------------------------------------------------------------------- // Mode + retry policy // --------------------------------------------------------------------------- /// Which host role this backfill is running against. /// /// Bundles two decisions that always move together: /// - whether to validate DIDs against the host (trust model), and /// - how patiently to retry transient page failures (criticality). #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub enum BackfillMode { /// Primary relay / upstream host. Trusted for DID listings; retried /// generously because a transient outage here will take the whole service /// down when the task exits. Relay, /// Untrusted PDS discovered via deep crawl. DIDs are validated against the /// host, and a few transient failures are enough to move on. DeepCrawl, } impl BackfillMode { /// True when listed DIDs must be independently verified to actually live /// on `host` before we trust their account status. fn validates_dids(self) -> bool { matches!(self, Self::DeepCrawl) } fn retry_policy(self) -> RetryPolicy { match self { Self::Relay => RetryPolicy::RELAY, Self::DeepCrawl => RetryPolicy::DEEP_CRAWL, } } } /// How many transient page failures to tolerate and how long to wait between. struct RetryPolicy { max_page_failures: u32, retry_delay: Duration, } impl RetryPolicy { /// Relay budget: 60 attempts × 30s = 30 minutes of retries before giving /// up and triggering service shutdown. Enough to absorb typical relay /// hiccups without flapping. const RELAY: Self = Self { max_page_failures: 60, retry_delay: Duration::from_secs(30), }; /// Deep-crawl budget: 3 × 10s. Fast give-up so dead PDSes don't monopolise /// crawl workers. const DEEP_CRAWL: Self = Self { max_page_failures: 3, retry_delay: Duration::from_secs(10), }; } // --------------------------------------------------------------------------- // Listed account state // --------------------------------------------------------------------------- /// Map the `active`/`status` fields from a listRepos `Repo` entry to a /// [`ListedAccountState`]. Follows the same mapping as firehose `#account` /// events in `account_event.rs`. fn classify_repo(active: Option, status: &Option>) -> RepoState { // active-as-fallback is safe: if active.unwrap_or(true) { return RepoState::Active; } let Some(jac_status) = status else { return RepoState::Active; }; match jac_status { RepoStatus::Takendown => RepoState::Takendown, RepoStatus::Suspended => RepoState::Suspended, RepoStatus::Deleted => RepoState::Deleted, RepoStatus::Deactivated => RepoState::Deactivated, RepoStatus::Desynchronized => RepoState::Desynchronized, RepoStatus::Throttled => RepoState::Throttled, RepoStatus::Other(_) => RepoState::Error, } } // --------------------------------------------------------------------------- // Main loop // --------------------------------------------------------------------------- /// Walk `listRepos` on `host` and enqueue new repos for resync. /// /// `mode` selects both the trust model (whether DIDs are validated against /// `host`) and the retry policy (how long we retry transient page failures /// before giving up). See [`BackfillMode`]. /// /// In `DeepCrawl` mode, DIDs are verified to actually live on `host` and /// non-matching ones are rejected. Validated DIDs are considered authoritative /// for account status — if the PDS says an account is deactivated, we trust /// it even if our local state is Active. In `Relay` mode we only write /// non-active status for accounts we haven't seen before, to avoid /// overwriting Active state with stale relay data (e.g. accounts that /// migrated off and appear deactivated on the old PDS). pub async fn run( host: Host, db: DbRef, client: ThrottledClient, token: CancellationToken, resolver: Arc, mode: BackfillMode, discovery_queue: Arc, ) -> Result { let validate = mode.validates_dids(); let retry_policy = mode.retry_policy(); let base: jacquard_common::url::Url = format!("https://{host}") .parse() .map_err(|e: jacquard_common::url::ParseError| crate::error::Error::Other(e.to_string()))?; // Resume from the last saved cursor (empty string → start from beginning). let mut cursor: Option = { let db = db.clone(); let host = host.clone(); tokio::task::spawn_blocking(move || get(&db, &host)).await?? } .map(|p| p.cursor); info!( host = %host, resume_cursor = cursor.as_deref().unwrap_or("(start)"), "backfill started" ); let mut total_queued: u64 = 0; loop { if token.is_cancelled() { return Ok(false); } let (dids, next_cursor) = match fetch_page( &host, &base, cursor.as_deref(), &client, &token, &retry_policy, ) .await { Some(page) => page, None => return Ok(false), }; // An absent or empty cursor both signal end-of-list. let next_cursor = next_cursor .filter(|next| !next.is_empty()) // treat an empty-string cursor as done .filter(|next| { let Some(current) = cursor else { return true; // none cursor on first request, allow next page }; if *next == current { warn!(host = %host, ?current, ?next, "cursor unchanged, possible crawl-trap, not requesting 'next' listRepos page."); // TODO: mark host as not trustworthy return false; } true }); let page_len = dids.len(); // For untrusted hosts (deep crawl), filter DIDs to those whose // resolved PDS actually matches this host. let dids = if validate { validate_dids(dids, &resolver, &host, &token).await } else { dids }; // Resolve each DID's actual PDS host for discovery queue routing. // Cache hits are free; misses fall back to the listed host. // // TODO: ...this is basically redundant with validate_dids now? let dids_with_hosts: Vec<(Did<'static>, Arc, PublicKey<'static>, RepoState)> = { let mut out = Vec::with_capacity(dids.len()); for (did, account_state) in dids { let Some(res) = token.run(resolver.resolve(&did)).await else { return Ok(false); // cancelled }; let resolved = match res { Ok(resolved) => resolved, Err(e) => { error!(did = %did, error = %e, "failed to resolve host for validated did; skipping"); continue; } }; out.push(( did, resolved.pds.clone(), resolved.pubkey.clone(), account_state, )); } out }; let progress_cursor = next_cursor.clone().unwrap_or_default(); let (page_queued, page_inactive, new_items) = { let db = db.clone(); let host = host.clone(); tokio::task::spawn_blocking(move || { persist_page(&db, &host, dids_with_hosts, progress_cursor, validate) }) .await?? }; // Push newly-discovered repos into the in-memory discovery queue. // Each push may block if the queue is full (backpressure). for (did, pds, pubkey) in new_items { let Some(_) = token .run(discovery_queue.push(DiscoveryItem { did, pds, pubkey })) .await else { return Ok(false); }; } total_queued += page_queued; trace!( host = %host, page_repos = page_len, page_queued, page_inactive, total_queued, next_cursor = next_cursor.as_deref().unwrap_or("(done)"), "backfill page" ); // next page: only if we have a next-cursor (omgggggggg sorry everyone) if let Some(next) = next_cursor { cursor = Some(next.to_string()); continue; } // Reached end-of-list (or bailed on a trapping host): persist completion // so re-crawls skip this host. let db = db.clone(); let host_owned = host.clone(); tokio::task::spawn_blocking(move || { set( &db, &host_owned, &BackfillProgress { cursor: "".to_string(), completed_at: Some(crate::util::to_millis(SystemTime::now()).to_string()), }, ) }) .await??; info!(host = %host, total_queued, "backfill complete"); return Ok(true); } } // --------------------------------------------------------------------------- // Page fetch // --------------------------------------------------------------------------- /// Fetch one `listRepos` page with retry logic. /// /// Returns `None` if the host gives a 4xx, exceeds `policy.max_page_failures` /// transient errors, or the token is cancelled (including mid-request). async fn fetch_page( host: &Host, base: &jacquard_common::url::Url, cursor: Option<&str>, client: &ThrottledClient, token: &CancellationToken, policy: &RetryPolicy, ) -> Option<(Vec<(Did<'static>, RepoState)>, Option)> { let req = ListRepos { cursor: cursor.map(Into::into), limit: Some(PAGE_LIMIT), }; let mut failures: u32 = 0; loop { let result = match token.run(client.xrpc(base.clone()).send(&req)).await? { Err(e) => { let is_client_err = matches!( e.kind(), ClientErrorKind::Http { status } if status.is_client_error() ); if is_client_err { warn!(error = %e, host = %host, "listRepos failed with client error; giving up on this host"); return None; } warn!(error = %e, host = %host, "listRepos request failed"); None } Ok(resp) => match resp.parse() { Ok(out) => { let next = out.cursor.as_deref().map(str::to_owned); let listed = out .repos .into_iter() .map(|r| { let state = classify_repo(r.active, &r.status); (r.did.into_static(), state) }) .collect::>(); Some((listed, next)) } Err(e) => { warn!(error = %e, host = %host, "listRepos response parse failed"); None } }, }; match result { Some(page) => return Some(page), None => { failures += 1; if failures >= policy.max_page_failures { warn!(host = %host, failures, max = policy.max_page_failures, "listRepos page failed too many times; giving up on this host"); return None; } if !token.sleep(policy.retry_delay).await { return None; } } } } } // --------------------------------------------------------------------------- // DID validation // --------------------------------------------------------------------------- /// Filter `dids` to those whose resolved PDS endpoint matches `host`. /// /// Returns early with whatever has been validated so far if the token is cancelled. async fn validate_dids( dids: Vec<(Did<'static>, RepoState)>, resolver: &Resolver, host: &Host, token: &CancellationToken, ) -> Vec<(Did<'static>, RepoState)> { let host_str = host.to_string(); let mut valid = Vec::with_capacity(dids.len()); for (did, account_state) in dids { let Some(r) = token.run(resolver.resolve(&did)).await else { break; }; match r { Ok(resolved) if resolved.pds.host_str() == Some(host_str.as_str()) => { valid.push((did, account_state)); } Ok(resolved) => { metrics::counter!("lightrail_backfill_did_rejected_total", "reason" => "pds_mismatch") .increment(1); trace!(did = %did, resolved_pds = %resolved.pds, expected = %host, "DID resolves to different PDS; skipping"); } Err(e) => { metrics::counter!("lightrail_backfill_did_rejected_total", "reason" => "resolution_failed") .increment(1); trace!(did = %did, error = %e, "DID resolution failed during backfill validation; skipping"); } }; } valid } // --------------------------------------------------------------------------- // Storage commit // --------------------------------------------------------------------------- type DidWithPds = (Did<'static>, Arc, PublicKey<'static>); /// Insert newly-seen DIDs into the repo table and persist the backfill cursor. /// /// Active accounts are returned for the caller to push into the discovery /// queue. Non-active accounts have their status written directly (no resync /// needed). When `authoritative` is false (relay backfill), non-active status /// is only written for accounts that have no existing Active record — we don't /// overwrite a locally-Active account with stale relay data. /// /// Returns `(active_count, inactive_count, new_active_items)`. fn persist_page( db: &DbRef, host: &Host, items: Vec<(Did<'static>, Arc, PublicKey<'static>, RepoState)>, progress_cursor: String, authoritative: bool, ) -> Result<(u64, u64, Vec)> { let mut new_active: Vec = Vec::new(); let mut inactive_count: u64 = 0; for (did, pds, pubkey, repo_state) in items { if let Some(deactiveated_account_status) = repo_state.to_account_inactive() { if write_inactive_status( db, &did, deactiveated_account_status, repo_state, authoritative, )? { inactive_count += 1; } } else { let newly_inserted = storage::repo::ensure_repo(db, &did)?; if newly_inserted { db.stats.repos_queued_total.fetch_add(1, Ordering::Relaxed); new_active.push((did, pds, pubkey)); } } } // Persist progress before advancing so a crash during the next // page restarts from here, not the beginning. set( db, host, &BackfillProgress { cursor: progress_cursor, completed_at: None, }, )?; Ok((new_active.len() as u64, inactive_count, new_active)) } /// Write a non-active account status from a listRepos entry. /// /// Returns `true` if the status was written, `false` if skipped. /// /// When `authoritative` is true (deep crawl — PDS confirmed via identity /// resolution), the status is always written. When false (relay backfill), /// we only write if the existing local status is not Active, to avoid /// overwriting with stale data (e.g. an account that migrated off this PDS /// and shows as deactivated on the old host). fn write_inactive_status( db: &DbRef, did: &Did<'_>, status: AccountStatus, state: RepoState, authoritative: bool, ) -> Result { let existing = storage::repo::get_status(db, did)?; match existing { Some(ref existing_status) if existing_status.is_active() && !authoritative => { debug!( did = %did, listed_status = status.as_str(), "skipping non-active status from non-authoritative source; local account is active" ); return Ok(false); } Some(ref existing_status) if *existing_status == status => { return Ok(false); // already matches } _ => {} } let info = RepoInfo { state, status: status.clone(), error: None, }; storage::repo::put_info(db, did, &info)?; metrics::counter!( "lightrail_backfill_inactive_total", "status" => status.as_str() ) .increment(1); Ok(true) }