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 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 #[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 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 #[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 #[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 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 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 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 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 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 pub(in crate::db) fn get(&self, key: &RawDataStoreKey) -> Option<RawRow> {
547 self.read(key).into_owned()
548 }
549
550 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 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 #[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 #[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 #[must_use]
625 pub(in crate::db) const fn generation(&self) -> u64 {
626 self.generation
627 }
628
629 #[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 #[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 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 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 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 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 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 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 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 pub(in crate::db) fn memory_bytes(&self) -> u64 {
830 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 #[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;