From 434e43fc6240be24957bf876bedf6db10c919ed5 Mon Sep 17 00:00:00 2001 From: DanConwayDev Date: Mon, 10 Aug 2026 15:20:38 +0000 Subject: [PATCH] fix(purgatory): pause background Git fan-out under pressure A cold archive sync under a two-CPU quota started event-directed Git work against many independent domains. The existing five-per-domain throttle remained individually polite but admitted roughly 149 service tasks collectively. Same-process NIP-11 latency rose from 0.3 ms idle to 0.8-5.0+ seconds, metrics scrapes timed out, and WebSocket handshakes intermittently failed. Add one resource-pressure admission gate around complete fetch_oids passes. A one-second probe allowance is sized at two passes per effective CPU, using Linux cgroup v2 cpu.max or available processors elsewhere. Healthy windows grow the allowance by 20% or one, whichever is greater, without a lasting ceiling. This gradual discovery lets retrospective pressure counters reflect each preceding step. New CPU throttling or memory-high pressure reduces new admission to one while already-running work drains. All event-directed purgatory recovery uses this background gate, regardless of event origin. Direct Git pushes remain foreground and satisfy purgatory without the outbound recovery path. Existing per-domain concurrency and request-rate limits remain the unconditional remote DoS protection. This deliberately avoids a permanent aggregate limit, configuration, multi-class scheduling, or a general feedback controller. The brief probe window exists only so the kernel can produce the pressure signal before a cold-start burst escapes. Validation: nix develop -c cargo test --lib purgatory::sync::throttle (14 passed); git diff --check. Production reproduction began 2026-08-10 15:14:20 UTC on the disposable archive with CPUQuota=200%, MemoryHigh=12G, and MemoryMax=16G. --- docs/explanation/architecture.md | 16 ++ src/purgatory/sync/functions.rs | 8 +- src/purgatory/sync/throttle.rs | 300 +++++++++++++++++++++++++++++++ 3 files changed, 323 insertions(+), 1 deletion(-) diff --git a/docs/explanation/architecture.md b/docs/explanation/architecture.md index 3737c64..f7d25c8 100644 --- a/docs/explanation/architecture.md +++ b/docs/explanation/architecture.md @@ -310,6 +310,22 @@ pub struct Purgatory { - Defers matching expiry while a concrete background Git sync is actively running; queued/backoff work receives no extension +6. **Pressure-aware Background Git Fan-out**: event-directed recovery begins + with a one-second probe allowance of two remote Git passes per effective CPU + - Healthy observation windows grow the allowance by 20% or one pass, + whichever is greater, avoiding a lasting cap while giving the + retrospective pressure signal time to reflect each preceding step + - New CPU throttling or memory-high pressure reduces new admission to one + pass while already-running work drains naturally + - Linux services read their cgroup v2 constraints and pressure counters; + other environments fall back to processor count and healthy-window growth + - Existing per-domain limits remain responsible for being polite to each + remote server + - The probe and pressure gate prevent a cold sync spanning many domains from + spawning enough Git processes to starve this relay's HTTP and WebSocket + service + - Foreground Git requests do not draw from this background-only limit + #### Data Types See [`types.rs`](../../src/purgatory/types.rs) for complete definitions: diff --git a/src/purgatory/sync/functions.rs b/src/purgatory/sync/functions.rs index b8b3948..750b00e 100644 --- a/src/purgatory/sync/functions.rs +++ b/src/purgatory/sync/functions.rs @@ -367,7 +367,13 @@ pub async fn sync_identifier_from_url( return 0; } - // Perform the fetch with throttle tracking + // Per-domain limits alone do not bound aggregate work when a cold sync + // discovers many domains simultaneously. Wait for the process-wide permit + // before occupying a domain slot, then retain it for the complete remote + // Git pass so subprocess fan-out stays bounded. + let _background_permit = throttle_manager.acquire_global_request().await; + + // Perform the fetch with per-domain throttle tracking. throttle_manager.start_request(&domain); let fetch_result = ctx.fetch_oids(&target_repo, url, &needed_oids).await; throttle_manager.complete_request(&domain); diff --git a/src/purgatory/sync/throttle.rs b/src/purgatory/sync/throttle.rs index 7f8f636..07319b4 100644 --- a/src/purgatory/sync/throttle.rs +++ b/src/purgatory/sync/throttle.rs @@ -20,12 +20,231 @@ use indexmap::IndexMap; use std::collections::{HashSet, VecDeque}; use std::sync::{Arc, Mutex, OnceLock}; use std::time::{Duration, Instant}; +use tokio::sync::Notify; use tracing::debug; use super::context::SyncContext; use super::functions::{sync_identifier_from_url, sync_identifier_next_url}; use crate::sync::naughty_list::NaughtyListTracker; +/// The first observation window admits a modest probe rather than launching a +/// cold-sync burst before the kernel has produced a pressure signal. Healthy +/// windows grow this allowance without a ceiling, so an unconstrained host +/// progressively returns to effectively unrestricted background concurrency. +const PURGATORY_GIT_FETCHES_PER_CPU: usize = 2; +const MIN_PROBE_CONCURRENCY: usize = 2; +const PRESSURE_CONCURRENCY: usize = 1; +const PRESSURE_OBSERVATION_INTERVAL: Duration = Duration::from_secs(1); +const PRESSURE_RECHECK_INTERVAL: Duration = Duration::from_millis(100); +// Sustained CPU throttling can alternate healthy and pressured samples. Keep +// the operational signal without turning that expected oscillation into spam. +const PRESSURE_LOG_INTERVAL: Duration = Duration::from_secs(60); + +fn grow_healthy_allowance(current: usize) -> usize { + // Pressure counters describe work that has already run. Grow gradually so + // each one-second observation has a chance to expose the previous step. + current.saturating_add((current / 5).max(1)) +} + +fn parse_cpu_max(value: &str) -> Option { + let mut fields = value.split_whitespace(); + let quota = fields.next()?; + let period = fields.next()?.parse::().ok()?; + if quota == "max" || period == 0 { + return None; + } + Some(quota.parse::().ok()? as f64 / period as f64) +} + +#[cfg(target_os = "linux")] +fn cgroup_cpu_quota() -> Option { + let membership = std::fs::read_to_string("/proc/self/cgroup").ok()?; + let relative = membership.lines().find_map(|line| { + let mut fields = line.splitn(3, ':'); + match (fields.next(), fields.next(), fields.next()) { + (Some("0"), Some(""), Some(path)) => Some(path.trim_start_matches('/')), + _ => None, + } + })?; + let cpu_max = std::fs::read_to_string( + std::path::Path::new("/sys/fs/cgroup") + .join(relative) + .join("cpu.max"), + ) + .ok()?; + parse_cpu_max(&cpu_max) +} + +#[cfg(not(target_os = "linux"))] +fn cgroup_cpu_quota() -> Option { + None +} + +fn purgatory_git_probe_limit() -> usize { + let effective_cpus = cgroup_cpu_quota().unwrap_or_else(|| { + std::thread::available_parallelism() + .map(usize::from) + .unwrap_or(1) as f64 + }); + (effective_cpus.ceil() as usize * PURGATORY_GIT_FETCHES_PER_CPU).max(MIN_PROBE_CONCURRENCY) +} + +#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)] +struct ResourceSnapshot { + throttled_usec: u64, + memory_high_events: u64, + memory_current: u64, + memory_high: Option, +} + +impl ResourceSnapshot { + fn pressured_since(self, previous: Self) -> bool { + self.throttled_usec > previous.throttled_usec + || self.memory_high_events > previous.memory_high_events + || self + .memory_high + .is_some_and(|high| high > 0 && self.memory_current >= high * 9 / 10) + } +} + +#[cfg(target_os = "linux")] +fn cgroup_resource_snapshot() -> Option { + let membership = std::fs::read_to_string("/proc/self/cgroup").ok()?; + let relative = membership.lines().find_map(|line| { + let mut fields = line.splitn(3, ':'); + match (fields.next(), fields.next(), fields.next()) { + (Some("0"), Some(""), Some(path)) => Some(path.trim_start_matches('/')), + _ => None, + } + })?; + let directory = std::path::Path::new("/sys/fs/cgroup").join(relative); + let stat = std::fs::read_to_string(directory.join("cpu.stat")).ok()?; + let throttled_usec = named_stat(&stat, "throttled_usec").unwrap_or(0); + let memory_events = std::fs::read_to_string(directory.join("memory.events")).ok()?; + let memory_high_events = named_stat(&memory_events, "high").unwrap_or(0); + let memory_current = read_u64(directory.join("memory.current")).unwrap_or(0); + let memory_high = std::fs::read_to_string(directory.join("memory.high")) + .ok() + .and_then(|value| value.trim().parse().ok()); + Some(ResourceSnapshot { + throttled_usec, + memory_high_events, + memory_current, + memory_high, + }) +} + +#[cfg(not(target_os = "linux"))] +fn cgroup_resource_snapshot() -> Option { + None +} + +fn named_stat(contents: &str, name: &str) -> Option { + contents.lines().find_map(|line| { + let mut fields = line.split_whitespace(); + (fields.next()? == name).then(|| fields.next()?.parse().ok())? + }) +} + +fn read_u64(path: impl AsRef) -> Option { + std::fs::read_to_string(path).ok()?.trim().parse().ok() +} + +#[derive(Debug)] +struct PressureGateState { + active: usize, + allowance: usize, + last_adjustment: Instant, + last_pressure_log: Option, + last_snapshot: Option, +} + +#[derive(Debug)] +struct BackgroundPressureGate { + probe_allowance: usize, + state: Mutex, + notify: Notify, +} + +impl BackgroundPressureGate { + fn new() -> Arc { + let probe_allowance = purgatory_git_probe_limit(); + Arc::new(Self { + probe_allowance, + state: Mutex::new(PressureGateState { + active: 0, + allowance: probe_allowance, + last_adjustment: Instant::now(), + last_pressure_log: None, + last_snapshot: cgroup_resource_snapshot(), + }), + notify: Notify::new(), + }) + } + + async fn acquire(self: &Arc) -> BackgroundPressurePermit { + loop { + let notified = self.notify.notified(); + let snapshot = cgroup_resource_snapshot(); + let admitted = { + let mut state = self.state.lock().unwrap(); + let now = Instant::now(); + let pressured = snapshot + .zip(state.last_snapshot) + .is_some_and(|(current, previous)| current.pressured_since(previous)); + if snapshot.is_some() { + state.last_snapshot = snapshot; + } + if pressured { + if state + .last_pressure_log + .is_none_or(|last| now.duration_since(last) >= PRESSURE_LOG_INTERVAL) + { + tracing::warn!( + active = state.active, + "Resource pressure paused new background Git fan-out" + ); + state.last_pressure_log = Some(now); + } + state.allowance = PRESSURE_CONCURRENCY; + state.last_adjustment = now; + } else if now.duration_since(state.last_adjustment) >= PRESSURE_OBSERVATION_INTERVAL + { + state.allowance = grow_healthy_allowance(state.allowance); + state.last_adjustment = now; + tracing::debug!( + allowance = state.allowance, + "Raised background Git allowance after a healthy resource window" + ); + } + if state.active < state.allowance { + state.active += 1; + true + } else { + false + } + }; + if admitted { + return BackgroundPressurePermit { gate: self.clone() }; + } + let _ = tokio::time::timeout(PRESSURE_RECHECK_INTERVAL, notified).await; + } + } +} + +pub(super) struct BackgroundPressurePermit { + gate: Arc, +} + +impl Drop for BackgroundPressurePermit { + fn drop(&mut self) { + let mut state = self.gate.state.lock().unwrap(); + state.active = state.active.saturating_sub(1); + drop(state); + self.gate.notify.notify_waiters(); + } +} + /// State for an identifier waiting in a domain's queue. /// /// Tracks which URLs from this domain have been tried and whether @@ -263,6 +482,9 @@ pub struct ThrottleManager { /// Maximum requests per minute per domain. max_per_minute_per_domain: u32, + /// Resource-pressure admission gate shared across every remote domain. + pressure_gate: Arc, + /// Sync context for processing queued identifiers. /// Set once at startup via `set_context()`. ctx: OnceLock>, @@ -279,15 +501,30 @@ impl ThrottleManager { /// * `max_concurrent` - Maximum concurrent in-flight requests per domain /// * `max_per_minute` - Maximum requests per 60-second window per domain pub fn new(max_concurrent: u32, max_per_minute: u32) -> Self { + let pressure_gate = BackgroundPressureGate::new(); + tracing::info!( + probe_allowance = pressure_gate.probe_allowance, + cpu_quota = ?cgroup_cpu_quota(), + "Configured resource-aware background Git admission" + ); Self { throttles: DashMap::new(), max_concurrent_per_domain: max_concurrent, max_per_minute_per_domain: max_per_minute, + pressure_gate, ctx: OnceLock::new(), git_naughty_list: OnceLock::new(), } } + /// Wait for process-wide background Git capacity. + /// + /// The owned permit spans the complete `fetch_oids` pass (ls-remote plus + /// any fetches) and releases automatically on every return or cancellation. + pub(super) async fn acquire_global_request(&self) -> BackgroundPressurePermit { + self.pressure_gate.acquire().await + } + /// Set the sync context (called once at startup). /// /// The context is used for processing queued identifiers when capacity @@ -769,6 +1006,69 @@ mod tests { assert!(manager.throttles.contains_key("example.com")); } + #[tokio::test] + async fn initial_probe_bounds_work_before_pressure_is_observable() { + let manager = Arc::new(ThrottleManager::new(5, 100)); + let probe_allowance = purgatory_git_probe_limit(); + let mut permits = Vec::new(); + for _ in 0..probe_allowance { + permits.push(manager.acquire_global_request().await); + } + + let waiting_manager = manager.clone(); + let mut waiting = + tokio::spawn(async move { waiting_manager.acquire_global_request().await }); + assert!( + tokio::time::timeout(Duration::from_millis(50), &mut waiting) + .await + .is_err(), + "the initial probe must wait before launching an unmeasured burst" + ); + + permits.pop(); + let _permit = tokio::time::timeout(Duration::from_secs(1), waiting) + .await + .expect("released capacity should wake a waiter") + .expect("waiter task should finish"); + } + + #[test] + fn cpu_quota_scales_background_git_capacity() { + assert_eq!(parse_cpu_max("200000 100000\n"), Some(2.0)); + assert_eq!(parse_cpu_max("150000 100000\n"), Some(1.5)); + assert_eq!(parse_cpu_max("max 100000\n"), None); + } + + #[test] + fn healthy_allowance_grows_gradually_without_a_ceiling() { + assert_eq!(grow_healthy_allowance(4), 5); + assert_eq!(grow_healthy_allowance(5), 6); + assert_eq!(grow_healthy_allowance(10), 12); + assert_eq!(grow_healthy_allowance(100), 120); + assert_eq!(grow_healthy_allowance(usize::MAX), usize::MAX); + } + + #[test] + fn resource_snapshot_detects_new_pressure_only() { + let previous = ResourceSnapshot { + throttled_usec: 10, + memory_high_events: 2, + memory_current: 50, + memory_high: Some(100), + }; + assert!(!previous.pressured_since(previous)); + assert!(ResourceSnapshot { + throttled_usec: 11, + ..previous + } + .pressured_since(previous)); + assert!(ResourceSnapshot { + memory_current: 90, + ..previous + } + .pressured_since(previous)); + } + #[test] fn has_queued_work_reflects_queue_state() { let manager = ThrottleManager::new(5, 100);