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>,
}
const WINDOW_DELTAS: f64 = 3.0;
const SETTLE_DELTAS: f64 = 1.5;
const MIN_WINDOW_S: f64 = 0.25;
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,
}
}
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 poll(&mut self, now: f64, delta: f64) -> Ramp {
if let Some(n) = self.settled {
return Ramp::Settled(n);
}
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;
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_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());
}
}