1use crate::{
7 db::{
8 Db,
9 codec::{
10 finalize_hash_sha256, new_hash_sha256_prefixed, write_hash_len_u32, write_hash_str_u32,
11 write_hash_tag_u8, write_hash_u64,
12 },
13 commit::{
14 AcceptedSchemaPublication, DatabaseControlOp, database_incarnation_id,
15 ensure_recovery_admitted, publish_accepted_schema_candidates_with_application_record,
16 publish_accepted_schema_candidates_with_database_control,
17 publish_generated_row_local_abort_with_application_record,
18 },
19 data::DataStore,
20 index::{IndexState, IndexStore},
21 integrity::DatabaseIncarnationId,
22 registry::{
23 StoreAllocationIdentity, StoreAllocationIdentityCapability, StoreCommitParticipation,
24 StoreDurability, StoreHandle, StoreRecoveryCapability, StoreRelationSourceCapability,
25 StoreRelationTargetCapability, StoreRuntimeStorageMode, StoreSchemaMetadataCapability,
26 },
27 relation::prove_empty_reverse_relation_domain,
28 schema::ensure_schema_migration_ready_for_schema_changes,
29 schema::{
30 AcceptedSchemaRevision, AcceptedSchemaRevisionBundle, CandidateSchemaRevision,
31 ConstraintActivationKind, ConstraintActivationState, ConstraintId, ConstraintOrigin,
32 ConstraintValidationPhase, ConstraintValidationProgress, ExistingProposalStore,
33 MAX_IDENTITY_STATE_RECORDS_PER_DATABASE, ProposalStoreTarget, SchemaApplicationRecord,
34 SchemaApplicationRecordOp, SchemaChangeActivation, SchemaChangeJob, SchemaChangeJobId,
35 SchemaChangeOutcome, SchemaChangeProgress, SchemaChangeProgressStatus,
36 SchemaChangeReceipt, SchemaChangeValidationPhase, StagedUserIndexDomainError,
37 UnpublishedRowLocalValidation, advance_accepted_row_local_constraint_activation,
38 constraint_validation_finding_output, derive_schema_change_job_id,
39 load_schema_application_record_read_only, lower_existing_schema_proposal,
40 lower_generated_existing_schema_proposal, lower_initial_schema_proposal,
41 prove_empty_user_index_domain, validate_unpublished_row_local_candidate_bounded,
42 with_schema_application_store,
43 },
44 },
45 error::InternalError,
46 traits::CanisterKind,
47 types::EntityTag,
48};
49use candid::CandidType;
50use icydb_schema::{
51 ExpectedAcceptedHead, ExpectedSchemaFingerprint, SchemaProposal, SchemaProposalDigest,
52 SchemaSubmissionKey, TargetDatabaseIdentity, TargetStoreIdentity,
53};
54use serde::Deserialize;
55use sha2::Digest;
56use std::cell::Cell;
57#[cfg(feature = "migration")]
58use std::collections::BTreeMap;
59
60#[cfg(feature = "migration")]
61use crate::db::commit::{
62 RecoveryProgress, StartupRecoveryFailure, continue_recovery_with_failure_authority,
63};
64
65#[cfg(feature = "migration")]
66use crate::db::schema::{
67 PersistedSchemaMigrationEntity, PersistedSchemaMigrationFindingKind,
68 PersistedSchemaMigrationIndex, PersistedSchemaMigrationPhase,
69 PersistedSchemaMigrationTransition, SchemaMigrationCommand, SchemaMigrationEntityTransition,
70 SchemaMigrationFinding, SchemaMigrationFindingKind, SchemaMigrationPhase,
71 SchemaMigrationReceipt, SchemaMigrationRecord, SchemaMigrationRecordOp,
72 SchemaMigrationStatusPage, SchemaMigrationStatusRequest,
73 live_schema_checkpoint::{load_entity_source_lineage_catalog, load_schema_migration_record},
74 migration_execution::{
75 cleanup_migration_staging_page, final_validate_migration_page,
76 migration_derived_domain_count, publish_migration_rewrite_page, rewrite_migration_page,
77 },
78 migration_lineage::{
79 AcceptedEntitySourceLineage, AcceptedEntitySourceLineageCatalog,
80 AcceptedEntitySourceLineageState, EntitySourceLineageCatalogOp,
81 },
82 migration_planner::{
83 PlannedEntitySourceLineage, SchemaMigrationPlanningError, plan_entity_source_adoption,
84 plan_initial_entity_source_lineage, plan_schema_migration,
85 },
86 migration_validation::{stage_migration_index_entries, validate_migration_page},
87};
88
89#[cfg(feature = "migration")]
90use icydb_diagnostic_code::SchemaMigrationCode;
91#[cfg(feature = "migration")]
92use icydb_schema::{EntitySourceKey, SchemaMigrationPlanDigest};
93
94const DATABASE_TARGET_FINGERPRINT_PROFILE: &[u8] = b"icydb.schema-target.database.v1";
95const STORE_TARGET_FINGERPRINT_PROFILE: &[u8] = b"icydb.schema-target.store.v1";
96const ACCEPTED_DATABASE_HEAD_FINGERPRINT_PROFILE: &[u8] = b"icydb.accepted-schema.database-head.v1";
97#[cfg(feature = "migration")]
98const SCHEMA_MIGRATION_SUBMISSION_PROFILE: &[u8] = b"icydb.schema-migration.submission.v1";
99
100#[derive(Clone, Copy)]
101struct GeneratedDatabaseIdentityCacheEntry {
102 registry: usize,
103 incarnation: DatabaseIncarnationId,
104 identity: TargetDatabaseIdentity,
105}
106
107thread_local! {
108 static GENERATED_DATABASE_IDENTITY: Cell<Option<GeneratedDatabaseIdentityCacheEntry>> =
112 const { Cell::new(None) };
113}
114
115#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
123pub struct SchemaApplicationStore {
124 path: String,
125 identity: TargetStoreIdentity,
126}
127
128impl SchemaApplicationStore {
129 #[must_use]
131 pub const fn path(&self) -> &str {
132 self.path.as_str()
133 }
134
135 #[must_use]
137 pub const fn identity(&self) -> TargetStoreIdentity {
138 self.identity
139 }
140}
141
142#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
151pub struct SchemaApplicationTarget {
152 database_identity: TargetDatabaseIdentity,
153 accepted_head: ExpectedAcceptedHead,
154 stores: Vec<SchemaApplicationStore>,
155}
156
157impl SchemaApplicationTarget {
158 #[must_use]
160 pub const fn database_identity(&self) -> TargetDatabaseIdentity {
161 self.database_identity
162 }
163
164 #[must_use]
166 pub const fn accepted_head(&self) -> &ExpectedAcceptedHead {
167 &self.accepted_head
168 }
169
170 #[must_use]
172 pub const fn stores(&self) -> &[SchemaApplicationStore] {
173 self.stores.as_slice()
174 }
175}
176
177#[derive(Clone, Copy)]
185struct StoreApplicationAuthority {
186 path: &'static str,
187 handle: StoreHandle,
188}
189
190struct PendingApplicationAbort {
192 authority: StoreApplicationAuthority,
193 current: AcceptedSchemaRevisionBundle,
194 entity_tag: EntityTag,
195 constraint_id: ConstraintId,
196 remove_validation_job: bool,
197}
198
199#[derive(Clone, Copy, Debug, Eq, PartialEq)]
207struct AcceptedStoreHead {
208 revision: u64,
209 fingerprint: [u8; 32],
210}
211
212#[derive(Clone)]
214struct DirectGeneratedRowLocalProof {
215 candidate_index: usize,
216 store: StoreHandle,
217 store_path: &'static str,
218 entity_tag: crate::types::EntityTag,
219 entity_path: String,
220 constraint_id: ConstraintId,
221 historical_rows: u64,
222}
223
224#[derive(Clone)]
226struct PendingGeneratedRowLocalConstraint {
227 proof: DirectGeneratedRowLocalProof,
228}
229
230struct LoweredApplication {
232 current_bundles: Vec<Option<AcceptedSchemaRevisionBundle>>,
233 candidates: Vec<CandidateSchemaRevision>,
234 pending: Option<PendingGeneratedRowLocalConstraint>,
235}
236
237pub(in crate::db) fn schema_application_target<C: CanisterKind>(
239 db: &Db<C>,
240) -> Result<SchemaApplicationTarget, InternalError> {
241 ensure_recovery_admitted(db)?;
242 let incarnation = database_incarnation_id()?;
243 let mut stores = db.with_store_registry(|registry| {
244 registry
245 .iter()
246 .map(|(path, handle)| StoreApplicationAuthority { path, handle })
247 .collect::<Vec<_>>()
248 });
249 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.path.cmp(right.path));
250
251 let database_identity = derive_database_identity(incarnation.to_bytes(), stores.as_slice());
252 let mut accepted_heads = Vec::with_capacity(stores.len());
253 let mut application_stores = Vec::with_capacity(stores.len());
254 for store in &stores {
255 let root = store
256 .handle
257 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_root)?
258 .map(|selection| AcceptedStoreHead {
259 revision: selection.root().revision().get(),
260 fingerprint: selection.root().fingerprint().as_bytes(),
261 });
262 accepted_heads.push((store.path, root));
263 application_stores.push(SchemaApplicationStore {
264 path: store.path.to_string(),
265 identity: derive_store_identity(database_identity, store),
266 });
267 }
268
269 Ok(SchemaApplicationTarget {
270 database_identity,
271 accepted_head: derive_accepted_head(accepted_heads.as_slice()),
272 stores: application_stores,
273 })
274}
275
276pub(in crate::db) fn schema_application_receipt<C: CanisterKind>(
279 db: &Db<C>,
280 database_identity: TargetDatabaseIdentity,
281 submission_key: &SchemaSubmissionKey,
282) -> Result<Option<SchemaChangeReceipt>, InternalError> {
283 ensure_recovery_admitted(db)?;
284 with_schema_application_store(|store| {
285 store
286 .load(database_identity, submission_key)
287 .map(|record| record.map(|record| record.receipt().clone()))
288 })
289}
290
291fn exact_schema_application_receipt(
292 proposal: &SchemaProposal,
293 proposal_digest: SchemaProposalDigest,
294) -> Result<Option<SchemaChangeReceipt>, InternalError> {
295 let Some(record) = with_schema_application_store(|store| {
296 store.load(proposal.target_database(), proposal.submission_key())
297 })?
298 else {
299 return Ok(None);
300 };
301 let receipt = record.receipt();
302 if !receipt.is_exact_submission(
303 proposal.target_database(),
304 proposal.submission_key(),
305 proposal_digest,
306 proposal.expected_head(),
307 ) {
308 return Err(InternalError::schema_application_conflict());
309 }
310 Ok(Some(receipt.clone()))
311}
312
313pub(in crate::db) fn continue_schema_application<C: CanisterKind>(
316 db: &Db<C>,
317 job_id: SchemaChangeJobId,
318 acknowledged_receipt: Option<u64>,
319) -> Result<SchemaChangeProgress, InternalError> {
320 ensure_recovery_admitted(db)?;
321 ensure_schema_migration_ready_for_schema_changes()?;
322 let record = with_schema_application_store(|store| store.load_job(job_id))?
323 .ok_or_else(InternalError::schema_application_conflict)?;
324 let target = schema_application_target(db)?;
325 if target.database_identity() != record.receipt().database_identity() {
326 return Err(InternalError::schema_application_conflict());
327 }
328 let candidate_head = match record.receipt().outcome() {
329 SchemaChangeOutcome::Pending {
330 job,
331 candidate_head,
332 } if job.id() == job_id => candidate_head,
333 SchemaChangeOutcome::Applied { .. } => {
334 return Ok(SchemaChangeProgress::new(
335 record.receipt().clone(),
336 SchemaChangeProgressStatus::Applied,
337 ));
338 }
339 SchemaChangeOutcome::Aborted { .. } => {
340 return Ok(SchemaChangeProgress::new(
341 record.receipt().clone(),
342 SchemaChangeProgressStatus::Aborted,
343 ));
344 }
345 _ => return Err(InternalError::store_corruption()),
346 };
347 let [activation] = record.activations() else {
348 return Err(InternalError::store_corruption());
349 };
350 let authorities = application_authorities(db);
351 let authority = authorities
352 .iter()
353 .find(|authority| {
354 derive_store_identity(target.database_identity(), authority) == activation.store()
355 })
356 .ok_or_else(InternalError::store_corruption)?;
357 let entity_tag = EntityTag::new(activation.entity_tag());
358 let constraint_id = ConstraintId::new(activation.constraint_id())
359 .ok_or_else(InternalError::store_corruption)?;
360 let bundle = authority
361 .handle
362 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
363 .ok_or_else(InternalError::store_corruption)?;
364 if bundle.store_path() != authority.path {
365 return Err(InternalError::store_corruption());
366 }
367 let snapshot = bundle
368 .entity_snapshots()
369 .get(&entity_tag)
370 .ok_or_else(InternalError::store_corruption)?;
371
372 let accepted = snapshot
373 .constraint_catalog()
374 .constraints()
375 .iter()
376 .any(|constraint| {
377 constraint.id() == constraint_id
378 && constraint.origin() == ConstraintOrigin::Generated
379 && matches!(
380 constraint.kind(),
381 crate::db::schema::AcceptedConstraintKind::Check { .. }
382 | crate::db::schema::AcceptedConstraintKind::TargetedRule { .. }
383 )
384 });
385 let pending = snapshot.constraint_catalog().activation(constraint_id);
386 if accepted && pending.is_none() {
387 return finalize_schema_application(
388 db,
389 &record,
390 candidate_head,
391 SchemaChangeProgressStatus::Applied,
392 );
393 }
394 let pending = pending.ok_or_else(InternalError::store_corruption)?;
395 if pending.origin() != ConstraintOrigin::Generated
396 || !matches!(
397 pending.kind(),
398 ConstraintActivationKind::Check { .. } | ConstraintActivationKind::TargetedRule { .. }
399 )
400 {
401 return Err(InternalError::store_corruption());
402 }
403 let entity_path = snapshot.entity_path().to_string();
404 let progress = advance_accepted_row_local_constraint_activation(
405 authority.handle,
406 authority.path,
407 entity_tag,
408 entity_path.as_str(),
409 constraint_id,
410 acknowledged_receipt,
411 )?;
412 let status = schema_change_progress_status(snapshot, entity_tag, constraint_id, progress)?;
413 if status == SchemaChangeProgressStatus::Applied {
414 finalize_schema_application(db, &record, candidate_head, status)
415 } else {
416 Ok(SchemaChangeProgress::new(record.receipt().clone(), status))
417 }
418}
419
420pub(in crate::db) fn abort_schema_application<C: CanisterKind>(
426 db: &Db<C>,
427 job_id: SchemaChangeJobId,
428 acknowledged_receipt: Option<u64>,
429) -> Result<SchemaChangeProgress, InternalError> {
430 ensure_recovery_admitted(db)?;
431 ensure_schema_migration_ready_for_schema_changes()?;
432 let record = with_schema_application_store(|store| store.load_job(job_id))?
433 .ok_or_else(InternalError::schema_application_conflict)?;
434 let target = schema_application_target(db)?;
435 if target.database_identity() != record.receipt().database_identity() {
436 return Err(InternalError::schema_application_conflict());
437 }
438 match record.receipt().outcome() {
439 SchemaChangeOutcome::Applied { .. } => {
440 return Ok(SchemaChangeProgress::new(
441 record.receipt().clone(),
442 SchemaChangeProgressStatus::Applied,
443 ));
444 }
445 SchemaChangeOutcome::Aborted { .. } => {
446 return Ok(SchemaChangeProgress::new(
447 record.receipt().clone(),
448 SchemaChangeProgressStatus::Aborted,
449 ));
450 }
451 SchemaChangeOutcome::Pending { job, .. } if job.id() == job_id => {}
452 SchemaChangeOutcome::NoOp { .. } | SchemaChangeOutcome::Pending { .. } => {
453 return Err(InternalError::store_corruption());
454 }
455 }
456
457 let authorities = application_authorities(db);
458 let abort = prepare_pending_application_abort(
459 target.database_identity(),
460 &record,
461 authorities.as_slice(),
462 acknowledged_receipt,
463 )?;
464 let candidate = aborted_generated_row_local_candidate(
465 &abort.current,
466 abort.entity_tag,
467 abort.constraint_id,
468 )?;
469 let accepted_head =
470 accepted_head_after_candidates(authorities.as_slice(), std::slice::from_ref(&candidate))?;
471 let receipt = SchemaChangeReceipt::new(
472 record.receipt().database_identity(),
473 record.receipt().submission_key().clone(),
474 record.receipt().proposal_digest(),
475 record.receipt().prior_head().clone(),
476 SchemaChangeOutcome::Aborted { accepted_head },
477 )?;
478 let terminal = SchemaApplicationRecord::new(receipt.clone(), Vec::new())?;
479 let operation = SchemaApplicationRecordOp::replace(&record, &terminal)?;
480 if abort.remove_validation_job {
481 publish_generated_row_local_abort_with_application_record(
482 abort.authority.path,
483 abort.authority.handle,
484 abort.current.revision(),
485 &candidate,
486 abort.entity_tag,
487 abort.constraint_id,
488 operation,
489 )?;
490 } else {
491 publish_accepted_schema_candidates_with_application_record(
492 vec![AcceptedSchemaPublication::new(
493 abort.authority.path,
494 abort.authority.handle,
495 abort.current.revision(),
496 &candidate,
497 )],
498 operation,
499 )?;
500 }
501 Ok(SchemaChangeProgress::new(
502 receipt,
503 SchemaChangeProgressStatus::Aborted,
504 ))
505}
506
507fn prepare_pending_application_abort(
508 database_identity: TargetDatabaseIdentity,
509 record: &SchemaApplicationRecord,
510 authorities: &[StoreApplicationAuthority],
511 acknowledged_receipt: Option<u64>,
512) -> Result<PendingApplicationAbort, InternalError> {
513 let [activation] = record.activations() else {
514 return Err(InternalError::store_corruption());
515 };
516 let authority = authorities
517 .iter()
518 .copied()
519 .find(|authority| derive_store_identity(database_identity, authority) == activation.store())
520 .ok_or_else(InternalError::store_corruption)?;
521 let entity_tag = EntityTag::new(activation.entity_tag());
522 let constraint_id = ConstraintId::new(activation.constraint_id())
523 .ok_or_else(InternalError::store_corruption)?;
524 let current = authority
525 .handle
526 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
527 .ok_or_else(InternalError::store_corruption)?;
528 if current.store_path() != authority.path {
529 return Err(InternalError::store_corruption());
530 }
531 let pending = current
532 .entity_snapshots()
533 .get(&entity_tag)
534 .and_then(|snapshot| snapshot.constraint_catalog().activation(constraint_id))
535 .filter(|pending| {
536 pending.origin() == ConstraintOrigin::Generated
537 && matches!(
538 pending.kind(),
539 ConstraintActivationKind::Check { .. }
540 | ConstraintActivationKind::TargetedRule { .. }
541 )
542 })
543 .ok_or_else(InternalError::store_corruption)?;
544 let remove_validation_job = pending_generated_row_local_job_retirement(
545 authority,
546 entity_tag,
547 constraint_id,
548 pending.state(),
549 acknowledged_receipt,
550 )?;
551 Ok(PendingApplicationAbort {
552 authority,
553 current,
554 entity_tag,
555 constraint_id,
556 remove_validation_job,
557 })
558}
559
560fn pending_generated_row_local_job_retirement(
561 authority: StoreApplicationAuthority,
562 entity_tag: EntityTag,
563 constraint_id: ConstraintId,
564 state: ConstraintActivationState,
565 acknowledged_receipt: Option<u64>,
566) -> Result<bool, InternalError> {
567 let job = authority
568 .handle
569 .with_schema(|store| store.constraint_validation_job(entity_tag, constraint_id))?;
570 match state {
571 ConstraintActivationState::EnforcingNewWrites => {
572 if acknowledged_receipt.is_some() || job.is_some() {
573 return Err(InternalError::schema_application_conflict());
574 }
575 Ok(false)
576 }
577 ConstraintActivationState::Validating => {
578 let mut job = job.ok_or_else(InternalError::store_corruption)?;
579 if !job.acknowledge_receipt(acknowledged_receipt) {
580 return Err(InternalError::schema_application_conflict());
581 }
582 Ok(true)
583 }
584 }
585}
586
587fn aborted_generated_row_local_candidate(
588 current: &AcceptedSchemaRevisionBundle,
589 entity_tag: EntityTag,
590 constraint_id: ConstraintId,
591) -> Result<CandidateSchemaRevision, InternalError> {
592 let snapshot = current
593 .entity_snapshots()
594 .get(&entity_tag)
595 .cloned()
596 .ok_or_else(InternalError::store_corruption)?;
597 let _activation = snapshot
598 .constraint_catalog()
599 .activation(constraint_id)
600 .filter(|activation| {
601 activation.origin() == ConstraintOrigin::Generated
602 && matches!(
603 activation.kind(),
604 ConstraintActivationKind::Check { .. }
605 | ConstraintActivationKind::TargetedRule { .. }
606 )
607 })
608 .ok_or_else(InternalError::store_corruption)?;
609 let catalog = snapshot
610 .constraint_catalog()
611 .clone()
612 .with_aborted_activation(constraint_id)
613 .map_err(|_| InternalError::store_invariant())?;
614 let accepted_identity_remains = catalog
615 .constraints()
616 .iter()
617 .any(|constraint| constraint.id() == constraint_id);
618 let mut snapshots = current.entity_snapshots().clone();
619 snapshots.insert(entity_tag, snapshot.with_constraint_catalog(catalog));
620 let mut source_bindings = current.source_bindings().clone();
621 if !accepted_identity_remains {
622 source_bindings.remove_constraint_identity(entity_tag, constraint_id)?;
623 }
624 let revision = current
625 .revision()
626 .checked_next()
627 .ok_or_else(InternalError::store_unsupported)?;
628 let bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
629 revision,
630 current.store_path(),
631 current.enum_catalog().clone(),
632 current.composite_catalog().clone(),
633 source_bindings,
634 snapshots,
635 )?;
636 CandidateSchemaRevision::new(bundle)
637}
638
639pub(in crate::db) fn apply_schema<C: CanisterKind>(
642 db: &Db<C>,
643 proposal: &SchemaProposal,
644) -> Result<SchemaChangeReceipt, InternalError> {
645 apply_schema_with_contract::<C, true>(db, proposal)
646}
647
648pub(in crate::db) fn apply_generated_schema<C: CanisterKind>(
651 db: &Db<C>,
652 proposal: &SchemaProposal,
653) -> Result<SchemaChangeReceipt, InternalError> {
654 if !proposal.removals().is_empty() {
655 return Err(InternalError::store_invariant());
656 }
657 apply_schema_with_contract::<C, false>(db, proposal)
658}
659
660fn apply_schema_with_contract<C: CanisterKind, const ALLOW_REMOVALS: bool>(
661 db: &Db<C>,
662 proposal: &SchemaProposal,
663) -> Result<SchemaChangeReceipt, InternalError> {
664 ensure_recovery_admitted(db)?;
665 ensure_schema_migration_ready_for_schema_changes()?;
666 let proposal_digest = proposal
667 .digest()
668 .map_err(|_| InternalError::store_unsupported())?;
669 if let Some(receipt) = exact_schema_application_receipt(proposal, proposal_digest)? {
670 return Ok(receipt);
671 }
672
673 let target = schema_application_target(db)?;
674 if target.database_identity() != proposal.target_database()
675 || target.accepted_head() != proposal.expected_head()
676 {
677 return Err(InternalError::schema_application_conflict());
678 }
679
680 preflight_ordinary_source_application(db, proposal, &target)?;
681
682 let authorities = application_authorities(db);
683 let LoweredApplication {
684 current_bundles,
685 candidates,
686 pending,
687 } = lower_application_candidates::<ALLOW_REMOVALS>(&target, proposal, authorities.as_slice())?;
688 validate_database_identity_state_capacity(
689 authorities.as_slice(),
690 candidates.as_slice(),
691 database_incarnation_id()?,
692 )?;
693 let accepted_head = if let Some(pending) = pending.as_ref() {
694 let final_candidates =
695 final_candidates_for_pending_row_local_constraint(&candidates, pending)?;
696 accepted_head_after_candidates(authorities.as_slice(), final_candidates.as_slice())?
697 } else if candidates.is_empty() {
698 target.accepted_head().clone()
699 } else {
700 accepted_head_after_candidates(authorities.as_slice(), candidates.as_slice())?
701 };
702 #[cfg(feature = "migration")]
703 let outcome_head = accepted_head.clone();
704 #[cfg(not(feature = "migration"))]
705 let outcome_head = accepted_head;
706 let outcome = if pending.is_some() {
707 let job_id = derive_schema_change_job_id(
708 target.database_identity(),
709 proposal.submission_key(),
710 proposal_digest,
711 target.accepted_head(),
712 )?;
713 SchemaChangeOutcome::Pending {
714 job: SchemaChangeJob::new(job_id),
715 candidate_head: outcome_head,
716 }
717 } else if candidates.is_empty() {
718 SchemaChangeOutcome::NoOp {
719 accepted_head: outcome_head,
720 }
721 } else {
722 SchemaChangeOutcome::Applied {
723 accepted_head: outcome_head,
724 }
725 };
726 let receipt = SchemaChangeReceipt::new(
727 target.database_identity(),
728 proposal.submission_key().clone(),
729 proposal_digest,
730 target.accepted_head().clone(),
731 outcome,
732 )?;
733 let activations = match pending {
734 Some(pending) => {
735 let authority = authorities
736 .iter()
737 .find(|authority| authority.path == pending.proof.store_path)
738 .ok_or_else(InternalError::store_invariant)?;
739 vec![SchemaChangeActivation::new(
740 derive_store_identity(target.database_identity(), authority),
741 pending.proof.entity_tag.value(),
742 pending.proof.constraint_id.get(),
743 )?]
744 }
745 None => Vec::new(),
746 };
747 let record = SchemaApplicationRecord::new(receipt.clone(), activations)?;
748 let operation = SchemaApplicationRecordOp::insert(&record)?;
749 #[cfg(feature = "migration")]
750 let database_control = attach_ordinary_lineage_publication(
751 proposal,
752 target.accepted_head(),
753 &accepted_head,
754 candidates.as_slice(),
755 operation,
756 )?;
757 #[cfg(not(feature = "migration"))]
758 let database_control = vec![DatabaseControlOp::SchemaApplication(operation)];
759 let publications =
760 application_publications(authorities.as_slice(), ¤t_bundles, &candidates)?;
761 publish_accepted_schema_candidates_with_database_control(publications, database_control)?;
762 Ok(receipt)
763}
764
765fn preflight_ordinary_source_application<C: CanisterKind>(
766 db: &Db<C>,
767 proposal: &SchemaProposal,
768 target: &SchemaApplicationTarget,
769) -> Result<(), InternalError> {
770 if proposal.migration().is_none()
771 || matches!(target.accepted_head(), ExpectedAcceptedHead::Empty)
772 {
773 return Ok(());
774 }
775 #[cfg(feature = "migration")]
776 if current_proposal_lineage_is_applied(db, proposal, target.accepted_head())? {
777 return Ok(());
778 }
779 #[cfg(feature = "migration")]
780 preflight_unpublished_schema_migration(target, proposal, db)?;
781 #[cfg(not(feature = "migration"))]
782 let _ = db;
783 Err(InternalError::store_unsupported())
784}
785
786#[cfg(feature = "migration")]
790pub(in crate::db) fn migrate_schema<C: CanisterKind>(
791 db: &Db<C>,
792 proposal: &SchemaProposal,
793 command: SchemaMigrationCommand,
794) -> Result<SchemaMigrationStatusPage, InternalError> {
795 ensure_recovery_admitted(db)?;
796 match command {
797 SchemaMigrationCommand::Adopt {
798 expected_database,
799 expected_head,
800 } => adopt_entity_source_lineage(db, proposal, expected_database, &expected_head),
801 SchemaMigrationCommand::Advance {
802 expected_database,
803 expected_head,
804 expected_plan,
805 acknowledged_finding_page,
806 } => advance_metadata_schema_migration(
807 db,
808 proposal,
809 expected_database,
810 &expected_head,
811 expected_plan,
812 acknowledged_finding_page,
813 ),
814 SchemaMigrationCommand::Abort {
815 expected_database,
816 expected_head,
817 expected_plan,
818 } => {
819 let plan = proposal.migration().ok_or_else(|| {
820 InternalError::schema_migration(SchemaMigrationCode::MissingMigration)
821 })?;
822 if proposal.target_database() != expected_database || plan.digest() != expected_plan {
823 return Err(InternalError::schema_migration(
824 SchemaMigrationCode::PlanChanged,
825 ));
826 }
827 if let Some(record) = exact_active_migration_record(
828 proposal,
829 expected_database,
830 &expected_head,
831 expected_plan,
832 )? {
833 let target = schema_application_target(db)?;
834 validate_active_migration_target(&record, &target)?;
835 if record.phase() == PersistedSchemaMigrationPhase::Applied
836 || record.phase() == PersistedSchemaMigrationPhase::Aborted
837 {
838 return active_migration_status(proposal, &target, &record);
839 }
840 if !record.phase().abortable() {
841 return Err(InternalError::schema_migration(
842 SchemaMigrationCode::AbortTooLate,
843 ));
844 }
845 let planned = recompile_active_physical_migration(db, proposal, &record)?;
846 let authorities = application_authorities(db);
847 let mut store_identities = BTreeMap::new();
848 for authority in &authorities {
849 store_identities.insert(
850 authority.path,
851 derive_store_identity(record.database_identity(), authority),
852 );
853 }
854 let (progress, exhausted) = cleanup_migration_staging_page(
855 db,
856 &planned,
857 record.progress(),
858 &store_identities,
859 )?;
860 let phase = if exhausted {
861 PersistedSchemaMigrationPhase::Aborted
862 } else {
863 record.phase()
864 };
865 let advanced = record.transition(phase, progress)?;
866 let operation = SchemaMigrationRecordOp::replace(&record, &advanced)?;
867 publish_accepted_schema_candidates_with_database_control(
868 Vec::new(),
869 vec![DatabaseControlOp::SchemaMigration(operation)],
870 )?;
871 return active_migration_status(proposal, &target, &advanced);
872 }
873 let target = exact_migration_target(db, expected_database, &expected_head)?;
874 let status = schema_migration_status_for_target(db, proposal, &target)?;
875 if status.phase() == SchemaMigrationPhase::Applied {
876 Ok(status)
877 } else {
878 Err(InternalError::schema_migration(
881 SchemaMigrationCode::MissingMigration,
882 ))
883 }
884 }
885 }
886}
887
888#[cfg(feature = "migration")]
890pub(in crate::db) fn schema_migration_status<C: CanisterKind>(
891 db: &Db<C>,
892 proposal: &SchemaProposal,
893 request: &SchemaMigrationStatusRequest,
894) -> Result<SchemaMigrationStatusPage, InternalError> {
895 ensure_recovery_admitted(db)?;
896 if !request.validate() || request.cursor().is_some() {
897 return Err(InternalError::cursor_invalid_continuation());
898 }
899 let target = schema_application_target(db)?;
900 schema_migration_status_for_target(db, proposal, &target)
901}
902
903#[cfg(feature = "migration")]
907pub(in crate::db) fn defer_generated_schema_application_for_prepared_migration<C: CanisterKind>(
908 db: &Db<C>,
909 proposal: &SchemaProposal,
910) -> Result<bool, InternalError> {
911 ensure_recovery_admitted(db)?;
912 let Some(record) = load_schema_migration_record()? else {
913 return Ok(false);
914 };
915 if matches!(
916 record.phase(),
917 PersistedSchemaMigrationPhase::Applied | PersistedSchemaMigrationPhase::Aborted
918 ) {
919 return Ok(false);
920 }
921 validate_active_migration_deployment(proposal, &record)?;
922 let target = schema_application_target(db)?;
923 validate_active_migration_target(&record, &target)?;
924 match record.phase() {
925 PersistedSchemaMigrationPhase::Prepared => Ok(true),
926 PersistedSchemaMigrationPhase::Validating
927 | PersistedSchemaMigrationPhase::ReadyToRewrite
928 | PersistedSchemaMigrationPhase::RewritingRows
929 | PersistedSchemaMigrationPhase::RebuildingIndexes
930 | PersistedSchemaMigrationPhase::FinalValidation
931 | PersistedSchemaMigrationPhase::Publishing
932 | PersistedSchemaMigrationPhase::Rejected => Err(InternalError::schema_migration(
933 SchemaMigrationCode::MigrationInProgress,
934 )),
935 PersistedSchemaMigrationPhase::Applied | PersistedSchemaMigrationPhase::Aborted => {
936 Ok(false)
937 }
938 }
939}
940
941#[cfg(feature = "migration")]
942fn adopt_entity_source_lineage<C: CanisterKind>(
943 db: &Db<C>,
944 proposal: &SchemaProposal,
945 expected_database: TargetDatabaseIdentity,
946 expected_head: &ExpectedAcceptedHead,
947) -> Result<SchemaMigrationStatusPage, InternalError> {
948 if proposal.migration().is_some() || proposal.target_database() != expected_database {
949 return Err(InternalError::schema_migration(
950 SchemaMigrationCode::PlanChanged,
951 ));
952 }
953 let proposal_digest = proposal
954 .digest()
955 .map_err(|_| InternalError::store_unsupported())?;
956 let submission_key = migration_submission_key(None)?;
957 if let Some(record) = load_exact_migration_record(
958 expected_database,
959 &submission_key,
960 proposal_digest,
961 expected_head,
962 )? {
963 let replay_target = exact_migration_replay_target(db, expected_database, &record)?;
964 return schema_migration_status_for_target(db, proposal, &replay_target);
965 }
966 let target = exact_migration_target(db, expected_database, expected_head)?;
967
968 let authorities = application_authorities(db);
969 let current_bundles = load_current_application_bundles(authorities.as_slice())?;
970 let stores = existing_proposal_stores(
971 target.database_identity(),
972 authorities.as_slice(),
973 current_bundles.as_slice(),
974 );
975 let stored_before = load_entity_source_lineage_catalog()?;
976 let before = stored_before.clone().unwrap_or_default();
977 let planned = plan_entity_source_adoption(proposal, stores.as_slice(), &before)
978 .map_err(schema_migration_planning_error)?;
979 let after = lineage_after_planned(&before, planned.as_slice(), expected_head)?;
980 let receipt = SchemaChangeReceipt::new(
981 expected_database,
982 submission_key,
983 proposal_digest,
984 expected_head.clone(),
985 SchemaChangeOutcome::NoOp {
986 accepted_head: expected_head.clone(),
987 },
988 )?;
989 let record = SchemaApplicationRecord::new(receipt, Vec::new())?;
990 let operation = SchemaApplicationRecordOp::insert(&record)?;
991 let lineage = EntitySourceLineageCatalogOp::replace(stored_before.as_ref(), &after)?;
992 publish_accepted_schema_candidates_with_database_control(
993 Vec::new(),
994 vec![
995 DatabaseControlOp::SchemaApplication(operation),
996 DatabaseControlOp::EntitySourceLineage(lineage),
997 ],
998 )?;
999 schema_migration_status_for_target(db, proposal, &target)
1000}
1001
1002#[cfg(feature = "migration")]
1003#[expect(
1004 clippy::too_many_lines,
1005 reason = "one migration entry point keeps preparation and exact replay ordering visible"
1006)]
1007fn advance_metadata_schema_migration<C: CanisterKind>(
1008 db: &Db<C>,
1009 proposal: &SchemaProposal,
1010 expected_database: TargetDatabaseIdentity,
1011 expected_head: &ExpectedAcceptedHead,
1012 expected_plan: SchemaMigrationPlanDigest,
1013 acknowledged_finding_page: Option<u64>,
1014) -> Result<SchemaMigrationStatusPage, InternalError> {
1015 let plan = proposal
1016 .migration()
1017 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::MissingMigration))?;
1018 if proposal.target_database() != expected_database || plan.digest() != expected_plan {
1019 return Err(InternalError::schema_migration(
1020 SchemaMigrationCode::PlanChanged,
1021 ));
1022 }
1023 let proposal_digest = proposal
1024 .digest()
1025 .map_err(|_| InternalError::store_unsupported())?;
1026 if let Some(record) =
1027 exact_active_migration_record(proposal, expected_database, expected_head, expected_plan)?
1028 {
1029 let target = schema_application_target(db)?;
1030 validate_active_migration_target(&record, &target)?;
1031 return advance_active_schema_migration(
1032 db,
1033 proposal,
1034 &target,
1035 &record,
1036 acknowledged_finding_page,
1037 );
1038 }
1039 if acknowledged_finding_page.is_some() {
1040 return Err(InternalError::schema_migration(
1041 SchemaMigrationCode::CandidateMismatch,
1042 ));
1043 }
1044 let submission_key = migration_submission_key(Some(expected_plan))?;
1045 if let Some(record) = load_exact_migration_record(
1046 expected_database,
1047 &submission_key,
1048 proposal_digest,
1049 expected_head,
1050 )? {
1051 let replay_target = exact_migration_replay_target(db, expected_database, &record)?;
1052 return schema_migration_status_for_target(db, proposal, &replay_target);
1053 }
1054 exact_migration_target(db, expected_database, expected_head)?;
1055
1056 let authorities = application_authorities(db);
1057 let current_bundles = load_current_application_bundles(authorities.as_slice())?;
1058 let stores = existing_proposal_stores(
1059 expected_database,
1060 authorities.as_slice(),
1061 current_bundles.as_slice(),
1062 );
1063 let before = load_entity_source_lineage_catalog()?
1064 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::Unadopted))?;
1065 let planned = plan_schema_migration(proposal, stores.as_slice(), &before)
1066 .map_err(schema_migration_planning_error)?;
1067 let mut candidates = planned.candidates().to_vec();
1068 let pending = if planned.requires_physical_validation() {
1069 None
1074 } else {
1075 preflight_existing_application(
1076 authorities.as_slice(),
1077 current_bundles.as_slice(),
1078 &mut candidates,
1079 )?
1080 };
1081 if pending.is_some() {
1082 return Err(InternalError::schema_migration(
1083 SchemaMigrationCode::MigrationInProgress,
1084 ));
1085 }
1086 if candidates.is_empty() || planned.lineage().is_empty() {
1087 return Err(InternalError::schema_migration(
1088 SchemaMigrationCode::EmptyEntityVersionBump,
1089 ));
1090 }
1091 validate_database_identity_state_capacity(
1092 authorities.as_slice(),
1093 candidates.as_slice(),
1094 database_incarnation_id()?,
1095 )?;
1096 let accepted_head =
1097 accepted_head_after_candidates(authorities.as_slice(), candidates.as_slice())?;
1098 if planned.requires_physical_validation() {
1099 ensure_physical_migration_stores_are_journaled(db, &planned)?;
1100 let record = prepared_physical_schema_migration(
1101 proposal,
1102 &planned,
1103 candidates.as_slice(),
1104 expected_database,
1105 expected_head,
1106 &accepted_head,
1107 proposal_digest,
1108 expected_plan,
1109 stores.as_slice(),
1110 )?;
1111 let prior = load_schema_migration_record()?;
1115 if prior.is_some()
1116 && continue_recovery_with_failure_authority(db)
1117 .map_err(StartupRecoveryFailure::into_error)?
1118 == RecoveryProgress::Pending
1119 {
1120 let target = schema_application_target(db)?;
1121 return schema_migration_status_for_target(db, proposal, &target);
1122 }
1123 let operation = match prior.as_ref() {
1124 Some(prior) => SchemaMigrationRecordOp::replace(prior, &record)?,
1125 None => SchemaMigrationRecordOp::insert(&record)?,
1126 };
1127 publish_accepted_schema_candidates_with_database_control(
1128 Vec::new(),
1129 vec![DatabaseControlOp::SchemaMigration(operation)],
1130 )?;
1131 let target = schema_application_target(db)?;
1132 validate_active_migration_target(&record, &target)?;
1133 return active_migration_status(proposal, &target, &record);
1134 }
1135 let after = lineage_after_planned(&before, planned.lineage(), &accepted_head)?;
1136 let receipt = SchemaChangeReceipt::new(
1137 expected_database,
1138 submission_key,
1139 proposal_digest,
1140 expected_head.clone(),
1141 SchemaChangeOutcome::Applied { accepted_head },
1142 )?;
1143 let record = SchemaApplicationRecord::new(receipt, Vec::new())?;
1144 let operation = SchemaApplicationRecordOp::insert(&record)?;
1145 let lineage = EntitySourceLineageCatalogOp::replace(Some(&before), &after)?;
1146 let publications = application_publications(
1147 authorities.as_slice(),
1148 current_bundles.as_slice(),
1149 candidates.as_slice(),
1150 )?;
1151 publish_accepted_schema_candidates_with_database_control(
1152 publications,
1153 vec![
1154 DatabaseControlOp::SchemaApplication(operation),
1155 DatabaseControlOp::EntitySourceLineage(lineage),
1156 ],
1157 )?;
1158 let applied_target = schema_application_target(db)?;
1159 schema_migration_status_for_target(db, proposal, &applied_target)
1160}
1161
1162#[cfg(feature = "migration")]
1163fn ensure_physical_migration_stores_are_journaled<C: CanisterKind>(
1164 db: &Db<C>,
1165 planned: &crate::db::schema::migration_planner::PlannedSchemaMigration,
1166) -> Result<(), InternalError> {
1167 for program in planned.programs() {
1168 let store = db.store_handle(program.store_path())?;
1169 if store.storage_capabilities().recovery()
1170 != StoreRecoveryCapability::StableBasePlusJournalReplay
1171 {
1172 return Err(InternalError::schema_migration(
1173 SchemaMigrationCode::PhysicalRunnerMissing,
1174 ));
1175 }
1176 }
1177 Ok(())
1178}
1179
1180#[cfg(feature = "migration")]
1181#[expect(
1182 clippy::too_many_arguments,
1183 reason = "migration preparation binds every immutable deployment and candidate identity"
1184)]
1185fn prepared_physical_schema_migration(
1186 proposal: &SchemaProposal,
1187 planned: &crate::db::schema::migration_planner::PlannedSchemaMigration,
1188 candidates: &[CandidateSchemaRevision],
1189 database_identity: TargetDatabaseIdentity,
1190 accepted_before: &ExpectedAcceptedHead,
1191 candidate_head: &ExpectedAcceptedHead,
1192 submission_digest: SchemaProposalDigest,
1193 plan_digest: SchemaMigrationPlanDigest,
1194 stores: &[ExistingProposalStore<'_>],
1195) -> Result<SchemaMigrationRecord, InternalError> {
1196 let plan = proposal
1197 .migration()
1198 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::MissingMigration))?;
1199 let transitions = plan
1200 .transitions()
1201 .iter()
1202 .map(|transition| {
1203 PersistedSchemaMigrationTransition::try_new(
1204 transition.entity().clone(),
1205 transition.from().get(),
1206 transition
1207 .from()
1208 .get()
1209 .checked_add(1)
1210 .ok_or_else(InternalError::store_invariant)?,
1211 )
1212 })
1213 .collect::<Result<Vec<_>, InternalError>>()?;
1214 let entities = planned
1215 .lineage()
1216 .iter()
1217 .map(|entity| {
1218 PersistedSchemaMigrationEntity::try_new(
1219 entity.store(),
1220 entity.entity(),
1221 entity.digest(),
1222 )
1223 })
1224 .collect::<Result<Vec<_>, InternalError>>()?;
1225 let mut staged_indexes = Vec::new();
1226 for candidate in candidates {
1227 let store = stores
1228 .iter()
1229 .find(|store| store.path == candidate.store_path())
1230 .ok_or_else(InternalError::store_invariant)?;
1231 for (entity, snapshot) in candidate.bundle().entity_snapshots() {
1232 let before = store.bundle.entity_snapshots().get(entity);
1233 for index in snapshot
1234 .indexes()
1235 .iter()
1236 .filter(|index| {
1237 before
1238 .and_then(|before| {
1239 before
1240 .indexes()
1241 .iter()
1242 .find(|old| old.schema_id() == index.schema_id())
1243 })
1244 .is_none_or(|old| old.physical_generation() != index.physical_generation())
1245 })
1246 .chain(snapshot.candidate_indexes())
1247 {
1248 staged_indexes.push(PersistedSchemaMigrationIndex::try_new(
1249 store.identity,
1250 *entity,
1251 u64::from(index.schema_id().get()),
1252 index.physical_generation(),
1253 )?);
1254 }
1255 }
1256 }
1257 icydb_schema::compact_sort_unstable_by(&mut staged_indexes, Ord::cmp);
1258 staged_indexes.dedup();
1259 SchemaMigrationRecord::prepared(
1260 database_identity,
1261 accepted_before.clone(),
1262 candidate_head.clone(),
1263 submission_digest,
1264 plan_digest,
1265 transitions,
1266 entities,
1267 staged_indexes,
1268 )
1269}
1270
1271#[cfg(feature = "migration")]
1272#[expect(
1273 clippy::too_many_lines,
1274 reason = "the closed phase match keeps every durable migration transition and publication boundary exhaustive"
1275)]
1276fn advance_active_schema_migration<C: CanisterKind>(
1277 db: &Db<C>,
1278 proposal: &SchemaProposal,
1279 target: &SchemaApplicationTarget,
1280 record: &SchemaMigrationRecord,
1281 acknowledged_finding_page: Option<u64>,
1282) -> Result<SchemaMigrationStatusPage, InternalError> {
1283 match record.phase() {
1284 PersistedSchemaMigrationPhase::Prepared => {
1285 if acknowledged_finding_page.is_some() {
1286 return Err(InternalError::schema_migration(
1287 SchemaMigrationCode::CandidateMismatch,
1288 ));
1289 }
1290 let validating = record.transition(
1291 PersistedSchemaMigrationPhase::Validating,
1292 record.progress().clone(),
1293 )?;
1294 publish_migration_record_replacement(record, &validating)?;
1295 active_migration_status(proposal, target, &validating)
1296 }
1297 PersistedSchemaMigrationPhase::Validating => {
1298 if acknowledged_finding_page.is_some() {
1299 return Err(InternalError::schema_migration(
1300 SchemaMigrationCode::CandidateMismatch,
1301 ));
1302 }
1303 let planned = recompile_active_physical_migration(db, proposal, record)?;
1304 let page = validate_migration_page(db, &planned, record.progress())?;
1305 let (progress, staged_entries, exhausted) = page.into_parts();
1306 let phase = if progress.findings().is_empty() {
1307 if exhausted {
1308 PersistedSchemaMigrationPhase::ReadyToRewrite
1309 } else {
1310 PersistedSchemaMigrationPhase::Validating
1311 }
1312 } else {
1313 PersistedSchemaMigrationPhase::Rejected
1314 };
1315 if progress.findings().is_empty() {
1316 stage_migration_index_entries(staged_entries)?;
1320 }
1321 let advanced = record.transition(phase, progress)?;
1322 publish_migration_record_replacement(record, &advanced)?;
1323 active_migration_status(proposal, target, &advanced)
1324 }
1325 PersistedSchemaMigrationPhase::Rejected => {
1326 if acknowledged_finding_page.is_some()
1327 && acknowledged_finding_page != record.progress().finding_page()
1328 {
1329 return Err(InternalError::schema_migration(
1330 SchemaMigrationCode::CandidateMismatch,
1331 ));
1332 }
1333 active_migration_status(proposal, target, record)
1334 }
1335 PersistedSchemaMigrationPhase::ReadyToRewrite => {
1336 if acknowledged_finding_page.is_some() {
1337 return Err(InternalError::schema_migration(
1338 SchemaMigrationCode::CandidateMismatch,
1339 ));
1340 }
1341 let progress = record.progress().begin_row_phase()?;
1342 let rewriting =
1343 record.transition(PersistedSchemaMigrationPhase::RewritingRows, progress)?;
1344 publish_migration_record_replacement(record, &rewriting)?;
1345 active_migration_status(proposal, target, &rewriting)
1346 }
1347 PersistedSchemaMigrationPhase::RewritingRows => {
1348 if acknowledged_finding_page.is_some() {
1349 return Err(InternalError::schema_migration(
1350 SchemaMigrationCode::CandidateMismatch,
1351 ));
1352 }
1353 let planned = recompile_active_physical_migration(db, proposal, record)?;
1354 let page =
1355 rewrite_migration_page(db, &planned, record.progress(), record.plan_digest())?;
1356 let (progress, effects, exhausted) = page.into_parts();
1357 if progress.rows_rewritten() > progress.rows_validated()
1358 || (exhausted && progress.rows_rewritten() != progress.rows_validated())
1359 {
1360 return Err(InternalError::schema_migration(
1361 SchemaMigrationCode::ProgressCorrupt,
1362 ));
1363 }
1364 let phase = if exhausted {
1365 PersistedSchemaMigrationPhase::RebuildingIndexes
1366 } else {
1367 PersistedSchemaMigrationPhase::RewritingRows
1368 };
1369 let advanced = record.transition(phase, progress)?;
1370 let operation = SchemaMigrationRecordOp::replace(record, &advanced)?;
1371 publish_migration_rewrite_page(effects, operation)?;
1372 active_migration_status(proposal, target, &advanced)
1373 }
1374 PersistedSchemaMigrationPhase::RebuildingIndexes => {
1375 if acknowledged_finding_page.is_some() {
1376 return Err(InternalError::schema_migration(
1377 SchemaMigrationCode::CandidateMismatch,
1378 ));
1379 }
1380 let planned = recompile_active_physical_migration(db, proposal, record)?;
1381 let rebuilt = migration_derived_domain_count(db, &planned)?;
1382 let progress = record
1383 .progress()
1384 .begin_row_phase()?
1385 .with_index_progress(None, rebuilt)?;
1386 let validating =
1387 record.transition(PersistedSchemaMigrationPhase::FinalValidation, progress)?;
1388 publish_migration_record_replacement(record, &validating)?;
1389 active_migration_status(proposal, target, &validating)
1390 }
1391 PersistedSchemaMigrationPhase::FinalValidation => {
1392 if acknowledged_finding_page.is_some() {
1393 return Err(InternalError::schema_migration(
1394 SchemaMigrationCode::CandidateMismatch,
1395 ));
1396 }
1397 let planned = recompile_active_physical_migration(db, proposal, record)?;
1398 let page = final_validate_migration_page(db, &planned, record.progress())?;
1399 let (progress, exhausted) = page.into_parts();
1400 let phase = if exhausted {
1401 PersistedSchemaMigrationPhase::Publishing
1402 } else {
1403 PersistedSchemaMigrationPhase::FinalValidation
1404 };
1405 let advanced = record.transition(phase, progress)?;
1406 publish_migration_record_replacement(record, &advanced)?;
1407 active_migration_status(proposal, target, &advanced)
1408 }
1409 PersistedSchemaMigrationPhase::Publishing => {
1410 if acknowledged_finding_page.is_some() {
1411 return Err(InternalError::schema_migration(
1412 SchemaMigrationCode::CandidateMismatch,
1413 ));
1414 }
1415 publish_completed_physical_migration(db, proposal, record)
1416 }
1417 PersistedSchemaMigrationPhase::Applied | PersistedSchemaMigrationPhase::Aborted => {
1418 if acknowledged_finding_page.is_some() {
1419 return Err(InternalError::schema_migration(
1420 SchemaMigrationCode::CandidateMismatch,
1421 ));
1422 }
1423 active_migration_status(proposal, target, record)
1424 }
1425 }
1426}
1427
1428#[cfg(feature = "migration")]
1429fn recompile_active_physical_migration<C: CanisterKind>(
1430 db: &Db<C>,
1431 proposal: &SchemaProposal,
1432 record: &SchemaMigrationRecord,
1433) -> Result<crate::db::schema::migration_planner::PlannedSchemaMigration, InternalError> {
1434 let authorities = application_authorities(db);
1435 let current_bundles = load_current_application_bundles(authorities.as_slice())?;
1436 let stores = existing_proposal_stores(
1437 record.database_identity(),
1438 authorities.as_slice(),
1439 current_bundles.as_slice(),
1440 );
1441 let lineage = load_entity_source_lineage_catalog()?
1442 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::Unadopted))?;
1443 let planned = plan_schema_migration(proposal, stores.as_slice(), &lineage)
1444 .map_err(schema_migration_planning_error)?;
1445 if !planned.requires_physical_validation() {
1446 return Err(InternalError::schema_migration(
1447 SchemaMigrationCode::PlanChanged,
1448 ));
1449 }
1450 let candidate_head =
1451 accepted_head_after_candidates(authorities.as_slice(), planned.candidates())?;
1452 if &candidate_head != record.candidate_head() {
1453 return Err(InternalError::schema_migration(
1454 SchemaMigrationCode::CandidateMismatch,
1455 ));
1456 }
1457 Ok(planned)
1458}
1459
1460#[cfg(feature = "migration")]
1461fn publish_completed_physical_migration<C: CanisterKind>(
1462 db: &Db<C>,
1463 proposal: &SchemaProposal,
1464 record: &SchemaMigrationRecord,
1465) -> Result<SchemaMigrationStatusPage, InternalError> {
1466 let target = schema_application_target(db)?;
1467 validate_active_migration_target(record, &target)?;
1468 let planned = recompile_active_physical_migration(db, proposal, record)?;
1469 let authorities = application_authorities(db);
1470 let current_bundles = load_current_application_bundles(authorities.as_slice())?;
1471 let candidate_head =
1472 accepted_head_after_candidates(authorities.as_slice(), planned.candidates())?;
1473 if &candidate_head != record.candidate_head() {
1474 return Err(InternalError::schema_migration(
1475 SchemaMigrationCode::PublicationRaceLost,
1476 ));
1477 }
1478 let lineage_before = load_entity_source_lineage_catalog()?
1479 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::Unadopted))?;
1480 let lineage_after =
1481 lineage_after_planned(&lineage_before, planned.lineage(), record.candidate_head())?;
1482 let receipt = SchemaChangeReceipt::new(
1483 record.database_identity(),
1484 migration_submission_key(Some(record.plan_digest()))?,
1485 record.submission_digest(),
1486 record.accepted_before().clone(),
1487 SchemaChangeOutcome::Applied {
1488 accepted_head: record.candidate_head().clone(),
1489 },
1490 )?;
1491 let application = SchemaApplicationRecord::new(receipt, Vec::new())?;
1492 let application = SchemaApplicationRecordOp::insert(&application)?;
1493 let lineage = EntitySourceLineageCatalogOp::replace(Some(&lineage_before), &lineage_after)?;
1494 let applied = record.transition(
1495 PersistedSchemaMigrationPhase::Applied,
1496 record.progress().clone(),
1497 )?;
1498 let migration = SchemaMigrationRecordOp::replace(record, &applied)?;
1499 let publications = application_publications(
1500 authorities.as_slice(),
1501 current_bundles.as_slice(),
1502 planned.candidates(),
1503 )?;
1504 publish_accepted_schema_candidates_with_database_control(
1505 publications,
1506 vec![
1507 DatabaseControlOp::SchemaApplication(application),
1508 DatabaseControlOp::EntitySourceLineage(lineage),
1509 DatabaseControlOp::SchemaMigration(migration),
1510 ],
1511 )?;
1512 db.mark_all_registered_index_stores_ready()?;
1513 let applied_target = schema_application_target(db)?;
1514 active_migration_status(proposal, &applied_target, &applied)
1515}
1516
1517#[cfg(feature = "migration")]
1518fn publish_migration_record_replacement(
1519 before: &SchemaMigrationRecord,
1520 after: &SchemaMigrationRecord,
1521) -> Result<(), InternalError> {
1522 let operation = SchemaMigrationRecordOp::replace(before, after)?;
1523 publish_accepted_schema_candidates_with_database_control(
1524 Vec::new(),
1525 vec![DatabaseControlOp::SchemaMigration(operation)],
1526 )
1527}
1528
1529#[cfg(feature = "migration")]
1530fn attach_ordinary_lineage_publication(
1531 proposal: &SchemaProposal,
1532 prior_head: &ExpectedAcceptedHead,
1533 accepted_head: &ExpectedAcceptedHead,
1534 candidates: &[CandidateSchemaRevision],
1535 operation: SchemaApplicationRecordOp,
1536) -> Result<Vec<DatabaseControlOp>, InternalError> {
1537 let mut operations = vec![DatabaseControlOp::SchemaApplication(operation)];
1538 let stored_before = load_entity_source_lineage_catalog()?;
1539 let planned = if matches!(prior_head, ExpectedAcceptedHead::Empty) {
1540 plan_initial_entity_source_lineage(proposal, candidates)
1541 .map_err(schema_migration_planning_error)?
1542 } else {
1543 Vec::new()
1544 };
1545 if planned.is_empty() && (stored_before.is_none() || prior_head == accepted_head) {
1546 return Ok(operations);
1547 }
1548 let before = stored_before.clone().unwrap_or_default();
1549 let after = lineage_after_planned(&before, planned.as_slice(), accepted_head)?;
1550 if before == after {
1551 return Ok(operations);
1552 }
1553 operations.push(DatabaseControlOp::EntitySourceLineage(
1554 EntitySourceLineageCatalogOp::replace(stored_before.as_ref(), &after)?,
1555 ));
1556 Ok(operations)
1557}
1558
1559#[cfg(feature = "migration")]
1560fn exact_migration_target<C: CanisterKind>(
1561 db: &Db<C>,
1562 expected_database: TargetDatabaseIdentity,
1563 expected_head: &ExpectedAcceptedHead,
1564) -> Result<SchemaApplicationTarget, InternalError> {
1565 let target = schema_application_target(db)?;
1566 if target.database_identity() != expected_database || target.accepted_head() != expected_head {
1567 return Err(InternalError::schema_migration(
1568 SchemaMigrationCode::StaleAcceptedHead,
1569 ));
1570 }
1571 Ok(target)
1572}
1573
1574#[cfg(feature = "migration")]
1575fn exact_migration_replay_target<C: CanisterKind>(
1576 db: &Db<C>,
1577 expected_database: TargetDatabaseIdentity,
1578 record: &SchemaApplicationRecord,
1579) -> Result<SchemaApplicationTarget, InternalError> {
1580 let target = schema_application_target(db)?;
1581 if target.database_identity() != expected_database
1582 || target.accepted_head() != migration_record_accepted_head(record)?
1583 {
1584 return Err(InternalError::schema_migration(
1585 SchemaMigrationCode::PlanChanged,
1586 ));
1587 }
1588 Ok(target)
1589}
1590
1591#[cfg(feature = "migration")]
1592fn schema_migration_status_for_target<C: CanisterKind>(
1593 db: &Db<C>,
1594 proposal: &SchemaProposal,
1595 target: &SchemaApplicationTarget,
1596) -> Result<SchemaMigrationStatusPage, InternalError> {
1597 if let Some(record) = load_schema_migration_record()?
1598 && (!record.phase().terminal()
1599 || proposal
1600 .digest()
1601 .map_err(|_| InternalError::store_unsupported())?
1602 == record.submission_digest())
1603 {
1604 validate_active_migration_deployment(proposal, &record)?;
1605 validate_active_migration_target(&record, target)?;
1606 return active_migration_status(proposal, target, &record);
1607 }
1608 let lineage = load_entity_source_lineage_catalog()?.unwrap_or_default();
1609 let plan_digest = proposal
1610 .migration()
1611 .map(icydb_schema::SchemaMigrationPlan::digest);
1612 let transitions = migration_transitions(proposal)?;
1613 let submission_key = migration_submission_key(plan_digest)?;
1614 let terminal = load_migration_record_for_status(target.database_identity(), &submission_key)?
1615 .map(|record| public_migration_receipt(&record, plan_digest))
1616 .transpose()?;
1617 let unadopted = lineage.entries().is_empty()
1618 || lineage
1619 .entries()
1620 .values()
1621 .any(|entry| matches!(entry.state(), AcceptedEntitySourceLineageState::Unadopted));
1622 let applied = current_proposal_lineage_is_applied(db, proposal, target.accepted_head())?;
1623 let phase = if unadopted {
1624 SchemaMigrationPhase::Unadopted
1625 } else if proposal.migration().is_none() {
1626 SchemaMigrationPhase::Adopted
1627 } else if applied {
1628 SchemaMigrationPhase::Applied
1629 } else {
1630 SchemaMigrationPhase::Idle
1631 };
1632 let terminal = terminal.filter(|receipt| {
1633 receipt.accepted_head() == target.accepted_head()
1634 && receipt.plan_digest() == plan_digest
1635 && matches!(
1636 phase,
1637 SchemaMigrationPhase::Adopted | SchemaMigrationPhase::Applied
1638 )
1639 });
1640 Ok(SchemaMigrationStatusPage::new(
1641 target.database_identity(),
1642 target.accepted_head().clone(),
1643 plan_digest,
1644 phase,
1645 transitions,
1646 0,
1647 0,
1648 0,
1649 Vec::new(),
1650 None,
1651 terminal,
1652 ))
1653}
1654
1655#[cfg(feature = "migration")]
1656fn exact_active_migration_record(
1657 proposal: &SchemaProposal,
1658 expected_database: TargetDatabaseIdentity,
1659 expected_head: &ExpectedAcceptedHead,
1660 expected_plan: SchemaMigrationPlanDigest,
1661) -> Result<Option<SchemaMigrationRecord>, InternalError> {
1662 let Some(record) = load_schema_migration_record()? else {
1663 return Ok(None);
1664 };
1665 if record.phase().terminal()
1668 && (record.accepted_before() != expected_head
1669 || record.plan_digest() != expected_plan
1670 || proposal
1671 .digest()
1672 .map_err(|_| InternalError::store_unsupported())?
1673 != record.submission_digest())
1674 {
1675 return Ok(None);
1676 }
1677 if record.database_identity() != expected_database
1678 || record.accepted_before() != expected_head
1679 || record.plan_digest() != expected_plan
1680 {
1681 return Err(InternalError::schema_migration(
1682 SchemaMigrationCode::PlanChanged,
1683 ));
1684 }
1685 validate_active_migration_deployment(proposal, &record)?;
1686 Ok(Some(record))
1687}
1688
1689#[cfg(feature = "migration")]
1690fn validate_active_migration_deployment(
1691 proposal: &SchemaProposal,
1692 record: &SchemaMigrationRecord,
1693) -> Result<(), InternalError> {
1694 let plan = proposal
1695 .migration()
1696 .ok_or_else(|| InternalError::schema_migration(SchemaMigrationCode::PlanChanged))?;
1697 let proposal_digest = proposal
1698 .digest()
1699 .map_err(|_| InternalError::store_unsupported())?;
1700 if proposal.target_database() != record.database_identity()
1701 || plan.digest() != record.plan_digest()
1702 || proposal_digest != record.submission_digest()
1703 {
1704 return Err(InternalError::schema_migration(
1705 SchemaMigrationCode::PlanChanged,
1706 ));
1707 }
1708 Ok(())
1709}
1710
1711#[cfg(feature = "migration")]
1712fn validate_active_migration_target(
1713 record: &SchemaMigrationRecord,
1714 target: &SchemaApplicationTarget,
1715) -> Result<(), InternalError> {
1716 let expected_head = if record.phase() == PersistedSchemaMigrationPhase::Applied {
1717 record.candidate_head()
1718 } else {
1719 record.accepted_before()
1720 };
1721 if target.database_identity() != record.database_identity()
1722 || target.accepted_head() != expected_head
1723 {
1724 return Err(InternalError::schema_migration(
1725 SchemaMigrationCode::PlanChanged,
1726 ));
1727 }
1728 Ok(())
1729}
1730
1731#[cfg(feature = "migration")]
1732fn active_migration_status(
1733 proposal: &SchemaProposal,
1734 target: &SchemaApplicationTarget,
1735 record: &SchemaMigrationRecord,
1736) -> Result<SchemaMigrationStatusPage, InternalError> {
1737 validate_active_migration_deployment(proposal, record)?;
1738 validate_active_migration_target(record, target)?;
1739 let transitions = record
1740 .transitions()
1741 .iter()
1742 .map(|transition| {
1743 SchemaMigrationEntityTransition::new(
1744 transition.entity().clone(),
1745 Some(transition.predecessor_version()),
1746 transition.target_version(),
1747 )
1748 })
1749 .collect();
1750 let findings = record
1751 .progress()
1752 .findings()
1753 .iter()
1754 .map(|finding| {
1755 let kind = match finding.kind() {
1756 PersistedSchemaMigrationFindingKind::Transform => {
1757 SchemaMigrationFindingKind::Transform
1758 }
1759 PersistedSchemaMigrationFindingKind::UniqueIndex => {
1760 SchemaMigrationFindingKind::UniqueIndex
1761 }
1762 PersistedSchemaMigrationFindingKind::Relation => {
1763 SchemaMigrationFindingKind::Relation
1764 }
1765 PersistedSchemaMigrationFindingKind::Constraint => {
1766 SchemaMigrationFindingKind::Constraint
1767 }
1768 };
1769 SchemaMigrationFinding::new(
1770 kind,
1771 finding.entity().value(),
1772 finding.primary_key().to_vec(),
1773 )
1774 })
1775 .collect();
1776 let phase = match record.phase() {
1777 PersistedSchemaMigrationPhase::Prepared => SchemaMigrationPhase::Prepared,
1778 PersistedSchemaMigrationPhase::Validating => SchemaMigrationPhase::Validating,
1779 PersistedSchemaMigrationPhase::ReadyToRewrite => SchemaMigrationPhase::ReadyToRewrite,
1780 PersistedSchemaMigrationPhase::RewritingRows => SchemaMigrationPhase::RewritingRows,
1781 PersistedSchemaMigrationPhase::RebuildingIndexes => SchemaMigrationPhase::RebuildingIndexes,
1782 PersistedSchemaMigrationPhase::FinalValidation => SchemaMigrationPhase::FinalValidation,
1783 PersistedSchemaMigrationPhase::Publishing => SchemaMigrationPhase::Publishing,
1784 PersistedSchemaMigrationPhase::Applied => SchemaMigrationPhase::Applied,
1785 PersistedSchemaMigrationPhase::Rejected => SchemaMigrationPhase::Rejected,
1786 PersistedSchemaMigrationPhase::Aborted => SchemaMigrationPhase::Aborted,
1787 };
1788 let terminal_receipt = (record.phase() == PersistedSchemaMigrationPhase::Applied).then(|| {
1789 SchemaMigrationReceipt::new(
1790 record.database_identity(),
1791 Some(record.plan_digest()),
1792 record.accepted_before().clone(),
1793 record.candidate_head().clone(),
1794 )
1795 });
1796 Ok(SchemaMigrationStatusPage::new(
1797 record.database_identity(),
1798 target.accepted_head().clone(),
1799 Some(record.plan_digest()),
1800 phase,
1801 transitions,
1802 record.progress().rows_validated(),
1803 record.progress().rows_rewritten(),
1804 record.progress().indexes_rebuilt(),
1805 findings,
1806 None,
1807 terminal_receipt,
1808 ))
1809}
1810
1811#[cfg(feature = "migration")]
1812fn load_current_application_bundles(
1813 authorities: &[StoreApplicationAuthority],
1814) -> Result<Vec<Option<AcceptedSchemaRevisionBundle>>, InternalError> {
1815 authorities
1816 .iter()
1817 .map(|authority| {
1818 authority
1819 .handle
1820 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)
1821 })
1822 .collect()
1823}
1824
1825#[cfg(feature = "migration")]
1826fn existing_proposal_stores<'a>(
1827 database_identity: TargetDatabaseIdentity,
1828 authorities: &[StoreApplicationAuthority],
1829 bundles: &'a [Option<AcceptedSchemaRevisionBundle>],
1830) -> Vec<ExistingProposalStore<'a>> {
1831 authorities
1832 .iter()
1833 .zip(bundles)
1834 .filter_map(|(authority, bundle)| {
1835 bundle.as_ref().map(|bundle| ExistingProposalStore {
1836 path: authority.path,
1837 identity: derive_store_identity(database_identity, authority),
1838 bundle,
1839 })
1840 })
1841 .collect()
1842}
1843
1844#[cfg(feature = "migration")]
1845fn lineage_after_planned(
1846 before: &AcceptedEntitySourceLineageCatalog,
1847 planned: &[PlannedEntitySourceLineage],
1848 accepted_head: &ExpectedAcceptedHead,
1849) -> Result<AcceptedEntitySourceLineageCatalog, InternalError> {
1850 let mut entries = BTreeMap::new();
1851 for (key, entry) in before.entries() {
1852 let next = match entry.state() {
1853 AcceptedEntitySourceLineageState::Unadopted => {
1854 AcceptedEntitySourceLineage::unadopted(accepted_head.clone())?
1855 }
1856 AcceptedEntitySourceLineageState::Adopted {
1857 version,
1858 source_digest,
1859 } => AcceptedEntitySourceLineage::adopted(
1860 accepted_head.clone(),
1861 *version,
1862 *source_digest,
1863 )?,
1864 };
1865 entries.insert(*key, next);
1866 }
1867 for next in planned {
1868 entries.insert(
1869 (next.store(), next.entity()),
1870 AcceptedEntitySourceLineage::adopted(
1871 accepted_head.clone(),
1872 next.version(),
1873 next.digest(),
1874 )?,
1875 );
1876 }
1877 AcceptedEntitySourceLineageCatalog::try_new(entries)
1878}
1879
1880#[cfg(feature = "migration")]
1881fn schema_migration_planning_error(error: SchemaMigrationPlanningError) -> InternalError {
1882 let reason = match error {
1883 SchemaMigrationPlanningError::Preparation(error) => return error,
1884 SchemaMigrationPlanningError::Unadopted => SchemaMigrationCode::Unadopted,
1885 SchemaMigrationPlanningError::MissingMigration => SchemaMigrationCode::MissingMigration,
1886 SchemaMigrationPlanningError::VersionGap => SchemaMigrationCode::VersionGap,
1887 SchemaMigrationPlanningError::Downgrade => SchemaMigrationCode::Downgrade,
1888 SchemaMigrationPlanningError::EmptyEntityVersionBump => {
1889 SchemaMigrationCode::EmptyEntityVersionBump
1890 }
1891 SchemaMigrationPlanningError::StaleAcceptedHead => SchemaMigrationCode::StaleAcceptedHead,
1892 SchemaMigrationPlanningError::UnknownFromObject => SchemaMigrationCode::UnknownFromObject,
1893 SchemaMigrationPlanningError::UnknownToObject => SchemaMigrationCode::UnknownToObject,
1894 SchemaMigrationPlanningError::KindMismatch => SchemaMigrationCode::KindMismatch,
1895 SchemaMigrationPlanningError::IdentityConflict => SchemaMigrationCode::IdentityConflict,
1896 SchemaMigrationPlanningError::UnexplainedSchemaDifference => {
1897 SchemaMigrationCode::UnexplainedSchemaDifference
1898 }
1899 SchemaMigrationPlanningError::UnsupportedTransform => {
1900 SchemaMigrationCode::UnsupportedTransform
1901 }
1902 SchemaMigrationPlanningError::RekeyedCatalogInvalid
1903 | SchemaMigrationPlanningError::CandidateMismatch => SchemaMigrationCode::CandidateMismatch,
1904 SchemaMigrationPlanningError::CorruptLineage => SchemaMigrationCode::ProgressCorrupt,
1905 };
1906 InternalError::schema_migration(reason)
1907}
1908
1909#[cfg(feature = "migration")]
1910fn migration_submission_key(
1911 plan_digest: Option<SchemaMigrationPlanDigest>,
1912) -> Result<SchemaSubmissionKey, InternalError> {
1913 let mut hasher = new_hash_sha256_prefixed(SCHEMA_MIGRATION_SUBMISSION_PROFILE);
1914 match plan_digest {
1915 None => write_hash_tag_u8(&mut hasher, 0),
1916 Some(digest) => {
1917 write_hash_tag_u8(&mut hasher, 1);
1918 hasher.update(digest.to_bytes());
1919 }
1920 }
1921 let digest = finalize_hash_sha256(hasher);
1922 let mut encoded = String::with_capacity(80);
1923 encoded.push_str("migration/");
1924 for byte in digest {
1925 use std::fmt::Write as _;
1926 write!(&mut encoded, "{byte:02x}").map_err(|_| InternalError::store_invariant())?;
1927 }
1928 SchemaSubmissionKey::try_new(encoded).map_err(|_| InternalError::store_invariant())
1929}
1930
1931#[cfg(feature = "migration")]
1932fn load_migration_record_for_status(
1933 database_identity: TargetDatabaseIdentity,
1934 submission_key: &SchemaSubmissionKey,
1935) -> Result<Option<SchemaApplicationRecord>, InternalError> {
1936 let record =
1937 with_schema_application_store(|store| store.load(database_identity, submission_key))?;
1938 if record
1939 .as_ref()
1940 .is_some_and(|record| record.receipt().database_identity() != database_identity)
1941 {
1942 return Err(InternalError::schema_migration(
1943 SchemaMigrationCode::ProgressCorrupt,
1944 ));
1945 }
1946 Ok(record)
1947}
1948
1949#[cfg(feature = "migration")]
1950fn load_exact_migration_record(
1951 database_identity: TargetDatabaseIdentity,
1952 submission_key: &SchemaSubmissionKey,
1953 proposal_digest: SchemaProposalDigest,
1954 prior_head: &ExpectedAcceptedHead,
1955) -> Result<Option<SchemaApplicationRecord>, InternalError> {
1956 let Some(record) =
1957 with_schema_application_store(|store| store.load(database_identity, submission_key))?
1958 else {
1959 return Ok(None);
1960 };
1961 if !record.receipt().is_exact_submission(
1962 database_identity,
1963 submission_key,
1964 proposal_digest,
1965 prior_head,
1966 ) {
1967 return Err(InternalError::schema_migration(
1968 SchemaMigrationCode::PlanChanged,
1969 ));
1970 }
1971 Ok(Some(record))
1972}
1973
1974#[cfg(feature = "migration")]
1975fn public_migration_receipt(
1976 record: &SchemaApplicationRecord,
1977 plan_digest: Option<SchemaMigrationPlanDigest>,
1978) -> Result<SchemaMigrationReceipt, InternalError> {
1979 let accepted_head = migration_record_accepted_head(record)?.clone();
1980 Ok(SchemaMigrationReceipt::new(
1981 record.receipt().database_identity(),
1982 plan_digest,
1983 record.receipt().prior_head().clone(),
1984 accepted_head,
1985 ))
1986}
1987
1988#[cfg(feature = "migration")]
1989fn migration_record_accepted_head(
1990 record: &SchemaApplicationRecord,
1991) -> Result<&ExpectedAcceptedHead, InternalError> {
1992 match record.receipt().outcome() {
1993 SchemaChangeOutcome::NoOp { accepted_head }
1994 | SchemaChangeOutcome::Applied { accepted_head } => Ok(accepted_head),
1995 SchemaChangeOutcome::Pending { .. } | SchemaChangeOutcome::Aborted { .. } => Err(
1996 InternalError::schema_migration(SchemaMigrationCode::ProgressCorrupt),
1997 ),
1998 }
1999}
2000
2001#[cfg(feature = "migration")]
2002fn migration_transitions(
2003 proposal: &SchemaProposal,
2004) -> Result<Vec<SchemaMigrationEntityTransition>, InternalError> {
2005 if let Some(plan) = proposal.migration() {
2006 return plan
2007 .transitions()
2008 .iter()
2009 .map(|transition| {
2010 let target = proposal_entity(proposal, transition.entity())?;
2011 Ok(SchemaMigrationEntityTransition::new(
2012 transition.entity().clone(),
2013 Some(transition.from().get()),
2014 target.version().get(),
2015 ))
2016 })
2017 .collect();
2018 }
2019 proposal
2020 .fragments()
2021 .iter()
2022 .flat_map(icydb_schema::SchemaFragment::entities)
2023 .map(|entity| {
2024 Ok(SchemaMigrationEntityTransition::new(
2025 entity.source_key().clone(),
2026 None,
2027 entity.version().get(),
2028 ))
2029 })
2030 .collect()
2031}
2032
2033#[cfg(feature = "migration")]
2034fn proposal_entity<'a>(
2035 proposal: &'a SchemaProposal,
2036 source: &EntitySourceKey,
2037) -> Result<&'a icydb_schema::EntityFragment, InternalError> {
2038 proposal
2039 .fragments()
2040 .iter()
2041 .flat_map(icydb_schema::SchemaFragment::entities)
2042 .find(|entity| entity.source_key() == source)
2043 .ok_or_else(InternalError::store_invariant)
2044}
2045
2046#[cfg(feature = "migration")]
2047fn current_proposal_lineage_is_applied<C: CanisterKind>(
2048 db: &Db<C>,
2049 proposal: &SchemaProposal,
2050 accepted_head: &ExpectedAcceptedHead,
2051) -> Result<bool, InternalError> {
2052 let Some(lineage) = load_entity_source_lineage_catalog()? else {
2053 return Ok(false);
2054 };
2055 let entities = proposal
2056 .fragments()
2057 .iter()
2058 .flat_map(icydb_schema::SchemaFragment::entities)
2059 .collect::<Vec<_>>();
2060 if entities.len() != lineage.entries().len() {
2061 return Ok(false);
2062 }
2063 let target = proposal.target_database();
2064 let authorities = application_authorities(db);
2065 for entity in entities {
2066 let source = entity.source_key();
2067 let digest = proposal
2068 .entity_source_digest(source)
2069 .map_err(|_| InternalError::store_invariant())?;
2070 let mut matched = false;
2071 for authority in &authorities {
2072 let store_identity = derive_store_identity(target, authority);
2073 let entity_tag = authority
2074 .handle
2075 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
2076 .and_then(|bundle| bundle.source_bindings().entity(source));
2077 let Some(entity_tag) = entity_tag else {
2078 continue;
2079 };
2080 let Some(entry) = lineage.get(store_identity, entity_tag) else {
2081 return Ok(false);
2082 };
2083 matched = entry.accepted_head() == accepted_head
2084 && matches!(
2085 entry.state(),
2086 AcceptedEntitySourceLineageState::Adopted { version, source_digest }
2087 if version.get() == entity.version().get() && *source_digest == digest
2088 );
2089 break;
2090 }
2091 if !matched {
2092 return Ok(false);
2093 }
2094 }
2095 Ok(true)
2096}
2097
2098#[cfg(feature = "migration")]
2099fn preflight_unpublished_schema_migration<C: CanisterKind>(
2100 target: &SchemaApplicationTarget,
2101 proposal: &SchemaProposal,
2102 db: &Db<C>,
2103) -> Result<(), InternalError> {
2104 let authorities = application_authorities(db);
2105 let current_bundles = authorities
2106 .iter()
2107 .map(|authority| {
2108 authority
2109 .handle
2110 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)
2111 })
2112 .collect::<Result<Vec<_>, InternalError>>()?;
2113 let stores = authorities
2114 .iter()
2115 .zip(¤t_bundles)
2116 .filter_map(|(authority, bundle)| {
2117 bundle.as_ref().map(|bundle| ExistingProposalStore {
2118 path: authority.path,
2119 identity: derive_store_identity(target.database_identity(), authority),
2120 bundle,
2121 })
2122 })
2123 .collect::<Vec<_>>();
2124 let lineage = load_entity_source_lineage_catalog()?.unwrap_or_default();
2125 let planned = plan_schema_migration(proposal, stores.as_slice(), &lineage)
2126 .map_err(schema_migration_planning_error)?;
2127 if planned.candidates().is_empty() || planned.lineage().is_empty() {
2128 return Err(InternalError::store_invariant());
2129 }
2130 for next in planned.lineage() {
2131 let current = lineage
2132 .get(next.store(), next.entity())
2133 .ok_or_else(InternalError::store_invariant)?;
2134 let AcceptedEntitySourceLineageState::Adopted {
2135 version,
2136 source_digest,
2137 } = current.state()
2138 else {
2139 return Err(InternalError::store_invariant());
2140 };
2141 let expected_version = version
2142 .get()
2143 .checked_add(1)
2144 .ok_or_else(InternalError::store_invariant)?;
2145 if next.version().get() != expected_version || next.digest() == *source_digest {
2146 return Err(InternalError::store_invariant());
2147 }
2148 }
2149 Ok(())
2150}
2151
2152fn lower_application_candidates<const ALLOW_REMOVALS: bool>(
2153 target: &SchemaApplicationTarget,
2154 proposal: &SchemaProposal,
2155 authorities: &[StoreApplicationAuthority],
2156) -> Result<LoweredApplication, InternalError> {
2157 let current_bundles = authorities
2158 .iter()
2159 .map(|authority| {
2160 authority
2161 .handle
2162 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)
2163 })
2164 .collect::<Result<Vec<_>, InternalError>>()?;
2165 let initial_application = matches!(target.accepted_head(), ExpectedAcceptedHead::Empty);
2166 let mut candidates = match target.accepted_head() {
2167 ExpectedAcceptedHead::Empty => {
2168 let stores = authorities
2169 .iter()
2170 .map(|authority| ProposalStoreTarget {
2171 path: authority.path,
2172 identity: derive_store_identity(target.database_identity(), authority),
2173 })
2174 .collect::<Vec<_>>();
2175 let candidates = lower_initial_schema_proposal(proposal, stores.as_slice())?;
2176 #[cfg(feature = "migration")]
2177 {
2178 let planned = plan_initial_entity_source_lineage(proposal, &candidates)
2179 .map_err(schema_migration_planning_error)?;
2180 if planned.len()
2181 != proposal
2182 .fragments()
2183 .iter()
2184 .map(|fragment| fragment.entities().len())
2185 .sum::<usize>()
2186 {
2187 return Err(InternalError::store_invariant());
2188 }
2189 }
2190 candidates
2191 }
2192 ExpectedAcceptedHead::Exact { .. }
2193 if proposal.fragments().is_empty() && proposal.removals().is_empty() =>
2194 {
2195 Vec::new()
2196 }
2197 ExpectedAcceptedHead::Exact { .. } => {
2198 let stores = authorities
2199 .iter()
2200 .zip(¤t_bundles)
2201 .filter_map(|(authority, bundle)| {
2202 bundle.as_ref().map(|bundle| ExistingProposalStore {
2203 path: authority.path,
2204 identity: derive_store_identity(target.database_identity(), authority),
2205 bundle,
2206 })
2207 })
2208 .collect::<Vec<_>>();
2209 if ALLOW_REMOVALS {
2210 lower_existing_schema_proposal(proposal, stores.as_slice())?
2211 } else {
2212 lower_generated_existing_schema_proposal(proposal, stores.as_slice())?
2213 }
2214 }
2215 };
2216 let pending = if initial_application {
2217 preflight_initial_application(authorities, &candidates)?;
2218 None
2219 } else {
2220 preflight_existing_application(authorities, ¤t_bundles, &mut candidates)?
2221 };
2222 Ok(LoweredApplication {
2223 current_bundles,
2224 candidates,
2225 pending,
2226 })
2227}
2228
2229fn validate_database_identity_state_capacity(
2230 authorities: &[StoreApplicationAuthority],
2231 candidates: &[CandidateSchemaRevision],
2232 incarnation: crate::db::integrity::DatabaseIncarnationId,
2233) -> Result<(), InternalError> {
2234 let mut total = 0usize;
2235 for authority in authorities {
2236 let count = match candidates
2237 .iter()
2238 .find(|candidate| candidate.store_path() == authority.path)
2239 {
2240 Some(candidate) => authority.handle.with_schema(|store| {
2241 store.projected_identity_state_count(incarnation, candidate)
2242 })?,
2243 None => authority
2244 .handle
2245 .with_schema(|store| store.identity_state_inventory_for_integrity(incarnation))?
2246 .len(),
2247 };
2248 total = include_identity_state_count(total, count)?;
2249 }
2250 Ok(())
2251}
2252
2253fn include_identity_state_count(total: usize, count: usize) -> Result<usize, InternalError> {
2254 let total = total
2255 .checked_add(count)
2256 .ok_or_else(InternalError::identity_state_capacity_exhausted)?;
2257 if total > MAX_IDENTITY_STATE_RECORDS_PER_DATABASE {
2258 return Err(InternalError::identity_state_capacity_exhausted());
2259 }
2260 Ok(total)
2261}
2262
2263fn preflight_initial_application(
2264 authorities: &[StoreApplicationAuthority],
2265 candidates: &[crate::db::schema::CandidateSchemaRevision],
2266) -> Result<(), InternalError> {
2267 for candidate in candidates {
2268 let authority = authorities
2269 .iter()
2270 .find(|authority| authority.path == candidate.store_path())
2271 .ok_or_else(InternalError::store_invariant)?;
2272 if authority.handle.with_data(DataStore::len) != 0
2273 || authority.handle.index_state() != IndexState::Ready
2274 || !authority.handle.with_index(IndexStore::is_empty)
2275 {
2276 return Err(InternalError::store_unsupported());
2277 }
2278 }
2279 Ok(())
2280}
2281
2282fn preflight_existing_application(
2289 authorities: &[StoreApplicationAuthority],
2290 current_bundles: &[Option<crate::db::schema::AcceptedSchemaRevisionBundle>],
2291 candidates: &mut [CandidateSchemaRevision],
2292) -> Result<Option<PendingGeneratedRowLocalConstraint>, InternalError> {
2293 require_empty_physical_entity_removal(authorities, current_bundles, candidates)?;
2294 require_empty_physical_field_removals(authorities, current_bundles, candidates)?;
2295 require_empty_physical_index_removals(authorities, current_bundles, candidates)?;
2296 require_empty_physical_relation_removals(authorities, current_bundles, candidates)?;
2297 let proofs = generated_row_local_constraint_proofs(authorities, current_bundles, candidates)?;
2298 if proofs
2299 .iter()
2300 .filter(|proof| proof.historical_rows != 0)
2301 .count()
2302 > 1
2303 {
2304 return Err(InternalError::store_unsupported());
2305 }
2306
2307 let mut pending = None;
2308 for candidate_index in 0..candidates.len() {
2309 let candidate = candidates
2310 .get(candidate_index)
2311 .cloned()
2312 .ok_or_else(InternalError::store_invariant)?;
2313 let candidate_proofs = proofs
2314 .iter()
2315 .filter(|proof| proof.candidate_index == candidate_index)
2316 .collect::<Vec<_>>();
2317 if candidate_proofs.is_empty() {
2318 continue;
2319 }
2320
2321 let mut snapshots = candidate.bundle().entity_snapshots().clone();
2322 for proof in candidate_proofs {
2323 let mut promote = true;
2324 if proof.historical_rows != 0 {
2325 match validate_unpublished_row_local_candidate_bounded(
2326 proof.store,
2327 proof.store_path,
2328 proof.entity_tag,
2329 proof.entity_path.as_str(),
2330 &candidate,
2331 proof.constraint_id,
2332 )? {
2333 UnpublishedRowLocalValidation::Complete { .. } => {}
2334 UnpublishedRowLocalValidation::Incomplete => {
2335 if proof.store.storage_capabilities().recovery()
2336 != StoreRecoveryCapability::StableBasePlusJournalReplay
2337 || pending.is_some()
2338 {
2339 return Err(InternalError::store_unsupported());
2340 }
2341 pending = Some(PendingGeneratedRowLocalConstraint {
2342 proof: (*proof).clone(),
2343 });
2344 promote = false;
2345 }
2346 }
2347 }
2348 if !promote {
2349 continue;
2350 }
2351 let snapshot = snapshots
2352 .get(&proof.entity_tag)
2353 .cloned()
2354 .ok_or_else(InternalError::store_invariant)?;
2355 let catalog = snapshot
2356 .constraint_catalog()
2357 .clone()
2358 .with_directly_validated_activation(proof.constraint_id)
2359 .map_err(|_| InternalError::store_invariant())?;
2360 snapshots.insert(proof.entity_tag, snapshot.with_constraint_catalog(catalog));
2361 }
2362 let bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
2363 candidate.revision(),
2364 candidate.bundle().store_path(),
2365 candidate.bundle().enum_catalog().clone(),
2366 candidate.bundle().composite_catalog().clone(),
2367 candidate.bundle().source_bindings().clone(),
2368 snapshots,
2369 )?;
2370 candidates[candidate_index] = CandidateSchemaRevision::new(bundle)?;
2371 }
2372 Ok(pending)
2373}
2374
2375fn require_empty_physical_entity_removal(
2382 authorities: &[StoreApplicationAuthority],
2383 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2384 candidates: &[CandidateSchemaRevision],
2385) -> Result<(), InternalError> {
2386 let mut removed_entity = None;
2387 for candidate in candidates {
2388 let (position, source_authority) = authorities
2389 .iter()
2390 .enumerate()
2391 .find(|(_, authority)| authority.path == candidate.store_path())
2392 .ok_or_else(InternalError::store_invariant)?;
2393 let current = current_bundles
2394 .get(position)
2395 .and_then(Option::as_ref)
2396 .ok_or_else(InternalError::store_invariant)?;
2397 let removed = current
2398 .entity_snapshots()
2399 .iter()
2400 .filter(|(entity_tag, _)| {
2401 !candidate
2402 .bundle()
2403 .entity_snapshots()
2404 .contains_key(entity_tag)
2405 })
2406 .collect::<Vec<_>>();
2407 if removed.is_empty() {
2408 continue;
2409 }
2410 let [(entity_tag, snapshot)] = removed.as_slice() else {
2411 return Err(InternalError::store_unsupported());
2412 };
2413 let entity_tag = **entity_tag;
2414 let snapshot = *snapshot;
2415 if removed_entity.is_some()
2416 || current.entity_snapshots().len()
2417 != candidate
2418 .bundle()
2419 .entity_snapshots()
2420 .len()
2421 .saturating_add(1)
2422 {
2423 return Err(InternalError::store_unsupported());
2424 }
2425 require_exact_empty_entity(source_authority.handle, entity_tag)?;
2426 source_authority
2427 .handle
2428 .with_index(|store| prove_empty_user_index_domain(store, entity_tag))
2429 .map_err(StagedUserIndexDomainError::into_internal_error)?;
2430 for relation in snapshot.relations() {
2431 let target_store = accepted_entity_store_for_path(
2432 authorities,
2433 current_bundles,
2434 relation.target_path(),
2435 )?;
2436 target_store.with_index(|store| {
2437 prove_empty_reverse_relation_domain(store, entity_tag, snapshot, relation)
2438 })?;
2439 }
2440 removed_entity = Some(snapshot.entity_path());
2441 }
2442
2443 let Some(removed_path) = removed_entity else {
2444 return Ok(());
2445 };
2446 for (position, authority) in authorities.iter().enumerate() {
2447 let after = candidates
2448 .iter()
2449 .find(|candidate| candidate.store_path() == authority.path)
2450 .map(CandidateSchemaRevision::bundle)
2451 .or_else(|| current_bundles.get(position).and_then(Option::as_ref));
2452 let Some(after) = after else {
2453 continue;
2454 };
2455 if after
2456 .entity_snapshots()
2457 .values()
2458 .flat_map(crate::db::schema::PersistedSchemaSnapshot::relations)
2459 .any(|relation| relation.target_path() == removed_path)
2460 {
2461 return Err(InternalError::store_unsupported());
2462 }
2463 }
2464 Ok(())
2465}
2466
2467fn require_empty_physical_relation_removals(
2470 authorities: &[StoreApplicationAuthority],
2471 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2472 candidates: &[CandidateSchemaRevision],
2473) -> Result<(), InternalError> {
2474 for candidate in candidates {
2475 let (position, source_authority) = authorities
2476 .iter()
2477 .enumerate()
2478 .find(|(_, authority)| authority.path == candidate.store_path())
2479 .ok_or_else(InternalError::store_invariant)?;
2480 let current = current_bundles
2481 .get(position)
2482 .and_then(Option::as_ref)
2483 .ok_or_else(InternalError::store_invariant)?;
2484 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2485 let before = current
2486 .entity_snapshots()
2487 .get(entity_tag)
2488 .ok_or_else(InternalError::store_invariant)?;
2489 let removed = before
2490 .relations()
2491 .iter()
2492 .filter(|relation| {
2493 !after
2494 .relations()
2495 .iter()
2496 .any(|candidate| candidate.id() == relation.id())
2497 })
2498 .collect::<Vec<_>>();
2499 if removed.is_empty() {
2500 continue;
2501 }
2502 let added = after.relations().iter().any(|relation| {
2503 !before
2504 .relations()
2505 .iter()
2506 .any(|accepted| accepted.id() == relation.id())
2507 });
2508 let [removed] = removed.as_slice() else {
2509 return Err(InternalError::store_unsupported());
2510 };
2511 if added || before.relations().len() != after.relations().len().saturating_add(1) {
2512 return Err(InternalError::store_unsupported());
2513 }
2514 require_exact_empty_entity(source_authority.handle, *entity_tag)?;
2515 let target_store = accepted_entity_store_for_path(
2516 authorities,
2517 current_bundles,
2518 removed.target_path(),
2519 )?;
2520 target_store.with_index(|store| {
2521 prove_empty_reverse_relation_domain(store, *entity_tag, before, removed)
2522 })?;
2523 }
2524 }
2525 Ok(())
2526}
2527
2528fn accepted_entity_store_for_path(
2529 authorities: &[StoreApplicationAuthority],
2530 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2531 entity_path: &str,
2532) -> Result<StoreHandle, InternalError> {
2533 let mut resolved = None;
2534 for (position, bundle) in current_bundles.iter().enumerate() {
2535 let Some(bundle) = bundle else {
2536 continue;
2537 };
2538 if !bundle
2539 .entity_snapshots()
2540 .values()
2541 .any(|snapshot| snapshot.entity_path() == entity_path)
2542 {
2543 continue;
2544 }
2545 if resolved.is_some() {
2546 return Err(InternalError::store_invariant());
2547 }
2548 resolved = authorities.get(position).map(|authority| authority.handle);
2549 }
2550 resolved.ok_or_else(InternalError::store_unsupported)
2551}
2552
2553fn require_empty_physical_index_removals(
2557 authorities: &[StoreApplicationAuthority],
2558 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2559 candidates: &[CandidateSchemaRevision],
2560) -> Result<(), InternalError> {
2561 for candidate in candidates {
2562 let (position, authority) = authorities
2563 .iter()
2564 .enumerate()
2565 .find(|(_, authority)| authority.path == candidate.store_path())
2566 .ok_or_else(InternalError::store_invariant)?;
2567 let current = current_bundles
2568 .get(position)
2569 .and_then(Option::as_ref)
2570 .ok_or_else(InternalError::store_invariant)?;
2571 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2572 let before = current
2573 .entity_snapshots()
2574 .get(entity_tag)
2575 .ok_or_else(InternalError::store_invariant)?;
2576 if before.indexes().len() == after.indexes().len() {
2577 continue;
2578 }
2579 if before.indexes().len() != after.indexes().len().saturating_add(1) {
2580 return Err(InternalError::store_unsupported());
2581 }
2582 require_exact_empty_entity(authority.handle, *entity_tag)?;
2583 authority
2584 .handle
2585 .with_index(|store| prove_empty_user_index_domain(store, *entity_tag))
2586 .map_err(StagedUserIndexDomainError::into_internal_error)?;
2587 }
2588 }
2589 Ok(())
2590}
2591
2592fn require_empty_physical_field_removals(
2595 authorities: &[StoreApplicationAuthority],
2596 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2597 candidates: &[CandidateSchemaRevision],
2598) -> Result<(), InternalError> {
2599 for candidate in candidates {
2600 let (position, authority) = authorities
2601 .iter()
2602 .enumerate()
2603 .find(|(_, authority)| authority.path == candidate.store_path())
2604 .ok_or_else(InternalError::store_invariant)?;
2605 let current = current_bundles
2606 .get(position)
2607 .and_then(Option::as_ref)
2608 .ok_or_else(InternalError::store_invariant)?;
2609 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2610 let before = current
2611 .entity_snapshots()
2612 .get(entity_tag)
2613 .ok_or_else(InternalError::store_invariant)?;
2614 if before.row_layout() == after.row_layout() {
2615 continue;
2616 }
2617 if before.fields().len() != after.fields().len().saturating_add(1) {
2618 return Err(InternalError::store_unsupported());
2619 }
2620 require_exact_empty_entity(authority.handle, *entity_tag)?;
2621 }
2622 }
2623 Ok(())
2624}
2625
2626fn require_exact_empty_entity(
2630 store: StoreHandle,
2631 entity_tag: EntityTag,
2632) -> Result<(), InternalError> {
2633 require_exact_empty_entity_count(store.exact_entity_count(entity_tag))
2634}
2635
2636fn require_exact_empty_entity_count(count: Option<u64>) -> Result<(), InternalError> {
2637 let count = count.ok_or_else(InternalError::store_corruption)?;
2638 if count != 0 {
2639 return Err(InternalError::store_unsupported());
2640 }
2641
2642 Ok(())
2643}
2644
2645fn generated_row_local_constraint_proofs(
2646 authorities: &[StoreApplicationAuthority],
2647 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2648 candidates: &[CandidateSchemaRevision],
2649) -> Result<Vec<DirectGeneratedRowLocalProof>, InternalError> {
2650 let mut proofs = Vec::new();
2651 for (candidate_index, candidate) in candidates.iter().enumerate() {
2652 let (position, authority) = authorities
2653 .iter()
2654 .enumerate()
2655 .find(|(_, authority)| authority.path == candidate.store_path())
2656 .ok_or_else(InternalError::store_invariant)?;
2657 let current = current_bundles
2658 .get(position)
2659 .and_then(Option::as_ref)
2660 .ok_or_else(InternalError::store_invariant)?;
2661 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2662 let before = current
2663 .entity_snapshots()
2664 .get(entity_tag)
2665 .ok_or_else(InternalError::store_invariant)?;
2666 for constraint_id in added_generated_row_local_activations(before, after) {
2667 let historical_rows = authority
2668 .handle
2669 .exact_entity_count(*entity_tag)
2670 .ok_or_else(InternalError::store_corruption)?;
2671 proofs.push(DirectGeneratedRowLocalProof {
2672 candidate_index,
2673 store: authority.handle,
2674 store_path: authority.path,
2675 entity_tag: *entity_tag,
2676 entity_path: after.entity_path().to_string(),
2677 constraint_id,
2678 historical_rows,
2679 });
2680 }
2681 }
2682 }
2683 Ok(proofs)
2684}
2685
2686fn added_generated_row_local_activations(
2687 before: &crate::db::schema::PersistedSchemaSnapshot,
2688 after: &crate::db::schema::PersistedSchemaSnapshot,
2689) -> Vec<ConstraintId> {
2690 after
2691 .constraint_activations()
2692 .iter()
2693 .filter(|candidate| {
2694 candidate.origin() == ConstraintOrigin::Generated
2695 && matches!(
2696 candidate.kind(),
2697 ConstraintActivationKind::Check { .. }
2698 | ConstraintActivationKind::TargetedRule { .. }
2699 )
2700 && !before
2701 .constraint_activations()
2702 .iter()
2703 .any(|accepted| accepted.id() == candidate.id())
2704 })
2705 .map(crate::db::schema::ConstraintActivationSnapshot::id)
2706 .collect()
2707}
2708
2709fn final_candidates_for_pending_row_local_constraint(
2710 candidates: &[CandidateSchemaRevision],
2711 pending: &PendingGeneratedRowLocalConstraint,
2712) -> Result<Vec<CandidateSchemaRevision>, InternalError> {
2713 let mut final_candidates = candidates.to_vec();
2714 let candidate = final_candidates
2715 .get(pending.proof.candidate_index)
2716 .cloned()
2717 .ok_or_else(InternalError::store_invariant)?;
2718 if candidate.store_path() != pending.proof.store_path {
2719 return Err(InternalError::store_invariant());
2720 }
2721 let mut snapshots = candidate.bundle().entity_snapshots().clone();
2722 let snapshot = snapshots
2723 .get(&pending.proof.entity_tag)
2724 .cloned()
2725 .ok_or_else(InternalError::store_invariant)?;
2726 let catalog = snapshot
2727 .constraint_catalog()
2728 .clone()
2729 .with_directly_validated_activation(pending.proof.constraint_id)
2730 .map_err(|_| InternalError::store_invariant())?;
2731 snapshots.insert(
2732 pending.proof.entity_tag,
2733 snapshot.with_constraint_catalog(catalog),
2734 );
2735 let final_revision = candidate
2736 .revision()
2737 .checked_next()
2738 .and_then(AcceptedSchemaRevision::checked_next)
2739 .ok_or_else(InternalError::store_unsupported)?;
2740 let bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
2741 final_revision,
2742 candidate.bundle().store_path(),
2743 candidate.bundle().enum_catalog().clone(),
2744 candidate.bundle().composite_catalog().clone(),
2745 candidate.bundle().source_bindings().clone(),
2746 snapshots,
2747 )?;
2748 final_candidates[pending.proof.candidate_index] = CandidateSchemaRevision::new(bundle)?;
2749 Ok(final_candidates)
2750}
2751
2752fn schema_change_progress_status(
2753 snapshot: &crate::db::schema::PersistedSchemaSnapshot,
2754 entity_tag: EntityTag,
2755 constraint_id: ConstraintId,
2756 progress: ConstraintValidationProgress,
2757) -> Result<SchemaChangeProgressStatus, InternalError> {
2758 match progress {
2759 ConstraintValidationProgress::Started => Ok(SchemaChangeProgressStatus::Started),
2760 ConstraintValidationProgress::Advanced {
2761 phase,
2762 rows_scanned,
2763 } => Ok(SchemaChangeProgressStatus::Advanced {
2764 phase: schema_change_validation_phase(phase),
2765 rows_scanned,
2766 }),
2767 ConstraintValidationProgress::Findings {
2768 receipt,
2769 phase,
2770 rows_scanned,
2771 } => {
2772 let activation = snapshot
2773 .constraint_catalog()
2774 .activation(constraint_id)
2775 .ok_or_else(InternalError::store_corruption)?;
2776 let fingerprint =
2777 crate::db::schema::accepted_schema_cache_fingerprint_for_persisted_snapshot(
2778 snapshot,
2779 )?;
2780 let findings = receipt
2781 .findings()
2782 .iter()
2783 .map(|finding| {
2784 constraint_validation_finding_output(
2785 fingerprint,
2786 entity_tag,
2787 activation,
2788 finding,
2789 )
2790 })
2791 .collect::<Result<Vec<_>, InternalError>>()?;
2792 Ok(SchemaChangeProgressStatus::Findings {
2793 phase: schema_change_validation_phase(phase),
2794 rows_scanned,
2795 page_sequence: receipt.page_sequence(),
2796 findings,
2797 })
2798 }
2799 ConstraintValidationProgress::Restarted { rows_scanned } => {
2800 Ok(SchemaChangeProgressStatus::Restarted { rows_scanned })
2801 }
2802 ConstraintValidationProgress::Promoted { .. } => Ok(SchemaChangeProgressStatus::Applied),
2803 }
2804}
2805
2806const fn schema_change_validation_phase(
2807 phase: ConstraintValidationPhase,
2808) -> SchemaChangeValidationPhase {
2809 match phase {
2810 ConstraintValidationPhase::Forward => SchemaChangeValidationPhase::Forward,
2811 ConstraintValidationPhase::Verify => SchemaChangeValidationPhase::Verify,
2812 }
2813}
2814
2815fn finalize_schema_application<C: CanisterKind>(
2816 db: &Db<C>,
2817 record: &SchemaApplicationRecord,
2818 candidate_head: &ExpectedAcceptedHead,
2819 status: SchemaChangeProgressStatus,
2820) -> Result<SchemaChangeProgress, InternalError> {
2821 if schema_application_target(db)?.accepted_head() != candidate_head {
2822 return Err(InternalError::schema_application_conflict());
2823 }
2824 let receipt = SchemaChangeReceipt::new(
2825 record.receipt().database_identity(),
2826 record.receipt().submission_key().clone(),
2827 record.receipt().proposal_digest(),
2828 record.receipt().prior_head().clone(),
2829 SchemaChangeOutcome::Applied {
2830 accepted_head: candidate_head.clone(),
2831 },
2832 )?;
2833 let terminal = SchemaApplicationRecord::new(receipt.clone(), Vec::new())?;
2834 let operation = SchemaApplicationRecordOp::replace(record, &terminal)?;
2835 publish_accepted_schema_candidates_with_application_record(Vec::new(), operation)?;
2836 Ok(SchemaChangeProgress::new(receipt, status))
2837}
2838
2839fn application_publications<'a>(
2840 authorities: &[StoreApplicationAuthority],
2841 current_bundles: &[Option<crate::db::schema::AcceptedSchemaRevisionBundle>],
2842 candidates: &'a [crate::db::schema::CandidateSchemaRevision],
2843) -> Result<Vec<AcceptedSchemaPublication<'a>>, InternalError> {
2844 candidates
2845 .iter()
2846 .map(|candidate| {
2847 let (position, authority) = authorities
2848 .iter()
2849 .enumerate()
2850 .find(|(_, authority)| authority.path == candidate.store_path())
2851 .ok_or_else(InternalError::store_invariant)?;
2852 let expected_revision = current_bundles[position].as_ref().map_or(
2853 AcceptedSchemaRevision::NONE,
2854 crate::db::schema::AcceptedSchemaRevisionBundle::revision,
2855 );
2856 Ok(AcceptedSchemaPublication::new(
2857 authority.path,
2858 authority.handle,
2859 expected_revision,
2860 candidate,
2861 ))
2862 })
2863 .collect()
2864}
2865
2866fn application_authorities<C: CanisterKind>(db: &Db<C>) -> Vec<StoreApplicationAuthority> {
2867 let mut authorities = db.with_store_registry(|registry| {
2868 registry
2869 .iter()
2870 .map(|(path, handle)| StoreApplicationAuthority { path, handle })
2871 .collect::<Vec<_>>()
2872 });
2873 icydb_schema::compact_sort_unstable_by(&mut authorities, |left, right| {
2874 left.path.cmp(right.path)
2875 });
2876 authorities
2877}
2878
2879fn accepted_head_after_candidates(
2880 authorities: &[StoreApplicationAuthority],
2881 candidates: &[crate::db::schema::CandidateSchemaRevision],
2882) -> Result<ExpectedAcceptedHead, InternalError> {
2883 let heads = authorities
2884 .iter()
2885 .map(|authority| {
2886 let candidate = candidates
2887 .iter()
2888 .find(|candidate| candidate.store_path() == authority.path);
2889 let head = match candidate {
2890 Some(candidate) => Some(AcceptedStoreHead {
2891 revision: candidate.revision().get(),
2892 fingerprint: candidate.root().fingerprint().as_bytes(),
2893 }),
2894 None => authority
2895 .handle
2896 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_root)?
2897 .map(|selection| AcceptedStoreHead {
2898 revision: selection.root().revision().get(),
2899 fingerprint: selection.root().fingerprint().as_bytes(),
2900 }),
2901 };
2902 Ok((authority.path, head))
2903 })
2904 .collect::<Result<Vec<_>, InternalError>>()?;
2905 Ok(derive_accepted_head(heads.as_slice()))
2906}
2907
2908fn derive_database_identity(
2909 incarnation: [u8; 16],
2910 stores: &[StoreApplicationAuthority],
2911) -> TargetDatabaseIdentity {
2912 let mut hasher = new_hash_sha256_prefixed(DATABASE_TARGET_FINGERPRINT_PROFILE);
2913 hasher.update(incarnation);
2914 write_hash_len_u32(&mut hasher, stores.len());
2915 for store in stores {
2916 write_store_authority(&mut hasher, store);
2917 }
2918 TargetDatabaseIdentity::from_bytes(finalize_hash_sha256(hasher))
2919}
2920
2921fn derive_store_identity(
2922 database_identity: TargetDatabaseIdentity,
2923 store: &StoreApplicationAuthority,
2924) -> TargetStoreIdentity {
2925 let mut hasher = new_hash_sha256_prefixed(STORE_TARGET_FINGERPRINT_PROFILE);
2926 hasher.update(database_identity.to_bytes());
2927 write_store_authority(&mut hasher, store);
2928 TargetStoreIdentity::from_bytes(finalize_hash_sha256(hasher))
2929}
2930
2931fn derive_accepted_head(stores: &[(&str, Option<AcceptedStoreHead>)]) -> ExpectedAcceptedHead {
2932 let Some(revision) = stores
2933 .iter()
2934 .filter_map(|(_, head)| head.map(|head| head.revision))
2935 .max()
2936 else {
2937 return ExpectedAcceptedHead::Empty;
2938 };
2939
2940 let mut hasher = new_hash_sha256_prefixed(ACCEPTED_DATABASE_HEAD_FINGERPRINT_PROFILE);
2941 write_hash_len_u32(&mut hasher, stores.len());
2942 for (path, head) in stores {
2943 write_hash_str_u32(&mut hasher, path);
2944 match head {
2945 None => write_hash_tag_u8(&mut hasher, 0),
2946 Some(head) => {
2947 write_hash_tag_u8(&mut hasher, 1);
2948 write_hash_u64(&mut hasher, head.revision);
2949 hasher.update(head.fingerprint);
2950 }
2951 }
2952 }
2953
2954 ExpectedAcceptedHead::Exact {
2955 revision,
2956 fingerprint: ExpectedSchemaFingerprint::from_bytes(finalize_hash_sha256(hasher)),
2957 }
2958}
2959
2960pub(in crate::db) fn generated_schema_reconciled(
2961 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
2962 incarnation: DatabaseIncarnationId,
2963 submission_key: &str,
2964) -> Result<(bool, ExpectedAcceptedHead), InternalError> {
2965 let submission_key = SchemaSubmissionKey::try_new(submission_key.to_string())
2966 .map_err(|_| InternalError::store_invariant())?;
2967 let (database_identity, accepted_head) = generated_schema_authority(registry, incarnation)?;
2968 let reconciled = generated_submission_is_reconciled(database_identity, &submission_key)?;
2969 Ok((reconciled, accepted_head))
2970}
2971
2972pub(in crate::db) fn generated_schema_is_reconciled(
2973 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
2974 incarnation: DatabaseIncarnationId,
2975 submission_key: &str,
2976) -> Result<bool, InternalError> {
2977 let submission_key = SchemaSubmissionKey::try_new(submission_key.to_string())
2978 .map_err(|_| InternalError::store_invariant())?;
2979 let database_identity = generated_database_identity(registry, incarnation);
2980 generated_submission_is_reconciled(database_identity, &submission_key)
2981}
2982
2983fn generated_submission_is_reconciled(
2984 database_identity: TargetDatabaseIdentity,
2985 submission_key: &SchemaSubmissionKey,
2986) -> Result<bool, InternalError> {
2987 Ok(
2988 load_schema_application_record_read_only(database_identity, submission_key)?.is_some_and(
2989 |record| match record.receipt().outcome() {
2990 SchemaChangeOutcome::NoOp { .. } | SchemaChangeOutcome::Applied { .. } => true,
2998 SchemaChangeOutcome::Pending { .. } | SchemaChangeOutcome::Aborted { .. } => false,
2999 },
3000 ),
3001 )
3002}
3003
3004pub(in crate::db) fn generated_schema_authority(
3005 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3006 incarnation: DatabaseIncarnationId,
3007) -> Result<(TargetDatabaseIdentity, ExpectedAcceptedHead), InternalError> {
3008 let database_identity = generated_database_identity(registry, incarnation);
3009 let mut stores = store_application_authorities(registry);
3010 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.path.cmp(right.path));
3011 let heads = stores
3012 .iter()
3013 .map(|store| {
3014 let head = store
3015 .handle
3016 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_root)?
3017 .map(|selection| AcceptedStoreHead {
3018 revision: selection.root().revision().get(),
3019 fingerprint: selection.root().fingerprint().as_bytes(),
3020 });
3021 Ok((store.path, head))
3022 })
3023 .collect::<Result<Vec<_>, InternalError>>()?;
3024 let accepted_head = derive_accepted_head(heads.as_slice());
3025 Ok((database_identity, accepted_head))
3026}
3027
3028fn generated_database_identity(
3029 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3030 incarnation: DatabaseIncarnationId,
3031) -> TargetDatabaseIdentity {
3032 let registry_key = std::ptr::from_ref(registry).cast::<()>() as usize;
3033 let cached = GENERATED_DATABASE_IDENTITY
3034 .with(Cell::get)
3035 .and_then(|entry| {
3036 (entry.registry == registry_key && entry.incarnation == incarnation)
3037 .then_some(entry.identity)
3038 });
3039 if let Some(identity) = cached {
3040 return identity;
3041 }
3042
3043 let mut stores = store_application_authorities(registry);
3044 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.path.cmp(right.path));
3045 let identity = derive_database_identity(incarnation.to_bytes(), stores.as_slice());
3046 GENERATED_DATABASE_IDENTITY.set(Some(GeneratedDatabaseIdentityCacheEntry {
3047 registry: registry_key,
3048 incarnation,
3049 identity,
3050 }));
3051 identity
3052}
3053
3054fn store_application_authorities(
3055 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3056) -> Vec<StoreApplicationAuthority> {
3057 registry.with(|registry| {
3058 registry
3059 .iter()
3060 .map(|(path, handle)| StoreApplicationAuthority { path, handle })
3061 .collect()
3062 })
3063}
3064
3065fn write_store_authority(hasher: &mut sha2::Sha256, store: &StoreApplicationAuthority) {
3066 write_hash_str_u32(hasher, store.path);
3067 write_storage_capabilities(hasher, store.handle);
3068 for allocation in [
3069 store.handle.data_allocation(),
3070 store.handle.index_allocation(),
3071 store.handle.schema_allocation(),
3072 store.handle.journal_allocation(),
3073 ] {
3074 write_allocation_identity(hasher, allocation);
3075 }
3076}
3077
3078fn write_storage_capabilities(hasher: &mut sha2::Sha256, store: StoreHandle) {
3079 let capabilities = store.storage_capabilities();
3080 write_hash_tag_u8(
3081 hasher,
3082 match capabilities.storage_mode() {
3083 StoreRuntimeStorageMode::Heap => 0,
3084 StoreRuntimeStorageMode::Journaled => 1,
3085 },
3086 );
3087 write_hash_tag_u8(
3088 hasher,
3089 match capabilities.allocation_identity() {
3090 StoreAllocationIdentityCapability::Present => 0,
3091 StoreAllocationIdentityCapability::Absent => 1,
3092 },
3093 );
3094 write_hash_tag_u8(
3095 hasher,
3096 match capabilities.durability() {
3097 StoreDurability::Durable => 0,
3098 StoreDurability::Volatile => 1,
3099 },
3100 );
3101 write_hash_tag_u8(
3102 hasher,
3103 match capabilities.recovery() {
3104 StoreRecoveryCapability::StableBasePlusJournalReplay => 0,
3105 StoreRecoveryCapability::None => 1,
3106 },
3107 );
3108 write_hash_tag_u8(
3109 hasher,
3110 match capabilities.commit_participation() {
3111 StoreCommitParticipation::Durable => 0,
3112 StoreCommitParticipation::LiveOnly => 1,
3113 },
3114 );
3115 write_hash_tag_u8(
3116 hasher,
3117 match capabilities.schema_metadata() {
3118 StoreSchemaMetadataCapability::LiveRebuiltMetadata => 0,
3119 StoreSchemaMetadataCapability::CanonicalStableHistoryPlusJournalTail => 1,
3120 },
3121 );
3122 write_hash_tag_u8(
3123 hasher,
3124 match capabilities.relation_source() {
3125 StoreRelationSourceCapability::DurableSource => 0,
3126 StoreRelationSourceCapability::LiveSource => 1,
3127 },
3128 );
3129 write_hash_tag_u8(
3130 hasher,
3131 match capabilities.relation_target() {
3132 StoreRelationTargetCapability::DurableTarget => 0,
3133 StoreRelationTargetCapability::VolatileTarget => 1,
3134 },
3135 );
3136}
3137
3138fn write_allocation_identity(
3139 hasher: &mut sha2::Sha256,
3140 allocation: Option<StoreAllocationIdentity>,
3141) {
3142 match allocation {
3143 None => write_hash_tag_u8(hasher, 0),
3144 Some(allocation) => {
3145 write_hash_tag_u8(hasher, 1);
3146 write_hash_tag_u8(hasher, allocation.memory_id());
3147 write_hash_str_u32(hasher, allocation.stable_key());
3148 }
3149 }
3150}
3151
3152#[cfg(test)]
3153mod tests {
3154 mod collection_relations;
3155
3156 #[cfg(feature = "migration")]
3157 mod nested_migration;
3158
3159 use super::{
3160 AcceptedSchemaPublication, AcceptedStoreHead, DirectGeneratedRowLocalProof,
3161 PendingGeneratedRowLocalConstraint, abort_schema_application,
3162 aborted_generated_row_local_candidate, accepted_head_after_candidates,
3163 application_authorities, apply_schema, continue_schema_application, derive_accepted_head,
3164 derive_schema_change_job_id, final_candidates_for_pending_row_local_constraint,
3165 generated_database_identity, include_identity_state_count, lower_existing_schema_proposal,
3166 lower_initial_schema_proposal, publish_accepted_schema_candidates_with_application_record,
3167 require_exact_empty_entity_count, schema_application_target,
3168 };
3169 use crate::{
3170 db::{
3171 DatabaseStartupState, Db, GeneratedStartupDriverStep,
3172 commit::{
3173 RecoveryProgress, continue_recovery, database_incarnation_id,
3174 forget_recovered_domain_for_tests,
3175 },
3176 data::DataStore,
3177 drive_generated_startup_recovery_page,
3178 index::IndexStore,
3179 journal::JournalTailStore,
3180 observe_generated_startup_state,
3181 registry::{
3182 StoreAllocationIdentities, StoreAllocationIdentity, StoreHandle, StoreRegistry,
3183 StoreRuntimeStorageCapabilities,
3184 },
3185 schema::{
3186 AcceptedConstraintKind, AcceptedRuleOperation, AcceptedSchemaRevisionBundle,
3187 CandidateSchemaRevision, ConstraintOrigin, ConstraintValidationJob,
3188 ExistingProposalStore, ProposalStoreTarget, SchemaApplicationRecord,
3189 SchemaApplicationRecordOp, SchemaChangeActivation, SchemaChangeJob,
3190 SchemaChangeOutcome, SchemaChangeProgressStatus, SchemaStore,
3191 cardinality_build::{
3192 CardinalityBuildAuthority, CardinalityGenerationPageOutcome,
3193 drive_cardinality_generation_page,
3194 },
3195 },
3196 },
3197 error::{ErrorClass, ErrorOrigin},
3198 testing::test_memory,
3199 traits::{CanisterKind, Path},
3200 };
3201 use ic_memory::RuntimeMemory;
3202 use ic_memory::ic_stable_structures::DefaultMemoryImpl;
3203 use icydb_schema::{
3204 ConstraintFragment, ConstraintSourceKey, DeclaredEntityVersion, EntityFragment,
3205 EntitySourceKey, EntityStoreAssignment, ExpectedAcceptedHead, ExpectedSchemaFingerprint,
3206 FieldFragment, FieldInsertPolicy, FieldSourceKey, FieldType, NamedTypeFragment,
3207 RuleSourceKey, ScalarLiteral, ScalarType, SchemaCapability, SchemaFragment, SchemaName,
3208 SchemaProposal, SchemaSubmissionKey, SourceCheckExpr, SourceCheckInstruction,
3209 SourceRuleOperation, TargetDatabaseIdentity, TargetStoreIdentity, TargetedRuleFragment,
3210 TypeSourceKey,
3211 };
3212 use std::cell::RefCell;
3213
3214 fn drive_startup_recovery_to_completion<C: CanisterKind>(db: &Db<C>) {
3215 for _ in 0..1_024 {
3216 match continue_recovery(db).expect("test startup recovery page should succeed") {
3217 RecoveryProgress::Complete => return,
3218 RecoveryProgress::Pending => {}
3219 }
3220 }
3221 panic!("test startup recovery should complete within 1,024 bounded pages");
3222 }
3223
3224 fn drive_cardinality_to_ready(store: StoreHandle) {
3225 let journal = store
3226 .journal_tail_store()
3227 .expect("cardinality fixture store should be journaled");
3228 for _ in 0..8 {
3229 let outcome = store
3230 .with_data(|data| {
3231 store.with_index(|index| {
3232 store.with_schema_mut(|schema| {
3233 drive_cardinality_generation_page(data, index, schema, |schema| {
3234 CardinalityBuildAuthority::derive(
3235 schema,
3236 database_incarnation_id()?,
3237 store.allocation_identities(),
3238 journal.with_borrow(JournalTailStore::fold_watermark)?,
3239 )
3240 })
3241 })
3242 })
3243 })
3244 .expect("bounded cardinality generation should advance");
3245 if outcome == CardinalityGenerationPageOutcome::Quiescent {
3246 return;
3247 }
3248 }
3249 panic!("cardinality generation should become Ready within eight bounded pages");
3250 }
3251
3252 #[cfg(feature = "migration")]
3253 use crate::db::schema::SchemaChangeReceipt;
3254 use crate::{
3255 db::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell},
3256 value::InputValue,
3257 };
3258 #[cfg(feature = "migration")]
3259 use icydb_schema::{
3260 EntityMigration, IndexFragment, IndexKeyFragment, RelationDeleteAction, RelationFragment,
3261 SchemaMigrationPlan, SchemaMigrationTransform, SchemaProposalDigest,
3262 };
3263
3264 fn version_one() -> DeclaredEntityVersion {
3265 DeclaredEntityVersion::try_new(1).expect("fixture version should admit")
3266 }
3267
3268 const ABORT_STORE_PATH: &str = "schema_application_tests::AbortStore";
3269 const EVOLUTION_STORE_PATH: &str = "schema_application_tests::EvolutionStore";
3270 #[cfg(feature = "migration")]
3271 const MIGRATION_STORE_PATH: &str = "schema_application_tests::MigrationStore";
3272 #[cfg(feature = "migration")]
3273 const MIGRATION_EXECUTION_STORE_PATH: &str =
3274 "schema_application_tests::MigrationExecutionStore";
3275 #[cfg(feature = "migration")]
3276 const MIGRATION_FINDING_STORE_PATH: &str = "schema_application_tests::MigrationFindingStore";
3277
3278 #[test]
3279 fn database_identity_state_capacity_combines_store_inventories_exactly() {
3280 let below = include_identity_state_count(0, 65_535)
3281 .expect("the first store inventory should remain below the database cap");
3282 let exact = include_identity_state_count(below, 1)
3283 .expect("the combined database boundary should admit");
3284 assert_eq!(exact, 65_536);
3285
3286 let error = include_identity_state_count(exact, 1)
3287 .expect_err("the next owner in another store must reject");
3288 assert_eq!(error.class(), ErrorClass::Unsupported);
3289 assert_eq!(error.origin(), ErrorOrigin::Identity);
3290 }
3291
3292 #[test]
3293 fn generated_database_identity_cache_is_bound_to_the_incarnation() {
3294 let first_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x41);
3295 let second_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x42);
3296 let first = generated_database_identity(&ABORT_REGISTRY, first_incarnation);
3297 assert_eq!(
3298 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3299 first,
3300 );
3301 let second = generated_database_identity(&ABORT_REGISTRY, second_incarnation);
3302 assert_ne!(second, first);
3303 assert_eq!(
3304 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3305 first,
3306 );
3307 }
3308
3309 #[test]
3310 fn exact_empty_entity_proof_distinguishes_corruption_from_non_empty_input() {
3311 let corrupt = require_exact_empty_entity_count(None)
3312 .expect_err("uninspectable cardinality must fail closed");
3313 assert_eq!(corrupt.class(), ErrorClass::Corruption);
3314
3315 let non_empty = require_exact_empty_entity_count(Some(1))
3316 .expect_err("non-empty cardinality must reject removal");
3317 assert_eq!(non_empty.class(), ErrorClass::Unsupported);
3318 assert!(require_exact_empty_entity_count(Some(0)).is_ok());
3319 }
3320
3321 thread_local! {
3322 static ABORT_DATA_MEMORY: RuntimeMemory<DefaultMemoryImpl> = test_memory(180);
3323 static ABORT_INDEX_MEMORY: RuntimeMemory<DefaultMemoryImpl> = test_memory(181);
3324 static ABORT_SCHEMA_MEMORY: RuntimeMemory<DefaultMemoryImpl> = test_memory(182);
3325 static ABORT_JOURNAL_MEMORY: RuntimeMemory<DefaultMemoryImpl> = test_memory(183);
3326 static ABORT_DATA: RefCell<DataStore> =
3327 ABORT_DATA_MEMORY.with(|memory| {
3328 RefCell::new(DataStore::init_journaled(memory.clone()))
3329 });
3330 static ABORT_INDEX: RefCell<IndexStore> =
3331 ABORT_INDEX_MEMORY.with(|memory| {
3332 RefCell::new(IndexStore::init_journaled(memory.clone()))
3333 });
3334 static ABORT_SCHEMA: RefCell<SchemaStore> =
3335 ABORT_SCHEMA_MEMORY.with(|memory| {
3336 RefCell::new(SchemaStore::init_journaled(memory.clone()))
3337 });
3338 static ABORT_JOURNAL: RefCell<JournalTailStore> =
3339 ABORT_JOURNAL_MEMORY.with(|memory| {
3340 RefCell::new(JournalTailStore::init(memory.clone()))
3341 });
3342 static ABORT_REGISTRY: StoreRegistry = {
3343 let mut registry = StoreRegistry::new();
3344 registry.register_journaled_store(
3345 ABORT_STORE_PATH,
3346 &ABORT_DATA,
3347 &ABORT_INDEX,
3348 &ABORT_SCHEMA,
3349 &ABORT_JOURNAL,
3350 StoreAllocationIdentities::new_journaled(
3351 StoreAllocationIdentity::new(180, "icydb.test.application_abort.data.v1"),
3352 StoreAllocationIdentity::new(181, "icydb.test.application_abort.index.v1"),
3353 StoreAllocationIdentity::new(182, "icydb.test.application_abort.schema.v1"),
3354 StoreAllocationIdentity::new(183, "icydb.test.application_abort.journal.v1"),
3355 ),
3356 StoreRuntimeStorageCapabilities::journaled(),
3357 ).expect("abort journaled store should register");
3358 registry
3359 };
3360 }
3361
3362 #[cfg(feature = "migration")]
3363 thread_local! {
3364 static MIGRATION_EXECUTION_DATA: RefCell<DataStore> =
3365 RefCell::new(DataStore::init_journaled(test_memory(210)));
3366 static MIGRATION_EXECUTION_INDEX: RefCell<IndexStore> =
3367 RefCell::new(IndexStore::init_journaled(test_memory(211)));
3368 static MIGRATION_EXECUTION_SCHEMA: RefCell<SchemaStore> =
3369 RefCell::new(SchemaStore::init_journaled(test_memory(212)));
3370 static MIGRATION_EXECUTION_JOURNAL: RefCell<JournalTailStore> =
3371 RefCell::new(JournalTailStore::init(test_memory(213)));
3372 static MIGRATION_EXECUTION_REGISTRY: StoreRegistry = {
3373 let mut registry = StoreRegistry::new();
3374 registry.register_journaled_store(
3375 MIGRATION_EXECUTION_STORE_PATH,
3376 &MIGRATION_EXECUTION_DATA,
3377 &MIGRATION_EXECUTION_INDEX,
3378 &MIGRATION_EXECUTION_SCHEMA,
3379 &MIGRATION_EXECUTION_JOURNAL,
3380 StoreAllocationIdentities::new_journaled(
3381 StoreAllocationIdentity::new(210, "icydb.test.migration_execution.data.v1"),
3382 StoreAllocationIdentity::new(211, "icydb.test.migration_execution.index.v1"),
3383 StoreAllocationIdentity::new(212, "icydb.test.migration_execution.schema.v1"),
3384 StoreAllocationIdentity::new(213, "icydb.test.migration_execution.journal.v1"),
3385 ),
3386 StoreRuntimeStorageCapabilities::journaled(),
3387 ).expect("migration execution store should register");
3388 registry
3389 };
3390 }
3391
3392 #[cfg(feature = "migration")]
3393 thread_local! {
3394 static MIGRATION_DATA: RefCell<DataStore> =
3395 RefCell::new(DataStore::init_journaled(test_memory(200)));
3396 static MIGRATION_INDEX: RefCell<IndexStore> =
3397 RefCell::new(IndexStore::init_journaled(test_memory(201)));
3398 static MIGRATION_SCHEMA: RefCell<SchemaStore> =
3399 RefCell::new(SchemaStore::init_journaled(test_memory(202)));
3400 static MIGRATION_JOURNAL: RefCell<JournalTailStore> =
3401 RefCell::new(JournalTailStore::init(test_memory(203)));
3402 static MIGRATION_REGISTRY: StoreRegistry = {
3403 let mut registry = StoreRegistry::new();
3404 registry.register_journaled_store(
3405 MIGRATION_STORE_PATH,
3406 &MIGRATION_DATA,
3407 &MIGRATION_INDEX,
3408 &MIGRATION_SCHEMA,
3409 &MIGRATION_JOURNAL,
3410 StoreAllocationIdentities::new_journaled(
3411 StoreAllocationIdentity::new(200, "icydb.test.migration_validation.data.v1"),
3412 StoreAllocationIdentity::new(201, "icydb.test.migration_validation.index.v1"),
3413 StoreAllocationIdentity::new(202, "icydb.test.migration_validation.schema.v1"),
3414 StoreAllocationIdentity::new(203, "icydb.test.migration_validation.journal.v1"),
3415 ),
3416 StoreRuntimeStorageCapabilities::journaled(),
3417 ).expect("migration validation store should register");
3418 registry
3419 };
3420 }
3421
3422 #[cfg(feature = "migration")]
3423 thread_local! {
3424 static MIGRATION_FINDING_DATA: RefCell<DataStore> =
3425 RefCell::new(DataStore::init_journaled(test_memory(206)));
3426 static MIGRATION_FINDING_INDEX: RefCell<IndexStore> =
3427 RefCell::new(IndexStore::init_journaled(test_memory(207)));
3428 static MIGRATION_FINDING_SCHEMA: RefCell<SchemaStore> =
3429 RefCell::new(SchemaStore::init_journaled(test_memory(208)));
3430 static MIGRATION_FINDING_JOURNAL: RefCell<JournalTailStore> =
3431 RefCell::new(JournalTailStore::init(test_memory(209)));
3432 static MIGRATION_FINDING_REGISTRY: StoreRegistry = {
3433 let mut registry = StoreRegistry::new();
3434 registry.register_journaled_store(
3435 MIGRATION_FINDING_STORE_PATH,
3436 &MIGRATION_FINDING_DATA,
3437 &MIGRATION_FINDING_INDEX,
3438 &MIGRATION_FINDING_SCHEMA,
3439 &MIGRATION_FINDING_JOURNAL,
3440 StoreAllocationIdentities::new_journaled(
3441 StoreAllocationIdentity::new(206, "icydb.test.migration_finding.data.v1"),
3442 StoreAllocationIdentity::new(207, "icydb.test.migration_finding.index.v1"),
3443 StoreAllocationIdentity::new(208, "icydb.test.migration_finding.schema.v1"),
3444 StoreAllocationIdentity::new(209, "icydb.test.migration_finding.journal.v1"),
3445 ),
3446 StoreRuntimeStorageCapabilities::journaled(),
3447 ).expect("migration finding store should register");
3448 registry
3449 };
3450 }
3451
3452 thread_local! {
3453 static EVOLUTION_DATA: RefCell<DataStore> =
3454 RefCell::new(DataStore::init_journaled(test_memory(192)));
3455 static EVOLUTION_INDEX: RefCell<IndexStore> =
3456 RefCell::new(IndexStore::init_journaled(test_memory(193)));
3457 static EVOLUTION_SCHEMA: RefCell<SchemaStore> =
3458 RefCell::new(SchemaStore::init_journaled(test_memory(194)));
3459 static EVOLUTION_JOURNAL: RefCell<JournalTailStore> =
3460 RefCell::new(JournalTailStore::init(test_memory(195)));
3461 static EVOLUTION_REGISTRY: StoreRegistry = {
3462 let mut registry = StoreRegistry::new();
3463 registry.register_journaled_store(
3464 EVOLUTION_STORE_PATH,
3465 &EVOLUTION_DATA,
3466 &EVOLUTION_INDEX,
3467 &EVOLUTION_SCHEMA,
3468 &EVOLUTION_JOURNAL,
3469 StoreAllocationIdentities::new_journaled(
3470 StoreAllocationIdentity::new(192, "icydb.test.rule_evolution.data.v1"),
3471 StoreAllocationIdentity::new(193, "icydb.test.rule_evolution.index.v1"),
3472 StoreAllocationIdentity::new(194, "icydb.test.rule_evolution.schema.v1"),
3473 StoreAllocationIdentity::new(195, "icydb.test.rule_evolution.journal.v1"),
3474 ),
3475 StoreRuntimeStorageCapabilities::journaled(),
3476 ).expect("rule-evolution journaled store should register");
3477 registry
3478 };
3479 }
3480
3481 struct AbortCanister;
3482
3483 impl Path for AbortCanister {
3484 const PATH: &'static str = "schema_application_tests::AbortCanister";
3485 }
3486
3487 impl CanisterKind for AbortCanister {
3488 const COMMIT_MEMORY_ID: u8 = 184;
3489 const COMMIT_STABLE_KEY: &'static str = "icydb.test.application_abort.commit.v1";
3490 const STARTUP_MEMORY_ID: u8 = 186;
3491 const STARTUP_STABLE_KEY: &'static str = "icydb.test.application_abort.startup.control.v1";
3492 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 185;
3493 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3494 "icydb.test.application_abort.integrity.v1";
3495 }
3496
3497 struct EvolutionCanister;
3498
3499 impl Path for EvolutionCanister {
3500 const PATH: &'static str = "schema_application_tests::EvolutionCanister";
3501 }
3502
3503 impl CanisterKind for EvolutionCanister {
3504 const COMMIT_MEMORY_ID: u8 = 196;
3505 const COMMIT_STABLE_KEY: &'static str = "icydb.test.rule_evolution.commit.v1";
3506 const STARTUP_MEMORY_ID: u8 = 198;
3507 const STARTUP_STABLE_KEY: &'static str = "icydb.test.rule_evolution.startup.control.v1";
3508 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 197;
3509 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3510 "icydb.test.rule_evolution.integrity.v1";
3511 }
3512
3513 #[cfg(feature = "migration")]
3514 struct MigrationCanister;
3515
3516 #[cfg(feature = "migration")]
3517 impl Path for MigrationCanister {
3518 const PATH: &'static str = "schema_application_tests::MigrationCanister";
3519 }
3520
3521 #[cfg(feature = "migration")]
3522 impl CanisterKind for MigrationCanister {
3523 const COMMIT_MEMORY_ID: u8 = 204;
3524 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_validation.commit.v1";
3525 const STARTUP_MEMORY_ID: u8 = 206;
3526 const STARTUP_STABLE_KEY: &'static str =
3527 "icydb.test.migration_validation.startup.control.v1";
3528 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 205;
3529 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3530 "icydb.test.migration_validation.integrity.v1";
3531 }
3532
3533 #[cfg(feature = "migration")]
3534 struct MigrationExecutionCanister;
3535
3536 #[cfg(feature = "migration")]
3537 impl Path for MigrationExecutionCanister {
3538 const PATH: &'static str = "schema_application_tests::MigrationExecutionCanister";
3539 }
3540
3541 #[cfg(feature = "migration")]
3542 impl CanisterKind for MigrationExecutionCanister {
3543 const COMMIT_MEMORY_ID: u8 = 214;
3544 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_execution.commit.v1";
3545 const STARTUP_MEMORY_ID: u8 = 216;
3546 const STARTUP_STABLE_KEY: &'static str =
3547 "icydb.test.migration_execution.startup.control.v1";
3548 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 215;
3549 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3550 "icydb.test.migration_execution.integrity.v1";
3551 }
3552
3553 #[cfg(feature = "migration")]
3554 struct MigrationFindingCanister;
3555
3556 #[cfg(feature = "migration")]
3557 impl Path for MigrationFindingCanister {
3558 const PATH: &'static str = "schema_application_tests::MigrationFindingCanister";
3559 }
3560
3561 #[cfg(feature = "migration")]
3562 impl CanisterKind for MigrationFindingCanister {
3563 const COMMIT_MEMORY_ID: u8 = 210;
3564 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_finding.commit.v1";
3565 const STARTUP_MEMORY_ID: u8 = 212;
3566 const STARTUP_STABLE_KEY: &'static str = "icydb.test.migration_finding.startup.control.v1";
3567 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 211;
3568 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3569 "icydb.test.migration_finding.integrity.v1";
3570 }
3571
3572 fn name(value: &str) -> SchemaName {
3573 SchemaName::try_new(value).expect("test schema name should admit")
3574 }
3575
3576 fn generated_check_proposal(
3577 expected_head: ExpectedAcceptedHead,
3578 submission_key: &str,
3579 include_check: bool,
3580 database: TargetDatabaseIdentity,
3581 store: TargetStoreIdentity,
3582 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3583 let entity_source = EntitySourceKey::try_new("Item").expect("entity source should admit");
3584 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3585 let score_source = FieldSourceKey::try_new("score").expect("score source should admit");
3586 let check_source =
3587 ConstraintSourceKey::try_new("score_non_negative").expect("check source should admit");
3588 let check = SourceCheckExpr::try_new(vec![
3589 SourceCheckInstruction::Field(score_source),
3590 SourceCheckInstruction::Literal(ScalarLiteral::Int(0)),
3591 SourceCheckInstruction::GreaterThanOrEqual,
3592 ])
3593 .expect("check expression should admit");
3594 let constraints = include_check
3595 .then(|| ConstraintFragment::check(name("score_non_negative"), check))
3596 .into_iter()
3597 .collect();
3598 let entity = EntityFragment::try_new(
3599 name("Item"),
3600 version_one(),
3601 vec![
3602 FieldFragment::new(
3603 name("id"),
3604 FieldType::Scalar(ScalarType::Nat64),
3605 false,
3606 FieldInsertPolicy::Required,
3607 None,
3608 ),
3609 FieldFragment::new(
3610 name("score"),
3611 FieldType::Scalar(ScalarType::Int64),
3612 false,
3613 FieldInsertPolicy::Required,
3614 None,
3615 ),
3616 ],
3617 vec![id_source],
3618 Vec::new(),
3619 Vec::new(),
3620 constraints,
3621 )
3622 .expect("entity should admit");
3623 let proposal = SchemaProposal::try_compose(
3624 vec![SchemaCapability::ACCEPTED_CHECKS],
3625 database,
3626 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3627 expected_head,
3628 vec![
3629 SchemaFragment::try_new(vec![entity], Vec::new())
3630 .expect("schema fragment should admit"),
3631 ],
3632 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3633 Vec::new(),
3634 None,
3635 )
3636 .expect("schema proposal should compose");
3637 (proposal, entity_source, check_source)
3638 }
3639
3640 #[cfg(feature = "migration")]
3641 #[derive(Clone, Copy, Eq, PartialEq)]
3642 enum ValidationMigrationShape {
3643 Clean,
3644 AllFindingFamilies,
3645 }
3646
3647 #[cfg(feature = "migration")]
3648 #[expect(
3649 clippy::too_many_lines,
3650 reason = "the fixture keeps both predecessor and candidate source contracts adjacent"
3651 )]
3652 fn validation_migration_proposal(
3653 shape: ValidationMigrationShape,
3654 current: bool,
3655 expected_head: ExpectedAcceptedHead,
3656 database: TargetDatabaseIdentity,
3657 store: TargetStoreIdentity,
3658 ) -> SchemaProposal {
3659 let entity_source = EntitySourceKey::try_new("MigratingItem")
3660 .expect("migration entity source should admit");
3661 let old_value =
3662 FieldSourceKey::try_new("old_value").expect("predecessor field source should admit");
3663 let current_value =
3664 FieldSourceKey::try_new("value").expect("candidate field source should admit");
3665 let target_entity = EntitySourceKey::try_new("MigrationTarget")
3666 .expect("migration target source should admit");
3667 let target_id = FieldSourceKey::try_new("id").expect("target id source should admit");
3668 let constraint = SourceCheckExpr::try_new(vec![
3669 SourceCheckInstruction::Field(current_value.clone()),
3670 SourceCheckInstruction::Literal(ScalarLiteral::Nat(8)),
3671 SourceCheckInstruction::LessThanOrEqual,
3672 ])
3673 .expect("candidate check should admit");
3674 let findings = shape == ValidationMigrationShape::AllFindingFamilies;
3675 let entity = EntityFragment::try_new(
3676 name("MigratingItem"),
3677 DeclaredEntityVersion::try_new(if current { 2 } else { 1 })
3678 .expect("migration version should admit"),
3679 vec![
3680 FieldFragment::new(
3681 name("id"),
3682 FieldType::Scalar(ScalarType::Nat64),
3683 false,
3684 FieldInsertPolicy::Required,
3685 None,
3686 ),
3687 FieldFragment::new(
3688 name(if current { "value" } else { "old_value" }),
3689 FieldType::Scalar(if current {
3690 ScalarType::Nat8
3691 } else {
3692 ScalarType::Int64
3693 }),
3694 false,
3695 FieldInsertPolicy::Required,
3696 None,
3697 ),
3698 ],
3699 vec![FieldSourceKey::try_new("id").expect("id source should admit")],
3700 current
3701 .then(|| {
3702 IndexFragment::try_new(
3703 name("value_unique"),
3704 vec![IndexKeyFragment::Field(current_value.clone())],
3705 true,
3706 None,
3707 )
3708 .expect("candidate index should admit")
3709 })
3710 .into_iter()
3711 .collect(),
3712 (current && findings)
3713 .then(|| {
3714 RelationFragment::try_new(
3715 name("value_target"),
3716 icydb_schema::RelationSourceFragment::direct(vec![current_value.clone()]),
3717 target_entity.clone(),
3718 vec![target_id.clone()],
3719 RelationDeleteAction::Restrict,
3720 )
3721 .expect("candidate relation should admit")
3722 })
3723 .into_iter()
3724 .collect(),
3725 (current && findings)
3726 .then(|| ConstraintFragment::check(name("value_at_most_eight"), constraint))
3727 .into_iter()
3728 .collect(),
3729 )
3730 .expect("migration entity should admit");
3731 let target = EntityFragment::try_new(
3732 name("MigrationTarget"),
3733 version_one(),
3734 vec![FieldFragment::new(
3735 name("id"),
3736 FieldType::Scalar(ScalarType::Nat8),
3737 false,
3738 FieldInsertPolicy::Required,
3739 None,
3740 )],
3741 vec![target_id],
3742 Vec::new(),
3743 Vec::new(),
3744 Vec::new(),
3745 )
3746 .expect("migration relation target should admit");
3747 let migration = current.then(|| {
3748 SchemaMigrationPlan::try_new(vec![
3749 EntityMigration::try_new(
3750 entity_source.clone(),
3751 DeclaredEntityVersion::try_new(1).expect("predecessor should admit"),
3752 None,
3753 Vec::new(),
3754 vec![SchemaMigrationTransform::CheckedCast {
3755 from: old_value.clone(),
3756 to: current_value,
3757 target: ScalarType::Nat8,
3758 }],
3759 )
3760 .expect("migration transition should admit"),
3761 ])
3762 .expect("migration plan should admit")
3763 });
3764 let mut capabilities = Vec::new();
3765 if current && findings {
3766 capabilities.extend([
3767 SchemaCapability::ACCEPTED_CHECKS,
3768 SchemaCapability::SECONDARY_INDEXES,
3769 SchemaCapability::RESTRICTIVE_RELATIONS,
3770 ]);
3771 }
3772 if migration.is_some() {
3773 capabilities.push(SchemaCapability::VERSIONED_MIGRATIONS);
3774 }
3775 let mut entities = vec![entity];
3776 let mut assignments = vec![EntityStoreAssignment::new(entity_source.clone(), store)];
3777 if findings {
3778 entities.push(target);
3779 assignments.push(EntityStoreAssignment::new(target_entity, store));
3780 }
3781 SchemaProposal::try_compose(
3782 capabilities,
3783 database,
3784 SchemaSubmissionKey::try_new(if current {
3785 "migration-validation-v2"
3786 } else {
3787 "migration-validation-v1"
3788 })
3789 .expect("submission should admit"),
3790 expected_head,
3791 vec![
3792 SchemaFragment::try_new(entities, Vec::new())
3793 .expect("migration fragment should admit"),
3794 ],
3795 assignments,
3796 current
3797 .then_some(icydb_schema::SchemaRemoval::Field {
3798 entity: entity_source,
3799 field: old_value,
3800 })
3801 .into_iter()
3802 .collect(),
3803 migration,
3804 )
3805 .expect("migration proposal should compose")
3806 }
3807
3808 fn targeted_rule_proposal(
3809 expected_head: ExpectedAcceptedHead,
3810 submission_key: &str,
3811 operation: SourceRuleOperation,
3812 database: TargetDatabaseIdentity,
3813 store: TargetStoreIdentity,
3814 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3815 let entity_source =
3816 EntitySourceKey::try_new("Measured").expect("entity source should admit");
3817 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3818 let value_source = FieldSourceKey::try_new("value").expect("value source should admit");
3819 let value_type = TypeSourceKey::try_new("Measure").expect("type source should admit");
3820 let rule_source = RuleSourceKey::try_new("limit").expect("rule source should admit");
3821 let constraint_source =
3822 ConstraintSourceKey::for_targeted_field_rule(&value_source, &value_type, &rule_source);
3823 let entity = EntityFragment::try_new(
3824 name("Measured"),
3825 version_one(),
3826 vec![
3827 FieldFragment::new(
3828 name("id"),
3829 FieldType::Scalar(ScalarType::Nat64),
3830 false,
3831 FieldInsertPolicy::Required,
3832 None,
3833 ),
3834 FieldFragment::new(
3835 name("value"),
3836 FieldType::Named(value_type.clone()),
3837 false,
3838 FieldInsertPolicy::Required,
3839 None,
3840 ),
3841 ],
3842 vec![id_source],
3843 Vec::new(),
3844 Vec::new(),
3845 vec![ConstraintFragment::targeted_rule(
3846 TargetedRuleFragment::new(value_source, value_type, name("limit"), operation),
3847 )],
3848 )
3849 .expect("targeted entity should admit");
3850 let proposal = SchemaProposal::try_compose(
3851 vec![SchemaCapability::ACCEPTED_CHECKS],
3852 database,
3853 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3854 expected_head,
3855 vec![
3856 SchemaFragment::try_new(
3857 vec![entity],
3858 vec![NamedTypeFragment::newtype(
3859 name("Measure"),
3860 FieldType::Scalar(ScalarType::Nat8),
3861 )],
3862 )
3863 .expect("schema fragment should admit"),
3864 ],
3865 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3866 Vec::new(),
3867 None,
3868 )
3869 .expect("schema proposal should compose");
3870 (proposal, entity_source, constraint_source)
3871 }
3872
3873 #[test]
3874 fn database_head_is_empty_only_when_every_store_root_is_absent() {
3875 assert_eq!(
3876 derive_accepted_head(&[("test::A", None), ("test::B", None)]),
3877 ExpectedAcceptedHead::Empty,
3878 );
3879 }
3880
3881 #[test]
3882 fn database_head_covers_store_path_revision_fingerprint_and_absence() {
3883 let first = derive_accepted_head(&[
3884 (
3885 "test::A",
3886 Some(AcceptedStoreHead {
3887 revision: 3,
3888 fingerprint: [0x11; 32],
3889 }),
3890 ),
3891 ("test::B", None),
3892 ]);
3893 let changed_fingerprint = derive_accepted_head(&[
3894 (
3895 "test::A",
3896 Some(AcceptedStoreHead {
3897 revision: 3,
3898 fingerprint: [0x12; 32],
3899 }),
3900 ),
3901 ("test::B", None),
3902 ]);
3903 let changed_absence = derive_accepted_head(&[
3904 (
3905 "test::A",
3906 Some(AcceptedStoreHead {
3907 revision: 3,
3908 fingerprint: [0x11; 32],
3909 }),
3910 ),
3911 (
3912 "test::B",
3913 Some(AcceptedStoreHead {
3914 revision: 1,
3915 fingerprint: [0x22; 32],
3916 }),
3917 ),
3918 ]);
3919
3920 assert_ne!(first, changed_fingerprint);
3921 assert_ne!(first, changed_absence);
3922 assert!(matches!(
3923 first,
3924 ExpectedAcceptedHead::Exact { revision: 3, .. }
3925 ));
3926 }
3927
3928 #[test]
3929 #[allow(
3930 clippy::too_many_lines,
3931 reason = "the end-to-end catalog assertion is clearer as one lifecycle test"
3932 )]
3933 fn generated_check_abort_retires_source_identity_and_allows_fresh_reproposal() {
3934 let database = TargetDatabaseIdentity::from_bytes([0x71; 32]);
3935 let store = TargetStoreIdentity::from_bytes([0x72; 32]);
3936 let (initial, entity_source, _) = generated_check_proposal(
3937 ExpectedAcceptedHead::Empty,
3938 "abort-initial",
3939 false,
3940 database,
3941 store,
3942 );
3943 let initial_candidate = lower_initial_schema_proposal(
3944 &initial,
3945 &[ProposalStoreTarget {
3946 path: "abort::Store",
3947 identity: store,
3948 }],
3949 )
3950 .expect("initial proposal should lower")
3951 .pop()
3952 .expect("initial proposal should produce one candidate");
3953 let (with_check, _, check_source) = generated_check_proposal(
3954 ExpectedAcceptedHead::Exact {
3955 revision: 1,
3956 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x73; 32]),
3957 },
3958 "abort-add-check",
3959 true,
3960 database,
3961 store,
3962 );
3963 let pending_candidate = lower_existing_schema_proposal(
3964 &with_check,
3965 &[ExistingProposalStore {
3966 path: "abort::Store",
3967 identity: store,
3968 bundle: initial_candidate.bundle(),
3969 }],
3970 )
3971 .expect("generated check should lower")
3972 .pop()
3973 .expect("generated check should produce one candidate");
3974 let entity_tag = pending_candidate
3975 .bundle()
3976 .source_bindings_for_tests()
3977 .entity(&entity_source)
3978 .expect("entity source should remain bound");
3979 let constraint_id = pending_candidate
3980 .bundle()
3981 .source_bindings_for_tests()
3982 .constraint(entity_tag, &check_source)
3983 .expect("generated check source should bind");
3984 let pending_snapshot = pending_candidate
3985 .bundle()
3986 .entity_snapshots()
3987 .get(&entity_tag)
3988 .expect("pending entity should exist");
3989 let activation = pending_snapshot
3990 .constraint_catalog()
3991 .activation(constraint_id)
3992 .expect("generated check should remain an activation");
3993 assert_eq!(activation.origin(), ConstraintOrigin::Generated);
3994
3995 let aborted = aborted_generated_row_local_candidate(
3996 pending_candidate.bundle(),
3997 entity_tag,
3998 constraint_id,
3999 )
4000 .expect("generated check abort should build one catalog-native candidate");
4001 let aborted_snapshot = aborted
4002 .bundle()
4003 .entity_snapshots()
4004 .get(&entity_tag)
4005 .expect("aborted entity should remain");
4006 assert!(
4007 aborted_snapshot
4008 .constraint_catalog()
4009 .activation(constraint_id)
4010 .is_none(),
4011 );
4012 assert_eq!(aborted_snapshot.row_layout(), pending_snapshot.row_layout());
4013 assert!(
4014 aborted
4015 .bundle()
4016 .source_bindings_for_tests()
4017 .constraint(entity_tag, &check_source)
4018 .is_none(),
4019 );
4020
4021 let reproposed = lower_existing_schema_proposal(
4022 &with_check,
4023 &[ExistingProposalStore {
4024 path: "abort::Store",
4025 identity: store,
4026 bundle: aborted.bundle(),
4027 }],
4028 )
4029 .expect("aborted generated check should be independently reproposable")
4030 .pop()
4031 .expect("reproposal should produce one candidate");
4032 let replacement_id = reproposed
4033 .bundle()
4034 .source_bindings_for_tests()
4035 .constraint(entity_tag, &check_source)
4036 .expect("reproposal should bind a fresh constraint identity");
4037 assert!(
4038 replacement_id > constraint_id,
4039 "aborted accepted IDs must remain retired",
4040 );
4041 }
4042
4043 #[test]
4044 fn targeted_rule_edit_abort_keeps_prior_accepted_semantics_and_source_identity() {
4045 let database = TargetDatabaseIdentity::from_bytes([0x81; 32]);
4046 let store = TargetStoreIdentity::from_bytes([0x82; 32]);
4047 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4048 ExpectedAcceptedHead::Empty,
4049 "targeted-abort-initial",
4050 SourceRuleOperation::NumericRangeInclusive {
4051 min: ScalarLiteral::Nat(0),
4052 max: ScalarLiteral::Nat(10),
4053 },
4054 database,
4055 store,
4056 );
4057 let initial_candidate = lower_initial_schema_proposal(
4058 &initial,
4059 &[ProposalStoreTarget {
4060 path: "abort::TargetedStore",
4061 identity: store,
4062 }],
4063 )
4064 .expect("initial targeted proposal should lower")
4065 .pop()
4066 .expect("initial targeted proposal should produce one candidate");
4067 let initial_bundle = initial_candidate.bundle();
4068 let entity_tag = initial_bundle
4069 .source_bindings_for_tests()
4070 .entity(&entity_source)
4071 .expect("entity source should bind");
4072 let constraint_id = initial_bundle
4073 .source_bindings_for_tests()
4074 .constraint(entity_tag, &constraint_source)
4075 .expect("targeted source should bind");
4076 let high_water = initial_bundle.entity_snapshots()[&entity_tag]
4077 .constraint_id_allocator()
4078 .high_water();
4079 let (edited, _, _) = targeted_rule_proposal(
4080 ExpectedAcceptedHead::Exact {
4081 revision: initial_bundle.revision().get(),
4082 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x83; 32]),
4083 },
4084 "targeted-abort-edit",
4085 SourceRuleOperation::NumericMaximumInclusive {
4086 value: ScalarLiteral::Nat(8),
4087 },
4088 database,
4089 store,
4090 );
4091 let staged = lower_existing_schema_proposal(
4092 &edited,
4093 &[ExistingProposalStore {
4094 path: "abort::TargetedStore",
4095 identity: store,
4096 bundle: initial_bundle,
4097 }],
4098 )
4099 .expect("targeted semantic edit should stage")
4100 .pop()
4101 .expect("targeted semantic edit should produce one candidate");
4102 let aborted =
4103 aborted_generated_row_local_candidate(staged.bundle(), entity_tag, constraint_id)
4104 .expect("targeted semantic edit should abort through catalog authority");
4105 let snapshot = &aborted.bundle().entity_snapshots()[&entity_tag];
4106
4107 assert!(
4108 snapshot
4109 .constraint_catalog()
4110 .activation(constraint_id)
4111 .is_none()
4112 );
4113 assert_eq!(snapshot.constraint_id_allocator().high_water(), high_water);
4114 assert_eq!(
4115 aborted
4116 .bundle()
4117 .source_bindings_for_tests()
4118 .constraint(entity_tag, &constraint_source),
4119 Some(constraint_id),
4120 );
4121 assert!(snapshot.constraints().iter().any(|constraint| {
4122 constraint.id() == constraint_id
4123 && matches!(
4124 constraint.kind(),
4125 AcceptedConstraintKind::TargetedRule { operation, .. }
4126 if matches!(
4127 operation.as_ref(),
4128 AcceptedRuleOperation::NumericRangeInclusive { .. }
4129 )
4130 )
4131 }));
4132 }
4133
4134 #[test]
4135 #[allow(
4136 clippy::too_many_lines,
4137 reason = "the staged publication, recovery, and promotion assertions form one lifecycle"
4138 )]
4139 fn generated_startup_driver_resumes_pending_activation_and_promotes_without_source_model() {
4140 let db = Db::<EvolutionCanister>::new(
4141 &EVOLUTION_REGISTRY,
4142 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4143 );
4144 drive_startup_recovery_to_completion(&db);
4145 let empty_target =
4146 schema_application_target(&db).expect("empty evolution target should issue");
4147 let store_identity = empty_target
4148 .stores()
4149 .first()
4150 .expect("evolution store should register")
4151 .identity();
4152 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4153 empty_target.accepted_head().clone(),
4154 "targeted-recovery-initial",
4155 SourceRuleOperation::NumericRangeInclusive {
4156 min: ScalarLiteral::Nat(0),
4157 max: ScalarLiteral::Nat(10),
4158 },
4159 empty_target.database_identity(),
4160 store_identity,
4161 );
4162 assert!(matches!(
4163 apply_schema(&db, &initial)
4164 .expect("initial targeted proposal should publish")
4165 .outcome(),
4166 SchemaChangeOutcome::Applied { .. },
4167 ));
4168 drive_startup_recovery_to_completion(&db);
4169 drive_cardinality_to_ready(
4170 db.store_handle(EVOLUTION_STORE_PATH)
4171 .expect("evolution store should resolve"),
4172 );
4173
4174 let direct_target =
4175 schema_application_target(&db).expect("direct evolution target should issue");
4176 let (direct_edit, _, _) = targeted_rule_proposal(
4177 direct_target.accepted_head().clone(),
4178 "targeted-direct-edit",
4179 SourceRuleOperation::NumericMaximumInclusive {
4180 value: ScalarLiteral::Nat(8),
4181 },
4182 direct_target.database_identity(),
4183 store_identity,
4184 );
4185 assert!(matches!(
4186 apply_schema(&db, &direct_edit)
4187 .expect("empty-domain semantic edit should publish directly")
4188 .outcome(),
4189 SchemaChangeOutcome::Applied { .. },
4190 ));
4191 let store = db
4192 .store_handle(EVOLUTION_STORE_PATH)
4193 .expect("evolution store should resolve");
4194 let direct = store
4195 .with_schema(SchemaStore::current_accepted_schema_bundle)
4196 .expect("directly edited bundle should remain readable")
4197 .expect("directly edited bundle should exist");
4198 let entity_tag = direct
4199 .source_bindings_for_tests()
4200 .entity(&entity_source)
4201 .expect("entity source should remain bound");
4202 let constraint_id = direct
4203 .source_bindings_for_tests()
4204 .constraint(entity_tag, &constraint_source)
4205 .expect("direct edit should preserve constraint identity");
4206 assert!(
4207 direct.entity_snapshots()[&entity_tag]
4208 .constraint_catalog()
4209 .activation(constraint_id)
4210 .is_none()
4211 );
4212 assert!(
4213 direct.entity_snapshots()[&entity_tag]
4214 .constraints()
4215 .iter()
4216 .any(|constraint| {
4217 constraint.id() == constraint_id
4218 && matches!(
4219 constraint.kind(),
4220 AcceptedConstraintKind::TargetedRule { operation, .. }
4221 if matches!(
4222 operation.as_ref(),
4223 AcceptedRuleOperation::NumericMaximumInclusive { .. }
4224 )
4225 )
4226 })
4227 );
4228
4229 let target = schema_application_target(&db).expect("staged evolution target should issue");
4230 let (edited, _, _) = targeted_rule_proposal(
4231 target.accepted_head().clone(),
4232 "targeted-recovery-edit",
4233 SourceRuleOperation::MultipleOf {
4234 divisor: ScalarLiteral::Nat(2),
4235 },
4236 target.database_identity(),
4237 store_identity,
4238 );
4239 let current = store
4240 .with_schema(SchemaStore::current_accepted_schema_bundle)
4241 .expect("accepted evolution bundle should remain readable")
4242 .expect("directly edited evolution bundle should exist");
4243 let staged = lower_existing_schema_proposal(
4244 &edited,
4245 &[ExistingProposalStore {
4246 path: EVOLUTION_STORE_PATH,
4247 identity: store_identity,
4248 bundle: ¤t,
4249 }],
4250 )
4251 .expect("targeted edit should stage")
4252 .pop()
4253 .expect("targeted edit should produce one candidate");
4254 assert_eq!(
4255 staged
4256 .bundle()
4257 .source_bindings_for_tests()
4258 .constraint(entity_tag, &constraint_source),
4259 Some(constraint_id),
4260 );
4261 let proof = DirectGeneratedRowLocalProof {
4262 candidate_index: 0,
4263 store,
4264 store_path: EVOLUTION_STORE_PATH,
4265 entity_tag,
4266 entity_path: staged.bundle().entity_snapshots()[&entity_tag]
4267 .entity_path()
4268 .to_string(),
4269 constraint_id,
4270 historical_rows: 0,
4271 };
4272 let final_candidates = final_candidates_for_pending_row_local_constraint(
4273 std::slice::from_ref(&staged),
4274 &PendingGeneratedRowLocalConstraint { proof },
4275 )
4276 .expect("final semantic replacement should derive without source input");
4277 let authorities = application_authorities(&db);
4278 let candidate_head =
4279 accepted_head_after_candidates(authorities.as_slice(), &final_candidates)
4280 .expect("final candidate head should derive");
4281 let digest = edited.digest().expect("proposal digest should derive");
4282 let job_id = derive_schema_change_job_id(
4283 target.database_identity(),
4284 edited.submission_key(),
4285 digest,
4286 target.accepted_head(),
4287 )
4288 .expect("job identity should derive");
4289 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4290 target.database_identity(),
4291 edited.submission_key().clone(),
4292 digest,
4293 target.accepted_head().clone(),
4294 SchemaChangeOutcome::Pending {
4295 job: SchemaChangeJob::new(job_id),
4296 candidate_head,
4297 },
4298 )
4299 .expect("pending replacement receipt should admit");
4300 let record = SchemaApplicationRecord::new(
4301 receipt,
4302 vec![
4303 SchemaChangeActivation::new(
4304 store_identity,
4305 entity_tag.value(),
4306 constraint_id.get(),
4307 )
4308 .expect("replacement activation should admit"),
4309 ],
4310 )
4311 .expect("pending replacement record should admit");
4312 let operation =
4313 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4314 publish_accepted_schema_candidates_with_application_record(
4315 vec![AcceptedSchemaPublication::new(
4316 EVOLUTION_STORE_PATH,
4317 store,
4318 current.revision(),
4319 &staged,
4320 )],
4321 operation,
4322 )
4323 .expect("staged replacement and record should publish atomically");
4324
4325 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4326 drive_startup_recovery_to_completion(&db);
4327
4328 let recovered = store
4329 .with_schema(SchemaStore::current_accepted_schema_bundle)
4330 .expect("recovered staged bundle should decode")
4331 .expect("recovered staged bundle should exist");
4332 let recovered_snapshot = recovered.entity_snapshots()[&entity_tag].clone();
4333 let validating_catalog = recovered_snapshot
4334 .constraint_catalog()
4335 .clone()
4336 .with_validation_started(constraint_id)
4337 .expect("recovered replacement should enter validation");
4338 let mut validating_snapshots = recovered.entity_snapshots().clone();
4339 validating_snapshots.insert(
4340 entity_tag,
4341 recovered_snapshot.with_constraint_catalog(validating_catalog),
4342 );
4343 let validating_bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
4344 recovered
4345 .revision()
4346 .checked_next()
4347 .expect("validation revision should remain available"),
4348 recovered.store_path(),
4349 recovered.enum_catalog().clone(),
4350 recovered.composite_catalog().clone(),
4351 recovered.source_bindings_for_tests().clone(),
4352 validating_snapshots,
4353 )
4354 .expect("validating replacement bundle should close");
4355 let validating_candidate = CandidateSchemaRevision::new(validating_bundle)
4356 .expect("validating replacement candidate should encode");
4357 let validating_activation = validating_candidate.bundle().entity_snapshots()[&entity_tag]
4358 .constraint_catalog()
4359 .activation(constraint_id)
4360 .expect("validating replacement activation should remain present");
4361 let validation_job = ConstraintValidationJob::start(
4362 entity_tag,
4363 validating_candidate.bundle().entity_snapshots()[&entity_tag]
4364 .entity_path()
4365 .to_string(),
4366 validating_activation,
4367 None,
4368 )
4369 .expect("validating replacement job should derive from accepted state");
4370 store
4371 .with_schema(|schema| {
4372 schema.validate_live_activation_transition(validating_candidate.bundle())?;
4373 schema.validate_constraint_validation_job_closure_with_change(
4374 validating_candidate.bundle(),
4375 Some(&validation_job),
4376 None,
4377 )
4378 })
4379 .expect("validating replacement transition and job should close");
4380 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4381 let startup_session =
4382 crate::db::DbSession::<EvolutionCanister>::new(&EVOLUTION_REGISTRY, &startup_root);
4383
4384 assert_eq!(
4385 drive_generated_startup_recovery_page(
4386 &startup_session,
4387 &EVOLUTION_REGISTRY,
4388 edited.submission_key().as_str(),
4389 )
4390 .expect("generated startup should begin pending validation"),
4391 GeneratedStartupDriverStep::Recovering,
4392 );
4393 let mut terminal = false;
4394 for _ in 0..8 {
4395 match drive_generated_startup_recovery_page(
4396 &startup_session,
4397 &EVOLUTION_REGISTRY,
4398 edited.submission_key().as_str(),
4399 )
4400 .expect("generated startup should advance pending validation")
4401 {
4402 GeneratedStartupDriverStep::Recovering => {}
4403 GeneratedStartupDriverStep::Terminal => {
4404 terminal = true;
4405 break;
4406 }
4407 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4408 panic!("an exact pending receipt must resume instead of being resubmitted")
4409 }
4410 }
4411 }
4412 assert!(
4413 terminal,
4414 "empty historical domain should promote within bounded startup steps"
4415 );
4416 assert_eq!(
4417 observe_generated_startup_state::<EvolutionCanister>(
4418 &EVOLUTION_REGISTRY,
4419 edited.submission_key().as_str(),
4420 ),
4421 Ok(DatabaseStartupState::Ready),
4422 );
4423 let applied = super::exact_schema_application_receipt(
4424 &edited,
4425 edited
4426 .digest()
4427 .expect("proposal digest should remain stable"),
4428 )
4429 .expect("terminal generated receipt should remain readable")
4430 .expect("terminal generated receipt should remain present");
4431 assert!(matches!(
4432 applied.outcome(),
4433 SchemaChangeOutcome::Applied { .. }
4434 ));
4435 let promoted = store
4436 .with_schema(SchemaStore::current_accepted_schema_bundle)
4437 .expect("promoted bundle should remain readable")
4438 .expect("promoted bundle should exist");
4439 let snapshot = &promoted.entity_snapshots()[&entity_tag];
4440 assert!(
4441 snapshot
4442 .constraint_catalog()
4443 .activation(constraint_id)
4444 .is_none()
4445 );
4446 assert_eq!(
4447 promoted
4448 .source_bindings_for_tests()
4449 .constraint(entity_tag, &constraint_source),
4450 Some(constraint_id),
4451 );
4452 assert!(snapshot.constraints().iter().any(|constraint| {
4453 constraint.id() == constraint_id
4454 && matches!(
4455 constraint.kind(),
4456 AcceptedConstraintKind::TargetedRule { operation, .. }
4457 if matches!(
4458 operation.as_ref(),
4459 AcceptedRuleOperation::MultipleOf { .. }
4460 )
4461 )
4462 }));
4463 }
4464
4465 #[test]
4466 #[allow(
4467 clippy::too_many_lines,
4468 reason = "the durable pending job, startup failure, and retained finding assertions form one scenario"
4469 )]
4470 fn generated_startup_driver_persists_e223_for_a_retained_historical_finding() {
4471 let db = Db::<AbortCanister>::new(
4472 &ABORT_REGISTRY,
4473 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4474 );
4475 drive_startup_recovery_to_completion(&db);
4476 let empty_target =
4477 schema_application_target(&db).expect("empty application target should issue");
4478 let store_identity = empty_target
4479 .stores()
4480 .first()
4481 .expect("abort store should be registered")
4482 .identity();
4483 let (initial, _, _) = generated_check_proposal(
4484 empty_target.accepted_head().clone(),
4485 "startup-finding-initial",
4486 false,
4487 empty_target.database_identity(),
4488 store_identity,
4489 );
4490 apply_schema(&db, &initial).expect("initial generated schema should publish");
4491
4492 let root = crate::db::RequestExecutionRoot::__new_runtime_root();
4493 let session = DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &root);
4494 let rows = (1..=257)
4495 .map(|id| DynamicMutation::Insert {
4496 entity: "Item".to_string(),
4497 patch: DynamicStructuralPatch::new(vec![
4498 (
4499 "id".to_string(),
4500 DynamicWriteCell::Value(InputValue::nat64(id)),
4501 ),
4502 (
4503 "score".to_string(),
4504 DynamicWriteCell::Value(InputValue::int64(if id == 257 { -1 } else { 1 })),
4505 ),
4506 ]),
4507 })
4508 .collect();
4509 session
4510 .execute_trusted_dynamic_mutation_batch(rows)
4511 .expect("historical finding fixture rows should commit as one legal batch");
4512 drive_startup_recovery_to_completion(&db);
4513 drive_cardinality_to_ready(
4514 db.store_handle(ABORT_STORE_PATH)
4515 .expect("abort store should resolve"),
4516 );
4517
4518 let target =
4519 schema_application_target(&db).expect("existing application target should issue");
4520 let (with_check, _, _) = generated_check_proposal(
4521 target.accepted_head().clone(),
4522 "startup-finding-pending",
4523 true,
4524 target.database_identity(),
4525 store_identity,
4526 );
4527 let pending = apply_schema(&db, &with_check)
4528 .expect("the first clean page should admit durable continuation");
4529 let SchemaChangeOutcome::Pending { job, .. } = pending.outcome() else {
4530 panic!("a 257-row domain must exceed the 256-row direct proof page")
4531 };
4532
4533 let mut terminal = false;
4534 for _ in 0..8 {
4535 match drive_generated_startup_recovery_page(
4536 &session,
4537 &ABORT_REGISTRY,
4538 with_check.submission_key().as_str(),
4539 )
4540 .expect("generated startup should retain a typed finding failure")
4541 {
4542 GeneratedStartupDriverStep::Recovering => {}
4543 GeneratedStartupDriverStep::Terminal => {
4544 terminal = true;
4545 break;
4546 }
4547 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4548 panic!("an exact pending receipt must not be resubmitted")
4549 }
4550 }
4551 }
4552 assert!(terminal, "the retained finding should become terminal");
4553 let failure = observe_generated_startup_state::<AbortCanister>(
4554 &ABORT_REGISTRY,
4555 with_check.submission_key().as_str(),
4556 )
4557 .expect_err("the retained finding must remain durably observable");
4558 assert_eq!(
4559 failure.kind(),
4560 crate::db::StartupFailureKind::SchemaReconciliation,
4561 );
4562 assert_eq!(
4563 failure.diagnostic().error_code(),
4564 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_CONSTRAINT_VIOLATION,
4565 );
4566 assert_eq!(
4567 ABORT_DATA.with(|store| store.borrow().len()),
4568 257,
4569 "terminal startup publication must not change historical rows",
4570 );
4571 assert!(matches!(
4572 continue_schema_application(&db, job.id(), None)
4573 .expect("the retained finding page should replay exactly")
4574 .status(),
4575 SchemaChangeProgressStatus::Findings { findings, .. } if !findings.is_empty(),
4576 ));
4577 }
4578
4579 #[test]
4580 #[allow(
4581 clippy::too_many_lines,
4582 reason = "the journaled abort, replay, and recovery assertions form one scenario"
4583 )]
4584 fn pending_generated_check_abort_is_atomic_terminal_and_replayable() {
4585 let db = Db::<AbortCanister>::new(
4586 &ABORT_REGISTRY,
4587 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4588 );
4589 drive_startup_recovery_to_completion(&db);
4590 let empty_target =
4591 schema_application_target(&db).expect("empty application target should issue");
4592 let store_identity = empty_target
4593 .stores()
4594 .first()
4595 .expect("abort store should be registered")
4596 .identity();
4597 let (initial, entity_source, _) = generated_check_proposal(
4598 empty_target.accepted_head().clone(),
4599 "abort-runtime-initial",
4600 false,
4601 empty_target.database_identity(),
4602 store_identity,
4603 );
4604 assert!(matches!(
4605 apply_schema(&db, &initial)
4606 .expect("initial application should publish")
4607 .outcome(),
4608 SchemaChangeOutcome::Applied { .. },
4609 ));
4610
4611 let target =
4612 schema_application_target(&db).expect("existing application target should issue");
4613 let (with_check, _, check_source) = generated_check_proposal(
4614 target.accepted_head().clone(),
4615 "abort-runtime-pending",
4616 true,
4617 target.database_identity(),
4618 store_identity,
4619 );
4620 let store = db
4621 .store_handle(ABORT_STORE_PATH)
4622 .expect("abort store should resolve");
4623 let current = store
4624 .with_schema(SchemaStore::current_accepted_schema_bundle)
4625 .expect("accepted bundle should remain readable")
4626 .expect("initial accepted bundle should exist");
4627 let pending_candidate = lower_existing_schema_proposal(
4628 &with_check,
4629 &[ExistingProposalStore {
4630 path: ABORT_STORE_PATH,
4631 identity: store_identity,
4632 bundle: ¤t,
4633 }],
4634 )
4635 .expect("pending generated check should lower")
4636 .pop()
4637 .expect("pending generated check should produce one candidate");
4638 let entity_tag = pending_candidate
4639 .bundle()
4640 .source_bindings_for_tests()
4641 .entity(&entity_source)
4642 .expect("entity source should bind");
4643 let constraint_id = pending_candidate
4644 .bundle()
4645 .source_bindings_for_tests()
4646 .constraint(entity_tag, &check_source)
4647 .expect("generated check source should bind");
4648 let digest = with_check.digest().expect("proposal digest should derive");
4649 let job_id = derive_schema_change_job_id(
4650 target.database_identity(),
4651 with_check.submission_key(),
4652 digest,
4653 target.accepted_head(),
4654 )
4655 .expect("job identity should derive");
4656 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4657 target.database_identity(),
4658 with_check.submission_key().clone(),
4659 digest,
4660 target.accepted_head().clone(),
4661 SchemaChangeOutcome::Pending {
4662 job: SchemaChangeJob::new(job_id),
4663 candidate_head: ExpectedAcceptedHead::Exact {
4664 revision: pending_candidate.revision().get().saturating_add(2),
4665 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x76; 32]),
4666 },
4667 },
4668 )
4669 .expect("pending receipt should admit");
4670 let record = SchemaApplicationRecord::new(
4671 receipt,
4672 vec![
4673 SchemaChangeActivation::new(
4674 store_identity,
4675 entity_tag.value(),
4676 constraint_id.get(),
4677 )
4678 .expect("application activation should admit"),
4679 ],
4680 )
4681 .expect("pending application record should admit");
4682 let operation =
4683 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4684 publish_accepted_schema_candidates_with_application_record(
4685 vec![AcceptedSchemaPublication::new(
4686 ABORT_STORE_PATH,
4687 store,
4688 current.revision(),
4689 &pending_candidate,
4690 )],
4691 operation,
4692 )
4693 .expect("pending candidate and record should publish atomically");
4694
4695 let started = continue_schema_application(&db, job_id, None)
4696 .expect("first continuation should durably start validation");
4697 assert_eq!(started.status(), &SchemaChangeProgressStatus::Started);
4698 let progress =
4699 abort_schema_application(&db, job_id, None).expect("pending application should abort");
4700 assert_eq!(progress.status(), &SchemaChangeProgressStatus::Aborted);
4701 assert!(matches!(
4702 progress.receipt().outcome(),
4703 SchemaChangeOutcome::Aborted { .. },
4704 ));
4705 let replay =
4706 abort_schema_application(&db, job_id, None).expect("terminal abort should replay");
4707 assert_eq!(replay, progress);
4708 assert_eq!(
4709 continue_schema_application(&db, job_id, None)
4710 .expect("continuation after abort should replay terminal state"),
4711 progress,
4712 );
4713
4714 let aborted = store
4715 .with_schema(SchemaStore::current_accepted_schema_bundle)
4716 .expect("accepted bundle should remain readable")
4717 .expect("aborted accepted bundle should exist");
4718 assert!(
4719 aborted
4720 .entity_snapshots()
4721 .get(&entity_tag)
4722 .expect("entity should remain after abort")
4723 .constraint_catalog()
4724 .activation(constraint_id)
4725 .is_none(),
4726 );
4727 assert!(
4728 aborted
4729 .source_bindings_for_tests()
4730 .constraint(entity_tag, &check_source)
4731 .is_none(),
4732 );
4733 assert!(
4734 store
4735 .with_schema(|schema| {
4736 schema.constraint_validation_job(entity_tag, constraint_id)
4737 })
4738 .expect("validation-job storage should remain readable")
4739 .is_none(),
4740 );
4741
4742 ABORT_DATA.with(|store| {
4743 ABORT_DATA_MEMORY.with(|memory| {
4744 *store.borrow_mut() = DataStore::init_journaled(memory.clone());
4745 });
4746 });
4747 ABORT_INDEX.with(|store| {
4748 ABORT_INDEX_MEMORY.with(|memory| {
4749 *store.borrow_mut() = IndexStore::init_journaled(memory.clone());
4750 });
4751 });
4752 ABORT_SCHEMA.with(|store| {
4753 ABORT_SCHEMA_MEMORY.with(|memory| {
4754 *store.borrow_mut() = SchemaStore::init_journaled(memory.clone());
4755 });
4756 });
4757 ABORT_JOURNAL.with(|store| {
4758 ABORT_JOURNAL_MEMORY.with(|memory| {
4759 *store.borrow_mut() = JournalTailStore::init(memory.clone());
4760 });
4761 });
4762 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4763 drive_startup_recovery_to_completion(&db);
4764 assert_eq!(
4765 abort_schema_application(&db, job_id, None)
4766 .expect("recovered terminal abort should replay"),
4767 progress,
4768 );
4769 assert!(
4770 store
4771 .with_schema(|schema| {
4772 schema.constraint_validation_job(entity_tag, constraint_id)
4773 })
4774 .expect("recovered validation-job storage should remain readable")
4775 .is_none(),
4776 );
4777 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4778 let startup_session =
4779 crate::db::DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &startup_root);
4780 ABORT_JOURNAL.with(|journal| {
4781 let journal = journal.borrow();
4782 assert!(
4783 journal
4784 .validate_current_tail_authority()
4785 .expect("recovered abort tail control should remain readable")
4786 .is_empty(),
4787 "recovered abort tail control should close exactly",
4788 );
4789 });
4790 assert_eq!(
4791 drive_generated_startup_recovery_page(
4792 &startup_session,
4793 &ABORT_REGISTRY,
4794 with_check.submission_key().as_str(),
4795 )
4796 .expect("an exact aborted generated submission should publish terminal startup state"),
4797 GeneratedStartupDriverStep::Terminal,
4798 );
4799 let failure = observe_generated_startup_state::<AbortCanister>(
4800 &ABORT_REGISTRY,
4801 with_check.submission_key().as_str(),
4802 )
4803 .expect_err("an aborted generated submission must not remain retryable forever");
4804 assert_eq!(
4805 failure.kind(),
4806 crate::db::StartupFailureKind::SchemaReconciliation,
4807 );
4808 assert_eq!(
4809 failure.diagnostic().error_code(),
4810 icydb_diagnostic_code::ErrorCode::RUNTIME_CONFLICT,
4811 );
4812 }
4813
4814 #[cfg(feature = "migration")]
4815 #[test]
4816 fn migration_planning_failures_retain_typed_public_classification() {
4817 use super::schema_migration_planning_error;
4818 use crate::db::schema::migration_planner::SchemaMigrationPlanningError;
4819 use icydb_diagnostic_code::{DiagnosticDetail, SchemaMigrationCode};
4820
4821 for (error, reason) in [
4822 (
4823 SchemaMigrationPlanningError::Unadopted,
4824 SchemaMigrationCode::Unadopted,
4825 ),
4826 (
4827 SchemaMigrationPlanningError::MissingMigration,
4828 SchemaMigrationCode::MissingMigration,
4829 ),
4830 (
4831 SchemaMigrationPlanningError::VersionGap,
4832 SchemaMigrationCode::VersionGap,
4833 ),
4834 (
4835 SchemaMigrationPlanningError::Downgrade,
4836 SchemaMigrationCode::Downgrade,
4837 ),
4838 (
4839 SchemaMigrationPlanningError::EmptyEntityVersionBump,
4840 SchemaMigrationCode::EmptyEntityVersionBump,
4841 ),
4842 (
4843 SchemaMigrationPlanningError::StaleAcceptedHead,
4844 SchemaMigrationCode::StaleAcceptedHead,
4845 ),
4846 (
4847 SchemaMigrationPlanningError::UnknownFromObject,
4848 SchemaMigrationCode::UnknownFromObject,
4849 ),
4850 (
4851 SchemaMigrationPlanningError::UnknownToObject,
4852 SchemaMigrationCode::UnknownToObject,
4853 ),
4854 (
4855 SchemaMigrationPlanningError::KindMismatch,
4856 SchemaMigrationCode::KindMismatch,
4857 ),
4858 (
4859 SchemaMigrationPlanningError::IdentityConflict,
4860 SchemaMigrationCode::IdentityConflict,
4861 ),
4862 (
4863 SchemaMigrationPlanningError::UnexplainedSchemaDifference,
4864 SchemaMigrationCode::UnexplainedSchemaDifference,
4865 ),
4866 (
4867 SchemaMigrationPlanningError::UnsupportedTransform,
4868 SchemaMigrationCode::UnsupportedTransform,
4869 ),
4870 (
4871 SchemaMigrationPlanningError::RekeyedCatalogInvalid,
4872 SchemaMigrationCode::CandidateMismatch,
4873 ),
4874 (
4875 SchemaMigrationPlanningError::CandidateMismatch,
4876 SchemaMigrationCode::CandidateMismatch,
4877 ),
4878 (
4879 SchemaMigrationPlanningError::CorruptLineage,
4880 SchemaMigrationCode::ProgressCorrupt,
4881 ),
4882 ] {
4883 let diagnostic = schema_migration_planning_error(error).diagnostic();
4884 assert_eq!(
4885 diagnostic.detail(),
4886 Some(&DiagnosticDetail::SchemaMigration { reason }),
4887 );
4888 assert_eq!(diagnostic.code(), reason.diagnostic_code());
4889 }
4890 }
4891
4892 #[cfg(feature = "migration")]
4893 #[test]
4894 fn migration_preparation_failure_retains_original_budget_diagnostic() {
4895 use crate::db::{
4896 executor::budget::MaintenanceConstructionBudget,
4897 query::construction::ConstructionBudget,
4898 schema::migration_planner::SchemaMigrationPlanningError,
4899 };
4900 use icydb_diagnostic_code::DiagnosticExecutionBudgetResource as Resource;
4901
4902 let budget =
4903 MaintenanceConstructionBudget::with_limit_for_tests(Resource::TemporaryBytes, 0);
4904 let error = budget.charge(Resource::TemporaryBytes, 1).unwrap_err();
4905 let expected = error.diagnostic();
4906 let result = super::schema_migration_planning_error(
4907 SchemaMigrationPlanningError::Preparation(error),
4908 );
4909 assert_eq!(result.diagnostic(), expected);
4910 }
4911
4912 #[cfg(feature = "migration")]
4913 #[test]
4914 #[expect(
4915 clippy::too_many_lines,
4916 reason = "the validation replay, staging, and unchanged-row assertions form one scenario"
4917 )]
4918 fn physical_migration_validation_is_bounded_staged_and_does_not_rewrite_rows() {
4919 use std::convert::Infallible;
4920
4921 use super::{defer_generated_schema_application_for_prepared_migration, migrate_schema};
4922 use crate::db::{
4923 data::StoreVisit,
4924 index::{IndexEntryValue, IndexId, IndexKey, IndexKeyKind},
4925 key_taxonomy::{PrimaryKeyComponent, PrimaryKeyValue},
4926 schema::{SchemaMigrationCommand, SchemaMigrationPhase},
4927 };
4928 use crate::types::EntityTag;
4929
4930 let db = Db::<MigrationCanister>::new(
4931 &MIGRATION_REGISTRY,
4932 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4933 );
4934 drive_startup_recovery_to_completion(&db);
4935 let initial_target = schema_application_target(&db).expect("initial target should issue");
4936 let store_identity = initial_target
4937 .stores()
4938 .first()
4939 .expect("migration store should exist")
4940 .identity();
4941 let initial = validation_migration_proposal(
4942 ValidationMigrationShape::Clean,
4943 false,
4944 initial_target.accepted_head().clone(),
4945 initial_target.database_identity(),
4946 store_identity,
4947 );
4948 apply_schema(&db, &initial).expect("initial schema should publish");
4949
4950 let session = DbSession::<MigrationCanister>::new(
4951 &MIGRATION_REGISTRY,
4952 &crate::db::RequestExecutionRoot::__new_runtime_root(),
4953 );
4954 for (id, value) in [(1, 7), (2, 8)] {
4955 session
4956 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4957 entity: "MigratingItem".to_string(),
4958 patch: DynamicStructuralPatch::new(vec![
4959 (
4960 "id".to_string(),
4961 DynamicWriteCell::Value(InputValue::nat64(id)),
4962 ),
4963 (
4964 "old_value".to_string(),
4965 DynamicWriteCell::Value(InputValue::int64(value)),
4966 ),
4967 ]),
4968 })
4969 .expect("predecessor row should insert");
4970 }
4971 let store = db
4972 .store_handle(MIGRATION_STORE_PATH)
4973 .expect("migration store should resolve");
4974 let row_bytes = || {
4975 store.with_data(|data| {
4976 let mut rows = Vec::new();
4977 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
4978 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
4979 Ok(StoreVisit::Continue)
4980 });
4981 result.expect("infallible row visit should complete");
4982 rows
4983 })
4984 };
4985 let before_rows = row_bytes();
4986
4987 let target = schema_application_target(&db).expect("migration target should issue");
4988 let proposal = validation_migration_proposal(
4989 ValidationMigrationShape::Clean,
4990 true,
4991 target.accepted_head().clone(),
4992 target.database_identity(),
4993 store_identity,
4994 );
4995 let plan = proposal
4996 .migration()
4997 .expect("migration plan should exist")
4998 .digest();
4999 let command = || SchemaMigrationCommand::Advance {
5000 expected_database: target.database_identity(),
5001 expected_head: target.accepted_head().clone(),
5002 expected_plan: plan,
5003 acknowledged_finding_page: None,
5004 };
5005 assert_eq!(
5006 migrate_schema(&db, &proposal, command())
5007 .unwrap_or_else(|error| {
5008 panic!(
5009 "physical migration should prepare: {:?}",
5010 error.diagnostic()
5011 )
5012 })
5013 .phase(),
5014 SchemaMigrationPhase::Prepared,
5015 );
5016 assert_eq!(
5017 migrate_schema(&db, &proposal, command())
5018 .expect("physical migration should enter validation")
5019 .phase(),
5020 SchemaMigrationPhase::Validating,
5021 );
5022 let record = super::load_schema_migration_record()
5023 .expect("migration record should remain readable")
5024 .expect("validating migration record should exist");
5025 let planned = super::recompile_active_physical_migration(&db, &proposal, &record)
5026 .expect("the exact active plan should recompile");
5027 for _ in 0..2 {
5028 let page = super::validate_migration_page(&db, &planned, record.progress())
5029 .expect("the same validation page should remain replayable");
5030 let (progress, staged, exhausted) = page.into_parts();
5031 assert!(progress.findings().is_empty());
5032 assert!(exhausted);
5033 super::stage_migration_index_entries(staged)
5034 .expect("staging before a cursor marker should be idempotent");
5035 }
5036 assert_eq!(
5037 store.with_index(IndexStore::len),
5038 2,
5039 "replaying an uncheckpointed page must retain one exact staged key per row",
5040 );
5041 let ready =
5042 migrate_schema(&db, &proposal, command()).expect("bounded validation should complete");
5043 assert_eq!(ready.phase(), SchemaMigrationPhase::ReadyToRewrite);
5044 assert_eq!(ready.rows_validated(), 2);
5045 assert!(ready.findings().is_empty());
5046 assert_eq!(row_bytes(), before_rows, "validation must not rewrite rows");
5047 assert_eq!(
5048 store.with_index(IndexStore::len),
5049 2,
5050 "the isolated candidate unique generation should be durably staged",
5051 );
5052 store.with_index_mut(|index| {
5053 for ordinal in 0..513_u64 {
5054 let component = ordinal.to_be_bytes();
5055 let key = IndexKey::new_from_components_with_primary_key_value(
5056 &IndexId::new(EntityTag::new(2), 0),
5057 IndexKeyKind::User,
5058 &[component],
5059 &PrimaryKeyValue::from(PrimaryKeyComponent::Nat64(ordinal)),
5060 )
5061 .expect("unrelated abort-scan key should build")
5062 .to_raw()
5063 .expect("unrelated abort-scan key should encode");
5064 index.insert(key, IndexEntryValue::presence());
5065 }
5066 });
5067 let abort = || SchemaMigrationCommand::Abort {
5068 expected_database: target.database_identity(),
5069 expected_head: target.accepted_head().clone(),
5070 expected_plan: plan,
5071 };
5072 let cleaning = migrate_schema(&db, &proposal, abort())
5073 .expect("the first bounded abort cleanup page should publish");
5074 assert_eq!(cleaning.phase(), SchemaMigrationPhase::ReadyToRewrite);
5075 assert_eq!(store.with_index(IndexStore::len), 513);
5076 let aborted =
5077 migrate_schema(&db, &proposal, abort()).expect("pre-rewrite migration should abort");
5078 assert_eq!(aborted.phase(), SchemaMigrationPhase::Aborted);
5079 assert_eq!(
5080 store.with_index(IndexStore::len),
5081 513,
5082 "abort must remove only planner-invisible candidate generations",
5083 );
5084 assert_eq!(
5085 row_bytes(),
5086 before_rows,
5087 "abort must retain predecessor rows"
5088 );
5089 assert!(
5090 !defer_generated_schema_application_for_prepared_migration(&db, &proposal)
5091 .expect("terminal aborted record must not block generated startup"),
5092 );
5093 }
5094
5095 #[cfg(feature = "migration")]
5096 #[test]
5097 #[expect(
5098 clippy::too_many_lines,
5099 reason = "the interrupted rewrite, recovery, final proof, and publication form one scenario"
5100 )]
5101 fn physical_migration_rewrite_recovers_and_publishes_one_complete_candidate() {
5102 use super::{
5103 defer_generated_schema_application_for_prepared_migration, migrate_schema,
5104 schema_migration_status_for_target,
5105 };
5106 use crate::db::{
5107 data::{CanonicalSlotReader, DecodedDataStoreKey, StoreVisit, StructuralSlotReader},
5108 schema::{
5109 MigrationRewriteInterruption, SchemaMigrationCommand, SchemaMigrationPhase,
5110 ensure_schema_migration_ready_for_ordinary_operations,
5111 interrupt_next_migration_rewrite_at,
5112 },
5113 };
5114 use crate::error::InternalError;
5115
5116 let db = Db::<MigrationExecutionCanister>::new(
5117 &MIGRATION_EXECUTION_REGISTRY,
5118 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5119 );
5120 drive_startup_recovery_to_completion(&db);
5121 let initial_target = schema_application_target(&db).expect("initial target should issue");
5122 let store_identity = initial_target
5123 .stores()
5124 .first()
5125 .expect("migration execution store should exist")
5126 .identity();
5127 let initial = validation_migration_proposal(
5128 ValidationMigrationShape::Clean,
5129 false,
5130 initial_target.accepted_head().clone(),
5131 initial_target.database_identity(),
5132 store_identity,
5133 );
5134 apply_schema(&db, &initial).expect("initial schema should publish");
5135 let session = DbSession::<MigrationExecutionCanister>::new(
5136 &MIGRATION_EXECUTION_REGISTRY,
5137 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5138 );
5139 for (id, value) in [(1, 7), (2, 8), (3, 9)] {
5140 session
5141 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5142 entity: "MigratingItem".to_string(),
5143 patch: DynamicStructuralPatch::new(vec![
5144 (
5145 "id".to_string(),
5146 DynamicWriteCell::Value(InputValue::nat64(id)),
5147 ),
5148 (
5149 "old_value".to_string(),
5150 DynamicWriteCell::Value(InputValue::int64(value)),
5151 ),
5152 ]),
5153 })
5154 .expect("predecessor row should insert");
5155 }
5156 let target = schema_application_target(&db).expect("migration target should issue");
5157 let proposal = validation_migration_proposal(
5158 ValidationMigrationShape::Clean,
5159 true,
5160 target.accepted_head().clone(),
5161 target.database_identity(),
5162 store_identity,
5163 );
5164 let plan = proposal
5165 .migration()
5166 .expect("migration plan should exist")
5167 .digest();
5168 let command = || SchemaMigrationCommand::Advance {
5169 expected_database: target.database_identity(),
5170 expected_head: target.accepted_head().clone(),
5171 expected_plan: plan,
5172 acknowledged_finding_page: None,
5173 };
5174 for expected in [
5175 SchemaMigrationPhase::Prepared,
5176 SchemaMigrationPhase::Validating,
5177 SchemaMigrationPhase::ReadyToRewrite,
5178 SchemaMigrationPhase::RewritingRows,
5179 ] {
5180 assert_eq!(
5181 migrate_schema(&db, &proposal, command())
5182 .expect("migration phase should advance")
5183 .phase(),
5184 expected,
5185 );
5186 }
5187
5188 for interruption in [
5189 MigrationRewriteInterruption::MarkerPersisted,
5190 MigrationRewriteInterruption::JournalPublished,
5191 MigrationRewriteInterruption::PhysicalApplied,
5192 ] {
5193 interrupt_next_migration_rewrite_at(interruption);
5194 migrate_schema(&db, &proposal, command())
5195 .expect_err("injected interruption should retain the rewrite marker");
5196
5197 forget_recovered_domain_for_tests(&db)
5198 .expect("upgrade should reset recovery ownership");
5199 drive_startup_recovery_to_completion(&db);
5200 }
5201
5202 let rebuilding = schema_migration_status_for_target(
5203 &db,
5204 &proposal,
5205 &schema_application_target(&db).expect("recovered target should issue"),
5206 )
5207 .expect("recovered status should remain readable");
5208 assert_eq!(rebuilding.phase(), SchemaMigrationPhase::RebuildingIndexes);
5209 assert_eq!(rebuilding.rows_rewritten(), 3);
5210 assert_eq!(
5211 migrate_schema(&db, &proposal, command())
5212 .expect("derived generations should complete")
5213 .phase(),
5214 SchemaMigrationPhase::FinalValidation,
5215 );
5216 assert_eq!(
5217 migrate_schema(&db, &proposal, command())
5218 .expect("final validation should complete")
5219 .phase(),
5220 SchemaMigrationPhase::Publishing,
5221 );
5222 let applied = migrate_schema(&db, &proposal, command())
5223 .expect("candidate publication should complete atomically");
5224 assert_eq!(applied.phase(), SchemaMigrationPhase::Applied);
5225 assert_eq!(applied.rows_rewritten(), 3);
5226 assert_eq!(applied.indexes_rebuilt(), 1);
5227 assert_ne!(applied.accepted_head(), target.accepted_head());
5228 let terminal_target = schema_application_target(&db).expect("terminal target should issue");
5229 let terminal_proposal = validation_migration_proposal(
5230 ValidationMigrationShape::Clean,
5231 true,
5232 terminal_target.accepted_head().clone(),
5233 terminal_target.database_identity(),
5234 store_identity,
5235 );
5236 assert!(
5237 !defer_generated_schema_application_for_prepared_migration(&db, &terminal_proposal,)
5238 .expect("terminal record must not block generated startup"),
5239 );
5240
5241 let store = db
5242 .store_handle(MIGRATION_EXECUTION_STORE_PATH)
5243 .expect("migration execution store should resolve");
5244 let runtime = db
5245 .accepted_runtime_entity_for_path("MigratingItem")
5246 .expect("published candidate entity should resolve");
5247 let selection = store
5248 .with_schema(|schema| {
5249 schema.current_accepted_catalog_selection(
5250 runtime.entity_tag(),
5251 runtime.entity_path(),
5252 runtime.store_path(),
5253 )
5254 })
5255 .expect("candidate selection should remain readable")
5256 .expect("candidate selection should exist");
5257 let contract = crate::db::data::AcceptedStructuralRowAuthority::from_catalog_selection(
5258 runtime.entity_path(),
5259 &selection,
5260 )
5261 .expect("candidate row authority should compile")
5262 .into_row_contract();
5263 let mut values = Vec::new();
5264 store
5265 .with_data(|data| {
5266 data.visit_entries(|key, row| {
5267 let decoded = DecodedDataStoreKey::try_from_raw(key)
5268 .expect("rewritten key should decode");
5269 let reader =
5270 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
5271 row, &contract,
5272 )
5273 .expect("rewritten row should use the candidate layout");
5274 reader
5275 .validate_primary_key(&decoded)
5276 .expect("rewritten row and key should remain bound");
5277 values.push(
5278 reader
5279 .required_value_by_contract(1)
5280 .expect("candidate value slot should decode"),
5281 );
5282 Ok::<StoreVisit, InternalError>(StoreVisit::Continue)
5283 })
5284 })
5285 .expect("rewritten row scan should complete");
5286 assert_eq!(
5287 values,
5288 vec![
5289 crate::value::Value::Nat64(7),
5290 crate::value::Value::Nat64(8),
5291 crate::value::Value::Nat64(9),
5292 ],
5293 );
5294 assert_eq!(store.with_index(IndexStore::len), 3);
5295 let accepted = store
5296 .with_schema(SchemaStore::current_accepted_schema_bundle)
5297 .expect("published candidate bundle should remain readable")
5298 .expect("published candidate bundle should exist");
5299 let entity_source = EntitySourceKey::try_new("MigratingItem")
5300 .expect("migration entity source should admit");
5301 let entity_tag = accepted
5302 .source_bindings_for_tests()
5303 .entity(&entity_source)
5304 .expect("candidate entity source should remain bound");
5305 let old_value =
5306 FieldSourceKey::try_new("old_value").expect("predecessor source should admit");
5307 let current_value =
5308 FieldSourceKey::try_new("value").expect("candidate source should admit");
5309 assert_eq!(
5310 accepted
5311 .source_bindings_for_tests()
5312 .field(entity_tag, &old_value),
5313 None,
5314 );
5315 assert!(
5316 accepted
5317 .source_bindings_for_tests()
5318 .field(entity_tag, ¤t_value)
5319 .is_some(),
5320 );
5321 ensure_schema_migration_ready_for_ordinary_operations()
5322 .expect("terminal publication must clear the database-wide gate");
5323 }
5324
5325 #[cfg(feature = "migration")]
5326 #[test]
5327 #[expect(
5328 clippy::too_many_lines,
5329 reason = "all four finding families share one ordered historical scan fixture"
5330 )]
5331 fn physical_migration_validation_reports_every_typed_finding_family_without_writes() {
5332 use std::convert::Infallible;
5333
5334 use super::migrate_schema;
5335 use crate::db::{
5336 data::StoreVisit,
5337 schema::{SchemaMigrationCommand, SchemaMigrationFindingKind, SchemaMigrationPhase},
5338 };
5339
5340 let db = Db::<MigrationFindingCanister>::new(
5341 &MIGRATION_FINDING_REGISTRY,
5342 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5343 );
5344 drive_startup_recovery_to_completion(&db);
5345 let initial_target = schema_application_target(&db).expect("initial target should issue");
5346 let store_identity = initial_target
5347 .stores()
5348 .first()
5349 .expect("migration finding store should exist")
5350 .identity();
5351 let initial = validation_migration_proposal(
5352 ValidationMigrationShape::AllFindingFamilies,
5353 false,
5354 initial_target.accepted_head().clone(),
5355 initial_target.database_identity(),
5356 store_identity,
5357 );
5358 apply_schema(&db, &initial).expect("initial finding schema should publish");
5359
5360 let session = DbSession::<MigrationFindingCanister>::new(
5361 &MIGRATION_FINDING_REGISTRY,
5362 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5363 );
5364 session
5365 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5366 entity: "MigrationTarget".to_string(),
5367 patch: DynamicStructuralPatch::new(vec![(
5368 "id".to_string(),
5369 DynamicWriteCell::Value(InputValue::nat64(7)),
5370 )]),
5371 })
5372 .expect("relation target should insert");
5373 for (id, value) in [(1, 9), (2, 8), (3, 7), (4, 7), (5, 300)] {
5374 session
5375 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5376 entity: "MigratingItem".to_string(),
5377 patch: DynamicStructuralPatch::new(vec![
5378 (
5379 "id".to_string(),
5380 DynamicWriteCell::Value(InputValue::nat64(id)),
5381 ),
5382 (
5383 "old_value".to_string(),
5384 DynamicWriteCell::Value(InputValue::int64(value)),
5385 ),
5386 ]),
5387 })
5388 .expect("predecessor finding row should insert");
5389 }
5390 let store = db
5391 .store_handle(MIGRATION_FINDING_STORE_PATH)
5392 .expect("migration finding store should resolve");
5393 let row_bytes = || {
5394 store.with_data(|data| {
5395 let mut rows = Vec::new();
5396 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
5397 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
5398 Ok(StoreVisit::Continue)
5399 });
5400 result.expect("infallible row visit should complete");
5401 rows
5402 })
5403 };
5404 let before_rows = row_bytes();
5405
5406 let target = schema_application_target(&db).expect("migration target should issue");
5407 let proposal = validation_migration_proposal(
5408 ValidationMigrationShape::AllFindingFamilies,
5409 true,
5410 target.accepted_head().clone(),
5411 target.database_identity(),
5412 store_identity,
5413 );
5414 let plan = proposal
5415 .migration()
5416 .expect("migration plan should exist")
5417 .digest();
5418 let command = || SchemaMigrationCommand::Advance {
5419 expected_database: target.database_identity(),
5420 expected_head: target.accepted_head().clone(),
5421 expected_plan: plan,
5422 acknowledged_finding_page: None,
5423 };
5424 assert_eq!(
5425 migrate_schema(&db, &proposal, command())
5426 .expect("finding migration should prepare")
5427 .phase(),
5428 SchemaMigrationPhase::Prepared,
5429 );
5430 assert_eq!(
5431 migrate_schema(&db, &proposal, command())
5432 .expect("finding migration should enter validation")
5433 .phase(),
5434 SchemaMigrationPhase::Validating,
5435 );
5436 let rejected =
5437 migrate_schema(&db, &proposal, command()).expect("validation should report findings");
5438 assert_eq!(rejected.phase(), SchemaMigrationPhase::Rejected);
5439 assert_eq!(rejected.rows_validated(), 5);
5440 assert_eq!(
5441 rejected
5442 .findings()
5443 .iter()
5444 .map(crate::db::schema::SchemaMigrationFinding::kind)
5445 .collect::<Vec<_>>(),
5446 vec![
5447 SchemaMigrationFindingKind::Constraint,
5448 SchemaMigrationFindingKind::Relation,
5449 SchemaMigrationFindingKind::UniqueIndex,
5450 SchemaMigrationFindingKind::Transform,
5451 ],
5452 );
5453 assert_eq!(
5454 row_bytes(),
5455 before_rows,
5456 "rejected validation must not rewrite accepted rows"
5457 );
5458 assert_eq!(
5459 store.with_index(IndexStore::len),
5460 0,
5461 "a rejected page must not publish any staged generation"
5462 );
5463 }
5464
5465 #[cfg(feature = "migration")]
5466 #[test]
5467 fn exact_migration_retry_binds_the_terminal_head_not_the_predecessor_head() {
5468 use super::exact_migration_replay_target;
5469
5470 let db = Db::<EvolutionCanister>::new(
5471 &EVOLUTION_REGISTRY,
5472 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5473 );
5474 drive_startup_recovery_to_completion(&db);
5475 let initial_target = schema_application_target(&db).expect("initial target should issue");
5476 let (proposal, _, _) = generated_check_proposal(
5477 initial_target.accepted_head().clone(),
5478 "migration-retry-initial",
5479 false,
5480 initial_target.database_identity(),
5481 initial_target
5482 .stores()
5483 .first()
5484 .expect("test store should exist")
5485 .identity(),
5486 );
5487 apply_schema(&db, &proposal).expect("initial schema should publish");
5488 let current_target = schema_application_target(&db).expect("current target should issue");
5489 assert_ne!(
5490 current_target.accepted_head(),
5491 initial_target.accepted_head(),
5492 );
5493
5494 let receipt = SchemaChangeReceipt::new(
5495 current_target.database_identity(),
5496 SchemaSubmissionKey::try_new("migration/retry")
5497 .expect("migration submission should admit"),
5498 SchemaProposalDigest::from_bytes([0x77; 32]),
5499 initial_target.accepted_head().clone(),
5500 SchemaChangeOutcome::Applied {
5501 accepted_head: current_target.accepted_head().clone(),
5502 },
5503 )
5504 .expect("terminal migration receipt should admit");
5505 let record = SchemaApplicationRecord::new(receipt, Vec::new())
5506 .expect("terminal migration record should admit");
5507
5508 assert_eq!(
5509 exact_migration_replay_target(&db, current_target.database_identity(), &record,)
5510 .expect("exact retry should resolve the terminal target"),
5511 current_target,
5512 );
5513 }
5514}