1use crate::db::{
7 commit::database_incarnation_id,
8 data::DataStore,
9 index::{IndexId, IndexKeyKind, IndexState, IndexStore, UserIndexPrefixCardinalityKey},
10 integrity::DatabaseIncarnationId,
11 journal::{FoldWatermark, JournalTailStore},
12 schema::{
13 SchemaStore,
14 cardinality_build::CardinalityBuildAuthority,
15 cardinality_generation::{
16 CardinalityAcceptedRootIdentity, CardinalityCountDigest, CardinalityGenerationState,
17 CardinalityStoreAllocationIdentity,
18 },
19 },
20};
21use crate::{error::InternalError, types::EntityTag};
22use candid::CandidType;
23use serde::Deserialize;
24use sha2::{Digest, Sha256};
25#[cfg(all(test, feature = "sql", feature = "diagnostics"))]
26use std::cell::Cell;
27use std::{cell::RefCell, thread::LocalKey};
28
29#[cfg(all(test, feature = "sql", feature = "diagnostics"))]
30thread_local! {
31 static EXACT_PREFIX_EVIDENCE_PROBE_CALLS: Cell<u64> = const { Cell::new(0) };
32 static EXACT_PREFIX_EVIDENCE_LIFECYCLE_READS: Cell<u64> = const { Cell::new(0) };
33}
34
35#[cfg(all(test, feature = "sql", feature = "diagnostics"))]
36pub(in crate::db) fn reset_exact_prefix_evidence_call_counts_for_tests() {
37 EXACT_PREFIX_EVIDENCE_PROBE_CALLS.with(|count| count.set(0));
38 EXACT_PREFIX_EVIDENCE_LIFECYCLE_READS.with(|count| count.set(0));
39}
40
41#[cfg(all(test, feature = "sql", feature = "diagnostics"))]
42pub(in crate::db) fn exact_prefix_evidence_call_counts_for_tests() -> (u64, u64) {
43 (
44 EXACT_PREFIX_EVIDENCE_PROBE_CALLS.with(Cell::get),
45 EXACT_PREFIX_EVIDENCE_LIFECYCLE_READS.with(Cell::get),
46 )
47}
48
49#[derive(Clone, Copy, Debug)]
59pub struct StoreHandle {
60 data: &'static LocalKey<RefCell<DataStore>>,
61 index: &'static LocalKey<RefCell<IndexStore>>,
62 schema: &'static LocalKey<RefCell<SchemaStore>>,
63 journal: Option<&'static LocalKey<RefCell<JournalTailStore>>>,
64 allocations: StoreAllocationIdentities,
65 cardinality_allocation: Option<CardinalityStoreAllocationIdentity>,
66 capabilities: StoreRuntimeStorageCapabilities,
67}
68
69enum ReadyCardinalityCountTargets<'a> {
70 Digests(&'a [CardinalityCountDigest]),
71 UserIndexPrefixes(&'a [UserIndexPrefixCardinalityKey]),
72}
73
74enum ReadyCardinalitySource {
75 Current {
76 database_incarnation: DatabaseIncarnationId,
77 },
78 Admitted {
79 database_incarnation: DatabaseIncarnationId,
80 accepted_root: CardinalityAcceptedRootIdentity,
81 fold_watermark: FoldWatermark,
82 },
83}
84
85#[derive(Clone, Copy, Debug, Eq, PartialEq)]
90pub(in crate::db) struct ExactPrefixCardinalityLifecycleStamp(
91 ExactPrefixCardinalityLifecycleIdentity,
92);
93
94#[derive(Clone, Copy, Debug, Eq, PartialEq)]
95enum ExactPrefixCardinalityLifecycleIdentity {
96 Volatile,
97 MissingDurableAuthority,
98 Corrupt,
99 Journaled {
100 header_digest: Option<[u8; 32]>,
101 cursor_present: bool,
102 delta_watermark: Option<FoldWatermark>,
103 },
104}
105
106#[derive(Clone, Debug, Eq, PartialEq)]
108pub(in crate::db) enum ExactUserIndexPrefixEvidence {
109 Exact(Vec<u64>),
110 Unavailable(ExactPrefixCardinalityLifecycleStamp),
111}
112
113impl ReadyCardinalityCountTargets<'_> {
114 const fn len(&self) -> usize {
115 match self {
116 Self::Digests(digests) => digests.len(),
117 Self::UserIndexPrefixes(keys) => keys.len(),
118 }
119 }
120}
121
122#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
126pub enum StoreRuntimeStorageMode {
127 #[default]
129 Heap,
130 Journaled,
132}
133
134impl StoreRuntimeStorageMode {
135 #[must_use]
137 pub const fn as_str(self) -> &'static str {
138 match self {
139 Self::Heap => "heap",
140 Self::Journaled => "journaled",
141 }
142 }
143}
144
145#[derive(CandidType, Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
147pub enum StoreAllocationIdentityCapability {
148 #[default]
150 Present,
151 Absent,
153}
154
155impl StoreAllocationIdentityCapability {
156 #[must_use]
158 pub const fn as_str(self) -> &'static str {
159 match self {
160 Self::Present => "present",
161 Self::Absent => "absent",
162 }
163 }
164}
165
166#[derive(CandidType, Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
168pub enum StoreDurability {
169 #[default]
171 Durable,
172 Volatile,
174}
175
176impl StoreDurability {
177 #[must_use]
179 pub const fn as_str(self) -> &'static str {
180 match self {
181 Self::Durable => "durable",
182 Self::Volatile => "volatile",
183 }
184 }
185}
186
187#[derive(CandidType, Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
189pub enum StoreRecoveryCapability {
190 #[default]
193 StableBasePlusJournalReplay,
194 None,
196}
197
198impl StoreRecoveryCapability {
199 #[must_use]
201 pub const fn as_str(self) -> &'static str {
202 match self {
203 Self::StableBasePlusJournalReplay => "stable-base-plus-journal-replay",
204 Self::None => "none",
205 }
206 }
207}
208
209#[derive(CandidType, Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
211pub enum StoreCommitParticipation {
212 #[default]
214 Durable,
215 LiveOnly,
217}
218
219impl StoreCommitParticipation {
220 #[must_use]
222 pub const fn as_str(self) -> &'static str {
223 match self {
224 Self::Durable => "durable",
225 Self::LiveOnly => "live-only",
226 }
227 }
228}
229
230#[derive(CandidType, Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
232pub enum StoreSchemaMetadataCapability {
233 LiveRebuiltMetadata,
236 #[default]
238 CanonicalStableHistoryPlusJournalTail,
239}
240
241impl StoreSchemaMetadataCapability {
242 #[must_use]
244 pub const fn as_str(self) -> &'static str {
245 match self {
246 Self::LiveRebuiltMetadata => "live-rebuilt-metadata",
247 Self::CanonicalStableHistoryPlusJournalTail => {
248 "canonical-stable-history-plus-journal-tail"
249 }
250 }
251 }
252}
253
254#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
256pub enum StoreRelationSourceCapability {
257 #[default]
259 DurableSource,
260 LiveSource,
262}
263
264#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
266pub enum StoreRelationTargetCapability {
267 #[default]
269 DurableTarget,
270 VolatileTarget,
272}
273
274#[derive(Clone, Copy, Debug, Default, Deserialize, Eq, PartialEq)]
278pub struct StoreRuntimeStorageCapabilities {
279 storage_mode: StoreRuntimeStorageMode,
280 allocation_identity: StoreAllocationIdentityCapability,
281 durability: StoreDurability,
282 recovery: StoreRecoveryCapability,
283 commit_participation: StoreCommitParticipation,
284 schema_metadata: StoreSchemaMetadataCapability,
285 relation_source: StoreRelationSourceCapability,
286 relation_target: StoreRelationTargetCapability,
287}
288
289impl StoreRuntimeStorageCapabilities {
290 #[must_use]
292 pub const fn heap() -> Self {
293 Self {
294 storage_mode: StoreRuntimeStorageMode::Heap,
295 allocation_identity: StoreAllocationIdentityCapability::Absent,
296 durability: StoreDurability::Volatile,
297 recovery: StoreRecoveryCapability::None,
298 commit_participation: StoreCommitParticipation::LiveOnly,
299 schema_metadata: StoreSchemaMetadataCapability::LiveRebuiltMetadata,
300 relation_source: StoreRelationSourceCapability::LiveSource,
301 relation_target: StoreRelationTargetCapability::VolatileTarget,
302 }
303 }
304
305 #[must_use]
307 pub const fn journaled() -> Self {
308 Self {
309 storage_mode: StoreRuntimeStorageMode::Journaled,
310 allocation_identity: StoreAllocationIdentityCapability::Present,
311 durability: StoreDurability::Durable,
312 recovery: StoreRecoveryCapability::StableBasePlusJournalReplay,
313 commit_participation: StoreCommitParticipation::Durable,
314 schema_metadata: StoreSchemaMetadataCapability::CanonicalStableHistoryPlusJournalTail,
315 relation_source: StoreRelationSourceCapability::DurableSource,
316 relation_target: StoreRelationTargetCapability::DurableTarget,
317 }
318 }
319
320 #[must_use]
322 pub const fn storage_mode(self) -> StoreRuntimeStorageMode {
323 self.storage_mode
324 }
325
326 #[must_use]
328 pub const fn allocation_identity(self) -> StoreAllocationIdentityCapability {
329 self.allocation_identity
330 }
331
332 #[must_use]
334 pub const fn durability(self) -> StoreDurability {
335 self.durability
336 }
337
338 #[must_use]
340 pub const fn recovery(self) -> StoreRecoveryCapability {
341 self.recovery
342 }
343
344 #[must_use]
346 pub const fn commit_participation(self) -> StoreCommitParticipation {
347 self.commit_participation
348 }
349
350 #[must_use]
352 pub const fn schema_metadata(self) -> StoreSchemaMetadataCapability {
353 self.schema_metadata
354 }
355
356 #[must_use]
358 pub const fn relation_source(self) -> StoreRelationSourceCapability {
359 self.relation_source
360 }
361
362 #[must_use]
364 pub const fn relation_target(self) -> StoreRelationTargetCapability {
365 self.relation_target
366 }
367}
368
369#[derive(Clone, Copy, Debug, Eq, PartialEq)]
376pub struct StoreAllocationIdentity {
377 memory_id: u8,
378 stable_key: &'static str,
379}
380
381impl StoreAllocationIdentity {
382 #[must_use]
384 pub const fn new(memory_id: u8, stable_key: &'static str) -> Self {
385 Self {
386 memory_id,
387 stable_key,
388 }
389 }
390
391 #[must_use]
393 pub const fn memory_id(self) -> u8 {
394 self.memory_id
395 }
396
397 #[must_use]
399 pub const fn stable_key(self) -> &'static str {
400 self.stable_key
401 }
402}
403
404#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
412pub struct StoreAllocationIdentities {
413 data: Option<StoreAllocationIdentity>,
414 index: Option<StoreAllocationIdentity>,
415 schema: Option<StoreAllocationIdentity>,
416 journal: Option<StoreAllocationIdentity>,
417}
418
419impl StoreAllocationIdentities {
420 #[must_use]
422 pub const fn absent() -> Self {
423 Self {
424 data: None,
425 index: None,
426 schema: None,
427 journal: None,
428 }
429 }
430
431 #[must_use]
433 pub const fn new_journaled(
434 data: StoreAllocationIdentity,
435 index: StoreAllocationIdentity,
436 schema: StoreAllocationIdentity,
437 journal: StoreAllocationIdentity,
438 ) -> Self {
439 Self {
440 data: Some(data),
441 index: Some(index),
442 schema: Some(schema),
443 journal: Some(journal),
444 }
445 }
446
447 #[must_use]
449 pub const fn data(self) -> Option<StoreAllocationIdentity> {
450 self.data
451 }
452
453 #[must_use]
455 pub const fn index(self) -> Option<StoreAllocationIdentity> {
456 self.index
457 }
458
459 #[must_use]
461 pub const fn schema(self) -> Option<StoreAllocationIdentity> {
462 self.schema
463 }
464
465 #[must_use]
467 pub const fn journal(self) -> Option<StoreAllocationIdentity> {
468 self.journal
469 }
470
471 #[must_use]
474 pub const fn allocation_identity_capability(self) -> Option<StoreAllocationIdentityCapability> {
475 match (self.data, self.index, self.schema) {
476 (Some(_), Some(_), Some(_)) => Some(StoreAllocationIdentityCapability::Present),
477 (None, None, None) if self.journal.is_none() => {
478 Some(StoreAllocationIdentityCapability::Absent)
479 }
480 _ => None,
481 }
482 }
483
484 #[must_use]
487 pub const fn matches_storage_capabilities(
488 self,
489 capabilities: StoreRuntimeStorageCapabilities,
490 ) -> bool {
491 match capabilities.storage_mode() {
492 StoreRuntimeStorageMode::Heap => {
493 self.data.is_none()
494 && self.index.is_none()
495 && self.schema.is_none()
496 && self.journal.is_none()
497 }
498 StoreRuntimeStorageMode::Journaled => {
499 self.data.is_some()
500 && self.index.is_some()
501 && self.schema.is_some()
502 && self.journal.is_some()
503 }
504 }
505 }
506}
507
508impl StoreHandle {
509 #[must_use]
511 pub const fn new(
512 data: &'static LocalKey<RefCell<DataStore>>,
513 index: &'static LocalKey<RefCell<IndexStore>>,
514 schema: &'static LocalKey<RefCell<SchemaStore>>,
515 allocations: StoreAllocationIdentities,
516 capabilities: StoreRuntimeStorageCapabilities,
517 ) -> Self {
518 Self {
519 data,
520 index,
521 schema,
522 journal: None,
523 allocations,
524 cardinality_allocation: None,
525 capabilities,
526 }
527 }
528
529 #[must_use]
531 pub fn new_journaled(
532 data: &'static LocalKey<RefCell<DataStore>>,
533 index: &'static LocalKey<RefCell<IndexStore>>,
534 schema: &'static LocalKey<RefCell<SchemaStore>>,
535 journal: &'static LocalKey<RefCell<JournalTailStore>>,
536 allocations: StoreAllocationIdentities,
537 capabilities: StoreRuntimeStorageCapabilities,
538 ) -> Self {
539 let cardinality_allocation = CardinalityStoreAllocationIdentity::derive(allocations).ok();
540 Self {
541 data,
542 index,
543 schema,
544 journal: Some(journal),
545 allocations,
546 cardinality_allocation,
547 capabilities,
548 }
549 }
550
551 pub fn with_data<R>(&self, f: impl FnOnce(&DataStore) -> R) -> R {
553 #[cfg(feature = "diagnostics")]
554 {
555 crate::db::physical_access::measure_physical_access_operation(|| {
556 self.data.with_borrow(f)
557 })
558 }
559
560 #[cfg(not(feature = "diagnostics"))]
561 {
562 self.data.with_borrow(f)
563 }
564 }
565
566 pub fn with_data_mut<R>(&self, f: impl FnOnce(&mut DataStore) -> R) -> R {
568 self.data.with_borrow_mut(f)
569 }
570
571 pub fn with_index<R>(&self, f: impl FnOnce(&IndexStore) -> R) -> R {
573 #[cfg(feature = "diagnostics")]
574 {
575 crate::db::physical_access::measure_physical_access_operation(|| {
576 self.index.with_borrow(f)
577 })
578 }
579
580 #[cfg(not(feature = "diagnostics"))]
581 {
582 self.index.with_borrow(f)
583 }
584 }
585
586 pub fn with_index_mut<R>(&self, f: impl FnOnce(&mut IndexStore) -> R) -> R {
588 self.index.with_borrow_mut(f)
589 }
590
591 pub fn with_schema<R>(&self, f: impl FnOnce(&SchemaStore) -> R) -> R {
593 self.schema.with_borrow(f)
594 }
595
596 pub fn with_schema_mut<R>(&self, f: impl FnOnce(&mut SchemaStore) -> R) -> R {
598 self.schema.with_borrow_mut(f)
599 }
600
601 #[must_use]
603 pub(in crate::db) fn exact_entity_count(&self, entity: EntityTag) -> Option<u64> {
604 if self.journal.is_none() {
605 return self.with_data(|store| store.exact_entity_count(entity));
606 }
607 let delta = self.with_data(|store| store.exact_entity_cardinality_delta(entity))?;
608 let digest = CardinalityCountDigest::for_entity(entity);
609 let base = self
610 .ready_cardinality_counts(&[digest], |authority| authority.accepts_entity(entity))
611 .ok()
612 .flatten()?
613 .into_iter()
614 .next()?;
615 apply_visible_cardinality_delta(base, delta)
616 }
617
618 #[must_use]
620 pub(in crate::db) fn exact_user_index_prefix_count(
621 &self,
622 data_generation: u64,
623 key_kind: IndexKeyKind,
624 index_id: IndexId,
625 components: &[Vec<u8>],
626 ) -> Option<u64> {
627 self.exact_user_index_prefix_counts(data_generation, key_kind, index_id, [components])?
628 .into_iter()
629 .next()
630 }
631
632 #[must_use]
634 pub(in crate::db) fn exact_user_index_prefix_counts<'a>(
635 &self,
636 data_generation: u64,
637 key_kind: IndexKeyKind,
638 index_id: IndexId,
639 component_prefixes: impl IntoIterator<Item = &'a [Vec<u8>]>,
640 ) -> Option<Vec<u64>> {
641 let component_prefixes = component_prefixes.into_iter().collect::<Vec<_>>();
642 if key_kind != IndexKeyKind::User {
643 return None;
644 }
645 if self.journal.is_none() {
646 return self.with_index(|store| {
647 component_prefixes
648 .iter()
649 .map(|components| {
650 store.exact_prefix_cardinality(
651 data_generation,
652 key_kind,
653 index_id,
654 components,
655 )
656 })
657 .collect()
658 });
659 }
660 let deltas = self.with_index(|store| {
661 component_prefixes
662 .iter()
663 .map(|components| {
664 store.exact_prefix_cardinality_delta(key_kind, index_id, components)
665 })
666 .collect::<Option<Vec<_>>>()
667 })?;
668 let digests = component_prefixes
669 .iter()
670 .map(|components| {
671 CardinalityCountDigest::for_user_index_prefix(index_id, components).ok()
672 })
673 .collect::<Option<Vec<_>>>()?;
674 let bases = self
675 .ready_cardinality_counts(&digests, |authority| {
676 component_prefixes.iter().all(|components| {
677 authority.accepts_user_index_prefix(index_id, components.len())
678 })
679 })
680 .ok()
681 .flatten()?;
682 bases
683 .into_iter()
684 .zip(deltas)
685 .map(|(base, delta)| apply_visible_cardinality_delta(base, delta))
686 .collect()
687 }
688
689 #[must_use]
691 pub(in crate::db) fn exact_user_index_prefix_key_counts(
692 &self,
693 data_generation: u64,
694 keys: &[UserIndexPrefixCardinalityKey],
695 ) -> Option<Vec<u64>> {
696 self.exact_user_index_prefix_key_counts_with_authority(data_generation, keys, None)
697 }
698
699 #[must_use]
701 pub(in crate::db) fn exact_user_index_prefix_key_counts_for_admitted_root(
702 &self,
703 database_incarnation: DatabaseIncarnationId,
704 accepted_root: CardinalityAcceptedRootIdentity,
705 data_generation: u64,
706 keys: &[UserIndexPrefixCardinalityKey],
707 ) -> Option<Vec<u64>> {
708 self.exact_user_index_prefix_key_counts_with_authority(
709 data_generation,
710 keys,
711 Some((database_incarnation, accepted_root)),
712 )
713 }
714
715 #[must_use]
720 pub(in crate::db) fn exact_user_index_prefix_evidence_for_admitted_root(
721 &self,
722 database_incarnation: DatabaseIncarnationId,
723 accepted_root: CardinalityAcceptedRootIdentity,
724 keys: &[UserIndexPrefixCardinalityKey],
725 ) -> ExactUserIndexPrefixEvidence {
726 #[cfg(all(test, feature = "sql", feature = "diagnostics"))]
727 EXACT_PREFIX_EVIDENCE_PROBE_CALLS.with(|count| count.set(count.get().saturating_add(1)));
728 let data_generation = self.with_data(DataStore::generation);
729 if let Some(counts) = self.exact_user_index_prefix_key_counts_for_admitted_root(
730 database_incarnation,
731 accepted_root,
732 data_generation,
733 keys,
734 ) {
735 return ExactUserIndexPrefixEvidence::Exact(counts);
736 }
737
738 ExactUserIndexPrefixEvidence::Unavailable(
739 self.exact_user_index_prefix_evidence_lifecycle_stamp(),
740 )
741 }
742
743 #[must_use]
745 pub(in crate::db) fn exact_user_index_prefix_evidence_lifecycle_stamp(
746 &self,
747 ) -> ExactPrefixCardinalityLifecycleStamp {
748 #[cfg(all(test, feature = "sql", feature = "diagnostics"))]
749 EXACT_PREFIX_EVIDENCE_LIFECYCLE_READS
750 .with(|count| count.set(count.get().saturating_add(1)));
751 if self.journal.is_none() {
752 return ExactPrefixCardinalityLifecycleStamp(
753 ExactPrefixCardinalityLifecycleIdentity::Volatile,
754 );
755 }
756 if self.cardinality_allocation.is_none() {
757 return ExactPrefixCardinalityLifecycleStamp(
758 ExactPrefixCardinalityLifecycleIdentity::MissingDurableAuthority,
759 );
760 }
761 let delta_watermark = self.with_index(IndexStore::exact_prefix_cardinality_delta_watermark);
762 match self.with_schema(SchemaStore::cardinality_generation_lifecycle_control) {
763 Ok((header, cursor_present)) => ExactPrefixCardinalityLifecycleStamp(
764 ExactPrefixCardinalityLifecycleIdentity::Journaled {
765 header_digest: header.map(|header| {
766 let mut hasher = Sha256::new();
767 hasher.update(b"icydb.cardinality-lifecycle-stamp.v1");
768 hasher.update(header.encode());
769 hasher.finalize().into()
770 }),
771 cursor_present,
772 delta_watermark,
773 },
774 ),
775 Err(_) => ExactPrefixCardinalityLifecycleStamp(
776 ExactPrefixCardinalityLifecycleIdentity::Corrupt,
777 ),
778 }
779 }
780
781 fn exact_user_index_prefix_key_counts_with_authority(
782 &self,
783 data_generation: u64,
784 keys: &[UserIndexPrefixCardinalityKey],
785 admitted: Option<(DatabaseIncarnationId, CardinalityAcceptedRootIdentity)>,
786 ) -> Option<Vec<u64>> {
787 if keys.is_empty() {
788 return None;
789 }
790 if self.journal.is_none() {
791 return self.with_index(|store| {
792 keys.iter()
793 .map(|key| {
794 store.exact_prefix_cardinality(
795 data_generation,
796 IndexKeyKind::User,
797 key.index_id(),
798 key.prefix_components(),
799 )
800 })
801 .collect()
802 });
803 }
804 let (delta_watermark, deltas) = self.with_index(|store| {
805 let watermark = store.exact_prefix_cardinality_delta_watermark()?;
806 keys.iter()
807 .map(|key| {
808 store.exact_prefix_cardinality_delta(
809 IndexKeyKind::User,
810 key.index_id(),
811 key.prefix_components(),
812 )
813 })
814 .collect::<Option<Vec<_>>>()
815 .map(|deltas| (watermark, deltas))
816 })?;
817 let accepts = |authority: &CardinalityBuildAuthority| {
818 keys.iter().all(|key| {
819 authority.accepts_user_index_prefix(key.index_id(), key.prefix_components().len())
820 })
821 };
822 let bases = match admitted {
823 Some((database_incarnation, accepted_root)) => self
824 .ready_cardinality_counts_for_source(
825 ReadyCardinalitySource::Admitted {
826 database_incarnation,
827 accepted_root,
828 fold_watermark: delta_watermark,
829 },
830 ReadyCardinalityCountTargets::UserIndexPrefixes(keys),
831 accepts,
832 ),
833 None => self.ready_cardinality_counts_for_targets(
834 ReadyCardinalityCountTargets::UserIndexPrefixes(keys),
835 accepts,
836 ),
837 }
838 .ok()
839 .flatten()?;
840 bases
841 .into_iter()
842 .zip(deltas)
843 .map(|(base, delta)| apply_visible_cardinality_delta(base, delta))
844 .collect()
845 }
846
847 #[must_use]
853 pub(in crate::db) fn user_index_prefix_family_has_ready_generation<'a, I>(
854 &self,
855 database_incarnation: DatabaseIncarnationId,
856 accepted_root: CardinalityAcceptedRootIdentity,
857 data_generation: u64,
858 key_kind: IndexKeyKind,
859 index_id: IndexId,
860 component_prefixes: I,
861 ) -> bool
862 where
863 I: Clone + IntoIterator<Item = &'a [Vec<u8>]>,
864 {
865 if key_kind != IndexKeyKind::User || component_prefixes.clone().into_iter().next().is_none()
866 {
867 return false;
868 }
869 if self.journal.is_none() {
870 return self.with_index(|store| {
871 component_prefixes.clone().into_iter().all(|components| {
872 store
873 .exact_prefix_cardinality(data_generation, key_kind, index_id, components)
874 .is_some()
875 })
876 });
877 }
878 let delta_watermark = self.with_index(|store| {
879 let watermark = store.exact_prefix_cardinality_delta_watermark()?;
880 component_prefixes
881 .clone()
882 .into_iter()
883 .all(|components| {
884 store
885 .exact_prefix_cardinality_delta(key_kind, index_id, components)
886 .is_some()
887 })
888 .then_some(watermark)
889 });
890 delta_watermark.is_some_and(|watermark| {
891 self.ready_cardinality_counts_for_source(
892 ReadyCardinalitySource::Admitted {
893 database_incarnation,
894 accepted_root,
895 fold_watermark: watermark,
896 },
897 ReadyCardinalityCountTargets::Digests(&[]),
898 |authority| {
899 component_prefixes.into_iter().all(|components| {
900 authority.accepts_user_index_prefix(index_id, components.len())
901 })
902 },
903 )
904 .is_ok_and(|counts| counts.is_some())
905 })
906 }
907
908 #[must_use]
910 pub(in crate::db) fn exact_user_index_prefix_count_sum<'a>(
911 &self,
912 data_generation: u64,
913 key_kind: IndexKeyKind,
914 index_id: IndexId,
915 component_prefixes: impl IntoIterator<Item = &'a [Vec<u8>]>,
916 stop_after: Option<u64>,
917 ) -> Option<u64> {
918 let component_prefixes = component_prefixes.into_iter().collect::<Vec<_>>();
919 if self.journal.is_none() {
920 return self.with_index(|store| {
921 store.exact_prefix_cardinality_sum(
922 data_generation,
923 key_kind,
924 index_id,
925 component_prefixes.iter().copied(),
926 stop_after,
927 )
928 });
929 }
930 let counts = self.exact_user_index_prefix_counts(
931 data_generation,
932 key_kind,
933 index_id,
934 component_prefixes.iter().copied(),
935 )?;
936 let mut total = 0_u64;
937 for count in counts {
938 total = total.checked_add(count)?;
939 if stop_after.is_some_and(|required| total >= required) {
940 break;
941 }
942 }
943 Some(total)
944 }
945
946 #[must_use]
948 pub(in crate::db) fn exact_user_index_child_prefixes_for_parent_set<'a>(
949 &self,
950 data_generation: u64,
951 index_id: IndexId,
952 parent_prefixes: impl IntoIterator<Item = &'a [Vec<u8>]>,
953 total_cap: usize,
954 ) -> Option<Vec<Vec<Vec<u8>>>> {
955 let mut parent_prefixes = parent_prefixes
956 .into_iter()
957 .map(<[Vec<u8>]>::to_vec)
958 .collect::<Vec<_>>();
959 if parent_prefixes.iter().any(Vec::is_empty) {
960 return None;
961 }
962 icydb_schema::compact_sort_unstable_by(&mut parent_prefixes, Ord::cmp);
963 parent_prefixes.dedup();
964 let child_prefixes = self.with_index(|store| {
965 store.exact_child_prefixes_for_parent_set(
966 data_generation,
967 IndexKeyKind::User,
968 index_id,
969 parent_prefixes.iter().map(Vec::as_slice),
970 total_cap,
971 )
972 })?;
973 if self.journal.is_none() {
974 return Some(child_prefixes);
975 }
976 let parent_count = parent_prefixes.len();
977 let counts = self.exact_user_index_prefix_counts(
978 data_generation,
979 IndexKeyKind::User,
980 index_id,
981 parent_prefixes
982 .iter()
983 .chain(&child_prefixes)
984 .map(Vec::as_slice),
985 )?;
986 let (parent_counts, child_counts) = counts.split_at(parent_count);
987 let parent_total = checked_cardinality_sum(parent_counts)?;
988 let child_total = checked_cardinality_sum(child_counts)?;
989 (parent_total == child_total).then_some(child_prefixes)
990 }
991
992 fn ready_cardinality_counts(
993 &self,
994 digests: &[CardinalityCountDigest],
995 accepts: impl FnOnce(&CardinalityBuildAuthority) -> bool,
996 ) -> Result<Option<Vec<u64>>, InternalError> {
997 self.ready_cardinality_counts_for_targets(
998 ReadyCardinalityCountTargets::Digests(digests),
999 accepts,
1000 )
1001 }
1002
1003 fn ready_cardinality_counts_for_targets(
1004 &self,
1005 targets: ReadyCardinalityCountTargets<'_>,
1006 accepts: impl FnOnce(&CardinalityBuildAuthority) -> bool,
1007 ) -> Result<Option<Vec<u64>>, InternalError> {
1008 let incarnation = database_incarnation_id()?;
1009 self.ready_cardinality_counts_for_source(
1010 ReadyCardinalitySource::Current {
1011 database_incarnation: incarnation,
1012 },
1013 targets,
1014 accepts,
1015 )
1016 }
1017
1018 fn ready_cardinality_counts_for_source(
1019 &self,
1020 source: ReadyCardinalitySource,
1021 targets: ReadyCardinalityCountTargets<'_>,
1022 accepts: impl FnOnce(&CardinalityBuildAuthority) -> bool,
1023 ) -> Result<Option<Vec<u64>>, InternalError> {
1024 let Some(journal) = self.journal else {
1025 return Ok(None);
1026 };
1027 let Some(allocation) = self.cardinality_allocation else {
1028 return Ok(None);
1029 };
1030 let (incarnation, accepted_root, watermark) = match source {
1031 ReadyCardinalitySource::Current {
1032 database_incarnation,
1033 } => (
1034 database_incarnation,
1035 None,
1036 journal.with_borrow(JournalTailStore::fold_watermark)?,
1037 ),
1038 ReadyCardinalitySource::Admitted {
1039 database_incarnation,
1040 accepted_root,
1041 fold_watermark,
1042 } => (database_incarnation, Some(accepted_root), fold_watermark),
1043 };
1044 self.with_schema(|schema| {
1045 let (header, cursor) = schema.cardinality_generation_control()?;
1046 let Some(header) = header else {
1047 return Ok(None);
1048 };
1049 if header.state() != CardinalityGenerationState::Ready || cursor.is_some() {
1050 return Ok(None);
1051 }
1052 let authority = match accepted_root {
1053 Some(root) => CardinalityBuildAuthority::derive_for_admitted_consumer_root(
1054 schema,
1055 incarnation,
1056 allocation,
1057 root,
1058 watermark,
1059 )?,
1060 None => CardinalityBuildAuthority::derive_for_current_consumer(
1061 schema,
1062 incarnation,
1063 allocation,
1064 watermark,
1065 )?,
1066 };
1067 let Some(authority) = authority else {
1068 return Ok(None);
1069 };
1070 if !accepts(&authority) {
1071 return Ok(None);
1072 }
1073 if header.validate_source(authority.source()).is_err() {
1074 return Ok(None);
1075 }
1076 if targets.len() != 0 && schema.cardinality_count_slot_is_empty(header.slot())? {
1077 return Ok(Some(vec![0; targets.len()]));
1078 }
1079 let counts = match targets {
1080 ReadyCardinalityCountTargets::Digests(digests) => digests
1081 .iter()
1082 .map(|digest| {
1083 schema
1084 .cardinality_count(header.slot(), header.generation(), *digest)
1085 .map(|count| count.unwrap_or(0))
1086 })
1087 .collect::<Result<Vec<_>, _>>()?,
1088 ReadyCardinalityCountTargets::UserIndexPrefixes(keys) => keys
1089 .iter()
1090 .map(|key| {
1091 let digest = CardinalityCountDigest::for_user_index_prefix(
1092 key.index_id(),
1093 key.prefix_components(),
1094 )?;
1095 schema
1096 .cardinality_count(header.slot(), header.generation(), digest)
1097 .map(|count| count.unwrap_or(0))
1098 })
1099 .collect::<Result<Vec<_>, _>>()?,
1100 };
1101 Ok(Some(counts))
1102 })
1103 }
1104
1105 #[must_use]
1107 pub(in crate::db) fn index_state(&self) -> IndexState {
1108 self.with_index(IndexStore::state)
1109 }
1110
1111 pub(in crate::db) fn access_state_revision(&self) -> Result<u64, crate::error::InternalError> {
1113 self.journal.map_or_else(
1114 || Ok(self.with_index(IndexStore::access_state_revision)),
1115 |journal| journal.with_borrow(JournalTailStore::access_state_revision),
1116 )
1117 }
1118
1119 pub(in crate::db) fn mark_index_building(&self) -> Result<(), crate::error::InternalError> {
1121 self.set_index_state(IndexState::Building)
1122 }
1123
1124 pub(in crate::db) fn mark_index_ready(&self) -> Result<(), crate::error::InternalError> {
1126 self.set_index_state(IndexState::Ready)
1127 }
1128
1129 fn set_index_state(&self, state: IndexState) -> Result<(), crate::error::InternalError> {
1130 if self.index_state() == state {
1131 return Ok(());
1132 }
1133 let revision = self.journal.map_or_else(
1134 || {
1135 self.with_index(IndexStore::access_state_revision)
1136 .checked_add(1)
1137 .ok_or_else(crate::error::InternalError::store_invariant)
1138 },
1139 |journal| journal.with_borrow_mut(JournalTailStore::advance_access_state_revision),
1140 )?;
1141 self.with_index_mut(|index| index.set_access_state(state, revision));
1142 Ok(())
1143 }
1144
1145 #[must_use]
1147 pub const fn data_store(&self) -> &'static LocalKey<RefCell<DataStore>> {
1148 self.data
1149 }
1150
1151 #[must_use]
1153 pub const fn index_store(&self) -> &'static LocalKey<RefCell<IndexStore>> {
1154 self.index
1155 }
1156
1157 #[must_use]
1159 pub const fn schema_store(&self) -> &'static LocalKey<RefCell<SchemaStore>> {
1160 self.schema
1161 }
1162
1163 #[must_use]
1165 pub const fn journal_tail_store(&self) -> Option<&'static LocalKey<RefCell<JournalTailStore>>> {
1166 self.journal
1167 }
1168
1169 #[must_use]
1172 pub const fn data_allocation(&self) -> Option<StoreAllocationIdentity> {
1173 self.allocations.data()
1174 }
1175
1176 #[must_use]
1179 pub const fn index_allocation(&self) -> Option<StoreAllocationIdentity> {
1180 self.allocations.index()
1181 }
1182
1183 #[must_use]
1186 pub const fn schema_allocation(&self) -> Option<StoreAllocationIdentity> {
1187 self.allocations.schema()
1188 }
1189
1190 #[must_use]
1193 pub const fn journal_allocation(&self) -> Option<StoreAllocationIdentity> {
1194 self.allocations.journal()
1195 }
1196
1197 #[must_use]
1199 pub(in crate::db) const fn allocation_identities(&self) -> StoreAllocationIdentities {
1200 self.allocations
1201 }
1202
1203 #[must_use]
1205 pub const fn storage_capabilities(&self) -> StoreRuntimeStorageCapabilities {
1206 self.capabilities
1207 }
1208}
1209
1210fn apply_visible_cardinality_delta(base: u64, delta: i64) -> Option<u64> {
1211 if delta >= 0 {
1212 base.checked_add(u64::try_from(delta).ok()?)
1213 } else {
1214 base.checked_sub(delta.unsigned_abs())
1215 }
1216}
1217
1218fn checked_cardinality_sum(counts: &[u64]) -> Option<u64> {
1219 counts
1220 .iter()
1221 .try_fold(0_u64, |total, count| total.checked_add(*count))
1222}