1use 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
37pub 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#[cfg(any(test, feature = "migration"))]
65pub(in crate::db) struct PreparedDataPositionPublication {
66 keys: Vec<RawDataStoreKey>,
67 position: JournalOverlayPosition,
68}
69
70pub(in crate::db) struct PreparedDataPositionRetirement {
72 entries: Vec<(RawDataStoreKey, PositionedOverlayRetirement)>,
73}
74
75pub(in crate::db) enum StoredRowRead<'a> {
81 Missing,
82 Borrowed(&'a RawRow),
83 Owned(RawRow),
84}
85
86impl StoredRowRead<'_> {
87 #[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 #[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#[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 #[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 #[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 entity_cardinality,
158 }
159 }
160
161 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 #[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 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 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 #[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 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 #[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 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 #[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 #[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 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 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 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 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 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 pub(in crate::db) fn get(&self, key: &RawDataStoreKey) -> Option<RawRow> {
525 self.read(key).into_owned()
526 }
527
528 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 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 #[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 #[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 #[must_use]
603 pub(in crate::db) const fn generation(&self) -> u64 {
604 self.generation
605 }
606
607 #[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 #[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 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 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 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 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 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 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 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 pub(in crate::db) fn memory_bytes(&self) -> u64 {
814 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 #[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;