mirror of
https://relay.ngit.dev/npub15qydau2hjma6ngxkl2cyar74wzyjshvl65za5k5rl69264ar2exs5cyejr/ngit-grasp.git
synced 2026-10-05 23:18:24 +00:00
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.
This commit is contained in:
@@ -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:
|
||||
|
||||
@@ -367,7 +367,13 @@ pub async fn sync_identifier_from_url<C: SyncContext + ?Sized>(
|
||||
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);
|
||||
|
||||
@@ -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<f64> {
|
||||
let mut fields = value.split_whitespace();
|
||||
let quota = fields.next()?;
|
||||
let period = fields.next()?.parse::<u64>().ok()?;
|
||||
if quota == "max" || period == 0 {
|
||||
return None;
|
||||
}
|
||||
Some(quota.parse::<u64>().ok()? as f64 / period as f64)
|
||||
}
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
fn cgroup_cpu_quota() -> Option<f64> {
|
||||
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<f64> {
|
||||
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<u64>,
|
||||
}
|
||||
|
||||
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<ResourceSnapshot> {
|
||||
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<ResourceSnapshot> {
|
||||
None
|
||||
}
|
||||
|
||||
fn named_stat(contents: &str, name: &str) -> Option<u64> {
|
||||
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<std::path::Path>) -> Option<u64> {
|
||||
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<Instant>,
|
||||
last_snapshot: Option<ResourceSnapshot>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
struct BackgroundPressureGate {
|
||||
probe_allowance: usize,
|
||||
state: Mutex<PressureGateState>,
|
||||
notify: Notify,
|
||||
}
|
||||
|
||||
impl BackgroundPressureGate {
|
||||
fn new() -> Arc<Self> {
|
||||
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<Self>) -> 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<BackgroundPressureGate>,
|
||||
}
|
||||
|
||||
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<BackgroundPressureGate>,
|
||||
|
||||
/// Sync context for processing queued identifiers.
|
||||
/// Set once at startup via `set_context()`.
|
||||
ctx: OnceLock<Arc<dyn SyncContext>>,
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user