mod support;
use std::net::{Ipv4Addr, SocketAddrV4};
use std::path::{Path, PathBuf};
use std::thread;
use std::time::{Duration, Instant};
use ringfire::replication::{Mirror, MirrorStart, MulticastConfig, ReplicaServer};
use ringfire::{RingConsumer, RingProducer};
use support::{elapsed_ns, latency_line, parse_arg, parse_flag, Pacer, SeqTracker};
#[derive(Debug, Clone, Copy)]
#[repr(C)]
struct Msg {
sent_ns: u64,
seq: u64,
_pad: [u8; 40],
}
fn connect_with_retry(peer: &str, path: &Path, iface: Ipv4Addr, unicast: bool) -> Mirror {
let deadline = Instant::now() + Duration::from_secs(60);
loop {
match Mirror::builder()
.start(MirrorStart::Latest)
.spin(true)
.interface(iface)
.unicast(unicast)
.connect(peer, path)
{
Ok(mirror) => return mirror,
Err(e) => {
assert!(
Instant::now() < deadline,
"could not connect to {}: {}",
peer,
e
);
thread::sleep(Duration::from_millis(200));
}
}
}
}
fn main() {
let args: Vec<String> = std::env::args().collect();
let mut role = String::new();
let mut bind = String::new();
let mut peer = String::new();
let mut rate = 100_000u64;
let mut seconds = 5u64;
let mut dir = PathBuf::from("/dev/shm");
let mut multicast: Option<SocketAddrV4> = None;
let mut iface = Ipv4Addr::UNSPECIFIED;
let mut linger_us: Option<u64> = None;
let mut warmup_ms = 1500u64;
let mut idle_ms = 1000u64;
let mut burst = 1u64;
let mut udp_port: Option<u16> = None;
let mut unicast = false;
let mut dup = 1u8;
let mut exit_on_close = false;
let mut i = 1;
while i + 1 < args.len() {
let v = args[i + 1].as_str();
match args[i].as_str() {
"--linger-us" => linger_us = Some(parse_arg("--linger-us", v)),
"--burst" => burst = parse_arg::<u64>("--burst", v).max(1),
"--udp" => udp_port = Some(parse_arg("--udp", v)),
"--dup" => dup = parse_arg("--dup", v),
"--unicast" => unicast = parse_flag(v),
"--warmup-ms" => warmup_ms = parse_arg("--warmup-ms", v),
"--idle-ms" => idle_ms = parse_arg("--idle-ms", v),
"--exit-on-close" => exit_on_close = parse_flag(v),
"--role" => role = v.to_string(),
"--bind" => bind = v.to_string(),
"--peer" => peer = v.to_string(),
"--rate" => rate = parse_arg("--rate", v),
"--seconds" => seconds = parse_arg("--seconds", v),
"--dir" => dir = PathBuf::from(v),
"--multicast" => multicast = Some(parse_arg("--multicast", v)),
"--iface" => iface = parse_arg("--iface", v),
_ => {}
}
i += 2;
}
assert!(
!role.is_empty() && !bind.is_empty() && !peer.is_empty(),
"--role, --bind and --peer are required"
);
assert!(rate > 0, "--rate must be positive");
let pid = std::process::id();
let out_path = dir.join(format!("ringfire_stress_{}_{}_out.shm", role, pid));
let in_path = dir.join(format!("ringfire_stress_{}_{}_in.shm", role, pid));
let _ = std::fs::remove_file(&out_path);
let _ = std::fs::remove_file(&in_path);
let mut out = RingProducer::<Msg>::create(&out_path, 1 << 18).unwrap();
let mut server = ReplicaServer::bind(&out_path, bind.as_str())
.unwrap()
.spin(true)
.linger(linger_us.map(Duration::from_micros));
if let Some(group) = multicast {
server = server.multicast(MulticastConfig::new(*group.ip(), group.port()).interface(iface));
}
if let Some(port) = udp_port {
server = server.unicast(port, 1400);
}
server = server.duplicate(dup);
server.spawn().unwrap();
let mut mirror = connect_with_retry(&peer, &in_path, iface, unicast);
eprintln!(
"{}: serving {} on {}, mirroring {} ({})",
role,
out_path.display(),
bind,
peer,
if mirror.is_multicast() {
"multicast"
} else {
"tcp"
}
);
let mut input = RingConsumer::<Msg>::attach(&in_path).unwrap();
let failed = match role.as_str() {
"ponger" => {
let mut echoed = 0u64;
let mut last_report = Instant::now();
loop {
let open = match mirror.step() {
Ok(open) => open,
Err(e) => {
eprintln!("ponger: mirror failed: {}", e);
break true;
}
};
if !open {
eprintln!(
"ponger: source closed after echoing {}; reader lapped {}, mirror datagrams {} naks {} retransmitted {} gaps {}{}",
echoed,
input.lapped_count(),
mirror.datagrams(),
mirror.naks(),
mirror.retransmitted(),
mirror.gaps(),
if exit_on_close { "" } else { "; reconnecting" }
);
if exit_on_close {
break false;
}
echoed = 0;
mirror = connect_with_retry(&peer, &in_path, iface, unicast);
input = RingConsumer::<Msg>::attach(&in_path).unwrap();
}
while let Some(msg) = input.try_recv() {
out.push(&msg);
echoed += 1;
}
if last_report.elapsed() >= Duration::from_secs(5) {
eprintln!(
"ponger: echoed {}, reader lapped {}, mirror datagrams {} naks {} retransmitted {} gaps {}",
echoed,
input.lapped_count(),
mirror.datagrams(),
mirror.naks(),
mirror.retransmitted(),
mirror.gaps()
);
last_report = Instant::now();
}
}
}
"pinger" => {
let handle = mirror.handle().unwrap();
let mirror_thread = thread::spawn(move || {
let result = mirror.run();
(
result,
mirror.datagrams(),
mirror.naks(),
mirror.retransmitted(),
mirror.gaps(),
)
});
let epoch = Instant::now();
let now_ns = move || epoch.elapsed().as_nanos() as u64;
let total = rate * seconds;
thread::sleep(Duration::from_millis(warmup_ms));
let mut pusher = Some(thread::spawn(move || {
let start = Instant::now();
let mut pace = Pacer::per_second(rate, burst);
for seq in 1..=total {
if (seq - 1) % burst == 0 {
pace.wait_next();
}
out.push(&Msg {
sent_ns: now_ns(),
seq,
_pad: [0; 40],
});
}
(out, start.elapsed())
}));
let mut rtt: Vec<i64> = Vec::with_capacity(total as usize);
let mut echoes = SeqTracker::new(total);
let idle = Duration::from_millis(idle_ms);
let mut last_arrival = Instant::now();
let mut pushed_done: Option<Duration> = None;
let mut out_ring = None;
let mut empty = 0u32;
loop {
if let Some(msg) = input.try_recv() {
let arrived = now_ns();
if echoes.record(msg.seq) {
rtt.push(elapsed_ns(msg.sent_ns, arrived));
}
empty = 0;
if pushed_done.is_some() {
last_arrival = Instant::now();
}
} else {
empty = empty.wrapping_add(1);
if empty.is_multiple_of(256) {
if pushed_done.is_none() && pusher.as_ref().is_some_and(|p| p.is_finished()) {
let (ring, elapsed) = pusher.take().unwrap().join().unwrap();
out_ring = Some(ring);
pushed_done = Some(elapsed);
last_arrival = Instant::now();
}
if pushed_done.is_some()
&& (echoes.unique() >= total || last_arrival.elapsed() > idle)
{
break;
}
}
core::hint::spin_loop();
}
}
let push_elapsed = pushed_done.unwrap_or_default();
drop(out_ring);
handle.shutdown().unwrap();
let (result, datagrams, naks, retransmitted, gaps) = mirror_thread.join().unwrap();
println!(
"rate {} msg/s (bursts of {}) x {} s: pushed {} in {:.2} s ({:.0} msg/s achieved), echoed {} ({:.3}%), lost {}, duplicates {}, disorder {}, reader lapped {}",
rate,
burst,
seconds,
total,
push_elapsed.as_secs_f64(),
total as f64 / push_elapsed.as_secs_f64().max(1e-9),
echoes.unique(),
echoes.unique() as f64 * 100.0 / total.max(1) as f64,
echoes.missing_in(1, total),
echoes.duplicates(),
echoes.out_of_order(),
input.lapped_count()
);
println!(" rtt {}", latency_line(&mut rtt));
println!(
" pinger mirror: datagrams {}, naks {}, retransmitted {}, gaps {}",
datagrams, naks, retransmitted, gaps
);
match result {
Ok(()) => false,
Err(e) => {
eprintln!("pinger: mirror failed: {}", e);
true
}
}
}
other => panic!("unknown role {}", other),
};
drop(input);
let _ = std::fs::remove_file(&in_path);
let _ = std::fs::remove_file(&out_path);
if failed {
std::process::exit(1);
}
}