1mod deep;
7mod derived;
8mod job;
9mod progress_codec;
10mod progress_store;
11mod proof;
12mod row;
13
14use crate::{
15 db::{
16 commit::ensure_recovery_admitted,
17 registry::{
18 StoreAllocationIdentities, StoreHandle, StoreRuntimeStorageCapabilities,
19 StoreRuntimeStorageMode,
20 },
21 schema::{
22 AcceptedInspectionPlan, IdentityStateLifecycle, MAX_IDENTITY_STATE_RECORDS_PER_DATABASE,
23 },
24 },
25 error::{ConstraintValuePath, ErrorClass, ErrorOrigin, InternalError},
26 traits::CanisterKind,
27};
28use candid::CandidType;
29use serde::Deserialize;
30use std::{
31 collections::BTreeMap,
32 sync::atomic::{AtomicU64, Ordering},
33};
34
35pub(in crate::db) use deep::{
36 abort_deep_integrity_job, continue_deep_integrity_job, run_next_integrity_retention_page,
37 start_deep_integrity_job,
38};
39pub(in crate::db) use derived::{
40 DerivedInspectionLimits, execute_index_integrity_page, execute_reverse_integrity_page,
41};
42pub use job::{
43 DeepIntegrityPage, DeepIntegrityPageStatus, IntegrityAbortReceipt, IntegrityAbortStatus,
44 IntegrityDeepError, IntegrityJobError, IntegrityJobId, IntegrityJobOwner, IntegrityJobReceipt,
45 IntegrityPendingTerminal, IntegritySubmissionKey, IntegrityTerminalOutcome,
46};
47pub(in crate::db) use job::{
48 IntegrityCheckpoint, IntegrityJob, IntegrityJobState, IntegrityReceiptEnvelope,
49 IntegrityReceiptReplayKey, MAX_INTEGRITY_IN_PROGRESS_PAGES,
50};
51#[cfg(any(feature = "sql", test))]
52pub(in crate::db) use progress_store::InsertMutationJobResult;
53#[cfg(feature = "sql")]
54pub(in crate::db) use progress_store::replace_mutation_progress_record_op;
55pub(in crate::db) use progress_store::{
56 MutationProgressRecordOp, apply_mutation_progress_record_op,
57 apply_preflighted_mutation_progress_record_op, preflight_mutation_progress_record_op,
58 verify_mutation_progress_record_op, with_mutation_progress_store,
59 with_resumable_progress_store,
60};
61pub use progress_store::{
62 ProgressJobFamily, ProgressJobInventory, ProgressJobInventoryRecord, ProgressJobLifecycle,
63};
64pub(in crate::db) use proof::{IntegrityProofVector, capture_integrity_proof_vector};
65pub(in crate::db) use row::{
66 PhysicalUnitCheckpoint, RowInspectionLimits, execute_row_integrity_page,
67};
68
69pub(in crate::db) const MAX_INTEGRITY_PATH_BYTES: usize = 4 * 1024;
70
71#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
79pub enum IntegrityCheckRequest {
80 Quick {
82 entity: IntegrityEntityIdentity,
84 },
85 DeepStart {
87 entity: IntegrityEntityIdentity,
89 submission_key: IntegritySubmissionKey,
91 },
92 DeepContinue {
94 job_id: IntegrityJobId,
96 acknowledged_sequence: u64,
98 },
99 DeepAbort {
101 job_id: IntegrityJobId,
103 },
104}
105
106impl IntegrityCheckRequest {
107 #[must_use]
109 pub const fn deep_continue(job_id: IntegrityJobId, acknowledged_sequence: u64) -> Self {
110 Self::DeepContinue {
111 job_id,
112 acknowledged_sequence,
113 }
114 }
115
116 #[must_use]
118 pub const fn deep_abort(job_id: IntegrityJobId) -> Self {
119 Self::DeepAbort { job_id }
120 }
121}
122
123#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
126pub enum IntegrityCheckResult {
127 Quick(QuickIntegrityResult),
129 Deep(IntegrityJobReceipt),
131}
132
133fn resolve_integrity_participating_stores<C: CanisterKind>(
135 db: &crate::db::Db<C>,
136 plan: &AcceptedInspectionPlan,
137) -> Result<BTreeMap<String, StoreHandle>, InternalError> {
138 let identity = plan.identity();
139 let source_store = db.store_handle(identity.store_path())?;
140 let mut stores = BTreeMap::from([(identity.store_path().to_string(), source_store)]);
141 for relation in plan.relation_inspection() {
142 stores
143 .entry(relation.target_store_path().to_string())
144 .or_insert_with(|| relation.target_store());
145 }
146 Ok(stores)
147}
148
149fn validate_quick_integrity_control<C: CanisterKind>(
150 db: &crate::db::Db<C>,
151 plan: &AcceptedInspectionPlan,
152 incarnation: DatabaseIncarnationId,
153) -> Result<Vec<IntegrityFinding>, InternalError> {
154 let identity = plan.identity();
155 let participating_stores = resolve_integrity_participating_stores(db, plan)?;
156
157 proof::validate_integrity_allocation_registry()?;
162 validate_quick_identity_control(db, incarnation)?;
163 let mut findings = Vec::new();
164 for (store_path, store) in &participating_stores {
165 if let Some(finding) = validate_quick_store_control(plan, store_path, *store)? {
166 findings.push(finding);
167 }
168 }
169 for ordinal in 0..plan.index_inspection().len() {
170 let _domain = plan
171 .index_inspection()
172 .domain(ordinal, identity.entity_tag())?;
173 }
174
175 Ok(findings)
176}
177
178fn validate_quick_identity_control<C: CanisterKind>(
179 db: &crate::db::Db<C>,
180 incarnation: DatabaseIncarnationId,
181) -> Result<(), InternalError> {
182 let mut stores = db.with_store_registry(|registry| registry.iter().collect::<Vec<_>>());
183 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.0.cmp(right.0));
184
185 let mut owners = BTreeMap::new();
186 let mut state_count = 0usize;
187 for (store_path, store) in stores {
188 let states = store.with_schema(|schema_store| {
189 schema_store.identity_state_inventory_for_integrity(incarnation)
190 })?;
191 state_count = state_count
192 .checked_add(states.len())
193 .ok_or_else(InternalError::identity_state_corruption)?;
194 if state_count > MAX_IDENTITY_STATE_RECORDS_PER_DATABASE {
195 return Err(InternalError::identity_state_corruption());
196 }
197
198 for state in states {
199 let owner = state.owner();
200 record_quick_identity_owner(&mut owners, store_path, &state)?;
201 if state.lifecycle() == IdentityStateLifecycle::Active {
202 let runtime_entity = db
203 .accepted_runtime_entity_for_tag(owner.entity_tag())
204 .map_err(|_| InternalError::identity_state_corruption())?;
205 if runtime_entity.store_path() != store_path {
206 return Err(InternalError::identity_state_corruption());
207 }
208 }
209 }
210 }
211
212 Ok(())
213}
214
215fn record_quick_identity_owner<'a>(
216 owners: &mut BTreeMap<(crate::types::EntityTag, crate::db::schema::FieldId), &'a str>,
217 store_path: &'a str,
218 state: &crate::db::schema::IdentityState,
219) -> Result<(), InternalError> {
220 let owner = state.owner();
221 let key = (owner.entity_tag(), owner.field_id());
222 if owners.insert(key, store_path).is_some() {
223 return Err(InternalError::identity_state_corruption());
224 }
225 Ok(())
226}
227
228fn validate_quick_store_control(
229 plan: &AcceptedInspectionPlan,
230 store_path: &str,
231 store: StoreHandle,
232) -> Result<Option<IntegrityFinding>, InternalError> {
233 let capabilities = store.storage_capabilities();
234 let allocations = store.allocation_identities();
235 match capabilities.storage_mode() {
236 StoreRuntimeStorageMode::Heap => {
237 if capabilities != StoreRuntimeStorageCapabilities::heap()
238 || allocations != StoreAllocationIdentities::absent()
239 || store.journal_tail_store().is_some()
240 {
241 return Err(InternalError::store_invariant());
242 }
243 Ok(None)
244 }
245 StoreRuntimeStorageMode::Journaled => {
246 if capabilities != StoreRuntimeStorageCapabilities::journaled()
247 || !allocations.matches_storage_capabilities(capabilities)
248 {
249 return Err(InternalError::store_invariant());
250 }
251 let journal = store
252 .journal_tail_store()
253 .ok_or_else(InternalError::store_invariant)?
254 .with_borrow(crate::db::journal::JournalTailStore::proof_identity)?;
255 if !journal.is_well_formed() {
256 return Ok(Some(quick_journal_control_finding(plan, store_path)));
257 }
258 Ok(None)
259 }
260 }
261}
262
263fn quick_journal_control_finding(
264 plan: &AcceptedInspectionPlan,
265 store_path: &str,
266) -> IntegrityFinding {
267 let error = InternalError::store_corruption();
268 IntegrityFinding {
269 diagnostic_code: error.diagnostic_code().error_code().raw(),
270 class: IntegrityFindingClass::Corruption,
271 severity: IntegritySeverity::Error,
272 kind: IntegrityFindingKind::JournalControlMismatch,
273 entity: IntegrityEntityIdentity::from_plan(plan),
274 store_path: store_path.to_string(),
275 phase: IntegrityPhase::QuickMetadata,
276 verifier_family: IntegrityVerifierFamily::JournalEnvelope,
277 physical_key: Vec::new(),
278 primary_key: None,
279 field_paths: Vec::new(),
280 value_path: None,
281 constraint_id: None,
282 constraint_name: None,
283 schema_index_id: None,
284 relation_id: None,
285 expected: Some("well-formed-journal-control".to_string()),
286 observed: Some("inconsistent-journal-control".to_string()),
287 }
288}
289
290fn relation_field_paths(plan: &AcceptedInspectionPlan, relation_id: u32) -> Vec<String> {
291 let snapshot = plan.snapshot().persisted_snapshot();
292 let Some(relation) = snapshot
293 .relations()
294 .iter()
295 .find(|relation| relation.id().get() == relation_id)
296 else {
297 return Vec::new();
298 };
299
300 relation
301 .source()
302 .root_field_ids()
303 .iter()
304 .filter_map(|field_id| {
305 snapshot
306 .fields()
307 .iter()
308 .find(|field| field.id() == *field_id)
309 .map(|field| field.name().to_string())
310 })
311 .collect()
312}
313
314const MAX_QUICK_RETURNED_FINDINGS: usize = 64;
315#[cfg(target_arch = "wasm32")]
316const DATABASE_INCARNATION_DOMAIN: &[u8] = b"icydb.database-incarnation.v1";
317static DATABASE_INCARNATION_SEQUENCE: AtomicU64 = AtomicU64::new(0);
318
319#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, Hash, PartialEq)]
326pub struct DatabaseIncarnationId([u8; 16]);
327
328impl DatabaseIncarnationId {
329 pub(crate) fn try_from_bytes(bytes: [u8; 16]) -> Result<Self, InternalError> {
331 if bytes == [0; 16] {
332 return Err(InternalError::database_incarnation_invalid());
333 }
334
335 Ok(Self(bytes))
336 }
337
338 #[must_use]
340 pub const fn to_bytes(self) -> [u8; 16] {
341 self.0
342 }
343
344 fn generate() -> Result<Self, InternalError> {
345 let sequence = DATABASE_INCARNATION_SEQUENCE
346 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |current| {
347 current.checked_add(1)
348 })
349 .map_err(|_| InternalError::database_incarnation_generation_failed())?
350 .checked_add(1)
351 .ok_or_else(InternalError::database_incarnation_generation_failed)?;
352
353 #[cfg(not(target_arch = "wasm32"))]
354 let bytes = {
355 let mut bytes = [0_u8; 16];
356 getrandom::fill(&mut bytes)
357 .map_err(|_| InternalError::database_incarnation_generation_failed())?;
358 bytes
359 };
360
361 #[cfg(target_arch = "wasm32")]
362 let bytes = {
363 use sha2::{Digest, Sha256};
364
365 let mut hasher = Sha256::new();
366 hasher.update(DATABASE_INCARNATION_DOMAIN);
367 hasher.update(ic_cdk::api::canister_self().as_slice());
368 hasher.update(ic_cdk::api::time().to_be_bytes());
369 hasher.update(sequence.to_be_bytes());
370 let digest = hasher.finalize();
371 let mut bytes = [0_u8; 16];
372 bytes.copy_from_slice(&digest[..16]);
373 bytes
374 };
375
376 let _ = sequence;
377 Self::try_from_bytes(bytes)
378 }
379
380 #[cfg(test)]
381 pub(crate) const fn for_tests(fill: u8) -> Self {
382 let mut bytes = [fill; 16];
383 if fill == 0 {
384 bytes[15] = 1;
385 }
386 Self(bytes)
387 }
388}
389
390#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
392pub struct IntegrityEntityIdentity {
393 entity_tag: u64,
394 entity_path: String,
395 store_path: String,
396}
397
398impl IntegrityEntityIdentity {
399 fn from_plan(plan: &AcceptedInspectionPlan) -> Self {
400 Self::from_accepted_identity(plan.identity_ref())
401 }
402
403 pub(in crate::db) fn from_accepted_identity(
404 identity: &crate::db::schema::AcceptedCatalogIdentity,
405 ) -> Self {
406 Self {
407 entity_tag: identity.entity_tag().value(),
408 entity_path: identity.entity_path().to_string(),
409 store_path: identity.store_path().to_string(),
410 }
411 }
412
413 pub(in crate::db) const fn validate(&self) -> Result<(), IntegrityJobError> {
414 if self.entity_tag == 0
415 || self.entity_path.is_empty()
416 || self.entity_path.len() > MAX_INTEGRITY_PATH_BYTES
417 || self.store_path.is_empty()
418 || self.store_path.len() > MAX_INTEGRITY_PATH_BYTES
419 {
420 return Err(IntegrityJobError::InvalidEntityIdentity);
421 }
422 Ok(())
423 }
424
425 #[must_use]
427 pub const fn entity_tag(&self) -> u64 {
428 self.entity_tag
429 }
430
431 #[must_use]
433 pub const fn entity_path(&self) -> &str {
434 self.entity_path.as_str()
435 }
436
437 #[must_use]
439 pub const fn store_path(&self) -> &str {
440 self.store_path.as_str()
441 }
442}
443
444#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
446pub enum IntegrityAuthorityClass {
447 Corruption,
449 IncompatiblePersistedFormat,
451 InvariantViolation,
453 Unsupported,
455 Internal,
457}
458
459#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
461pub enum IntegrityFindingClass {
462 Corruption,
464 IncompatiblePersistedFormat,
466 ResourceLimited,
468}
469
470#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
473pub enum IntegrityFindingKind {
474 MalformedDataKey,
476
477 MalformedRow,
479
480 OversizedRow,
482
483 InvalidFieldValue,
485
486 PrimaryKeyMismatch,
488
489 InvalidIdentityValue,
491
492 IdentityHighWaterExceeded,
494
495 ConstraintViolation,
497
498 MissingIndexEntry,
500
501 DivergentIndexEntry,
503
504 MalformedIndexEntry,
506
507 OrphanIndexEntry,
509
510 DuplicateUniqueIndexKey,
512
513 MissingRelationTarget,
515
516 MissingReverseRelationEntry,
518
519 DivergentReverseRelationEntry,
521
522 MalformedReverseRelationEntry,
524
525 OrphanReverseRelationEntry,
527
528 MalformedJournalBatch,
530
531 JournalSequenceGap,
533
534 DuplicateJournalBatchIdentity,
536
537 JournalControlMismatch,
539}
540
541#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
544pub enum IntegrityPhase {
545 QuickMetadata,
547
548 Rows,
550
551 IndexEntries,
553
554 ReverseRelations,
556
557 JournalTails,
559
560 FinalProofVectorCheck,
562}
563
564#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, Hash, Ord, PartialEq, PartialOrd)]
567pub enum IntegrityVerifierFamily {
568 DataKey,
570
571 RowEnvelope,
573
574 FieldValue,
576
577 PrimaryKey,
579
580 IdentityState,
582
583 ValidatedConstraints,
585
586 ForwardIndex,
588
589 IndexEntry,
591
592 UniqueIndex,
594
595 Relation,
597
598 ReverseRelationEntry,
600
601 JournalEnvelope,
603
604 JournalBatchIdentity,
606}
607
608#[derive(CandidType, Clone, Copy, Debug, Deserialize, Eq, PartialEq)]
610pub enum IntegritySeverity {
611 Error,
613 Advisory,
615}
616
617#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
621pub struct IntegrityFinding {
622 diagnostic_code: u16,
623 class: IntegrityFindingClass,
624 severity: IntegritySeverity,
625 kind: IntegrityFindingKind,
626 entity: IntegrityEntityIdentity,
627 store_path: String,
628 phase: IntegrityPhase,
629 verifier_family: IntegrityVerifierFamily,
630 physical_key: Vec<u8>,
631 primary_key: Option<Vec<u8>>,
632 field_paths: Vec<String>,
633 value_path: Option<Box<ConstraintValuePath>>,
634 constraint_id: Option<u32>,
635 constraint_name: Option<String>,
636 schema_index_id: Option<u32>,
637 relation_id: Option<u32>,
638 expected: Option<String>,
639 observed: Option<String>,
640}
641
642impl IntegrityFinding {
643 #[must_use]
645 pub const fn diagnostic_code(&self) -> u16 {
646 self.diagnostic_code
647 }
648
649 #[must_use]
651 pub const fn class(&self) -> IntegrityFindingClass {
652 self.class
653 }
654
655 #[must_use]
657 pub const fn severity(&self) -> IntegritySeverity {
658 self.severity
659 }
660
661 #[must_use]
663 pub const fn kind(&self) -> IntegrityFindingKind {
664 self.kind
665 }
666
667 #[must_use]
669 pub const fn entity(&self) -> &IntegrityEntityIdentity {
670 &self.entity
671 }
672
673 #[must_use]
675 pub const fn store_path(&self) -> &str {
676 self.store_path.as_str()
677 }
678
679 #[must_use]
681 pub const fn phase(&self) -> IntegrityPhase {
682 self.phase
683 }
684
685 #[must_use]
687 pub const fn verifier_family(&self) -> IntegrityVerifierFamily {
688 self.verifier_family
689 }
690
691 #[must_use]
693 pub const fn physical_key(&self) -> &[u8] {
694 self.physical_key.as_slice()
695 }
696
697 #[must_use]
699 pub fn primary_key(&self) -> Option<&[u8]> {
700 self.primary_key.as_deref()
701 }
702
703 #[must_use]
705 pub const fn field_paths(&self) -> &[String] {
706 self.field_paths.as_slice()
707 }
708
709 #[must_use]
711 pub fn value_path(&self) -> Option<&ConstraintValuePath> {
712 self.value_path.as_deref()
713 }
714
715 #[must_use]
717 pub const fn constraint_id(&self) -> Option<u32> {
718 self.constraint_id
719 }
720
721 #[must_use]
723 pub fn constraint_name(&self) -> Option<&str> {
724 self.constraint_name.as_deref()
725 }
726
727 #[must_use]
729 pub const fn schema_index_id(&self) -> Option<u32> {
730 self.schema_index_id
731 }
732
733 #[must_use]
735 pub const fn relation_id(&self) -> Option<u32> {
736 self.relation_id
737 }
738
739 #[must_use]
741 pub fn expected(&self) -> Option<&str> {
742 self.expected.as_deref()
743 }
744
745 #[must_use]
747 pub fn observed(&self) -> Option<&str> {
748 self.observed.as_deref()
749 }
750}
751
752#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
754pub struct IntegrityAuthorityDiagnostic {
755 diagnostic_code: u16,
756 class: IntegrityAuthorityClass,
757}
758
759impl IntegrityAuthorityDiagnostic {
760 pub(in crate::db) fn from_internal(error: &InternalError) -> Self {
761 let class = match error.class {
762 ErrorClass::Corruption => IntegrityAuthorityClass::Corruption,
763 ErrorClass::IncompatiblePersistedFormat => {
764 IntegrityAuthorityClass::IncompatiblePersistedFormat
765 }
766 ErrorClass::InvariantViolation => IntegrityAuthorityClass::InvariantViolation,
767 ErrorClass::Unsupported | ErrorClass::NotFound | ErrorClass::Conflict => {
768 IntegrityAuthorityClass::Unsupported
769 }
770 ErrorClass::Internal => IntegrityAuthorityClass::Internal,
771 };
772 Self {
773 diagnostic_code: error.diagnostic_code().error_code().raw(),
774 class,
775 }
776 }
777
778 #[must_use]
780 pub const fn diagnostic_code(&self) -> u16 {
781 self.diagnostic_code
782 }
783
784 #[must_use]
786 pub const fn class(&self) -> IntegrityAuthorityClass {
787 self.class
788 }
789}
790
791#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
793pub struct IntegrityResourceDiagnostic {
794 diagnostic_code: u16,
795}
796
797impl IntegrityResourceDiagnostic {
798 #[must_use]
800 pub const fn diagnostic_code(&self) -> u16 {
801 self.diagnostic_code
802 }
803}
804
805#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
807pub enum QuickIntegrityStatus {
808 CompleteClean,
810 CompleteWithFindings,
812 Uninspectable(IntegrityAuthorityDiagnostic),
814 ResourceLimited(IntegrityResourceDiagnostic),
816}
817
818#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
820pub struct QuickIntegrityResult {
821 entity: IntegrityEntityIdentity,
822 database_incarnation_id: DatabaseIncarnationId,
823 accepted_schema_version: u32,
824 accepted_schema_fingerprint: [u8; 16],
825 status: QuickIntegrityStatus,
826 total_findings: u64,
827 omitted_findings: u64,
828 findings: Vec<IntegrityFinding>,
829}
830
831impl QuickIntegrityResult {
832 #[must_use]
834 pub const fn entity(&self) -> &IntegrityEntityIdentity {
835 &self.entity
836 }
837
838 #[must_use]
840 pub const fn database_incarnation_id(&self) -> DatabaseIncarnationId {
841 self.database_incarnation_id
842 }
843
844 #[must_use]
846 pub const fn accepted_schema_version(&self) -> u32 {
847 self.accepted_schema_version
848 }
849
850 #[must_use]
852 pub const fn accepted_schema_fingerprint(&self) -> [u8; 16] {
853 self.accepted_schema_fingerprint
854 }
855
856 #[must_use]
858 pub const fn status(&self) -> &QuickIntegrityStatus {
859 &self.status
860 }
861
862 #[must_use]
864 pub const fn total_findings(&self) -> u64 {
865 self.total_findings
866 }
867
868 #[must_use]
870 pub const fn omitted_findings(&self) -> u64 {
871 self.omitted_findings
872 }
873
874 #[must_use]
876 pub const fn findings(&self) -> &[IntegrityFinding] {
877 self.findings.as_slice()
878 }
879}
880
881struct QuickIntegrityAccumulator {
882 total_findings: u64,
883 findings: Vec<IntegrityFinding>,
884}
885
886impl QuickIntegrityAccumulator {
887 const fn new() -> Self {
888 Self {
889 total_findings: 0,
890 findings: Vec::new(),
891 }
892 }
893
894 fn record(&mut self, finding: IntegrityFinding) -> Result<(), IntegrityResourceDiagnostic> {
895 self.total_findings =
896 self.total_findings
897 .checked_add(1)
898 .ok_or(IntegrityResourceDiagnostic {
899 diagnostic_code: icydb_diagnostic_code::ErrorCode::RUNTIME_INTERNAL.raw(),
900 })?;
901 if self.findings.len() < MAX_QUICK_RETURNED_FINDINGS {
902 self.findings.push(finding);
903 }
904 Ok(())
905 }
906
907 fn complete(
908 self,
909 plan: &AcceptedInspectionPlan,
910 incarnation: DatabaseIncarnationId,
911 ) -> Result<QuickIntegrityResult, InternalError> {
912 let status = if self.total_findings == 0 {
913 QuickIntegrityStatus::CompleteClean
914 } else {
915 QuickIntegrityStatus::CompleteWithFindings
916 };
917 let omitted_findings = self.omitted_findings()?;
918 let identity = plan.identity();
919
920 Ok(QuickIntegrityResult {
921 entity: IntegrityEntityIdentity::from_plan(plan),
922 database_incarnation_id: incarnation,
923 accepted_schema_version: identity.accepted_schema_version().get(),
924 accepted_schema_fingerprint: identity.accepted_schema_fingerprint(),
925 status,
926 total_findings: self.total_findings,
927 omitted_findings,
928 findings: self.findings,
929 })
930 }
931
932 fn resource_limited(
933 self,
934 plan: &AcceptedInspectionPlan,
935 incarnation: DatabaseIncarnationId,
936 diagnostic: IntegrityResourceDiagnostic,
937 ) -> Result<QuickIntegrityResult, InternalError> {
938 let omitted_findings = self.omitted_findings()?;
939 let identity = plan.identity();
940
941 Ok(QuickIntegrityResult {
942 entity: IntegrityEntityIdentity::from_plan(plan),
943 database_incarnation_id: incarnation,
944 accepted_schema_version: identity.accepted_schema_version().get(),
945 accepted_schema_fingerprint: identity.accepted_schema_fingerprint(),
946 status: QuickIntegrityStatus::ResourceLimited(diagnostic),
947 total_findings: self.total_findings,
948 omitted_findings,
949 findings: self.findings,
950 })
951 }
952
953 fn omitted_findings(&self) -> Result<u64, InternalError> {
954 let returned = u64::try_from(self.findings.len()).map_err(|_| {
955 InternalError::classified(ErrorClass::InvariantViolation, ErrorOrigin::Response)
956 })?;
957 self.total_findings.checked_sub(returned).ok_or_else(|| {
958 InternalError::classified(ErrorClass::InvariantViolation, ErrorOrigin::Response)
959 })
960 }
961}
962
963pub(in crate::db) fn uninspectable_quick_integrity(
964 identity: crate::db::schema::AcceptedCatalogIdentity,
965 incarnation: DatabaseIncarnationId,
966 error: &InternalError,
967) -> QuickIntegrityResult {
968 QuickIntegrityResult {
969 entity: IntegrityEntityIdentity::from_accepted_identity(&identity),
970 database_incarnation_id: incarnation,
971 accepted_schema_version: identity.accepted_schema_version().get(),
972 accepted_schema_fingerprint: identity.accepted_schema_fingerprint(),
973 status: QuickIntegrityStatus::Uninspectable(IntegrityAuthorityDiagnostic::from_internal(
974 error,
975 )),
976 total_findings: 0,
977 omitted_findings: 0,
978 findings: Vec::new(),
979 }
980}
981
982pub(in crate::db) fn execute_quick_integrity<C: CanisterKind>(
983 db: &crate::db::Db<C>,
984 plan: &AcceptedInspectionPlan,
985 incarnation: DatabaseIncarnationId,
986) -> Result<QuickIntegrityResult, InternalError> {
987 ensure_recovery_admitted(db)?;
988 let findings = match validate_quick_integrity_control(db, plan, incarnation) {
989 Ok(findings) => findings,
990 Err(error) => {
991 return Ok(uninspectable_quick_integrity(
992 plan.identity(),
993 incarnation,
994 &error,
995 ));
996 }
997 };
998 let mut accumulator = QuickIntegrityAccumulator::new();
999 for finding in findings {
1000 if let Err(diagnostic) = accumulator.record(finding) {
1001 return accumulator.resource_limited(plan, incarnation, diagnostic);
1002 }
1003 }
1004
1005 accumulator.complete(plan, incarnation)
1006}
1007
1008pub(crate) fn generate_database_incarnation_id() -> Result<DatabaseIncarnationId, InternalError> {
1009 DatabaseIncarnationId::generate()
1010}
1011
1012pub(crate) fn generate_cursor_authentication_key() -> Result<[u8; 32], InternalError> {
1018 let first = <crate::types::Ulid as crate::types::GenerateKey>::generate()?;
1019 let second = <crate::types::Ulid as crate::types::GenerateKey>::generate()?;
1020 let mut bytes = [0_u8; 32];
1021 bytes[..16].copy_from_slice(&first.to_bytes());
1022 bytes[16..].copy_from_slice(&second.to_bytes());
1023 if bytes == [0; 32] {
1024 return Err(InternalError::database_incarnation_generation_failed());
1025 }
1026
1027 Ok(bytes)
1028}
1029
1030#[cfg(test)]
1031mod tests {
1032 use super::*;
1033 use crate::{
1034 db::schema::{FieldStorageDecode, LeafCodec, ScalarCodec},
1035 db::{
1036 commit::CommitSchemaFingerprint,
1037 schema::{
1038 AcceptedCatalogIdentity, AcceptedCompositeCatalog, AcceptedFieldKind,
1039 AcceptedSchemaRevision, AcceptedSchemaSnapshot, AcceptedValueCatalogHandle,
1040 FieldId, IdentityState, IdentityStateOwner, PersistedFieldSnapshot,
1041 PersistedSchemaSnapshot, SchemaFieldSlot, SchemaInsertDefault, SchemaRowLayout,
1042 SchemaVersion, empty_accepted_enum_catalog_for_tests,
1043 },
1044 },
1045 types::EntityTag,
1046 };
1047
1048 fn plan() -> AcceptedInspectionPlan {
1049 let revision = AcceptedSchemaRevision::INITIAL;
1050 let identity = AcceptedCatalogIdentity::new(
1051 EntityTag::new(23),
1052 "tests::QuickEntity",
1053 "tests::QuickStore",
1054 revision,
1055 SchemaVersion::initial(),
1056 CommitSchemaFingerprint::from([0x44; 16]),
1057 );
1058 let snapshot = AcceptedSchemaSnapshot::new(PersistedSchemaSnapshot::new(
1059 SchemaVersion::initial(),
1060 "tests::QuickEntity".to_string(),
1061 "QuickEntity".to_string(),
1062 FieldId::new(1),
1063 SchemaRowLayout::initial(vec![(FieldId::new(1), SchemaFieldSlot::new(0))]),
1064 vec![PersistedFieldSnapshot::new_initial(
1065 FieldId::new(1),
1066 "id".to_string(),
1067 SchemaFieldSlot::new(0),
1068 AcceptedFieldKind::Nat64,
1069 Vec::new(),
1070 false,
1071 SchemaInsertDefault::None,
1072 FieldStorageDecode::ByKind,
1073 LeafCodec::Scalar(ScalarCodec::Nat64),
1074 )],
1075 ));
1076 let value_catalog = AcceptedValueCatalogHandle::new_for_tests(
1077 empty_accepted_enum_catalog_for_tests(),
1078 AcceptedCompositeCatalog::empty(),
1079 revision,
1080 );
1081
1082 AcceptedInspectionPlan::compile_relation_free_for_tests(identity, snapshot, value_catalog)
1083 .expect("accepted Quick plan should compile")
1084 }
1085
1086 fn finding(plan: &AcceptedInspectionPlan) -> IntegrityFinding {
1087 IntegrityFinding {
1088 diagnostic_code: icydb_diagnostic_code::ErrorCode::STORE_CORRUPTION.raw(),
1089 class: IntegrityFindingClass::Corruption,
1090 severity: IntegritySeverity::Error,
1091 kind: IntegrityFindingKind::MalformedRow,
1092 entity: IntegrityEntityIdentity::from_plan(plan),
1093 store_path: plan.identity().store_path().to_string(),
1094 phase: IntegrityPhase::Rows,
1095 verifier_family: IntegrityVerifierFamily::RowEnvelope,
1096 physical_key: vec![1],
1097 primary_key: None,
1098 field_paths: Vec::new(),
1099 value_path: None,
1100 constraint_id: None,
1101 constraint_name: None,
1102 schema_index_id: None,
1103 relation_id: None,
1104 expected: None,
1105 observed: None,
1106 }
1107 }
1108
1109 #[test]
1110 fn database_incarnation_rejects_zero_and_round_trips_current_bytes() {
1111 assert!(DatabaseIncarnationId::try_from_bytes([0; 16]).is_err());
1112
1113 let identity = DatabaseIncarnationId::for_tests(7);
1114 assert_eq!(
1115 DatabaseIncarnationId::try_from_bytes(identity.to_bytes())
1116 .expect("nonzero incarnation should decode"),
1117 identity,
1118 );
1119 }
1120
1121 #[test]
1122 fn integrity_finding_candid_preserves_targeted_constraint_path() {
1123 let plan = plan();
1124 let mut finding = finding(&plan);
1125 let path = ConstraintValuePath::new(vec![
1126 crate::error::ConstraintValuePathComponent::RootField { field_id: 1 },
1127 crate::error::ConstraintValuePathComponent::ListElement { index: 2 },
1128 ]);
1129 finding.kind = IntegrityFindingKind::ConstraintViolation;
1130 finding.value_path = Some(Box::new(path.clone()));
1131 finding.constraint_id = Some(7);
1132 finding.constraint_name = Some("nested_limit".to_string());
1133
1134 let bytes = candid::encode_one(&finding).expect("integrity finding should encode");
1135 let decoded: IntegrityFinding =
1136 candid::decode_one(&bytes).expect("integrity finding should decode");
1137 assert_eq!(decoded.value_path(), Some(&path));
1138 assert_eq!(decoded.constraint_id(), Some(7));
1139 assert_eq!(decoded.constraint_name(), Some("nested_limit"));
1140 }
1141
1142 #[test]
1143 fn quick_clean_result_binds_incarnation_and_accepted_plan_identity() {
1144 let plan = plan();
1145 let incarnation = DatabaseIncarnationId::for_tests(8);
1146 let result = QuickIntegrityAccumulator::new()
1147 .complete(&plan, incarnation)
1148 .expect("clean Quick accounting should remain valid");
1149
1150 assert_eq!(result.status(), &QuickIntegrityStatus::CompleteClean);
1151 assert_eq!(result.database_incarnation_id(), incarnation);
1152 assert_eq!(result.accepted_schema_version(), 1);
1153 assert_eq!(result.accepted_schema_fingerprint(), [0x44; 16]);
1154 assert_eq!(result.total_findings(), 0);
1155 assert_eq!(result.omitted_findings(), 0);
1156 }
1157
1158 #[test]
1159 fn quick_findings_keep_a_bounded_prefix_and_exact_omitted_count() {
1160 let plan = plan();
1161 let mut accumulator = QuickIntegrityAccumulator::new();
1162 for _ in 0..=MAX_QUICK_RETURNED_FINDINGS {
1163 accumulator
1164 .record(finding(&plan))
1165 .expect("bounded test finding count should fit");
1166 }
1167 let result = accumulator
1168 .complete(&plan, DatabaseIncarnationId::for_tests(9))
1169 .expect("one-over-cap Quick accounting should remain valid");
1170
1171 assert_eq!(result.status(), &QuickIntegrityStatus::CompleteWithFindings,);
1172 assert_eq!(result.total_findings(), 65);
1173 assert_eq!(result.findings().len(), MAX_QUICK_RETURNED_FINDINGS);
1174 assert_eq!(result.omitted_findings(), 1);
1175 assert_eq!(
1176 result.total_findings(),
1177 result.findings().len() as u64 + result.omitted_findings(),
1178 );
1179 }
1180
1181 #[test]
1182 fn quick_findings_at_the_exact_returned_cap_have_no_omissions() {
1183 let plan = plan();
1184 let mut accumulator = QuickIntegrityAccumulator::new();
1185 for _ in 0..MAX_QUICK_RETURNED_FINDINGS {
1186 accumulator
1187 .record(finding(&plan))
1188 .expect("exact-cap finding count should fit");
1189 }
1190 let result = accumulator
1191 .complete(&plan, DatabaseIncarnationId::for_tests(10))
1192 .expect("exact-cap Quick accounting should remain valid");
1193
1194 assert_eq!(result.total_findings(), 64);
1195 assert_eq!(result.findings().len(), MAX_QUICK_RETURNED_FINDINGS);
1196 assert_eq!(result.omitted_findings(), 0);
1197 }
1198
1199 #[test]
1200 fn quick_selected_authority_failure_is_not_a_clean_completion() {
1201 let plan = plan();
1202 let error = InternalError::accepted_row_constraint_program_corrupt();
1203 let result = uninspectable_quick_integrity(
1204 plan.identity(),
1205 DatabaseIncarnationId::for_tests(11),
1206 &error,
1207 );
1208
1209 assert!(matches!(
1210 result.status(),
1211 QuickIntegrityStatus::Uninspectable(IntegrityAuthorityDiagnostic {
1212 class: IntegrityAuthorityClass::Corruption,
1213 ..
1214 }),
1215 ));
1216 assert_eq!(result.total_findings(), 0);
1217 assert_eq!(result.omitted_findings(), 0);
1218 }
1219
1220 #[test]
1221 fn quick_identity_inventory_rejects_active_retired_owner_collision_first() {
1222 let incarnation = DatabaseIncarnationId::for_tests(12);
1223 let owner = IdentityStateOwner::try_new(incarnation, EntityTag::new(31), FieldId::new(1))
1224 .expect("identity owner should admit");
1225 let active = IdentityState::new_active(owner, AcceptedFieldKind::Nat64)
1226 .expect("active identity state should admit");
1227 let retired = active.retire().expect("active state should retire");
1228 let mut owners = BTreeMap::new();
1229
1230 record_quick_identity_owner(&mut owners, "tests::first", &active)
1231 .expect("the first owner should admit");
1232 let error = record_quick_identity_owner(&mut owners, "tests::second", &retired)
1233 .expect_err("an active/retired owner collision must reject");
1234
1235 assert_eq!(error.class(), ErrorClass::Corruption);
1236 assert_eq!(error.origin(), ErrorOrigin::Identity);
1237 }
1238}