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    #[cfg(test)]
457    fn retire_positioned_journal_effect(
458        &mut self,
459        key: &RawDataStoreKey,
460        position: JournalOverlayPosition,
461    ) -> Result<PositionedOverlayRetirement, crate::error::InternalError> {
462        let DataStoreBackend::Journaled { positions, .. } = &self.backend else {
463            return Err(crate::error::InternalError::store_invariant());
464        };
465        let retirement = positions.preflight_retirement(key, position)?;
466        let prepared = PreparedDataPositionRetirement {
467            entries: vec![(key.clone(), retirement)],
468        };
469        self.apply_prepared_position_retirement(prepared);
470        Ok(retirement)
471    }
472
473    /// Apply one folded journal row put into the canonical stable base.
474    pub(in crate::db) fn fold_recovered_journal_put(
475        &mut self,
476        key: RawDataStoreKey,
477        row: RawRow,
478    ) -> Result<Option<RawRow>, crate::error::InternalError> {
479        let DataStoreBackend::Journaled {
480            canonical,
481            live,
482            tombstones,
483            ..
484        } = &mut self.backend
485        else {
486            return Err(crate::error::InternalError::store_invariant());
487        };
488
489        let visible = !live.contains_key(&key) && !tombstones.contains(&key);
490        let cardinality_key = key.clone();
491        let previous = canonical.insert(key, row);
492        if visible {
493            self.entity_cardinality
494                .apply_insert(&cardinality_key, previous.as_ref());
495        } else if previous.is_none() {
496            self.apply_entity_overlay_delta(&cardinality_key, true, false);
497        }
498        self.bump_generation();
499
500        Ok(previous)
501    }
502
503    /// Apply one folded journal row delete into the canonical stable base.
504    pub(in crate::db) fn fold_recovered_journal_delete(
505        &mut self,
506        key: &RawDataStoreKey,
507    ) -> Result<Option<RawRow>, crate::error::InternalError> {
508        let DataStoreBackend::Journaled {
509            canonical,
510            live,
511            tombstones,
512            ..
513        } = &mut self.backend
514        else {
515            return Err(crate::error::InternalError::store_invariant());
516        };
517
518        let visible = !live.contains_key(key) && !tombstones.contains(key);
519        let previous = canonical.remove(key);
520        if visible {
521            self.entity_cardinality.apply_remove(key, previous.as_ref());
522        } else if previous.is_some() {
523            self.apply_entity_overlay_delta(key, false, true);
524        }
525        self.bump_generation();
526
527        Ok(previous)
528    }
529
530    /// Prove that recovered journal rows can be folded into canonical storage.
531    pub(in crate::db) fn preflight_fold_recovered_journal(
532        &self,
533    ) -> Result<(), crate::error::InternalError> {
534        match self.backend {
535            DataStoreBackend::Journaled { .. } => Ok(()),
536            DataStoreBackend::Heap(_) => Err(crate::error::InternalError::store_invariant()),
537        }
538    }
539
540    /// Load one row by raw key.
541    pub(in crate::db) fn get(&self, key: &RawDataStoreKey) -> Option<RawRow> {
542        self.read(key).into_owned()
543    }
544
545    /// Load one row from the canonical predecessor view.
546    ///
547    /// Online journal folding uses this view so a newer positioned live effect
548    /// cannot replace the predecessor evidence for the older batch being
549    /// canonicalized.
550    pub(in crate::db) fn get_canonical(&self, key: &RawDataStoreKey) -> Option<RawRow> {
551        match &self.backend {
552            DataStoreBackend::Heap(map) => map.get(key).cloned(),
553            DataStoreBackend::Journaled { canonical, .. } => canonical.get(key),
554        }
555    }
556
557    /// Read one visible row without cloning heap/live payloads.
558    pub(in crate::db) fn read<'a>(&'a self, key: &RawDataStoreKey) -> StoredRowRead<'a> {
559        #[cfg(test)]
560        record_data_store_get_call();
561
562        match &self.backend {
563            DataStoreBackend::Heap(map) => map
564                .get(key)
565                .map_or(StoredRowRead::Missing, StoredRowRead::Borrowed),
566            DataStoreBackend::Journaled {
567                canonical,
568                live,
569                tombstones,
570                ..
571            } => {
572                if tombstones.contains(key) {
573                    StoredRowRead::Missing
574                } else if let Some(row) = live.get(key) {
575                    StoredRowRead::Borrowed(row)
576                } else {
577                    canonical
578                        .get(key)
579                        .map_or(StoredRowRead::Missing, StoredRowRead::Owned)
580                }
581            }
582        }
583    }
584
585    /// Return whether one raw key exists without cloning the row payload.
586    #[must_use]
587    pub(in crate::db) fn contains(&self, key: &RawDataStoreKey) -> bool {
588        match &self.backend {
589            DataStoreBackend::Heap(map) => map.contains_key(key),
590            DataStoreBackend::Journaled {
591                canonical,
592                live,
593                tombstones,
594                ..
595            } => {
596                !tombstones.contains(key)
597                    && (live.contains_key(key) || canonical.get(key).is_some())
598            }
599        }
600    }
601
602    /// Return the number of stored rows without exposing the backing map.
603    #[must_use]
604    pub(in crate::db) fn len(&self) -> u64 {
605        match &self.backend {
606            DataStoreBackend::Heap(map) => u64::try_from(map.len()).unwrap_or(u64::MAX),
607            DataStoreBackend::Journaled { .. } => {
608                let mut count = 0_u64;
609                let _: Result<(), Infallible> = self.visit_entries(|_key, _row| {
610                    count = count.saturating_add(1);
611                    Ok(StoreVisit::Continue)
612                });
613                count
614            }
615        }
616    }
617
618    /// Return the row-store generation used to prove index metadata freshness.
619    #[must_use]
620    pub(in crate::db) const fn generation(&self) -> u64 {
621        self.generation
622    }
623
624    /// Return an exact current row count for one entity when store metadata is valid.
625    #[must_use]
626    pub(in crate::db) fn exact_entity_count(&self, entity: EntityTag) -> Option<u64> {
627        self.entity_cardinality.exact_count(entity)
628    }
629
630    /// Return the exact live-overlay delta from the canonical base for one entity.
631    #[must_use]
632    pub(in crate::db) fn exact_entity_cardinality_delta(&self, entity: EntityTag) -> Option<i64> {
633        match &self.backend {
634            DataStoreBackend::Heap(_) => Some(0),
635            DataStoreBackend::Journaled {
636                entity_cardinality_delta,
637                ..
638            } => entity_cardinality_delta.exact_delta(entity),
639        }
640    }
641
642    /// Visit raw row entries in canonical storage order.
643    pub(in crate::db) fn visit_entries<E>(
644        &self,
645        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
646    ) -> Result<(), E> {
647        match &self.backend {
648            DataStoreBackend::Heap(map) => {
649                for (key, row) in map {
650                    if visitor(key, row)?.should_stop() {
651                        break;
652                    }
653                }
654            }
655            DataStoreBackend::Journaled { .. } => Self::visit_journaled_entries_in_bounds(
656                &self.backend,
657                (Bound::Unbounded, Bound::Unbounded),
658                visitor,
659            )?,
660        }
661
662        Ok(())
663    }
664
665    /// Visit canonical predecessor rows after one exclusive physical checkpoint.
666    ///
667    /// Optional derived-evidence builders use this view so newer live overlays
668    /// cannot contaminate a generation bound to the canonical fold watermark.
669    pub(in crate::db) fn visit_canonical_entries_after(
670        &self,
671        checkpoint: Option<&RawDataStoreKey>,
672        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<bool, crate::error::InternalError>,
673    ) -> Result<(), crate::error::InternalError> {
674        let lower = checkpoint.map_or(Bound::Unbounded, Bound::Excluded);
675        match &self.backend {
676            DataStoreBackend::Heap(map) => {
677                for (key, row) in map.range((lower, Bound::Unbounded)) {
678                    if visitor(key, row)? {
679                        break;
680                    }
681                }
682            }
683            DataStoreBackend::Journaled { canonical, .. } => {
684                for entry in canonical.range((lower, Bound::Unbounded)) {
685                    if visitor(entry.key(), &entry.value())? {
686                        break;
687                    }
688                }
689            }
690        }
691        Ok(())
692    }
693
694    /// Return whether the canonical predecessor domain is physically empty.
695    ///
696    /// This bounded root observation never walks live overlays or materializes
697    /// row payloads.
698    pub(in crate::db) fn canonical_is_empty(&self) -> Result<bool, crate::error::InternalError> {
699        match &self.backend {
700            DataStoreBackend::Journaled { canonical, .. } => Ok(canonical.is_empty()),
701            DataStoreBackend::Heap(_) => Err(crate::error::InternalError::store_invariant()),
702        }
703    }
704
705    /// Visit raw row entries whose keys belong to the provided storage range.
706    pub(in crate::db) fn visit_range<E>(
707        &self,
708        key_range: impl RangeBounds<RawDataStoreKey>,
709        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
710    ) -> Result<(), E> {
711        let bounds = (key_range.start_bound(), key_range.end_bound());
712        match &self.backend {
713            DataStoreBackend::Heap(map) => {
714                for (key, row) in map.range(bounds) {
715                    if visitor(key, row)?.should_stop() {
716                        break;
717                    }
718                }
719            }
720            DataStoreBackend::Journaled { .. } => {
721                Self::visit_journaled_entries_in_bounds(&self.backend, bounds, visitor)?;
722            }
723        }
724
725        Ok(())
726    }
727
728    /// Visit one ascending row range while allowing the caller to stop after
729    /// seeing the key but before a stable row payload is materialized.
730    ///
731    /// Mixed journal overlays retain the ordinary key-then-point-read path;
732    /// one physical backing can keep the range iterator open for the complete
733    /// scan and avoid a second tree lookup per row.
734    pub(in crate::db) fn try_visit_range_with_row_preflight<E>(
735        &self,
736        key_range: impl RangeBounds<RawDataStoreKey>,
737        mut preflight: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
738        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
739    ) -> Result<Option<bool>, E> {
740        let bounds = (key_range.start_bound(), key_range.end_bound());
741        let mut stopped = false;
742        match &self.backend {
743            DataStoreBackend::Heap(map) => {
744                for (key, row) in map.range(bounds) {
745                    if preflight(key)?.should_stop() {
746                        stopped = true;
747                        break;
748                    }
749                    if visitor(key, row)?.should_stop() {
750                        stopped = true;
751                        break;
752                    }
753                }
754            }
755            DataStoreBackend::Journaled {
756                canonical,
757                live,
758                tombstones,
759                ..
760            } if canonical.is_empty() => {
761                for (key, row) in live.range(bounds) {
762                    if tombstones.contains(key) {
763                        continue;
764                    }
765                    if preflight(key)?.should_stop() {
766                        stopped = true;
767                        break;
768                    }
769                    if visitor(key, row)?.should_stop() {
770                        stopped = true;
771                        break;
772                    }
773                }
774            }
775            DataStoreBackend::Journaled {
776                canonical,
777                live,
778                tombstones,
779                ..
780            } if live.is_empty() && tombstones.is_empty() => {
781                for entry in canonical.range(bounds) {
782                    if preflight(entry.key())?.should_stop() {
783                        stopped = true;
784                        break;
785                    }
786                    if visitor(entry.key(), &entry.value())?.should_stop() {
787                        stopped = true;
788                        break;
789                    }
790                }
791            }
792            DataStoreBackend::Journaled { .. } => return Ok(None),
793        }
794
795        Ok(Some(!stopped))
796    }
797
798    /// Visit only raw keys in storage order without fetching row payloads.
799    ///
800    /// Primary-key access streams use this boundary to discover candidate
801    /// identities before the terminal row runtime decides whether the payload
802    /// is needed. Journaled traversal merges canonical and live keys while
803    /// preserving live overrides and tombstones without reading stable values.
804    pub(in crate::db) fn visit_key_range<E>(
805        &self,
806        key_range: impl RangeBounds<RawDataStoreKey>,
807        visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
808    ) -> Result<(), E> {
809        self.visit_keys_in_bounds(
810            (key_range.start_bound(), key_range.end_bound()),
811            false,
812            visitor,
813        )
814    }
815
816    /// Visit only raw keys in reverse storage order without fetching row payloads.
817    pub(in crate::db) fn visit_key_range_rev<E>(
818        &self,
819        key_range: impl RangeBounds<RawDataStoreKey>,
820        visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
821    ) -> Result<(), E> {
822        self.visit_keys_in_bounds(
823            (key_range.start_bound(), key_range.end_bound()),
824            true,
825            visitor,
826        )
827    }
828
829    /// Sum of bytes used by all stored rows.
830    pub(in crate::db) fn memory_bytes(&self) -> u64 {
831        // Report map footprint as key bytes + row bytes per entry.
832        let mut bytes = 0u64;
833        let _: Result<(), Infallible> = self.visit_entries(|key, row| {
834            bytes = bytes.saturating_add(key.as_bytes().len() as u64 + row.len() as u64);
835            Ok(StoreVisit::Continue)
836        });
837        bytes
838    }
839
840    const fn bump_generation(&mut self) {
841        self.generation = self.generation.saturating_add(1);
842    }
843
844    #[cfg(test)]
845    fn rebuild_entity_cardinality_from_entries(&mut self) {
846        let mut cardinality = EntityCardinality::empty();
847        let _: Result<(), Infallible> = self.visit_entries(|key, _row| {
848            cardinality.apply_present_key(key);
849            Ok(StoreVisit::Continue)
850        });
851        self.entity_cardinality = cardinality;
852    }
853
854    fn apply_entity_overlay_delta(
855        &mut self,
856        key: &RawDataStoreKey,
857        previous_present: bool,
858        next_present: bool,
859    ) {
860        let DataStoreBackend::Journaled {
861            entity_cardinality_delta,
862            ..
863        } = &mut self.backend
864        else {
865            return;
866        };
867        entity_cardinality_delta.apply_presence_transition(key, previous_present, next_present);
868    }
869
870    /// Return the monotonic perf-only count of stable row fetches seen by this process.
871    #[cfg(test)]
872    pub(in crate::db) fn current_get_call_count() -> u64 {
873        DATA_STORE_GET_CALL_COUNT.with(Cell::get)
874    }
875
876    fn visit_keys_in_bounds<E>(
877        &self,
878        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
879        reverse: bool,
880        mut visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
881    ) -> Result<(), E> {
882        match &self.backend {
883            DataStoreBackend::Heap(map) => {
884                if reverse {
885                    for (key, _row) in map.range(bounds).rev() {
886                        if visitor(key)?.should_stop() {
887                            break;
888                        }
889                    }
890                } else {
891                    for (key, _row) in map.range(bounds) {
892                        if visitor(key)?.should_stop() {
893                            break;
894                        }
895                    }
896                }
897            }
898            DataStoreBackend::Journaled { .. } => {
899                Self::visit_journaled_keys_in_bounds(&self.backend, bounds, reverse, visitor)?;
900            }
901        }
902
903        Ok(())
904    }
905
906    fn visit_journaled_keys_in_bounds<E>(
907        backend: &DataStoreBackend,
908        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
909        reverse: bool,
910        mut visitor: impl FnMut(&RawDataStoreKey) -> Result<StoreVisit, E>,
911    ) -> Result<(), E> {
912        let DataStoreBackend::Journaled {
913            canonical,
914            live,
915            tombstones,
916            ..
917        } = backend
918        else {
919            return Ok(());
920        };
921
922        if canonical.is_empty() {
923            if reverse {
924                for (key, _row) in live.range(bounds).rev() {
925                    if visitor(key)?.should_stop() {
926                        return Ok(());
927                    }
928                }
929            } else {
930                for (key, _row) in live.range(bounds) {
931                    if visitor(key)?.should_stop() {
932                        return Ok(());
933                    }
934                }
935            }
936            return Ok(());
937        }
938
939        if live.is_empty() && tombstones.is_empty() {
940            if reverse {
941                for entry in canonical.range(bounds).rev() {
942                    if visitor(entry.key())?.should_stop() {
943                        return Ok(());
944                    }
945                }
946            } else {
947                for entry in canonical.range(bounds) {
948                    if visitor(entry.key())?.should_stop() {
949                        return Ok(());
950                    }
951                }
952            }
953            return Ok(());
954        }
955
956        let direction = if reverse {
957            Direction::Desc
958        } else {
959            Direction::Asc
960        };
961        match direction {
962            Direction::Asc => {
963                for entry in ordered_overlay_entries(
964                    canonical.range(bounds),
965                    live.range(bounds),
966                    direction,
967                    |entry| entry.key(),
968                    |entry| entry.0,
969                    tombstones,
970                ) {
971                    let visit = match entry {
972                        OrderedOverlayEntry::Canonical(canonical_entry) => {
973                            visitor(canonical_entry.key())?
974                        }
975                        OrderedOverlayEntry::Live((key, _row)) => visitor(key)?,
976                    };
977                    if visit.should_stop() {
978                        return Ok(());
979                    }
980                }
981            }
982            Direction::Desc => {
983                for entry in ordered_overlay_entries(
984                    canonical.range(bounds).rev(),
985                    live.range(bounds).rev(),
986                    direction,
987                    |entry| entry.key(),
988                    |entry| entry.0,
989                    tombstones,
990                ) {
991                    let visit = match entry {
992                        OrderedOverlayEntry::Canonical(canonical_entry) => {
993                            visitor(canonical_entry.key())?
994                        }
995                        OrderedOverlayEntry::Live((key, _row)) => visitor(key)?,
996                    };
997                    if visit.should_stop() {
998                        return Ok(());
999                    }
1000                }
1001            }
1002        }
1003
1004        Ok(())
1005    }
1006
1007    fn visit_journaled_entries_in_bounds<E>(
1008        backend: &DataStoreBackend,
1009        bounds: (Bound<&RawDataStoreKey>, Bound<&RawDataStoreKey>),
1010        mut visitor: impl FnMut(&RawDataStoreKey, &RawRow) -> Result<StoreVisit, E>,
1011    ) -> Result<(), E> {
1012        let DataStoreBackend::Journaled {
1013            canonical,
1014            live,
1015            tombstones,
1016            ..
1017        } = backend
1018        else {
1019            return Ok(());
1020        };
1021
1022        if canonical.is_empty() {
1023            for (key, row) in live.range(bounds) {
1024                if visitor(key, row)?.should_stop() {
1025                    return Ok(());
1026                }
1027            }
1028            return Ok(());
1029        }
1030
1031        if live.is_empty() && tombstones.is_empty() {
1032            for entry in canonical.range(bounds) {
1033                if visitor(entry.key(), &entry.value())?.should_stop() {
1034                    return Ok(());
1035                }
1036            }
1037            return Ok(());
1038        }
1039
1040        for entry in ordered_overlay_entries(
1041            canonical.range(bounds),
1042            live.range(bounds),
1043            Direction::Asc,
1044            |entry| entry.key(),
1045            |entry| entry.0,
1046            tombstones,
1047        ) {
1048            let visit = match entry {
1049                OrderedOverlayEntry::Canonical(canonical_entry) => {
1050                    visitor(canonical_entry.key(), &canonical_entry.value())?
1051                }
1052                OrderedOverlayEntry::Live((key, row)) => visitor(key, row)?,
1053            };
1054            if visit.should_stop() {
1055                return Ok(());
1056            }
1057        }
1058
1059        Ok(())
1060    }
1061}
1062
1063#[derive(Clone, Debug)]
1064struct EntityCardinality {
1065    counts: HeapBTreeMap<EntityTag, u64>,
1066    decodable: bool,
1067}
1068
1069#[derive(Clone, Debug)]
1070struct EntityCardinalityDelta {
1071    counts: HeapBTreeMap<EntityTag, i64>,
1072    decodable: bool,
1073}
1074
1075impl EntityCardinalityDelta {
1076    const fn empty() -> Self {
1077        Self {
1078            counts: HeapBTreeMap::new(),
1079            decodable: true,
1080        }
1081    }
1082
1083    fn exact_delta(&self, entity: EntityTag) -> Option<i64> {
1084        self.decodable
1085            .then(|| self.counts.get(&entity).copied().unwrap_or(0))
1086    }
1087
1088    fn apply_presence_transition(
1089        &mut self,
1090        key: &RawDataStoreKey,
1091        previous_present: bool,
1092        next_present: bool,
1093    ) {
1094        if !self.decodable || previous_present == next_present {
1095            return;
1096        }
1097        let Some(entity) = key.entity_tag_prefix() else {
1098            self.counts.clear();
1099            self.decodable = false;
1100            return;
1101        };
1102        let delta = if next_present { 1_i64 } else { -1_i64 };
1103        let count = self.counts.entry(entity).or_insert(0);
1104        let Some(next) = count.checked_add(delta) else {
1105            self.counts.clear();
1106            self.decodable = false;
1107            return;
1108        };
1109        *count = next;
1110        if *count == 0 {
1111            self.counts.remove(&entity);
1112        }
1113    }
1114}
1115
1116impl EntityCardinality {
1117    const fn empty() -> Self {
1118        Self {
1119            counts: HeapBTreeMap::new(),
1120            decodable: true,
1121        }
1122    }
1123
1124    const fn unavailable() -> Self {
1125        Self {
1126            counts: HeapBTreeMap::new(),
1127            decodable: false,
1128        }
1129    }
1130
1131    fn exact_count(&self, entity: EntityTag) -> Option<u64> {
1132        self.decodable
1133            .then(|| self.counts.get(&entity).copied().unwrap_or(0))
1134    }
1135
1136    fn apply_insert(&mut self, key: &RawDataStoreKey, previous: Option<&RawRow>) {
1137        if previous.is_some() {
1138            return;
1139        }
1140        self.apply_present_key(key);
1141    }
1142
1143    fn apply_remove(&mut self, key: &RawDataStoreKey, previous: Option<&RawRow>) {
1144        if previous.is_none() {
1145            return;
1146        }
1147        self.apply_removed_key(key);
1148    }
1149
1150    fn apply_present_key(&mut self, key: &RawDataStoreKey) {
1151        if !self.decodable {
1152            return;
1153        }
1154        let Some(entity) = key.entity_tag_prefix() else {
1155            self.invalidate();
1156            return;
1157        };
1158
1159        let count = self.counts.entry(entity).or_insert(0);
1160        *count = count.saturating_add(1);
1161    }
1162
1163    fn apply_removed_key(&mut self, key: &RawDataStoreKey) {
1164        if !self.decodable {
1165            return;
1166        }
1167        let Some(entity) = key.entity_tag_prefix() else {
1168            self.invalidate();
1169            return;
1170        };
1171
1172        if let Some(count) = self.counts.get_mut(&entity) {
1173            *count = count.saturating_sub(1);
1174            if *count == 0 {
1175                self.counts.remove(&entity);
1176            }
1177        }
1178    }
1179
1180    fn invalidate(&mut self) {
1181        self.counts.clear();
1182        self.decodable = false;
1183    }
1184}
1185
1186#[cfg(test)]
1187mod tests;