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