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}