use std::sync::Arc;
use std::time::{Duration, Instant};
use nmbrs_metrics::controls::{ControlOrigin, ErasedControl};
macro_rules! gov_log {
($level:expr, $($arg:tt)*) => {
crate::observer::log_tagged(
$level,
crate::observer::EventTag::in_flight(
crate::observer::EventCategory::Throttle),
&format!($($arg)*),
)
};
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub enum Decision {
Down(f64),
Up(f64),
Hold,
}
pub const MEMORY_CLEAR_STREAK: u32 = 3;
pub fn memory_clear_step(streak: u32, evidence: bool) -> (u32, bool) {
if !evidence {
return (0, false);
}
let streak = streak + 1;
if streak >= MEMORY_CLEAR_STREAK {
(0, true)
} else {
(streak, false)
}
}
pub fn decide(
frac: f64,
current: f64,
high: f64,
low: f64,
floor: f64,
ceiling: f64,
last_bad: Option<f64>,
) -> Decision {
if frac > high {
let target = (current * (1.0 - frac).clamp(0.25, 0.9)).max(floor);
if target < current {
return Decision::Down(target);
}
return Decision::Hold;
}
if frac < low && current < ceiling {
let target = match last_bad {
None => (current * 2.0).max(current + 1.0),
Some(bad) => {
let safe = (bad * 0.75).max(floor);
let fast = (current * 1.5).min(safe);
if fast > current {
fast
} else {
current + (bad * 0.02).max(1.0)
}
}
}
.min(ceiling);
if target > current {
return Decision::Up(target);
}
}
Decision::Hold
}
pub struct ThrottleGovernor {
control: Arc<dyn ErasedControl>,
control_name: String,
phase_name: String,
high: f64,
low: f64,
floor: f64,
ceiling: f64,
start: f64,
window: Duration,
window_start: Instant,
base_success: u64,
base_failure: u64,
last_bad: Option<f64>,
clean_streak: u32,
}
impl ThrottleGovernor {
pub fn from_spec(
spec: &nmbrs_workload::model::ThrottleSpec,
component: Option<&Arc<std::sync::RwLock<nmbrs_metrics::component::Component>>>,
phase_name: &str,
authored_concurrency: usize,
authored_rate: Option<f64>,
) -> Option<Self> {
let Some(component) = component else {
gov_log!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': throttle: no component attached — \
governor disabled"
);
return None;
};
let control = component
.read()
.unwrap_or_else(|e| e.into_inner())
.find_control_erased_up(&spec.control);
let Some(control) = control else {
gov_log!(
crate::observer::LogLevel::Warn,
"phase '{phase_name}': throttle: control '{}' is not \
declared on this phase — governor disabled \
(`control: rate` needs a `rate:` on the phase)",
spec.control
);
return None;
};
let ceiling = match spec.control.as_str() {
"rate" => authored_rate.unwrap_or(f64::INFINITY),
_ => authored_concurrency as f64,
};
let start = spec.start.unwrap_or(spec.floor).clamp(
spec.floor,
if ceiling.is_finite() {
ceiling
} else {
f64::MAX
},
);
let window = nmbrs_workload::magnitude::parse_magnitude(&spec.window)
.map(Duration::from_secs_f64)
.or_else(|| {
crate::timeval::parse_time_ms(&spec.window)
.ok()
.map(Duration::from_millis)
})
.unwrap_or(Duration::from_secs(2));
let low = spec.low.unwrap_or(spec.high / 5.0);
gov_log!(
crate::observer::LogLevel::Info,
"phase '{phase_name}': throttle: governing '{}' — slow-start \
at {} toward ceiling {} (floor {}); back off above {:.1}% \
windowed attempt failure ({}), recover below {:.1}%",
spec.control,
fmt_val(start),
fmt_val(ceiling),
fmt_val(spec.floor),
spec.high * 100.0,
spec.window,
low * 100.0
);
let governor = Self {
control,
control_name: spec.control.clone(),
phase_name: phase_name.to_string(),
high: spec.high,
low,
floor: spec.floor,
ceiling,
start,
window,
window_start: Instant::now(),
base_success: 0,
base_failure: 0,
last_bad: None,
clean_streak: 0,
};
if start < ceiling {
governor.write(start);
}
Some(governor)
}
pub fn initial_concurrency(&self) -> Option<usize> {
(self.control_name == "concurrency").then_some((self.start.max(1.0)) as usize)
}
pub fn tick(&mut self, attempt_success: u64, attempt_failure: u64) {
if self.window_start.elapsed() < self.window {
return;
}
let d_success = attempt_success.saturating_sub(self.base_success);
let d_failure = attempt_failure.saturating_sub(self.base_failure);
self.window_start = Instant::now();
self.base_success = attempt_success;
self.base_failure = attempt_failure;
let d_total = d_success + d_failure;
if d_total == 0 {
return;
}
let frac = d_failure as f64 / d_total as f64;
let Some(current) = self.control.gauge_f64() else {
return;
};
if let Some(bad) = self.last_bad {
let evidence = frac < self.low && current >= bad;
let (streak, clears) = memory_clear_step(self.clean_streak, evidence);
self.clean_streak = streak;
if clears {
gov_log!(
crate::observer::LogLevel::Info,
"throttle: phase '{}': {} consecutive clean windows at or \
above prior congestion point {} (now at {}) — memory \
cleared, resuming climb",
self.phase_name,
MEMORY_CLEAR_STREAK,
fmt_val(bad),
fmt_val(current)
);
self.last_bad = None;
} else if evidence {
gov_log!(
crate::observer::LogLevel::Debug,
"throttle: phase '{}': clean window at {} ≥ prior \
congestion point {} ({streak}/{} toward clearing memory)",
self.phase_name,
fmt_val(current),
fmt_val(bad),
MEMORY_CLEAR_STREAK
);
}
}
match decide(
frac,
current,
self.high,
self.low,
self.floor,
self.ceiling,
self.last_bad,
) {
Decision::Down(target) => {
gov_log!(
crate::observer::LogLevel::Warn,
"throttle: phase '{}': windowed attempt failure {:.1}% \
({d_failure}/{d_total} over {:.1}s) > {:.1}% — {} {} → {}",
self.phase_name,
frac * 100.0,
self.window.as_secs_f64(),
self.high * 100.0,
self.control_name,
fmt_val(current),
fmt_val(target)
);
self.last_bad = Some(current);
self.write(target);
}
Decision::Up(target) => {
let mode = match self.last_bad {
None => "climbing",
Some(bad) if target < bad * 0.75 => "reclimbing",
Some(_) => "probing",
};
gov_log!(
crate::observer::LogLevel::Info,
"throttle: phase '{}': windowed attempt failure {:.1}% \
< {:.1}% — {mode} {} {} → {} (ceiling {})",
self.phase_name,
frac * 100.0,
self.low * 100.0,
self.control_name,
fmt_val(current),
fmt_val(target),
fmt_val(self.ceiling)
);
self.write(target);
}
Decision::Hold => {}
}
}
fn write(&self, target: f64) {
let control = self.control.clone();
let origin = ControlOrigin::Governor {
source: format!("throttle:{}", self.phase_name),
};
let name = self.control_name.clone();
let phase = self.phase_name.clone();
tokio::spawn(async move {
if let Err(e) = control.set_f64(target, origin).await {
gov_log!(
crate::observer::LogLevel::Warn,
"throttle: phase '{phase}': write {name}={target} \
failed: {e}"
);
}
});
}
}
fn fmt_val(v: f64) -> String {
if v.is_infinite() {
"∞".to_string()
} else if (v.fract()).abs() < 1e-9 {
format!("{}", v as i64)
} else {
format!("{v:.1}")
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn backoff_scales_with_severity() {
assert_eq!(
decide(1.0, 100.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Down(25.0)
);
assert_eq!(
decide(0.07, 100.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Down(90.0)
);
assert_eq!(
decide(0.5, 40.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Down(20.0)
);
assert_eq!(
decide(1.0, 5.0, 0.05, 0.01, 4.0, 100.0, None),
Decision::Down(4.0)
);
assert_eq!(
decide(1.0, 4.0, 0.05, 0.01, 4.0, 100.0, None),
Decision::Hold
);
}
#[test]
fn slow_start_doubles_while_clean() {
assert_eq!(
decide(0.0, 1.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Up(2.0)
);
assert_eq!(
decide(0.0, 8.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Up(16.0)
);
assert_eq!(
decide(0.0, 64.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Up(100.0)
);
assert_eq!(
decide(0.0, 100.0, 0.05, 0.01, 1.0, 100.0, None),
Decision::Hold
);
}
#[test]
fn congestion_memory_gates_the_reclimb() {
assert_eq!(
decide(0.0, 10.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
Decision::Up(15.0)
);
assert_eq!(
decide(0.0, 24.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
Decision::Up(30.0)
); assert_eq!(
decide(0.0, 30.0, 0.05, 0.01, 1.0, 100.0, Some(40.0)),
Decision::Up(31.0)
); assert_eq!(
decide(0.0, 40_000.0, 0.05, 0.01, 1.0, 100_000.0, Some(50_000.0)),
Decision::Up(41_000.0)
);
}
#[test]
fn memory_clears_on_a_sustained_clean_streak() {
assert_eq!(memory_clear_step(0, true), (1, false));
assert_eq!(memory_clear_step(1, true), (2, false));
assert_eq!(memory_clear_step(2, true), (0, true));
assert_eq!(memory_clear_step(2, false), (0, false));
assert_eq!(memory_clear_step(0, false), (0, false));
}
#[test]
fn dead_band_holds() {
assert_eq!(
decide(0.03, 50.0, 0.05, 0.01, 4.0, 100.0, None),
Decision::Hold
);
assert_eq!(
decide(0.03, 50.0, 0.05, 0.01, 4.0, 100.0, Some(60.0)),
Decision::Hold
);
}
}