use std::collections::HashMap;
use std::io::{BufRead, Write};
use std::time::{Duration, Instant};
use anyhow::{Context, Result, anyhow, bail};
use base64::Engine as _;
use serde::{Deserialize, Serialize};
use zenoh::Session;
use zenoh::sample::SampleKind;
use crate::ingest::{IngestRow, parse_row};
use crate::registry::SliceSet;
use crate::sub::{EventStream, FleetEvent, SampleView, StreamItem};
pub const ZREC_VERSION: u32 = 1;
fn b64(bytes: &[u8]) -> String {
base64::engine::general_purpose::STANDARD.encode(bytes)
}
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),
)
}
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")
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct ZrecHeader {
pub zrec: u32,
pub selectors: Vec<String>,
pub base: String,
pub captured_at: String,
}
pub struct ZrecWriter<W: Write> {
out: W,
epoch: Instant,
samples: u64,
dropped: u64,
}
impl<W: Write> ZrecWriter<W> {
pub fn new(mut out: W, header: &ZrecHeader) -> Result<Self> {
serde_json::to_writer(&mut out, header).context("write .zrec header")?;
out.write_all(b"\n").context("write .zrec header")?;
Ok(ZrecWriter {
out,
epoch: Instant::now(),
samples: 0,
dropped: 0,
})
}
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 obj = serde_json::json!({
"key": view.key,
"t": t_us,
});
if view.kind == SampleKind::Delete {
obj["delete"] = true.into();
} else {
obj["bytes"] = b64(&view.payload.to_bytes()).into();
}
if !view.encoding.is_empty() {
obj["encoding"] = view.encoding.clone().into();
}
if let Some(t) = view.timestamp {
obj["timestamp"] = t.to_string().into();
}
if let Some(profile) = zenkey::qos::QosProfile::ALL
.into_iter()
.find(|p| view.qos_matches(*p))
{
obj["qos"] = profile.name().into();
}
if let Some(a) = &view.attachment {
obj["attachment_b64"] = b64(&a.to_bytes()).into();
}
serde_json::to_writer(&mut self.out, &obj).context("write .zrec row")?;
self.out.write_all(b"\n").context("write .zrec row")?;
self.samples += 1;
Ok(())
}
pub fn write_dropped(&mut self, n: u64) -> Result<()> {
serde_json::to_writer(&mut self.out, &serde_json::json!({ "dropped": n }))
.context("write .zrec drop record")?;
self.out
.write_all(b"\n")
.context("write .zrec drop record")?;
self.dropped += n;
Ok(())
}
pub fn counts(&self) -> (u64, u64) {
(self.samples, self.dropped)
}
pub fn finish(mut self) -> Result<W> {
self.out.flush().context("flush .zrec")?;
Ok(self.out)
}
}
#[derive(Debug, Clone, Copy, Default)]
pub struct RecordBounds {
pub max_samples: Option<u64>,
pub max_duration: Option<Duration>,
}
pub async fn record<W: Write>(
events: &mut EventStream,
writer: &mut ZrecWriter<W>,
bounds: RecordBounds,
mut on_progress: impl FnMut(u64, u64),
) -> Result<()> {
let deadline = bounds.max_duration.map(|d| Instant::now() + d);
loop {
let (samples, _) = writer.counts();
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))) => {
writer.write_sample(&view)?;
}
Some(StreamItem::Dropped(n)) => {
writer.write_dropped(n)?;
}
Some(_) => continue,
None => return Ok(()),
}
let (samples, dropped) = writer.counts();
on_progress(samples, dropped);
}
}
#[derive(Debug, Clone, Serialize)]
pub struct RecordReport {
pub header: ZrecHeader,
#[serde(skip_serializing_if = "Option::is_none")]
pub out: Option<String>,
pub samples: u64,
pub dropped: u64,
pub duration_ms: u64,
}
#[derive(Debug, Clone)]
pub enum ZrecItem {
Sample {
row: IngestRow,
t_us: Option<u64>,
timestamp: Option<String>,
},
Dropped(u64),
}
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(|| anyhow!("empty file — not a .zrec (no header line)"))?
.context("read .zrec header")?;
let header: ZrecHeader = serde_json::from_str(&first)
.map_err(|e| anyhow!("line 1 is not a .zrec header: {e}"))?;
if header.zrec != ZREC_VERSION {
bail!(
"unsupported .zrec version {} (this reader speaks {})",
header.zrec,
ZREC_VERSION
);
}
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()
&& let Some(n) = v.get("dropped").and_then(serde_json::Value::as_u64)
{
return Some(Ok(ZrecItem::Dropped(n)));
}
return Some(match parse_row(&line) {
Ok(row) => {
let v: serde_json::Value = serde_json::from_str(&line).unwrap_or_default();
Ok(ZrecItem::Sample {
row,
t_us: v.get("t").and_then(serde_json::Value::as_u64),
timestamp: v
.get("timestamp")
.and_then(serde_json::Value::as_str)
.map(str::to_string),
})
}
Err(e) => Err(format!("line {}: {e}", self.line)),
});
}
}
}
pub enum ReplayTarget<'a> {
DryRun,
Bus {
session: &'a Session,
slices: Option<&'a SliceSet>,
},
}
#[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),
}
#[derive(Debug, Clone, Serialize)]
pub struct ReplayReport {
pub header: ZrecHeader,
pub dry_run: bool,
pub speed: f64,
pub published: u64,
pub tombstones: u64,
pub malformed: u64,
pub refused: u64,
pub capture_dropped: u64,
#[serde(skip_serializing_if = "Vec::is_empty")]
pub first_errors: Vec<String>,
}
pub async fn replay<R: BufRead>(
reader: &mut ZrecReader<R>,
target: ReplayTarget<'_>,
speed: f64,
i_know: bool,
default_qos: &str,
mut on_event: impl FnMut(ReplayEvent<'_>),
) -> Result<ReplayReport> {
if !(speed.is_finite() && speed > 0.0) {
bail!("--speed 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(),
};
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: HashMap<String, crate::write::Publication> = HashMap::new();
let mut prev_t: Option<u64> = None;
while let Some(item) = reader.next() {
let (row, t_us) = match item {
Ok(ZrecItem::Sample { row, t_us, .. }) => (row, t_us),
Ok(ZrecItem::Dropped(n)) => {
report.capture_dropped += n;
on_event(ReplayEvent::CaptureDropped(n));
continue;
}
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::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;
}
match &target {
ReplayTarget::DryRun => {
if row.delete {
on_event(ReplayEvent::WouldRetire { key: &row.key });
report.tombstones += 1;
} else {
on_event(ReplayEvent::WouldPut {
key: &row.key,
bytes: row.payload.len(),
encoding: row.encoding.as_deref(),
});
report.published += 1;
}
}
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_name = row.qos.as_deref().unwrap_or(default_qos);
let Some(qos) = zenkey::qos::QosProfile::from_name(qos_name) else {
let reason = format!("unknown QoS profile {qos_name:?}");
on_event(ReplayEvent::Malformed {
reason: reason.clone(),
});
record_err(&mut report, reason, false);
continue;
};
let publication = crate::write::declare_publication(
session,
&row.key,
qos,
row.encoding.as_deref(),
)
.await?;
e.insert(publication)
}
};
if row.delete {
publication.retire().await?;
report.tombstones += 1;
} else {
publication.send(row.payload, row.attachment).await?;
report.published += 1;
}
}
}
}
for (_, publication) in publications.drain() {
publication.undeclare().await?;
}
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(),
}
}
#[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":99,"selectors":[],"base":"","captured_at":"x"}"#;
let err = ZrecReader::new(future.as_bytes())
.err()
.unwrap()
.to_string();
assert!(err.contains("version 99"), "{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 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 = ZrecReader::new(body.as_bytes()).unwrap();
let mut would = Vec::new();
let report = replay(
&mut reader,
ReplayTarget::DryRun,
1.0,
false,
"refreshed",
|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 = ZrecReader::new(body.as_bytes()).unwrap();
let report = replay(
&mut reader,
ReplayTarget::DryRun,
1.0,
false,
"refreshed",
|_| {},
)
.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 = ZrecReader::new(body.as_bytes()).unwrap();
let err = replay(
&mut reader,
ReplayTarget::DryRun,
bad,
false,
"refreshed",
|_| {},
)
.await
.unwrap_err()
.to_string();
assert!(err.contains("speed"), "{err}");
}
}
}