Skip to main content

nmbrs_runtime/
throttle.rs

1// Copyright 2024-2026 Jonathan Shook
2// SPDX-License-Identifier: Apache-2.0
3
4//! Phase `throttle:` — the adaptive backpressure governor (SRD-83
5//! Part 9).
6//!
7//! A saturated target converts client overload into retry churn: the
8//! tries wrapper absorbs server rejections, `ok:%` stays green, and
9//! the only truthful signal is the ATTEMPT plane — the windowed
10//! attempt-failure fraction (`attempt_failure / resolved attempts`,
11//! the see-through-retries view).
12//!
13//! The governor assumes the MOST FRAGILE target by default and scales
14//! to robust ones, TCP-style:
15//!
16//! - **Slow-start.** The offered value begins at `start` (default:
17//!   `floor`) — the authored `concurrency:`/`rate:` is the CEILING,
18//!   not the opening offer. While no congestion has been seen, each
19//!   clean window DOUBLES the offer toward the ceiling: a robust
20//!   target climbs to full load in a handful of windows with zero
21//!   failures; a fragile one is never assaulted at all. A target
22//!   known to be robust at phase entry declares `start:` explicitly.
23//! - **Severity-proportional back-off.** Above `high`, the offer is
24//!   multiplied by `clamp(1 − frac, 0.25, 0.9)`: a marginal breach
25//!   trims gently (×0.9), total failure collapses fast (×0.25),
26//!   never below `floor`.
27//! - **Congestion memory.** Each back-off records the offer at which
28//!   failure was observed (`last_bad`). Recovery climbs ×1.5 through
29//!   the proven-safe zone (up to 75% of `last_bad`), then probes
30//!   ADDITIVELY (+max(1, 2% of `last_bad`) per clean window) — no
31//!   more marching multiplicatively back into the same wall.
32//!   [`MEMORY_CLEAR_STREAK`] CONSECUTIVE clean windows at-or-above
33//!   `last_bad` clear the memory (the target got healthier — warmed
34//!   caches, finished compactions), restoring the multiplicative
35//!   climb. One clean window is not evidence: the additive probes
36//!   keep stepping through the streak, so each window in it sits a
37//!   notch higher than the last, and a single quiet window at a
38//!   marginal congestion point can never re-arm the doubling climb
39//!   straight back into the wall.
40//!
41//! Windows are computed from counter DELTAS on the drain-loop tick —
42//! a true trailing window, never a lifetime average. Writes ride the
43//! push-on-set control path (`ControlOrigin::Governor`,
44//! confirmed-apply, spawned off the loop); every movement logs one
45//! line naming the signal — visible, never silent.
46//!
47//! Measurement honesty: the throttled steady state IS the
48//! measurement — the target's capacity at the declared failure
49//! bound. A load figure taken at high attempt-failure is a
50//! saturation artifact.
51
52use std::sync::Arc;
53use std::time::{Duration, Instant};
54
55use nmbrs_metrics::controls::{ControlOrigin, ErasedControl};
56
57/// All governor lines carry the `Throttle` category (in-flight —
58/// governance happens mid-body, attached to no boundary), so sinks
59/// and counters can dispatch on the axis instead of matching the
60/// rendered `throttle:` prefix.
61macro_rules! gov_log {
62    ($level:expr, $($arg:tt)*) => {
63        crate::observer::log_tagged(
64            $level,
65            crate::observer::EventTag::in_flight(
66                crate::observer::EventCategory::Throttle),
67            &format!($($arg)*),
68        )
69    };
70}
71
72/// What one window decided — pure, unit-testable.
73#[derive(Debug, Clone, Copy, PartialEq)]
74pub enum Decision {
75    /// Signal above `high`: back off to the contained value.
76    Down(f64),
77    /// Clean window with headroom: raise to the contained value.
78    Up(f64),
79    /// Hold (dead band, no traffic, or at a bound).
80    Hold,
81}
82
83/// Consecutive clean windows at-or-above the remembered congestion
84/// point required before the memory clears and the multiplicative
85/// climb resumes. Additive probing continues through the streak, so
86/// clearing means the target stayed clean across a rising run of
87/// offers, not one lucky window.
88pub const MEMORY_CLEAR_STREAK: u32 = 3;
89
90/// Pure congestion-memory step: fold one window's evidence (`clean
91/// at-or-above last_bad`) into the running streak. Returns the new
92/// streak and whether the memory clears on this window. Any window
93/// without that evidence — a breach, a dead-band hold, or a clean
94/// window still below the congestion point — resets the streak.
95pub fn memory_clear_step(streak: u32, evidence: bool) -> (u32, bool) {
96    if !evidence {
97        return (0, false);
98    }
99    let streak = streak + 1;
100    if streak >= MEMORY_CLEAR_STREAK {
101        (0, true)
102    } else {
103        (streak, false)
104    }
105}
106
107/// Pure governor step. `frac` is the windowed attempt-failure
108/// fraction, `current` the committed offer, `last_bad` the offer at
109/// which congestion was last observed (`None` = unexplored — slow
110/// start).
111pub fn decide(
112    frac: f64,
113    current: f64,
114    high: f64,
115    low: f64,
116    floor: f64,
117    ceiling: f64,
118    last_bad: Option<f64>,
119) -> Decision {
120    if frac > high {
121        // Severity-proportional multiplicative decrease: marginal
122        // breach trims ×0.9; total failure collapses ×0.25.
123        let target = (current * (1.0 - frac).clamp(0.25, 0.9)).max(floor);
124        if target < current {
125            return Decision::Down(target);
126        }
127        return Decision::Hold;
128    }
129    if frac < low && current < ceiling {
130        let target = match last_bad {
131            // Unexplored territory: slow-start doubling.
132            None => (current * 2.0).max(current + 1.0),
133            Some(bad) => {
134                // Fast reclimb through the proven-safe zone, then
135                // cautious additive probing toward the old wall.
136                let safe = (bad * 0.75).max(floor);
137                let fast = (current * 1.5).min(safe);
138                if fast > current {
139                    fast
140                } else {
141                    current + (bad * 0.02).max(1.0)
142                }
143            }
144        }
145        .min(ceiling);
146        if target > current {
147            return Decision::Up(target);
148        }
149    }
150    Decision::Hold
151}
152
153/// The per-phase governor: window bookkeeping over the activity's
154/// cumulative attempt counters plus the resolved control handle.
155pub struct ThrottleGovernor {
156    control: Arc<dyn ErasedControl>,
157    control_name: String,
158    phase_name: String,
159    high: f64,
160    low: f64,
161    floor: f64,
162    ceiling: f64,
163    start: f64,
164    window: Duration,
165    window_start: Instant,
166    base_success: u64,
167    base_failure: u64,
168    /// The offer at which congestion was last observed. `None` =
169    /// unexplored (slow-start regime).
170    last_bad: Option<f64>,
171    /// Consecutive clean windows observed at-or-above `last_bad`;
172    /// the memory clears at [`MEMORY_CLEAR_STREAK`].
173    clean_streak: u32,
174}
175
176impl ThrottleGovernor {
177    /// Resolve the governor from the phase's declared spec against the
178    /// attached component (where `Activity::attach_component` declared
179    /// the controls). Returns `None` — with a logged warning, never
180    /// silently — when the named control is not declared (e.g.
181    /// `control: rate` on a phase without `rate:`).
182    ///
183    /// Construction immediately publishes the slow-start offer to the
184    /// control (push-on-set), so the phase OPENS at `start`, not at
185    /// the authored ceiling; the caller also reads
186    /// [`Self::initial_concurrency`] to spawn the fiber pool at the
187    /// same offer.
188    pub fn from_spec(
189        spec: &nmbrs_workload::model::ThrottleSpec,
190        component: Option<&Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>>,
191        phase_name: &str,
192        authored_concurrency: usize,
193        authored_rate: Option<f64>,
194    ) -> Option<Self> {
195        let Some(component) = component else {
196            gov_log!(
197                crate::observer::LogLevel::Warn,
198                "phase '{phase_name}': throttle: no component attached — \
199                 governor disabled"
200            );
201            return None;
202        };
203        let control = component
204            .read()
205            .unwrap_or_else(|e| e.into_inner())
206            .find_control_erased_up(&spec.control);
207        let Some(control) = control else {
208            gov_log!(
209                crate::observer::LogLevel::Warn,
210                "phase '{phase_name}': throttle: control '{}' is not \
211                 declared on this phase — governor disabled \
212                 (`control: rate` needs a `rate:` on the phase)",
213                spec.control
214            );
215            return None;
216        };
217        let ceiling = match spec.control.as_str() {
218            "rate" => authored_rate.unwrap_or(f64::INFINITY),
219            _ => authored_concurrency as f64,
220        };
221        let start = spec.start.unwrap_or(spec.floor).clamp(
222            spec.floor,
223            if ceiling.is_finite() {
224                ceiling
225            } else {
226                f64::MAX
227            },
228        );
229        let window = nmbrs_workload::magnitude::parse_magnitude(&spec.window)
230            .map(Duration::from_secs_f64)
231            .or_else(|| {
232                crate::timeval::parse_time_ms(&spec.window)
233                    .ok()
234                    .map(Duration::from_millis)
235            })
236            .unwrap_or(Duration::from_secs(2));
237        let low = spec.low.unwrap_or(spec.high / 5.0);
238        gov_log!(
239            crate::observer::LogLevel::Info,
240            "phase '{phase_name}': throttle: governing '{}' — slow-start \
241             at {} toward ceiling {} (floor {}); back off above {:.1}% \
242             windowed attempt failure ({}), recover below {:.1}%",
243            spec.control,
244            fmt_val(start),
245            fmt_val(ceiling),
246            fmt_val(spec.floor),
247            spec.high * 100.0,
248            spec.window,
249            low * 100.0
250        );
251        let governor = Self {
252            control,
253            control_name: spec.control.clone(),
254            phase_name: phase_name.to_string(),
255            high: spec.high,
256            low,
257            floor: spec.floor,
258            ceiling,
259            start,
260            window,
261            window_start: Instant::now(),
262            base_success: 0,
263            base_failure: 0,
264            last_bad: None,
265            clean_streak: 0,
266        };
267        // Publish the opening offer so the control's committed value
268        // reflects the slow-start from the first cycle. (For a
269        // concurrency governor the caller ALSO spawns the pool at
270        // `initial_concurrency`, so the offer and the pool agree.)
271        if start < ceiling {
272            governor.write(start);
273        }
274        Some(governor)
275    }
276
277    /// The fiber count the activity should OPEN with when this
278    /// governor walks `concurrency` — the slow-start offer, not the
279    /// authored ceiling. `None` when the governor walks `rate` (the
280    /// pool spawns at the authored concurrency; the rate limiter
281    /// carries the slow-start instead).
282    pub fn initial_concurrency(&self) -> Option<usize> {
283        (self.control_name == "concurrency").then_some((self.start.max(1.0)) as usize)
284    }
285
286    /// Called every activity-loop pass with the CUMULATIVE attempt
287    /// counters; acts at most once per window. The control write is
288    /// push-on-set (spawned, confirmed-apply) — the loop never blocks.
289    pub fn tick(&mut self, attempt_success: u64, attempt_failure: u64) {
290        if self.window_start.elapsed() < self.window {
291            return;
292        }
293        let d_success = attempt_success.saturating_sub(self.base_success);
294        let d_failure = attempt_failure.saturating_sub(self.base_failure);
295        self.window_start = Instant::now();
296        self.base_success = attempt_success;
297        self.base_failure = attempt_failure;
298
299        let d_total = d_success + d_failure;
300        if d_total == 0 {
301            return;
302        }
303        let frac = d_failure as f64 / d_total as f64;
304        // The committed value is the truth to step from — external
305        // writers (TUI, web, control_set) are honored, not fought.
306        let Some(current) = self.control.gauge_f64() else {
307            return;
308        };
309
310        // A SUSTAINED run of clean windows at-or-above the remembered
311        // congestion point means the target got healthier (caches
312        // warm, compactions done): clear the memory and resume the
313        // multiplicative climb. One clean window at a marginal
314        // congestion point is noise, not evidence — the additive
315        // probes keep stepping through the streak, so clearing means
316        // the target stayed clean across a rising run of offers.
317        if let Some(bad) = self.last_bad {
318            let evidence = frac < self.low && current >= bad;
319            let (streak, clears) = memory_clear_step(self.clean_streak, evidence);
320            self.clean_streak = streak;
321            if clears {
322                gov_log!(
323                    crate::observer::LogLevel::Info,
324                    "throttle: phase '{}': {} consecutive clean windows at or \
325                     above prior congestion point {} (now at {}) — memory \
326                     cleared, resuming climb",
327                    self.phase_name,
328                    MEMORY_CLEAR_STREAK,
329                    fmt_val(bad),
330                    fmt_val(current)
331                );
332                self.last_bad = None;
333            } else if evidence {
334                gov_log!(
335                    crate::observer::LogLevel::Debug,
336                    "throttle: phase '{}': clean window at {} ≥ prior \
337                     congestion point {} ({streak}/{} toward clearing memory)",
338                    self.phase_name,
339                    fmt_val(current),
340                    fmt_val(bad),
341                    MEMORY_CLEAR_STREAK
342                );
343            }
344        }
345
346        match decide(
347            frac,
348            current,
349            self.high,
350            self.low,
351            self.floor,
352            self.ceiling,
353            self.last_bad,
354        ) {
355            Decision::Down(target) => {
356                gov_log!(
357                    crate::observer::LogLevel::Warn,
358                    "throttle: phase '{}': windowed attempt failure {:.1}% \
359                     ({d_failure}/{d_total} over {:.1}s) > {:.1}% — {} {} → {}",
360                    self.phase_name,
361                    frac * 100.0,
362                    self.window.as_secs_f64(),
363                    self.high * 100.0,
364                    self.control_name,
365                    fmt_val(current),
366                    fmt_val(target)
367                );
368                self.last_bad = Some(current);
369                self.write(target);
370            }
371            Decision::Up(target) => {
372                let mode = match self.last_bad {
373                    None => "climbing",
374                    Some(bad) if target < bad * 0.75 => "reclimbing",
375                    Some(_) => "probing",
376                };
377                gov_log!(
378                    crate::observer::LogLevel::Info,
379                    "throttle: phase '{}': windowed attempt failure {:.1}% \
380                     < {:.1}% — {mode} {} {} → {} (ceiling {})",
381                    self.phase_name,
382                    frac * 100.0,
383                    self.low * 100.0,
384                    self.control_name,
385                    fmt_val(current),
386                    fmt_val(target),
387                    fmt_val(self.ceiling)
388                );
389                self.write(target);
390            }
391            Decision::Hold => {}
392        }
393    }
394
395    fn write(&self, target: f64) {
396        let control = self.control.clone();
397        let origin = ControlOrigin::Governor {
398            source: format!("throttle:{}", self.phase_name),
399        };
400        let name = self.control_name.clone();
401        let phase = self.phase_name.clone();
402        tokio::spawn(async move {
403            if let Err(e) = control.set_f64(target, origin).await {
404                gov_log!(
405                    crate::observer::LogLevel::Warn,
406                    "throttle: phase '{phase}': write {name}={target} \
407                     failed: {e}"
408                );
409            }
410        });
411    }
412}
413
414fn fmt_val(v: f64) -> String {
415    if v.is_infinite() {
416        "∞".to_string()
417    } else if (v.fract()).abs() < 1e-9 {
418        format!("{}", v as i64)
419    } else {
420        format!("{v:.1}")
421    }
422}
423
424#[cfg(test)]
425mod tests {
426    use super::*;
427
428    /// Severity-proportional back-off: total failure collapses fast,
429    /// a marginal breach trims gently, the floor contains the walk.
430    #[test]
431    fn backoff_scales_with_severity() {
432        // 100% failure → ×0.25.
433        assert_eq!(
434            decide(1.0, 100.0, 0.05, 0.01, 1.0, 100.0, None),
435            Decision::Down(25.0)
436        );
437        // Marginal breach (7%) → ×0.9 trim.
438        assert_eq!(
439            decide(0.07, 100.0, 0.05, 0.01, 1.0, 100.0, None),
440            Decision::Down(90.0)
441        );
442        // Mid-severity (50%) → ×0.5.
443        assert_eq!(
444            decide(0.5, 40.0, 0.05, 0.01, 1.0, 100.0, None),
445            Decision::Down(20.0)
446        );
447        // The floor contains the collapse; at the floor, hold.
448        assert_eq!(
449            decide(1.0, 5.0, 0.05, 0.01, 4.0, 100.0, None),
450            Decision::Down(4.0)
451        );
452        assert_eq!(
453            decide(1.0, 4.0, 0.05, 0.01, 4.0, 100.0, None),
454            Decision::Hold
455        );
456    }
457
458    /// Unexplored territory (slow-start): clean windows DOUBLE the
459    /// offer toward the ceiling — a robust target reaches full load
460    /// in log2(ceiling/start) windows with zero failures.
461    #[test]
462    fn slow_start_doubles_while_clean() {
463        assert_eq!(
464            decide(0.0, 1.0, 0.05, 0.01, 1.0, 100.0, None),
465            Decision::Up(2.0)
466        );
467        assert_eq!(
468            decide(0.0, 8.0, 0.05, 0.01, 1.0, 100.0, None),
469            Decision::Up(16.0)
470        );
471        // Ceiling contains the climb; at the ceiling, hold.
472        assert_eq!(
473            decide(0.0, 64.0, 0.05, 0.01, 1.0, 100.0, None),
474            Decision::Up(100.0)
475        );
476        assert_eq!(
477            decide(0.0, 100.0, 0.05, 0.01, 1.0, 100.0, None),
478            Decision::Hold
479        );
480    }
481
482    /// After congestion at `last_bad`, recovery climbs ×1.5 only
483    /// through the proven-safe zone (75% of last_bad), then probes
484    /// additively — never a multiplicative march back into the wall.
485    #[test]
486    fn congestion_memory_gates_the_reclimb() {
487        // Fast reclimb below the safe zone, capped at it.
488        assert_eq!(
489            decide(0.0, 10.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
490            Decision::Up(15.0)
491        );
492        assert_eq!(
493            decide(0.0, 24.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
494            Decision::Up(30.0)
495        ); // capped at 40*0.75
496        // At/above the safe zone: additive probing only.
497        assert_eq!(
498            decide(0.0, 30.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
499            Decision::Up(31.0)
500        ); // + max(1, 40*0.02)
501        // Large-scale controls probe proportionally (+2% of bad).
502        assert_eq!(
503            decide(0.0, 40_000.0, 0.05, 0.01, 1.0, 100_000.0, Some(50_000.0)),
504            Decision::Up(41_000.0)
505        );
506    }
507
508    /// Congestion memory clears only on a sustained streak of clean
509    /// windows at-or-above the congestion point; any window without
510    /// that evidence resets the streak.
511    #[test]
512    fn memory_clears_on_a_sustained_clean_streak() {
513        assert_eq!(memory_clear_step(0, true), (1, false));
514        assert_eq!(memory_clear_step(1, true), (2, false));
515        assert_eq!(memory_clear_step(2, true), (0, true));
516        // A breach, a dead-band hold, or a clean window still below
517        // the congestion point resets the streak.
518        assert_eq!(memory_clear_step(2, false), (0, false));
519        assert_eq!(memory_clear_step(0, false), (0, false));
520    }
521
522    /// The dead band between `low` and `high` holds — no hunting
523    /// around the operating point.
524    #[test]
525    fn dead_band_holds() {
526        assert_eq!(
527            decide(0.03, 50.0, 0.05, 0.01, 4.0, 100.0, None),
528            Decision::Hold
529        );
530        assert_eq!(
531            decide(0.03, 50.0, 0.05, 0.01, 4.0, 100.0, Some(60.0)),
532            Decision::Hold
533        );
534    }
535}