use hya_core::{Action, Scheduler, Source};
struct Rng(u64);
impl Rng {
fn next_f64(&mut self) -> f64 {
let mut x = self.0;
x ^= x >> 12;
x ^= x << 25;
x ^= x >> 27;
self.0 = x;
(x.wrapping_mul(0x2545_F491_4F6C_DD1D) >> 11) as f64 / (1u64 << 53) as f64
}
fn normal(&mut self) -> f64 {
let s: f64 = (0..4).map(|_| self.next_f64()).sum();
(s - 2.0) * 1.732_050_8
}
}
#[derive(Clone, Copy)]
struct Flow {
ready_at: f64,
live: bool,
task_hi: u64,
task_off: u64,
weight: f64,
}
struct Params {
size: u64,
n: usize,
capacity: f64,
delta: f64,
tau: f64,
sigma: f64,
sigma_w: f64,
theta_scale: f64,
seed: u64,
collapse_at: Option<f64>,
collapse_to: f64,
}
struct Outcome {
makespan: f64,
requests: u64,
repairs: u64,
repair_progress: Vec<f64>,
repair_theta: Vec<f64>,
setup_waste: f64,
dup_bytes: u64,
}
fn run(p: &Params, zombies_bite: bool) -> Outcome {
let dt = 0.010;
let sources = vec![Source {
gamma_est: p.capacity / p.n as f64,
delta_est: p.delta,
..Default::default()
}];
let mut sched = Scheduler::new(p.size, sources, &[p.n]).with_theta_scale(p.theta_scale);
let mut rng = Rng(p.seed | 1);
let mut flows: Vec<Flow> = (0..p.n)
.map(|_| Flow {
ready_at: 0.0,
live: false,
weight: (p.sigma_w * rng.normal()).exp(),
task_hi: 0,
task_off: 0,
})
.collect();
let mut setup_waste = 0.0;
let mut dup_bytes = 0u64;
let mut collapsed = false;
let mut repair_progress: Vec<f64> = Vec::new();
let mut repair_theta: Vec<f64> = Vec::new();
let mut last_repairs = 0u64;
let mut now = 0.0;
let horizon = 100.0 * p.size as f64 / p.capacity;
while now < horizon {
for act in sched.tick(now) {
match act {
Action::Request { conn, range } => {
flows[conn].ready_at = now + p.delta;
flows[conn].live = true;
flows[conn].task_hi = range.hi;
flows[conn].task_off = range.lo;
}
Action::Cancel { conn } => flows[conn].live = false,
Action::Shrink { conn, hi } => {
if !zombies_bite {
flows[conn].task_hi = hi;
}
}
}
}
if sched.stats.repairs > last_repairs {
let k = (sched.stats.repairs - last_repairs) as usize;
let frac = sched.bytes_held() as f64 / p.size as f64;
let th = sched.theta_now(now);
for _ in 0..k {
repair_progress.push(frac);
repair_theta.push(th);
}
last_repairs = sched.stats.repairs;
}
if sched.is_complete() {
break;
}
if let Some(frac) = p.collapse_at {
let at = frac * p.size as f64 / p.capacity;
if now >= at && !collapsed {
flows[0].weight *= p.collapse_to;
collapsed = true;
}
}
let mut ramp = vec![0.0f64; p.n];
let mut zombie = vec![false; p.n];
for j in 0..p.n {
let sched_hi = sched.conn_range(j).map(|(_, _, hi)| hi);
if zombies_bite {
let live_socket = flows[j].live && flows[j].task_off < flows[j].task_hi;
let shrunk = match sched_hi {
Some(hi) => hi < flows[j].task_hi,
None => true,
};
if live_socket && shrunk && now >= flows[j].ready_at {
let age = now - flows[j].ready_at;
let base = 1.0 - (-age / p.tau).exp();
ramp[j] = (base * flows[j].weight).max(0.0);
zombie[j] = true;
continue;
}
}
let has_range = sched.conn_range(j).is_some();
if has_range && flows[j].live && now >= flows[j].ready_at {
let age = now - flows[j].ready_at;
let base = 1.0 - (-age / p.tau).exp();
let noise = (p.sigma * rng.normal()).exp();
ramp[j] = (base * noise * flows[j].weight).max(0.0);
} else if has_range && flows[j].live {
setup_waste += dt;
}
}
let total: f64 = ramp.iter().sum();
if total > 0.0 {
for j in 0..p.n {
if ramp[j] <= 0.0 {
continue;
}
let rate = p.capacity * ramp[j] / total;
let bytes = (rate * dt) as u64;
if bytes == 0 {
continue;
}
if zombie[j] {
let room = flows[j].task_hi.saturating_sub(flows[j].task_off);
let step = bytes.min(room);
flows[j].task_off += step;
dup_bytes += step;
if flows[j].task_off >= flows[j].task_hi {
flows[j].live = false;
}
continue;
}
if let Some((_, pos, _)) = sched.conn_range(j) {
sched.on_bytes_at(j, pos, bytes, now, dt);
flows[j].task_off = pos + bytes;
}
}
}
now += dt;
}
Outcome {
makespan: now,
requests: sched.stats.requests,
repairs: sched.stats.repairs,
setup_waste,
dup_bytes,
repair_progress,
repair_theta,
}
}
fn main() {
let out = std::env::args().nth(1);
let size = 5_300_000u64;
let capacity = 1_400_000.0; let oracle = size as f64 / capacity;
let seeds = 12u64;
let mut csv = String::from(
"n,zombie,seed,makespan_s,oracle_ratio,requests,repairs,dup_bytes,setup_waste_s\n",
);
println!(
"object {} bytes, shared capacity {:.0} B/s, fluid oracle {:.2}s, {} seeds",
size, capacity, oracle, seeds
);
println!("stationary transfer: the correct repair count is 0.");
println!("zombie=yes is the SHIPPED transport: a repair shrinks the victim's range");
println!("scheduler-side but emits no Cancel, and fetch_range loops on the `hi` it");
println!("captured at spawn, so the stolen span crosses the wire twice.\n");
println!(
"{:>3} {:>7} {:>9} {:>7} {:>8} {:>8} {:>9}",
"n", "zombie", "makespan", "ratio", "requests", "repairs", "dup_MB"
);
for &n in &[1usize, 2, 4, 8, 16] {
for &z in &[true, false] {
let (mut ms, mut rq, mut rp, mut db) = (0.0, 0u64, 0u64, 0u64);
let mut all_prog: Vec<f64> = Vec::new();
let mut all_theta: Vec<f64> = Vec::new();
for sd in 0..seeds {
let p = Params {
size,
n,
capacity,
delta: 0.12,
tau: 0.35,
sigma: 0.10,
sigma_w: 0.45,
theta_scale: 1.0,
seed: 0x5EED + sd * 7919,
collapse_at: None,
collapse_to: 1.0,
};
let o = run(&p, z);
csv.push_str(&format!(
"{},{},{},{:.4},{:.4},{},{},{},{:.4}\n",
n,
if z { 1 } else { 0 },
p.seed,
o.makespan,
o.makespan / oracle,
o.requests,
o.repairs,
o.dup_bytes,
o.setup_waste
));
all_prog.extend(o.repair_progress.iter().copied());
all_theta.extend(o.repair_theta.iter().copied());
ms += o.makespan;
rq += o.requests;
rp += o.repairs;
db += o.dup_bytes;
}
if !z && !all_prog.is_empty() {
let nr = all_prog.len() as f64;
let mean_p = all_prog.iter().sum::<f64>() / nr;
let late = all_prog.iter().filter(|x| **x > 0.9).count();
let mean_th = all_theta.iter().sum::<f64>() / nr;
let min_th = all_theta.iter().cloned().fold(f64::INFINITY, f64::min);
eprintln!(
" n={n}: {} repairs | mean progress at repair {:.2} | {} fired past 90% \
| theta mean {:.3}s min {:.4}s",
all_prog.len(),
mean_p,
late,
mean_th,
min_th
);
}
let f = seeds as f64;
println!(
"{:>3} {:>7} {:>8.2}s {:>7.3} {:>8.1} {:>8.1} {:>8.2}",
n,
if z { "yes" } else { "no" },
ms / f,
(ms / f) / oracle,
rq as f64 / f,
rp as f64 / f,
db as f64 / f / 1e6
);
}
}
println!("\ncollapse arm: connection 0 drops to 5% of its share at 30% progress");
println!(
"{:>3} {:>9} {:>7} {:>8} {:>8}",
"n", "makespan", "ratio", "requests", "repairs"
);
for &n in &[2usize, 4, 8] {
let (mut ms, mut rq, mut rp) = (0.0, 0u64, 0u64);
for sd in 0..seeds {
let o = run(
&Params {
size,
n,
capacity,
delta: 0.12,
tau: 0.35,
sigma: 0.10,
sigma_w: 0.45,
theta_scale: 1.0,
seed: 0x5EED + sd * 7919,
collapse_at: Some(0.30),
collapse_to: 0.05,
},
false,
);
ms += o.makespan;
rq += o.requests;
rp += o.repairs;
}
let f = seeds as f64;
println!(
"{:>3} {:>8.2}s {:>7.3} {:>8.1} {:>8.1}",
n,
ms / f,
(ms / f) / oracle,
rq as f64 / f,
rp as f64 / f
);
}
if let Some(path) = out {
std::fs::write(&path, csv).expect("write csv");
println!("\nwrote {}", path);
}
}