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 #[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 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 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 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 pub(in crate::db) fn get(&self, key: &RawDataStoreKey) -> Option<RawRow> {
542 self.read(key).into_owned()
543 }
544
545 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 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 #[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 #[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 #[must_use]
620 pub(in crate::db) const fn generation(&self) -> u64 {
621 self.generation
622 }
623
624 #[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 #[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 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 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 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 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 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 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 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 pub(in crate::db) fn memory_bytes(&self) -> u64 {
831 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 #[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;