Skip to main content

nautilus_event_store/markers/
verifier.rs

1// -------------------------------------------------------------------------------------------------
2//  Copyright (C) 2015-2026 Nautech Systems Pty Ltd. All rights reserved.
3//  https://nautechsystems.io
4//
5//  Licensed under the GNU Lesser General Public License Version 3.0 (the "License");
6//  You may not use this file except in compliance with the License.
7//  You may obtain a copy of the License at https://www.gnu.org/licenses/lgpl-3.0.en.html
8//
9//  Unless required by applicable law or agreed to in writing, software
10//  distributed under the License is distributed on an "AS IS" BASIS,
11//  WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12//  See the License for the specific language governing permissions and
13//  limitations under the License.
14// -------------------------------------------------------------------------------------------------
15
16//! Non-fatal integrity verifier for data marker sidecars.
17
18use 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/// Verifier for a single marker sidecar backend.
34#[derive(Debug, Default)]
35pub struct MarkerVerifier;
36
37impl MarkerVerifier {
38    /// Scans the marker sidecar and returns all non-fatal findings.
39    ///
40    /// `entry_high_watermark` comes from the sibling entry run and bounds marker
41    /// `event_seq_before` values.
42    ///
43    /// # Errors
44    ///
45    /// Returns [`EventStoreError`] when the backend cannot scan a marker table or manifest.
46    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/// Structured report produced by [`MarkerVerifier`].
75#[derive(Debug, Clone, PartialEq, Eq)]
76pub struct MarkerVerifyReport {
77    /// The id of the run this marker sidecar belongs to.
78    pub run_id: RunId,
79    /// The marker sidecar lifecycle status at verification time.
80    pub status: RunStatus,
81    /// Number of cursor snapshots scanned.
82    pub snapshots_scanned: u64,
83    /// Number of high-fidelity markers scanned.
84    pub hifi_scanned: u64,
85    /// Number of marker gaps scanned.
86    pub gaps_scanned: u64,
87    /// Number of stream dictionary entries scanned.
88    pub dict_entries_scanned: u64,
89    /// Every marker-sidecar integrity finding.
90    pub findings: Vec<MarkerFinding>,
91}
92
93impl MarkerVerifyReport {
94    /// Returns `true` when no marker findings were accumulated.
95    #[must_use]
96    pub fn is_clean(&self) -> bool {
97        self.findings.is_empty()
98    }
99}
100
101/// Durable marker table kind used in hash findings.
102#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
103pub enum MarkerRecordKind {
104    /// Cursor snapshot table.
105    Snapshot,
106    /// High-fidelity marker table.
107    HiFi,
108    /// Marker gap table.
109    Gap,
110    /// Stream dictionary table.
111    Dict,
112}
113
114/// Marker manifest count field used in count-mismatch findings.
115#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
116pub enum MarkerCountKind {
117    /// `snapshot_count`.
118    Snapshot,
119    /// `hifi_count`.
120    HiFi,
121    /// `gap_count`.
122    Gap,
123    /// `dict_count`.
124    Dict,
125}
126
127/// One non-fatal integrity finding from a marker verification scan.
128#[derive(Debug, Clone, PartialEq, Eq)]
129pub enum MarkerFinding {
130    /// A marker manifest count disagrees with the scanned table row count.
131    ManifestCountMismatch {
132        /// Which manifest count diverged.
133        kind: MarkerCountKind,
134        /// The count recorded in the marker manifest.
135        manifest_count: u64,
136        /// The count observed by scanning the marker table.
137        scanned_count: u64,
138    },
139    /// One or more marker sequence values were neither present nor covered by a gap.
140    MarkerSeqGap {
141        /// First missing marker sequence.
142        from_marker_seq: u64,
143        /// Last missing marker sequence.
144        to_marker_seq: u64,
145    },
146    /// A marker sequence value was covered more than once.
147    MarkerSeqOverlap {
148        /// First overlapped marker sequence.
149        from_marker_seq: u64,
150        /// Last overlapped marker sequence.
151        to_marker_seq: u64,
152    },
153    /// A stored gap has an invalid inclusive range.
154    InvalidMarkerGap {
155        /// The gap's first marker sequence.
156        from_marker_seq: u64,
157        /// The gap's last marker sequence.
158        to_marker_seq: u64,
159    },
160    /// `event_seq_before` decreased as `marker_seq` advanced.
161    EventSeqRegressed {
162        /// The marker sequence where the regression was observed.
163        marker_seq: u64,
164        /// The previous event sequence boundary.
165        previous_event_seq_before: u64,
166        /// The current event sequence boundary.
167        event_seq_before: u64,
168    },
169    /// `event_seq_before` exceeded the sibling entry run high-watermark.
170    EventSeqExceedsHighWatermark {
171        /// The marker sequence carrying the invalid event boundary.
172        marker_seq: u64,
173        /// The invalid event sequence boundary.
174        event_seq_before: u64,
175        /// The sibling entry run high-watermark.
176        high_watermark: u64,
177    },
178    /// A per-slot cursor count decreased between snapshots.
179    CursorCountRegressed {
180        /// The marker sequence carrying the regressed cursor.
181        marker_seq: u64,
182        /// The stream slot whose cursor regressed.
183        slot: StreamSlot,
184        /// The previous count for this slot.
185        previous_count: u64,
186        /// The current count for this slot.
187        count: u64,
188    },
189    /// A per-slot highest `ts_init` decreased between snapshots.
190    CursorTsInitRegressed {
191        /// The marker sequence carrying the regressed cursor.
192        marker_seq: u64,
193        /// The stream slot whose cursor regressed.
194        slot: StreamSlot,
195        /// The previous highest `ts_init` for this slot.
196        previous_ts_init_hi: UnixNanos,
197        /// The current highest `ts_init` for this slot.
198        ts_init_hi: UnixNanos,
199    },
200    /// A stored durable record hash did not match the recomputed canonical hash.
201    HashMismatch {
202        /// The table kind whose stored hash diverged.
203        record: MarkerRecordKind,
204        /// The marker sequence when the record kind carries one.
205        marker_seq: Option<u64>,
206        /// The stream slot when the record kind carries one.
207        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}