use std::io::Write;
use std::sync::{Arc, Mutex};
use std::time::Duration;
use zenkey::qos::QosProfile;
use zenkey_fleet::{
RecordBounds, ReplayTarget, ZREC_VERSION, ZrecHeader, ZrecSink, ZrecSource,
declare_publication, record, replay,
};
mod util;
use util::peer_pair;
#[derive(Clone, Default)]
struct SharedBuf(Arc<Mutex<Vec<u8>>>);
impl SharedBuf {
fn take(&self) -> Vec<u8> {
self.0.lock().expect("buffer lock").clone()
}
}
impl Write for SharedBuf {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
self.0.lock().expect("buffer lock").write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
#[derive(Clone)]
struct SlowBuf {
inner: SharedBuf,
per_write: Duration,
}
impl Write for SlowBuf {
fn write(&mut self, buf: &[u8]) -> std::io::Result<usize> {
std::thread::sleep(self.per_write);
self.inner.write(buf)
}
fn flush(&mut self) -> std::io::Result<()> {
Ok(())
}
}
fn header(selector: &str) -> ZrecHeader {
ZrecHeader {
zrec: ZREC_VERSION,
selectors: vec![selector.to_string()],
base: String::new(),
captured_at: "2026-08-12T00:00:00Z".to_string(),
}
}
const KEY: &str = "v1/h-aaaaaaaaaaaa/state/demo/health";
const SELECTOR: &str = "v1/h-aaaaaaaaaaaa/state/demo/**";
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_capture_replays_onto_a_second_bus_intact() {
let (a, b) = peer_pair().await;
let monitor = zenkey_fleet::Monitor::start(&b, zenkey_fleet::MonitorSpec::default())
.await
.expect("monitor");
let mut events = monitor.events();
monitor.watch(SELECTOR).await.expect("watch");
let publication =
declare_publication(&a, KEY, QosProfile::Transition, Some("application/json"))
.await
.expect("declare");
let matching = publication.matching_events().await.expect("matching");
assert!(
tokio::time::timeout(util::SETTLE, matching.recv())
.await
.expect("matching within 5s")
.expect("listener alive")
);
let binary: Vec<u8> = vec![0x00, 0xff, 0x01, 0xfe, 0x80];
publication
.send(br#"{"ok":true}"#.to_vec(), Some(b"who=test".to_vec()))
.await
.expect("send json");
publication
.send(binary.clone(), None)
.await
.expect("send binary");
publication.retire().await.expect("retire");
let buf = SharedBuf::default();
let sink = ZrecSink::spawn(buf.clone(), &header(SELECTOR))
.await
.expect("sink");
record(
&mut events,
&sink,
RecordBounds {
max_samples: Some(3),
max_duration: Some(Duration::from_secs(10)),
},
|_, _| {},
)
.await
.expect("record");
let (samples, dropped) = sink.finish().await.expect("finish");
assert_eq!(samples, 3, "two puts and a tombstone");
assert_eq!(dropped, 0);
let file = buf.take();
let (c, d) = peer_pair().await;
let replay_monitor = zenkey_fleet::Monitor::start(&d, zenkey_fleet::MonitorSpec::default())
.await
.expect("replay monitor");
let mut replayed = replay_monitor.events();
replay_monitor.watch(SELECTOR).await.expect("watch replay");
let gate = declare_publication(&c, KEY, QosProfile::Transition, None)
.await
.expect("gate");
let gate_matching = gate.matching_events().await.expect("gate matching");
assert!(
tokio::time::timeout(util::SETTLE, gate_matching.recv())
.await
.expect("gate matching within 5s")
.expect("listener alive")
);
gate.undeclare().await.expect("undeclare gate");
let mut reader = ZrecSource::spawn(std::io::Cursor::new(file))
.await
.expect("reader");
let report = replay(
&mut reader,
zenkey_fleet::ReplaySpec {
target: ReplayTarget::Bus {
session: &c,
slices: None,
},
speed: 1000.0,
i_know: false,
default_qos: zenkey::qos::QosProfile::Refreshed,
},
|_| {},
)
.await
.expect("replay");
assert_eq!(report.published, 2);
assert_eq!(report.tombstones, 1, "state-shaped delete needs no force");
assert_eq!(report.malformed, 0);
assert_eq!(report.refused, 0);
let mut views = Vec::new();
while views.len() < 3 {
let item = tokio::time::timeout(util::SETTLE, replayed.recv())
.await
.expect("replayed event within 5s")
.expect("stream alive");
if let zenkey_fleet::StreamItem::Event(zenkey_fleet::FleetEvent::Sample(s)) = item {
views.push(s);
}
}
assert!(views.iter().all(|v| v.key == KEY));
assert_eq!(views[0].payload.to_bytes().as_ref(), br#"{"ok":true}"#);
assert_eq!(views[0].encoding, "application/json");
assert_eq!(
views[0]
.attachment
.as_ref()
.expect("attachment survives the file")
.to_bytes()
.as_ref(),
b"who=test"
);
assert_eq!(
views[1].payload.to_bytes().as_ref(),
binary.as_slice(),
"binary payload is lossless through the bytes field"
);
assert_eq!(views[2].kind, zenoh::sample::SampleKind::Delete);
assert!(
views[0].qos_matches(QosProfile::Transition),
"the recorded profile name declared the replay publisher"
);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_lossy_capture_says_so_at_both_ends() {
let core = zenkey_fleet::MonitorCore::new(4);
let mut events = core.events();
for i in 0..32u8 {
core.ingest(
zenkey_fleet::SampleView {
key: KEY.to_string(),
payload: zenoh::bytes::ZBytes::from(vec![i]),
encoding: String::new(),
kind: zenoh::sample::SampleKind::Put,
timestamp: None,
stamped_by: None,
attachment: None,
priority: zenoh::qos::Priority::Data,
congestion_control: zenoh::qos::CongestionControl::Drop,
reliability: zenoh::qos::Reliability::BestEffort,
express: false,
source: None,
received: std::time::Instant::now(),
},
None,
);
}
let buf = SharedBuf::default();
let sink = ZrecSink::spawn(buf.clone(), &header(SELECTOR))
.await
.expect("sink");
record(
&mut events,
&sink,
RecordBounds {
max_samples: None,
max_duration: Some(Duration::from_secs(1)),
},
|_, _| {},
)
.await
.expect("record");
let (samples, dropped) = sink.finish().await.expect("finish");
assert!(dropped > 0, "a capacity-4 channel under 32 sends must lag");
assert!(samples > 0);
let file = buf.take();
let text = String::from_utf8(file.clone()).expect("a .zrec is text");
assert!(
text.lines().any(|l| l.contains("\"dropped\"")),
"the drop ledger is in the file, not only in memory"
);
let mut reader = ZrecSource::spawn(std::io::Cursor::new(file))
.await
.expect("reader");
let report = replay(
&mut reader,
zenkey_fleet::ReplaySpec {
target: ReplayTarget::DryRun,
speed: 1.0,
i_know: false,
default_qos: zenkey::qos::QosProfile::Refreshed,
},
|_| {},
)
.await
.expect("dry replay");
assert_eq!(
report.capture_dropped, dropped,
"the ledger survives replay"
);
assert_eq!(u64::from(report.dry_run), 1);
assert_eq!(report.published, samples);
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_slow_writer_does_not_become_the_captures_drops() {
const SAMPLES: u64 = 1_500;
const BURST: u64 = 150;
let core = zenkey_fleet::MonitorCore::new(1024); let mut events = core.events();
let buf = SharedBuf::default();
let sink = ZrecSink::spawn(
SlowBuf {
inner: buf.clone(),
per_write: Duration::from_micros(100),
},
&header(SELECTOR),
)
.await
.expect("sink");
let producer = {
let core = std::sync::Arc::clone(&core);
tokio::spawn(async move {
for i in 0..SAMPLES {
if i > 0 && i % BURST == 0 {
tokio::time::sleep(Duration::from_millis(5)).await;
}
core.ingest(
zenkey_fleet::SampleView {
key: KEY.to_string(),
payload: zenoh::bytes::ZBytes::from(i.to_le_bytes().to_vec()),
encoding: String::new(),
kind: zenoh::sample::SampleKind::Put,
timestamp: None,
stamped_by: None,
attachment: None,
priority: zenoh::qos::Priority::Data,
congestion_control: zenoh::qos::CongestionControl::Drop,
reliability: zenoh::qos::Reliability::BestEffort,
express: false,
source: None,
received: std::time::Instant::now(),
},
None,
);
tokio::task::yield_now().await;
}
})
};
record(
&mut events,
&sink,
RecordBounds {
max_samples: Some(SAMPLES),
max_duration: Some(Duration::from_secs(30)),
},
|_, _| {},
)
.await
.expect("record");
producer.await.expect("producer");
let draining = std::time::Instant::now();
let (samples, dropped) = sink.finish().await.expect("finish");
let drain_took = draining.elapsed();
assert_eq!(dropped, 0, "the writer's latency is not the bus's loss");
assert_eq!(samples, SAMPLES, "every sample reached the file");
assert!(
drain_took > Duration::from_millis(20),
"the writer kept up on its own, so this proves nothing: {drain_took:?}"
);
let text = String::from_utf8(buf.take()).expect("a .zrec is text");
assert!(
!text.lines().any(|l| l.contains("\"dropped\"")),
"no self-inflicted drop record"
);
}