1use std::collections::HashMap;
24use std::io::{BufRead, Write};
25use std::time::{Duration, Instant};
26
27use anyhow::{Context, Result, anyhow, bail};
28use base64::Engine as _;
29use serde::{Deserialize, Serialize};
30use zenoh::Session;
31use zenoh::sample::SampleKind;
32
33use crate::ingest::{IngestRow, parse_row};
34use crate::registry::SliceSet;
35use crate::sub::{EventStream, FleetEvent, SampleView, StreamItem};
36
37pub const ZREC_VERSION: u32 = 1;
39
40fn b64(bytes: &[u8]) -> String {
41 base64::engine::general_purpose::STANDARD.encode(bytes)
42}
43
44pub fn rfc3339_now() -> String {
48 rfc3339_from_unix(
49 std::time::SystemTime::now()
50 .duration_since(std::time::UNIX_EPOCH)
51 .map(|d| d.as_secs())
52 .unwrap_or(0),
53 )
54}
55
56fn rfc3339_from_unix(secs: u64) -> String {
57 let (days, rem) = (secs / 86_400, secs % 86_400);
58 let (h, m, s) = (rem / 3600, (rem % 3600) / 60, rem % 60);
59 let z = days as i64 + 719_468;
61 let era = z.div_euclid(146_097);
62 let doe = z.rem_euclid(146_097);
63 let yoe = (doe - doe / 1460 + doe / 36_524 - doe / 146_096) / 365;
64 let y = yoe + era * 400;
65 let doy = doe - (365 * yoe + yoe / 4 - yoe / 100);
66 let mp = (5 * doy + 2) / 153;
67 let d = doy - (153 * mp + 2) / 5 + 1;
68 let mo = if mp < 10 { mp + 3 } else { mp - 9 };
69 let y = if mo <= 2 { y + 1 } else { y };
70 format!("{y:04}-{mo:02}-{d:02}T{h:02}:{m:02}:{s:02}Z")
71}
72
73#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
78pub struct ZrecHeader {
79 pub zrec: u32,
81 pub selectors: Vec<String>,
86 pub base: String,
89 pub captured_at: String,
92}
93
94pub struct ZrecWriter<W: Write> {
99 out: W,
100 epoch: Instant,
105 samples: u64,
106 dropped: u64,
107}
108
109impl<W: Write> ZrecWriter<W> {
110 pub fn new(mut out: W, header: &ZrecHeader) -> Result<Self> {
112 serde_json::to_writer(&mut out, header).context("write .zrec header")?;
113 out.write_all(b"\n").context("write .zrec header")?;
114 Ok(ZrecWriter {
115 out,
116 epoch: Instant::now(),
117 samples: 0,
118 dropped: 0,
119 })
120 }
121
122 pub fn write_sample(&mut self, view: &SampleView) -> Result<()> {
131 let t_us = u64::try_from(
132 view.received
133 .saturating_duration_since(self.epoch)
134 .as_micros(),
135 )
136 .unwrap_or(u64::MAX);
137 let mut obj = serde_json::json!({
138 "key": view.key,
139 "t": t_us,
140 });
141 if view.kind == SampleKind::Delete {
142 obj["delete"] = true.into();
143 } else {
144 obj["bytes"] = b64(&view.payload.to_bytes()).into();
145 }
146 if !view.encoding.is_empty() {
147 obj["encoding"] = view.encoding.clone().into();
148 }
149 if let Some(t) = view.timestamp {
150 obj["timestamp"] = t.to_string().into();
151 }
152 if let Some(profile) = zenkey::qos::QosProfile::ALL
153 .into_iter()
154 .find(|p| view.qos_matches(*p))
155 {
156 obj["qos"] = profile.name().into();
157 }
158 if let Some(a) = &view.attachment {
159 obj["attachment_b64"] = b64(&a.to_bytes()).into();
160 }
161 serde_json::to_writer(&mut self.out, &obj).context("write .zrec row")?;
162 self.out.write_all(b"\n").context("write .zrec row")?;
163 self.samples += 1;
164 Ok(())
165 }
166
167 pub fn write_dropped(&mut self, n: u64) -> Result<()> {
169 serde_json::to_writer(&mut self.out, &serde_json::json!({ "dropped": n }))
170 .context("write .zrec drop record")?;
171 self.out
172 .write_all(b"\n")
173 .context("write .zrec drop record")?;
174 self.dropped += n;
175 Ok(())
176 }
177
178 pub fn counts(&self) -> (u64, u64) {
180 (self.samples, self.dropped)
181 }
182
183 pub fn finish(mut self) -> Result<W> {
185 self.out.flush().context("flush .zrec")?;
186 Ok(self.out)
187 }
188}
189
190#[derive(Debug, Clone, Copy, Default)]
194pub struct RecordBounds {
195 pub max_samples: Option<u64>,
197 pub max_duration: Option<Duration>,
199}
200
201pub async fn record<W: Write>(
211 events: &mut EventStream,
212 writer: &mut ZrecWriter<W>,
213 bounds: RecordBounds,
214 mut on_progress: impl FnMut(u64, u64),
215) -> Result<()> {
216 let deadline = bounds.max_duration.map(|d| Instant::now() + d);
217 loop {
218 let (samples, _) = writer.counts();
219 if bounds.max_samples.is_some_and(|max| samples >= max) {
220 return Ok(());
221 }
222 let item = match deadline {
223 Some(d) => {
224 let left = d.saturating_duration_since(Instant::now());
225 if left.is_zero() {
226 return Ok(());
227 }
228 match tokio::time::timeout(left, events.recv()).await {
229 Ok(item) => item,
230 Err(_) => return Ok(()),
231 }
232 }
233 None => events.recv().await,
234 };
235 match item {
236 Some(StreamItem::Event(FleetEvent::Sample(view))) => {
237 writer.write_sample(&view)?;
238 }
239 Some(StreamItem::Dropped(n)) => {
240 writer.write_dropped(n)?;
241 }
242 Some(_) => continue,
243 None => return Ok(()),
244 }
245 let (samples, dropped) = writer.counts();
246 on_progress(samples, dropped);
247 }
248}
249
250#[derive(Debug, Clone, Serialize)]
252pub struct RecordReport {
253 pub header: ZrecHeader,
255 #[serde(skip_serializing_if = "Option::is_none")]
257 pub out: Option<String>,
258 pub samples: u64,
260 pub dropped: u64,
263 pub duration_ms: u64,
265}
266
267#[derive(Debug, Clone)]
269pub enum ZrecItem {
270 Sample {
274 row: IngestRow,
275 t_us: Option<u64>,
276 timestamp: Option<String>,
277 },
278 Dropped(u64),
280}
281
282pub struct ZrecReader<R: BufRead> {
285 header: ZrecHeader,
286 lines: std::io::Lines<R>,
287 line: u64,
289}
290
291impl<R: BufRead> ZrecReader<R> {
292 pub fn new(source: R) -> Result<Self> {
296 let mut lines = source.lines();
297 let first = lines
298 .next()
299 .ok_or_else(|| anyhow!("empty file — not a .zrec (no header line)"))?
300 .context("read .zrec header")?;
301 let header: ZrecHeader = serde_json::from_str(&first)
302 .map_err(|e| anyhow!("line 1 is not a .zrec header: {e}"))?;
303 if header.zrec != ZREC_VERSION {
304 bail!(
305 "unsupported .zrec version {} (this reader speaks {})",
306 header.zrec,
307 ZREC_VERSION
308 );
309 }
310 Ok(ZrecReader {
311 header,
312 lines,
313 line: 1,
314 })
315 }
316
317 pub fn header(&self) -> &ZrecHeader {
318 &self.header
319 }
320
321 #[allow(clippy::should_implement_trait)] pub fn next(&mut self) -> Option<std::result::Result<ZrecItem, String>> {
326 loop {
327 let line = match self.lines.next()? {
328 Ok(l) => l,
329 Err(e) => {
330 self.line += 1;
331 return Some(Err(format!("line {}: read: {e}", self.line)));
332 }
333 };
334 self.line += 1;
335 if line.trim().is_empty() {
336 continue;
337 }
338 if let Ok(v) = serde_json::from_str::<serde_json::Value>(&line)
340 && v.get("key").is_none()
341 && let Some(n) = v.get("dropped").and_then(serde_json::Value::as_u64)
342 {
343 return Some(Ok(ZrecItem::Dropped(n)));
344 }
345 return Some(match parse_row(&line) {
346 Ok(row) => {
347 let v: serde_json::Value = serde_json::from_str(&line).unwrap_or_default();
348 Ok(ZrecItem::Sample {
349 row,
350 t_us: v.get("t").and_then(serde_json::Value::as_u64),
351 timestamp: v
352 .get("timestamp")
353 .and_then(serde_json::Value::as_str)
354 .map(str::to_string),
355 })
356 }
357 Err(e) => Err(format!("line {}: {e}", self.line)),
358 });
359 }
360 }
361}
362
363pub enum ReplayTarget<'a> {
365 DryRun,
368 Bus {
370 session: &'a Session,
371 slices: Option<&'a SliceSet>,
374 },
375}
376
377#[derive(Debug, Clone)]
380pub enum ReplayEvent<'a> {
381 WouldPut {
383 key: &'a str,
384 bytes: usize,
385 encoding: Option<&'a str>,
386 },
387 WouldRetire { key: &'a str },
389 Malformed { reason: String },
391 Refused { key: String, reason: String },
393 CaptureDropped(u64),
396}
397
398#[derive(Debug, Clone, Serialize)]
400pub struct ReplayReport {
401 pub header: ZrecHeader,
403 pub dry_run: bool,
404 pub speed: f64,
405 pub published: u64,
407 pub tombstones: u64,
409 pub malformed: u64,
411 pub refused: u64,
413 pub capture_dropped: u64,
416 #[serde(skip_serializing_if = "Vec::is_empty")]
418 pub first_errors: Vec<String>,
419}
420
421pub async fn replay<R: BufRead>(
432 reader: &mut ZrecReader<R>,
433 target: ReplayTarget<'_>,
434 speed: f64,
435 i_know: bool,
436 default_qos: &str,
437 mut on_event: impl FnMut(ReplayEvent<'_>),
438) -> Result<ReplayReport> {
439 if !(speed.is_finite() && speed > 0.0) {
440 bail!("--speed must be a positive number (got {speed})");
441 }
442 let base = reader.header().base.clone();
443 let mut report = ReplayReport {
444 header: reader.header().clone(),
445 dry_run: matches!(target, ReplayTarget::DryRun),
446 speed,
447 published: 0,
448 tombstones: 0,
449 malformed: 0,
450 refused: 0,
451 capture_dropped: 0,
452 first_errors: Vec::new(),
453 };
454 let record_err = |report: &mut ReplayReport, reason: String, refused: bool| {
455 if refused {
456 report.refused += 1;
457 } else {
458 report.malformed += 1;
459 }
460 if report.first_errors.len() < 3 {
461 report.first_errors.push(reason);
462 }
463 };
464 let mut publications: HashMap<String, crate::write::Publication> = HashMap::new();
465 let mut prev_t: Option<u64> = None;
466 while let Some(item) = reader.next() {
467 let (row, t_us) = match item {
468 Ok(ZrecItem::Sample { row, t_us, .. }) => (row, t_us),
469 Ok(ZrecItem::Dropped(n)) => {
470 report.capture_dropped += n;
471 on_event(ReplayEvent::CaptureDropped(n));
472 continue;
473 }
474 Err(reason) => {
475 on_event(ReplayEvent::Malformed {
476 reason: reason.clone(),
477 });
478 record_err(&mut report, reason, false);
479 continue;
480 }
481 };
482 let slices = match &target {
483 ReplayTarget::Bus { slices, .. } => *slices,
484 ReplayTarget::DryRun => None,
485 };
486 if row.delete
487 && let Err(e) = crate::write::check_retire(&base, &row.key, slices, i_know)
488 {
489 let reason = e.to_string();
490 on_event(ReplayEvent::Refused {
491 key: row.key.clone(),
492 reason: reason.clone(),
493 });
494 record_err(&mut report, format!("{}: {reason}", row.key), true);
495 continue;
496 }
497 match &target {
498 ReplayTarget::DryRun => {
499 if row.delete {
500 on_event(ReplayEvent::WouldRetire { key: &row.key });
501 report.tombstones += 1;
502 } else {
503 on_event(ReplayEvent::WouldPut {
504 key: &row.key,
505 bytes: row.payload.len(),
506 encoding: row.encoding.as_deref(),
507 });
508 report.published += 1;
509 }
510 }
511 ReplayTarget::Bus { session, .. } => {
512 if let (Some(prev), Some(t)) = (prev_t, t_us)
515 && t > prev
516 {
517 let delay = Duration::from_micros(t - prev).div_f64(speed);
518 tokio::time::sleep(delay).await;
519 }
520 if t_us.is_some() {
521 prev_t = t_us;
522 }
523 let publication = match publications.entry(row.key.clone()) {
524 std::collections::hash_map::Entry::Occupied(e) => e.into_mut(),
525 std::collections::hash_map::Entry::Vacant(e) => {
526 let qos_name = row.qos.as_deref().unwrap_or(default_qos);
527 let Some(qos) = zenkey::qos::QosProfile::from_name(qos_name) else {
528 let reason = format!("unknown QoS profile {qos_name:?}");
529 on_event(ReplayEvent::Malformed {
530 reason: reason.clone(),
531 });
532 record_err(&mut report, reason, false);
533 continue;
534 };
535 let publication = crate::write::declare_publication(
536 session,
537 &row.key,
538 qos,
539 row.encoding.as_deref(),
540 )
541 .await?;
542 e.insert(publication)
543 }
544 };
545 if row.delete {
546 publication.retire().await?;
547 report.tombstones += 1;
548 } else {
549 publication.send(row.payload, row.attachment).await?;
550 report.published += 1;
551 }
552 }
553 }
554 }
555 for (_, publication) in publications.drain() {
556 publication.undeclare().await?;
557 }
558 Ok(report)
559}
560
561#[cfg(test)]
562mod tests {
563 use super::*;
564
565 #[test]
567 fn the_wall_clock_formats_correctly() {
568 assert_eq!(rfc3339_from_unix(0), "1970-01-01T00:00:00Z");
569 assert_eq!(rfc3339_from_unix(951_782_400), "2000-02-29T00:00:00Z");
570 assert_eq!(rfc3339_from_unix(1_786_492_800), "2026-08-12T00:00:00Z");
571 assert!(!rfc3339_now().is_empty());
572 }
573
574 fn header() -> ZrecHeader {
575 ZrecHeader {
576 zrec: ZREC_VERSION,
577 selectors: vec!["v1/**".into()],
578 base: String::new(),
579 captured_at: "2026-08-12T00:00:00Z".into(),
580 }
581 }
582
583 #[test]
586 fn the_header_is_a_contract() {
587 let mut sink = Vec::new();
588 let writer = ZrecWriter::new(&mut sink, &header()).unwrap();
589 let _ = writer.finish().unwrap();
590 let reader = ZrecReader::new(sink.as_slice()).unwrap();
591 assert_eq!(reader.header(), &header());
592
593 let future = r#"{"zrec":99,"selectors":[],"base":"","captured_at":"x"}"#;
594 let err = ZrecReader::new(future.as_bytes())
595 .err()
596 .unwrap()
597 .to_string();
598 assert!(err.contains("version 99"), "{err}");
599
600 let not_zrec = r#"{"key":"v1/x","value":1}"#;
601 let err = ZrecReader::new(not_zrec.as_bytes())
602 .err()
603 .unwrap()
604 .to_string();
605 assert!(err.contains("header"), "{err}");
606 }
607
608 #[test]
611 fn drops_are_interleaved_facts() {
612 let body = format!(
613 "{}\n{}\n{}\n{}\n",
614 serde_json::to_string(&header()).unwrap(),
615 r#"{"key":"v1/h/state/p/a","t":0,"bytes":"AQ=="}"#,
616 r#"{"dropped":7}"#,
617 r#"{"key":"v1/h/state/p/a","t":1000,"bytes":"Ag=="}"#,
618 );
619 let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
620 assert!(matches!(reader.next(), Some(Ok(ZrecItem::Sample { .. }))));
621 assert!(matches!(reader.next(), Some(Ok(ZrecItem::Dropped(7)))));
622 assert!(matches!(
623 reader.next(),
624 Some(Ok(ZrecItem::Sample {
625 t_us: Some(1000),
626 ..
627 }))
628 ));
629 assert!(reader.next().is_none());
630 }
631
632 #[test]
635 fn malformed_lines_are_named_not_skipped() {
636 let body = format!(
637 "{}\nnot json\n{}\n",
638 serde_json::to_string(&header()).unwrap(),
639 r#"{"key":"v1/h/state/p/a","t":0,"bytes":"AQ=="}"#,
640 );
641 let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
642 let err = match reader.next() {
643 Some(Err(e)) => e,
644 other => panic!("expected a named error, got {other:?}"),
645 };
646 assert!(err.starts_with("line 2:"), "{err}");
647 assert!(matches!(reader.next(), Some(Ok(ZrecItem::Sample { .. }))));
648 }
649
650 #[tokio::test]
653 async fn a_dry_run_lists_and_publishes_nothing() {
654 let body = format!(
655 "{}\n{}\n{}\n{}\n",
656 serde_json::to_string(&header()).unwrap(),
657 r#"{"key":"v1/h-0123456789ab/state/p/health","t":0,"bytes":"eyJvayI6dHJ1ZX0=","encoding":"application/json"}"#,
658 r#"{"dropped":3}"#,
659 r#"{"key":"v1/h-0123456789ab/state/p/health","t":500000,"delete":true}"#,
660 );
661 let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
662 let mut would = Vec::new();
663 let report = replay(
664 &mut reader,
665 ReplayTarget::DryRun,
666 1.0,
667 false,
668 "refreshed",
669 |ev| {
670 would.push(format!("{ev:?}"));
671 },
672 )
673 .await
674 .unwrap();
675 assert!(report.dry_run);
676 assert_eq!(report.published, 1);
677 assert_eq!(report.tombstones, 1); assert_eq!(report.capture_dropped, 3);
679 assert_eq!(report.malformed, 0);
680 assert_eq!(would.len(), 3, "{would:?}");
681 }
682
683 #[tokio::test]
686 async fn replayed_tombstones_pass_the_retire_gate() {
687 let body = format!(
688 "{}\n{}\n",
689 serde_json::to_string(&header()).unwrap(),
690 r#"{"key":"v1/h-0123456789ab/telemetry/p/temp","t":0,"delete":true}"#,
691 );
692 let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
693 let report = replay(
694 &mut reader,
695 ReplayTarget::DryRun,
696 1.0,
697 false,
698 "refreshed",
699 |_| {},
700 )
701 .await
702 .unwrap();
703 assert_eq!(report.refused, 1);
704 assert_eq!(report.tombstones, 0);
705 assert!(
706 report.first_errors[0].contains("telemetry"),
707 "{:?}",
708 report.first_errors
709 );
710 }
711
712 #[tokio::test]
714 async fn speed_must_be_positive() {
715 let body = serde_json::to_string(&header()).unwrap() + "\n";
716 for bad in [0.0, -1.0, f64::NAN, f64::INFINITY] {
717 let mut reader = ZrecReader::new(body.as_bytes()).unwrap();
718 let err = replay(
719 &mut reader,
720 ReplayTarget::DryRun,
721 bad,
722 false,
723 "refreshed",
724 |_| {},
725 )
726 .await
727 .unwrap_err()
728 .to_string();
729 assert!(err.contains("speed"), "{err}");
730 }
731 }
732}