diff --git a/src/main.rs b/src/main.rs index ed3acf7..3c071a5 100644 --- a/src/main.rs +++ b/src/main.rs @@ -10,6 +10,7 @@ use tracing::{info, warn}; use lightrail::error::{Error, Result}; use lightrail::identity; use lightrail::storage; +use lightrail::sync::discovery_queue::DiscoveryQueue; use lightrail::sync::{self, backfill, firehose, resync}; use lightrail::util::TokenExt; @@ -159,6 +160,8 @@ async fn main() -> Result<()> { std::sync::Mutex::new(resync::dispatcher::DispatcherSnapshot::default()), ); + let discovery_queue = std::sync::Arc::new(DiscoveryQueue::new(4096, 36)); + let mut tasks: JoinSet> = JoinSet::new(); tasks.spawn({ @@ -189,6 +192,7 @@ async fn main() -> Result<()> { let client = client.clone(); let host = subscribe_host.clone(); let resolver = resolver.clone(); + let discovery_queue = discovery_queue.clone(); async move { if backfill::run( host, @@ -197,7 +201,7 @@ async fn main() -> Result<()> { token.clone(), resolver, false, - args.crawl_qps, + discovery_queue, ) .await .inspect(|_| info!("backfill done.")) @@ -216,6 +220,7 @@ async fn main() -> Result<()> { let client = client.clone(); let resolver = resolver.clone(); let dispatcher_state = dispatcher_state.clone(); + let discovery_queue = discovery_queue.clone(); async move { resync::dispatcher::run(resync::DispatcherConfig { resolver, @@ -227,6 +232,7 @@ async fn main() -> Result<()> { token, force_get_repo: args.heavy, state: dispatcher_state, + discovery_queue, }) .await .inspect(|_| info!("resync done.")) @@ -284,6 +290,7 @@ async fn main() -> Result<()> { let client = client.clone(); let host = subscribe_host.clone(); let resolver = resolver.clone(); + let discovery_queue = discovery_queue.clone(); async move { sync::deep_crawl::run( host, @@ -292,7 +299,7 @@ async fn main() -> Result<()> { args.max_deep_crawl_workers, token, resolver, - args.crawl_qps, + discovery_queue, ) .await .inspect(|_| info!("deep crawl done.")) diff --git a/src/sync/backfill.rs b/src/sync/backfill.rs index d11a0ea..3106862 100644 --- a/src/sync/backfill.rs +++ b/src/sync/backfill.rs @@ -6,18 +6,15 @@ //! (any state) are skipped — the dispatcher's retry mechanism handles repos //! that need re-syncing. -use std::collections::HashMap; -use std::num::NonZeroU32; use std::sync::Arc; use jacquard_api::com_atproto::sync::list_repos::ListRepos; use jacquard_common::{ error::ClientErrorKind, types::string::Did, - url::Host, + url::{Host, Url}, {IntoStatic, xrpc::XrpcExt}, }; -use reqwest::Url; use tokio::time::Duration; use tokio_util::sync::CancellationToken; use tracing::{error, info, trace, warn}; @@ -29,8 +26,8 @@ use crate::{ storage::{ self, DbRef, backfill_progress::{BackfillProgress, get, set}, - resync_queue::ResyncItem, }, + sync::discovery_queue::{DiscoveryItem, DiscoveryQueue}, util::TokenExt, }; use std::sync::atomic::Ordering; @@ -41,8 +38,6 @@ const PAGE_LIMIT: i64 = 500; const RETRY_DELAY_SECS: u64 = 10; /// Maximum consecutive transient failures before giving up on this host. const MAX_PAGE_FAILURES: u32 = 3; -/// Typical number of requests required to complete one repo resync -const REQUESTS_PER_RESYNC: u64 = 2; /// Walk `listRepos` on `host` and enqueue new repos for resync. /// @@ -57,7 +52,7 @@ pub async fn run( token: CancellationToken, resolver: Arc, validate: bool, - crawl_qps: NonZeroU32, + discovery_queue: Arc, ) -> Result { let base: jacquard_common::url::Url = format!("https://{host}") .parse() @@ -78,12 +73,6 @@ pub async fn run( ); let mut total_queued: u64 = 0; - // Per-host staggering: track last-scheduled timestamp so that items for the - // same host are spread across time at 1/crawl_qps intervals. - // This prevents the timestamp-ordered queue from bunching all items for - // a popular host together. - let host_interval = Duration::from_secs(1) / crawl_qps.get(); - let mut host_schedule: HashMap, SystemTime> = HashMap::new(); loop { if token.is_cancelled() { @@ -97,7 +86,6 @@ pub async fn run( }; let page_len = dids.len(); - let now = SystemTime::now(); // For untrusted hosts (deep crawl), filter DIDs to those whose // resolved PDS actually matches this host. @@ -130,26 +118,25 @@ pub async fn run( }; let progress_cursor = next_cursor.clone().unwrap_or_default(); - let (page_queued, schedule_back) = { + let (page_queued, new_items) = { let db = db.clone(); let host = host.clone(); - // Move the schedule into the blocking task so timestamps are - // only advanced for DIDs that are actually newly inserted. - let schedule = std::mem::take(&mut host_schedule); tokio::task::spawn_blocking(move || { - store_page( - &db, - &host, - dids_with_hosts, - progress_cursor, - now, - host_interval, - schedule, - ) + persist_page(&db, &host, dids_with_hosts, progress_cursor) }) .await?? }; - host_schedule = schedule_back; + + // Push newly-discovered repos into the in-memory discovery queue. + // Each push may block if the queue is full (backpressure). + for (did, pds) in new_items { + let Some(_) = token + .run(discovery_queue.push(DiscoveryItem { did, pds })) + .await + else { + return Ok(false); + }; + } total_queued += page_queued; @@ -181,7 +168,7 @@ pub async fn run( &host_owned, &BackfillProgress { cursor: "".to_string(), - completed_at: Some(crate::util::to_millis(now).to_string()), + completed_at: Some(crate::util::to_millis(SystemTime::now()).to_string()), }, ) }) @@ -305,43 +292,24 @@ async fn validate_dids( // Storage commit // --------------------------------------------------------------------------- -/// Enqueue newly-seen DIDs and persist the backfill cursor in one blocking task. -/// -/// Timestamps are staggered per PDS host so that the timestamp-ordered queue -/// naturally interleaves items across hosts. Only newly-inserted DIDs advance -/// the per-host schedule — already-known repos are skipped without leaving gaps. +type DidWithPds = (Did<'static>, Arc); + +/// Insert newly-seen DIDs into the repo table and persist the backfill cursor. /// -/// Returns `(count_inserted, updated_host_schedule)`. -fn store_page( +/// Returns `(count_inserted, new_items)` — the caller pushes `new_items` into +/// the in-memory discovery queue (async, with backpressure). +fn persist_page( db: &DbRef, host: &Host, - items: Vec<(Did<'static>, Arc)>, + items: Vec, progress_cursor: String, - now: SystemTime, - interval: Duration, - mut host_schedule: HashMap, SystemTime>, -) -> Result<(u64, HashMap, SystemTime>)> { - let mut count: u64 = 0; - let meta_interval = interval * REQUESTS_PER_RESYNC as u32; +) -> Result<(u64, Vec)> { + let mut new_items: Vec = Vec::new(); for (did, pds) in items { let newly_inserted = storage::repo::ensure_repo(db, &did)?; if newly_inserted { - let last = host_schedule.get(&pds).copied().unwrap_or(now); - let when = if last >= now { - last + meta_interval - } else { - now - }; - host_schedule.insert(pds, when); - let item = ResyncItem { - did, - retry_count: 0, - retry_reason: "backfill".to_string(), - commit_cbor: vec![], - }; - storage::resync_queue::enqueue(db, when, &item)?; db.stats.repos_queued_total.fetch_add(1, Ordering::Relaxed); - count += 1; + new_items.push((did, pds)); } } // Persist progress before advancing so a crash during the next @@ -354,5 +322,5 @@ fn store_page( completed_at: None, }, )?; - Ok((count, host_schedule)) + Ok((new_items.len() as u64, new_items)) } diff --git a/src/sync/deep_crawl.rs b/src/sync/deep_crawl.rs index 51eafb9..587d7bc 100644 --- a/src/sync/deep_crawl.rs +++ b/src/sync/deep_crawl.rs @@ -5,7 +5,6 @@ //! have already had a completed backfill are skipped — new repos on those PDSes //! arrive via the firehose. -use std::num::NonZeroU32; use std::sync::Arc; use jacquard_api::com_atproto::sync::list_hosts::ListHosts; @@ -19,7 +18,7 @@ use crate::{ http::ThrottledClient, identity::Resolver, storage::{DbRef, backfill_progress, list_hosts_cursor}, - sync::backfill, + sync::{backfill, discovery_queue::DiscoveryQueue}, util::TokenExt, }; @@ -40,7 +39,7 @@ pub async fn run( max_workers: usize, token: CancellationToken, resolver: Arc, - crawl_qps: NonZeroU32, + discovery_queue: Arc, ) -> Result<()> { info!(relay = %relay_host, "deep crawl started"); @@ -57,7 +56,7 @@ pub async fn run( max_workers, token.clone(), resolver.clone(), - crawl_qps, + discovery_queue.clone(), ) .await?; @@ -88,7 +87,7 @@ async fn crawl_pass( max_workers: usize, token: CancellationToken, resolver: Arc, - crawl_qps: NonZeroU32, + discovery_queue: Arc, ) -> Result<()> { let base: jacquard_common::url::Url = format!("https://{relay_host}") .parse() @@ -166,17 +165,10 @@ async fn crawl_pass( 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, - true, - crawl_qps, - ) - .await; + let outcome = + backfill::run(host2.clone(), db2, client2, child, resolver2, true, dq2).await; (host2, outcome) }); metrics::gauge!("lightrail_deep_crawl_workers").set(workers.len() as f64); diff --git a/src/sync/discovery_queue.rs b/src/sync/discovery_queue.rs new file mode 100644 index 0000000..8c3472a --- /dev/null +++ b/src/sync/discovery_queue.rs @@ -0,0 +1,412 @@ +//! Bounded, host-fair in-memory queue for discovery-sourced resync work. +//! +//! Backfill and deep-crawl tasks push newly-discovered repos here instead of +//! the on-disk resync queue. The dispatcher pops items round-robin across PDS +//! hosts, achieving fair throughput distribution without timestamp stagger +//! arithmetic. +//! +//! The queue is bounded in two dimensions: +//! - **Global capacity**: total items across all hosts; prevents unbounded +//! memory growth. +//! - **Per-host limit**: items for a single PDS; prevents one large PDS from +//! monopolising the queue and starving other hosts. +//! +//! When either limit is reached, producers block on [`DiscoveryQueue::push`] +//! until the dispatcher drains items, applying natural backpressure to crawl +//! tasks. +//! +//! On process restart, queued items are lost. This is acceptable because +//! backfill/deep-crawl cursors are persisted — the gap self-heals on the next +//! crawl pass. + +use std::collections::{HashMap, HashSet, VecDeque}; +use std::pin::pin; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; + +use jacquard_common::types::string::Did; +use jacquard_common::url::Url; + +/// A newly-discovered repo awaiting its initial resync. +pub struct DiscoveryItem { + pub did: Did<'static>, + pub pds: Arc, +} + +pub struct DiscoveryQueue { + inner: Mutex, + /// Wakes producers when space becomes available (after a pop). + space_available: tokio::sync::Notify, + /// Wakes the dispatcher when items arrive. + item_available: tokio::sync::Notify, + /// Maximum total items across all host queues. + capacity: usize, + /// Maximum items for a single PDS host. + per_host_limit: usize, + /// Total items, readable without the lock (kept in sync under the lock). + len: AtomicUsize, +} + +struct Inner { + /// Ordered list of PDS hosts for round-robin traversal. + hosts: Vec>, + /// Per-host FIFO queues of DIDs awaiting resync. + queues: HashMap, VecDeque>>, + /// Round-robin cursor: index into `hosts` for the next pop. + cursor: usize, + /// Total items across all queues. + len: usize, +} + +impl Inner { + fn can_push(&self, pds: &Arc, capacity: usize, per_host_limit: usize) -> bool { + if self.len >= capacity { + return false; + } + let host_len = self.queues.get(pds).map_or(0, |q| q.len()); + host_len < per_host_limit + } + + fn insert(&mut self, item: DiscoveryItem) { + if !self.queues.contains_key(&item.pds) { + self.hosts.push(item.pds.clone()); + self.queues.insert(item.pds.clone(), VecDeque::new()); + } + self.queues.get_mut(&item.pds).unwrap().push_back(item.did); + self.len += 1; + } +} + +impl DiscoveryQueue { + pub fn new(capacity: usize, per_host_limit: usize) -> Self { + Self { + inner: Mutex::new(Inner { + hosts: Vec::new(), + queues: HashMap::new(), + cursor: 0, + len: 0, + }), + space_available: tokio::sync::Notify::new(), + item_available: tokio::sync::Notify::new(), + capacity, + per_host_limit, + len: AtomicUsize::new(0), + } + } + + /// Push an item, blocking if the global capacity or per-host limit is + /// reached. + /// + /// Callers should `tokio::select!` this against their cancellation token + /// for clean shutdown. + pub async fn push(&self, item: DiscoveryItem) { + loop { + // Register for wake-up BEFORE checking, so a pop between check + // and await still wakes us. + let mut notified = pin!(self.space_available.notified()); + notified.as_mut().enable(); + + { + let mut inner = self.inner.lock().unwrap(); + if inner.can_push(&item.pds, self.capacity, self.per_host_limit) { + inner.insert(item); + self.len.store(inner.len, Ordering::Relaxed); + drop(inner); + self.item_available.notify_one(); + return; + } + } + + notified.await; + } + } + + /// Pop up to `n` items round-robin across hosts. + /// + /// Hosts in `skip_hosts` are left in the queue (e.g. cooling after 429). + /// Returns immediately with whatever is available (may be empty). + pub fn pop_batch(&self, n: usize, skip_hosts: &HashSet>) -> Vec { + let mut inner = self.inner.lock().unwrap(); + let mut result = Vec::with_capacity(n); + + if inner.hosts.is_empty() || n == 0 { + return result; + } + + let mut consecutive_skips = 0; + while result.len() < n && consecutive_skips <= inner.hosts.len() { + if inner.hosts.is_empty() { + break; + } + let idx = inner.cursor % inner.hosts.len(); + let host = inner.hosts[idx].clone(); + + if skip_hosts.contains(&host) { + inner.cursor = idx + 1; + consecutive_skips += 1; + continue; + } + + let queue = inner + .queues + .get_mut(&host) + .expect("host in ring must have a queue"); + let did = queue.pop_front().expect("host in ring must have items"); + let host_empty = queue.is_empty(); + inner.len -= 1; + consecutive_skips = 0; + + if host_empty { + inner.hosts.remove(idx); + inner.queues.remove(&host); + // Don't advance cursor — the next host slid into this index. + if !inner.hosts.is_empty() { + inner.cursor = idx % inner.hosts.len(); + } + } else { + inner.cursor = idx + 1; + } + + result.push(DiscoveryItem { did, pds: host }); + } + + let popped = result.len(); + self.len.store(inner.len, Ordering::Relaxed); + drop(inner); + + if popped > 0 { + self.space_available.notify_waiters(); + } + result + } + + /// Number of items across all host queues. + pub fn len(&self) -> usize { + self.len.load(Ordering::Relaxed) + } + + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + /// Returns a future that completes when an item is pushed. + /// + /// Use in `tokio::select!` alongside the idle-poll timer so the dispatcher + /// wakes immediately when discovery work arrives instead of waiting for + /// the next poll cycle. + pub fn notified(&self) -> tokio::sync::futures::Notified<'_> { + self.item_available.notified() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn url(s: &str) -> Arc { + Arc::new(s.parse().unwrap()) + } + + fn did(s: &str) -> Did<'static> { + Did::new_owned(s).unwrap() + } + + fn item(did_str: &str, pds_str: &str) -> DiscoveryItem { + DiscoveryItem { + did: did(did_str), + pds: url(pds_str), + } + } + + #[tokio::test] + async fn push_and_pop_single() { + let q = DiscoveryQueue::new(10, 10); + q.push(item("did:web:alice.test", "https://pds1.test")) + .await; + assert_eq!(q.len(), 1); + + let batch = q.pop_batch(10, &HashSet::new()); + assert_eq!(batch.len(), 1); + assert_eq!(batch[0].did.as_str(), "did:web:alice.test"); + assert_eq!(q.len(), 0); + } + + #[tokio::test] + async fn round_robin_across_hosts() { + let q = DiscoveryQueue::new(100, 100); + // Two items on pds1, two on pds2. + q.push(item("did:web:a1.test", "https://pds1.test")).await; + q.push(item("did:web:a2.test", "https://pds1.test")).await; + q.push(item("did:web:b1.test", "https://pds2.test")).await; + q.push(item("did:web:b2.test", "https://pds2.test")).await; + + let batch = q.pop_batch(4, &HashSet::new()); + assert_eq!(batch.len(), 4); + // Should alternate: pds1, pds2, pds1, pds2. + assert_eq!(batch[0].did.as_str(), "did:web:a1.test"); + assert_eq!(batch[1].did.as_str(), "did:web:b1.test"); + assert_eq!(batch[2].did.as_str(), "did:web:a2.test"); + assert_eq!(batch[3].did.as_str(), "did:web:b2.test"); + } + + #[tokio::test] + async fn skip_hosts_leaves_items_in_queue() { + let q = DiscoveryQueue::new(100, 100); + let pds1 = url("https://pds1.test"); + q.push(item("did:web:a1.test", "https://pds1.test")).await; + q.push(item("did:web:b1.test", "https://pds2.test")).await; + + let skip = HashSet::from([pds1]); + let batch = q.pop_batch(10, &skip); + assert_eq!(batch.len(), 1); + assert_eq!(batch[0].did.as_str(), "did:web:b1.test"); + // pds1 item is still queued. + assert_eq!(q.len(), 1); + } + + #[tokio::test] + async fn pop_batch_respects_limit() { + let q = DiscoveryQueue::new(100, 100); + for i in 0..10 { + q.push(item(&format!("did:web:user{i}.test"), "https://pds1.test")) + .await; + } + + let batch = q.pop_batch(3, &HashSet::new()); + assert_eq!(batch.len(), 3); + assert_eq!(q.len(), 7); + } + + #[tokio::test] + async fn empty_pop_returns_empty() { + let q = DiscoveryQueue::new(10, 10); + let batch = q.pop_batch(10, &HashSet::new()); + assert!(batch.is_empty()); + } + + /// Helper: yield to the runtime repeatedly so a spawned task has a + /// chance to make progress. If a push would complete, it will within + /// a few yields; if it's genuinely blocked, it stays blocked. + async fn yield_to_runtime() { + for _ in 0..5 { + tokio::task::yield_now().await; + } + } + + #[tokio::test] + async fn global_backpressure_blocks_producer() { + let q = Arc::new(DiscoveryQueue::new(2, 10)); + q.push(item("did:web:a.test", "https://pds1.test")).await; + q.push(item("did:web:b.test", "https://pds2.test")).await; + assert_eq!(q.len(), 2); + + // Third push should block (global capacity 2). + let q2 = q.clone(); + let handle = tokio::spawn(async move { + q2.push(item("did:web:c.test", "https://pds3.test")).await; + }); + + yield_to_runtime().await; + assert!(!handle.is_finished()); + + // Pop one to free capacity — push should complete. + q.pop_batch(1, &HashSet::new()); + yield_to_runtime().await; + assert!(handle.is_finished()); + handle.await.expect("push task shouldn't panic"); + assert_eq!(q.len(), 2); + } + + #[tokio::test] + async fn per_host_limit_blocks_producer() { + let q = Arc::new(DiscoveryQueue::new(100, 2)); + // Fill host pds1 to its per-host limit of 2. + q.push(item("did:web:a.test", "https://pds1.test")).await; + q.push(item("did:web:b.test", "https://pds1.test")).await; + + // Third push to same host should block despite global capacity. + let q2 = q.clone(); + let handle = tokio::spawn(async move { + q2.push(item("did:web:c.test", "https://pds1.test")).await; + }); + + yield_to_runtime().await; + assert!(!handle.is_finished()); + + // A different host should still be pushable. + q.push(item("did:web:d.test", "https://pds2.test")).await; + assert_eq!(q.len(), 3); + + // Pop one from pds1 to unblock the waiting producer. + let batch = q.pop_batch(1, &HashSet::new()); + assert_eq!(batch[0].pds.as_str(), "https://pds1.test/"); + yield_to_runtime().await; + assert!(handle.is_finished()); + handle.await.expect("push task shouldn't panic"); + assert_eq!(q.len(), 3); + } + + #[tokio::test] + async fn per_host_limit_prevents_monopolisation() { + let q = DiscoveryQueue::new(100, 3); + // Push 3 items for pds1 (at limit), then 3 for pds2. + for i in 0..3 { + q.push(item(&format!("did:web:a{i}.test"), "https://pds1.test")) + .await; + } + for i in 0..3 { + q.push(item(&format!("did:web:b{i}.test"), "https://pds2.test")) + .await; + } + assert_eq!(q.len(), 6); + + // Pop 6: should interleave thanks to round-robin. + let batch = q.pop_batch(6, &HashSet::new()); + assert_eq!(batch.len(), 6); + assert_eq!(batch[0].pds.as_str(), "https://pds1.test/"); + assert_eq!(batch[1].pds.as_str(), "https://pds2.test/"); + assert_eq!(batch[2].pds.as_str(), "https://pds1.test/"); + assert_eq!(batch[3].pds.as_str(), "https://pds2.test/"); + } + + #[tokio::test] + async fn cursor_persists_across_batches() { + let q = DiscoveryQueue::new(100, 100); + // 3 hosts, 1 item each. + q.push(item("did:web:a.test", "https://pds1.test")).await; + q.push(item("did:web:b.test", "https://pds2.test")).await; + q.push(item("did:web:c.test", "https://pds3.test")).await; + // Add second item to each. + q.push(item("did:web:a2.test", "https://pds1.test")).await; + q.push(item("did:web:b2.test", "https://pds2.test")).await; + q.push(item("did:web:c2.test", "https://pds3.test")).await; + + // Pop 2: should get pds1, pds2. Cursor now at pds3. + let batch1 = q.pop_batch(2, &HashSet::new()); + assert_eq!(batch1[0].did.as_str(), "did:web:a.test"); + assert_eq!(batch1[1].did.as_str(), "did:web:b.test"); + + // Pop 2 more: cursor resumes at pds3, then wraps to pds1. + let batch2 = q.pop_batch(2, &HashSet::new()); + assert_eq!(batch2[0].did.as_str(), "did:web:c.test"); + assert_eq!(batch2[1].did.as_str(), "did:web:a2.test"); + } + + #[tokio::test] + async fn host_removed_when_drained() { + let q = DiscoveryQueue::new(100, 100); + q.push(item("did:web:a.test", "https://pds1.test")).await; + q.push(item("did:web:b.test", "https://pds2.test")).await; + + // Drain pds1 by skipping pds2. + let pds2 = url("https://pds2.test"); + let batch = q.pop_batch(10, &HashSet::from([pds2])); + assert_eq!(batch.len(), 1); + + // Now un-skip pds2 and pop. + let batch = q.pop_batch(10, &HashSet::new()); + assert_eq!(batch.len(), 1); + assert_eq!(batch[0].did.as_str(), "did:web:b.test"); + assert!(q.is_empty()); + } +} diff --git a/src/sync/mod.rs b/src/sync/mod.rs index 5b388da..ca7003c 100644 --- a/src/sync/mod.rs +++ b/src/sync/mod.rs @@ -1,4 +1,5 @@ pub mod backfill; pub mod deep_crawl; +pub mod discovery_queue; pub mod firehose; pub mod resync; diff --git a/src/sync/resync/dispatcher.rs b/src/sync/resync/dispatcher.rs index e355e72..0784cc9 100644 --- a/src/sync/resync/dispatcher.rs +++ b/src/sync/resync/dispatcher.rs @@ -29,6 +29,7 @@ use crate::storage::{ repo::{AccountStatus, RepoInfo, RepoState}, resync_queue::ResyncItem, }; +use crate::sync::discovery_queue::DiscoveryQueue; use crate::util::TokenExt; /// Point-in-time view of the dispatcher's operational state, published for @@ -43,6 +44,8 @@ pub struct DispatcherSnapshot { pub hosts: Vec, /// Total number of worker tasks in the JoinSet. pub worker_count: usize, + /// Number of items in the in-memory discovery queue. + pub discovery_queue_depth: usize, } #[derive(Clone, Serialize)] @@ -85,6 +88,7 @@ pub struct DispatcherConfig { pub token: tokio_util::sync::CancellationToken, pub force_get_repo: bool, pub state: DispatcherState, + pub discovery_queue: Arc, } /// Run the resync dispatcher until the future is cancelled. @@ -103,6 +107,7 @@ pub async fn run( token, force_get_repo, state, + discovery_queue, }: DispatcherConfig, ) -> Result<()> { let mut busy: HashSet> = HashSet::new(); @@ -210,11 +215,68 @@ pub async fn run( } } + // Disk queue exhausted — fill remaining slots from the memory + // discovery queue (backfill / deep-crawl items). + if workers.len() < max_concurrent { + let now_inst = Instant::now(); + let cooling_set: HashSet> = cooling_hosts + .iter() + .filter(|&(_, &until)| now_inst < until) + .map(|(url, _)| Arc::clone(url)) + .collect(); + let batch = discovery_queue.pop_batch(max_concurrent - workers.len(), &cooling_set); + for dq_item in batch { + if busy.contains(&dq_item.did) { + continue; // already being resynced via disk queue + } + let did = dq_item.did.clone(); + // Transition to Resyncing so the firehose buffers commits. + if let Err(e) = + transition_state(db.clone(), did.clone(), RepoState::Resyncing, None).await + { + warn!(did = %did, error = %e, + "failed to transition discovery item to Resyncing; skipping"); + continue; + } + + busy.insert(did.clone()); + let item = ResyncItem { + did: dq_item.did, + retry_count: 0, + retry_reason: String::new(), + commit_cbor: vec![], + }; + let worker_stuff = AppStuff { + resolver: resolver.clone(), + client: client.clone(), + db: db.clone(), + token: token.clone(), + }; + let handle = workers.spawn(async move { + run_worker( + item, + worker_stuff, + describe_timeout, + get_repo_timeout, + force_get_repo, + ) + .await + }); + task_dids.insert(handle.id(), did.clone()); + metrics::gauge!("lightrail_resync_workers").set(workers.len() as f64); + trace!(did = %did, running = workers.len(), "spawned resync worker (discovery)"); + } + } + if workers.is_empty() { - // Nothing running, nothing claimable right now. Yield before polling. - if !token.sleep(IDLE_POLL).await { - return Ok(()); - }; + // Nothing running, nothing claimable right now. Wait for new + // items in the discovery queue or the next poll interval. + let notified = discovery_queue.notified(); + tokio::select! { + _ = token.cancelled() => return Ok(()), + _ = tokio::time::sleep(IDLE_POLL) => {}, + _ = notified => {}, + } continue; } @@ -339,6 +401,7 @@ pub async fn run( entries }, worker_count: workers.len(), + discovery_queue_depth: discovery_queue.len(), }; // Unwrap is safe: we never poison the lock (no panics while holding it). *state.lock().unwrap() = snap;