mod support;
use std::io::{Read, Write};
use std::net::{Ipv4Addr, SocketAddrV4, TcpListener, TcpStream, ToSocketAddrs};
use std::path::PathBuf;
use std::sync::atomic::{AtomicI64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::thread;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use ringfire::replication::{Mirror, MirrorStart, MulticastConfig, ReplicaServer};
use ringfire::{RingConsumer, RingProducer};
use support::{Pacer, SeqTracker, latency_line, parse_arg, parse_flag};
#[derive(Debug, Clone, Copy)]
#[repr(C)]
struct Msg {
sent_ns: u64,
seq: u64,
sent_mono_ns: u64,
_pad: [u8; 32],
}
fn wall_ns() -> i64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_nanos() as i64
}
fn report(label: &str, samples: &mut [i64], extra: &str) {
println!("{}: {} {}", label, latency_line(samples), extra);
}
type Echoes = Arc<Mutex<Vec<(String, Vec<i64>, SeqTracker)>>>;
type SentTimes = Arc<Vec<AtomicI64>>;
fn clock_service(
listener: TcpListener,
sent: SentTimes,
echoes: Echoes,
active: Arc<AtomicUsize>,
epoch: Instant,
) {
for stream in listener.incoming().flatten() {
let sent = sent.clone();
let echoes = echoes.clone();
let active = active.clone();
active.fetch_add(1, Ordering::SeqCst);
thread::spawn(move || {
let mut stream = stream;
stream.set_nodelay(true).ok();
let mut req = [0u8; 24];
let mut mine: Vec<i64> = Vec::new();
let mut seqs = SeqTracker::new(sent.len().saturating_sub(1) as u64);
let mut name = String::new();
while stream.read_exact(&mut req).is_ok() {
let kind = u64::from_le_bytes(req[0..8].try_into().unwrap());
let a = u64::from_le_bytes(req[8..16].try_into().unwrap());
match kind {
0 => {
let t2 = wall_ns();
let mut reply = [0u8; 24];
reply[0..8].copy_from_slice(&a.to_le_bytes());
reply[8..16].copy_from_slice(&t2.to_le_bytes());
reply[16..24].copy_from_slice(&wall_ns().to_le_bytes());
if stream.write_all(&reply).is_err() {
break;
}
}
1 => {
let now = epoch.elapsed().as_nanos() as i64;
if name.is_empty() {
name = String::from_utf8_lossy(&req[16..24])
.trim_end_matches('\0')
.to_string();
}
if seqs.record(a) {
let t = sent[a as usize].load(Ordering::Acquire);
if t != 0 {
mine.push(now - t);
}
}
}
_ => {}
}
}
if seqs.unique() + seqs.duplicates() + seqs.out_of_range() > 0 {
echoes.lock().unwrap().push((name, mine, seqs));
}
active.fetch_sub(1, Ordering::SeqCst);
});
}
}
fn connect_bounded(addr: &str, timeout: Duration) -> std::io::Result<TcpStream> {
let target = addr.to_socket_addrs()?.next().ok_or_else(|| {
std::io::Error::new(
std::io::ErrorKind::NotFound,
format!("no address for {}", addr),
)
})?;
let stream = TcpStream::connect_timeout(&target, timeout)?;
stream.set_read_timeout(Some(timeout))?;
stream.set_write_timeout(Some(timeout))?;
Ok(stream)
}
fn clock_offset(stream: &mut TcpStream, probes: usize) -> std::io::Result<(i64, i64)> {
let mut best_rtt = i64::MAX;
let mut best_offset = 0i64;
let mut reply = [0u8; 24];
let mut req = [0u8; 24];
let budget = Instant::now() + Duration::from_secs(2);
for _ in 0..probes {
if Instant::now() > budget {
break;
}
let t1 = wall_ns();
req[8..16].copy_from_slice(&t1.to_le_bytes());
stream.write_all(&req)?;
stream.read_exact(&mut reply)?;
let t4 = wall_ns();
let t2 = i64::from_le_bytes(reply[8..16].try_into().unwrap());
let t3 = i64::from_le_bytes(reply[16..24].try_into().unwrap());
let rtt = (t4 - t1) - (t3 - t2);
if rtt < best_rtt {
best_rtt = rtt;
best_offset = ((t2 - t1) + (t3 - t4)) / 2;
}
thread::sleep(Duration::from_micros(200));
}
if best_rtt == i64::MAX {
return Err(std::io::Error::new(
std::io::ErrorKind::TimedOut,
"no clock probe completed",
));
}
Ok((best_offset, best_rtt))
}
fn main() {
let args: Vec<String> = std::env::args().collect();
let mut role = String::new();
let mut bind = String::from("0.0.0.0:7400");
let mut clock = String::from("0.0.0.0:7402");
let mut source = String::new();
let mut ring: Option<PathBuf> = None;
let mut name = String::from("slave");
let mut rate = 1000u64;
let mut seconds = 5u64;
let mut burst = 1u64;
let mut warmup_ms = 3000u64;
let mut drain_secs = 30u64;
let mut multicast: Option<SocketAddrV4> = None;
let mut iface = Ipv4Addr::UNSPECIFIED;
let mut udp_port: Option<u16> = None;
let mut dup = 1u8;
let mut unicast = false;
let mut i = 1;
while i + 1 < args.len() {
let v = args[i + 1].as_str();
match args[i].as_str() {
"--role" => role = v.to_string(),
"--bind" => bind = v.to_string(),
"--clock" => clock = v.to_string(),
"--source" => source = v.to_string(),
"--ring" => ring = Some(PathBuf::from(v)),
"--name" => name = v.to_string(),
"--rate" => rate = parse_arg("--rate", v),
"--seconds" => seconds = parse_arg("--seconds", v),
"--burst" => burst = parse_arg::<u64>("--burst", v).max(1),
"--warmup-ms" => warmup_ms = parse_arg("--warmup-ms", v),
"--drain-secs" => drain_secs = parse_arg("--drain-secs", v),
"--multicast" => multicast = Some(parse_arg("--multicast", v)),
"--iface" => iface = parse_arg("--iface", v),
"--udp" => udp_port = Some(parse_arg("--udp", v)),
"--dup" => dup = parse_arg("--dup", v),
"--unicast" => unicast = parse_flag(v),
_ => {}
}
i += 2;
}
let ring = ring.unwrap_or_else(|| {
PathBuf::from(format!(
"/dev/shm/ringfire_stages_{}_{}.shm",
role,
std::process::id()
))
});
match role.as_str() {
"master" => {
assert!(rate > 0, "--rate must be positive");
let _ = std::fs::remove_file(&ring);
let mut producer = RingProducer::<Msg>::create(&ring, 1 << 18).unwrap();
let mut server = ReplicaServer::bind(&ring, bind.as_str())
.unwrap()
.spin(true);
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 total = rate * seconds;
let epoch = Instant::now();
let sent: SentTimes = Arc::new((0..=total).map(|_| AtomicI64::new(0)).collect());
let echoes: Echoes = Arc::new(Mutex::new(Vec::new()));
let active = Arc::new(AtomicUsize::new(0));
let listener = TcpListener::bind(clock.as_str()).unwrap();
let (sent_c, echoes_c, active_c) = (sent.clone(), echoes.clone(), active.clone());
thread::spawn(move || clock_service(listener, sent_c, echoes_c, active_c, epoch));
eprintln!(
"master: ring {} served on {}, clock on {}",
ring.display(),
bind,
clock
);
let mut local = RingConsumer::<Msg>::attach(&ring).unwrap();
let reader = thread::spawn(move || {
let mut samples = Vec::with_capacity(total as usize);
let mut seqs = SeqTracker::new(total);
let mut last_seen = Instant::now();
loop {
if let Some(msg) = local.try_recv() {
let now = epoch.elapsed().as_nanos() as i64;
if seqs.record(msg.seq) {
samples.push(now - msg.sent_mono_ns as i64);
}
last_seen = Instant::now();
if msg.seq >= total {
break;
}
} else {
if last_seen.elapsed() > Duration::from_secs(30) {
break;
}
core::hint::spin_loop();
}
}
(samples, seqs, local.lapped_count())
});
thread::sleep(Duration::from_millis(warmup_ms));
let mut pace = Pacer::per_second(rate, burst);
for seq in 1..=total {
if (seq - 1) % burst == 0 {
pace.wait_next();
}
let mono = epoch.elapsed().as_nanos() as i64;
sent[seq as usize].store(mono, Ordering::Release);
producer.push(&Msg {
sent_ns: wall_ns() as u64,
seq,
sent_mono_ns: mono as u64,
_pad: [0; 32],
});
}
let (mut samples, seqs, lapped) = reader.join().unwrap();
println!(
"master pushed {} at {} msg/s (bursts of {}), {} B slots",
total,
rate,
burst,
std::mem::size_of::<Msg>() + 8
);
report(
"stage master-local-consumer (push -> read, same host)",
&mut samples,
&format!("lapped={} {}", lapped, seqs.summary(1, total)),
);
let deadline = Instant::now() + Duration::from_secs(drain_secs);
while active.load(Ordering::SeqCst) > 0 && Instant::now() < deadline {
thread::sleep(Duration::from_millis(10));
}
let still_connected = active.load(Ordering::SeqCst);
if still_connected > 0 {
eprintln!(
"master: {} slave connection(s) still open after {} s; their round trips are not reported",
still_connected, drain_secs
);
}
let mut echoes = echoes.lock().unwrap();
echoes.sort_by(|a, b| a.0.cmp(&b.0));
for (name, rtt, seqs) in echoes.iter_mut() {
report(
&format!(
"stage round-trip via slave {} (push on master -> read on slave -> echo, master clock)",
name
),
rtt,
&seqs.summary(1, total),
);
}
drop(producer);
let _ = std::fs::remove_file(&ring);
}
"slave" => {
let mut echo = connect_bounded(clock.as_str(), Duration::from_secs(10)).unwrap();
echo.set_nodelay(true).unwrap();
let (offset, sync_rtt) =
clock_offset(&mut echo, 400).expect("initial clock synchronization failed");
let _ = std::fs::remove_file(&ring);
let mut mirror = Mirror::builder()
.start(MirrorStart::Latest)
.spin(true)
.interface(iface)
.unicast(unicast)
.connect(source.as_str(), &ring)
.unwrap();
let transport = if mirror.is_unicast() {
"udp unicast"
} else if mirror.is_multicast() {
"multicast"
} else {
"tcp"
};
let mut echo_req = [0u8; 24];
echo_req[0..8].copy_from_slice(&1u64.to_le_bytes());
for (i, b) in name.bytes().take(8).enumerate() {
echo_req[16 + i] = b;
}
let handle = mirror.handle().unwrap();
let mirror_thread = thread::spawn(move || {
let result = mirror.run();
(
result,
mirror.sequence(),
mirror.naks(),
mirror.retransmitted(),
mirror.gaps(),
)
});
let mut consumer = RingConsumer::<Msg>::attach(&ring).unwrap();
let mut samples: Vec<i64> = Vec::with_capacity(1 << 20);
let mut seqs = SeqTracker::new(1 << 28);
let mut echo_errors = 0u64;
let mut first = None;
let mut last_seen = Instant::now();
let start = Instant::now();
loop {
if let Some(msg) = consumer.try_recv() {
let now = wall_ns();
if seqs.record(msg.seq) {
samples.push(now + offset - msg.sent_ns as i64);
}
last_seen = Instant::now();
first.get_or_insert(last_seen);
echo_req[8..16].copy_from_slice(&msg.seq.to_le_bytes());
if echo.write_all(&echo_req).is_err() {
echo_errors += 1;
}
} else {
if (first.is_some() && last_seen.elapsed() > Duration::from_secs(2))
|| start.elapsed() > Duration::from_secs(seconds + 40)
{
break;
}
core::hint::spin_loop();
}
}
let (offset_end, sync_rtt_end) = if echo_errors == 0 {
clock_offset(&mut echo, 400).expect("final clock synchronization failed")
} else {
(offset, sync_rtt)
};
handle.shutdown().unwrap();
let (result, seq, naks, retransmitted, gaps) = mirror_thread.join().unwrap();
report(
&format!(
"stage slave {} ({}, push on master -> read on slave)",
name, transport
),
&mut samples,
&format!(
"{} (between first and last seen) last_seq={} mirror_seq={} lapped={} naks={} retransmitted={} gaps={} echo_errors={} clock_offset_us={:.1} sync_rtt_us={:.1} offset_change_us={:.1} (sync_rtt_end_us={:.1})",
seqs.summary(seqs.lowest(), seqs.highest()),
seqs.highest(),
seq,
consumer.lapped_count(),
naks,
retransmitted,
gaps,
echo_errors,
offset as f64 / 1000.0,
sync_rtt as f64 / 1000.0,
(offset_end - offset) as f64 / 1000.0,
sync_rtt_end as f64 / 1000.0
),
);
drop(consumer);
let _ = std::fs::remove_file(&ring);
if let Err(e) = result {
eprintln!("slave {}: mirror failed: {}", name, e);
std::process::exit(1);
}
}
other => panic!("unknown role {}", other),
}
}