use crate::{Admission, Admit};
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub enum Ramp {
Hold,
Raise(usize),
Settled(usize),
}
#[derive(Clone, Debug)]
pub struct ConcurrencyRamp {
adm: Admission,
level: usize,
max: usize,
window_open_at: f64,
window_ends_at: f64,
bytes: u64,
counting_since: f64,
settled: Option<usize>,
held_rate: Option<f64>,
delivering: Option<usize>,
warm_deadline: Option<f64>,
}
const WINDOW_DELTAS: f64 = 3.0;
const SETTLE_DELTAS: f64 = 1.5;
const MIN_WINDOW_S: f64 = 0.25;
const WARM_DELTAS: f64 = 8.0;
const MIN_WARM_S: f64 = 1.0;
const MAX_WARM_S: f64 = 8.0;
const MAX_WINDOW_S: f64 = 0.6;
impl ConcurrencyRamp {
pub fn starting_at(min_gain_frac: f64, start: usize, max: usize) -> Self {
let mut r = Self::new(min_gain_frac, max);
r.level = start.clamp(1, r.max);
r
}
pub fn new(min_gain_frac: f64, max: usize) -> Self {
Self {
adm: Admission::new(min_gain_frac, max),
level: 1,
max: max.max(1),
window_open_at: 0.0,
window_ends_at: 0.0,
bytes: 0,
counting_since: 0.0,
settled: None,
held_rate: None,
delivering: None,
warm_deadline: None,
}
}
pub fn arm_warmup(&mut self, now: f64, delta: f64) {
self.delivering = Some(0);
self.warm_deadline = Some(now + (delta * WARM_DELTAS).clamp(MIN_WARM_S, MAX_WARM_S));
}
pub fn note_delivering(&mut self, n: usize) {
if self.delivering.is_some() {
self.delivering = Some(n);
}
}
pub fn start(&mut self, now: f64, delta: f64) {
let settle = (delta * SETTLE_DELTAS).clamp(MIN_WINDOW_S, MAX_WINDOW_S);
let window = (delta * WINDOW_DELTAS).clamp(MIN_WINDOW_S, MAX_WINDOW_S);
self.window_open_at = now + settle;
self.window_ends_at = self.window_open_at + window;
self.counting_since = self.window_open_at;
self.bytes = 0;
}
pub fn observe(&mut self, bytes: u64, now: f64) {
if now >= self.window_open_at {
self.bytes += bytes;
}
}
pub fn settled(&self) -> Option<usize> {
self.settled
}
pub fn level(&self) -> usize {
self.level
}
pub fn clamp_max(&mut self, max: usize) {
let max = max.max(1);
if max >= self.max {
return;
}
self.max = max;
self.adm.clamp_max(max);
if self.level > max {
self.level = max;
self.held_rate = None;
}
if let Some(n) = self.settled {
self.settled = Some(n.min(max));
}
}
pub fn poll(&mut self, now: f64, delta: f64) -> Ramp {
if let Some(n) = self.settled {
return Ramp::Settled(n);
}
if let Some(deadline) = self.warm_deadline {
let warm = self.delivering.unwrap_or(usize::MAX) >= self.level;
if warm || now >= deadline {
self.warm_deadline = None;
}
self.held_rate = None;
self.start(now, delta);
return Ramp::Hold;
}
if now < self.window_ends_at {
return Ramp::Hold;
}
let span = (now - self.counting_since).max(1e-3);
let rate = self.bytes as f64 / span;
if self.bytes == 0 {
self.start(now, delta);
return Ramp::Hold;
}
if let Some(first) = self.held_rate.take() {
let _ = first;
return self.decide(rate, now, delta);
}
self.held_rate = Some(rate);
self.start(now, delta);
Ramp::Hold
}
fn decide(&mut self, rate: f64, now: f64, delta: f64) -> Ramp {
let trace = std::env::var_os("HYDRA_RAMP_TRACE").is_some();
if trace {
eprintln!(
"ramp: level={} rate={:.0} B/s at t={:.2}s delta={:.3}",
self.level, rate, now, delta
);
}
match self.adm.observe_at(self.level, rate) {
Admit::Stop => {
let n = self.adm.settled().unwrap_or(self.level).clamp(1, self.max);
self.settled = Some(n);
self.level = n;
Ramp::Settled(n)
}
Admit::Add if self.level < self.max => {
self.level = (self.level * 2).min(self.max);
self.held_rate = None;
if self.delivering.is_some() {
self.warm_deadline =
Some(now + (delta * WARM_DELTAS).clamp(MIN_WARM_S, MAX_WARM_S));
}
self.start(now, delta);
Ramp::Raise(self.level)
}
Admit::Add => {
self.settled = Some(self.level);
Ramp::Settled(self.level)
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn run(sat: usize, max: usize, per_conn: f64, delta: f64) -> usize {
let mut r = ConcurrencyRamp::new(0.15, max);
let mut now = 0.0;
r.start(now, delta);
for _ in 0..(max * 40) {
let rate = per_conn * r.level().min(sat) as f64;
let step = 0.05;
now += step;
r.observe((rate * step) as u64, now);
if let Ramp::Settled(n) = r.poll(now, delta) {
return n;
}
}
r.level()
}
#[test]
fn a_clamped_ceiling_settles_the_search_there() {
let mut r = ConcurrencyRamp::new(0.15, 8);
let delta = 0.12;
let mut now = 0.0;
r.start(now, delta);
for _ in 0..400 {
let rate = 400e3 * r.level() as f64;
now += 0.05;
r.observe((rate * 0.05) as u64, now);
if r.level() >= 4 {
break;
}
let _ = r.poll(now, delta);
}
assert!(
r.level() >= 4,
"the ramp must have climbed before being clamped"
);
r.clamp_max(2);
assert_eq!(
r.level(),
2,
"the clamp must take the level down with the ceiling"
);
let mut settled = None;
for _ in 0..400 {
let rate = 400e3 * r.level() as f64;
now += 0.05;
r.observe((rate * 0.05) as u64, now);
if let Ramp::Settled(n) = r.poll(now, delta) {
settled = Some(n);
break;
}
}
assert_eq!(
settled,
Some(2),
"a search clamped to 2 on a path with headroom must settle at 2, \
not keep asking for connections the origin has refused"
);
}
#[test]
fn a_saturated_path_settles_low() {
let n = run(1, 8, 1.4e6, 0.12);
assert!(
n <= 2,
"settled at {n} connections on a path that saturates at 1; \
the extra connections are pure setup cost"
);
}
#[test]
fn a_path_with_headroom_ramps_up() {
let n = run(6, 8, 400e3, 0.12);
assert!(
n >= 4,
"settled at {n} connections on a path that scales to 6; \
the ramp is leaving throughput on the table"
);
}
#[test]
fn full_concurrency_is_reached_in_bounded_time() {
for &delta in &[0.01f64, 0.12, 0.5, 1.0, 5.0] {
let mut r = ConcurrencyRamp::new(0.15, 8);
let mut now = 0.0;
r.start(now, delta);
let mut reached_at = None;
while now < 30.0 {
now += 0.02;
r.observe((2e6 * r.level() as f64 * 0.02) as u64, now);
let out = r.poll(now, delta);
if r.level() >= 8 {
reached_at = Some(now);
break;
}
if let Ramp::Settled(_) = out {
break;
}
}
let t = reached_at.unwrap_or(f64::INFINITY);
assert!(
t <= 10.0,
"took {t:.1}s to reach 8 connections at delta={delta}: a search that \
costs more clock than the transfer saves is a net loss"
);
}
}
#[test]
fn the_ceiling_is_never_exceeded() {
for max in [1usize, 2, 4] {
let n = run(64, max, 400e3, 0.12);
assert!(n <= max, "settled at {n} above ceiling {max}");
}
}
#[test]
fn an_empty_window_does_not_settle_the_search() {
let mut r = ConcurrencyRamp::new(0.15, 8);
r.start(0.0, 0.12);
let out = r.poll(10.0, 0.12);
assert_eq!(out, Ramp::Hold, "silence must not be read as saturation");
assert!(r.settled().is_none());
}
fn run_warming(
sat: usize,
max: usize,
per_conn: f64,
delta: f64,
handshake: f64,
slow_start: f64,
gate: bool,
) -> (usize, f64) {
let mut r = ConcurrencyRamp::new(0.15, max);
let mut now = 0.0;
r.start(now, delta);
if gate {
r.arm_warmup(now, delta);
}
let mut opened = vec![0.0f64];
let mut hs = vec![0.0f64];
let step = 0.02;
while now < 60.0 {
now += step;
let mut live = 0usize;
let mut rate = 0.0;
for (i, &t0) in opened.iter().enumerate() {
let since = now - t0 - hs[i];
if since >= 0.0 {
live += 1;
rate += per_conn * (since / slow_start).clamp(0.0, 1.0);
}
}
rate = rate.min(per_conn * sat as f64);
r.observe((rate * step) as u64, now);
if gate {
r.note_delivering(live);
}
match r.poll(now, delta) {
Ramp::Raise(n) => {
while opened.len() < n {
opened.push(now);
hs.push(handshake);
}
}
Ramp::Settled(n) => return (n, now),
Ramp::Hold => {}
}
}
(r.level(), now)
}
#[test]
fn a_slow_handshake_defeats_the_ungated_search() {
let (n, _) = run_warming(8, 8, 2.4e6, 0.4, 2.5, 2.0, false);
assert_eq!(
n, 1,
"expected the ungated ramp to be fooled by a 2.5 s handshake; it settled \
at {n}, so this test no longer covers the defect it was written for"
);
}
#[test]
fn the_warm_up_gate_finds_the_headroom_a_slow_handshake_hides() {
let (n, t) = run_warming(8, 8, 2.4e6, 0.4, 2.5, 2.0, true);
assert!(
n >= 4,
"settled at {n} on a path that scales to 8: the gate did not restore the \
measurement"
);
assert!(
t <= 30.0,
"took {t:.1}s to settle at {n}: patience is not free and must be bounded"
);
}
#[test]
fn the_gate_gives_up_on_a_connection_that_never_delivers() {
let (n, t) = run_warming(8, 8, 2.4e6, 0.4, 1e6, 2.0, true);
assert!(
n <= 2,
"settled at {n} on a path where only one connection ever delivered"
);
assert!(
t < 60.0,
"the search never settled: the warm-up deadline is not bounding the wait"
);
}
}