1use std::io::{BufRead, Write};
14
15use crate::report::{Snapshot, SnapshotRow, ZsnapHeader};
16use crate::{Error, Result};
17
18pub 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
28pub struct ZsnapWriter<W: Write> {
31 out: W,
32 rows: u64,
33}
34
35impl<W: Write> ZsnapWriter<W> {
36 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 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 pub fn rows(&self) -> u64 {
53 self.rows
54 }
55
56 pub fn finish(mut self) -> Result<W> {
58 self.out.flush().map_err(io)?;
59 Ok(self.out)
60 }
61}
62
63pub struct ZsnapReader<R: BufRead> {
66 header: ZsnapHeader,
67 lines: std::io::Lines<R>,
68 line: u64,
70}
71
72impl<R: BufRead> ZsnapReader<R> {
73 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 #[allow(clippy::should_implement_trait)] 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 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#[cfg(feature = "decode")]
143#[derive(Debug, Clone)]
144pub struct SnapshotSpec {
145 pub selectors: Vec<String>,
147 pub timeout: std::time::Duration,
149 pub max_replies: usize,
151 pub roster: bool,
154}
155
156#[cfg(feature = "decode")]
159#[derive(Debug, Clone)]
160pub struct Taken {
161 pub snapshot: Snapshot,
162 pub incomplete: Vec<String>,
163}
164
165#[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
301pub 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 #[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 #[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 #[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}