Skip to main content

icydb_core/db/data/
store.rs

1//! Module: data::store
2//! Responsibility: journaled-or-heap row storage behind the data-store boundary.
3//! Does not own: key/row validation policy beyond type boundaries.
4//! Boundary: commit/executor call into this layer after prevalidation.
5
6use crate::{
7    db::{
8        data::{CanonicalRow, RawDataStoreKey, RawRow},
9        direction::Direction,
10        ordered_overlay::{OrderedOverlayEntry, ordered_overlay_entries},
11        positioned_overlay::{
12            JournalOverlayPosition, PositionedOverlayMetadata, PositionedOverlayRetirement,
13        },
14    },
15    types::EntityTag,
16};
17use ic_memory::RuntimeMemory;
18use ic_memory::ic_stable_structures::{BTreeMap as StableBTreeMap, DefaultMemoryImpl};
19#[cfg(test)]
20use std::cell::Cell;
21use std::collections::{BTreeMap as HeapBTreeMap, BTreeSet};
22use std::convert::Infallible;
23use std::ops::{Bound, RangeBounds};
24
25#[cfg(test)]
26thread_local! {
27    static DATA_STORE_GET_CALL_COUNT: Cell<u64> = const { Cell::new(0) };
28}
29
30#[cfg(test)]
31fn record_data_store_get_call() {
32    DATA_STORE_GET_CALL_COUNT.with(|count| {
33        count.set(count.get().saturating_add(1));
34    });
35}
36
37///
38/// DataStore
39///
40/// Thin persistence wrapper over one journaled or heap BTreeMap.
41///
42/// Invariant: callers provide already-validated `RawDataStoreKey` and canonical row bytes.
43/// This type intentionally does not enforce commit-phase ordering.
44///
45
46pub struct DataStore {
47    backend: DataStoreBackend,
48    generation: u64,
49    entity_cardinality: EntityCardinality,
50}
51
52enum DataStoreBackend {
53    Heap(HeapBTreeMap<RawDataStoreKey, RawRow>),
54    Journaled {
55        canonical: StableBTreeMap<RawDataStoreKey, RawRow, RuntimeMemory<DefaultMemoryImpl>>,
56        live: HeapBTreeMap<RawDataStoreKey, RawRow>,
57        tombstones: BTreeSet<RawDataStoreKey>,
58        positions: PositionedOverlayMetadata<RawDataStoreKey>,
59        entity_cardinality_delta: EntityCardinalityDelta,
60    },
61}
62
63/// Preflighted provenance publication for direct journal record families.
64#[cfg(any(test, feature = "migration"))]
65pub(in crate::db) struct PreparedDataPositionPublication {
66    keys: Vec<RawDataStoreKey>,
67    position: JournalOverlayPosition,
68}
69
70/// Preflighted exact retirement for one complete journal batch.
71pub(in crate::db) struct PreparedDataPositionRetirement {
72    entries: Vec<(RawDataStoreKey, PositionedOverlayRetirement)>,
73}
74
75/// One visible row read that borrows heap/live state and owns stable state.
76///
77/// Callers that only need selected fields can evaluate a borrowed row while
78/// the store handle is active. Stable-structure reads remain owned because
79/// that backend cannot expose a value reference beyond its storage call.
80pub(in crate::db) enum StoredRowRead<'a> {
81    Missing,
82    Borrowed(&'a RawRow),
83    Owned(RawRow),
84}
85
86impl StoredRowRead<'_> {
87    /// Borrow the visible row regardless of its physical backing.
88    #[must_use]
89    pub(in crate::db) const fn as_row(&self) -> Option<&RawRow> {
90        match self {
91            Self::Missing => None,
92            Self::Borrowed(row) => Some(row),
93            Self::Owned(row) => Some(row),
94        }
95    }
96
97    /// Convert the visible row into the existing owned get contract.
98    #[must_use]
99    fn into_owned(self) -> Option<RawRow> {
100        match self {
101            Self::Missing => None,
102            Self::Borrowed(row) => Some(row.clone()),
103            Self::Owned(row) => Some(row),
104        }
105    }
106}
107
108/// Control-flow result for store traversal visitors.
109#[derive(Clone, Copy, Debug, Eq, PartialEq)]
110pub(in crate::db) enum StoreVisit {
111    Continue,
112    Stop,
113}
114
115impl StoreVisit {
116    const fn should_stop(self) -> bool {
117        matches!(self, Self::Stop)
118    }
119}
120
121impl DataStore {
122    /// Initialize a volatile heap-backed data store.
123    #[must_use]
124    pub const fn init_heap() -> Self {
125        Self {
126            backend: DataStoreBackend::Heap(HeapBTreeMap::new()),
127            generation: 0,
128            entity_cardinality: EntityCardinality::empty(),
129        }
130    }
131
132    /// Initialize a journaled cached-stable data store.
133    ///
134    /// Normal writes update only the live projection. The canonical stable map
135    /// is the future fold target and is not mutated by this wrapper's write
136    /// methods.
137    #[must_use]
138    pub fn init_journaled(memory: RuntimeMemory<DefaultMemoryImpl>) -> Self {
139        let canonical = StableBTreeMap::init(memory);
140        let entity_cardinality = if canonical.is_empty() {
141            EntityCardinality::empty()
142        } else {
143            EntityCardinality::unavailable()
144        };
145        Self {
146            backend: DataStoreBackend::Journaled {
147                canonical,
148                live: HeapBTreeMap::new(),
149                tombstones: BTreeSet::new(),
150                positions: PositionedOverlayMetadata::new(),
151                entity_cardinality_delta: EntityCardinalityDelta::empty(),
152            },
153            generation: 0,
154            // Stable rows remain authoritative after reinitialization. Exact
155            // zero cardinality is still known for an empty canonical map;
156            // populated maps remain unavailable without a startup scan.
157            entity_cardinality,
158        }
159    }
160
161    /// Insert or replace one row by raw key.
162    pub(in crate::db) fn insert(
163        &mut self,
164        key: RawDataStoreKey,
165        row: CanonicalRow,
166    ) -> Option<RawRow> {
167        let row = row.into_raw_row();
168        let previous_journaled = if matches!(self.backend, DataStoreBackend::Journaled { .. }) {
169            self.get(&key)
170        } else {
171            None
172        };
173        let cardinality_key = key.clone();
174        let previous = match &mut self.backend {
175            DataStoreBackend::Heap(map) => map.insert(key, row),
176            DataStoreBackend::Journaled {
177                live, tombstones, ..
178            } => {
179                tombstones.remove(&key);
180                live.insert(key, row);
181                previous_journaled
182            }
183        };
184        self.entity_cardinality
185            .apply_insert(&cardinality_key, previous.as_ref());
186        self.apply_entity_overlay_delta(&cardinality_key, previous.is_some(), true);
187        self.bump_generation();
188        previous
189    }
190
191    /// Insert one raw row directly for corruption-focused test setup only.
192    #[cfg(test)]
193    pub(in crate::db) fn insert_raw_for_test(
194        &mut self,
195        key: RawDataStoreKey,
196        row: RawRow,
197    ) -> Option<RawRow> {
198        let previous_journaled = if matches!(self.backend, DataStoreBackend::Journaled { .. }) {
199            self.get(&key)
200        } else {
201            None
202        };
203        let cardinality_key = key.clone();
204        let previous = match &mut self.backend {
205            DataStoreBackend::Heap(map) => map.insert(key, row),
206            DataStoreBackend::Journaled {
207                live, tombstones, ..
208            } => {
209                tombstones.remove(&key);
210                live.insert(key, row);
211                previous_journaled
212            }
213        };
214        self.entity_cardinality
215            .apply_insert(&cardinality_key, previous.as_ref());
216        self.apply_entity_overlay_delta(&cardinality_key, previous.is_some(), true);
217        self.bump_generation();
218        previous
219    }
220
221    /// Remove one row by raw key.
222    pub(in crate::db) fn remove(&mut self, key: &RawDataStoreKey) -> Option<RawRow> {
223        let previous_journaled = if matches!(self.backend, DataStoreBackend::Journaled { .. }) {
224            self.get(key)
225        } else {
226            None
227        };
228        let previous = match &mut self.backend {
229            DataStoreBackend::Heap(map) => map.remove(key),
230            DataStoreBackend::Journaled {
231                live, tombstones, ..
232            } => {
233                live.remove(key);
234                tombstones.insert(key.clone());
235                previous_journaled
236            }
237        };
238        self.entity_cardinality.apply_remove(key, previous.as_ref());
239        self.apply_entity_overlay_delta(key, previous.is_some(), false);
240        self.bump_generation();
241        previous
242    }
243
244    /// Reset the volatile projection for journaled recovery without mutating
245    /// the canonical stable base.
246    pub(in crate::db) fn reset_journaled_live_projection(
247        &mut self,
248    ) -> Result<(), crate::error::InternalError> {
249        let DataStoreBackend::Journaled {
250            canonical,
251            live,
252            tombstones,
253            positions,
254            entity_cardinality_delta,
255        } = &mut self.backend
256        else {
257            return Err(crate::error::InternalError::store_invariant());
258        };
259
260        live.clear();
261        tombstones.clear();
262        positions.clear();
263        *entity_cardinality_delta = EntityCardinalityDelta::empty();
264        self.entity_cardinality = if canonical.is_empty() {
265            EntityCardinality::empty()
266        } else {
267            EntityCardinality::unavailable()
268        };
269        self.bump_generation();
270
271        Ok(())
272    }
273
274    /// Apply one recovered journal row put into the volatile projection.
275    #[cfg(any(test, feature = "migration"))]
276    pub(in crate::db) fn apply_recovered_journal_put(
277        &mut self,
278        key: RawDataStoreKey,
279        row: RawRow,
280    ) -> Result<Option<RawRow>, crate::error::InternalError> {
281        let DataStoreBackend::Journaled {
282            canonical,
283            live,
284            tombstones,
285            ..
286        } = &mut self.backend
287        else {
288            return Err(crate::error::InternalError::store_invariant());
289        };
290
291        let previous = if tombstones.contains(&key) {
292            None
293        } else {
294            live.get(&key).cloned().or_else(|| canonical.get(&key))
295        };
296        tombstones.remove(&key);
297        let cardinality_key = key.clone();
298        live.insert(key, row);
299        self.entity_cardinality
300            .apply_insert(&cardinality_key, previous.as_ref());
301        self.apply_entity_overlay_delta(&cardinality_key, previous.is_some(), true);
302        self.bump_generation();
303
304        Ok(previous)
305    }
306
307    /// Publish one preflighted positioned row value or tombstone.
308    pub(in crate::db) fn publish_preflighted_journal_entry(
309        &mut self,
310        key: RawDataStoreKey,
311        row: Option<RawRow>,
312        position: JournalOverlayPosition,
313    ) -> Result<Option<RawRow>, crate::error::InternalError> {
314        let DataStoreBackend::Journaled {
315            canonical,
316            live,
317            tombstones,
318            positions,
319            entity_cardinality_delta,
320        } = &mut self.backend
321        else {
322            return Err(crate::error::InternalError::store_invariant());
323        };
324        let previous = if tombstones.contains(&key) {
325            None
326        } else {
327            live.get(&key).cloned().or_else(|| canonical.get(&key))
328        };
329        let next_present = row.is_some();
330        if let Some(row) = row {
331            tombstones.remove(&key);
332            live.insert(key.clone(), row);
333            self.entity_cardinality
334                .apply_insert(&key, previous.as_ref());
335        } else {
336            live.remove(&key);
337            tombstones.insert(key.clone());
338            self.entity_cardinality
339                .apply_remove(&key, previous.as_ref());
340        }
341        entity_cardinality_delta.apply_presence_transition(&key, previous.is_some(), next_present);
342        positions.publish_preflighted(key, position);
343        self.bump_generation();
344
345        Ok(previous)
346    }
347
348    /// Validate and publish one positioned row for direct store tests.
349    #[cfg(test)]
350    pub(in crate::db) fn publish_positioned_journal_entry(
351        &mut self,
352        key: RawDataStoreKey,
353        row: Option<RawRow>,
354        position: JournalOverlayPosition,
355    ) -> Result<Option<RawRow>, crate::error::InternalError> {
356        self.preflight_positioned_journal_entry(&key, position)?;
357        self.publish_preflighted_journal_entry(key, row, position)
358    }
359
360    /// Preflight row provenance before marker publication.
361    pub(in crate::db) fn preflight_positioned_journal_entry(
362        &self,
363        key: &RawDataStoreKey,
364        position: JournalOverlayPosition,
365    ) -> Result<(), crate::error::InternalError> {
366        let DataStoreBackend::Journaled { positions, .. } = &self.backend else {
367            return Err(crate::error::InternalError::store_invariant());
368        };
369        positions.preflight_publish(key, position)
370    }
371
372    /// Preflight direct row provenance before marker publication.
373    #[cfg(any(test, feature = "migration"))]
374    pub(in crate::db) fn prepare_position_publication(
375        &self,
376        keys: impl IntoIterator<Item = RawDataStoreKey>,
377        position: JournalOverlayPosition,
378    ) -> Result<PreparedDataPositionPublication, crate::error::InternalError> {
379        let DataStoreBackend::Journaled { positions, .. } = &self.backend else {
380            return Err(crate::error::InternalError::store_invariant());
381        };
382        let keys = keys.into_iter().collect::<BTreeSet<_>>();
383        for key in &keys {
384            positions.preflight_publish(key, position)?;
385        }
386        Ok(PreparedDataPositionPublication {
387            keys: keys.into_iter().collect(),
388            position,
389        })
390    }
391
392    /// Publish direct row provenance after its values have been applied.
393    #[cfg(any(test, feature = "migration"))]
394    pub(in crate::db) fn publish_prepared_positions(
395        &mut self,
396        prepared: PreparedDataPositionPublication,
397    ) {
398        let DataStoreBackend::Journaled { positions, .. } = &mut self.backend else {
399            debug_assert!(false, "preflighted row positions require a journaled store");
400            return;
401        };
402        for key in prepared.keys {
403            positions.publish_preflighted(key, prepared.position);
404        }
405    }
406
407    /// Preflight exact row-overlay retirement before canonical mutation.
408    pub(in crate::db) fn prepare_position_retirement(
409        &self,
410        keys: impl IntoIterator<Item = RawDataStoreKey>,
411        position: JournalOverlayPosition,
412    ) -> Result<PreparedDataPositionRetirement, crate::error::InternalError> {
413        let DataStoreBackend::Journaled { positions, .. } = &self.backend else {
414            return Err(crate::error::InternalError::store_invariant());
415        };
416        let entries = keys
417            .into_iter()
418            .collect::<BTreeSet<_>>()
419            .into_iter()
420            .map(|key| {
421                positions
422                    .preflight_retirement(&key, position)
423                    .map(|retirement| (key, retirement))
424            })
425            .collect::<Result<Vec<_>, _>>()?;
426        Ok(PreparedDataPositionRetirement { entries })
427    }
428
429    /// Retire only exact row overlays after canonical mutation succeeds.
430    pub(in crate::db) fn apply_prepared_position_retirement(
431        &mut self,
432        prepared: PreparedDataPositionRetirement,
433    ) {
434        let DataStoreBackend::Journaled {
435            live,
436            tombstones,
437            positions,
438            ..
439        } = &mut self.backend
440        else {
441            debug_assert!(
442                false,
443                "preflighted row retirement requires a journaled store"
444            );
445            return;
446        };
447        for (key, retirement) in prepared.entries {
448            if retirement == PositionedOverlayRetirement::Exact {
449                live.remove(&key);
450                tombstones.remove(&key);
451                positions.retire_preflighted(&key, retirement);
452            }
453        }
454    }
455
456    /// Apply one folded journal row put into the canonical stable base.
457    pub(in crate::db) fn fold_recovered_journal_put(
458        &mut self,
459        key: RawDataStoreKey,
460        row: RawRow,
461    ) -> Result<Option<RawRow>, crate::error::InternalError> {
462        let DataStoreBackend::Journaled {
463            canonical,
464            live,
465            tombstones,
466            ..
467        } = &mut self.backend
468        else {
469            return Err(crate::error::InternalError::store_invariant());
470        };
471
472        let visible = !live.contains_key(&key) && !tombstones.contains(&key);
473        let cardinality_key = key.clone();
474        let previous = canonical.insert(key, row);
475        if visible {
476            self.entity_cardinality
477                .apply_insert(&cardinality_key, previous.as_ref());
478        } else if previous.is_none() {
479            self.apply_entity_overlay_delta(&cardinality_key, true, false);
480        }
481        self.bump_generation();
482
483        Ok(previous)
484    }
485
486    /// Apply one folded journal row delete into the canonical stable base.
487    pub(in crate::db) fn fold_recovered_journal_delete(
488        &mut self,
489        key: &RawDataStoreKey,
490    ) -> Result<Option<RawRow>, crate::error::InternalError> {
491        let DataStoreBackend::Journaled {
492            canonical,
493            live,
494            tombstones,
495            ..
496        } = &mut self.backend
497        else {
498            return Err(crate::error::InternalError::store_invariant());
499        };
500
501        let visible = !live.contains_key(key) && !tombstones.contains(key);
502        let previous = canonical.remove(key);
503        if visible {
504            self.entity_cardinality.apply_remove(key, previous.as_ref());
505        } else if previous.is_some() {
506            self.apply_entity_overlay_delta(key, false, true);
507        }
508        self.bump_generation();
509
510        Ok(previous)
511    }
512
513    /// Prove that recovered journal rows can be folded into canonical storage.
514    pub(in crate::db) fn preflight_fold_recovered_journal(
515        &self,
516    ) -> Result<(), crate::error::InternalError> {
517        match self.backend {
518            DataStoreBackend::Journaled { .. } => Ok(()),
519            DataStoreBackend::Heap(_) => Err(crate::error::InternalError::store_invariant()),
520        }
521    }
522
523    /// Load one row by raw key.
524    pub(in crate::db) fn get(&self, key: &RawDataStoreKey) -> Option<RawRow> {
525        self.read(key).into_owned()
526    }
527
528    /// Load one row from the canonical predecessor view.
529    ///
530    /// Online journal folding uses this view so a newer positioned live effect
531    /// cannot replace the predecessor evidence for the older batch being
532    /// canonicalized.
533    pub(in crate::db) fn get_canonical(&self, key: &RawDataStoreKey) -> Option<RawRow> {
534        match &self.backend {
535            DataStoreBackend::Heap(map) => map.get(key).cloned(),
536            DataStoreBackend::Journaled { canonical, .. } => canonical.get(key),
537        }
538    }
539
540    /// Read one visible row without cloning heap/live payloads.
541    pub(in crate::db) fn read<'a>(&'a self, key: &RawDataStoreKey) -> StoredRowRead<'a> {
542        #[cfg(test)]
543        record_data_store_get_call();
544
545        match &self.backend {
546            DataStoreBackend::Heap(map) => map
547                .get(key)
548                .map_or(StoredRowRead::Missing, StoredRowRead::Borrowed),
549            DataStoreBackend::Journaled {
550                canonical,
551                live,
552                tombstones,
553                ..
554            } => {
555                if tombstones.contains(key) {
556                    StoredRowRead::Missing
557                } else if let Some(row) = live.get(key) {
558                    StoredRowRead::Borrowed(row)
559                } else {
560                    canonical
561                        .get(key)
562                        .map_or(StoredRowRead::Missing, StoredRowRead::Owned)
563                }
564            }
565        }
566    }
567
568    /// Return whether one raw key exists without cloning the row payload.
569    #[must_use]
570    pub(in crate::db) fn contains(&self, key: &RawDataStoreKey) -> bool {
571        match &self.backend {
572            DataStoreBackend::Heap(map) => map.contains_key(key),
573            DataStoreBackend::Journaled {
574                canonical,
575                live,
576                tombstones,
577                ..
578            } => {
579                !tombstones.contains(key)
580                    && (live.contains_key(key) || canonical.get(key).is_some())
581            }
582        }
583    }
584
585    /// Return the number of stored rows without exposing the backing map.
586    #[must_use]
587    pub(in crate::db) fn len(&self) -> u64 {
588        match &self.backend {
589            DataStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
590            DataStoreBackend::Journaled { .. } => {
591                let mut count = 0_u64;
592                let _: Result<(), Infallible> = self.visit_entries(|_key, _row| {
593                    count = count.saturating_add(1);
594                    Ok(StoreVisit::Continue)
595                });
596                count
597            }
598        }
599    }
600
601    /// Return the row-store generation used to prove index metadata freshness.
602    #[must_use]
603    pub(in crate::db) const fn generation(&self) -> u64 {
604        self.generation
605    }
606
607    /// Return an exact current row count for one entity when store metadata is valid.
608    #[must_use]
609    pub(in crate::db) fn exact_entity_count(&self, entity: EntityTag) -> Option<u64> {
610        self.entity_cardinality.exact_count(entity)
611    }
612
613    /// Return the exact live-overlay delta from the canonical base for one entity.
614    #[must_use]
615    pub(in crate::db) fn exact_entity_cardinality_delta(&self, entity: EntityTag) -> Option<i64> {
616        match &self.backend {
617            DataStoreBackend::Heap(_) => Some(0),
618            DataStoreBackend::Journaled {
619                entity_cardinality_delta,
620                ..
621            } => entity_cardinality_delta.exact_delta(entity),
622        }
623    }
624
625    /// Visit raw row entries in canonical storage order.
626    pub(in crate::db) fn visit_entries<E>(
627        &self,
628        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
629    ) -> Result<(), E> {
630        match &self.backend {
631            DataStoreBackend::Heap(map) => {
632                for (key, row) in map {
633                    if visitor(key, row)?.should_stop() {
634                        break;
635                    }
636                }
637            }
638            DataStoreBackend::Journaled { .. } => Self::visit_journaled_entries_in_bounds(
639                &self.backend,
640                (Bound::Unbounded, Bound::Unbounded),
641                visitor,
642            )?,
643        }
644
645        Ok(())
646    }
647
648    /// Visit canonical predecessor rows after one exclusive physical checkpoint.
649    ///
650    /// Optional derived-evidence builders use this view so newer live overlays
651    /// cannot contaminate a generation bound to the canonical fold watermark.
652    pub(in crate::db) fn visit_canonical_entries_after(
653        &self,
654        checkpoint: Option<&RawDataStoreKey>,
655        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<bool, crate::error::InternalError>,
656    ) -> Result<(), crate::error::InternalError> {
657        let lower = checkpoint.map_or(Bound::Unbounded, Bound::Excluded);
658        match &self.backend {
659            DataStoreBackend::Heap(map) => {
660                for (key, row) in map.range((lower, Bound::Unbounded)) {
661                    if visitor(key, row)? {
662                        break;
663                    }
664                }
665            }
666            DataStoreBackend::Journaled { canonical, .. } => {
667                for entry in canonical.range((lower, Bound::Unbounded)) {
668                    if visitor(entry.key(), &entry.value())? {
669                        break;
670                    }
671                }
672            }
673        }
674        Ok(())
675    }
676
677    /// Return whether the canonical predecessor domain is physically empty.
678    ///
679    /// This bounded root observation never walks live overlays or materializes
680    /// row payloads.
681    pub(in crate::db) fn canonical_is_empty(&self) -> Result<bool, crate::error::InternalError> {
682        match &self.backend {
683            DataStoreBackend::Journaled { canonical, .. } => Ok(canonical.is_empty()),
684            DataStoreBackend::Heap(_) => Err(crate::error::InternalError::store_invariant()),
685        }
686    }
687
688    /// Visit raw row entries whose keys belong to the provided storage range.
689    pub(in crate::db) fn visit_range<E>(
690        &self,
691        key_range: impl RangeBounds<RawDataStoreKey>,
692        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
693    ) -> Result<(), E> {
694        let bounds = (key_range.start_bound(), key_range.end_bound());
695        match &self.backend {
696            DataStoreBackend::Heap(map) => {
697                for (key, row) in map.range(bounds) {
698                    if visitor(key, row)?.should_stop() {
699                        break;
700                    }
701                }
702            }
703            DataStoreBackend::Journaled { .. } => {
704                Self::visit_journaled_entries_in_bounds(&self.backend, bounds, visitor)?;
705            }
706        }
707
708        Ok(())
709    }
710
711    /// Visit one ascending row range while allowing the caller to stop after
712    /// seeing the key but before a stable row payload is materialized.
713    ///
714    /// Mixed journal overlays retain the ordinary key-then-point-read path;
715    /// one physical backing can keep the range iterator open for the complete
716    /// scan and avoid a second tree lookup per row.
717    pub(in crate::db) fn try_visit_range_with_row_preflight<E>(
718        &self,
719        key_range: impl RangeBounds<RawDataStoreKey>,
720        mut preflight: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
721        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
722    ) -> Result<Option<bool>, E> {
723        let bounds = (key_range.start_bound(), key_range.end_bound());
724        let mut stopped = false;
725        match &self.backend {
726            DataStoreBackend::Heap(map) => {
727                for (key, row) in map.range(bounds) {
728                    if preflight(key)?.should_stop() {
729                        stopped = true;
730                        break;
731                    }
732                    if visitor(key, row)?.should_stop() {
733                        stopped = true;
734                        break;
735                    }
736                }
737            }
738            DataStoreBackend::Journaled {
739                canonical,
740                live,
741                tombstones,
742                ..
743            } if canonical.is_empty() => {
744                for (key, row) in live.range(bounds) {
745                    if tombstones.contains(key) {
746                        continue;
747                    }
748                    if preflight(key)?.should_stop() {
749                        stopped = true;
750                        break;
751                    }
752                    if visitor(key, row)?.should_stop() {
753                        stopped = true;
754                        break;
755                    }
756                }
757            }
758            DataStoreBackend::Journaled {
759                canonical,
760                live,
761                tombstones,
762                ..
763            } if live.is_empty() && tombstones.is_empty() => {
764                for entry in canonical.range(bounds) {
765                    if preflight(entry.key())?.should_stop() {
766                        stopped = true;
767                        break;
768                    }
769                    if visitor(entry.key(), &entry.value())?.should_stop() {
770                        stopped = true;
771                        break;
772                    }
773                }
774            }
775            DataStoreBackend::Journaled { .. } => return Ok(None),
776        }
777
778        Ok(Some(!stopped))
779    }
780
781    /// Visit only raw keys in storage order without fetching row payloads.
782    ///
783    /// Primary-key access streams use this boundary to discover candidate
784    /// identities before the terminal row runtime decides whether the payload
785    /// is needed. Journaled traversal merges canonical and live keys while
786    /// preserving live overrides and tombstones without reading stable values.
787    pub(in crate::db) fn visit_key_range<E>(
788        &self,
789        key_range: impl RangeBounds<RawDataStoreKey>,
790        visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
791    ) -> Result<(), E> {
792        self.visit_keys_in_bounds(
793            (key_range.start_bound(), key_range.end_bound()),
794            false,
795            visitor,
796        )
797    }
798
799    /// Visit only raw keys in reverse storage order without fetching row payloads.
800    pub(in crate::db) fn visit_key_range_rev<E>(
801        &self,
802        key_range: impl RangeBounds<RawDataStoreKey>,
803        visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
804    ) -> Result<(), E> {
805        self.visit_keys_in_bounds(
806            (key_range.start_bound(), key_range.end_bound()),
807            true,
808            visitor,
809        )
810    }
811
812    /// Sum of bytes used by all stored rows.
813    pub(in crate::db) fn memory_bytes(&self) -> u64 {
814        // Report map footprint as key bytes + row bytes per entry.
815        let mut bytes = 0u64;
816        let _: Result<(), Infallible> = self.visit_entries(|key, row| {
817            bytes = bytes.saturating_add(key.as_bytes().len() as u64 + row.len() as u64);
818            Ok(StoreVisit::Continue)
819        });
820        bytes
821    }
822
823    const fn bump_generation(&mut self) {
824        self.generation = self.generation.saturating_add(1);
825    }
826
827    #[cfg(test)]
828    fn rebuild_entity_cardinality_from_entries(&mut self) {
829        let mut cardinality = EntityCardinality::empty();
830        let _: Result<(), Infallible> = self.visit_entries(|key, _row| {
831            cardinality.apply_present_key(key);
832            Ok(StoreVisit::Continue)
833        });
834        self.entity_cardinality = cardinality;
835    }
836
837    fn apply_entity_overlay_delta(
838        &mut self,
839        key: &RawDataStoreKey,
840        previous_present: bool,
841        next_present: bool,
842    ) {
843        let DataStoreBackend::Journaled {
844            entity_cardinality_delta,
845            ..
846        } = &mut self.backend
847        else {
848            return;
849        };
850        entity_cardinality_delta.apply_presence_transition(key, previous_present, next_present);
851    }
852
853    /// Return the monotonic perf-only count of stable row fetches seen by this process.
854    #[cfg(test)]
855    pub(in crate::db) fn current_get_call_count() -> u64 {
856        DATA_STORE_GET_CALL_COUNT.with(Cell::get)
857    }
858
859    fn visit_keys_in_bounds<E>(
860        &self,
861        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
862        reverse: bool,
863        mut visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
864    ) -> Result<(), E> {
865        match &self.backend {
866            DataStoreBackend::Heap(map) => {
867                if reverse {
868                    for (key, _row) in map.range(bounds).rev() {
869                        if visitor(key)?.should_stop() {
870                            break;
871                        }
872                    }
873                } else {
874                    for (key, _row) in map.range(bounds) {
875                        if visitor(key)?.should_stop() {
876                            break;
877                        }
878                    }
879                }
880            }
881            DataStoreBackend::Journaled { .. } => {
882                Self::visit_journaled_keys_in_bounds(&self.backend, bounds, reverse, visitor)?;
883            }
884        }
885
886        Ok(())
887    }
888
889    fn visit_journaled_keys_in_bounds<E>(
890        backend: &DataStoreBackend,
891        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
892        reverse: bool,
893        mut visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
894    ) -> Result<(), E> {
895        let DataStoreBackend::Journaled {
896            canonical,
897            live,
898            tombstones,
899            ..
900        } = backend
901        else {
902            return Ok(());
903        };
904
905        if canonical.is_empty() {
906            if reverse {
907                for (key, _row) in live.range(bounds).rev() {
908                    if visitor(key)?.should_stop() {
909                        return Ok(());
910                    }
911                }
912            } else {
913                for (key, _row) in live.range(bounds) {
914                    if visitor(key)?.should_stop() {
915                        return Ok(());
916                    }
917                }
918            }
919            return Ok(());
920        }
921
922        if live.is_empty() && tombstones.is_empty() {
923            if reverse {
924                for entry in canonical.range(bounds).rev() {
925                    if visitor(entry.key())?.should_stop() {
926                        return Ok(());
927                    }
928                }
929            } else {
930                for entry in canonical.range(bounds) {
931                    if visitor(entry.key())?.should_stop() {
932                        return Ok(());
933                    }
934                }
935            }
936            return Ok(());
937        }
938
939        let direction = if reverse {
940            Direction::Desc
941        } else {
942            Direction::Asc
943        };
944        match direction {
945            Direction::Asc => {
946                for entry in ordered_overlay_entries(
947                    canonical.range(bounds),
948                    live.range(bounds),
949                    direction,
950                    |entry| entry.key(),
951                    |entry| entry.0,
952                    tombstones,
953                ) {
954                    let visit = match entry {
955                        OrderedOverlayEntry::Canonical(canonical_entry) => {
956                            visitor(canonical_entry.key())?
957                        }
958                        OrderedOverlayEntry::Live((key, _row)) => visitor(key)?,
959                    };
960                    if visit.should_stop() {
961                        return Ok(());
962                    }
963                }
964            }
965            Direction::Desc => {
966                for entry in ordered_overlay_entries(
967                    canonical.range(bounds).rev(),
968                    live.range(bounds).rev(),
969                    direction,
970                    |entry| entry.key(),
971                    |entry| entry.0,
972                    tombstones,
973                ) {
974                    let visit = match entry {
975                        OrderedOverlayEntry::Canonical(canonical_entry) => {
976                            visitor(canonical_entry.key())?
977                        }
978                        OrderedOverlayEntry::Live((key, _row)) => visitor(key)?,
979                    };
980                    if visit.should_stop() {
981                        return Ok(());
982                    }
983                }
984            }
985        }
986
987        Ok(())
988    }
989
990    fn visit_journaled_entries_in_bounds<E>(
991        backend: &DataStoreBackend,
992        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
993        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
994    ) -> Result<(), E> {
995        let DataStoreBackend::Journaled {
996            canonical,
997            live,
998            tombstones,
999            ..
1000        } = backend
1001        else {
1002            return Ok(());
1003        };
1004
1005        if canonical.is_empty() {
1006            for (key, row) in live.range(bounds) {
1007                if visitor(key, row)?.should_stop() {
1008                    return Ok(());
1009                }
1010            }
1011            return Ok(());
1012        }
1013
1014        if live.is_empty() && tombstones.is_empty() {
1015            for entry in canonical.range(bounds) {
1016                if visitor(entry.key(), &entry.value())?.should_stop() {
1017                    return Ok(());
1018                }
1019            }
1020            return Ok(());
1021        }
1022
1023        for entry in ordered_overlay_entries(
1024            canonical.range(bounds),
1025            live.range(bounds),
1026            Direction::Asc,
1027            |entry| entry.key(),
1028            |entry| entry.0,
1029            tombstones,
1030        ) {
1031            let visit = match entry {
1032                OrderedOverlayEntry::Canonical(canonical_entry) => {
1033                    visitor(canonical_entry.key(), &canonical_entry.value())?
1034                }
1035                OrderedOverlayEntry::Live((key, row)) => visitor(key, row)?,
1036            };
1037            if visit.should_stop() {
1038                return Ok(());
1039            }
1040        }
1041
1042        Ok(())
1043    }
1044}
1045
1046#[derive(Clone, Debug)]
1047struct EntityCardinality {
1048    counts: HeapBTreeMap<EntityTag, u64>,
1049    decodable: bool,
1050}
1051
1052#[derive(Clone, Debug)]
1053struct EntityCardinalityDelta {
1054    counts: HeapBTreeMap<EntityTag, i64>,
1055    decodable: bool,
1056}
1057
1058impl EntityCardinalityDelta {
1059    const fn empty() -> Self {
1060        Self {
1061            counts: HeapBTreeMap::new(),
1062            decodable: true,
1063        }
1064    }
1065
1066    fn exact_delta(&self, entity: EntityTag) -> Option<i64> {
1067        self.decodable
1068            .then(|| self.counts.get(&entity).copied().unwrap_or(0))
1069    }
1070
1071    fn apply_presence_transition(
1072        &mut self,
1073        key: &RawDataStoreKey,
1074        previous_present: bool,
1075        next_present: bool,
1076    ) {
1077        if !self.decodable || previous_present == next_present {
1078            return;
1079        }
1080        let Some(entity) = key.entity_tag_prefix() else {
1081            self.counts.clear();
1082            self.decodable = false;
1083            return;
1084        };
1085        let delta = if next_present { 1_i64 } else { -1_i64 };
1086        let count = self.counts.entry(entity).or_insert(0);
1087        let Some(next) = count.checked_add(delta) else {
1088            self.counts.clear();
1089            self.decodable = false;
1090            return;
1091        };
1092        *count = next;
1093        if *count == 0 {
1094            self.counts.remove(&entity);
1095        }
1096    }
1097}
1098
1099impl EntityCardinality {
1100    const fn empty() -> Self {
1101        Self {
1102            counts: HeapBTreeMap::new(),
1103            decodable: true,
1104        }
1105    }
1106
1107    const fn unavailable() -> Self {
1108        Self {
1109            counts: HeapBTreeMap::new(),
1110            decodable: false,
1111        }
1112    }
1113
1114    fn exact_count(&self, entity: EntityTag) -> Option<u64> {
1115        self.decodable
1116            .then(|| self.counts.get(&entity).copied().unwrap_or(0))
1117    }
1118
1119    fn apply_insert(&mut self, key: &RawDataStoreKey, previous: Option<&RawRow>) {
1120        if previous.is_some() {
1121            return;
1122        }
1123        self.apply_present_key(key);
1124    }
1125
1126    fn apply_remove(&mut self, key: &RawDataStoreKey, previous: Option<&RawRow>) {
1127        if previous.is_none() {
1128            return;
1129        }
1130        self.apply_removed_key(key);
1131    }
1132
1133    fn apply_present_key(&mut self, key: &RawDataStoreKey) {
1134        if !self.decodable {
1135            return;
1136        }
1137        let Some(entity) = key.entity_tag_prefix() else {
1138            self.invalidate();
1139            return;
1140        };
1141
1142        let count = self.counts.entry(entity).or_insert(0);
1143        *count = count.saturating_add(1);
1144    }
1145
1146    fn apply_removed_key(&mut self, key: &RawDataStoreKey) {
1147        if !self.decodable {
1148            return;
1149        }
1150        let Some(entity) = key.entity_tag_prefix() else {
1151            self.invalidate();
1152            return;
1153        };
1154
1155        if let Some(count) = self.counts.get_mut(&entity) {
1156            *count = count.saturating_sub(1);
1157            if *count == 0 {
1158                self.counts.remove(&entity);
1159            }
1160        }
1161    }
1162
1163    fn invalidate(&mut self) {
1164        self.counts.clear();
1165        self.decodable = false;
1166    }
1167}
1168
1169#[cfg(test)]
1170mod tests;