use std::collections::HashMap;
use std::io::{BufRead, Write};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, Instant};
use crate::{Error, Result};
use zenkey::qos::QosProfile;
use zenoh::Session;
use zenoh::sample::SampleKind;
use crate::bus::monitor::{EventStream, FleetEvent, SampleView, StreamItem};
use crate::model::registry::SliceSet;
use crate::report::{ReplayReport, SampleRow, Transition, ZrecHeader};
use crate::tape::ingest::{IngestRow, parse_row};
pub const ZREC_VERSION: u32 = 2;
pub const ZREC_READS: [u32; 2] = [1, 2];
pub const PREAMBLE_SKIP_REASON: &str = "state at capture start; re-stamping it republishes a \
snapshot over live state (RFC 13 §4.2) — pass --seed-state to mean it";
pub fn rfc3339_now() -> String {
rfc3339_from_unix(
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_secs())
.unwrap_or(0),
)
}
pub fn rfc3339_from_unix(secs: u64) -> String {
let (days, rem) = (secs / 86_400, secs % 86_400);
let (h, m, s) = (rem / 3600, (rem % 3600) / 60, rem % 60);
let z = days as i64 + 719_468;
let era = z.div_euclid(146_097);
let doe = z.rem_euclid(146_097);
let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
let y = yoe + era * 400;
let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
let mp = (5 * doy + 2) / 153;
let d = doy - (153 * mp + 2) / 5 + 1;
let mo = if mp < 10 { mp + 3 } else { mp - 9 };
let y = if mo <= 2 { y + 1 } else { y };
format!("{y:04}-{mo:02}-{d:02}T{h:02}:{m:02}:{s:02}Z")
}
pub struct ZrecWriter<W: Write> {
out: W,
epoch: Instant,
counts: SinkCounts,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct SinkCounts {
pub samples: u64,
pub dropped: u64,
pub preamble: u64,
pub triggers: u64,
}
impl<W: Write> ZrecWriter<W> {
pub fn new(out: W, header: &ZrecHeader) -> Result<Self> {
ZrecWriter::new_at(out, header, Instant::now())
}
pub fn new_at(mut out: W, header: &ZrecHeader, epoch: Instant) -> Result<Self> {
serde_json::to_writer(&mut out, header).map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e.into(),
})?;
out.write_all(b"\n").map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e,
})?;
Ok(ZrecWriter {
out,
epoch,
counts: SinkCounts::default(),
})
}
fn line(&mut self, line: &str) -> Result<()> {
self.out.write_all(line.as_bytes()).map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e,
})?;
self.out.write_all(b"\n").map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e,
})
}
fn row_of(view: &SampleView) -> SampleRow {
let mut row = SampleRow {
key: view.key.clone(),
..SampleRow::default()
}
.with_wire(view);
if view.kind != SampleKind::Delete {
row = row.with_payload_bytes(&view.payload.to_bytes());
}
if let Some(a) = &view.attachment {
row.attachment_b64 = Some(crate::tape::ingest::b64(&a.to_bytes()));
}
row
}
pub fn write_sample(&mut self, view: &SampleView) -> Result<()> {
let t_us = u64::try_from(
view.received
.saturating_duration_since(self.epoch)
.as_micros(),
)
.unwrap_or(u64::MAX);
let mut row = ZrecWriter::<W>::row_of(view);
row.t = Some(t_us);
self.line(&row.to_line())?;
self.counts.samples += 1;
Ok(())
}
pub fn write_preamble(&mut self, view: &SampleView) -> Result<()> {
let mut row = ZrecWriter::<W>::row_of(view);
row.t = Some(0);
row.preamble = Some(true);
self.line(&row.to_line())?;
self.counts.preamble += 1;
Ok(())
}
pub fn write_dropped(&mut self, n: u64) -> Result<()> {
let line =
serde_json::to_string(&serde_json::json!({ "dropped": n })).map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e.into(),
})?;
self.line(&line)?;
self.counts.dropped += n;
Ok(())
}
pub fn write_trigger(&mut self, transition: &Transition) -> Result<()> {
let line =
serde_json::to_string(&serde_json::json!({ "trigger": transition })).map_err(|e| {
Error::Io {
path: std::path::PathBuf::new(),
source: e.into(),
}
})?;
self.line(&line)?;
self.counts.triggers += 1;
Ok(())
}
pub fn counts(&self) -> SinkCounts {
self.counts
}
pub fn finish(mut self) -> Result<W> {
self.out.flush().map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e,
})?;
Ok(self.out)
}
}
const SINK_QUEUE: usize = 4096;
enum ZrecLine {
Sample(Arc<SampleView>),
Dropped(u64),
Preamble(Arc<SampleView>),
Trigger(Box<Transition>),
}
#[derive(Debug, Default)]
struct SinkState {
samples: AtomicU64,
dropped: AtomicU64,
preamble: AtomicU64,
triggers: AtomicU64,
failure: std::sync::Mutex<Option<String>>,
}
pub struct ZrecSink {
tx: tokio::sync::mpsc::Sender<ZrecLine>,
state: Arc<SinkState>,
writer: tokio::task::JoinHandle<Result<SinkCounts>>,
}
impl ZrecSink {
pub async fn spawn<W: Write + Send + 'static>(out: W, header: &ZrecHeader) -> Result<ZrecSink> {
ZrecSink::spawn_at(out, header, Instant::now()).await
}
pub async fn spawn_at<W: Write + Send + 'static>(
out: W,
header: &ZrecHeader,
epoch: Instant,
) -> Result<ZrecSink> {
let (tx, mut rx) = tokio::sync::mpsc::channel(SINK_QUEUE);
let (ready, opened) = tokio::sync::oneshot::channel();
let state = Arc::new(SinkState::default());
let header = header.clone();
let task_state = Arc::clone(&state);
let writer = tokio::task::spawn_blocking(move || {
let mut writer = match ZrecWriter::new_at(out, &header, epoch) {
Ok(w) => {
let _ = ready.send(None);
w
}
Err(e) => {
let _ = ready.send(Some(crate::one_line(&e)));
return Err(e);
}
};
while let Some(line) = rx.blocking_recv() {
let wrote = match line {
ZrecLine::Sample(view) => writer.write_sample(&view),
ZrecLine::Dropped(n) => writer.write_dropped(n),
ZrecLine::Preamble(view) => writer.write_preamble(&view),
ZrecLine::Trigger(t) => writer.write_trigger(&t),
};
if let Err(e) = wrote {
*task_state.failure.lock().expect("sink failure lock") =
Some(crate::one_line(&e));
return Err(e);
}
}
let counts = writer.counts();
writer.finish().map(|_| counts)
});
match opened.await {
Ok(None) => Ok(ZrecSink { tx, state, writer }),
Ok(Some(reason)) => Err(Error::Io {
path: std::path::PathBuf::new(),
source: std::io::Error::other(reason),
}),
Err(_) => Err(Error::Internal(
"the .zrec writer stopped before it opened".into(),
)),
}
}
pub async fn write_sample(&self, view: Arc<SampleView>) -> Result<()> {
self.send(ZrecLine::Sample(view)).await?;
self.state.samples.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub async fn write_dropped(&self, n: u64) -> Result<()> {
self.send(ZrecLine::Dropped(n)).await?;
self.state.dropped.fetch_add(n, Ordering::Relaxed);
Ok(())
}
pub async fn write_preamble(&self, view: Arc<SampleView>) -> Result<()> {
self.send(ZrecLine::Preamble(view)).await?;
self.state.preamble.fetch_add(1, Ordering::Relaxed);
Ok(())
}
pub async fn write_trigger(&self, transition: Transition) -> Result<()> {
self.send(ZrecLine::Trigger(Box::new(transition))).await?;
self.state.triggers.fetch_add(1, Ordering::Relaxed);
Ok(())
}
async fn send(&self, line: ZrecLine) -> Result<()> {
if self.tx.send(line).await.is_ok() {
return Ok(());
}
let failure = self
.state
.failure
.lock()
.expect("sink failure lock")
.clone();
Err(Error::Internal(
failure.unwrap_or_else(|| "the .zrec writer stopped".to_string()),
))
}
pub fn counts(&self) -> SinkCounts {
SinkCounts {
samples: self.state.samples.load(Ordering::Relaxed),
dropped: self.state.dropped.load(Ordering::Relaxed),
preamble: self.state.preamble.load(Ordering::Relaxed),
triggers: self.state.triggers.load(Ordering::Relaxed),
}
}
pub async fn finish(self) -> Result<SinkCounts> {
let ZrecSink { tx, state, writer } = self;
drop(tx);
drop(state);
writer
.await
.map_err(|e| Error::Internal(format!("the .zrec writer panicked: {e}")))?
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct RecordBounds {
pub max_samples: Option<u64>,
pub max_duration: Option<Duration>,
}
pub async fn record(
events: &mut EventStream,
sink: &ZrecSink,
bounds: RecordBounds,
mut on_progress: impl FnMut(u64, u64),
) -> Result<()> {
let deadline = bounds.max_duration.map(|d| Instant::now() + d);
loop {
let samples = sink.counts().samples;
if bounds.max_samples.is_some_and(|max| samples >= max) {
return Ok(());
}
let item = match deadline {
Some(d) => {
let left = d.saturating_duration_since(Instant::now());
if left.is_zero() {
return Ok(());
}
match tokio::time::timeout(left, events.recv()).await {
Ok(item) => item,
Err(_) => return Ok(()),
}
}
None => events.recv().await,
};
match item {
Some(StreamItem::Event(FleetEvent::Sample(view))) => {
sink.write_sample(view).await?;
}
Some(StreamItem::Dropped(n)) => {
sink.write_dropped(n).await?;
}
Some(_) => continue,
None => return Ok(()),
}
let counts = sink.counts();
on_progress(counts.samples, counts.dropped);
}
}
#[derive(Debug, Clone)]
pub enum ZrecItem {
Sample {
row: IngestRow,
t_us: Option<u64>,
timestamp: Option<String>,
source: Option<String>,
},
Dropped(u64),
Preamble {
row: IngestRow,
timestamp: Option<String>,
},
Trigger(Box<Transition>),
}
pub struct ZrecReader<R: BufRead> {
header: ZrecHeader,
lines: std::io::Lines<R>,
line: u64,
}
impl<R: BufRead> ZrecReader<R> {
pub fn new(source: R) -> Result<Self> {
let mut lines = source.lines();
let first = lines
.next()
.ok_or_else(|| Error::malformed(".zrec", "empty file — no header line"))?
.map_err(|e| Error::Io {
path: std::path::PathBuf::new(),
source: e,
})?;
let header: ZrecHeader = serde_json::from_str(&first)
.map_err(|e| Error::malformed_with(".zrec line 1", "is not a header", e))?;
if !ZREC_READS.contains(&header.zrec) {
return Err(Error::malformed(
".zrec",
format!(
"unsupported version {} (this reader speaks {ZREC_VERSION} and reads {})",
header.zrec,
ZREC_READS
.iter()
.filter(|v| **v != ZREC_VERSION)
.map(u32::to_string)
.collect::<Vec<_>>()
.join(", ")
),
));
}
Ok(ZrecReader {
header,
lines,
line: 1,
})
}
pub fn header(&self) -> &ZrecHeader {
&self.header
}
#[allow(clippy::should_implement_trait)] pub fn next(&mut self) -> Option<std::result::Result<ZrecItem, String>> {
loop {
let line = match self.lines.next()? {
Ok(l) => l,
Err(e) => {
self.line += 1;
return Some(Err(format!("line {}: read: {e}", self.line)));
}
};
self.line += 1;
if line.trim().is_empty() {
continue;
}
if let Ok(v) = serde_json::from_str::<serde_json::Value>(&line)
&& v.get("key").is_none()
{
if let Some(n) = v.get("dropped").and_then(serde_json::Value::as_u64) {
return Some(Ok(ZrecItem::Dropped(n)));
}
if let Some(t) = v.get("trigger") {
return Some(
serde_json::from_value::<Transition>(t.clone())
.map(|t| ZrecItem::Trigger(Box::new(t)))
.map_err(|e| format!("line {}: trigger record: {e}", self.line)),
);
}
}
return Some(match parse_row(&line) {
Ok(row) => {
let v: serde_json::Value = serde_json::from_str(&line).unwrap_or_default();
let timestamp = v
.get("timestamp")
.and_then(serde_json::Value::as_str)
.map(str::to_string);
if v.get("preamble").and_then(serde_json::Value::as_bool) == Some(true) {
Ok(ZrecItem::Preamble { row, timestamp })
} else {
Ok(ZrecItem::Sample {
row,
t_us: v.get("t").and_then(serde_json::Value::as_u64),
timestamp,
source: v
.get("source")
.and_then(serde_json::Value::as_str)
.map(str::to_string),
})
}
}
Err(e) => Err(format!("line {}: {e}", self.line)),
});
}
}
}
pub struct ZrecSource {
header: ZrecHeader,
rx: tokio::sync::mpsc::Receiver<std::result::Result<ZrecItem, String>>,
}
impl ZrecSource {
pub async fn spawn<R: BufRead + Send + 'static>(source: R) -> Result<ZrecSource> {
let (tx, rx) = tokio::sync::mpsc::channel(SINK_QUEUE);
let (ready, opened) = tokio::sync::oneshot::channel();
tokio::task::spawn_blocking(move || {
let mut reader = match ZrecReader::new(source) {
Ok(r) => r,
Err(e) => {
let _ = ready.send(Err(e));
return;
}
};
if ready.send(Ok(reader.header().clone())).is_err() {
return;
}
while let Some(item) = reader.next() {
if tx.blocking_send(item).is_err() {
return;
}
}
});
match opened.await {
Ok(header) => Ok(ZrecSource {
header: header?,
rx,
}),
Err(_) => Err(Error::Internal(
"the .zrec reader stopped before it opened".into(),
)),
}
}
pub fn header(&self) -> &ZrecHeader {
&self.header
}
pub async fn next(&mut self) -> Option<std::result::Result<ZrecItem, String>> {
self.rx.recv().await
}
}
pub enum ReplayTarget<'a> {
DryRun,
Bus {
session: &'a Session,
slices: Option<&'a SliceSet>,
},
}
pub struct ReplaySpec<'a> {
pub target: ReplayTarget<'a>,
pub speed: f64,
pub i_know: bool,
pub default_qos: QosProfile,
pub seed_state: bool,
}
#[derive(Debug, Clone)]
pub enum ReplayEvent<'a> {
WouldPut {
key: &'a str,
bytes: usize,
encoding: Option<&'a str>,
},
WouldRetire { key: &'a str },
Malformed { reason: String },
Refused { key: String, reason: String },
CaptureDropped(u64),
PreambleSkipped { key: &'a str, reason: &'static str },
Trigger(&'a Transition),
}
#[derive(Default)]
struct Publications(HashMap<String, crate::bus::write::Publication>);
impl std::ops::Deref for Publications {
type Target = HashMap<String, crate::bus::write::Publication>;
fn deref(&self) -> &Self::Target {
&self.0
}
}
impl std::ops::DerefMut for Publications {
fn deref_mut(&mut self) -> &mut Self::Target {
&mut self.0
}
}
impl Publications {
async fn close(mut self) -> Result<()> {
crate::bus::teardown::drain_undeclare(self.0.drain().collect(), |p| {
crate::bus::write::Publication::undeclare(p)
})
.await
}
}
impl Drop for Publications {
fn drop(&mut self) {
if self.0.is_empty() {
return;
}
let Ok(runtime) = tokio::runtime::Handle::try_current() else {
return;
};
let declared: Vec<(String, crate::bus::write::Publication)> = self.0.drain().collect();
runtime.spawn(async move {
for (key, publication) in declared {
if let Err(e) = publication.undeclare().await {
tracing::warn!(key = %key, "undeclare after a cancelled replay: {e}");
}
}
});
}
}
pub async fn replay(
reader: &mut ZrecSource,
spec: ReplaySpec<'_>,
mut on_event: impl FnMut(ReplayEvent<'_>),
) -> Result<ReplayReport> {
let ReplaySpec {
target,
speed,
i_know,
default_qos,
seed_state,
} = spec;
if !(speed.is_finite() && speed > 0.0) {
return Err(Error::unaskable(
"--speed",
format!("must be a positive number (got {speed})"),
));
}
let base = reader.header().base.clone();
let mut report = ReplayReport {
header: reader.header().clone(),
dry_run: matches!(target, ReplayTarget::DryRun),
speed,
published: 0,
tombstones: 0,
malformed: 0,
refused: 0,
capture_dropped: 0,
first_errors: Vec::new(),
preamble_skipped: 0,
preamble_seeded: 0,
triggers: 0,
};
let record_err = |report: &mut ReplayReport, reason: String, refused: bool| {
if refused {
report.refused += 1;
} else {
report.malformed += 1;
}
if report.first_errors.len() < 3 {
report.first_errors.push(reason);
}
};
let mut publications = Publications::default();
let mut prev_t: Option<u64> = None;
let mut fatal: Option<Error> = None;
while let Some(item) = reader.next().await {
let (row, t_us, seeding) = match item {
Ok(ZrecItem::Sample { row, t_us, .. }) => (row, t_us, false),
Ok(ZrecItem::Dropped(n)) => {
report.capture_dropped += n;
on_event(ReplayEvent::CaptureDropped(n));
continue;
}
Ok(ZrecItem::Trigger(t)) => {
report.triggers += 1;
on_event(ReplayEvent::Trigger(&t));
continue;
}
Ok(ZrecItem::Preamble { row, .. }) => {
if !seed_state {
report.preamble_skipped += 1;
on_event(ReplayEvent::PreambleSkipped {
key: &row.key,
reason: PREAMBLE_SKIP_REASON,
});
continue;
}
(row, Some(0), true)
}
Err(reason) => {
on_event(ReplayEvent::Malformed {
reason: reason.clone(),
});
record_err(&mut report, reason, false);
continue;
}
};
let slices = match &target {
ReplayTarget::Bus { slices, .. } => *slices,
ReplayTarget::DryRun => None,
};
if row.delete
&& let Err(e) = crate::bus::write::check_retire(&base, &row.key, slices, i_know)
{
let reason = e.to_string();
on_event(ReplayEvent::Refused {
key: row.key.clone(),
reason: reason.clone(),
});
record_err(&mut report, format!("{}: {reason}", row.key), true);
continue;
}
let count_put = |report: &mut ReplayReport, delete: bool| match (seeding, delete) {
(true, _) => report.preamble_seeded += 1,
(false, true) => report.tombstones += 1,
(false, false) => report.published += 1,
};
match &target {
ReplayTarget::DryRun => {
if row.delete {
on_event(ReplayEvent::WouldRetire { key: &row.key });
} else {
on_event(ReplayEvent::WouldPut {
key: &row.key,
bytes: row.payload.len(),
encoding: row.encoding.as_deref(),
});
}
count_put(&mut report, row.delete);
}
ReplayTarget::Bus { session, .. } => {
if let (Some(prev), Some(t)) = (prev_t, t_us)
&& t > prev
{
let delay = Duration::from_micros(t - prev).div_f64(speed);
tokio::time::sleep(delay).await;
}
if t_us.is_some() {
prev_t = t_us;
}
let publication = match publications.entry(row.key.clone()) {
std::collections::hash_map::Entry::Occupied(e) => e.into_mut(),
std::collections::hash_map::Entry::Vacant(e) => {
let qos = match &row.qos {
None => default_qos,
Some(name) => match zenkey::qos::QosProfile::from_name(name) {
Some(qos) => qos,
None => {
let reason = format!("unknown QoS profile {name:?}");
on_event(ReplayEvent::Malformed {
reason: reason.clone(),
});
record_err(&mut report, reason, false);
continue;
}
},
};
let publication = match crate::bus::write::declare_publication(
session,
&row.key,
qos,
row.encoding.as_deref(),
)
.await
{
Ok(p) => p,
Err(e) => {
fatal = Some(e);
break;
}
};
e.insert(publication)
}
};
let delete = row.delete;
let sent = if delete {
publication.retire().await
} else {
publication.send(row.payload, row.attachment).await
};
match sent {
Ok(()) => count_put(&mut report, delete),
Err(e) => {
fatal = Some(e);
break;
}
}
}
}
}
let closed = publications.close().await;
if let Some(e) = fatal {
return Err(e);
}
closed?;
Ok(report)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn the_wall_clock_formats_correctly() {
assert_eq!(rfc3339_from_unix(0), "1970-01-01T00:00:00Z");
assert_eq!(rfc3339_from_unix(951_782_400), "2000-02-29T00:00:00Z");
assert_eq!(rfc3339_from_unix(1_786_492_800), "2026-08-12T00:00:00Z");
assert!(!rfc3339_now().is_empty());
}
fn header() -> ZrecHeader {
ZrecHeader {
zrec: ZREC_VERSION,
selectors: vec!["v1/**".into()],
base: String::new(),
captured_at: "2026-08-12T00:00:00Z".into(),
preamble: None,
pre_roll: None,
}
}
async fn source_of(body: &str) -> ZrecSource {
ZrecSource::spawn(std::io::Cursor::new(body.as_bytes().to_vec()))
.await
.expect("a .zrec header")
}
#[test]
fn the_header_is_a_contract() {
let mut sink = Vec::new();
let writer = ZrecWriter::new(&mut sink, &header()).unwrap();
let _ = writer.finish().unwrap();
let reader = ZrecReader::new(sink.as_slice()).unwrap();
assert_eq!(reader.header(), &header());
let future = r#"{"zrec":3,"selectors":[],"base":"","captured_at":"x"}"#;
let err = ZrecReader::new(future.as_bytes())
.err()
.unwrap()
.to_string();
assert!(err.contains("version 3"), "{err}");
assert!(err.contains("speaks 2 and reads 1"), "{err}");
let not_zrec = r#"{"key":"v1/x","value":1}"#;
let err = ZrecReader::new(not_zrec.as_bytes())
.err()
.unwrap()
.to_string();
assert!(err.contains("header"), "{err}");
}
#[test]
fn a_version_one_body_reads_under_the_version_two_reader() {
let body = concat!(
r#"{"zrec":1,"selectors":["v1/**"],"base":"","captured_at":"2026-08-12T00:00:00Z"}"#,
"\n",
r#"{"key":"v1/h/state/p/a","t":0,"bytes":"AQ=="}"#,
"\n",
r#"{"dropped":2}"#,
"\n",
r#"{"key":"v1/h/state/p/a","t":1000,"delete":true}"#,
"\n",
);
let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
assert_eq!(reader.header().zrec, 1);
assert_eq!(reader.header().preamble, None);
assert_eq!(reader.header().pre_roll, None);
assert!(matches!(
reader.next(),
Some(Ok(ZrecItem::Sample { t_us: Some(0), .. }))
));
assert!(matches!(reader.next(), Some(Ok(ZrecItem::Dropped(2)))));
assert!(matches!(
reader.next(),
Some(Ok(ZrecItem::Sample {
t_us: Some(1000),
..
}))
));
assert!(reader.next().is_none());
}
#[test]
fn version_two_lines_read_back_by_kind() {
use crate::report::{CondState, PreRollInfo, PreambleInfo, PreambleSemantics};
let epoch = Instant::now();
let header = ZrecHeader {
preamble: Some(PreambleInfo {
count: 1,
collected_over_s: 0.25,
selectors: vec!["v1/*/state/**".into()],
semantics: PreambleSemantics::AbsentFromWindow,
incomplete: 0,
failed: vec!["v1/*/telemetry/**".into()],
}),
pre_roll: Some(PreRollInfo {
asked_s: 30.0,
covered_s: 12.5,
watched: vec!["v1/**".into()],
evicted: 0,
expired: 40,
}),
..header()
};
let stamp = zenoh::time::Timestamp::new(
zenoh::time::NTP64::from(Duration::from_secs(1_700_000_000)),
zenoh::time::TimestampId::try_from([7u8; 16]).unwrap(),
);
let fetched = crate::bus::monitor::SampleView {
key: "v1/h-0123456789ab/state/p/config".into(),
payload: zenoh::bytes::ZBytes::from(vec![9u8]),
encoding: String::new(),
kind: SampleKind::Put,
timestamp: Some(stamp),
stamped_by: None,
attachment: None,
priority: zenoh::qos::Priority::DEFAULT,
congestion_control: zenoh::qos::CongestionControl::DEFAULT,
reliability: zenoh::qos::Reliability::DEFAULT,
express: false,
source: None,
received: epoch + Duration::from_secs(5),
};
let fired = Transition {
rule: "silent-for v1/h-0123456789ab/state/p/health 0.7".into(),
from: Some(CondState::Ok),
to: CondState::Firing,
at: "2026-09-06T00:00:00Z".into(),
evidence: "no sample for 0.7s, on a drop-free observer".into(),
};
let mut sink = Vec::new();
let mut w = ZrecWriter::new_at(&mut sink, &header, epoch).unwrap();
w.write_preamble(&fetched).unwrap();
w.write_trigger(&fired).unwrap();
assert_eq!(
w.counts(),
SinkCounts {
samples: 0,
dropped: 0,
preamble: 1,
triggers: 1
},
"the kinds are counted apart"
);
let _ = w.finish().unwrap();
let text = String::from_utf8(sink.clone()).unwrap();
assert!(text.contains(r#""preamble":true"#), "{text}");
assert!(text.contains(r#""t":0"#), "{text}");
assert!(text.contains(r#"{"trigger":{"#), "{text}");
let mut reader = ZrecReader::new(sink.as_slice()).unwrap();
assert_eq!(reader.header(), &header);
match reader.next() {
Some(Ok(ZrecItem::Preamble { row, timestamp })) => {
assert_eq!(row.key, fetched.key);
assert_eq!(row.payload, vec![9u8]);
assert_eq!(timestamp.as_deref(), Some(stamp.to_string().as_str()));
}
other => panic!("expected a preamble row, got {other:?}"),
}
match reader.next() {
Some(Ok(ZrecItem::Trigger(t))) => assert_eq!(*t, fired),
other => panic!("expected the trigger record, got {other:?}"),
}
assert!(reader.next().is_none());
}
#[tokio::test]
async fn a_dry_run_skips_the_preamble_by_default_and_seeds_it_on_request() {
let body = format!(
"{}\n{}\n{}\n{}\n",
serde_json::to_string(&header()).unwrap(),
r#"{"key":"v1/h-0123456789ab/state/p/config","t":0,"preamble":true,"bytes":"CQ=="}"#,
r#"{"key":"v1/h-0123456789ab/state/p/health","t":250000,"bytes":"eyJvayI6dHJ1ZX0="}"#,
r#"{"trigger":{"rule":"silent-for k 0.7","from":"ok","to":"firing","at":"x","evidence":"e"}}"#,
);
let spec = |seed_state| ReplaySpec {
target: ReplayTarget::DryRun,
speed: 1.0,
i_know: false,
default_qos: QosProfile::Refreshed,
seed_state,
};
let mut reader = source_of(&body).await;
let mut skipped = Vec::new();
let mut triggers = 0;
let report = replay(&mut reader, spec(false), |ev| match ev {
ReplayEvent::PreambleSkipped { key, reason } => {
skipped.push((key.to_string(), reason));
}
ReplayEvent::Trigger(t) => {
assert_eq!(t.to, crate::report::CondState::Firing);
triggers += 1;
}
_ => {}
})
.await
.unwrap();
assert_eq!(report.preamble_skipped, 1);
assert_eq!(report.preamble_seeded, 0);
assert_eq!(report.published, 1);
assert_eq!(report.triggers, 1);
assert_eq!(triggers, 1);
assert_eq!(skipped.len(), 1);
assert_eq!(skipped[0].0, "v1/h-0123456789ab/state/p/config");
assert!(skipped[0].1.contains("--seed-state"), "{}", skipped[0].1);
assert!(skipped[0].1.contains("RFC 13 §4.2"), "{}", skipped[0].1);
let mut reader = source_of(&body).await;
let mut would_put = 0;
let report = replay(&mut reader, spec(true), |ev| {
if matches!(ev, ReplayEvent::WouldPut { .. }) {
would_put += 1;
}
})
.await
.unwrap();
assert_eq!(report.preamble_seeded, 1);
assert_eq!(report.preamble_skipped, 0);
assert_eq!(report.published, 1, "the observed row, not the seed");
assert_eq!(would_put, 2, "both rows are listed as puts");
}
#[test]
fn an_injected_epoch_preserves_a_window_written_after_the_fact() {
let epoch = Instant::now();
let view = |t_ms: u64| crate::bus::monitor::SampleView {
key: "v1/h-0123456789ab/state/p/a".into(),
payload: zenoh::bytes::ZBytes::from(vec![1u8]),
encoding: String::new(),
kind: SampleKind::Put,
timestamp: None,
stamped_by: None,
attachment: None,
priority: zenoh::qos::Priority::DEFAULT,
congestion_control: zenoh::qos::CongestionControl::DEFAULT,
reliability: zenoh::qos::Reliability::DEFAULT,
express: false,
source: None,
received: epoch + Duration::from_millis(t_ms),
};
let mut sink = Vec::new();
let mut w = ZrecWriter::new_at(&mut sink, &header(), epoch).unwrap();
w.write_sample(&view(0)).unwrap();
w.write_sample(&view(1500)).unwrap();
let _ = w.finish().unwrap();
let mut reader = ZrecReader::new(sink.as_slice()).unwrap();
let t_of = |item| match item {
Some(Ok(ZrecItem::Sample { t_us, .. })) => t_us,
other => panic!("expected a sample, got {other:?}"),
};
assert_eq!(t_of(reader.next()), Some(0));
assert_eq!(
t_of(reader.next()),
Some(1_500_000),
"the offset the ring preserved, not a saturated zero"
);
}
#[test]
fn drops_are_interleaved_facts() {
let body = format!(
"{}\n{}\n{}\n{}\n",
serde_json::to_string(&header()).unwrap(),
r#"{"key":"v1/h/state/p/a","t":0,"bytes":"AQ=="}"#,
r#"{"dropped":7}"#,
r#"{"key":"v1/h/state/p/a","t":1000,"bytes":"Ag=="}"#,
);
let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
assert!(matches!(reader.next(), Some(Ok(ZrecItem::Sample { .. }))));
assert!(matches!(reader.next(), Some(Ok(ZrecItem::Dropped(7)))));
assert!(matches!(
reader.next(),
Some(Ok(ZrecItem::Sample {
t_us: Some(1000),
..
}))
));
assert!(reader.next().is_none());
}
#[test]
fn malformed_lines_are_named_not_skipped() {
let body = format!(
"{}\nnot json\n{}\n",
serde_json::to_string(&header()).unwrap(),
r#"{"key":"v1/h/state/p/a","t":0,"bytes":"AQ=="}"#,
);
let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
let err = match reader.next() {
Some(Err(e)) => e,
other => panic!("expected a named error, got {other:?}"),
};
assert!(err.starts_with("line 2:"), "{err}");
assert!(matches!(reader.next(), Some(Ok(ZrecItem::Sample { .. }))));
}
#[tokio::test]
async fn a_dry_run_lists_and_publishes_nothing() {
let body = format!(
"{}\n{}\n{}\n{}\n",
serde_json::to_string(&header()).unwrap(),
r#"{"key":"v1/h-0123456789ab/state/p/health","t":0,"bytes":"eyJvayI6dHJ1ZX0=","encoding":"application/json"}"#,
r#"{"dropped":3}"#,
r#"{"key":"v1/h-0123456789ab/state/p/health","t":500000,"delete":true}"#,
);
let mut reader = source_of(&body).await;
let mut would = Vec::new();
let report = replay(
&mut reader,
ReplaySpec {
target: ReplayTarget::DryRun,
speed: 1.0,
i_know: false,
default_qos: QosProfile::Refreshed,
seed_state: false,
},
|ev| {
would.push(format!("{ev:?}"));
},
)
.await
.unwrap();
assert!(report.dry_run);
assert_eq!(report.published, 1);
assert_eq!(report.tombstones, 1); assert_eq!(report.capture_dropped, 3);
assert_eq!(report.malformed, 0);
assert_eq!(would.len(), 3, "{would:?}");
}
#[tokio::test]
async fn replayed_tombstones_pass_the_retire_gate() {
let body = format!(
"{}\n{}\n",
serde_json::to_string(&header()).unwrap(),
r#"{"key":"v1/h-0123456789ab/telemetry/p/temp","t":0,"delete":true}"#,
);
let mut reader = source_of(&body).await;
let report = replay(
&mut reader,
ReplaySpec {
target: ReplayTarget::DryRun,
speed: 1.0,
i_know: false,
default_qos: QosProfile::Refreshed,
seed_state: false,
},
|_| {},
)
.await
.unwrap();
assert_eq!(report.refused, 1);
assert_eq!(report.tombstones, 0);
assert!(
report.first_errors[0].contains("telemetry"),
"{:?}",
report.first_errors
);
}
#[tokio::test]
async fn speed_must_be_positive() {
let body = serde_json::to_string(&header()).unwrap() + "\n";
for bad in [0.0, -1.0, f64::NAN, f64::INFINITY] {
let mut reader = source_of(&body).await;
let err = replay(
&mut reader,
ReplaySpec {
target: ReplayTarget::DryRun,
speed: bad,
i_know: false,
default_qos: QosProfile::Refreshed,
seed_state: false,
},
|_| {},
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("speed"), "{err}");
}
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn a_row_the_bus_refuses_tears_down_and_still_reports_itself() {
let session = crate::bus::session::open(&[], &[], false)
.await
.expect("a standalone peer");
let good = SampleRow {
key: "v1/h-aaaaaaaaaaaa/state/demo/health".into(),
..SampleRow::default()
}
.with_payload_bytes(b"{}");
let bad = SampleRow {
key: "v1//nowhere".into(),
..SampleRow::default()
}
.with_payload_bytes(b"{}");
let body = format!(
"{}\n{}\n{}\n",
serde_json::to_string(&header()).unwrap(),
good.to_line(),
bad.to_line(),
);
let mut reader = source_of(&body).await;
let err = replay(
&mut reader,
ReplaySpec {
target: ReplayTarget::Bus {
session: &session,
slices: None,
},
speed: 1000.0,
i_know: false,
default_qos: QosProfile::Transition,
seed_state: false,
},
|_| {},
)
.await
.expect_err("the bus refused the second row")
.to_string();
assert!(err.contains("nowhere"), "{err}");
session.close().await.expect("close the session");
}
}