diff --git a/docs/configuration.md b/docs/configuration.md index fd77eb0..027ab5e 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -48,8 +48,8 @@ hydrant is configured via environment variables, all prefixed with `HYDRANT_` (e | :--- | :--- | :--- | | `CRAWLER_URLS` | relay hosts (full network), `https://lightrail.microcosm.blue` (filter) | comma-separated list of `[mode::]url` crawler sources. mode is `relay` or `by_collection`; bare URLs use the default mode. set to empty string to disable crawling | | `ENABLE_CRAWLER` | `true` if full network or sources configured | whether to actively query the network for unknown repositories | -| `CRAWLER_MAX_PENDING_REPOS` | `2000` | max pending repos before the crawler pauses | -| `CRAWLER_RESUME_PENDING_REPOS` | `1000` | pending-repo count at which the crawler resumes | +| `CRAWLER_MAX_PENDING_REPOS` | `2000` (`25000` full network) | persistent discovery lookahead retained in the pending queue before the crawler pauses; this does not increase active backfill task concurrency | +| `CRAWLER_RESUME_PENDING_REPOS` | `1000` (`12500` full network) | pending-repo count at which discovery resumes | ## backfill & identity diff --git a/src/config.rs b/src/config.rs index 54f755b..e14a32b 100644 --- a/src/config.rs +++ b/src/config.rs @@ -137,7 +137,8 @@ pub struct Config { /// whether to run the network crawler. `None` defers to the default for the current mode. /// set via `HYDRANT_ENABLE_CRAWLER`. pub enable_crawler: Option, - /// maximum number of repos allowed in the backfill pending queue before the crawler pauses. + /// maximum number of discovered repos retained as persistent backfill lookahead before the + /// crawler pauses. backfill task concurrency remains bounded separately. /// set via `HYDRANT_CRAWLER_MAX_PENDING_REPOS`. pub crawler_max_pending_repos: usize, /// pending queue size at which the crawler resumes after being paused. @@ -390,6 +391,8 @@ impl Config { plc_urls: vec![Url::parse("https://plc.directory").unwrap()], firehose_workers: 24, backfill_concurrency_limit: 64, + crawler_max_pending_repos: 25_000, + crawler_resume_pending_repos: 12_500, crawler_sources: vec![CrawlerSource { url: Url::parse("wss://relay.fire.hose.cam/").unwrap(), mode: CrawlerMode::ListRepos, @@ -405,6 +408,19 @@ impl Config { } } +#[cfg(test)] +mod tests { + use super::Config; + + #[test] + fn full_network_keeps_a_large_bounded_discovery_window() { + let config = Config::full_network(); + assert_eq!(config.crawler_max_pending_repos, 25_000); + assert_eq!(config.crawler_resume_pending_repos, 12_500); + assert!(config.crawler_resume_pending_repos < config.crawler_max_pending_repos); + } +} + macro_rules! config_line { ($f:expr, $label:expr, $value:expr) => { writeln!($f, " {: CapacityDecision { + if pending > max_pending as u64 { + if was_throttled { + CapacityDecision::StayThrottled + } else { + CapacityDecision::EnterThrottle + } + } else if was_throttled && pending > resume_pending as u64 { + CapacityDecision::StayThrottled + } else if was_throttled { + CapacityDecision::ReleaseThrottle + } else { + CapacityDecision::Proceed + } +} + /// a cursor write to include atomically with the batch. pub(crate) struct CursorUpdate { pub(super) key: Vec, @@ -86,13 +115,19 @@ impl CrawlerWorker { /// blocks until the pending queue has capacity. mirrors the hysteresis logic /// that was previously duplicated in each crawler loop: /// - above `max_pending`: hard stop, poll every 5s - /// - between `resume_pending` and `max_pending`: cooldown until below `resume_pending` + /// - after crossing `max_pending`: cooldown until below `resume_pending` /// - below `resume_pending`: proceed (or release throttle) async fn wait_for_capacity(&mut self) { loop { let pending = self.state.db.get_count("pending").await; - if pending > self.max_pending as u64 { - if !self.was_throttled { + match capacity_decision( + pending, + self.max_pending, + self.resume_pending, + self.was_throttled, + ) { + CapacityDecision::Proceed => break, + CapacityDecision::EnterThrottle => { debug!( pending, max = self.max_pending, @@ -100,35 +135,17 @@ impl CrawlerWorker { ); self.was_throttled = true; self.stats.set_throttled(true); + tokio::time::sleep(Duration::from_secs(5)).await; } - tokio::time::sleep(Duration::from_secs(5)).await; - } else if pending > self.resume_pending as u64 { - if !self.was_throttled { - debug!( - pending, - resume = self.resume_pending, - "throttling: entering cooldown" - ); - self.was_throttled = true; - self.stats.set_throttled(true); - } - loop { + CapacityDecision::StayThrottled => { tokio::time::sleep(Duration::from_secs(5)).await; - if self.state.db.get_count("pending").await <= self.resume_pending as u64 { - break; - } } - self.was_throttled = false; - self.stats.set_throttled(false); - info!("throttling released"); - break; - } else { - if self.was_throttled { + CapacityDecision::ReleaseThrottle => { self.was_throttled = false; self.stats.set_throttled(false); info!("throttling released"); + break; } - break; } } } @@ -441,6 +458,46 @@ mod tests { use bytes::Bytes; use miette::IntoDiagnostic; + #[test] + fn pending_between_resume_and_max_does_not_start_throttling() { + assert_eq!( + capacity_decision(75, 100, 50, false), + CapacityDecision::Proceed + ); + assert_eq!( + capacity_decision(100, 100, 50, false), + CapacityDecision::Proceed + ); + } + + #[test] + fn pending_above_max_starts_throttling() { + assert_eq!( + capacity_decision(101, 100, 50, false), + CapacityDecision::EnterThrottle + ); + } + + #[test] + fn throttling_stays_active_until_resume_boundary() { + assert_eq!( + capacity_decision(100, 100, 50, true), + CapacityDecision::StayThrottled + ); + assert_eq!( + capacity_decision(51, 100, 50, true), + CapacityDecision::StayThrottled + ); + assert_eq!( + capacity_decision(50, 100, 50, true), + CapacityDecision::ReleaseThrottle + ); + assert_eq!( + capacity_decision(0, 100, 50, true), + CapacityDecision::ReleaseThrottle + ); + } + fn did() -> Did<'static> { Did::new_static("did:plc:aaaaaaaaaaaaaaaaaaaaaaaa").unwrap() }