Skip to main content

zenkey_fleet/tape/
snapshot.rs

1//! `.zsnap` — a snapshot kept on disk (RFC 13 §4.4, v1.34; #219): the
2//! writer, the reader, and — behind `decode` — the collection itself.
3//!
4//! The file is the `.zrec` dialect's sibling: newline-delimited JSON, line 1
5//! the header, one row per key after it, payloads lossless as base64
6//! `bytes`. A reader that does not know the version refuses the file rather
7//! than guessing, with the same wording [`crate::ZrecReader`] uses.
8//!
9//! What deliberately does not exist here: a *replay* of a snapshot. A
10//! snapshot is read, never published (RFC 13 §4.4); seeding a fleet from a
11//! file is `--seed-state` over a version-2 `.zrec` preamble.
12
13use std::io::{BufRead, Write};
14
15use crate::report::{Snapshot, SnapshotRow, ZsnapHeader};
16use crate::{Error, Result};
17
18/// The current `.zsnap` format version, written into every header.
19pub const ZSNAP_VERSION: u32 = 1;
20
21fn io(e: std::io::Error) -> Error {
22    Error::Io {
23        path: std::path::PathBuf::new(),
24        source: e,
25    }
26}
27
28/// A `.zsnap` writer over any byte sink: header first, then rows. Wrap the
29/// sink in a `BufWriter` — the writer emits line-at-a-time.
30pub struct ZsnapWriter<W: Write> {
31    out: W,
32    rows: u64,
33}
34
35impl<W: Write> ZsnapWriter<W> {
36    /// Write the header line and hand back a row writer.
37    pub fn new(mut out: W, header: &ZsnapHeader) -> Result<Self> {
38        serde_json::to_writer(&mut out, header).map_err(|e| io(e.into()))?;
39        out.write_all(b"\n").map_err(io)?;
40        Ok(ZsnapWriter { out, rows: 0 })
41    }
42
43    /// Write one row.
44    pub fn write_row(&mut self, row: &SnapshotRow) -> Result<()> {
45        serde_json::to_writer(&mut self.out, row).map_err(|e| io(e.into()))?;
46        self.out.write_all(b"\n").map_err(io)?;
47        self.rows += 1;
48        Ok(())
49    }
50
51    /// Rows written so far.
52    pub fn rows(&self) -> u64 {
53        self.rows
54    }
55
56    /// Flush and hand the sink back.
57    pub fn finish(mut self) -> Result<W> {
58        self.out.flush().map_err(io)?;
59        Ok(self.out)
60    }
61}
62
63/// A `.zsnap` reader over any buffered byte source: header up front, then
64/// one row per line — bounded memory, like the writer.
65pub struct ZsnapReader<R: BufRead> {
66    header: ZsnapHeader,
67    lines: std::io::Lines<R>,
68    /// 1-based number of the last line handed out (the header is line 1).
69    line: u64,
70}
71
72impl<R: BufRead> ZsnapReader<R> {
73    /// Parse the header line. A file without one is not a `.zsnap`; a
74    /// version this reader does not speak is refused, never guessed at
75    /// (RFC 13 §4.1's rule, which §4.4 inherits).
76    pub fn new(source: R) -> Result<Self> {
77        let mut lines = source.lines();
78        let first = lines
79            .next()
80            .ok_or_else(|| Error::malformed(".zsnap", "empty file — no header line"))?
81            .map_err(io)?;
82        let header: ZsnapHeader = serde_json::from_str(&first)
83            .map_err(|e| Error::malformed_with(".zsnap line 1", "is not a header", e))?;
84        if header.zsnap != ZSNAP_VERSION {
85            return Err(Error::malformed(
86                ".zsnap",
87                format!(
88                    "unsupported version {} (this reader speaks {ZSNAP_VERSION})",
89                    header.zsnap
90                ),
91            ));
92        }
93        Ok(ZsnapReader {
94            header,
95            lines,
96            line: 1,
97        })
98    }
99
100    pub fn header(&self) -> &ZsnapHeader {
101        &self.header
102    }
103
104    /// The next row, or `Err` naming the line and the reason — a malformed
105    /// row is reported, never silently skipped. `None` ends the file.
106    #[allow(clippy::should_implement_trait)] // fallible, line-numbered next
107    pub fn next(&mut self) -> Option<std::result::Result<SnapshotRow, String>> {
108        loop {
109            let line = match self.lines.next()? {
110                Ok(l) => l,
111                Err(e) => {
112                    self.line += 1;
113                    return Some(Err(format!("line {}: read: {e}", self.line)));
114                }
115            };
116            self.line += 1;
117            if line.trim().is_empty() {
118                continue;
119            }
120            return Some(
121                serde_json::from_str::<SnapshotRow>(&line)
122                    .map_err(|e| format!("line {}: {e}", self.line)),
123            );
124        }
125    }
126
127    /// The whole file, or the first malformed line as an error — a diff over
128    /// a file with a hole in it would be a diff of a document nobody wrote.
129    pub fn read_all(mut self) -> Result<Snapshot> {
130        let mut rows = Vec::new();
131        while let Some(row) = self.next() {
132            rows.push(row.map_err(|e| Error::malformed(".zsnap", e))?);
133        }
134        Ok(Snapshot {
135            header: self.header,
136            rows,
137        })
138    }
139}
140
141/// What to collect (`decode`).
142#[cfg(feature = "decode")]
143#[derive(Debug, Clone)]
144pub struct SnapshotSpec {
145    /// Full wire selectors, one fan-in GET each.
146    pub selectors: Vec<String>,
147    /// Reply-wait per GET, and the roster ask's bound.
148    pub timeout: std::time::Duration,
149    /// Replies kept per GET (#339); what the bound cost rides the header.
150    pub max_replies: usize,
151    /// Whether to ask the liveliness roster. Off, every holder is
152    /// `unattributed { reason: "roster not asked" }` — cheaper, and honest.
153    pub roster: bool,
154}
155
156/// What [`take_snapshot`] hands back: the snapshot, and the selectors whose
157/// GET could not be issued at all (asked, never put — RFC 09 §5.1 O5).
158#[cfg(feature = "decode")]
159#[derive(Debug, Clone)]
160pub struct Taken {
161    pub snapshot: Snapshot,
162    pub incomplete: Vec<String>,
163}
164
165/// Collect a snapshot (RFC 13 §4.4).
166///
167/// The watchdog's discipline first: the schema store is pre-warmed for
168/// every producer the registry names and then sealed, so no decode inside
169/// the collection reaches for the bus ([`crate::prewarm`],
170/// [`crate::SchemaStore::seal`]). Then the roster ask and one
171/// [`snapshot_get`](crate::bus::query::snapshot_get) per selector run
172/// **together** — one [`crate::GetOpts`] across the GETs, so `elided`
173/// sums. The replies fold per key last-writer-wins
174/// ([`crate::model::snapshot::fold_latest`]), and each kept value becomes a
175/// row: registration from [`crate::KeyFacts`], verdict from
176/// [`crate::decode_sample`], holder from the roster.
177///
178/// `collection_span_s` is measured from before the first ask to after the
179/// last row — the roster ask included. That is the honest number: a fan-in
180/// GET is collected *over* a span, and every rendering states it.
181#[cfg(feature = "decode")]
182pub async fn take_snapshot(
183    fleet: &crate::Fleet<'_>,
184    slices: Option<&crate::SliceSet>,
185    store: &crate::SchemaStore,
186    spec: &SnapshotSpec,
187) -> Result<Taken> {
188    use crate::model::snapshot::{fold_latest, holder_of, registration_of, stamper_of, verdict_of};
189    use crate::report::{Asked, VerdictWire};
190    use zenoh::sample::SampleKind;
191
192    let started = std::time::Instant::now();
193    let collected_at = crate::tape::record::rfc3339_now();
194
195    crate::model::decode::prewarm(fleet, store, slices).await;
196    let _sealed = store.seal();
197
198    let opts = crate::GetOpts::new(spec.timeout).max_replies(spec.max_replies);
199    let gets = futures_util::future::join_all(spec.selectors.iter().map(|selector| {
200        let opts = &opts;
201        async move {
202            (
203                selector.clone(),
204                crate::bus::query::snapshot_get(fleet.session(), selector, opts).await,
205            )
206        }
207    }));
208    let roster = async {
209        if spec.roster {
210            crate::bus::roster::roster(fleet, spec.timeout)
211                .await
212                .map(Some)
213        } else {
214            Ok(None)
215        }
216    };
217    let (replies, roster) = tokio::join!(gets, roster);
218    let roster = roster?;
219
220    let mut values = Vec::new();
221    let mut errors = 0u64;
222    let mut incomplete = Vec::new();
223    for (selector, replies) in replies {
224        match replies {
225            Ok(r) => {
226                errors += r.errors;
227                values.extend(r.values);
228            }
229            Err(e) => {
230                tracing::warn!(selector, error = %e, "snapshot GET could not be issued");
231                incomplete.push(selector);
232            }
233        }
234    }
235    let answered = values.len() as u64;
236    let (kept, superseded) = fold_latest(values);
237
238    let base = fleet.base();
239    let mut rows = Vec::with_capacity(kept.len());
240    for (key, (view, replier)) in kept {
241        let mut facts = crate::KeyFacts::project(base, &key);
242        if let Some(slices) = slices {
243            facts.resolve(slices);
244        }
245        let delete = view.kind == SampleKind::Delete;
246        let (bytes, verdict) = if delete {
247            (
248                None,
249                VerdictWire::NotValidated {
250                    reason: "tombstone".into(),
251                },
252            )
253        } else {
254            let payload = view.payload.to_bytes();
255            let encoding = (!view.encoding.is_empty()).then_some(view.encoding.as_str());
256            let decoded =
257                crate::decode_sample(fleet, store, slices, &key, encoding, &payload).await;
258            (
259                Some(crate::tape::ingest::b64(&payload)),
260                verdict_of(&decoded.verdict),
261            )
262        };
263        let holder = holder_of(base, &key, &view, replier, roster.as_ref());
264        rows.push(SnapshotRow {
265            key,
266            delete,
267            bytes,
268            encoding: (!view.encoding.is_empty()).then(|| view.encoding.clone()),
269            timestamp: view.timestamp.map(|t| t.to_string()),
270            stamper: view.stamped_by.as_ref().map(stamper_of),
271            source: view.source.map(|s| format!("{}:{}#{}", s.zid, s.eid, s.sn)),
272            source_zid: replier.map(|z| z.to_string()),
273            registration: registration_of(&facts),
274            verdict,
275            holder,
276        });
277    }
278
279    let header = ZsnapHeader {
280        zsnap: ZSNAP_VERSION,
281        selectors: spec.selectors.clone(),
282        base: base.to_string(),
283        collected_at,
284        collection_span_s: started.elapsed().as_secs_f64(),
285        asked: spec.selectors.len() as u64,
286        answered,
287        elided: opts.elided(),
288        errors,
289        superseded,
290        roster: match &roster {
291            Some(r) => Asked::Asked(r.len()),
292            None => Asked::NotAsked,
293        },
294    };
295    Ok(Taken {
296        snapshot: Snapshot { header, rows },
297        incomplete,
298    })
299}
300
301/// The report `zenctl snapshot` renders, counted off the rows.
302pub fn report_of(
303    snapshot: &Snapshot,
304    out: Option<String>,
305    incomplete: Vec<String>,
306) -> crate::report::SnapshotReport {
307    use crate::report::Holder;
308    let mut report = crate::report::SnapshotReport {
309        header: snapshot.header.clone(),
310        out,
311        live: 0,
312        storage_only: 0,
313        unattributed: 0,
314        incomplete,
315    };
316    for row in &snapshot.rows {
317        match row.holder {
318            Holder::Live { .. } => report.live += 1,
319            Holder::StorageOnly { .. } => report.storage_only += 1,
320            Holder::Unattributed { .. } => report.unattributed += 1,
321        }
322    }
323    report
324}
325
326#[cfg(test)]
327mod tests {
328    use super::*;
329    use crate::report::{Asked, Holder, RegistrationWire, VerdictWire};
330
331    fn header() -> ZsnapHeader {
332        ZsnapHeader {
333            zsnap: ZSNAP_VERSION,
334            selectors: vec!["v1/**".into()],
335            base: String::new(),
336            collected_at: "2026-09-06T00:00:00Z".into(),
337            collection_span_s: 0.75,
338            asked: 1,
339            answered: 3,
340            elided: 0,
341            errors: 1,
342            superseded: 1,
343            roster: Asked::Asked(1),
344        }
345    }
346
347    fn row(key: &str) -> SnapshotRow {
348        SnapshotRow {
349            key: key.into(),
350            delete: false,
351            bytes: Some("e30=".into()),
352            encoding: Some("application/json".into()),
353            timestamp: Some("7f00...".into()),
354            stamper: None,
355            source: None,
356            source_zid: Some("ab12".into()),
357            registration: RegistrationWire::RegistryNotLoaded,
358            verdict: VerdictWire::NotValidated {
359                reason: "no_registry".into(),
360            },
361            holder: Holder::StorageOnly {
362                origin: "h-aaaaaaaaaaaa".into(),
363            },
364        }
365    }
366
367    /// Writer → reader is the identity on the whole document.
368    #[test]
369    fn a_snapshot_round_trips_through_the_file() {
370        let rows = vec![
371            row("v1/h-aaaaaaaaaaaa/state/p/a"),
372            row("v1/h-aaaaaaaaaaaa/state/p/b"),
373        ];
374        let mut sink = Vec::new();
375        let mut w = ZsnapWriter::new(&mut sink, &header()).unwrap();
376        for r in &rows {
377            w.write_row(r).unwrap();
378        }
379        assert_eq!(w.rows(), 2);
380        w.finish().unwrap();
381
382        let snapshot = ZsnapReader::new(sink.as_slice())
383            .unwrap()
384            .read_all()
385            .unwrap();
386        assert_eq!(
387            snapshot,
388            Snapshot {
389                header: header(),
390                rows
391            }
392        );
393    }
394
395    /// A versioned reader refuses what it cannot speak rather than guessing,
396    /// in the `.zrec` reader's words; and a row is not a header.
397    #[test]
398    fn the_header_is_a_contract() {
399        let future = r#"{"zsnap":99,"selectors":[],"base":"","collected_at":"x","collection_span_s":0,"asked":0,"answered":0}"#;
400        let err = ZsnapReader::new(future.as_bytes())
401            .err()
402            .unwrap()
403            .to_string();
404        assert!(err.contains("unsupported version 99"), "{err}");
405
406        let not_a_header = r#"{"key":"v1/x","delete":false}"#;
407        let err = ZsnapReader::new(not_a_header.as_bytes())
408            .err()
409            .unwrap()
410            .to_string();
411        assert!(err.contains("is not a header"), "{err}");
412
413        let err = ZsnapReader::new("".as_bytes()).err().unwrap().to_string();
414        assert!(err.contains("no header"), "{err}");
415    }
416
417    /// A malformed row names its line and stops the read — a diff over a
418    /// document with a hole is not a diff of anything.
419    #[test]
420    fn a_malformed_row_is_reported_by_line() {
421        let body = format!(
422            "{}\n{}\nnot json\n",
423            serde_json::to_string(&header()).unwrap(),
424            serde_json::to_string(&row("v1/x")).unwrap()
425        );
426        let err = ZsnapReader::new(body.as_bytes())
427            .unwrap()
428            .read_all()
429            .err()
430            .unwrap()
431            .to_string();
432        assert!(err.contains("line 3"), "{err}");
433    }
434
435    #[test]
436    fn the_report_counts_holders() {
437        let mut live = row("v1/h-aaaaaaaaaaaa/state/p/a");
438        live.holder = Holder::Live {
439            origin: "h-aaaaaaaaaaaa".into(),
440            answered_by: crate::report::AnsweredBy::Stamper,
441        };
442        let snapshot = Snapshot {
443            header: header(),
444            rows: vec![live, row("v1/h-aaaaaaaaaaaa/state/p/b")],
445        };
446        let r = report_of(&snapshot, Some("a.zsnap".into()), vec![]);
447        assert_eq!((r.live, r.storage_only, r.unattributed), (1, 1, 0));
448    }
449}