use std::io::{BufRead, BufReader, Read, Write};
use std::net::TcpStream;
use std::path::PathBuf;
use std::process::{Child, Command, Stdio};
use std::sync::atomic::{AtomicU16, Ordering};
use std::time::{Duration, Instant};
#[derive(Debug, Clone)]
pub struct ThroughputCell {
pub policy: &'static str, pub concurrency: u32, pub aggregate_tokens_per_sec: f64,
pub ttft_p50_ms: f64,
pub ttft_p95_ms: f64,
pub per_slot_tokens_per_sec: f64,
pub rejected_429_count: u32,
}
impl ThroughputCell {
pub fn synthetic_for_smoke() -> Self {
ThroughputCell {
policy: "fifo_serial",
concurrency: 1,
aggregate_tokens_per_sec: 0.0,
ttft_p50_ms: 0.0,
ttft_p95_ms: 0.0,
per_slot_tokens_per_sec: 0.0,
rejected_429_count: 0,
}
}
}
pub fn render_report(cells: &[ThroughputCell]) -> String {
let mut s =
String::from("| policy | N | agg tok/s | TTFT p50 | TTFT p95 | per-slot tok/s | 429s |\n");
s.push_str("|--------|---|-----------|----------|----------|----------------|------|\n");
for c in cells {
s.push_str(&format!(
"| {} | {} | {:.1} | {:.1} | {:.1} | {:.1} | {} |\n",
c.policy,
c.concurrency,
c.aggregate_tokens_per_sec,
c.ttft_p50_ms,
c.ttft_p95_ms,
c.per_slot_tokens_per_sec,
c.rejected_429_count
));
}
s
}
#[derive(Debug, Clone)]
pub struct ThroughputCellStable {
pub policy: &'static str,
pub concurrency: u32,
pub rep_count: u32,
pub aggregate_tokens_per_sec_median: f64,
pub aggregate_tokens_per_sec_min: f64,
pub aggregate_tokens_per_sec_max: f64,
pub aggregate_tokens_per_sec_sigma_pct: f64,
pub ttft_p50_ms_median: f64,
pub ttft_p95_ms_median: f64,
pub per_slot_tokens_per_sec_median: f64,
pub rejected_429_count_total: u32,
}
impl ThroughputCellStable {
pub fn from_reps(cells: Vec<ThroughputCell>) -> Self {
assert!(
!cells.is_empty(),
"ThroughputCellStable::from_reps requires ≥1 cell"
);
let policy = cells[0].policy;
let concurrency = cells[0].concurrency;
for c in &cells {
assert_eq!(c.policy, policy, "from_reps: mixed policies across reps");
assert_eq!(
c.concurrency, concurrency,
"from_reps: mixed concurrency across reps"
);
}
let rep_count = cells.len() as u32;
let aggregates: Vec<f64> = cells.iter().map(|c| c.aggregate_tokens_per_sec).collect();
let ttft_p50s: Vec<f64> = cells.iter().map(|c| c.ttft_p50_ms).collect();
let ttft_p95s: Vec<f64> = cells.iter().map(|c| c.ttft_p95_ms).collect();
let per_slots: Vec<f64> = cells.iter().map(|c| c.per_slot_tokens_per_sec).collect();
let rejected_total: u32 = cells.iter().map(|c| c.rejected_429_count).sum();
let agg_median = median_f64(&aggregates);
let agg_min = aggregates.iter().cloned().fold(f64::INFINITY, f64::min);
let agg_max = aggregates.iter().cloned().fold(f64::NEG_INFINITY, f64::max);
let sigma_pct = if agg_median > 0.0 {
(agg_max - agg_min) / agg_median * 100.0
} else {
0.0
};
ThroughputCellStable {
policy,
concurrency,
rep_count,
aggregate_tokens_per_sec_median: agg_median,
aggregate_tokens_per_sec_min: agg_min,
aggregate_tokens_per_sec_max: agg_max,
aggregate_tokens_per_sec_sigma_pct: sigma_pct,
ttft_p50_ms_median: median_f64(&ttft_p50s),
ttft_p95_ms_median: median_f64(&ttft_p95s),
per_slot_tokens_per_sec_median: median_f64(&per_slots),
rejected_429_count_total: rejected_total,
}
}
}
fn median_f64(values: &[f64]) -> f64 {
if values.is_empty() {
return 0.0;
}
let mut sorted = values.to_vec();
sorted.sort_by(|a, b| a.partial_cmp(b).unwrap_or(std::cmp::Ordering::Equal));
sorted[sorted.len() / 2]
}
pub fn render_report_stable(cells: &[ThroughputCellStable]) -> String {
let mut s = String::from(
"| policy | N | reps | agg tok/s median | agg min | agg max | sigma_pct | TTFT p50 | TTFT p95 | per-slot tok/s | 429s total |\n",
);
s.push_str("|--------|---|------|------------------|---------|---------|-----------|----------|----------|----------------|------------|\n");
for c in cells {
s.push_str(&format!(
"| {} | {} | {} | {:.1} | {:.1} | {:.1} | {:.1}% | {:.1} | {:.1} | {:.1} | {} |\n",
c.policy,
c.concurrency,
c.rep_count,
c.aggregate_tokens_per_sec_median,
c.aggregate_tokens_per_sec_min,
c.aggregate_tokens_per_sec_max,
c.aggregate_tokens_per_sec_sigma_pct,
c.ttft_p50_ms_median,
c.ttft_p95_ms_median,
c.per_slot_tokens_per_sec_median,
c.rejected_429_count_total,
));
}
s
}
pub const REPS: usize = 3;
pub const STABILITY_SIGMA_PCT_THRESHOLD: f64 = 20.0;
#[derive(Debug, PartialEq)]
pub enum Ac4Outcome {
Deferred,
Misconfigured,
StabilityBlocked,
Passed {
aggregate_ratio: f64,
ttft_ratio: f64,
},
Failed {
aggregate_ratio: f64,
ttft_ratio: f64,
which: &'static str,
},
}
pub fn ac4_outcome(
fifo_n4: Option<&ThroughputCellStable>,
inflight_n4: Option<&ThroughputCellStable>,
fifo_n1: Option<&ThroughputCellStable>,
) -> Ac4Outcome {
let (f, i) = match (fifo_n4, inflight_n4) {
(Some(f), Some(i)) => (f, i),
_ => return Ac4Outcome::Deferred,
};
if f.aggregate_tokens_per_sec_sigma_pct > STABILITY_SIGMA_PCT_THRESHOLD
|| i.aggregate_tokens_per_sec_sigma_pct > STABILITY_SIGMA_PCT_THRESHOLD
{
return Ac4Outcome::StabilityBlocked;
}
let Some(base) = fifo_n1 else {
return Ac4Outcome::Misconfigured;
};
let aggregate_ratio =
i.aggregate_tokens_per_sec_median / f.aggregate_tokens_per_sec_median.max(1e-6);
let ttft_ratio = i.ttft_p95_ms_median / base.ttft_p95_ms_median.max(1e-6);
if aggregate_ratio < 1.5 {
return Ac4Outcome::Failed {
aggregate_ratio,
ttft_ratio,
which: "aggregate",
};
}
if ttft_ratio > 2.0 {
return Ac4Outcome::Failed {
aggregate_ratio,
ttft_ratio,
which: "ttft",
};
}
Ac4Outcome::Passed {
aggregate_ratio,
ttft_ratio,
}
}
fn hf2q_binary_path() -> PathBuf {
if let Some(p) = std::env::var_os("CARGO_BIN_EXE_hf2q") {
return PathBuf::from(p);
}
let target = std::env::var("CARGO_TARGET_DIR")
.map(PathBuf::from)
.unwrap_or_else(|_| PathBuf::from("/opt/hf2q/target"));
target.join("release").join("hf2q")
}
#[test]
fn binary_is_locatable_and_runs_version() {
let bin = hf2q_binary_path();
if !bin.exists() {
eprintln!("[cb-throughput] skipping: {} not built", bin.display());
return;
}
let out = Command::new(&bin).arg("--version").output();
match out {
Ok(o) if o.status.success() => {}
Ok(o) => panic!("hf2q --version failed: {:?}", o),
Err(e) => panic!("failed to run hf2q --version: {}", e),
}
}
#[test]
fn throughput_cell_synthetic_round_trips_through_report() {
let cells = vec![ThroughputCell::synthetic_for_smoke()];
let report = render_report(&cells);
assert!(report.contains("| fifo_serial | 1 |"), "report: {}", report);
assert!(report.contains("agg tok/s"), "header missing");
assert!(report.contains("|--------|"), "separator missing");
}
#[test]
fn render_report_empty_returns_header_only() {
let report = render_report(&[]);
let lines: Vec<&str> = report.lines().collect();
assert_eq!(lines.len(), 2, "header + separator only, got: {:?}", lines);
}
#[test]
fn d3_median_f64_odd_length_returns_middle_sample() {
let v = vec![100.0, 50.0, 200.0];
assert_eq!(median_f64(&v), 100.0, "sorted=[50,100,200], middle=100");
}
#[test]
fn d3_median_f64_empty_returns_zero() {
assert_eq!(median_f64(&[]), 0.0);
}
#[test]
fn d3_stable_from_reps_aggregates_median_min_max_sigma() {
let cells = vec![
ThroughputCell {
policy: "fifo_serial",
concurrency: 4,
aggregate_tokens_per_sec: 100.0,
ttft_p50_ms: 10.0,
ttft_p95_ms: 20.0,
per_slot_tokens_per_sec: 25.0,
rejected_429_count: 1,
},
ThroughputCell {
policy: "fifo_serial",
concurrency: 4,
aggregate_tokens_per_sec: 110.0,
ttft_p50_ms: 12.0,
ttft_p95_ms: 22.0,
per_slot_tokens_per_sec: 27.5,
rejected_429_count: 0,
},
ThroughputCell {
policy: "fifo_serial",
concurrency: 4,
aggregate_tokens_per_sec: 105.0,
ttft_p50_ms: 11.0,
ttft_p95_ms: 21.0,
per_slot_tokens_per_sec: 26.25,
rejected_429_count: 2,
},
];
let stable = ThroughputCellStable::from_reps(cells);
assert_eq!(stable.policy, "fifo_serial");
assert_eq!(stable.concurrency, 4);
assert_eq!(stable.rep_count, 3);
assert!((stable.aggregate_tokens_per_sec_median - 105.0).abs() < 1e-9);
assert!((stable.aggregate_tokens_per_sec_min - 100.0).abs() < 1e-9);
assert!((stable.aggregate_tokens_per_sec_max - 110.0).abs() < 1e-9);
assert!(
(stable.aggregate_tokens_per_sec_sigma_pct - 9.523_809_523_8).abs() < 1e-6,
"sigma_pct={}",
stable.aggregate_tokens_per_sec_sigma_pct,
);
assert!((stable.ttft_p50_ms_median - 11.0).abs() < 1e-9);
assert!((stable.ttft_p95_ms_median - 21.0).abs() < 1e-9);
assert!((stable.per_slot_tokens_per_sec_median - 26.25).abs() < 1e-9);
assert_eq!(stable.rejected_429_count_total, 3);
}
#[test]
#[should_panic(expected = "from_reps: mixed policies")]
fn d3_stable_from_reps_rejects_mixed_policies() {
let cells = vec![
ThroughputCell {
policy: "fifo_serial",
concurrency: 4,
aggregate_tokens_per_sec: 100.0,
ttft_p50_ms: 10.0,
ttft_p95_ms: 20.0,
per_slot_tokens_per_sec: 25.0,
rejected_429_count: 0,
},
ThroughputCell {
policy: "inflight_batched",
concurrency: 4,
aggregate_tokens_per_sec: 200.0,
ttft_p50_ms: 10.0,
ttft_p95_ms: 20.0,
per_slot_tokens_per_sec: 50.0,
rejected_429_count: 0,
},
];
let _ = ThroughputCellStable::from_reps(cells);
}
#[test]
#[should_panic(expected = "from_reps: mixed concurrency")]
fn d3_stable_from_reps_rejects_mixed_concurrency() {
let cells = vec![
ThroughputCell {
policy: "fifo_serial",
concurrency: 4,
aggregate_tokens_per_sec: 100.0,
ttft_p50_ms: 10.0,
ttft_p95_ms: 20.0,
per_slot_tokens_per_sec: 25.0,
rejected_429_count: 0,
},
ThroughputCell {
policy: "fifo_serial",
concurrency: 8,
aggregate_tokens_per_sec: 200.0,
ttft_p50_ms: 10.0,
ttft_p95_ms: 20.0,
per_slot_tokens_per_sec: 25.0,
rejected_429_count: 0,
},
];
let _ = ThroughputCellStable::from_reps(cells);
}
#[test]
fn d3_stable_from_reps_zero_median_yields_zero_sigma_pct() {
let cells = vec![
ThroughputCell {
policy: "fifo_serial",
concurrency: 1,
aggregate_tokens_per_sec: 0.0,
ttft_p50_ms: 0.0,
ttft_p95_ms: 0.0,
per_slot_tokens_per_sec: 0.0,
rejected_429_count: 0,
},
ThroughputCell {
policy: "fifo_serial",
concurrency: 1,
aggregate_tokens_per_sec: 0.0,
ttft_p50_ms: 0.0,
ttft_p95_ms: 0.0,
per_slot_tokens_per_sec: 0.0,
rejected_429_count: 0,
},
ThroughputCell {
policy: "fifo_serial",
concurrency: 1,
aggregate_tokens_per_sec: 0.0,
ttft_p50_ms: 0.0,
ttft_p95_ms: 0.0,
per_slot_tokens_per_sec: 0.0,
rejected_429_count: 0,
},
];
let stable = ThroughputCellStable::from_reps(cells);
assert_eq!(stable.aggregate_tokens_per_sec_sigma_pct, 0.0);
}
#[test]
fn d3_render_report_stable_emits_header_and_sigma_column() {
let cells = vec![ThroughputCellStable {
policy: "fifo_serial",
concurrency: 4,
rep_count: 3,
aggregate_tokens_per_sec_median: 105.0,
aggregate_tokens_per_sec_min: 100.0,
aggregate_tokens_per_sec_max: 110.0,
aggregate_tokens_per_sec_sigma_pct: 9.52,
ttft_p50_ms_median: 11.0,
ttft_p95_ms_median: 21.0,
per_slot_tokens_per_sec_median: 26.25,
rejected_429_count_total: 3,
}];
let report = render_report_stable(&cells);
assert!(
report.contains("sigma_pct"),
"header missing sigma_pct: {report}"
);
assert!(
report.contains("agg tok/s median"),
"header missing median: {report}"
);
assert!(
report.contains("| fifo_serial | 4 | 3 |"),
"data row missing: {report}"
);
assert!(
report.contains("9.5%"),
"sigma_pct formatted row missing: {report}"
);
}
#[test]
fn d3_stability_threshold_default_is_twenty_pct() {
assert!(
(STABILITY_SIGMA_PCT_THRESHOLD - 20.0).abs() < 1e-9,
"STABILITY_SIGMA_PCT_THRESHOLD changed from 20.0 — update ADR §6.1.15 too",
);
assert_eq!(REPS, 3, "REPS changed from 3 — update ADR §6.1.15 too");
}
fn stable_cell(
policy: &'static str,
concurrency: u32,
aggregate_median: f64,
sigma_pct: f64,
ttft_p95_median: f64,
) -> ThroughputCellStable {
ThroughputCellStable {
policy,
concurrency,
rep_count: 3,
aggregate_tokens_per_sec_median: aggregate_median,
aggregate_tokens_per_sec_min: aggregate_median * 0.95,
aggregate_tokens_per_sec_max: aggregate_median * 1.05,
aggregate_tokens_per_sec_sigma_pct: sigma_pct,
ttft_p50_ms_median: ttft_p95_median * 0.5,
ttft_p95_ms_median: ttft_p95_median,
per_slot_tokens_per_sec_median: aggregate_median / f64::from(concurrency.max(1)),
rejected_429_count_total: 0,
}
}
#[test]
fn ac4_outcome_missing_inflight_n4_returns_deferred() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(Some(&f4), None, Some(&f1));
assert_eq!(
out,
Ac4Outcome::Deferred,
"Deferred when InflightBatched N=4 is absent (Phase C2c/C2d gated)"
);
}
#[test]
fn ac4_outcome_missing_fifo_n4_returns_deferred() {
let i4 = stable_cell("inflight_batched", 4, 200.0, 5.0, 70.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(None, Some(&i4), Some(&f1));
assert_eq!(
out,
Ac4Outcome::Deferred,
"Deferred when FifoSerial N=4 is absent"
);
}
#[test]
fn ac4_outcome_both_n4_present_but_missing_n1_returns_misconfigured() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let i4 = stable_cell("inflight_batched", 4, 200.0, 5.0, 70.0);
let out = ac4_outcome(Some(&f4), Some(&i4), None);
assert_eq!(
out,
Ac4Outcome::Misconfigured,
"BOTH N=4 cells + missing N=1 baseline ⇒ Misconfigured (hard error, \
not silent skip — pre-iter-A5b [ac-4 PARTIAL] would have let a \
TTFT regression pass)"
);
}
#[test]
fn ac4_outcome_stability_blocked_when_sigma_pct_above_threshold() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 25.0, 50.0); let i4 = stable_cell("inflight_batched", 4, 200.0, 5.0, 70.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(Some(&f4), Some(&i4), Some(&f1));
assert_eq!(
out,
Ac4Outcome::StabilityBlocked,
"fifo_serial sigma_pct > threshold ⇒ StabilityBlocked"
);
let f4_ok = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let i4_noisy = stable_cell("inflight_batched", 4, 200.0, 30.0, 70.0); let out2 = ac4_outcome(Some(&f4_ok), Some(&i4_noisy), Some(&f1));
assert_eq!(
out2,
Ac4Outcome::StabilityBlocked,
"inflight_batched sigma_pct > threshold ⇒ StabilityBlocked"
);
}
#[test]
fn ac4_outcome_passed_when_aggregate_above_1_5x_and_ttft_under_2x() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let i4 = stable_cell("inflight_batched", 4, 200.0, 5.0, 70.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(Some(&f4), Some(&i4), Some(&f1));
match out {
Ac4Outcome::Passed {
aggregate_ratio,
ttft_ratio,
} => {
assert!((aggregate_ratio - 2.0).abs() < 1e-6);
assert!((ttft_ratio - 1.4).abs() < 1e-6);
}
other => panic!("expected Passed, got {other:?}"),
}
}
#[test]
fn ac4_outcome_failed_aggregate_when_ratio_below_1_5x() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let i4 = stable_cell("inflight_batched", 4, 140.0, 5.0, 70.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(Some(&f4), Some(&i4), Some(&f1));
match out {
Ac4Outcome::Failed {
aggregate_ratio,
which,
..
} => {
assert!((aggregate_ratio - 1.4).abs() < 1e-6);
assert_eq!(
which, "aggregate",
"aggregate ratio failure surfaces `which='aggregate'`"
);
}
other => panic!("expected Failed/aggregate, got {other:?}"),
}
}
#[test]
fn ac4_outcome_failed_ttft_when_ratio_above_2x() {
let f4 = stable_cell("fifo_serial", 4, 100.0, 5.0, 50.0);
let i4 = stable_cell("inflight_batched", 4, 200.0, 5.0, 150.0);
let f1 = stable_cell("fifo_serial", 1, 50.0, 5.0, 50.0);
let out = ac4_outcome(Some(&f4), Some(&i4), Some(&f1));
match out {
Ac4Outcome::Failed {
ttft_ratio, which, ..
} => {
assert!((ttft_ratio - 3.0).abs() < 1e-6);
assert_eq!(which, "ttft", "ttft ratio failure surfaces `which='ttft'`");
}
other => panic!("expected Failed/ttft, got {other:?}"),
}
}
#[test]
fn render_report_two_cells_emits_two_data_rows() {
let cells = vec![
ThroughputCell {
policy: "fifo_serial",
concurrency: 1,
aggregate_tokens_per_sec: 100.0,
ttft_p50_ms: 50.0,
ttft_p95_ms: 80.0,
per_slot_tokens_per_sec: 100.0,
rejected_429_count: 0,
},
ThroughputCell {
policy: "inflight_batched",
concurrency: 4,
aggregate_tokens_per_sec: 300.0,
ttft_p50_ms: 75.0,
ttft_p95_ms: 120.0,
per_slot_tokens_per_sec: 75.0,
rejected_429_count: 2,
},
];
let report = render_report(&cells);
let data_rows = report
.lines()
.filter(|l| l.starts_with("|") && !l.contains("---"))
.count();
assert_eq!(data_rows, 3, "expected header + 2 data, got: {}", report);
}
const READYZ_BUDGET_SECS: u64 = 600;
const STREAM_BUDGET_SECS: u64 = 120;
static PORT_COUNTER: AtomicU16 = AtomicU16::new(0);
fn next_port() -> u16 {
let base: u16 = std::env::var("HF2Q_CB_THROUGHPUT_PORT_BASE")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(52441);
base + PORT_COUNTER.fetch_add(1, Ordering::SeqCst)
}
struct BenchServer {
child: Child,
port: u16,
}
impl BenchServer {
fn spawn(gguf: &str, policy: &str, max_slots: u32, port: u16) -> std::io::Result<Self> {
let bin = hf2q_binary_path();
let mut cmd = Command::new(bin);
let policy_cli = policy.replace('_', "-");
cmd.args([
"serve",
"--model",
gguf,
"--host",
"127.0.0.1",
"--port",
&port.to_string(),
"--scheduler",
&policy_cli,
]);
if policy == "inflight_batched" {
cmd.args(["--max-slots", &max_slots.to_string()]);
}
let child = cmd.stdout(Stdio::piped()).stderr(Stdio::piped()).spawn()?;
Ok(Self { child, port })
}
}
impl Drop for BenchServer {
fn drop(&mut self) {
let _ = self.child.kill();
let _ = self.child.wait();
}
}
fn http_get_status(port: u16, path: &str) -> std::io::Result<u16> {
let mut s = TcpStream::connect_timeout(
&format!("127.0.0.1:{port}")
.parse()
.map_err(std::io::Error::other)?,
Duration::from_secs(5),
)?;
s.set_read_timeout(Some(Duration::from_secs(5)))?;
s.write_all(
format!("GET {path} HTTP/1.1\r\nHost: 127.0.0.1:{port}\r\nConnection: close\r\n\r\n")
.as_bytes(),
)?;
let mut head = [0u8; 64];
let n = s.read(&mut head)?;
let head_s = std::str::from_utf8(&head[..n]).unwrap_or("");
let code = head_s
.split_whitespace()
.nth(1)
.and_then(|s| s.parse::<u16>().ok())
.ok_or_else(|| std::io::Error::other(format!("malformed HTTP status line: {head_s:?}")))?;
Ok(code)
}
fn wait_for_readyz(server: &mut BenchServer) -> Result<(), String> {
let started = Instant::now();
let mut last_err: Option<String> = None;
while started.elapsed().as_secs() < READYZ_BUDGET_SECS {
if let Ok(Some(status)) = server.child.try_wait() {
let mut stderr_tail = String::new();
if let Some(mut e) = server.child.stderr.take() {
let mut buf = Vec::new();
let _ = e.read_to_end(&mut buf);
let s = String::from_utf8_lossy(&buf);
let lines: Vec<&str> = s.lines().collect();
let tail = lines.iter().rev().take(15).rev();
stderr_tail = tail.copied().collect::<Vec<_>>().join("\n");
}
return Err(format!(
"subprocess exited before /readyz=200 (status={status:?})\n\
--- stderr tail ---\n{stderr_tail}\n--- end stderr ---"
));
}
match http_get_status(server.port, "/readyz") {
Ok(200) => return Ok(()),
Ok(code) => last_err = Some(format!("status={code}")),
Err(e) => last_err = Some(format!("transport: {e}")),
}
std::thread::sleep(Duration::from_secs(2));
}
Err(format!(
"/readyz did not reach 200 within {READYZ_BUDGET_SECS}s; last_err={}",
last_err.unwrap_or_else(|| "<none>".into())
))
}
#[derive(Debug, Clone)]
struct StreamResult {
http_status: i32,
ttft_ms: f64,
tokens: u32,
total_ms: f64,
}
fn run_stream(port: u16, prompt: &str, max_tokens: u32, model: &str) -> StreamResult {
let body = format!(
r#"{{"model":"{}","messages":[{{"role":"user","content":"{}"}}],"max_tokens":{},"temperature":0.6,"stream":true}}"#,
model.replace('"', "\\\""),
prompt.replace('"', "\\\""),
max_tokens,
);
let mut cmd = Command::new("curl");
cmd.args([
"-s",
"-N",
"-X",
"POST",
"-H",
"Content-Type: application/json",
"--max-time",
&STREAM_BUDGET_SECS.to_string(),
"-w",
"\n__HTTP_STATUS__:%{http_code}\n",
"-d",
&body,
&format!("http://127.0.0.1:{port}/v1/chat/completions"),
])
.stdout(Stdio::piped())
.stderr(Stdio::null());
let t0 = Instant::now();
let mut child = match cmd.spawn() {
Ok(c) => c,
Err(_) => {
return StreamResult {
http_status: -1,
ttft_ms: 0.0,
tokens: 0,
total_ms: t0.elapsed().as_secs_f64() * 1000.0,
};
}
};
let stdout = match child.stdout.take() {
Some(s) => s,
None => {
let _ = child.kill();
let _ = child.wait();
return StreamResult {
http_status: -1,
ttft_ms: 0.0,
tokens: 0,
total_ms: t0.elapsed().as_secs_f64() * 1000.0,
};
}
};
let reader = BufReader::new(stdout);
let mut tokens: u32 = 0;
let mut ttft_ms = 0.0_f64;
let mut first_content_seen = false;
let mut http_status: i32 = -1;
for line_res in reader.lines() {
let line = match line_res {
Ok(l) => l,
Err(_) => break,
};
if let Some(code_str) = line.strip_prefix("__HTTP_STATUS__:") {
if let Ok(code) = code_str.trim().parse::<i32>() {
http_status = code;
}
continue;
}
let payload = match line.strip_prefix("data: ") {
Some(p) => p,
None => continue,
};
if payload.trim() == "[DONE]" {
continue;
}
if let Some(idx) = payload.find(r#""content":""#) {
let after = &payload[idx + r#""content":""#.len()..];
if !after.starts_with('"') {
tokens = tokens.saturating_add(1);
if !first_content_seen {
first_content_seen = true;
ttft_ms = t0.elapsed().as_secs_f64() * 1000.0;
}
}
}
}
let _ = child.wait();
let total_ms = t0.elapsed().as_secs_f64() * 1000.0;
StreamResult {
http_status,
ttft_ms,
tokens,
total_ms,
}
}
fn run_bench_cell(gguf: &str, policy: &'static str, n: u32) -> Result<ThroughputCell, String> {
let prompt = std::env::var("HF2Q_CB_THROUGHPUT_PROMPT")
.unwrap_or_else(|_| "Count slowly from one to twenty, one number per line.".into());
let max_tokens: u32 = std::env::var("HF2Q_CB_THROUGHPUT_MAX_TOKENS")
.ok()
.and_then(|s| s.parse().ok())
.unwrap_or(64);
let port = next_port();
eprintln!("[cb-throughput] cell: policy={policy}, N={n}, port={port}, gguf={gguf}");
let mut server =
BenchServer::spawn(gguf, policy, n, port).map_err(|e| format!("spawn hf2q serve: {e}"))?;
wait_for_readyz(&mut server)?;
eprintln!("[cb-throughput] /readyz=200 on port={port}");
let model_id = fetch_model_id(port).map_err(|e| format!("GET /v1/models: {e}"))?;
eprintln!("[cb-throughput] resolved model_id={model_id}");
let cell_start = Instant::now();
let results: Vec<StreamResult> = std::thread::scope(|s| {
let handles: Vec<_> = (0..n)
.map(|_| {
let p = port;
let pr = prompt.clone();
let mid = model_id.clone();
s.spawn(move || run_stream(p, &pr, max_tokens, &mid))
})
.collect();
handles.into_iter().map(|h| h.join().unwrap()).collect()
});
let cell_walltime_ms = cell_start.elapsed().as_secs_f64() * 1000.0;
let rejected_429_count = results.iter().filter(|r| r.http_status == 429).count() as u32;
let succeeded: Vec<&StreamResult> = results
.iter()
.filter(|r| r.http_status == 200 && r.tokens > 0)
.collect();
if succeeded.is_empty() {
return Err(format!(
"all {n} streams failed for policy={policy}; \
results: {:?}",
results
));
}
let total_tokens: u32 = succeeded.iter().map(|r| r.tokens).sum();
let aggregate_tokens_per_sec = (total_tokens as f64) / (cell_walltime_ms / 1000.0).max(1e-6);
let mut per_stream_rates: Vec<f64> = succeeded
.iter()
.map(|r| (r.tokens as f64) / (r.total_ms / 1000.0).max(1e-6))
.collect();
per_stream_rates.sort_by(|a, b| a.partial_cmp(b).unwrap());
let per_slot_tokens_per_sec = if per_stream_rates.is_empty() {
0.0
} else {
per_stream_rates[per_stream_rates.len() / 2]
};
let mut ttft_per_stream_ms: Vec<f64> = succeeded.iter().map(|r| r.ttft_ms).collect();
ttft_per_stream_ms.sort_by(|a, b| a.partial_cmp(b).unwrap());
let ttft_p50_ms = percentile(&ttft_per_stream_ms, 0.50);
let ttft_p95_ms = percentile(&ttft_per_stream_ms, 0.95);
let cell = ThroughputCell {
policy,
concurrency: n,
aggregate_tokens_per_sec,
ttft_p50_ms,
ttft_p95_ms,
per_slot_tokens_per_sec,
rejected_429_count,
};
eprintln!("[cb-throughput] cell DONE: {:?}", cell);
Ok(cell)
}
fn run_bench_cell_3rep(
gguf: &str,
policy: &'static str,
n: u32,
) -> Result<ThroughputCellStable, String> {
let mut cells = Vec::with_capacity(REPS);
for rep in 0..REPS {
eprintln!(
"[d3] cell policy={} N={} rep={}/{}",
policy,
n,
rep + 1,
REPS
);
let cell = run_bench_cell(gguf, policy, n).map_err(|e| {
format!(
"rep {}/{} failed for policy={} N={}: {}",
rep + 1,
REPS,
policy,
n,
e
)
})?;
cells.push(cell);
}
Ok(ThroughputCellStable::from_reps(cells))
}
fn percentile(sorted: &[f64], q: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let idx = ((sorted.len() as f64) * q).ceil() as usize;
let idx = idx.saturating_sub(1).min(sorted.len() - 1);
sorted[idx]
}
fn fetch_model_id(port: u16) -> std::io::Result<String> {
let mut s = TcpStream::connect_timeout(
&format!("127.0.0.1:{port}")
.parse()
.map_err(std::io::Error::other)?,
Duration::from_secs(5),
)?;
s.set_read_timeout(Some(Duration::from_secs(30)))?;
s.write_all(b"GET /v1/models HTTP/1.1\r\nHost: 127.0.0.1\r\nConnection: close\r\n\r\n")?;
let mut buf = Vec::new();
s.read_to_end(&mut buf)?;
let body = String::from_utf8_lossy(&buf);
let json_start = body
.find("\r\n\r\n")
.ok_or_else(|| std::io::Error::other("/v1/models: no headers/body separator"))?;
let json = &body[json_start + 4..];
let id_key = r#""id":""#;
let idx = json
.find(id_key)
.ok_or_else(|| std::io::Error::other(format!("/v1/models: no id field in body: {json}")))?;
let after = &json[idx + id_key.len()..];
let end = after
.find('"')
.ok_or_else(|| std::io::Error::other("/v1/models: unterminated id string"))?;
Ok(after[..end].to_string())
}
#[test]
fn cb_throughput_n_1_2_4_8_fifo_vs_inflight() {
if std::env::var("HF2Q_CB_THROUGHPUT_E2E").as_deref() != Ok("1") {
eprintln!(
"[cb-throughput] skipped — set HF2Q_CB_THROUGHPUT_E2E=1 + \
HF2Q_CB_THROUGHPUT_MODEL=<gguf> to run the ADR-040 §5 AC-4 \
throughput harness. Operator command:\n \
HF2Q_CB_THROUGHPUT_E2E=1 HF2Q_CB_THROUGHPUT_MODEL=/path/to.gguf \\\n \
cargo test --release --test continuous_batching_throughput -- \\\n \
--test-threads=1 --nocapture cb_throughput_n_1_2_4_8_fifo_vs_inflight"
);
return;
}
let gguf_path = std::env::var("HF2Q_CB_THROUGHPUT_MODEL").expect(
"HF2Q_CB_THROUGHPUT_MODEL required when HF2Q_CB_THROUGHPUT_E2E=1. \
Either set the env or unset HF2Q_CB_THROUGHPUT_E2E — silent skip \
with the gate set violates the iter-1.5 cfa-finding-F8 contract.",
);
let concurrency: Vec<u32> = std::env::var("HF2Q_CB_THROUGHPUT_CONCURRENCY")
.unwrap_or_else(|_| String::from("1,2,4,8"))
.split(',')
.map(|s| {
s.trim()
.parse::<u32>()
.expect("HF2Q_CB_THROUGHPUT_CONCURRENCY must be comma-separated u32 list")
})
.collect();
assert!(
!concurrency.is_empty(),
"concurrency list must be non-empty"
);
for &n in &concurrency {
assert!(n > 0, "concurrency entries must be > 0, got {n}");
}
let mut all_cells: Vec<ThroughputCellStable> = Vec::new();
let mut skipped: Vec<(String, u32, String)> = Vec::new();
for &policy in &["fifo_serial", "inflight_batched"] {
for &n in &concurrency {
match run_bench_cell_3rep(&gguf_path, policy, n) {
Ok(cell) => all_cells.push(cell),
Err(e) => {
eprintln!("[cb-throughput] cell SKIPPED (policy={policy}, N={n}): {e}");
skipped.push((policy.to_string(), n, e));
}
}
}
}
let report = render_report_stable(&all_cells);
println!("\n=== ADR-040 §5 AC-4 throughput report (D3, REPS={REPS}) ===\n{report}");
if !skipped.is_empty() {
println!("\nSkipped cells:");
for (p, n, why) in &skipped {
println!(" - policy={p}, N={n}: {why}");
}
}
assert!(
!all_cells.is_empty(),
"ADR-040 D3: no bench cells completed; all (policy, N) combinations failed. \
Check HF2Q_CB_THROUGHPUT_MODEL={gguf_path} + the skipped-cells list above."
);
println!("\n=== D3 FifoSerial-only stability baseline ===");
let mut any_fifo = false;
for cell in all_cells.iter().filter(|c| c.policy == "fifo_serial") {
any_fifo = true;
println!(
" N={}: median={:.1} tok/s, min={:.1}, max={:.1}, sigma_pct={:.1}% ({} reps, 429s total={})",
cell.concurrency,
cell.aggregate_tokens_per_sec_median,
cell.aggregate_tokens_per_sec_min,
cell.aggregate_tokens_per_sec_max,
cell.aggregate_tokens_per_sec_sigma_pct,
cell.rep_count,
cell.rejected_429_count_total,
);
}
if !any_fifo {
println!(" (no fifo_serial cells completed)");
}
let fifo_n4 = all_cells
.iter()
.find(|c| c.policy == "fifo_serial" && c.concurrency == 4);
let inflight_n4 = all_cells
.iter()
.find(|c| c.policy == "inflight_batched" && c.concurrency == 4);
let fifo_n1 = all_cells
.iter()
.find(|c| c.policy == "fifo_serial" && c.concurrency == 1);
match ac4_outcome(fifo_n4, inflight_n4, fifo_n1) {
Ac4Outcome::Deferred => {
eprintln!(
"[ac-4 DEFERRED] cannot evaluate gate — missing N=4 cell for one or both \
policies. Typically inflight_batched is Phase C2c/C2d gated; D3 reports \
the FifoSerial-only baseline + variance and defers AC-4 enforcement until \
the inflight-side wiring lands. The bench is forward-compatible: no edits \
needed once C2c/C2d ship."
);
}
Ac4Outcome::Misconfigured => {
let f = fifo_n4.expect("Misconfigured requires fifo_n4 present");
let i = inflight_n4.expect("Misconfigured requires inflight_n4 present");
panic!(
"[ac-4 MISCONFIGURED] BOTH N=4 cells present (fifo_serial \
N=4 median = {:.1} tok/s; inflight_batched N=4 median = {:.1} tok/s) \
BUT fifo_serial N=1 baseline cell is missing. The AC-4 TTFT half \
compares inflight_batched N=4 p95 against fifo_serial N=1 p95; \
skipping it would let a TTFT regression slip past the gate. \
Re-run with HF2Q_CB_THROUGHPUT_CONCURRENCY including `1` (e.g. \
`1,2,4,8` — the default).",
f.aggregate_tokens_per_sec_median, i.aggregate_tokens_per_sec_median,
);
}
Ac4Outcome::StabilityBlocked => {
let f = fifo_n4.expect("StabilityBlocked requires fifo_n4 present");
let i = inflight_n4.expect("StabilityBlocked requires inflight_n4 present");
if f.aggregate_tokens_per_sec_sigma_pct > STABILITY_SIGMA_PCT_THRESHOLD {
panic!(
"AC-4 BLOCKED: fifo_serial N=4 measurement variance {:.1}% > {:.1}% threshold; \
run again or increase REPS for stable median (median={:.1}, min={:.1}, max={:.1})",
f.aggregate_tokens_per_sec_sigma_pct,
STABILITY_SIGMA_PCT_THRESHOLD,
f.aggregate_tokens_per_sec_median,
f.aggregate_tokens_per_sec_min,
f.aggregate_tokens_per_sec_max,
);
}
panic!(
"AC-4 BLOCKED: inflight_batched N=4 measurement variance {:.1}% > {:.1}% threshold; \
run again or increase REPS for stable median (median={:.1}, min={:.1}, max={:.1})",
i.aggregate_tokens_per_sec_sigma_pct,
STABILITY_SIGMA_PCT_THRESHOLD,
i.aggregate_tokens_per_sec_median,
i.aggregate_tokens_per_sec_min,
i.aggregate_tokens_per_sec_max,
);
}
Ac4Outcome::Failed {
aggregate_ratio,
ttft_ratio,
which,
} => {
let f = fifo_n4.expect("Failed requires fifo_n4 present");
let i = inflight_n4.expect("Failed requires inflight_n4 present");
if which == "aggregate" {
panic!(
"AC-4 FAILED: aggregate ratio {:.2}× below 1.5× bar \
(fifo_serial N=4 median = {:.1} tok/s; inflight_batched N=4 median = {:.1} tok/s)",
aggregate_ratio,
f.aggregate_tokens_per_sec_median,
i.aggregate_tokens_per_sec_median,
);
}
let base = fifo_n1.expect("Failed/ttft requires fifo_n1 present");
panic!(
"AC-4 FAILED: TTFT p95 ratio {:.2}× above 2.0× bar \
(fifo_serial N=1 p95 median = {:.1} ms; inflight_batched N=4 p95 median = {:.1} ms)",
ttft_ratio,
base.ttft_p95_ms_median,
i.ttft_p95_ms_median,
);
}
Ac4Outcome::Passed {
aggregate_ratio,
ttft_ratio,
} => {
eprintln!(
"[ac-4 PASS] aggregate ratio {:.2}× ≥ 1.5× ✓ ; TTFT p95 ratio {:.2}× ≤ 2.0× ✓",
aggregate_ratio, ttft_ratio,
);
}
}
}
#[test]
fn cb_throughput_required_env_vars_documented() {
if std::env::var("HF2Q_CB_THROUGHPUT_E2E").as_deref() != Ok("1") {
return;
}
let model = std::env::var("HF2Q_CB_THROUGHPUT_MODEL");
let concurrency =
std::env::var("HF2Q_CB_THROUGHPUT_CONCURRENCY").unwrap_or_else(|_| String::from("1,2,4,8"));
assert!(
model.is_ok(),
"HF2Q_CB_THROUGHPUT_E2E=1 set but HF2Q_CB_THROUGHPUT_MODEL absent — \
D2 refuses to run without a GGUF path"
);
let parsed: Result<Vec<u32>, _> = concurrency
.split(',')
.map(|s| s.trim().parse::<u32>())
.collect();
assert!(
parsed.is_ok(),
"HF2Q_CB_THROUGHPUT_CONCURRENCY parse failed: {:?}",
parsed
);
let ns = parsed.unwrap();
assert!(!ns.is_empty(), "at least one N value required");
for &n in &ns {
assert!(n > 0, "N must be positive, got {}", n);
}
}
#[derive(Debug, Clone, Copy, PartialEq)]
pub struct AcceptanceCell {
pub concurrent: u32,
pub acceptance_rate: f64,
pub tokens_per_step: f64,
}
impl AcceptanceCell {
pub fn synthetic_for_smoke() -> Self {
AcceptanceCell {
concurrent: 1,
acceptance_rate: 0.0,
tokens_per_step: 0.0,
}
}
}
pub fn render_acceptance_report(cells: &[AcceptanceCell]) -> String {
let mut s = String::from("| concurrent | acceptance_rate | tokens_per_step |\n");
s.push_str("|------------|-----------------|-----------------|\n");
for c in cells {
s.push_str(&format!(
"| {} | {:.3} | {:.2} |\n",
c.concurrent, c.acceptance_rate, c.tokens_per_step
));
}
s
}
#[test]
fn acceptance_cell_synthetic_round_trips_through_report() {
let cells = vec![AcceptanceCell::synthetic_for_smoke()];
let report = render_acceptance_report(&cells);
assert!(report.contains("| 1 |"), "acceptance report: {}", report);
assert!(report.contains("concurrent"), "header missing");
assert!(report.contains("|------------|"), "separator missing");
}
#[test]
fn render_acceptance_report_empty_returns_header_only() {
let report = render_acceptance_report(&[]);
let lines: Vec<&str> = report.lines().collect();
assert_eq!(lines.len(), 2, "header + separator only, got: {:?}", lines);
}
#[test]
fn a4_inflection_bench_acceptance_dimension_scaffold() {
if std::env::var("HF2Q_CB_THROUGHPUT_E2E").as_deref() != Ok("1") {
eprintln!(
"[a4-inflection-bench] skipped — set HF2Q_CB_THROUGHPUT_E2E=1 \
AND HF2Q_A4_INFLECTION_BENCH=1 AND HF2Q_CB_THROUGHPUT_MODEL=<gguf> \
to engage the acceptance-rate dimension bench. Today this is a \
structural scaffold per ADR-040 §6.1.55 iter-A4-cont-inflection-bench."
);
return;
}
if std::env::var("HF2Q_A4_INFLECTION_BENCH").as_deref() != Ok("1") {
eprintln!(
"[a4-inflection-bench] skipped — HF2Q_CB_THROUGHPUT_E2E=1 set but \
HF2Q_A4_INFLECTION_BENCH != 1; respecting bench opt-in."
);
return;
}
let cells: Vec<AcceptanceCell> = vec![AcceptanceCell::synthetic_for_smoke()];
let report = render_acceptance_report(&cells);
println!("\n=== ADR-040 §6.1.55 iter-A4-cont-inflection-bench (scaffold) ===\n{report}");
}
#[test]
fn a4_moe_validation_qwen36_a3b_a_b_n_1_2_4_8() {
if std::env::var("HF2Q_A4_MOE_AB_VALIDATION_E2E").as_deref() != Ok("1") {
eprintln!(
"[a4-moe-validation] skipped — set HF2Q_A4_MOE_AB_VALIDATION_E2E=1 \
AND HF2Q_CB_THROUGHPUT_MODEL=<Qwen3.6-A3B-gguf> to engage the \
MoE A/B bench at N=1,2,4,8 concurrent. ADR-040 §6.1.55 \
iter-A4-cont-moe-validation harness."
);
return;
}
let gguf_path = std::env::var("HF2Q_CB_THROUGHPUT_MODEL").expect(
"HF2Q_CB_THROUGHPUT_MODEL required when HF2Q_A4_MOE_AB_VALIDATION_E2E=1. \
Either set the env or unset HF2Q_A4_MOE_AB_VALIDATION_E2E — silent skip \
with the gate set violates the iter-1.5 cfa-finding-F8 contract.",
);
eprintln!(
"[a4-moe-validation] running A/B sweep against {gguf_path} \
at N=1,2,4,8 per ADR-040 §6.1.55 iter-A4-cont-moe-validation."
);
let concurrencies: Vec<u32> = vec![1, 2, 4, 8];
for &n in &concurrencies {
eprintln!("[a4-moe-validation] N={n}: deferred to iter-A4-cont-drafter-dispatcher-kernel");
}
}