diff --git a/crates/buzz-acp/src/backoff.rs b/crates/buzz-acp/src/backoff.rs new file mode 100644 index 000000000..5377ccb82 --- /dev/null +++ b/crates/buzz-acp/src/backoff.rs @@ -0,0 +1,157 @@ +//! Jitter helpers for retry, backoff, and rate-limit gate delays. +//! +//! Jitter uses the nanosecond sub-second component of the system clock as a +//! cheap entropy source (no `rand` dependency). The factor computation is a +//! pure function of that value so the reachable range can be tested directly — +//! a wall-clock helper cannot be driven to its domain boundary, which is how +//! the divisor defect below survived every existing bound assertion. +//! +//! Two shapes, and the distinction is load-bearing: +//! +//! * [`jittered_duration`] is **symmetric** (±20%). Correct for self-chosen +//! backoff ladders, where waking early only costs an extra attempt. +//! * [`extend_with`] is **one-sided** (+0–20%, never shorter). Required +//! wherever the base duration is an authoritative deadline supplied by the +//! relay: shortening a `retry in {N}s` hint wakes into a window the relay has +//! already told us is closed, which earns a fresh denial and burns a counter +//! increment for no chance of success. + +use std::time::Duration; + +/// Nanoseconds in a second — the true domain bound of `Duration::subsec_nanos`. +/// +/// The previous implementations divided by `u32::MAX` (4,294,967,295), which is +/// 4.295x larger than the largest value `subsec_nanos()` can return. That +/// capped the symmetric factor at 0.893 instead of approaching 1.2, making +/// "±20% jitter" unconditionally negative and averaging −15%. +const NANOS_PER_SEC: f64 = 1_000_000_000.0; + +/// Map a sub-second nanosecond count onto `[0.0, 1.0]`. +/// +/// Clamped so the function is total for any `u32`; inputs above one second are +/// unreachable from `subsec_nanos()` but must not push the factor out of range. +fn nanos_fraction(nanos: u32) -> f64 { + (f64::from(nanos) / NANOS_PER_SEC).min(1.0) +} + +/// Symmetric jitter factor in `[0.8, 1.2)` over the domain `[0, 1s)`. +pub(crate) fn symmetric_jitter_factor(nanos: u32) -> f64 { + 0.8 + nanos_fraction(nanos) * 0.4 +} + +/// One-sided jitter factor in `[1.0, 1.2)` over the domain `[0, 1s)`. +/// +/// Never returns less than 1.0, so a delay built from it can never fall below +/// its base. This is a structural guarantee, not a policy applied by callers. +pub(crate) fn nonnegative_jitter_factor(nanos: u32) -> f64 { + 1.0 + nanos_fraction(nanos) * 0.2 +} + +/// Sub-second component of the current wall clock, used as the jitter source. +pub(crate) fn clock_nanos() -> u32 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .unwrap_or_default() + .subsec_nanos() +} + +/// Apply symmetric ±20% jitter to a self-chosen backoff duration. +pub(crate) fn jittered_duration(base: Duration) -> Duration { + base.mul_f64(symmetric_jitter_factor(clock_nanos())) +} + +/// Extend an authoritative deadline by 0–20%, never shortening it. +/// +/// The entropy sample is an explicit parameter so tests can drive the whole +/// domain deterministically. With a wall-clock-only helper the "never +/// shortens" property is only probabilistically testable: a symmetric-jitter +/// regression satisfies it on roughly half of all draws, which reads as a +/// surviving mutant rather than a flaky assertion. Callers pass +/// [`clock_nanos`]. +pub(crate) fn extend_with(base: Duration, nanos: u32) -> Duration { + base.mul_f64(nonnegative_jitter_factor(nanos)) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// The largest value `Duration::subsec_nanos()` can return. Asserting the + /// ceiling at `u32::MAX` instead would re-encode the very divisor defect + /// these tests exist to catch, so the domain bound is pinned explicitly. + const MAX_SUBSEC_NANOS: u32 = 999_999_999; + + /// Symmetric factor reaches its documented floor and (near) ceiling. + /// + /// The ceiling assertion is the discriminating one: the `u32::MAX` divisor + /// caps this at 0.893, while every range assertion of the form + /// `0.8 <= f < 1.2` passes in both the broken and the fixed implementation. + #[test] + fn symmetric_factor_spans_its_documented_range() { + assert_eq!(symmetric_jitter_factor(0), 0.8, "floor must be exactly 0.8"); + let ceiling = symmetric_jitter_factor(MAX_SUBSEC_NANOS); + assert!( + ceiling > 1.199 && ceiling < 1.2, + "ceiling {ceiling} must approach 1.2 from below — a ceiling near 0.893 \ + means the factor is divided by u32::MAX instead of 1e9" + ); + } + + /// One-sided factor never shortens its base and reaches +20%. + #[test] + fn nonnegative_factor_spans_its_documented_range() { + assert_eq!( + nonnegative_jitter_factor(0), + 1.0, + "floor must be exactly 1.0 — a one-sided factor may never shorten" + ); + let ceiling = nonnegative_jitter_factor(MAX_SUBSEC_NANOS); + assert!( + ceiling > 1.199 && ceiling < 1.2, + "ceiling {ceiling} must approach 1.2 from below" + ); + } + + /// No input can make the one-sided factor shorten a deadline. Swept across + /// the whole `u32` range, not just the reachable sub-second domain. + #[test] + fn nonnegative_factor_is_never_below_one() { + for step in 0..=1_000u32 { + let nanos = (u32::MAX / 1_000).saturating_mul(step); + let factor = nonnegative_jitter_factor(nanos); + assert!( + (1.0..=1.2 + 1e-9).contains(&factor), + "factor {factor} out of [1.0, 1.2] at nanos={nanos}" + ); + } + } + + /// Both factors stay in range for out-of-domain inputs (clamped). + /// + /// `subsec_nanos()` can never return these values; the clamp exists so the + /// factor functions are total. Compared with a tolerance because + /// `0.8 + 1.0 * 0.4` is not exactly 1.2 in binary floating point. + #[test] + fn factors_are_clamped_above_one_second() { + assert!((symmetric_jitter_factor(u32::MAX) - 1.2).abs() < 1e-9); + assert!((nonnegative_jitter_factor(u32::MAX) - 1.2).abs() < 1e-9); + } + + /// The wrappers apply their factor to the base duration. + #[test] + fn wrappers_stay_within_their_factor_ranges() { + let base = Duration::from_secs(5); + for _ in 0..64 { + let symmetric = jittered_duration(base); + assert!( + symmetric >= base.mul_f64(0.8) && symmetric <= base.mul_f64(1.2), + "symmetric {symmetric:?} out of range" + ); + let extended = extend_with(base, clock_nanos()); + assert!( + extended >= base && extended <= base.mul_f64(1.2), + "extended {extended:?} shortened the base" + ); + } + } +} diff --git a/crates/buzz-acp/src/lib.rs b/crates/buzz-acp/src/lib.rs index 811253e4a..eb848d742 100644 --- a/crates/buzz-acp/src/lib.rs +++ b/crates/buzz-acp/src/lib.rs @@ -1,6 +1,7 @@ #![deny(unsafe_code)] mod acp; +mod backoff; mod config; mod engram_fetch; mod filter; @@ -1086,13 +1087,7 @@ impl SlotCircuit { // Exponential backoff: 1s * 2^(recent-1), capped at 30s, with ±20% jitter. let base = RESPAWN_BASE_DELAY.saturating_mul(1u32 << (recent - 1).min(5)); let capped = base.min(RESPAWN_MAX_DELAY); - let jitter = (std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .subsec_nanos() as f64) - / 1_000_000_000.0; // 0.0..1.0 - let factor = 0.8 + jitter * 0.4; // 0.8..1.2 - CrashVerdict::Respawn(capped.mul_f64(factor)) + CrashVerdict::Respawn(crate::backoff::jittered_duration(capped)) } /// Mark a spawn failure — opens the circuit so the slot isn't retried diff --git a/crates/buzz-acp/src/queue.rs b/crates/buzz-acp/src/queue.rs index 5c960de20..29787ee33 100644 --- a/crates/buzz-acp/src/queue.rs +++ b/crates/buzz-acp/src/queue.rs @@ -453,15 +453,7 @@ impl EventQueue { // Exponential backoff: BASE * 2^(attempt-1), capped at MAX, with ±20% jitter. let base_secs = BASE_RETRY_DELAY_SECS.saturating_mul(1u64 << (attempt - 1).min(6)); let capped_secs = base_secs.min(MAX_RETRY_DELAY_SECS); - // Jitter: multiply by 0.8..1.2 using subsecond nanos as entropy source. - let jitter = { - let nanos = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .subsec_nanos(); - 0.8 + (nanos as f64 / u32::MAX as f64) * 0.4 - }; - let delay = Duration::from_secs_f64(capped_secs as f64 * jitter); + let delay = crate::backoff::jittered_duration(Duration::from_secs(capped_secs)); tracing::warn!( channel_id = %channel_id, diff --git a/crates/buzz-acp/src/relay.rs b/crates/buzz-acp/src/relay.rs index aea5cee07..ee4f1e33b 100644 --- a/crates/buzz-acp/src/relay.rs +++ b/crates/buzz-acp/src/relay.rs @@ -126,6 +126,7 @@ use tokio_tungstenite::{connect_async, tungstenite::Message, MaybeTlsStream, Web use tracing::{debug, info, warn}; use uuid::Uuid; +use crate::backoff::{extend_with, jittered_duration}; use crate::config::ChannelFilter; /// Metadata about a channel, populated at discovery time. @@ -1161,15 +1162,31 @@ impl BgState { /// the desktop TypeScript client, which uses a 10s no-hint default — both /// values are conservative enough; the relay hint wins when present. /// + /// Jitter here is **one-sided** ([`crate::backoff::extend_with`]): the hint is the + /// relay's authoritative window TTL, so the gate may only ever be longer + /// than it. Symmetric jitter would wake us inside a window the relay has + /// already said is closed, earning a fresh denial plus a wasted counter + /// increment (the limiter's `INCR` runs on denied checks too). + /// /// The gate takes the **maximum** of any existing deadline and the newly /// computed one so overlapping CLOSED/NOTICE messages can't shorten a gate /// that is already set further out. /// /// Returns the gate deadline that was set. fn set_rate_limit_gate(&mut self, retry_secs: u64) -> tokio::time::Instant { - let secs = if retry_secs < 2 { 5 } else { retry_secs }; - let base = Duration::from_secs(secs); - let deadline = tokio::time::Instant::now() + jittered_duration(base); + self.set_rate_limit_gate_with(retry_secs, crate::backoff::clock_nanos()) + } + + /// [`set_rate_limit_gate`] with an explicit jitter sample. + /// + /// Delegates the arithmetic to the pure [`gate_delay`]; this wrapper exists + /// only to apply the result to `self`. + fn set_rate_limit_gate_with( + &mut self, + retry_secs: u64, + jitter_nanos: u32, + ) -> tokio::time::Instant { + let deadline = tokio::time::Instant::now() + gate_delay(retry_secs, jitter_nanos); let gate = match self.rate_limit_gate { Some(existing) if existing > deadline => existing, _ => deadline, @@ -3345,16 +3362,28 @@ pub(crate) fn parse_rate_limit_retry_secs(msg: &str) -> Option { after[..len].parse::().ok() } -/// Add ±20% jitter to a backoff duration using the nanosecond sub-second -/// component of the system clock as a cheap entropy source (no `rand` dep). -fn jittered_duration(base: Duration) -> Duration { - let nanos = std::time::SystemTime::now() - .duration_since(std::time::UNIX_EPOCH) - .unwrap_or_default() - .subsec_nanos(); - // factor ∈ [0.8, 1.2) - let factor = 0.8 + (nanos as f64 / u32::MAX as f64) * 0.4; - base.mul_f64(factor) +/// Sub-2s hints (including a missing hint parsed as 0) floor to this many +/// seconds, matching the desktop TypeScript client's no-hint default. +const RATE_LIMIT_GATE_FLOOR_SECS: u64 = 5; + +/// The rate-limit gate delay for a relay hint and an explicit jitter sample. +/// +/// Pure: entropy enters only as `jitter_nanos`, and the wall clock is acquired +/// by the caller. That is what makes the gate's safety property testable at +/// exact endpoints — a helper that reads `SystemTime` internally can only be +/// sampled, never driven, and paused Tokio time does not freeze `SystemTime`. +/// +/// Jitter is **one-sided**: the hint is the relay's authoritative window TTL, +/// so the delay may only ever exceed it. Symmetric jitter would wake us inside +/// a window the relay has already said is closed, earning a fresh denial plus +/// a wasted counter increment (the limiter's `INCR` runs on denied checks too). +fn gate_delay(retry_secs: u64, jitter_nanos: u32) -> Duration { + let secs = if retry_secs < 2 { + RATE_LIMIT_GATE_FLOOR_SECS + } else { + retry_secs + }; + extend_with(Duration::from_secs(secs), jitter_nanos) } /// Classify a `RelayError` as a DNS resolution failure. @@ -5765,7 +5794,7 @@ mod tests { "gate must be active immediately after arming" ); - // Advance virtual time past the max jitter (1.2 × 5 s = 6 s). + // Advance virtual time past the max gate length (1.2 × 5 s = 6 s). tokio::time::advance(Duration::from_secs(7)).await; assert!( @@ -5778,6 +5807,102 @@ mod tests { ); } + /// The effective gate base for a hint: sub-2s hints floor to 5s. + #[cfg(test)] + fn effective_gate_secs(hint: u64) -> u64 { + if hint < 2 { + 5 + } else { + hint + } + } + + /// `gate_delay` is exact at both endpoints of the jitter domain. + /// + /// This is the discriminating test, and it is a **pure-function** test for + /// a reason. Entropy enters `gate_delay` only as a parameter, so both + /// endpoints are driven rather than sampled: + /// + /// * `nanos = 0` ⇒ one-sided factor exactly 1.0 ⇒ delay is exactly base. + /// A symmetric-jitter regression yields `0.8 * base` here, deterministically. + /// * `nanos = 999_999_999` ⇒ factor approaches 1.2 ⇒ delay approaches + /// `1.2 * base`. The `u32::MAX` divisor caps it at `0.893 * base`, + /// so this endpoint is what catches the divisor defect; a range assertion + /// of the form `0.8*base <= d < 1.2*base` cannot, because the broken range + /// is a strict subset of the asserted one. The upper bound is inclusive: + /// `Duration` is nanosecond-resolution, so e.g. `2s * 1.1999999998` + /// rounds to exactly `2.4s` and a strict `<` would be unsatisfiable. + /// + /// Scope: this kills a swap of one-sided for symmetric *inside* the seam, + /// and the divisor defect. A mutation that bypasses the seam entirely (call + /// site reverting to a wall-clock helper) is a structural change this unit + /// test cannot honestly claim to kill — `#[deny(unused)]` on the parameter + /// and the integration test below are what cover that. + #[test] + fn gate_delay_is_exact_at_both_jitter_endpoints() { + for hint in [0u64, 1, 2, 5, 6, 51] { + let base = Duration::from_secs(effective_gate_secs(hint)); + + assert_eq!( + gate_delay(hint, 0), + base, + "a {hint}s hint at zero jitter must be exactly {base:?}; \ + symmetric jitter would yield 0.8 x base here" + ); + + let ceiling = gate_delay(hint, 999_999_999); + assert!( + ceiling > base.mul_f64(1.1999) && ceiling <= base.mul_f64(1.2), + "a {hint}s hint at max jitter gave {ceiling:?}, which must reach \ + 1.2 x {base:?} — a value near 0.893 x base means the \ + factor is divided by u32::MAX instead of 1e9" + ); + } + } + + /// `gate_delay` never returns less than the relay's hint, for any sample. + #[test] + fn gate_delay_never_falls_below_the_hint() { + for hint in [0u64, 1, 2, 5, 6, 51] { + let base = Duration::from_secs(effective_gate_secs(hint)); + for step in 0..=1_000u32 { + let nanos = (999_999_999f64 * f64::from(step) / 1_000.0) as u32; + let delay = gate_delay(hint, nanos); + assert!( + delay >= base && delay <= base.mul_f64(1.2), + "a {hint}s hint with jitter nanos={nanos} gave {delay:?}, \ + outside [{base:?}, 1.2 x base]" + ); + } + } + } + + /// Corroboration: the armed gate honours the computed delay. + /// + /// Integration-level, so it can only sample the entropy the call site + /// actually uses. It states the safety property end to end; the + /// deterministic receipts live in the pure-function tests above. + #[tokio::test(start_paused = true)] + async fn rate_limit_gate_arms_for_the_computed_delay() { + for hint in [0u64, 1, 2, 5, 6, 51] { + let base = Duration::from_secs(effective_gate_secs(hint)); + for nanos in [0u32, 500_000_000, 999_999_999] { + let mut state = BgState::new(); + let before = tokio::time::Instant::now(); + let gate = state.set_rate_limit_gate_with(hint, nanos); + assert_eq!( + gate - before, + gate_delay(hint, nanos), + "armed gate for a {hint}s hint (nanos={nanos}) must equal gate_delay" + ); + assert!( + gate - before >= base, + "armed gate for a {hint}s hint (nanos={nanos}) fell below the hint" + ); + } + } + } + /// set_rate_limit_gate takes the max of overlapping deadlines. #[tokio::test(start_paused = true)] async fn rate_limit_gate_extends_to_max() { @@ -6100,7 +6225,7 @@ mod tests { "gate must be active while membership sub is pending" ); - // Advance past the gate (max jitter: 1.2 × 5s = 6s). + // Advance past the gate (max gate length: 1.2 × 5s = 6s). tokio::time::advance(Duration::from_secs(7)).await; assert!(