//! Deep-crawl orchestration: discover PDS hosts via `listHosts`, then crawl //! each one's repos via the existing `backfill::run()` worker. //! //! Runs alongside (not instead of) the normal relay-level backfill. Hosts that //! have already had a completed backfill are skipped — new repos on those PDSes //! arrive via the firehose. use std::sync::Arc; use jacquard_api::com_atproto::sync::list_hosts::ListHosts; use jacquard_common::{url::Host, xrpc::XrpcExt}; use tokio::{task::JoinSet, time::Duration}; use tokio_util::sync::CancellationToken; use tracing::{info, warn}; use crate::{ error::Result, http::ThrottledClient, identity::Resolver, storage::{DbRef, backfill_progress, list_hosts_cursor}, sync::{backfill, discovery_queue::DiscoveryQueue}, util::TokenExt, }; const PAGE_LIMIT: i64 = 500; /// Delay between retry attempts after a failed page request. const RETRY_DELAY_MS: u64 = 10_000; /// Maximum number of consecutive page failures before abandoning the pass. const MAX_PAGE_FAILURES: u32 = 3; /// How long to wait between full passes through all listHosts pages. const REPOLL_MS: u64 = 20 * 60 * 60_000; /// Discover PDS hosts via `listHosts` on `relay_host` and crawl each one, /// re-polling every 20 hours after each full pass. pub async fn run( relay_host: Host, db: DbRef, client: ThrottledClient, max_workers: usize, token: CancellationToken, resolver: Arc, discovery_queue: Arc, ) -> Result<()> { info!(relay = %relay_host, "deep crawl started"); loop { if token.is_cancelled() { return Ok(()); } info!(relay = %relay_host, "deep crawl pass started"); crawl_pass( relay_host.clone(), db.clone(), client.clone(), max_workers, token.clone(), resolver.clone(), discovery_queue.clone(), ) .await?; if token.is_cancelled() { return Ok(()); } info!( relay = %relay_host, repoll_ms = REPOLL_MS, "deep crawl pass complete; re-polling in 20 hours" ); if !token.sleep(Duration::from_millis(REPOLL_MS)).await { return Ok(()); } } } /// Page through the full `listHosts` feed once, spawning a `backfill::run` /// worker for each PDS host not yet crawled. Workers run concurrently up to /// `max_workers`; the cursor is persisted after each page so a restart resumes /// mid-pass rather than from the beginning. async fn crawl_pass( relay_host: Host, db: DbRef, client: ThrottledClient, max_workers: usize, token: CancellationToken, resolver: Arc, discovery_queue: Arc, ) -> Result<()> { let base: jacquard_common::url::Url = format!("https://{relay_host}") .parse() .map_err(|e: jacquard_common::url::ParseError| crate::error::Error::Other(e.to_string()))?; // Resume from the last saved cursor. let mut cursor: Option = { let db = db.clone(); let host = relay_host.clone(); tokio::task::spawn_blocking(move || list_hosts_cursor::get(&db, &host)).await?? }; let mut workers: JoinSet<(Host, Result)> = JoinSet::new(); loop { if token.is_cancelled() { drain_workers(&mut workers).await; return Ok(()); } let (hosts, next_cursor) = match fetch_hosts_page(&relay_host, &base, cursor.as_deref(), &client, &token).await { Some(page) => page, None => { drain_workers(&mut workers).await; return Ok(()); } }; // Spawn a per-PDS backfill worker for each new host on this page. for hostname in hosts { if token.is_cancelled() { drain_workers(&mut workers).await; return Ok(()); } let pds_host = match Host::parse(&hostname) { Ok(h) => h, Err(e) => { warn!(hostname = %hostname, error = %e, "failed to parse PDS hostname; skipping"); metrics::counter!("lightrail_deep_crawl_pds_discovered_total", "outcome" => "parse_error") .increment(1); continue; } }; // Skip hosts that have already been fully crawled. let already_done = { let db = db.clone(); let h = pds_host.clone(); tokio::task::spawn_blocking(move || backfill_progress::get(&db, &h)).await?? }; if already_done.and_then(|p| p.completed_at).is_some() { metrics::counter!("lightrail_deep_crawl_pds_discovered_total", "outcome" => "skipped") .increment(1); continue; } // If at capacity, wait for one worker to finish before spawning. if workers.len() >= max_workers { tokio::select! { _ = token.cancelled() => { drain_workers(&mut workers).await; return Ok(()); } Some(result) = workers.join_next() => { log_worker_result(result); metrics::gauge!("lightrail_deep_crawl_workers").set(workers.len() as f64); } } } let db2 = db.clone(); let client2 = client.clone(); let child = token.child_token(); let host2 = pds_host.clone(); let resolver2 = resolver.clone(); let dq2 = discovery_queue.clone(); workers.spawn(async move { let outcome = backfill::run( host2.clone(), db2, client2, child, resolver2, backfill::BackfillMode::DeepCrawl, dq2, ) .await; (host2, outcome) }); metrics::gauge!("lightrail_deep_crawl_workers").set(workers.len() as f64); metrics::counter!("lightrail_deep_crawl_pds_discovered_total", "outcome" => "spawned") .increment(1); info!(pds = %pds_host, "spawned deep crawl backfill worker"); } // Persist the cursor after processing this page. let cursor_to_save = next_cursor.clone().unwrap_or_default(); { let db = db.clone(); let host = relay_host.clone(); tokio::task::spawn_blocking(move || { list_hosts_cursor::set(&db, &host, &cursor_to_save) }) .await??; } match next_cursor { Some(c) => cursor = Some(c), None => { // End of listHosts — drain remaining workers then return. drain_workers(&mut workers).await; return Ok(()); } } } } // --------------------------------------------------------------------------- // Page fetch // --------------------------------------------------------------------------- /// Fetch one `listHosts` page with retry logic. /// /// Returns `None` if the host exceeds `MAX_PAGE_FAILURES` transient errors or /// the token is cancelled (including mid-request). async fn fetch_hosts_page( relay_host: &Host, base: &jacquard_common::url::Url, cursor: Option<&str>, client: &ThrottledClient, token: &CancellationToken, ) -> Option<(Vec, Option)> { let req = ListHosts { cursor: cursor.map(Into::into), limit: Some(PAGE_LIMIT), }; for attempt in 0..MAX_PAGE_FAILURES { if attempt > 0 && !token.sleep(Duration::from_millis(RETRY_DELAY_MS)).await { return None; } let Ok(res) = token .run(client.xrpc(base.clone()).send(&req)) .await? .inspect_err(|e| info!(error = %e, relay = %relay_host, "listHosts request failed")) else { continue; }; let Ok(out) = res.parse().inspect_err( |e| info!(error = %e, relay = %relay_host, "listHosts response parse failed"), ) else { continue; }; let next = out.cursor.as_deref().map(str::to_owned); let hostnames = out .hosts .into_iter() .map(|h| h.hostname.to_string()) .collect::>(); return Some((hostnames, next)); } warn!( relay = %relay_host, "listHosts page failed {MAX_PAGE_FAILURES} times; abandoning pass" ); None } // --------------------------------------------------------------------------- // Worker pool helpers // --------------------------------------------------------------------------- async fn drain_workers(workers: &mut JoinSet<(Host, Result)>) { while let Some(result) = workers.join_next().await { log_worker_result(result); } metrics::gauge!("lightrail_deep_crawl_workers").set(0.0); } fn log_worker_result(result: std::result::Result<(Host, Result), tokio::task::JoinError>) { match result { Ok((host, Ok(true))) => { metrics::counter!("lightrail_deep_crawl_pds_completed_total", "outcome" => "completed") .increment(1); info!(pds = %host, "deep crawl backfill worker completed"); } Ok((host, Ok(false))) => { metrics::counter!("lightrail_deep_crawl_pds_completed_total", "outcome" => "cancelled") .increment(1); info!(pds = %host, "deep crawl backfill worker cancelled"); } Ok((host, Err(e))) => { metrics::counter!("lightrail_deep_crawl_pds_completed_total", "outcome" => "failed") .increment(1); warn!(pds = %host, error = %e, "deep crawl backfill worker failed"); } Err(e) => { metrics::counter!("lightrail_deep_crawl_pds_completed_total", "outcome" => "panicked") .increment(1); warn!(error = %e, "deep crawl backfill worker panicked"); } } }