1use std::cmp::Ordering;
19
20use ahash::AHashMap;
21use nautilus_core::UnixNanos;
22
23use crate::{
24 error::EventStoreError,
25 manifest::{RunId, RunStatus},
26 markers::{
27 DataCursorSnapshot, HiFiMarker, MarkerBackend, MarkerGap, MarkerManifest, StreamCursor,
28 StreamDictEntry, StreamSlot, compute_dict_hash, compute_gap_hash, compute_hifi_hash,
29 compute_marker_hash,
30 },
31};
32
33#[derive(Debug, Default)]
35pub struct MarkerVerifier;
36
37impl MarkerVerifier {
38 pub fn scan(
47 backend: &dyn MarkerBackend,
48 entry_high_watermark: u64,
49 ) -> Result<MarkerVerifyReport, EventStoreError> {
50 let manifest = backend.manifest()?;
51 let mut findings = Vec::new();
52 let snapshots = read_snapshots(backend, &mut findings)?;
53 let hifi = read_hifi(backend, &mut findings)?;
54 let gaps = read_gaps(backend, &mut findings)?;
55 let dict = read_dict(backend, &mut findings)?;
56
57 check_manifest_counts(&manifest, &snapshots, &hifi, &gaps, &dict, &mut findings);
58 check_marker_sequence(&snapshots, &hifi, &gaps, &mut findings);
59 check_event_seq(&snapshots, &hifi, entry_high_watermark, &mut findings);
60 check_cursor_monotonicity(&snapshots, &mut findings);
61
62 Ok(MarkerVerifyReport {
63 run_id: manifest.run_id,
64 status: manifest.status,
65 snapshots_scanned: snapshots.len() as u64,
66 hifi_scanned: hifi.len() as u64,
67 gaps_scanned: gaps.len() as u64,
68 dict_entries_scanned: dict.len() as u64,
69 findings,
70 })
71 }
72}
73
74#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct MarkerVerifyReport {
77 pub run_id: RunId,
79 pub status: RunStatus,
81 pub snapshots_scanned: u64,
83 pub hifi_scanned: u64,
85 pub gaps_scanned: u64,
87 pub dict_entries_scanned: u64,
89 pub findings: Vec<MarkerFinding>,
91}
92
93impl MarkerVerifyReport {
94 #[must_use]
96 pub fn is_clean(&self) -> bool {
97 self.findings.is_empty()
98 }
99}
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
103pub enum MarkerRecordKind {
104 Snapshot,
106 HiFi,
108 Gap,
110 Dict,
112}
113
114#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
116pub enum MarkerCountKind {
117 Snapshot,
119 HiFi,
121 Gap,
123 Dict,
125}
126
127#[derive(Debug, Clone, PartialEq, Eq)]
129pub enum MarkerFinding {
130 ManifestCountMismatch {
132 kind: MarkerCountKind,
134 manifest_count: u64,
136 scanned_count: u64,
138 },
139 MarkerSeqGap {
141 from_marker_seq: u64,
143 to_marker_seq: u64,
145 },
146 MarkerSeqOverlap {
148 from_marker_seq: u64,
150 to_marker_seq: u64,
152 },
153 InvalidMarkerGap {
155 from_marker_seq: u64,
157 to_marker_seq: u64,
159 },
160 EventSeqRegressed {
162 marker_seq: u64,
164 previous_event_seq_before: u64,
166 event_seq_before: u64,
168 },
169 EventSeqExceedsHighWatermark {
171 marker_seq: u64,
173 event_seq_before: u64,
175 high_watermark: u64,
177 },
178 CursorCountRegressed {
180 marker_seq: u64,
182 slot: StreamSlot,
184 previous_count: u64,
186 count: u64,
188 },
189 CursorTsInitRegressed {
191 marker_seq: u64,
193 slot: StreamSlot,
195 previous_ts_init_hi: UnixNanos,
197 ts_init_hi: UnixNanos,
199 },
200 HashMismatch {
202 record: MarkerRecordKind,
204 marker_seq: Option<u64>,
206 slot: Option<StreamSlot>,
208 },
209}
210
211#[derive(Debug, Clone)]
212struct ScannedRecord<T> {
213 record: T,
214}
215
216#[derive(Debug, Clone, Copy)]
217struct SequencedMarker {
218 marker_seq: u64,
219 event_seq_before: u64,
220}
221
222fn read_snapshots(
223 backend: &dyn MarkerBackend,
224 findings: &mut Vec<MarkerFinding>,
225) -> Result<Vec<ScannedRecord<DataCursorSnapshot>>, EventStoreError> {
226 if let Some(stored) = backend.scan_snapshot_records()? {
227 let out = stored
228 .into_iter()
229 .map(|stored| {
230 check_hash(
231 stored.hash,
232 compute_marker_hash(&stored.record),
233 MarkerRecordKind::Snapshot,
234 Some(stored.record.marker_seq),
235 None,
236 findings,
237 );
238 ScannedRecord {
239 record: stored.record,
240 }
241 })
242 .collect();
243 return Ok(out);
244 }
245
246 Ok(backend
247 .scan_snapshots()?
248 .into_iter()
249 .map(|record| ScannedRecord { record })
250 .collect())
251}
252
253fn read_hifi(
254 backend: &dyn MarkerBackend,
255 findings: &mut Vec<MarkerFinding>,
256) -> Result<Vec<ScannedRecord<HiFiMarker>>, EventStoreError> {
257 if let Some(stored) = backend.scan_hifi_records()? {
258 let out = stored
259 .into_iter()
260 .map(|stored| {
261 check_hash(
262 stored.hash,
263 compute_hifi_hash(&stored.record),
264 MarkerRecordKind::HiFi,
265 Some(stored.record.marker_seq),
266 Some(stored.record.slot),
267 findings,
268 );
269 ScannedRecord {
270 record: stored.record,
271 }
272 })
273 .collect();
274 return Ok(out);
275 }
276
277 Ok(backend
278 .scan_hifi()?
279 .into_iter()
280 .map(|record| ScannedRecord { record })
281 .collect())
282}
283
284fn read_gaps(
285 backend: &dyn MarkerBackend,
286 findings: &mut Vec<MarkerFinding>,
287) -> Result<Vec<ScannedRecord<MarkerGap>>, EventStoreError> {
288 if let Some(stored) = backend.scan_gap_records()? {
289 let out = stored
290 .into_iter()
291 .map(|stored| {
292 check_hash(
293 stored.hash,
294 compute_gap_hash(&stored.record),
295 MarkerRecordKind::Gap,
296 None,
297 None,
298 findings,
299 );
300 ScannedRecord {
301 record: stored.record,
302 }
303 })
304 .collect();
305 return Ok(out);
306 }
307
308 Ok(backend
309 .scan_gaps()?
310 .into_iter()
311 .map(|record| ScannedRecord { record })
312 .collect())
313}
314
315fn read_dict(
316 backend: &dyn MarkerBackend,
317 findings: &mut Vec<MarkerFinding>,
318) -> Result<Vec<ScannedRecord<StreamDictEntry>>, EventStoreError> {
319 if let Some(stored) = backend.scan_dict_records()? {
320 let out = stored
321 .into_iter()
322 .map(|stored| {
323 check_hash(
324 stored.hash,
325 compute_dict_hash(&stored.record),
326 MarkerRecordKind::Dict,
327 None,
328 Some(stored.record.slot),
329 findings,
330 );
331 ScannedRecord {
332 record: stored.record,
333 }
334 })
335 .collect();
336 return Ok(out);
337 }
338
339 Ok(backend
340 .scan_dict()?
341 .into_iter()
342 .map(|record| ScannedRecord { record })
343 .collect())
344}
345
346fn check_hash(
347 stored_hash: [u8; 32],
348 computed_hash: [u8; 32],
349 record: MarkerRecordKind,
350 marker_seq: Option<u64>,
351 slot: Option<StreamSlot>,
352 findings: &mut Vec<MarkerFinding>,
353) {
354 if stored_hash != computed_hash {
355 findings.push(MarkerFinding::HashMismatch {
356 record,
357 marker_seq,
358 slot,
359 });
360 }
361}
362
363fn check_manifest_counts(
364 manifest: &MarkerManifest,
365 snapshots: &[ScannedRecord<DataCursorSnapshot>],
366 hifi: &[ScannedRecord<HiFiMarker>],
367 gaps: &[ScannedRecord<MarkerGap>],
368 dict: &[ScannedRecord<StreamDictEntry>],
369 findings: &mut Vec<MarkerFinding>,
370) {
371 check_manifest_count(
372 MarkerCountKind::Snapshot,
373 manifest.snapshot_count,
374 snapshots.len() as u64,
375 findings,
376 );
377 check_manifest_count(
378 MarkerCountKind::HiFi,
379 manifest.hifi_count,
380 hifi.len() as u64,
381 findings,
382 );
383 check_manifest_count(
384 MarkerCountKind::Gap,
385 manifest.gap_count,
386 gaps.len() as u64,
387 findings,
388 );
389 check_manifest_count(
390 MarkerCountKind::Dict,
391 manifest.dict_count,
392 dict.len() as u64,
393 findings,
394 );
395}
396
397fn check_manifest_count(
398 kind: MarkerCountKind,
399 manifest_count: u64,
400 scanned_count: u64,
401 findings: &mut Vec<MarkerFinding>,
402) {
403 if manifest_count != scanned_count {
404 findings.push(MarkerFinding::ManifestCountMismatch {
405 kind,
406 manifest_count,
407 scanned_count,
408 });
409 }
410}
411
412fn check_marker_sequence(
413 snapshots: &[ScannedRecord<DataCursorSnapshot>],
414 hifi: &[ScannedRecord<HiFiMarker>],
415 gaps: &[ScannedRecord<MarkerGap>],
416 findings: &mut Vec<MarkerFinding>,
417) {
418 let mut coverages = Vec::new();
419
420 for snapshot in snapshots {
421 coverages.push((snapshot.record.marker_seq, snapshot.record.marker_seq));
422 }
423
424 for marker in hifi {
425 coverages.push((marker.record.marker_seq, marker.record.marker_seq));
426 }
427
428 for gap in gaps {
429 if gap.record.from_marker_seq > gap.record.to_marker_seq {
430 findings.push(MarkerFinding::InvalidMarkerGap {
431 from_marker_seq: gap.record.from_marker_seq,
432 to_marker_seq: gap.record.to_marker_seq,
433 });
434 continue;
435 }
436 coverages.push((gap.record.from_marker_seq, gap.record.to_marker_seq));
437 }
438
439 coverages.sort_unstable_by_key(|(from, to)| (*from, *to));
440 let mut expected = 1_u64;
441
442 for (from, to) in coverages {
443 match from.cmp(&expected) {
444 Ordering::Greater => {
445 findings.push(MarkerFinding::MarkerSeqGap {
446 from_marker_seq: expected,
447 to_marker_seq: from - 1,
448 });
449 expected = to.saturating_add(1);
450 }
451 Ordering::Less => {
452 findings.push(MarkerFinding::MarkerSeqOverlap {
453 from_marker_seq: from,
454 to_marker_seq: to.min(expected - 1),
455 });
456
457 if to >= expected {
458 expected = to.saturating_add(1);
459 }
460 }
461 Ordering::Equal => {
462 expected = to.saturating_add(1);
463 }
464 }
465 }
466}
467
468fn check_event_seq(
469 snapshots: &[ScannedRecord<DataCursorSnapshot>],
470 hifi: &[ScannedRecord<HiFiMarker>],
471 entry_high_watermark: u64,
472 findings: &mut Vec<MarkerFinding>,
473) {
474 let mut markers = Vec::with_capacity(snapshots.len() + hifi.len());
475
476 markers.extend(snapshots.iter().map(|snapshot| SequencedMarker {
477 marker_seq: snapshot.record.marker_seq,
478 event_seq_before: snapshot.record.event_seq_before,
479 }));
480 markers.extend(hifi.iter().map(|marker| SequencedMarker {
481 marker_seq: marker.record.marker_seq,
482 event_seq_before: marker.record.event_seq_before,
483 }));
484 markers.sort_unstable_by_key(|marker| marker.marker_seq);
485
486 let mut previous_event_seq_before = None;
487
488 for marker in markers {
489 if marker.event_seq_before > entry_high_watermark {
490 findings.push(MarkerFinding::EventSeqExceedsHighWatermark {
491 marker_seq: marker.marker_seq,
492 event_seq_before: marker.event_seq_before,
493 high_watermark: entry_high_watermark,
494 });
495 }
496
497 if let Some(previous) = previous_event_seq_before
498 && marker.event_seq_before < previous
499 {
500 findings.push(MarkerFinding::EventSeqRegressed {
501 marker_seq: marker.marker_seq,
502 previous_event_seq_before: previous,
503 event_seq_before: marker.event_seq_before,
504 });
505 }
506
507 previous_event_seq_before = Some(marker.event_seq_before);
508 }
509}
510
511fn check_cursor_monotonicity(
512 snapshots: &[ScannedRecord<DataCursorSnapshot>],
513 findings: &mut Vec<MarkerFinding>,
514) {
515 let mut ordered: Vec<&ScannedRecord<DataCursorSnapshot>> = snapshots.iter().collect();
516 ordered.sort_unstable_by_key(|snapshot| snapshot.record.marker_seq);
517 let mut cursors: AHashMap<StreamSlot, StreamCursor> = AHashMap::new();
518
519 for snapshot in ordered {
520 for cursor in &snapshot.record.advanced {
521 if let Some(previous) = cursors.get(&cursor.slot) {
522 if cursor.count < previous.count {
523 findings.push(MarkerFinding::CursorCountRegressed {
524 marker_seq: snapshot.record.marker_seq,
525 slot: cursor.slot,
526 previous_count: previous.count,
527 count: cursor.count,
528 });
529 }
530
531 if cursor.ts_init_hi < previous.ts_init_hi {
532 findings.push(MarkerFinding::CursorTsInitRegressed {
533 marker_seq: snapshot.record.marker_seq,
534 slot: cursor.slot,
535 previous_ts_init_hi: previous.ts_init_hi,
536 ts_init_hi: cursor.ts_init_hi,
537 });
538 }
539 }
540
541 cursors.insert(cursor.slot, cursor.clone());
542 }
543 }
544}
545
546#[cfg(test)]
547mod tests {
548 use rstest::rstest;
549
550 use super::*;
551 use crate::{
552 manifest::RunStatus,
553 markers::{
554 DataClass, DataCursorSnapshot, MarkerBackend, MarkerGap, MarkerGapReason,
555 MarkerManifest, MemoryMarkerBackend, StreamCursor, StreamDictEntry, compute_gap_hash,
556 compute_hifi_hash, compute_marker_hash,
557 },
558 };
559
560 fn manifest() -> MarkerManifest {
561 MarkerManifest {
562 run_id: "1700000000-marker-verifier".to_string(),
563 enabled_classes: vec![DataClass::Quote],
564 high_fidelity: false,
565 snapshot_count: 0,
566 hifi_count: 0,
567 gap_count: 0,
568 dict_count: 0,
569 status: RunStatus::Running,
570 }
571 }
572
573 fn snapshot(
574 marker_seq: u64,
575 event_seq_before: u64,
576 ts_init_hi: u64,
577 count: u64,
578 ) -> DataCursorSnapshot {
579 DataCursorSnapshot {
580 marker_seq,
581 event_seq_before,
582 ts_init: UnixNanos::from(1_700_000_000_000_000_000 + marker_seq),
583 advanced: vec![StreamCursor {
584 slot: 0,
585 ts_init_hi: UnixNanos::from(ts_init_hi),
586 count,
587 }],
588 }
589 }
590
591 fn hifi(marker_seq: u64, event_seq_before: u64, slot: StreamSlot) -> HiFiMarker {
592 HiFiMarker {
593 marker_seq,
594 event_seq_before,
595 slot,
596 ts_event: UnixNanos::from(1_700_000_000_000_000_100 + marker_seq),
597 ts_init: UnixNanos::from(1_700_000_000_000_000_200 + marker_seq),
598 same_ts_ordinal: 0,
599 record_fingerprint: [7; 32],
600 }
601 }
602
603 fn dict(slot: StreamSlot) -> StreamDictEntry {
604 StreamDictEntry {
605 slot,
606 data_cls: DataClass::Quote,
607 identifier: "ETHUSDT.BINANCE".to_string(),
608 }
609 }
610
611 fn open_backend() -> MemoryMarkerBackend {
612 let mut backend = MemoryMarkerBackend::new();
613 backend.open_run(manifest()).expect("open run");
614 backend
615 }
616
617 #[derive(Debug)]
618 struct ManifestOverrideMarkerBackend {
619 inner: MemoryMarkerBackend,
620 manifest_override: MarkerManifest,
621 }
622
623 impl ManifestOverrideMarkerBackend {
624 fn new(inner: MemoryMarkerBackend, manifest_override: MarkerManifest) -> Self {
625 Self {
626 inner,
627 manifest_override,
628 }
629 }
630 }
631
632 impl MarkerBackend for ManifestOverrideMarkerBackend {
633 fn open_run(&mut self, manifest: MarkerManifest) -> Result<(), EventStoreError> {
634 self.inner.open_run(manifest)
635 }
636
637 fn append_snapshot(
638 &mut self,
639 snapshot: &DataCursorSnapshot,
640 hash: [u8; 32],
641 ) -> Result<(), EventStoreError> {
642 self.inner.append_snapshot(snapshot, hash)
643 }
644
645 fn append_hifi(
646 &mut self,
647 marker: &HiFiMarker,
648 hash: [u8; 32],
649 ) -> Result<(), EventStoreError> {
650 self.inner.append_hifi(marker, hash)
651 }
652
653 fn append_gap(&mut self, gap: &MarkerGap, hash: [u8; 32]) -> Result<(), EventStoreError> {
654 self.inner.append_gap(gap, hash)
655 }
656
657 fn put_dict(
658 &mut self,
659 entry: &StreamDictEntry,
660 hash: [u8; 32],
661 ) -> Result<(), EventStoreError> {
662 self.inner.put_dict(entry, hash)
663 }
664
665 fn scan_snapshots(&self) -> Result<Vec<DataCursorSnapshot>, EventStoreError> {
666 self.inner.scan_snapshots()
667 }
668
669 fn scan_snapshot_records(
670 &self,
671 ) -> Result<
672 Option<Vec<crate::markers::StoredMarkerRecord<DataCursorSnapshot>>>,
673 EventStoreError,
674 > {
675 self.inner.scan_snapshot_records()
676 }
677
678 fn scan_hifi(&self) -> Result<Vec<HiFiMarker>, EventStoreError> {
679 self.inner.scan_hifi()
680 }
681
682 fn scan_hifi_records(
683 &self,
684 ) -> Result<Option<Vec<crate::markers::StoredMarkerRecord<HiFiMarker>>>, EventStoreError>
685 {
686 self.inner.scan_hifi_records()
687 }
688
689 fn scan_gaps(&self) -> Result<Vec<MarkerGap>, EventStoreError> {
690 self.inner.scan_gaps()
691 }
692
693 fn scan_gap_records(
694 &self,
695 ) -> Result<Option<Vec<crate::markers::StoredMarkerRecord<MarkerGap>>>, EventStoreError>
696 {
697 self.inner.scan_gap_records()
698 }
699
700 fn scan_dict(&self) -> Result<Vec<StreamDictEntry>, EventStoreError> {
701 self.inner.scan_dict()
702 }
703
704 fn scan_dict_records(
705 &self,
706 ) -> Result<Option<Vec<crate::markers::StoredMarkerRecord<StreamDictEntry>>>, EventStoreError>
707 {
708 self.inner.scan_dict_records()
709 }
710
711 fn seal(&mut self, status: RunStatus) -> Result<(), EventStoreError> {
712 self.inner.seal(status)
713 }
714
715 fn manifest(&self) -> Result<MarkerManifest, EventStoreError> {
716 Ok(self.manifest_override.clone())
717 }
718 }
719
720 #[derive(Debug)]
721 struct DecodedOnlyMarkerBackend {
722 inner: MemoryMarkerBackend,
723 }
724
725 impl MarkerBackend for DecodedOnlyMarkerBackend {
726 fn open_run(&mut self, manifest: MarkerManifest) -> Result<(), EventStoreError> {
727 self.inner.open_run(manifest)
728 }
729
730 fn append_snapshot(
731 &mut self,
732 snapshot: &DataCursorSnapshot,
733 hash: [u8; 32],
734 ) -> Result<(), EventStoreError> {
735 self.inner.append_snapshot(snapshot, hash)
736 }
737
738 fn append_hifi(
739 &mut self,
740 marker: &HiFiMarker,
741 hash: [u8; 32],
742 ) -> Result<(), EventStoreError> {
743 self.inner.append_hifi(marker, hash)
744 }
745
746 fn append_gap(&mut self, gap: &MarkerGap, hash: [u8; 32]) -> Result<(), EventStoreError> {
747 self.inner.append_gap(gap, hash)
748 }
749
750 fn put_dict(
751 &mut self,
752 entry: &StreamDictEntry,
753 hash: [u8; 32],
754 ) -> Result<(), EventStoreError> {
755 self.inner.put_dict(entry, hash)
756 }
757
758 fn scan_snapshots(&self) -> Result<Vec<DataCursorSnapshot>, EventStoreError> {
759 self.inner.scan_snapshots()
760 }
761
762 fn scan_hifi(&self) -> Result<Vec<HiFiMarker>, EventStoreError> {
763 self.inner.scan_hifi()
764 }
765
766 fn scan_gaps(&self) -> Result<Vec<MarkerGap>, EventStoreError> {
767 self.inner.scan_gaps()
768 }
769
770 fn scan_dict(&self) -> Result<Vec<StreamDictEntry>, EventStoreError> {
771 self.inner.scan_dict()
772 }
773
774 fn seal(&mut self, status: RunStatus) -> Result<(), EventStoreError> {
775 self.inner.seal(status)
776 }
777
778 fn manifest(&self) -> Result<MarkerManifest, EventStoreError> {
779 self.inner.manifest()
780 }
781 }
782
783 #[rstest]
784 fn decoded_only_backend_uses_verifier_fallbacks() {
785 let mut backend = DecodedOnlyMarkerBackend {
786 inner: open_backend(),
787 };
788 let snapshot = snapshot(1, 1, 100, 1);
789 let hifi = hifi(2, 2, 1);
790 let gap = MarkerGap {
791 from_marker_seq: 3,
792 to_marker_seq: 4,
793 reason: MarkerGapReason::Overflow,
794 };
795 let dict = dict(1);
796 backend
797 .append_snapshot(&snapshot, compute_marker_hash(&snapshot))
798 .expect("append snapshot");
799 backend
800 .append_hifi(&hifi, compute_hifi_hash(&hifi))
801 .expect("append hifi");
802 backend
803 .append_gap(&gap, compute_gap_hash(&gap))
804 .expect("append gap");
805 backend
806 .put_dict(&dict, compute_dict_hash(&dict))
807 .expect("put dict");
808
809 let report = MarkerVerifier::scan(&backend, 2).expect("scan");
810
811 assert_eq!(
812 backend.scan_snapshot_records().expect("snapshot records"),
813 None,
814 );
815 assert_eq!(backend.scan_hifi_records().expect("hifi records"), None);
816 assert_eq!(backend.scan_gap_records().expect("gap records"), None);
817 assert_eq!(backend.scan_dict_records().expect("dict records"), None);
818 assert_eq!(
819 report,
820 MarkerVerifyReport {
821 run_id: "1700000000-marker-verifier".to_string(),
822 status: RunStatus::Running,
823 snapshots_scanned: 1,
824 hifi_scanned: 1,
825 gaps_scanned: 1,
826 dict_entries_scanned: 1,
827 findings: Vec::new(),
828 },
829 );
830 }
831
832 #[rstest]
833 fn contiguous_marker_seq_passes() {
834 let mut backend = open_backend();
835 let s1 = snapshot(1, 1, 100, 1);
836 let s2 = snapshot(2, 2, 200, 2);
837 backend
838 .append_snapshot(&s1, compute_marker_hash(&s1))
839 .expect("append s1");
840 backend
841 .append_snapshot(&s2, compute_marker_hash(&s2))
842 .expect("append s2");
843
844 let report = MarkerVerifier::scan(&backend, 2).expect("scan");
845
846 assert!(report.is_clean(), "findings was: {:?}", report.findings);
847 }
848
849 #[rstest]
850 fn hifi_markers_participate_in_marker_seq_coverage() {
851 let mut backend = open_backend();
852 let s1 = snapshot(1, 1, 100, 1);
853 let m2 = hifi(2, 1, 0);
854 let s3 = snapshot(3, 2, 200, 2);
855 backend
856 .append_snapshot(&s1, compute_marker_hash(&s1))
857 .expect("append s1");
858 backend
859 .append_hifi(&m2, compute_hifi_hash(&m2))
860 .expect("append hifi");
861 backend
862 .append_snapshot(&s3, compute_marker_hash(&s3))
863 .expect("append s3");
864
865 let report = MarkerVerifier::scan(&backend, 2).expect("scan");
866
867 assert!(report.is_clean(), "findings was: {:?}", report.findings);
868 }
869
870 #[rstest]
871 fn manifest_count_mismatch_is_corrupt() {
872 let mut backend = open_backend();
873 let s1 = snapshot(1, 1, 100, 1);
874 backend
875 .append_snapshot(&s1, compute_marker_hash(&s1))
876 .expect("append s1");
877 let mut manifest_override = backend.manifest().expect("manifest");
878 manifest_override.snapshot_count = 2;
879 let backend = ManifestOverrideMarkerBackend::new(backend, manifest_override);
880
881 let report = MarkerVerifier::scan(&backend, 1).expect("scan");
882
883 assert!(
884 report.findings.iter().any(|finding| matches!(
885 finding,
886 MarkerFinding::ManifestCountMismatch {
887 kind: MarkerCountKind::Snapshot,
888 manifest_count: 2,
889 scanned_count: 1,
890 }
891 )),
892 "findings was: {:?}",
893 report.findings,
894 );
895 }
896
897 #[rstest]
898 fn missing_marker_seq_without_gap_is_corrupt() {
899 let mut backend = open_backend();
900 let s1 = snapshot(1, 1, 100, 1);
901 let s3 = snapshot(3, 2, 200, 2);
902 backend
903 .append_snapshot(&s1, compute_marker_hash(&s1))
904 .expect("append s1");
905 backend
906 .append_snapshot(&s3, compute_marker_hash(&s3))
907 .expect("append s3");
908
909 let report = MarkerVerifier::scan(&backend, 3).expect("scan");
910
911 assert!(
912 report.findings.iter().any(|finding| matches!(
913 finding,
914 MarkerFinding::MarkerSeqGap {
915 from_marker_seq: 2,
916 to_marker_seq: 2,
917 }
918 )),
919 "findings was: {:?}",
920 report.findings,
921 );
922 }
923
924 #[rstest]
925 fn invalid_marker_gap_is_corrupt() {
926 let mut backend = open_backend();
927 let gap = MarkerGap {
928 from_marker_seq: 4,
929 to_marker_seq: 2,
930 reason: MarkerGapReason::Overflow,
931 };
932 backend
933 .append_gap(&gap, compute_gap_hash(&gap))
934 .expect("append invalid gap");
935
936 let report = MarkerVerifier::scan(&backend, 0).expect("scan");
937
938 assert!(
939 report.findings.iter().any(|finding| matches!(
940 finding,
941 MarkerFinding::InvalidMarkerGap {
942 from_marker_seq: 4,
943 to_marker_seq: 2,
944 }
945 )),
946 "findings was: {:?}",
947 report.findings,
948 );
949 }
950
951 #[rstest]
952 fn overlapping_marker_coverage_is_corrupt() {
953 let mut backend = open_backend();
954 let s1 = snapshot(1, 1, 100, 1);
955 let gap = MarkerGap {
956 from_marker_seq: 1,
957 to_marker_seq: 1,
958 reason: MarkerGapReason::Overflow,
959 };
960 backend
961 .append_snapshot(&s1, compute_marker_hash(&s1))
962 .expect("append s1");
963 backend
964 .append_gap(&gap, compute_gap_hash(&gap))
965 .expect("append overlapping gap");
966
967 let report = MarkerVerifier::scan(&backend, 1).expect("scan");
968
969 assert!(
970 report.findings.iter().any(|finding| matches!(
971 finding,
972 MarkerFinding::MarkerSeqOverlap {
973 from_marker_seq: 1,
974 to_marker_seq: 1,
975 }
976 )),
977 "findings was: {:?}",
978 report.findings,
979 );
980 }
981
982 #[rstest]
983 fn non_monotonic_count_is_corrupt() {
984 let mut backend = open_backend();
985 let s1 = snapshot(1, 1, 100, 10);
986 let s2 = snapshot(2, 2, 200, 9);
987 backend
988 .append_snapshot(&s1, compute_marker_hash(&s1))
989 .expect("append s1");
990 backend
991 .append_snapshot(&s2, compute_marker_hash(&s2))
992 .expect("append s2");
993
994 let report = MarkerVerifier::scan(&backend, 2).expect("scan");
995
996 assert!(
997 report.findings.iter().any(|finding| matches!(
998 finding,
999 MarkerFinding::CursorCountRegressed {
1000 marker_seq: 2,
1001 slot: 0,
1002 previous_count: 10,
1003 count: 9,
1004 }
1005 )),
1006 "findings was: {:?}",
1007 report.findings,
1008 );
1009 }
1010
1011 #[rstest]
1012 fn non_monotonic_ts_init_hi_is_corrupt() {
1013 let mut backend = open_backend();
1014 let s1 = snapshot(1, 1, 200, 1);
1015 let s2 = snapshot(2, 2, 199, 2);
1016 backend
1017 .append_snapshot(&s1, compute_marker_hash(&s1))
1018 .expect("append s1");
1019 backend
1020 .append_snapshot(&s2, compute_marker_hash(&s2))
1021 .expect("append s2");
1022
1023 let report = MarkerVerifier::scan(&backend, 2).expect("scan");
1024
1025 assert!(
1026 report.findings.iter().any(|finding| matches!(
1027 finding,
1028 MarkerFinding::CursorTsInitRegressed {
1029 marker_seq: 2,
1030 slot: 0,
1031 previous_ts_init_hi,
1032 ts_init_hi,
1033 } if *previous_ts_init_hi == UnixNanos::from(200)
1034 && *ts_init_hi == UnixNanos::from(199)
1035 )),
1036 "findings was: {:?}",
1037 report.findings,
1038 );
1039 }
1040
1041 #[rstest]
1042 fn event_seq_before_regression_is_corrupt() {
1043 let mut backend = open_backend();
1044 let s1 = snapshot(1, 10, 100, 1);
1045 let s2 = snapshot(2, 9, 200, 2);
1046 backend
1047 .append_snapshot(&s1, compute_marker_hash(&s1))
1048 .expect("append s1");
1049 backend
1050 .append_snapshot(&s2, compute_marker_hash(&s2))
1051 .expect("append s2");
1052
1053 let report = MarkerVerifier::scan(&backend, 10).expect("scan");
1054
1055 assert!(
1056 report.findings.iter().any(|finding| matches!(
1057 finding,
1058 MarkerFinding::EventSeqRegressed {
1059 marker_seq: 2,
1060 previous_event_seq_before: 10,
1061 event_seq_before: 9,
1062 }
1063 )),
1064 "findings was: {:?}",
1065 report.findings,
1066 );
1067 }
1068
1069 #[rstest]
1070 fn event_seq_before_exceeding_high_watermark_is_corrupt() {
1071 let mut backend = open_backend();
1072 let s1 = snapshot(1, 11, 100, 1);
1073 backend
1074 .append_snapshot(&s1, compute_marker_hash(&s1))
1075 .expect("append s1");
1076
1077 let report = MarkerVerifier::scan(&backend, 10).expect("scan");
1078
1079 assert!(
1080 report.findings.iter().any(|finding| matches!(
1081 finding,
1082 MarkerFinding::EventSeqExceedsHighWatermark {
1083 marker_seq: 1,
1084 event_seq_before: 11,
1085 high_watermark: 10,
1086 }
1087 )),
1088 "findings was: {:?}",
1089 report.findings,
1090 );
1091 }
1092
1093 #[rstest]
1094 fn non_snapshot_record_hash_mismatches_are_corrupt() {
1095 let mut backend = open_backend();
1096 let s1 = snapshot(1, 1, 100, 1);
1097 let m2 = hifi(2, 1, 0);
1098 let g3 = MarkerGap {
1099 from_marker_seq: 3,
1100 to_marker_seq: 3,
1101 reason: MarkerGapReason::Overflow,
1102 };
1103 let d0 = dict(0);
1104 backend
1105 .append_snapshot(&s1, compute_marker_hash(&s1))
1106 .expect("append s1");
1107 backend
1108 .append_hifi(&m2, [0xBB; 32])
1109 .expect("append bad hifi hash");
1110 backend
1111 .append_gap(&g3, [0xCC; 32])
1112 .expect("append bad gap hash");
1113 backend
1114 .put_dict(&d0, [0xDD; 32])
1115 .expect("put bad dict hash");
1116
1117 let report = MarkerVerifier::scan(&backend, 1).expect("scan");
1118
1119 assert!(
1120 report.findings.iter().any(|finding| matches!(
1121 finding,
1122 MarkerFinding::HashMismatch {
1123 record: MarkerRecordKind::HiFi,
1124 marker_seq: Some(2),
1125 slot: Some(0),
1126 }
1127 )),
1128 "findings was: {:?}",
1129 report.findings,
1130 );
1131 assert!(
1132 report.findings.iter().any(|finding| matches!(
1133 finding,
1134 MarkerFinding::HashMismatch {
1135 record: MarkerRecordKind::Gap,
1136 marker_seq: None,
1137 slot: None,
1138 }
1139 )),
1140 "findings was: {:?}",
1141 report.findings,
1142 );
1143 assert!(
1144 report.findings.iter().any(|finding| matches!(
1145 finding,
1146 MarkerFinding::HashMismatch {
1147 record: MarkerRecordKind::Dict,
1148 marker_seq: None,
1149 slot: Some(0),
1150 }
1151 )),
1152 "findings was: {:?}",
1153 report.findings,
1154 );
1155 }
1156
1157 #[rstest]
1158 fn bad_record_hash_is_corrupt() {
1159 let mut backend = open_backend();
1160 let s1 = snapshot(1, 1, 100, 1);
1161 backend
1162 .append_snapshot(&s1, [0xAA; 32])
1163 .expect("append bad hash");
1164
1165 let report = MarkerVerifier::scan(&backend, 1).expect("scan");
1166
1167 assert!(
1168 report.findings.iter().any(|finding| matches!(
1169 finding,
1170 MarkerFinding::HashMismatch {
1171 record: MarkerRecordKind::Snapshot,
1172 marker_seq: Some(1),
1173 slot: None,
1174 }
1175 )),
1176 "findings was: {:?}",
1177 report.findings,
1178 );
1179 }
1180
1181 #[rstest]
1182 fn marker_gap_covers_missing_marker_seq() {
1183 let mut backend = open_backend();
1184 let s1 = snapshot(1, 1, 100, 1);
1185 let s3 = snapshot(3, 2, 200, 2);
1186 let gap = MarkerGap {
1187 from_marker_seq: 2,
1188 to_marker_seq: 2,
1189 reason: MarkerGapReason::Overflow,
1190 };
1191 backend
1192 .append_snapshot(&s1, compute_marker_hash(&s1))
1193 .expect("append s1");
1194 backend
1195 .append_gap(&gap, compute_gap_hash(&gap))
1196 .expect("append gap");
1197 backend
1198 .append_snapshot(&s3, compute_marker_hash(&s3))
1199 .expect("append s3");
1200
1201 let report = MarkerVerifier::scan(&backend, 3).expect("scan");
1202
1203 assert!(report.is_clean(), "findings was: {:?}", report.findings);
1204 }
1205}