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::Unadopted => SchemaMigrationCode::Unadopted,
1884 SchemaMigrationPlanningError::MissingMigration => SchemaMigrationCode::MissingMigration,
1885 SchemaMigrationPlanningError::VersionGap => SchemaMigrationCode::VersionGap,
1886 SchemaMigrationPlanningError::Downgrade => SchemaMigrationCode::Downgrade,
1887 SchemaMigrationPlanningError::EmptyEntityVersionBump => {
1888 SchemaMigrationCode::EmptyEntityVersionBump
1889 }
1890 SchemaMigrationPlanningError::StaleAcceptedHead => SchemaMigrationCode::StaleAcceptedHead,
1891 SchemaMigrationPlanningError::UnknownFromObject => SchemaMigrationCode::UnknownFromObject,
1892 SchemaMigrationPlanningError::UnknownToObject => SchemaMigrationCode::UnknownToObject,
1893 SchemaMigrationPlanningError::KindMismatch => SchemaMigrationCode::KindMismatch,
1894 SchemaMigrationPlanningError::IdentityConflict => SchemaMigrationCode::IdentityConflict,
1895 SchemaMigrationPlanningError::UnexplainedSchemaDifference => {
1896 SchemaMigrationCode::UnexplainedSchemaDifference
1897 }
1898 SchemaMigrationPlanningError::UnsupportedTransform => {
1899 SchemaMigrationCode::UnsupportedTransform
1900 }
1901 SchemaMigrationPlanningError::RekeyedCatalogInvalid
1902 | SchemaMigrationPlanningError::CandidateMismatch => SchemaMigrationCode::CandidateMismatch,
1903 SchemaMigrationPlanningError::CorruptLineage => SchemaMigrationCode::ProgressCorrupt,
1904 };
1905 InternalError::schema_migration(reason)
1906}
1907
1908#[cfg(feature = "migration")]
1909fn migration_submission_key(
1910 plan_digest: Option<SchemaMigrationPlanDigest>,
1911) -> Result<SchemaSubmissionKey, InternalError> {
1912 let mut hasher = new_hash_sha256_prefixed(SCHEMA_MIGRATION_SUBMISSION_PROFILE);
1913 match plan_digest {
1914 None => write_hash_tag_u8(&mut hasher, 0),
1915 Some(digest) => {
1916 write_hash_tag_u8(&mut hasher, 1);
1917 hasher.update(digest.to_bytes());
1918 }
1919 }
1920 let digest = finalize_hash_sha256(hasher);
1921 let mut encoded = String::with_capacity(80);
1922 encoded.push_str("migration/");
1923 for byte in digest {
1924 use std::fmt::Write as _;
1925 write!(&mut encoded, "{byte:02x}").map_err(|_| InternalError::store_invariant())?;
1926 }
1927 SchemaSubmissionKey::try_new(encoded).map_err(|_| InternalError::store_invariant())
1928}
1929
1930#[cfg(feature = "migration")]
1931fn load_migration_record_for_status(
1932 database_identity: TargetDatabaseIdentity,
1933 submission_key: &SchemaSubmissionKey,
1934) -> Result<Option<SchemaApplicationRecord>, InternalError> {
1935 let record =
1936 with_schema_application_store(|store| store.load(database_identity, submission_key))?;
1937 if record
1938 .as_ref()
1939 .is_some_and(|record| record.receipt().database_identity() != database_identity)
1940 {
1941 return Err(InternalError::schema_migration(
1942 SchemaMigrationCode::ProgressCorrupt,
1943 ));
1944 }
1945 Ok(record)
1946}
1947
1948#[cfg(feature = "migration")]
1949fn load_exact_migration_record(
1950 database_identity: TargetDatabaseIdentity,
1951 submission_key: &SchemaSubmissionKey,
1952 proposal_digest: SchemaProposalDigest,
1953 prior_head: &ExpectedAcceptedHead,
1954) -> Result<Option<SchemaApplicationRecord>, InternalError> {
1955 let Some(record) =
1956 with_schema_application_store(|store| store.load(database_identity, submission_key))?
1957 else {
1958 return Ok(None);
1959 };
1960 if !record.receipt().is_exact_submission(
1961 database_identity,
1962 submission_key,
1963 proposal_digest,
1964 prior_head,
1965 ) {
1966 return Err(InternalError::schema_migration(
1967 SchemaMigrationCode::PlanChanged,
1968 ));
1969 }
1970 Ok(Some(record))
1971}
1972
1973#[cfg(feature = "migration")]
1974fn public_migration_receipt(
1975 record: &SchemaApplicationRecord,
1976 plan_digest: Option<SchemaMigrationPlanDigest>,
1977) -> Result<SchemaMigrationReceipt, InternalError> {
1978 let accepted_head = migration_record_accepted_head(record)?.clone();
1979 Ok(SchemaMigrationReceipt::new(
1980 record.receipt().database_identity(),
1981 plan_digest,
1982 record.receipt().prior_head().clone(),
1983 accepted_head,
1984 ))
1985}
1986
1987#[cfg(feature = "migration")]
1988fn migration_record_accepted_head(
1989 record: &SchemaApplicationRecord,
1990) -> Result<&ExpectedAcceptedHead, InternalError> {
1991 match record.receipt().outcome() {
1992 SchemaChangeOutcome::NoOp { accepted_head }
1993 | SchemaChangeOutcome::Applied { accepted_head } => Ok(accepted_head),
1994 SchemaChangeOutcome::Pending { .. } | SchemaChangeOutcome::Aborted { .. } => Err(
1995 InternalError::schema_migration(SchemaMigrationCode::ProgressCorrupt),
1996 ),
1997 }
1998}
1999
2000#[cfg(feature = "migration")]
2001fn migration_transitions(
2002 proposal: &SchemaProposal,
2003) -> Result<Vec<SchemaMigrationEntityTransition>, InternalError> {
2004 if let Some(plan) = proposal.migration() {
2005 return plan
2006 .transitions()
2007 .iter()
2008 .map(|transition| {
2009 let target = proposal_entity(proposal, transition.entity())?;
2010 Ok(SchemaMigrationEntityTransition::new(
2011 transition.entity().clone(),
2012 Some(transition.from().get()),
2013 target.version().get(),
2014 ))
2015 })
2016 .collect();
2017 }
2018 proposal
2019 .fragments()
2020 .iter()
2021 .flat_map(icydb_schema::SchemaFragment::entities)
2022 .map(|entity| {
2023 Ok(SchemaMigrationEntityTransition::new(
2024 entity.source_key().clone(),
2025 None,
2026 entity.version().get(),
2027 ))
2028 })
2029 .collect()
2030}
2031
2032#[cfg(feature = "migration")]
2033fn proposal_entity<'a>(
2034 proposal: &'a SchemaProposal,
2035 source: &EntitySourceKey,
2036) -> Result<&'a icydb_schema::EntityFragment, InternalError> {
2037 proposal
2038 .fragments()
2039 .iter()
2040 .flat_map(icydb_schema::SchemaFragment::entities)
2041 .find(|entity| entity.source_key() == source)
2042 .ok_or_else(InternalError::store_invariant)
2043}
2044
2045#[cfg(feature = "migration")]
2046fn current_proposal_lineage_is_applied<C: CanisterKind>(
2047 db: &Db<C>,
2048 proposal: &SchemaProposal,
2049 accepted_head: &ExpectedAcceptedHead,
2050) -> Result<bool, InternalError> {
2051 let Some(lineage) = load_entity_source_lineage_catalog()? else {
2052 return Ok(false);
2053 };
2054 let entities = proposal
2055 .fragments()
2056 .iter()
2057 .flat_map(icydb_schema::SchemaFragment::entities)
2058 .collect::<Vec<_>>();
2059 if entities.len() != lineage.entries().len() {
2060 return Ok(false);
2061 }
2062 let target = proposal.target_database();
2063 let authorities = application_authorities(db);
2064 for entity in entities {
2065 let source = entity.source_key();
2066 let digest = proposal
2067 .entity_source_digest(source)
2068 .map_err(|_| InternalError::store_invariant())?;
2069 let mut matched = false;
2070 for authority in &authorities {
2071 let store_identity = derive_store_identity(target, authority);
2072 let entity_tag = authority
2073 .handle
2074 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)?
2075 .and_then(|bundle| bundle.source_bindings().entity(source));
2076 let Some(entity_tag) = entity_tag else {
2077 continue;
2078 };
2079 let Some(entry) = lineage.get(store_identity, entity_tag) else {
2080 return Ok(false);
2081 };
2082 matched = entry.accepted_head() == accepted_head
2083 && matches!(
2084 entry.state(),
2085 AcceptedEntitySourceLineageState::Adopted { version, source_digest }
2086 if version.get() == entity.version().get() && *source_digest == digest
2087 );
2088 break;
2089 }
2090 if !matched {
2091 return Ok(false);
2092 }
2093 }
2094 Ok(true)
2095}
2096
2097#[cfg(feature = "migration")]
2098fn preflight_unpublished_schema_migration<C: CanisterKind>(
2099 target: &SchemaApplicationTarget,
2100 proposal: &SchemaProposal,
2101 db: &Db<C>,
2102) -> Result<(), InternalError> {
2103 let authorities = application_authorities(db);
2104 let current_bundles = authorities
2105 .iter()
2106 .map(|authority| {
2107 authority
2108 .handle
2109 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)
2110 })
2111 .collect::<Result<Vec<_>, InternalError>>()?;
2112 let stores = authorities
2113 .iter()
2114 .zip(¤t_bundles)
2115 .filter_map(|(authority, bundle)| {
2116 bundle.as_ref().map(|bundle| ExistingProposalStore {
2117 path: authority.path,
2118 identity: derive_store_identity(target.database_identity(), authority),
2119 bundle,
2120 })
2121 })
2122 .collect::<Vec<_>>();
2123 let lineage = load_entity_source_lineage_catalog()?.unwrap_or_default();
2124 let planned = plan_schema_migration(proposal, stores.as_slice(), &lineage)
2125 .map_err(schema_migration_planning_error)?;
2126 if planned.candidates().is_empty() || planned.lineage().is_empty() {
2127 return Err(InternalError::store_invariant());
2128 }
2129 for next in planned.lineage() {
2130 let current = lineage
2131 .get(next.store(), next.entity())
2132 .ok_or_else(InternalError::store_invariant)?;
2133 let AcceptedEntitySourceLineageState::Adopted {
2134 version,
2135 source_digest,
2136 } = current.state()
2137 else {
2138 return Err(InternalError::store_invariant());
2139 };
2140 let expected_version = version
2141 .get()
2142 .checked_add(1)
2143 .ok_or_else(InternalError::store_invariant)?;
2144 if next.version().get() != expected_version || next.digest() == *source_digest {
2145 return Err(InternalError::store_invariant());
2146 }
2147 }
2148 Ok(())
2149}
2150
2151fn lower_application_candidates<const ALLOW_REMOVALS: bool>(
2152 target: &SchemaApplicationTarget,
2153 proposal: &SchemaProposal,
2154 authorities: &[StoreApplicationAuthority],
2155) -> Result<LoweredApplication, InternalError> {
2156 let current_bundles = authorities
2157 .iter()
2158 .map(|authority| {
2159 authority
2160 .handle
2161 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_bundle)
2162 })
2163 .collect::<Result<Vec<_>, InternalError>>()?;
2164 let initial_application = matches!(target.accepted_head(), ExpectedAcceptedHead::Empty);
2165 let mut candidates = match target.accepted_head() {
2166 ExpectedAcceptedHead::Empty => {
2167 let stores = authorities
2168 .iter()
2169 .map(|authority| ProposalStoreTarget {
2170 path: authority.path,
2171 identity: derive_store_identity(target.database_identity(), authority),
2172 })
2173 .collect::<Vec<_>>();
2174 let candidates = lower_initial_schema_proposal(proposal, stores.as_slice())?;
2175 #[cfg(feature = "migration")]
2176 {
2177 let planned = plan_initial_entity_source_lineage(proposal, &candidates)
2178 .map_err(schema_migration_planning_error)?;
2179 if planned.len()
2180 != proposal
2181 .fragments()
2182 .iter()
2183 .map(|fragment| fragment.entities().len())
2184 .sum::<usize>()
2185 {
2186 return Err(InternalError::store_invariant());
2187 }
2188 }
2189 candidates
2190 }
2191 ExpectedAcceptedHead::Exact { .. }
2192 if proposal.fragments().is_empty() && proposal.removals().is_empty() =>
2193 {
2194 Vec::new()
2195 }
2196 ExpectedAcceptedHead::Exact { .. } => {
2197 let stores = authorities
2198 .iter()
2199 .zip(¤t_bundles)
2200 .filter_map(|(authority, bundle)| {
2201 bundle.as_ref().map(|bundle| ExistingProposalStore {
2202 path: authority.path,
2203 identity: derive_store_identity(target.database_identity(), authority),
2204 bundle,
2205 })
2206 })
2207 .collect::<Vec<_>>();
2208 if ALLOW_REMOVALS {
2209 lower_existing_schema_proposal(proposal, stores.as_slice())?
2210 } else {
2211 lower_generated_existing_schema_proposal(proposal, stores.as_slice())?
2212 }
2213 }
2214 };
2215 let pending = if initial_application {
2216 preflight_initial_application(authorities, &candidates)?;
2217 None
2218 } else {
2219 preflight_existing_application(authorities, ¤t_bundles, &mut candidates)?
2220 };
2221 Ok(LoweredApplication {
2222 current_bundles,
2223 candidates,
2224 pending,
2225 })
2226}
2227
2228fn validate_database_identity_state_capacity(
2229 authorities: &[StoreApplicationAuthority],
2230 candidates: &[CandidateSchemaRevision],
2231 incarnation: crate::db::integrity::DatabaseIncarnationId,
2232) -> Result<(), InternalError> {
2233 let mut total = 0usize;
2234 for authority in authorities {
2235 let count = match candidates
2236 .iter()
2237 .find(|candidate| candidate.store_path() == authority.path)
2238 {
2239 Some(candidate) => authority.handle.with_schema(|store| {
2240 store.projected_identity_state_count(incarnation, candidate)
2241 })?,
2242 None => authority
2243 .handle
2244 .with_schema(|store| store.identity_state_inventory_for_integrity(incarnation))?
2245 .len(),
2246 };
2247 total = include_identity_state_count(total, count)?;
2248 }
2249 Ok(())
2250}
2251
2252fn include_identity_state_count(total: usize, count: usize) -> Result<usize, InternalError> {
2253 let total = total
2254 .checked_add(count)
2255 .ok_or_else(InternalError::identity_state_capacity_exhausted)?;
2256 if total > MAX_IDENTITY_STATE_RECORDS_PER_DATABASE {
2257 return Err(InternalError::identity_state_capacity_exhausted());
2258 }
2259 Ok(total)
2260}
2261
2262fn preflight_initial_application(
2263 authorities: &[StoreApplicationAuthority],
2264 candidates: &[crate::db::schema::CandidateSchemaRevision],
2265) -> Result<(), InternalError> {
2266 for candidate in candidates {
2267 let authority = authorities
2268 .iter()
2269 .find(|authority| authority.path == candidate.store_path())
2270 .ok_or_else(InternalError::store_invariant)?;
2271 if authority.handle.with_data(DataStore::len) != 0
2272 || authority.handle.index_state() != IndexState::Ready
2273 || !authority.handle.with_index(IndexStore::is_empty)
2274 {
2275 return Err(InternalError::store_unsupported());
2276 }
2277 }
2278 Ok(())
2279}
2280
2281fn preflight_existing_application(
2288 authorities: &[StoreApplicationAuthority],
2289 current_bundles: &[Option<crate::db::schema::AcceptedSchemaRevisionBundle>],
2290 candidates: &mut [CandidateSchemaRevision],
2291) -> Result<Option<PendingGeneratedRowLocalConstraint>, InternalError> {
2292 require_empty_physical_entity_removal(authorities, current_bundles, candidates)?;
2293 require_empty_physical_field_removals(authorities, current_bundles, candidates)?;
2294 require_empty_physical_index_removals(authorities, current_bundles, candidates)?;
2295 require_empty_physical_relation_removals(authorities, current_bundles, candidates)?;
2296 let proofs = generated_row_local_constraint_proofs(authorities, current_bundles, candidates)?;
2297 if proofs
2298 .iter()
2299 .filter(|proof| proof.historical_rows != 0)
2300 .count()
2301 > 1
2302 {
2303 return Err(InternalError::store_unsupported());
2304 }
2305
2306 let mut pending = None;
2307 for candidate_index in 0..candidates.len() {
2308 let candidate = candidates
2309 .get(candidate_index)
2310 .cloned()
2311 .ok_or_else(InternalError::store_invariant)?;
2312 let candidate_proofs = proofs
2313 .iter()
2314 .filter(|proof| proof.candidate_index == candidate_index)
2315 .collect::<Vec<_>>();
2316 if candidate_proofs.is_empty() {
2317 continue;
2318 }
2319
2320 let mut snapshots = candidate.bundle().entity_snapshots().clone();
2321 for proof in candidate_proofs {
2322 let mut promote = true;
2323 if proof.historical_rows != 0 {
2324 match validate_unpublished_row_local_candidate_bounded(
2325 proof.store,
2326 proof.store_path,
2327 proof.entity_tag,
2328 proof.entity_path.as_str(),
2329 &candidate,
2330 proof.constraint_id,
2331 )? {
2332 UnpublishedRowLocalValidation::Complete { .. } => {}
2333 UnpublishedRowLocalValidation::Incomplete => {
2334 if proof.store.storage_capabilities().recovery()
2335 != StoreRecoveryCapability::StableBasePlusJournalReplay
2336 || pending.is_some()
2337 {
2338 return Err(InternalError::store_unsupported());
2339 }
2340 pending = Some(PendingGeneratedRowLocalConstraint {
2341 proof: (*proof).clone(),
2342 });
2343 promote = false;
2344 }
2345 }
2346 }
2347 if !promote {
2348 continue;
2349 }
2350 let snapshot = snapshots
2351 .get(&proof.entity_tag)
2352 .cloned()
2353 .ok_or_else(InternalError::store_invariant)?;
2354 let catalog = snapshot
2355 .constraint_catalog()
2356 .clone()
2357 .with_directly_validated_activation(proof.constraint_id)
2358 .map_err(|_| InternalError::store_invariant())?;
2359 snapshots.insert(proof.entity_tag, snapshot.with_constraint_catalog(catalog));
2360 }
2361 let bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
2362 candidate.revision(),
2363 candidate.bundle().store_path(),
2364 candidate.bundle().enum_catalog().clone(),
2365 candidate.bundle().composite_catalog().clone(),
2366 candidate.bundle().source_bindings().clone(),
2367 snapshots,
2368 )?;
2369 candidates[candidate_index] = CandidateSchemaRevision::new(bundle)?;
2370 }
2371 Ok(pending)
2372}
2373
2374fn require_empty_physical_entity_removal(
2381 authorities: &[StoreApplicationAuthority],
2382 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2383 candidates: &[CandidateSchemaRevision],
2384) -> Result<(), InternalError> {
2385 let mut removed_entity = None;
2386 for candidate in candidates {
2387 let (position, source_authority) = authorities
2388 .iter()
2389 .enumerate()
2390 .find(|(_, authority)| authority.path == candidate.store_path())
2391 .ok_or_else(InternalError::store_invariant)?;
2392 let current = current_bundles
2393 .get(position)
2394 .and_then(Option::as_ref)
2395 .ok_or_else(InternalError::store_invariant)?;
2396 let removed = current
2397 .entity_snapshots()
2398 .iter()
2399 .filter(|(entity_tag, _)| {
2400 !candidate
2401 .bundle()
2402 .entity_snapshots()
2403 .contains_key(entity_tag)
2404 })
2405 .collect::<Vec<_>>();
2406 if removed.is_empty() {
2407 continue;
2408 }
2409 let [(entity_tag, snapshot)] = removed.as_slice() else {
2410 return Err(InternalError::store_unsupported());
2411 };
2412 let entity_tag = **entity_tag;
2413 let snapshot = *snapshot;
2414 if removed_entity.is_some()
2415 || current.entity_snapshots().len()
2416 != candidate
2417 .bundle()
2418 .entity_snapshots()
2419 .len()
2420 .saturating_add(1)
2421 {
2422 return Err(InternalError::store_unsupported());
2423 }
2424 require_exact_empty_entity(source_authority.handle, entity_tag)?;
2425 source_authority
2426 .handle
2427 .with_index(|store| prove_empty_user_index_domain(store, entity_tag))
2428 .map_err(StagedUserIndexDomainError::into_internal_error)?;
2429 for relation in snapshot.relations() {
2430 let target_store = accepted_entity_store_for_path(
2431 authorities,
2432 current_bundles,
2433 relation.target_path(),
2434 )?;
2435 target_store.with_index(|store| {
2436 prove_empty_reverse_relation_domain(store, entity_tag, snapshot, relation)
2437 })?;
2438 }
2439 removed_entity = Some(snapshot.entity_path());
2440 }
2441
2442 let Some(removed_path) = removed_entity else {
2443 return Ok(());
2444 };
2445 for (position, authority) in authorities.iter().enumerate() {
2446 let after = candidates
2447 .iter()
2448 .find(|candidate| candidate.store_path() == authority.path)
2449 .map(CandidateSchemaRevision::bundle)
2450 .or_else(|| current_bundles.get(position).and_then(Option::as_ref));
2451 let Some(after) = after else {
2452 continue;
2453 };
2454 if after
2455 .entity_snapshots()
2456 .values()
2457 .flat_map(crate::db::schema::PersistedSchemaSnapshot::relations)
2458 .any(|relation| relation.target_path() == removed_path)
2459 {
2460 return Err(InternalError::store_unsupported());
2461 }
2462 }
2463 Ok(())
2464}
2465
2466fn require_empty_physical_relation_removals(
2469 authorities: &[StoreApplicationAuthority],
2470 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2471 candidates: &[CandidateSchemaRevision],
2472) -> Result<(), InternalError> {
2473 for candidate in candidates {
2474 let (position, source_authority) = authorities
2475 .iter()
2476 .enumerate()
2477 .find(|(_, authority)| authority.path == candidate.store_path())
2478 .ok_or_else(InternalError::store_invariant)?;
2479 let current = current_bundles
2480 .get(position)
2481 .and_then(Option::as_ref)
2482 .ok_or_else(InternalError::store_invariant)?;
2483 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2484 let before = current
2485 .entity_snapshots()
2486 .get(entity_tag)
2487 .ok_or_else(InternalError::store_invariant)?;
2488 let removed = before
2489 .relations()
2490 .iter()
2491 .filter(|relation| {
2492 !after
2493 .relations()
2494 .iter()
2495 .any(|candidate| candidate.id() == relation.id())
2496 })
2497 .collect::<Vec<_>>();
2498 if removed.is_empty() {
2499 continue;
2500 }
2501 let added = after.relations().iter().any(|relation| {
2502 !before
2503 .relations()
2504 .iter()
2505 .any(|accepted| accepted.id() == relation.id())
2506 });
2507 let [removed] = removed.as_slice() else {
2508 return Err(InternalError::store_unsupported());
2509 };
2510 if added || before.relations().len() != after.relations().len().saturating_add(1) {
2511 return Err(InternalError::store_unsupported());
2512 }
2513 require_exact_empty_entity(source_authority.handle, *entity_tag)?;
2514 let target_store = accepted_entity_store_for_path(
2515 authorities,
2516 current_bundles,
2517 removed.target_path(),
2518 )?;
2519 target_store.with_index(|store| {
2520 prove_empty_reverse_relation_domain(store, *entity_tag, before, removed)
2521 })?;
2522 }
2523 }
2524 Ok(())
2525}
2526
2527fn accepted_entity_store_for_path(
2528 authorities: &[StoreApplicationAuthority],
2529 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2530 entity_path: &str,
2531) -> Result<StoreHandle, InternalError> {
2532 let mut resolved = None;
2533 for (position, bundle) in current_bundles.iter().enumerate() {
2534 let Some(bundle) = bundle else {
2535 continue;
2536 };
2537 if !bundle
2538 .entity_snapshots()
2539 .values()
2540 .any(|snapshot| snapshot.entity_path() == entity_path)
2541 {
2542 continue;
2543 }
2544 if resolved.is_some() {
2545 return Err(InternalError::store_invariant());
2546 }
2547 resolved = authorities.get(position).map(|authority| authority.handle);
2548 }
2549 resolved.ok_or_else(InternalError::store_unsupported)
2550}
2551
2552fn require_empty_physical_index_removals(
2556 authorities: &[StoreApplicationAuthority],
2557 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2558 candidates: &[CandidateSchemaRevision],
2559) -> Result<(), InternalError> {
2560 for candidate in candidates {
2561 let (position, authority) = authorities
2562 .iter()
2563 .enumerate()
2564 .find(|(_, authority)| authority.path == candidate.store_path())
2565 .ok_or_else(InternalError::store_invariant)?;
2566 let current = current_bundles
2567 .get(position)
2568 .and_then(Option::as_ref)
2569 .ok_or_else(InternalError::store_invariant)?;
2570 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2571 let before = current
2572 .entity_snapshots()
2573 .get(entity_tag)
2574 .ok_or_else(InternalError::store_invariant)?;
2575 if before.indexes().len() == after.indexes().len() {
2576 continue;
2577 }
2578 if before.indexes().len() != after.indexes().len().saturating_add(1) {
2579 return Err(InternalError::store_unsupported());
2580 }
2581 require_exact_empty_entity(authority.handle, *entity_tag)?;
2582 authority
2583 .handle
2584 .with_index(|store| prove_empty_user_index_domain(store, *entity_tag))
2585 .map_err(StagedUserIndexDomainError::into_internal_error)?;
2586 }
2587 }
2588 Ok(())
2589}
2590
2591fn require_empty_physical_field_removals(
2594 authorities: &[StoreApplicationAuthority],
2595 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2596 candidates: &[CandidateSchemaRevision],
2597) -> Result<(), InternalError> {
2598 for candidate in candidates {
2599 let (position, authority) = authorities
2600 .iter()
2601 .enumerate()
2602 .find(|(_, authority)| authority.path == candidate.store_path())
2603 .ok_or_else(InternalError::store_invariant)?;
2604 let current = current_bundles
2605 .get(position)
2606 .and_then(Option::as_ref)
2607 .ok_or_else(InternalError::store_invariant)?;
2608 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2609 let before = current
2610 .entity_snapshots()
2611 .get(entity_tag)
2612 .ok_or_else(InternalError::store_invariant)?;
2613 if before.row_layout() == after.row_layout() {
2614 continue;
2615 }
2616 if before.fields().len() != after.fields().len().saturating_add(1) {
2617 return Err(InternalError::store_unsupported());
2618 }
2619 require_exact_empty_entity(authority.handle, *entity_tag)?;
2620 }
2621 }
2622 Ok(())
2623}
2624
2625fn require_exact_empty_entity(
2629 store: StoreHandle,
2630 entity_tag: EntityTag,
2631) -> Result<(), InternalError> {
2632 require_exact_empty_entity_count(store.exact_entity_count(entity_tag))
2633}
2634
2635fn require_exact_empty_entity_count(count: Option<u64>) -> Result<(), InternalError> {
2636 let count = count.ok_or_else(InternalError::store_corruption)?;
2637 if count != 0 {
2638 return Err(InternalError::store_unsupported());
2639 }
2640
2641 Ok(())
2642}
2643
2644fn generated_row_local_constraint_proofs(
2645 authorities: &[StoreApplicationAuthority],
2646 current_bundles: &[Option<AcceptedSchemaRevisionBundle>],
2647 candidates: &[CandidateSchemaRevision],
2648) -> Result<Vec<DirectGeneratedRowLocalProof>, InternalError> {
2649 let mut proofs = Vec::new();
2650 for (candidate_index, candidate) in candidates.iter().enumerate() {
2651 let (position, authority) = authorities
2652 .iter()
2653 .enumerate()
2654 .find(|(_, authority)| authority.path == candidate.store_path())
2655 .ok_or_else(InternalError::store_invariant)?;
2656 let current = current_bundles
2657 .get(position)
2658 .and_then(Option::as_ref)
2659 .ok_or_else(InternalError::store_invariant)?;
2660 for (entity_tag, after) in candidate.bundle().entity_snapshots() {
2661 let before = current
2662 .entity_snapshots()
2663 .get(entity_tag)
2664 .ok_or_else(InternalError::store_invariant)?;
2665 for constraint_id in added_generated_row_local_activations(before, after) {
2666 let historical_rows = authority
2667 .handle
2668 .exact_entity_count(*entity_tag)
2669 .ok_or_else(InternalError::store_corruption)?;
2670 proofs.push(DirectGeneratedRowLocalProof {
2671 candidate_index,
2672 store: authority.handle,
2673 store_path: authority.path,
2674 entity_tag: *entity_tag,
2675 entity_path: after.entity_path().to_string(),
2676 constraint_id,
2677 historical_rows,
2678 });
2679 }
2680 }
2681 }
2682 Ok(proofs)
2683}
2684
2685fn added_generated_row_local_activations(
2686 before: &crate::db::schema::PersistedSchemaSnapshot,
2687 after: &crate::db::schema::PersistedSchemaSnapshot,
2688) -> Vec<ConstraintId> {
2689 after
2690 .constraint_activations()
2691 .iter()
2692 .filter(|candidate| {
2693 candidate.origin() == ConstraintOrigin::Generated
2694 && matches!(
2695 candidate.kind(),
2696 ConstraintActivationKind::Check { .. }
2697 | ConstraintActivationKind::TargetedRule { .. }
2698 )
2699 && !before
2700 .constraint_activations()
2701 .iter()
2702 .any(|accepted| accepted.id() == candidate.id())
2703 })
2704 .map(crate::db::schema::ConstraintActivationSnapshot::id)
2705 .collect()
2706}
2707
2708fn final_candidates_for_pending_row_local_constraint(
2709 candidates: &[CandidateSchemaRevision],
2710 pending: &PendingGeneratedRowLocalConstraint,
2711) -> Result<Vec<CandidateSchemaRevision>, InternalError> {
2712 let mut final_candidates = candidates.to_vec();
2713 let candidate = final_candidates
2714 .get(pending.proof.candidate_index)
2715 .cloned()
2716 .ok_or_else(InternalError::store_invariant)?;
2717 if candidate.store_path() != pending.proof.store_path {
2718 return Err(InternalError::store_invariant());
2719 }
2720 let mut snapshots = candidate.bundle().entity_snapshots().clone();
2721 let snapshot = snapshots
2722 .get(&pending.proof.entity_tag)
2723 .cloned()
2724 .ok_or_else(InternalError::store_invariant)?;
2725 let catalog = snapshot
2726 .constraint_catalog()
2727 .clone()
2728 .with_directly_validated_activation(pending.proof.constraint_id)
2729 .map_err(|_| InternalError::store_invariant())?;
2730 snapshots.insert(
2731 pending.proof.entity_tag,
2732 snapshot.with_constraint_catalog(catalog),
2733 );
2734 let final_revision = candidate
2735 .revision()
2736 .checked_next()
2737 .and_then(AcceptedSchemaRevision::checked_next)
2738 .ok_or_else(InternalError::store_unsupported)?;
2739 let bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
2740 final_revision,
2741 candidate.bundle().store_path(),
2742 candidate.bundle().enum_catalog().clone(),
2743 candidate.bundle().composite_catalog().clone(),
2744 candidate.bundle().source_bindings().clone(),
2745 snapshots,
2746 )?;
2747 final_candidates[pending.proof.candidate_index] = CandidateSchemaRevision::new(bundle)?;
2748 Ok(final_candidates)
2749}
2750
2751fn schema_change_progress_status(
2752 snapshot: &crate::db::schema::PersistedSchemaSnapshot,
2753 entity_tag: EntityTag,
2754 constraint_id: ConstraintId,
2755 progress: ConstraintValidationProgress,
2756) -> Result<SchemaChangeProgressStatus, InternalError> {
2757 match progress {
2758 ConstraintValidationProgress::Started => Ok(SchemaChangeProgressStatus::Started),
2759 ConstraintValidationProgress::Advanced {
2760 phase,
2761 rows_scanned,
2762 } => Ok(SchemaChangeProgressStatus::Advanced {
2763 phase: schema_change_validation_phase(phase),
2764 rows_scanned,
2765 }),
2766 ConstraintValidationProgress::Findings {
2767 receipt,
2768 phase,
2769 rows_scanned,
2770 } => {
2771 let activation = snapshot
2772 .constraint_catalog()
2773 .activation(constraint_id)
2774 .ok_or_else(InternalError::store_corruption)?;
2775 let fingerprint =
2776 crate::db::schema::accepted_schema_cache_fingerprint_for_persisted_snapshot(
2777 snapshot,
2778 )?;
2779 let findings = receipt
2780 .findings()
2781 .iter()
2782 .map(|finding| {
2783 constraint_validation_finding_output(
2784 fingerprint,
2785 entity_tag,
2786 activation,
2787 finding,
2788 )
2789 })
2790 .collect::<Result<Vec<_>, InternalError>>()?;
2791 Ok(SchemaChangeProgressStatus::Findings {
2792 phase: schema_change_validation_phase(phase),
2793 rows_scanned,
2794 page_sequence: receipt.page_sequence(),
2795 findings,
2796 })
2797 }
2798 ConstraintValidationProgress::Restarted { rows_scanned } => {
2799 Ok(SchemaChangeProgressStatus::Restarted { rows_scanned })
2800 }
2801 ConstraintValidationProgress::Promoted { .. } => Ok(SchemaChangeProgressStatus::Applied),
2802 }
2803}
2804
2805const fn schema_change_validation_phase(
2806 phase: ConstraintValidationPhase,
2807) -> SchemaChangeValidationPhase {
2808 match phase {
2809 ConstraintValidationPhase::Forward => SchemaChangeValidationPhase::Forward,
2810 ConstraintValidationPhase::Verify => SchemaChangeValidationPhase::Verify,
2811 }
2812}
2813
2814fn finalize_schema_application<C: CanisterKind>(
2815 db: &Db<C>,
2816 record: &SchemaApplicationRecord,
2817 candidate_head: &ExpectedAcceptedHead,
2818 status: SchemaChangeProgressStatus,
2819) -> Result<SchemaChangeProgress, InternalError> {
2820 if schema_application_target(db)?.accepted_head() != candidate_head {
2821 return Err(InternalError::schema_application_conflict());
2822 }
2823 let receipt = SchemaChangeReceipt::new(
2824 record.receipt().database_identity(),
2825 record.receipt().submission_key().clone(),
2826 record.receipt().proposal_digest(),
2827 record.receipt().prior_head().clone(),
2828 SchemaChangeOutcome::Applied {
2829 accepted_head: candidate_head.clone(),
2830 },
2831 )?;
2832 let terminal = SchemaApplicationRecord::new(receipt.clone(), Vec::new())?;
2833 let operation = SchemaApplicationRecordOp::replace(record, &terminal)?;
2834 publish_accepted_schema_candidates_with_application_record(Vec::new(), operation)?;
2835 Ok(SchemaChangeProgress::new(receipt, status))
2836}
2837
2838fn application_publications<'a>(
2839 authorities: &[StoreApplicationAuthority],
2840 current_bundles: &[Option<crate::db::schema::AcceptedSchemaRevisionBundle>],
2841 candidates: &'a [crate::db::schema::CandidateSchemaRevision],
2842) -> Result<Vec<AcceptedSchemaPublication<'a>>, InternalError> {
2843 candidates
2844 .iter()
2845 .map(|candidate| {
2846 let (position, authority) = authorities
2847 .iter()
2848 .enumerate()
2849 .find(|(_, authority)| authority.path == candidate.store_path())
2850 .ok_or_else(InternalError::store_invariant)?;
2851 let expected_revision = current_bundles[position].as_ref().map_or(
2852 AcceptedSchemaRevision::NONE,
2853 crate::db::schema::AcceptedSchemaRevisionBundle::revision,
2854 );
2855 Ok(AcceptedSchemaPublication::new(
2856 authority.path,
2857 authority.handle,
2858 expected_revision,
2859 candidate,
2860 ))
2861 })
2862 .collect()
2863}
2864
2865fn application_authorities<C: CanisterKind>(db: &Db<C>) -> Vec<StoreApplicationAuthority> {
2866 let mut authorities = db.with_store_registry(|registry| {
2867 registry
2868 .iter()
2869 .map(|(path, handle)| StoreApplicationAuthority { path, handle })
2870 .collect::<Vec<_>>()
2871 });
2872 icydb_schema::compact_sort_unstable_by(&mut authorities, |left, right| {
2873 left.path.cmp(right.path)
2874 });
2875 authorities
2876}
2877
2878fn accepted_head_after_candidates(
2879 authorities: &[StoreApplicationAuthority],
2880 candidates: &[crate::db::schema::CandidateSchemaRevision],
2881) -> Result<ExpectedAcceptedHead, InternalError> {
2882 let heads = authorities
2883 .iter()
2884 .map(|authority| {
2885 let candidate = candidates
2886 .iter()
2887 .find(|candidate| candidate.store_path() == authority.path);
2888 let head = match candidate {
2889 Some(candidate) => Some(AcceptedStoreHead {
2890 revision: candidate.revision().get(),
2891 fingerprint: candidate.root().fingerprint().as_bytes(),
2892 }),
2893 None => authority
2894 .handle
2895 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_root)?
2896 .map(|selection| AcceptedStoreHead {
2897 revision: selection.root().revision().get(),
2898 fingerprint: selection.root().fingerprint().as_bytes(),
2899 }),
2900 };
2901 Ok((authority.path, head))
2902 })
2903 .collect::<Result<Vec<_>, InternalError>>()?;
2904 Ok(derive_accepted_head(heads.as_slice()))
2905}
2906
2907fn derive_database_identity(
2908 incarnation: [u8; 16],
2909 stores: &[StoreApplicationAuthority],
2910) -> TargetDatabaseIdentity {
2911 let mut hasher = new_hash_sha256_prefixed(DATABASE_TARGET_FINGERPRINT_PROFILE);
2912 hasher.update(incarnation);
2913 write_hash_len_u32(&mut hasher, stores.len());
2914 for store in stores {
2915 write_store_authority(&mut hasher, store);
2916 }
2917 TargetDatabaseIdentity::from_bytes(finalize_hash_sha256(hasher))
2918}
2919
2920fn derive_store_identity(
2921 database_identity: TargetDatabaseIdentity,
2922 store: &StoreApplicationAuthority,
2923) -> TargetStoreIdentity {
2924 let mut hasher = new_hash_sha256_prefixed(STORE_TARGET_FINGERPRINT_PROFILE);
2925 hasher.update(database_identity.to_bytes());
2926 write_store_authority(&mut hasher, store);
2927 TargetStoreIdentity::from_bytes(finalize_hash_sha256(hasher))
2928}
2929
2930fn derive_accepted_head(stores: &[(&str, Option<AcceptedStoreHead>)]) -> ExpectedAcceptedHead {
2931 let Some(revision) = stores
2932 .iter()
2933 .filter_map(|(_, head)| head.map(|head| head.revision))
2934 .max()
2935 else {
2936 return ExpectedAcceptedHead::Empty;
2937 };
2938
2939 let mut hasher = new_hash_sha256_prefixed(ACCEPTED_DATABASE_HEAD_FINGERPRINT_PROFILE);
2940 write_hash_len_u32(&mut hasher, stores.len());
2941 for (path, head) in stores {
2942 write_hash_str_u32(&mut hasher, path);
2943 match head {
2944 None => write_hash_tag_u8(&mut hasher, 0),
2945 Some(head) => {
2946 write_hash_tag_u8(&mut hasher, 1);
2947 write_hash_u64(&mut hasher, head.revision);
2948 hasher.update(head.fingerprint);
2949 }
2950 }
2951 }
2952
2953 ExpectedAcceptedHead::Exact {
2954 revision,
2955 fingerprint: ExpectedSchemaFingerprint::from_bytes(finalize_hash_sha256(hasher)),
2956 }
2957}
2958
2959pub(in crate::db) fn generated_schema_reconciled(
2960 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
2961 incarnation: DatabaseIncarnationId,
2962 submission_key: &str,
2963) -> Result<(bool, ExpectedAcceptedHead), InternalError> {
2964 let submission_key = SchemaSubmissionKey::try_new(submission_key.to_string())
2965 .map_err(|_| InternalError::store_invariant())?;
2966 let (database_identity, accepted_head) = generated_schema_authority(registry, incarnation)?;
2967 let reconciled = generated_submission_is_reconciled(database_identity, &submission_key)?;
2968 Ok((reconciled, accepted_head))
2969}
2970
2971pub(in crate::db) fn generated_schema_is_reconciled(
2972 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
2973 incarnation: DatabaseIncarnationId,
2974 submission_key: &str,
2975) -> Result<bool, InternalError> {
2976 let submission_key = SchemaSubmissionKey::try_new(submission_key.to_string())
2977 .map_err(|_| InternalError::store_invariant())?;
2978 let database_identity = generated_database_identity(registry, incarnation);
2979 generated_submission_is_reconciled(database_identity, &submission_key)
2980}
2981
2982fn generated_submission_is_reconciled(
2983 database_identity: TargetDatabaseIdentity,
2984 submission_key: &SchemaSubmissionKey,
2985) -> Result<bool, InternalError> {
2986 Ok(
2987 load_schema_application_record_read_only(database_identity, submission_key)?.is_some_and(
2988 |record| match record.receipt().outcome() {
2989 SchemaChangeOutcome::NoOp { .. } | SchemaChangeOutcome::Applied { .. } => true,
2997 SchemaChangeOutcome::Pending { .. } | SchemaChangeOutcome::Aborted { .. } => false,
2998 },
2999 ),
3000 )
3001}
3002
3003pub(in crate::db) fn generated_schema_authority(
3004 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3005 incarnation: DatabaseIncarnationId,
3006) -> Result<(TargetDatabaseIdentity, ExpectedAcceptedHead), InternalError> {
3007 let database_identity = generated_database_identity(registry, incarnation);
3008 let mut stores = store_application_authorities(registry);
3009 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.path.cmp(right.path));
3010 let heads = stores
3011 .iter()
3012 .map(|store| {
3013 let head = store
3014 .handle
3015 .with_schema(crate::db::schema::SchemaStore::current_accepted_schema_root)?
3016 .map(|selection| AcceptedStoreHead {
3017 revision: selection.root().revision().get(),
3018 fingerprint: selection.root().fingerprint().as_bytes(),
3019 });
3020 Ok((store.path, head))
3021 })
3022 .collect::<Result<Vec<_>, InternalError>>()?;
3023 let accepted_head = derive_accepted_head(heads.as_slice());
3024 Ok((database_identity, accepted_head))
3025}
3026
3027fn generated_database_identity(
3028 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3029 incarnation: DatabaseIncarnationId,
3030) -> TargetDatabaseIdentity {
3031 let registry_key = std::ptr::from_ref(registry).cast::<()>() as usize;
3032 let cached = GENERATED_DATABASE_IDENTITY
3033 .with(Cell::get)
3034 .and_then(|entry| {
3035 (entry.registry == registry_key && entry.incarnation == incarnation)
3036 .then_some(entry.identity)
3037 });
3038 if let Some(identity) = cached {
3039 return identity;
3040 }
3041
3042 let mut stores = store_application_authorities(registry);
3043 icydb_schema::compact_sort_unstable_by(&mut stores, |left, right| left.path.cmp(right.path));
3044 let identity = derive_database_identity(incarnation.to_bytes(), stores.as_slice());
3045 GENERATED_DATABASE_IDENTITY.set(Some(GeneratedDatabaseIdentityCacheEntry {
3046 registry: registry_key,
3047 incarnation,
3048 identity,
3049 }));
3050 identity
3051}
3052
3053fn store_application_authorities(
3054 registry: &'static std::thread::LocalKey<crate::db::registry::StoreRegistry>,
3055) -> Vec<StoreApplicationAuthority> {
3056 registry.with(|registry| {
3057 registry
3058 .iter()
3059 .map(|(path, handle)| StoreApplicationAuthority { path, handle })
3060 .collect()
3061 })
3062}
3063
3064fn write_store_authority(hasher: &mut sha2::Sha256, store: &StoreApplicationAuthority) {
3065 write_hash_str_u32(hasher, store.path);
3066 write_storage_capabilities(hasher, store.handle);
3067 for allocation in [
3068 store.handle.data_allocation(),
3069 store.handle.index_allocation(),
3070 store.handle.schema_allocation(),
3071 store.handle.journal_allocation(),
3072 ] {
3073 write_allocation_identity(hasher, allocation);
3074 }
3075}
3076
3077fn write_storage_capabilities(hasher: &mut sha2::Sha256, store: StoreHandle) {
3078 let capabilities = store.storage_capabilities();
3079 write_hash_tag_u8(
3080 hasher,
3081 match capabilities.storage_mode() {
3082 StoreRuntimeStorageMode::Heap => 0,
3083 StoreRuntimeStorageMode::Journaled => 1,
3084 },
3085 );
3086 write_hash_tag_u8(
3087 hasher,
3088 match capabilities.allocation_identity() {
3089 StoreAllocationIdentityCapability::Present => 0,
3090 StoreAllocationIdentityCapability::Absent => 1,
3091 },
3092 );
3093 write_hash_tag_u8(
3094 hasher,
3095 match capabilities.durability() {
3096 StoreDurability::Durable => 0,
3097 StoreDurability::Volatile => 1,
3098 },
3099 );
3100 write_hash_tag_u8(
3101 hasher,
3102 match capabilities.recovery() {
3103 StoreRecoveryCapability::StableBasePlusJournalReplay => 0,
3104 StoreRecoveryCapability::None => 1,
3105 },
3106 );
3107 write_hash_tag_u8(
3108 hasher,
3109 match capabilities.commit_participation() {
3110 StoreCommitParticipation::Durable => 0,
3111 StoreCommitParticipation::LiveOnly => 1,
3112 },
3113 );
3114 write_hash_tag_u8(
3115 hasher,
3116 match capabilities.schema_metadata() {
3117 StoreSchemaMetadataCapability::LiveRebuiltMetadata => 0,
3118 StoreSchemaMetadataCapability::CanonicalStableHistoryPlusJournalTail => 1,
3119 },
3120 );
3121 write_hash_tag_u8(
3122 hasher,
3123 match capabilities.relation_source() {
3124 StoreRelationSourceCapability::DurableSource => 0,
3125 StoreRelationSourceCapability::LiveSource => 1,
3126 },
3127 );
3128 write_hash_tag_u8(
3129 hasher,
3130 match capabilities.relation_target() {
3131 StoreRelationTargetCapability::DurableTarget => 0,
3132 StoreRelationTargetCapability::VolatileTarget => 1,
3133 },
3134 );
3135}
3136
3137fn write_allocation_identity(
3138 hasher: &mut sha2::Sha256,
3139 allocation: Option<StoreAllocationIdentity>,
3140) {
3141 match allocation {
3142 None => write_hash_tag_u8(hasher, 0),
3143 Some(allocation) => {
3144 write_hash_tag_u8(hasher, 1);
3145 write_hash_tag_u8(hasher, allocation.memory_id());
3146 write_hash_str_u32(hasher, allocation.stable_key());
3147 }
3148 }
3149}
3150
3151#[cfg(test)]
3152mod tests {
3153 #[cfg(feature = "migration")]
3154 mod nested_migration;
3155
3156 use super::{
3157 AcceptedSchemaPublication, AcceptedStoreHead, DirectGeneratedRowLocalProof,
3158 PendingGeneratedRowLocalConstraint, abort_schema_application,
3159 aborted_generated_row_local_candidate, accepted_head_after_candidates,
3160 application_authorities, apply_schema, continue_schema_application, derive_accepted_head,
3161 derive_schema_change_job_id, final_candidates_for_pending_row_local_constraint,
3162 generated_database_identity, include_identity_state_count, lower_existing_schema_proposal,
3163 lower_initial_schema_proposal, publish_accepted_schema_candidates_with_application_record,
3164 require_exact_empty_entity_count, schema_application_target,
3165 };
3166 use crate::{
3167 db::{
3168 DatabaseStartupState, Db, GeneratedStartupDriverStep,
3169 commit::{
3170 RecoveryProgress, continue_recovery, database_incarnation_id,
3171 forget_recovered_domain_for_tests,
3172 },
3173 data::DataStore,
3174 drive_generated_startup_recovery_page,
3175 index::IndexStore,
3176 journal::JournalTailStore,
3177 observe_generated_startup_state,
3178 registry::{
3179 StoreAllocationIdentities, StoreAllocationIdentity, StoreHandle, StoreRegistry,
3180 StoreRuntimeStorageCapabilities,
3181 },
3182 schema::{
3183 AcceptedConstraintKind, AcceptedRuleOperation, AcceptedSchemaRevisionBundle,
3184 CandidateSchemaRevision, ConstraintOrigin, ConstraintValidationJob,
3185 ExistingProposalStore, ProposalStoreTarget, SchemaApplicationRecord,
3186 SchemaApplicationRecordOp, SchemaChangeActivation, SchemaChangeJob,
3187 SchemaChangeOutcome, SchemaChangeProgressStatus, SchemaStore,
3188 cardinality_build::{
3189 CardinalityBuildAuthority, CardinalityGenerationPageOutcome,
3190 drive_cardinality_generation_page,
3191 },
3192 },
3193 },
3194 error::{ErrorClass, ErrorOrigin},
3195 testing::test_memory,
3196 traits::{CanisterKind, Path},
3197 };
3198 use ic_stable_structures::{DefaultMemoryImpl, memory_manager::VirtualMemory};
3199 use icydb_schema::{
3200 ConstraintFragment, ConstraintSourceKey, DeclaredEntityVersion, EntityFragment,
3201 EntitySourceKey, EntityStoreAssignment, ExpectedAcceptedHead, ExpectedSchemaFingerprint,
3202 FieldFragment, FieldInsertPolicy, FieldSourceKey, FieldType, NamedTypeFragment,
3203 RuleSourceKey, ScalarLiteral, ScalarType, SchemaCapability, SchemaFragment, SchemaName,
3204 SchemaProposal, SchemaSubmissionKey, SourceCheckExpr, SourceCheckInstruction,
3205 SourceRuleOperation, TargetDatabaseIdentity, TargetStoreIdentity, TargetedRuleFragment,
3206 TypeSourceKey,
3207 };
3208 use std::cell::RefCell;
3209
3210 fn drive_startup_recovery_to_completion<C: CanisterKind>(db: &Db<C>) {
3211 for _ in 0..1_024 {
3212 match continue_recovery(db).expect("test startup recovery page should succeed") {
3213 RecoveryProgress::Complete => return,
3214 RecoveryProgress::Pending => {}
3215 }
3216 }
3217 panic!("test startup recovery should complete within 1,024 bounded pages");
3218 }
3219
3220 fn drive_cardinality_to_ready(store: StoreHandle) {
3221 let journal = store
3222 .journal_tail_store()
3223 .expect("cardinality fixture store should be journaled");
3224 for _ in 0..8 {
3225 let outcome = store
3226 .with_data(|data| {
3227 store.with_index(|index| {
3228 store.with_schema_mut(|schema| {
3229 drive_cardinality_generation_page(data, index, schema, |schema| {
3230 CardinalityBuildAuthority::derive(
3231 schema,
3232 database_incarnation_id()?,
3233 store.allocation_identities(),
3234 journal.with_borrow(JournalTailStore::fold_watermark)?,
3235 )
3236 })
3237 })
3238 })
3239 })
3240 .expect("bounded cardinality generation should advance");
3241 if outcome == CardinalityGenerationPageOutcome::Quiescent {
3242 return;
3243 }
3244 }
3245 panic!("cardinality generation should become Ready within eight bounded pages");
3246 }
3247
3248 #[cfg(feature = "migration")]
3249 use crate::db::schema::SchemaChangeReceipt;
3250 use crate::{
3251 db::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell},
3252 value::InputValue,
3253 };
3254 #[cfg(feature = "migration")]
3255 use icydb_schema::{
3256 EntityMigration, IndexFragment, IndexKeyFragment, RelationDeleteAction, RelationFragment,
3257 SchemaMigrationPlan, SchemaMigrationTransform, SchemaProposalDigest,
3258 };
3259
3260 fn version_one() -> DeclaredEntityVersion {
3261 DeclaredEntityVersion::try_new(1).expect("fixture version should admit")
3262 }
3263
3264 const ABORT_STORE_PATH: &str = "schema_application_tests::AbortStore";
3265 const EVOLUTION_STORE_PATH: &str = "schema_application_tests::EvolutionStore";
3266 #[cfg(feature = "migration")]
3267 const MIGRATION_STORE_PATH: &str = "schema_application_tests::MigrationStore";
3268 #[cfg(feature = "migration")]
3269 const MIGRATION_EXECUTION_STORE_PATH: &str =
3270 "schema_application_tests::MigrationExecutionStore";
3271 #[cfg(feature = "migration")]
3272 const MIGRATION_FINDING_STORE_PATH: &str = "schema_application_tests::MigrationFindingStore";
3273
3274 #[test]
3275 fn database_identity_state_capacity_combines_store_inventories_exactly() {
3276 let below = include_identity_state_count(0, 65_535)
3277 .expect("the first store inventory should remain below the database cap");
3278 let exact = include_identity_state_count(below, 1)
3279 .expect("the combined database boundary should admit");
3280 assert_eq!(exact, 65_536);
3281
3282 let error = include_identity_state_count(exact, 1)
3283 .expect_err("the next owner in another store must reject");
3284 assert_eq!(error.class(), ErrorClass::Unsupported);
3285 assert_eq!(error.origin(), ErrorOrigin::Identity);
3286 }
3287
3288 #[test]
3289 fn generated_database_identity_cache_is_bound_to_the_incarnation() {
3290 let first_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x41);
3291 let second_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x42);
3292 let first = generated_database_identity(&ABORT_REGISTRY, first_incarnation);
3293 assert_eq!(
3294 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3295 first,
3296 );
3297 let second = generated_database_identity(&ABORT_REGISTRY, second_incarnation);
3298 assert_ne!(second, first);
3299 assert_eq!(
3300 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3301 first,
3302 );
3303 }
3304
3305 #[test]
3306 fn exact_empty_entity_proof_distinguishes_corruption_from_non_empty_input() {
3307 let corrupt = require_exact_empty_entity_count(None)
3308 .expect_err("uninspectable cardinality must fail closed");
3309 assert_eq!(corrupt.class(), ErrorClass::Corruption);
3310
3311 let non_empty = require_exact_empty_entity_count(Some(1))
3312 .expect_err("non-empty cardinality must reject removal");
3313 assert_eq!(non_empty.class(), ErrorClass::Unsupported);
3314 assert!(require_exact_empty_entity_count(Some(0)).is_ok());
3315 }
3316
3317 thread_local! {
3318 static ABORT_DATA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(180);
3319 static ABORT_INDEX_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(181);
3320 static ABORT_SCHEMA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(182);
3321 static ABORT_JOURNAL_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(183);
3322 static ABORT_DATA: RefCell<DataStore> =
3323 ABORT_DATA_MEMORY.with(|memory| {
3324 RefCell::new(DataStore::init_journaled(memory.clone()))
3325 });
3326 static ABORT_INDEX: RefCell<IndexStore> =
3327 ABORT_INDEX_MEMORY.with(|memory| {
3328 RefCell::new(IndexStore::init_journaled(memory.clone()))
3329 });
3330 static ABORT_SCHEMA: RefCell<SchemaStore> =
3331 ABORT_SCHEMA_MEMORY.with(|memory| {
3332 RefCell::new(SchemaStore::init_journaled(memory.clone()))
3333 });
3334 static ABORT_JOURNAL: RefCell<JournalTailStore> =
3335 ABORT_JOURNAL_MEMORY.with(|memory| {
3336 RefCell::new(JournalTailStore::init(memory.clone()))
3337 });
3338 static ABORT_REGISTRY: StoreRegistry = {
3339 let mut registry = StoreRegistry::new();
3340 registry.register_journaled_store(
3341 ABORT_STORE_PATH,
3342 &ABORT_DATA,
3343 &ABORT_INDEX,
3344 &ABORT_SCHEMA,
3345 &ABORT_JOURNAL,
3346 StoreAllocationIdentities::new_journaled(
3347 StoreAllocationIdentity::new(180, "icydb.test.application_abort.data.v1"),
3348 StoreAllocationIdentity::new(181, "icydb.test.application_abort.index.v1"),
3349 StoreAllocationIdentity::new(182, "icydb.test.application_abort.schema.v1"),
3350 StoreAllocationIdentity::new(183, "icydb.test.application_abort.journal.v1"),
3351 ),
3352 StoreRuntimeStorageCapabilities::journaled(),
3353 ).expect("abort journaled store should register");
3354 registry
3355 };
3356 }
3357
3358 #[cfg(feature = "migration")]
3359 thread_local! {
3360 static MIGRATION_EXECUTION_DATA: RefCell<DataStore> =
3361 RefCell::new(DataStore::init_journaled(test_memory(210)));
3362 static MIGRATION_EXECUTION_INDEX: RefCell<IndexStore> =
3363 RefCell::new(IndexStore::init_journaled(test_memory(211)));
3364 static MIGRATION_EXECUTION_SCHEMA: RefCell<SchemaStore> =
3365 RefCell::new(SchemaStore::init_journaled(test_memory(212)));
3366 static MIGRATION_EXECUTION_JOURNAL: RefCell<JournalTailStore> =
3367 RefCell::new(JournalTailStore::init(test_memory(213)));
3368 static MIGRATION_EXECUTION_REGISTRY: StoreRegistry = {
3369 let mut registry = StoreRegistry::new();
3370 registry.register_journaled_store(
3371 MIGRATION_EXECUTION_STORE_PATH,
3372 &MIGRATION_EXECUTION_DATA,
3373 &MIGRATION_EXECUTION_INDEX,
3374 &MIGRATION_EXECUTION_SCHEMA,
3375 &MIGRATION_EXECUTION_JOURNAL,
3376 StoreAllocationIdentities::new_journaled(
3377 StoreAllocationIdentity::new(210, "icydb.test.migration_execution.data.v1"),
3378 StoreAllocationIdentity::new(211, "icydb.test.migration_execution.index.v1"),
3379 StoreAllocationIdentity::new(212, "icydb.test.migration_execution.schema.v1"),
3380 StoreAllocationIdentity::new(213, "icydb.test.migration_execution.journal.v1"),
3381 ),
3382 StoreRuntimeStorageCapabilities::journaled(),
3383 ).expect("migration execution store should register");
3384 registry
3385 };
3386 }
3387
3388 #[cfg(feature = "migration")]
3389 thread_local! {
3390 static MIGRATION_DATA: RefCell<DataStore> =
3391 RefCell::new(DataStore::init_journaled(test_memory(200)));
3392 static MIGRATION_INDEX: RefCell<IndexStore> =
3393 RefCell::new(IndexStore::init_journaled(test_memory(201)));
3394 static MIGRATION_SCHEMA: RefCell<SchemaStore> =
3395 RefCell::new(SchemaStore::init_journaled(test_memory(202)));
3396 static MIGRATION_JOURNAL: RefCell<JournalTailStore> =
3397 RefCell::new(JournalTailStore::init(test_memory(203)));
3398 static MIGRATION_REGISTRY: StoreRegistry = {
3399 let mut registry = StoreRegistry::new();
3400 registry.register_journaled_store(
3401 MIGRATION_STORE_PATH,
3402 &MIGRATION_DATA,
3403 &MIGRATION_INDEX,
3404 &MIGRATION_SCHEMA,
3405 &MIGRATION_JOURNAL,
3406 StoreAllocationIdentities::new_journaled(
3407 StoreAllocationIdentity::new(200, "icydb.test.migration_validation.data.v1"),
3408 StoreAllocationIdentity::new(201, "icydb.test.migration_validation.index.v1"),
3409 StoreAllocationIdentity::new(202, "icydb.test.migration_validation.schema.v1"),
3410 StoreAllocationIdentity::new(203, "icydb.test.migration_validation.journal.v1"),
3411 ),
3412 StoreRuntimeStorageCapabilities::journaled(),
3413 ).expect("migration validation store should register");
3414 registry
3415 };
3416 }
3417
3418 #[cfg(feature = "migration")]
3419 thread_local! {
3420 static MIGRATION_FINDING_DATA: RefCell<DataStore> =
3421 RefCell::new(DataStore::init_journaled(test_memory(206)));
3422 static MIGRATION_FINDING_INDEX: RefCell<IndexStore> =
3423 RefCell::new(IndexStore::init_journaled(test_memory(207)));
3424 static MIGRATION_FINDING_SCHEMA: RefCell<SchemaStore> =
3425 RefCell::new(SchemaStore::init_journaled(test_memory(208)));
3426 static MIGRATION_FINDING_JOURNAL: RefCell<JournalTailStore> =
3427 RefCell::new(JournalTailStore::init(test_memory(209)));
3428 static MIGRATION_FINDING_REGISTRY: StoreRegistry = {
3429 let mut registry = StoreRegistry::new();
3430 registry.register_journaled_store(
3431 MIGRATION_FINDING_STORE_PATH,
3432 &MIGRATION_FINDING_DATA,
3433 &MIGRATION_FINDING_INDEX,
3434 &MIGRATION_FINDING_SCHEMA,
3435 &MIGRATION_FINDING_JOURNAL,
3436 StoreAllocationIdentities::new_journaled(
3437 StoreAllocationIdentity::new(206, "icydb.test.migration_finding.data.v1"),
3438 StoreAllocationIdentity::new(207, "icydb.test.migration_finding.index.v1"),
3439 StoreAllocationIdentity::new(208, "icydb.test.migration_finding.schema.v1"),
3440 StoreAllocationIdentity::new(209, "icydb.test.migration_finding.journal.v1"),
3441 ),
3442 StoreRuntimeStorageCapabilities::journaled(),
3443 ).expect("migration finding store should register");
3444 registry
3445 };
3446 }
3447
3448 thread_local! {
3449 static EVOLUTION_DATA: RefCell<DataStore> =
3450 RefCell::new(DataStore::init_journaled(test_memory(192)));
3451 static EVOLUTION_INDEX: RefCell<IndexStore> =
3452 RefCell::new(IndexStore::init_journaled(test_memory(193)));
3453 static EVOLUTION_SCHEMA: RefCell<SchemaStore> =
3454 RefCell::new(SchemaStore::init_journaled(test_memory(194)));
3455 static EVOLUTION_JOURNAL: RefCell<JournalTailStore> =
3456 RefCell::new(JournalTailStore::init(test_memory(195)));
3457 static EVOLUTION_REGISTRY: StoreRegistry = {
3458 let mut registry = StoreRegistry::new();
3459 registry.register_journaled_store(
3460 EVOLUTION_STORE_PATH,
3461 &EVOLUTION_DATA,
3462 &EVOLUTION_INDEX,
3463 &EVOLUTION_SCHEMA,
3464 &EVOLUTION_JOURNAL,
3465 StoreAllocationIdentities::new_journaled(
3466 StoreAllocationIdentity::new(192, "icydb.test.rule_evolution.data.v1"),
3467 StoreAllocationIdentity::new(193, "icydb.test.rule_evolution.index.v1"),
3468 StoreAllocationIdentity::new(194, "icydb.test.rule_evolution.schema.v1"),
3469 StoreAllocationIdentity::new(195, "icydb.test.rule_evolution.journal.v1"),
3470 ),
3471 StoreRuntimeStorageCapabilities::journaled(),
3472 ).expect("rule-evolution journaled store should register");
3473 registry
3474 };
3475 }
3476
3477 struct AbortCanister;
3478
3479 impl Path for AbortCanister {
3480 const PATH: &'static str = "schema_application_tests::AbortCanister";
3481 }
3482
3483 impl CanisterKind for AbortCanister {
3484 const COMMIT_MEMORY_ID: u8 = 184;
3485 const COMMIT_STABLE_KEY: &'static str = "icydb.test.application_abort.commit.v1";
3486 const STARTUP_MEMORY_ID: u8 = 186;
3487 const STARTUP_STABLE_KEY: &'static str = "icydb.test.application_abort.startup.control.v1";
3488 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 185;
3489 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3490 "icydb.test.application_abort.integrity.v1";
3491 }
3492
3493 struct EvolutionCanister;
3494
3495 impl Path for EvolutionCanister {
3496 const PATH: &'static str = "schema_application_tests::EvolutionCanister";
3497 }
3498
3499 impl CanisterKind for EvolutionCanister {
3500 const COMMIT_MEMORY_ID: u8 = 196;
3501 const COMMIT_STABLE_KEY: &'static str = "icydb.test.rule_evolution.commit.v1";
3502 const STARTUP_MEMORY_ID: u8 = 198;
3503 const STARTUP_STABLE_KEY: &'static str = "icydb.test.rule_evolution.startup.control.v1";
3504 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 197;
3505 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3506 "icydb.test.rule_evolution.integrity.v1";
3507 }
3508
3509 #[cfg(feature = "migration")]
3510 struct MigrationCanister;
3511
3512 #[cfg(feature = "migration")]
3513 impl Path for MigrationCanister {
3514 const PATH: &'static str = "schema_application_tests::MigrationCanister";
3515 }
3516
3517 #[cfg(feature = "migration")]
3518 impl CanisterKind for MigrationCanister {
3519 const COMMIT_MEMORY_ID: u8 = 204;
3520 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_validation.commit.v1";
3521 const STARTUP_MEMORY_ID: u8 = 206;
3522 const STARTUP_STABLE_KEY: &'static str =
3523 "icydb.test.migration_validation.startup.control.v1";
3524 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 205;
3525 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3526 "icydb.test.migration_validation.integrity.v1";
3527 }
3528
3529 #[cfg(feature = "migration")]
3530 struct MigrationExecutionCanister;
3531
3532 #[cfg(feature = "migration")]
3533 impl Path for MigrationExecutionCanister {
3534 const PATH: &'static str = "schema_application_tests::MigrationExecutionCanister";
3535 }
3536
3537 #[cfg(feature = "migration")]
3538 impl CanisterKind for MigrationExecutionCanister {
3539 const COMMIT_MEMORY_ID: u8 = 214;
3540 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_execution.commit.v1";
3541 const STARTUP_MEMORY_ID: u8 = 216;
3542 const STARTUP_STABLE_KEY: &'static str =
3543 "icydb.test.migration_execution.startup.control.v1";
3544 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 215;
3545 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3546 "icydb.test.migration_execution.integrity.v1";
3547 }
3548
3549 #[cfg(feature = "migration")]
3550 struct MigrationFindingCanister;
3551
3552 #[cfg(feature = "migration")]
3553 impl Path for MigrationFindingCanister {
3554 const PATH: &'static str = "schema_application_tests::MigrationFindingCanister";
3555 }
3556
3557 #[cfg(feature = "migration")]
3558 impl CanisterKind for MigrationFindingCanister {
3559 const COMMIT_MEMORY_ID: u8 = 210;
3560 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_finding.commit.v1";
3561 const STARTUP_MEMORY_ID: u8 = 212;
3562 const STARTUP_STABLE_KEY: &'static str = "icydb.test.migration_finding.startup.control.v1";
3563 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 211;
3564 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3565 "icydb.test.migration_finding.integrity.v1";
3566 }
3567
3568 fn name(value: &str) -> SchemaName {
3569 SchemaName::try_new(value).expect("test schema name should admit")
3570 }
3571
3572 fn generated_check_proposal(
3573 expected_head: ExpectedAcceptedHead,
3574 submission_key: &str,
3575 include_check: bool,
3576 database: TargetDatabaseIdentity,
3577 store: TargetStoreIdentity,
3578 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3579 let entity_source = EntitySourceKey::try_new("Item").expect("entity source should admit");
3580 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3581 let score_source = FieldSourceKey::try_new("score").expect("score source should admit");
3582 let check_source =
3583 ConstraintSourceKey::try_new("score_non_negative").expect("check source should admit");
3584 let check = SourceCheckExpr::try_new(vec![
3585 SourceCheckInstruction::Field(score_source),
3586 SourceCheckInstruction::Literal(ScalarLiteral::Int(0)),
3587 SourceCheckInstruction::GreaterThanOrEqual,
3588 ])
3589 .expect("check expression should admit");
3590 let constraints = include_check
3591 .then(|| ConstraintFragment::check(name("score_non_negative"), check))
3592 .into_iter()
3593 .collect();
3594 let entity = EntityFragment::try_new(
3595 name("Item"),
3596 version_one(),
3597 vec![
3598 FieldFragment::new(
3599 name("id"),
3600 FieldType::Scalar(ScalarType::Nat64),
3601 false,
3602 FieldInsertPolicy::Required,
3603 None,
3604 ),
3605 FieldFragment::new(
3606 name("score"),
3607 FieldType::Scalar(ScalarType::Int64),
3608 false,
3609 FieldInsertPolicy::Required,
3610 None,
3611 ),
3612 ],
3613 vec![id_source],
3614 Vec::new(),
3615 Vec::new(),
3616 constraints,
3617 )
3618 .expect("entity should admit");
3619 let proposal = SchemaProposal::try_compose(
3620 vec![SchemaCapability::ACCEPTED_CHECKS],
3621 database,
3622 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3623 expected_head,
3624 vec![
3625 SchemaFragment::try_new(vec![entity], Vec::new())
3626 .expect("schema fragment should admit"),
3627 ],
3628 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3629 Vec::new(),
3630 None,
3631 )
3632 .expect("schema proposal should compose");
3633 (proposal, entity_source, check_source)
3634 }
3635
3636 #[cfg(feature = "migration")]
3637 #[derive(Clone, Copy, Eq, PartialEq)]
3638 enum ValidationMigrationShape {
3639 Clean,
3640 AllFindingFamilies,
3641 }
3642
3643 #[cfg(feature = "migration")]
3644 #[expect(
3645 clippy::too_many_lines,
3646 reason = "the fixture keeps both predecessor and candidate source contracts adjacent"
3647 )]
3648 fn validation_migration_proposal(
3649 shape: ValidationMigrationShape,
3650 current: bool,
3651 expected_head: ExpectedAcceptedHead,
3652 database: TargetDatabaseIdentity,
3653 store: TargetStoreIdentity,
3654 ) -> SchemaProposal {
3655 let entity_source = EntitySourceKey::try_new("MigratingItem")
3656 .expect("migration entity source should admit");
3657 let old_value =
3658 FieldSourceKey::try_new("old_value").expect("predecessor field source should admit");
3659 let current_value =
3660 FieldSourceKey::try_new("value").expect("candidate field source should admit");
3661 let target_entity = EntitySourceKey::try_new("MigrationTarget")
3662 .expect("migration target source should admit");
3663 let target_id = FieldSourceKey::try_new("id").expect("target id source should admit");
3664 let constraint = SourceCheckExpr::try_new(vec![
3665 SourceCheckInstruction::Field(current_value.clone()),
3666 SourceCheckInstruction::Literal(ScalarLiteral::Nat(8)),
3667 SourceCheckInstruction::LessThanOrEqual,
3668 ])
3669 .expect("candidate check should admit");
3670 let findings = shape == ValidationMigrationShape::AllFindingFamilies;
3671 let entity = EntityFragment::try_new(
3672 name("MigratingItem"),
3673 DeclaredEntityVersion::try_new(if current { 2 } else { 1 })
3674 .expect("migration version should admit"),
3675 vec![
3676 FieldFragment::new(
3677 name("id"),
3678 FieldType::Scalar(ScalarType::Nat64),
3679 false,
3680 FieldInsertPolicy::Required,
3681 None,
3682 ),
3683 FieldFragment::new(
3684 name(if current { "value" } else { "old_value" }),
3685 FieldType::Scalar(if current {
3686 ScalarType::Nat8
3687 } else {
3688 ScalarType::Int64
3689 }),
3690 false,
3691 FieldInsertPolicy::Required,
3692 None,
3693 ),
3694 ],
3695 vec![FieldSourceKey::try_new("id").expect("id source should admit")],
3696 current
3697 .then(|| {
3698 IndexFragment::try_new(
3699 name("value_unique"),
3700 vec![IndexKeyFragment::Field(current_value.clone())],
3701 true,
3702 None,
3703 )
3704 .expect("candidate index should admit")
3705 })
3706 .into_iter()
3707 .collect(),
3708 (current && findings)
3709 .then(|| {
3710 RelationFragment::try_new(
3711 name("value_target"),
3712 icydb_schema::RelationSourceFragment::direct(vec![current_value.clone()]),
3713 target_entity.clone(),
3714 vec![target_id.clone()],
3715 RelationDeleteAction::Restrict,
3716 )
3717 .expect("candidate relation should admit")
3718 })
3719 .into_iter()
3720 .collect(),
3721 (current && findings)
3722 .then(|| ConstraintFragment::check(name("value_at_most_eight"), constraint))
3723 .into_iter()
3724 .collect(),
3725 )
3726 .expect("migration entity should admit");
3727 let target = EntityFragment::try_new(
3728 name("MigrationTarget"),
3729 version_one(),
3730 vec![FieldFragment::new(
3731 name("id"),
3732 FieldType::Scalar(ScalarType::Nat8),
3733 false,
3734 FieldInsertPolicy::Required,
3735 None,
3736 )],
3737 vec![target_id],
3738 Vec::new(),
3739 Vec::new(),
3740 Vec::new(),
3741 )
3742 .expect("migration relation target should admit");
3743 let migration = current.then(|| {
3744 SchemaMigrationPlan::try_new(vec![
3745 EntityMigration::try_new(
3746 entity_source.clone(),
3747 DeclaredEntityVersion::try_new(1).expect("predecessor should admit"),
3748 None,
3749 Vec::new(),
3750 vec![SchemaMigrationTransform::CheckedCast {
3751 from: old_value.clone(),
3752 to: current_value,
3753 target: ScalarType::Nat8,
3754 }],
3755 )
3756 .expect("migration transition should admit"),
3757 ])
3758 .expect("migration plan should admit")
3759 });
3760 let mut capabilities = Vec::new();
3761 if current && findings {
3762 capabilities.extend([
3763 SchemaCapability::ACCEPTED_CHECKS,
3764 SchemaCapability::SECONDARY_INDEXES,
3765 SchemaCapability::RESTRICTIVE_RELATIONS,
3766 ]);
3767 }
3768 if migration.is_some() {
3769 capabilities.push(SchemaCapability::VERSIONED_MIGRATIONS);
3770 }
3771 let mut entities = vec![entity];
3772 let mut assignments = vec![EntityStoreAssignment::new(entity_source.clone(), store)];
3773 if findings {
3774 entities.push(target);
3775 assignments.push(EntityStoreAssignment::new(target_entity, store));
3776 }
3777 SchemaProposal::try_compose(
3778 capabilities,
3779 database,
3780 SchemaSubmissionKey::try_new(if current {
3781 "migration-validation-v2"
3782 } else {
3783 "migration-validation-v1"
3784 })
3785 .expect("submission should admit"),
3786 expected_head,
3787 vec![
3788 SchemaFragment::try_new(entities, Vec::new())
3789 .expect("migration fragment should admit"),
3790 ],
3791 assignments,
3792 current
3793 .then_some(icydb_schema::SchemaRemoval::Field {
3794 entity: entity_source,
3795 field: old_value,
3796 })
3797 .into_iter()
3798 .collect(),
3799 migration,
3800 )
3801 .expect("migration proposal should compose")
3802 }
3803
3804 fn targeted_rule_proposal(
3805 expected_head: ExpectedAcceptedHead,
3806 submission_key: &str,
3807 operation: SourceRuleOperation,
3808 database: TargetDatabaseIdentity,
3809 store: TargetStoreIdentity,
3810 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3811 let entity_source =
3812 EntitySourceKey::try_new("Measured").expect("entity source should admit");
3813 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3814 let value_source = FieldSourceKey::try_new("value").expect("value source should admit");
3815 let value_type = TypeSourceKey::try_new("Measure").expect("type source should admit");
3816 let rule_source = RuleSourceKey::try_new("limit").expect("rule source should admit");
3817 let constraint_source =
3818 ConstraintSourceKey::for_targeted_field_rule(&value_source, &value_type, &rule_source);
3819 let entity = EntityFragment::try_new(
3820 name("Measured"),
3821 version_one(),
3822 vec![
3823 FieldFragment::new(
3824 name("id"),
3825 FieldType::Scalar(ScalarType::Nat64),
3826 false,
3827 FieldInsertPolicy::Required,
3828 None,
3829 ),
3830 FieldFragment::new(
3831 name("value"),
3832 FieldType::Named(value_type.clone()),
3833 false,
3834 FieldInsertPolicy::Required,
3835 None,
3836 ),
3837 ],
3838 vec![id_source],
3839 Vec::new(),
3840 Vec::new(),
3841 vec![ConstraintFragment::targeted_rule(
3842 TargetedRuleFragment::new(value_source, value_type, name("limit"), operation),
3843 )],
3844 )
3845 .expect("targeted entity should admit");
3846 let proposal = SchemaProposal::try_compose(
3847 vec![SchemaCapability::ACCEPTED_CHECKS],
3848 database,
3849 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3850 expected_head,
3851 vec![
3852 SchemaFragment::try_new(
3853 vec![entity],
3854 vec![NamedTypeFragment::newtype(
3855 name("Measure"),
3856 FieldType::Scalar(ScalarType::Nat8),
3857 )],
3858 )
3859 .expect("schema fragment should admit"),
3860 ],
3861 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3862 Vec::new(),
3863 None,
3864 )
3865 .expect("schema proposal should compose");
3866 (proposal, entity_source, constraint_source)
3867 }
3868
3869 #[test]
3870 fn database_head_is_empty_only_when_every_store_root_is_absent() {
3871 assert_eq!(
3872 derive_accepted_head(&[("test::A", None), ("test::B", None)]),
3873 ExpectedAcceptedHead::Empty,
3874 );
3875 }
3876
3877 #[test]
3878 fn database_head_covers_store_path_revision_fingerprint_and_absence() {
3879 let first = derive_accepted_head(&[
3880 (
3881 "test::A",
3882 Some(AcceptedStoreHead {
3883 revision: 3,
3884 fingerprint: [0x11; 32],
3885 }),
3886 ),
3887 ("test::B", None),
3888 ]);
3889 let changed_fingerprint = derive_accepted_head(&[
3890 (
3891 "test::A",
3892 Some(AcceptedStoreHead {
3893 revision: 3,
3894 fingerprint: [0x12; 32],
3895 }),
3896 ),
3897 ("test::B", None),
3898 ]);
3899 let changed_absence = derive_accepted_head(&[
3900 (
3901 "test::A",
3902 Some(AcceptedStoreHead {
3903 revision: 3,
3904 fingerprint: [0x11; 32],
3905 }),
3906 ),
3907 (
3908 "test::B",
3909 Some(AcceptedStoreHead {
3910 revision: 1,
3911 fingerprint: [0x22; 32],
3912 }),
3913 ),
3914 ]);
3915
3916 assert_ne!(first, changed_fingerprint);
3917 assert_ne!(first, changed_absence);
3918 assert!(matches!(
3919 first,
3920 ExpectedAcceptedHead::Exact { revision: 3, .. }
3921 ));
3922 }
3923
3924 #[test]
3925 #[allow(
3926 clippy::too_many_lines,
3927 reason = "the end-to-end catalog assertion is clearer as one lifecycle test"
3928 )]
3929 fn generated_check_abort_retires_source_identity_and_allows_fresh_reproposal() {
3930 let database = TargetDatabaseIdentity::from_bytes([0x71; 32]);
3931 let store = TargetStoreIdentity::from_bytes([0x72; 32]);
3932 let (initial, entity_source, _) = generated_check_proposal(
3933 ExpectedAcceptedHead::Empty,
3934 "abort-initial",
3935 false,
3936 database,
3937 store,
3938 );
3939 let initial_candidate = lower_initial_schema_proposal(
3940 &initial,
3941 &[ProposalStoreTarget {
3942 path: "abort::Store",
3943 identity: store,
3944 }],
3945 )
3946 .expect("initial proposal should lower")
3947 .pop()
3948 .expect("initial proposal should produce one candidate");
3949 let (with_check, _, check_source) = generated_check_proposal(
3950 ExpectedAcceptedHead::Exact {
3951 revision: 1,
3952 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x73; 32]),
3953 },
3954 "abort-add-check",
3955 true,
3956 database,
3957 store,
3958 );
3959 let pending_candidate = lower_existing_schema_proposal(
3960 &with_check,
3961 &[ExistingProposalStore {
3962 path: "abort::Store",
3963 identity: store,
3964 bundle: initial_candidate.bundle(),
3965 }],
3966 )
3967 .expect("generated check should lower")
3968 .pop()
3969 .expect("generated check should produce one candidate");
3970 let entity_tag = pending_candidate
3971 .bundle()
3972 .source_bindings_for_tests()
3973 .entity(&entity_source)
3974 .expect("entity source should remain bound");
3975 let constraint_id = pending_candidate
3976 .bundle()
3977 .source_bindings_for_tests()
3978 .constraint(entity_tag, &check_source)
3979 .expect("generated check source should bind");
3980 let pending_snapshot = pending_candidate
3981 .bundle()
3982 .entity_snapshots()
3983 .get(&entity_tag)
3984 .expect("pending entity should exist");
3985 let activation = pending_snapshot
3986 .constraint_catalog()
3987 .activation(constraint_id)
3988 .expect("generated check should remain an activation");
3989 assert_eq!(activation.origin(), ConstraintOrigin::Generated);
3990
3991 let aborted = aborted_generated_row_local_candidate(
3992 pending_candidate.bundle(),
3993 entity_tag,
3994 constraint_id,
3995 )
3996 .expect("generated check abort should build one catalog-native candidate");
3997 let aborted_snapshot = aborted
3998 .bundle()
3999 .entity_snapshots()
4000 .get(&entity_tag)
4001 .expect("aborted entity should remain");
4002 assert!(
4003 aborted_snapshot
4004 .constraint_catalog()
4005 .activation(constraint_id)
4006 .is_none(),
4007 );
4008 assert_eq!(aborted_snapshot.row_layout(), pending_snapshot.row_layout());
4009 assert!(
4010 aborted
4011 .bundle()
4012 .source_bindings_for_tests()
4013 .constraint(entity_tag, &check_source)
4014 .is_none(),
4015 );
4016
4017 let reproposed = lower_existing_schema_proposal(
4018 &with_check,
4019 &[ExistingProposalStore {
4020 path: "abort::Store",
4021 identity: store,
4022 bundle: aborted.bundle(),
4023 }],
4024 )
4025 .expect("aborted generated check should be independently reproposable")
4026 .pop()
4027 .expect("reproposal should produce one candidate");
4028 let replacement_id = reproposed
4029 .bundle()
4030 .source_bindings_for_tests()
4031 .constraint(entity_tag, &check_source)
4032 .expect("reproposal should bind a fresh constraint identity");
4033 assert!(
4034 replacement_id > constraint_id,
4035 "aborted accepted IDs must remain retired",
4036 );
4037 }
4038
4039 #[test]
4040 fn targeted_rule_edit_abort_keeps_prior_accepted_semantics_and_source_identity() {
4041 let database = TargetDatabaseIdentity::from_bytes([0x81; 32]);
4042 let store = TargetStoreIdentity::from_bytes([0x82; 32]);
4043 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4044 ExpectedAcceptedHead::Empty,
4045 "targeted-abort-initial",
4046 SourceRuleOperation::NumericRangeInclusive {
4047 min: ScalarLiteral::Nat(0),
4048 max: ScalarLiteral::Nat(10),
4049 },
4050 database,
4051 store,
4052 );
4053 let initial_candidate = lower_initial_schema_proposal(
4054 &initial,
4055 &[ProposalStoreTarget {
4056 path: "abort::TargetedStore",
4057 identity: store,
4058 }],
4059 )
4060 .expect("initial targeted proposal should lower")
4061 .pop()
4062 .expect("initial targeted proposal should produce one candidate");
4063 let initial_bundle = initial_candidate.bundle();
4064 let entity_tag = initial_bundle
4065 .source_bindings_for_tests()
4066 .entity(&entity_source)
4067 .expect("entity source should bind");
4068 let constraint_id = initial_bundle
4069 .source_bindings_for_tests()
4070 .constraint(entity_tag, &constraint_source)
4071 .expect("targeted source should bind");
4072 let high_water = initial_bundle.entity_snapshots()[&entity_tag]
4073 .constraint_id_allocator()
4074 .high_water();
4075 let (edited, _, _) = targeted_rule_proposal(
4076 ExpectedAcceptedHead::Exact {
4077 revision: initial_bundle.revision().get(),
4078 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x83; 32]),
4079 },
4080 "targeted-abort-edit",
4081 SourceRuleOperation::NumericMaximumInclusive {
4082 value: ScalarLiteral::Nat(8),
4083 },
4084 database,
4085 store,
4086 );
4087 let staged = lower_existing_schema_proposal(
4088 &edited,
4089 &[ExistingProposalStore {
4090 path: "abort::TargetedStore",
4091 identity: store,
4092 bundle: initial_bundle,
4093 }],
4094 )
4095 .expect("targeted semantic edit should stage")
4096 .pop()
4097 .expect("targeted semantic edit should produce one candidate");
4098 let aborted =
4099 aborted_generated_row_local_candidate(staged.bundle(), entity_tag, constraint_id)
4100 .expect("targeted semantic edit should abort through catalog authority");
4101 let snapshot = &aborted.bundle().entity_snapshots()[&entity_tag];
4102
4103 assert!(
4104 snapshot
4105 .constraint_catalog()
4106 .activation(constraint_id)
4107 .is_none()
4108 );
4109 assert_eq!(snapshot.constraint_id_allocator().high_water(), high_water);
4110 assert_eq!(
4111 aborted
4112 .bundle()
4113 .source_bindings_for_tests()
4114 .constraint(entity_tag, &constraint_source),
4115 Some(constraint_id),
4116 );
4117 assert!(snapshot.constraints().iter().any(|constraint| {
4118 constraint.id() == constraint_id
4119 && matches!(
4120 constraint.kind(),
4121 AcceptedConstraintKind::TargetedRule { operation, .. }
4122 if matches!(
4123 operation.as_ref(),
4124 AcceptedRuleOperation::NumericRangeInclusive { .. }
4125 )
4126 )
4127 }));
4128 }
4129
4130 #[test]
4131 #[allow(
4132 clippy::too_many_lines,
4133 reason = "the staged publication, recovery, and promotion assertions form one lifecycle"
4134 )]
4135 fn generated_startup_driver_resumes_pending_activation_and_promotes_without_source_model() {
4136 let db = Db::<EvolutionCanister>::new(
4137 &EVOLUTION_REGISTRY,
4138 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4139 );
4140 drive_startup_recovery_to_completion(&db);
4141 let empty_target =
4142 schema_application_target(&db).expect("empty evolution target should issue");
4143 let store_identity = empty_target
4144 .stores()
4145 .first()
4146 .expect("evolution store should register")
4147 .identity();
4148 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4149 empty_target.accepted_head().clone(),
4150 "targeted-recovery-initial",
4151 SourceRuleOperation::NumericRangeInclusive {
4152 min: ScalarLiteral::Nat(0),
4153 max: ScalarLiteral::Nat(10),
4154 },
4155 empty_target.database_identity(),
4156 store_identity,
4157 );
4158 assert!(matches!(
4159 apply_schema(&db, &initial)
4160 .expect("initial targeted proposal should publish")
4161 .outcome(),
4162 SchemaChangeOutcome::Applied { .. },
4163 ));
4164 drive_startup_recovery_to_completion(&db);
4165 drive_cardinality_to_ready(
4166 db.store_handle(EVOLUTION_STORE_PATH)
4167 .expect("evolution store should resolve"),
4168 );
4169
4170 let direct_target =
4171 schema_application_target(&db).expect("direct evolution target should issue");
4172 let (direct_edit, _, _) = targeted_rule_proposal(
4173 direct_target.accepted_head().clone(),
4174 "targeted-direct-edit",
4175 SourceRuleOperation::NumericMaximumInclusive {
4176 value: ScalarLiteral::Nat(8),
4177 },
4178 direct_target.database_identity(),
4179 store_identity,
4180 );
4181 assert!(matches!(
4182 apply_schema(&db, &direct_edit)
4183 .expect("empty-domain semantic edit should publish directly")
4184 .outcome(),
4185 SchemaChangeOutcome::Applied { .. },
4186 ));
4187 let store = db
4188 .store_handle(EVOLUTION_STORE_PATH)
4189 .expect("evolution store should resolve");
4190 let direct = store
4191 .with_schema(SchemaStore::current_accepted_schema_bundle)
4192 .expect("directly edited bundle should remain readable")
4193 .expect("directly edited bundle should exist");
4194 let entity_tag = direct
4195 .source_bindings_for_tests()
4196 .entity(&entity_source)
4197 .expect("entity source should remain bound");
4198 let constraint_id = direct
4199 .source_bindings_for_tests()
4200 .constraint(entity_tag, &constraint_source)
4201 .expect("direct edit should preserve constraint identity");
4202 assert!(
4203 direct.entity_snapshots()[&entity_tag]
4204 .constraint_catalog()
4205 .activation(constraint_id)
4206 .is_none()
4207 );
4208 assert!(
4209 direct.entity_snapshots()[&entity_tag]
4210 .constraints()
4211 .iter()
4212 .any(|constraint| {
4213 constraint.id() == constraint_id
4214 && matches!(
4215 constraint.kind(),
4216 AcceptedConstraintKind::TargetedRule { operation, .. }
4217 if matches!(
4218 operation.as_ref(),
4219 AcceptedRuleOperation::NumericMaximumInclusive { .. }
4220 )
4221 )
4222 })
4223 );
4224
4225 let target = schema_application_target(&db).expect("staged evolution target should issue");
4226 let (edited, _, _) = targeted_rule_proposal(
4227 target.accepted_head().clone(),
4228 "targeted-recovery-edit",
4229 SourceRuleOperation::MultipleOf {
4230 divisor: ScalarLiteral::Nat(2),
4231 },
4232 target.database_identity(),
4233 store_identity,
4234 );
4235 let current = store
4236 .with_schema(SchemaStore::current_accepted_schema_bundle)
4237 .expect("accepted evolution bundle should remain readable")
4238 .expect("directly edited evolution bundle should exist");
4239 let staged = lower_existing_schema_proposal(
4240 &edited,
4241 &[ExistingProposalStore {
4242 path: EVOLUTION_STORE_PATH,
4243 identity: store_identity,
4244 bundle: ¤t,
4245 }],
4246 )
4247 .expect("targeted edit should stage")
4248 .pop()
4249 .expect("targeted edit should produce one candidate");
4250 assert_eq!(
4251 staged
4252 .bundle()
4253 .source_bindings_for_tests()
4254 .constraint(entity_tag, &constraint_source),
4255 Some(constraint_id),
4256 );
4257 let proof = DirectGeneratedRowLocalProof {
4258 candidate_index: 0,
4259 store,
4260 store_path: EVOLUTION_STORE_PATH,
4261 entity_tag,
4262 entity_path: staged.bundle().entity_snapshots()[&entity_tag]
4263 .entity_path()
4264 .to_string(),
4265 constraint_id,
4266 historical_rows: 0,
4267 };
4268 let final_candidates = final_candidates_for_pending_row_local_constraint(
4269 std::slice::from_ref(&staged),
4270 &PendingGeneratedRowLocalConstraint { proof },
4271 )
4272 .expect("final semantic replacement should derive without source input");
4273 let authorities = application_authorities(&db);
4274 let candidate_head =
4275 accepted_head_after_candidates(authorities.as_slice(), &final_candidates)
4276 .expect("final candidate head should derive");
4277 let digest = edited.digest().expect("proposal digest should derive");
4278 let job_id = derive_schema_change_job_id(
4279 target.database_identity(),
4280 edited.submission_key(),
4281 digest,
4282 target.accepted_head(),
4283 )
4284 .expect("job identity should derive");
4285 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4286 target.database_identity(),
4287 edited.submission_key().clone(),
4288 digest,
4289 target.accepted_head().clone(),
4290 SchemaChangeOutcome::Pending {
4291 job: SchemaChangeJob::new(job_id),
4292 candidate_head,
4293 },
4294 )
4295 .expect("pending replacement receipt should admit");
4296 let record = SchemaApplicationRecord::new(
4297 receipt,
4298 vec![
4299 SchemaChangeActivation::new(
4300 store_identity,
4301 entity_tag.value(),
4302 constraint_id.get(),
4303 )
4304 .expect("replacement activation should admit"),
4305 ],
4306 )
4307 .expect("pending replacement record should admit");
4308 let operation =
4309 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4310 publish_accepted_schema_candidates_with_application_record(
4311 vec![AcceptedSchemaPublication::new(
4312 EVOLUTION_STORE_PATH,
4313 store,
4314 current.revision(),
4315 &staged,
4316 )],
4317 operation,
4318 )
4319 .expect("staged replacement and record should publish atomically");
4320
4321 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4322 drive_startup_recovery_to_completion(&db);
4323
4324 let recovered = store
4325 .with_schema(SchemaStore::current_accepted_schema_bundle)
4326 .expect("recovered staged bundle should decode")
4327 .expect("recovered staged bundle should exist");
4328 let recovered_snapshot = recovered.entity_snapshots()[&entity_tag].clone();
4329 let validating_catalog = recovered_snapshot
4330 .constraint_catalog()
4331 .clone()
4332 .with_validation_started(constraint_id)
4333 .expect("recovered replacement should enter validation");
4334 let mut validating_snapshots = recovered.entity_snapshots().clone();
4335 validating_snapshots.insert(
4336 entity_tag,
4337 recovered_snapshot.with_constraint_catalog(validating_catalog),
4338 );
4339 let validating_bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
4340 recovered
4341 .revision()
4342 .checked_next()
4343 .expect("validation revision should remain available"),
4344 recovered.store_path(),
4345 recovered.enum_catalog().clone(),
4346 recovered.composite_catalog().clone(),
4347 recovered.source_bindings_for_tests().clone(),
4348 validating_snapshots,
4349 )
4350 .expect("validating replacement bundle should close");
4351 let validating_candidate = CandidateSchemaRevision::new(validating_bundle)
4352 .expect("validating replacement candidate should encode");
4353 let validating_activation = validating_candidate.bundle().entity_snapshots()[&entity_tag]
4354 .constraint_catalog()
4355 .activation(constraint_id)
4356 .expect("validating replacement activation should remain present");
4357 let validation_job = ConstraintValidationJob::start(
4358 entity_tag,
4359 validating_candidate.bundle().entity_snapshots()[&entity_tag]
4360 .entity_path()
4361 .to_string(),
4362 validating_activation,
4363 None,
4364 )
4365 .expect("validating replacement job should derive from accepted state");
4366 store
4367 .with_schema(|schema| {
4368 schema.validate_live_activation_transition(validating_candidate.bundle())?;
4369 schema.validate_constraint_validation_job_closure_with_change(
4370 validating_candidate.bundle(),
4371 Some(&validation_job),
4372 None,
4373 )
4374 })
4375 .expect("validating replacement transition and job should close");
4376 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4377 let startup_session =
4378 crate::db::DbSession::<EvolutionCanister>::new(&EVOLUTION_REGISTRY, &startup_root);
4379
4380 assert_eq!(
4381 drive_generated_startup_recovery_page(
4382 &startup_session,
4383 &EVOLUTION_REGISTRY,
4384 edited.submission_key().as_str(),
4385 )
4386 .expect("generated startup should begin pending validation"),
4387 GeneratedStartupDriverStep::Recovering,
4388 );
4389 let mut terminal = false;
4390 for _ in 0..8 {
4391 match drive_generated_startup_recovery_page(
4392 &startup_session,
4393 &EVOLUTION_REGISTRY,
4394 edited.submission_key().as_str(),
4395 )
4396 .expect("generated startup should advance pending validation")
4397 {
4398 GeneratedStartupDriverStep::Recovering => {}
4399 GeneratedStartupDriverStep::Terminal => {
4400 terminal = true;
4401 break;
4402 }
4403 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4404 panic!("an exact pending receipt must resume instead of being resubmitted")
4405 }
4406 }
4407 }
4408 assert!(
4409 terminal,
4410 "empty historical domain should promote within bounded startup steps"
4411 );
4412 assert_eq!(
4413 observe_generated_startup_state::<EvolutionCanister>(
4414 &EVOLUTION_REGISTRY,
4415 edited.submission_key().as_str(),
4416 ),
4417 Ok(DatabaseStartupState::Ready),
4418 );
4419 let applied = super::exact_schema_application_receipt(
4420 &edited,
4421 edited
4422 .digest()
4423 .expect("proposal digest should remain stable"),
4424 )
4425 .expect("terminal generated receipt should remain readable")
4426 .expect("terminal generated receipt should remain present");
4427 assert!(matches!(
4428 applied.outcome(),
4429 SchemaChangeOutcome::Applied { .. }
4430 ));
4431 let promoted = store
4432 .with_schema(SchemaStore::current_accepted_schema_bundle)
4433 .expect("promoted bundle should remain readable")
4434 .expect("promoted bundle should exist");
4435 let snapshot = &promoted.entity_snapshots()[&entity_tag];
4436 assert!(
4437 snapshot
4438 .constraint_catalog()
4439 .activation(constraint_id)
4440 .is_none()
4441 );
4442 assert_eq!(
4443 promoted
4444 .source_bindings_for_tests()
4445 .constraint(entity_tag, &constraint_source),
4446 Some(constraint_id),
4447 );
4448 assert!(snapshot.constraints().iter().any(|constraint| {
4449 constraint.id() == constraint_id
4450 && matches!(
4451 constraint.kind(),
4452 AcceptedConstraintKind::TargetedRule { operation, .. }
4453 if matches!(
4454 operation.as_ref(),
4455 AcceptedRuleOperation::MultipleOf { .. }
4456 )
4457 )
4458 }));
4459 }
4460
4461 #[test]
4462 #[allow(
4463 clippy::too_many_lines,
4464 reason = "the durable pending job, startup failure, and retained finding assertions form one scenario"
4465 )]
4466 fn generated_startup_driver_persists_e223_for_a_retained_historical_finding() {
4467 let db = Db::<AbortCanister>::new(
4468 &ABORT_REGISTRY,
4469 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4470 );
4471 drive_startup_recovery_to_completion(&db);
4472 let empty_target =
4473 schema_application_target(&db).expect("empty application target should issue");
4474 let store_identity = empty_target
4475 .stores()
4476 .first()
4477 .expect("abort store should be registered")
4478 .identity();
4479 let (initial, _, _) = generated_check_proposal(
4480 empty_target.accepted_head().clone(),
4481 "startup-finding-initial",
4482 false,
4483 empty_target.database_identity(),
4484 store_identity,
4485 );
4486 apply_schema(&db, &initial).expect("initial generated schema should publish");
4487
4488 let root = crate::db::RequestExecutionRoot::__new_runtime_root();
4489 let session = DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &root);
4490 let rows = (1..=257)
4491 .map(|id| DynamicMutation::Insert {
4492 entity: "Item".to_string(),
4493 patch: DynamicStructuralPatch::new(vec![
4494 (
4495 "id".to_string(),
4496 DynamicWriteCell::Value(InputValue::nat64(id)),
4497 ),
4498 (
4499 "score".to_string(),
4500 DynamicWriteCell::Value(InputValue::int64(if id == 257 { -1 } else { 1 })),
4501 ),
4502 ]),
4503 })
4504 .collect();
4505 session
4506 .execute_trusted_dynamic_mutation_batch(rows)
4507 .expect("historical finding fixture rows should commit as one legal batch");
4508 drive_startup_recovery_to_completion(&db);
4509 drive_cardinality_to_ready(
4510 db.store_handle(ABORT_STORE_PATH)
4511 .expect("abort store should resolve"),
4512 );
4513
4514 let target =
4515 schema_application_target(&db).expect("existing application target should issue");
4516 let (with_check, _, _) = generated_check_proposal(
4517 target.accepted_head().clone(),
4518 "startup-finding-pending",
4519 true,
4520 target.database_identity(),
4521 store_identity,
4522 );
4523 let pending = apply_schema(&db, &with_check)
4524 .expect("the first clean page should admit durable continuation");
4525 let SchemaChangeOutcome::Pending { job, .. } = pending.outcome() else {
4526 panic!("a 257-row domain must exceed the 256-row direct proof page")
4527 };
4528
4529 let mut terminal = false;
4530 for _ in 0..8 {
4531 match drive_generated_startup_recovery_page(
4532 &session,
4533 &ABORT_REGISTRY,
4534 with_check.submission_key().as_str(),
4535 )
4536 .expect("generated startup should retain a typed finding failure")
4537 {
4538 GeneratedStartupDriverStep::Recovering => {}
4539 GeneratedStartupDriverStep::Terminal => {
4540 terminal = true;
4541 break;
4542 }
4543 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4544 panic!("an exact pending receipt must not be resubmitted")
4545 }
4546 }
4547 }
4548 assert!(terminal, "the retained finding should become terminal");
4549 let failure = observe_generated_startup_state::<AbortCanister>(
4550 &ABORT_REGISTRY,
4551 with_check.submission_key().as_str(),
4552 )
4553 .expect_err("the retained finding must remain durably observable");
4554 assert_eq!(
4555 failure.kind(),
4556 crate::db::StartupFailureKind::SchemaReconciliation,
4557 );
4558 assert_eq!(
4559 failure.diagnostic().error_code(),
4560 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_CONSTRAINT_VIOLATION,
4561 );
4562 assert_eq!(
4563 ABORT_DATA.with(|store| store.borrow().len()),
4564 257,
4565 "terminal startup publication must not change historical rows",
4566 );
4567 assert!(matches!(
4568 continue_schema_application(&db, job.id(), None)
4569 .expect("the retained finding page should replay exactly")
4570 .status(),
4571 SchemaChangeProgressStatus::Findings { findings, .. } if !findings.is_empty(),
4572 ));
4573 }
4574
4575 #[test]
4576 #[allow(
4577 clippy::too_many_lines,
4578 reason = "the journaled abort, replay, and recovery assertions form one scenario"
4579 )]
4580 fn pending_generated_check_abort_is_atomic_terminal_and_replayable() {
4581 let db = Db::<AbortCanister>::new(
4582 &ABORT_REGISTRY,
4583 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4584 );
4585 drive_startup_recovery_to_completion(&db);
4586 let empty_target =
4587 schema_application_target(&db).expect("empty application target should issue");
4588 let store_identity = empty_target
4589 .stores()
4590 .first()
4591 .expect("abort store should be registered")
4592 .identity();
4593 let (initial, entity_source, _) = generated_check_proposal(
4594 empty_target.accepted_head().clone(),
4595 "abort-runtime-initial",
4596 false,
4597 empty_target.database_identity(),
4598 store_identity,
4599 );
4600 assert!(matches!(
4601 apply_schema(&db, &initial)
4602 .expect("initial application should publish")
4603 .outcome(),
4604 SchemaChangeOutcome::Applied { .. },
4605 ));
4606
4607 let target =
4608 schema_application_target(&db).expect("existing application target should issue");
4609 let (with_check, _, check_source) = generated_check_proposal(
4610 target.accepted_head().clone(),
4611 "abort-runtime-pending",
4612 true,
4613 target.database_identity(),
4614 store_identity,
4615 );
4616 let store = db
4617 .store_handle(ABORT_STORE_PATH)
4618 .expect("abort store should resolve");
4619 let current = store
4620 .with_schema(SchemaStore::current_accepted_schema_bundle)
4621 .expect("accepted bundle should remain readable")
4622 .expect("initial accepted bundle should exist");
4623 let pending_candidate = lower_existing_schema_proposal(
4624 &with_check,
4625 &[ExistingProposalStore {
4626 path: ABORT_STORE_PATH,
4627 identity: store_identity,
4628 bundle: ¤t,
4629 }],
4630 )
4631 .expect("pending generated check should lower")
4632 .pop()
4633 .expect("pending generated check should produce one candidate");
4634 let entity_tag = pending_candidate
4635 .bundle()
4636 .source_bindings_for_tests()
4637 .entity(&entity_source)
4638 .expect("entity source should bind");
4639 let constraint_id = pending_candidate
4640 .bundle()
4641 .source_bindings_for_tests()
4642 .constraint(entity_tag, &check_source)
4643 .expect("generated check source should bind");
4644 let digest = with_check.digest().expect("proposal digest should derive");
4645 let job_id = derive_schema_change_job_id(
4646 target.database_identity(),
4647 with_check.submission_key(),
4648 digest,
4649 target.accepted_head(),
4650 )
4651 .expect("job identity should derive");
4652 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4653 target.database_identity(),
4654 with_check.submission_key().clone(),
4655 digest,
4656 target.accepted_head().clone(),
4657 SchemaChangeOutcome::Pending {
4658 job: SchemaChangeJob::new(job_id),
4659 candidate_head: ExpectedAcceptedHead::Exact {
4660 revision: pending_candidate.revision().get().saturating_add(2),
4661 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x76; 32]),
4662 },
4663 },
4664 )
4665 .expect("pending receipt should admit");
4666 let record = SchemaApplicationRecord::new(
4667 receipt,
4668 vec![
4669 SchemaChangeActivation::new(
4670 store_identity,
4671 entity_tag.value(),
4672 constraint_id.get(),
4673 )
4674 .expect("application activation should admit"),
4675 ],
4676 )
4677 .expect("pending application record should admit");
4678 let operation =
4679 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4680 publish_accepted_schema_candidates_with_application_record(
4681 vec![AcceptedSchemaPublication::new(
4682 ABORT_STORE_PATH,
4683 store,
4684 current.revision(),
4685 &pending_candidate,
4686 )],
4687 operation,
4688 )
4689 .expect("pending candidate and record should publish atomically");
4690
4691 let started = continue_schema_application(&db, job_id, None)
4692 .expect("first continuation should durably start validation");
4693 assert_eq!(started.status(), &SchemaChangeProgressStatus::Started);
4694 let progress =
4695 abort_schema_application(&db, job_id, None).expect("pending application should abort");
4696 assert_eq!(progress.status(), &SchemaChangeProgressStatus::Aborted);
4697 assert!(matches!(
4698 progress.receipt().outcome(),
4699 SchemaChangeOutcome::Aborted { .. },
4700 ));
4701 let replay =
4702 abort_schema_application(&db, job_id, None).expect("terminal abort should replay");
4703 assert_eq!(replay, progress);
4704 assert_eq!(
4705 continue_schema_application(&db, job_id, None)
4706 .expect("continuation after abort should replay terminal state"),
4707 progress,
4708 );
4709
4710 let aborted = store
4711 .with_schema(SchemaStore::current_accepted_schema_bundle)
4712 .expect("accepted bundle should remain readable")
4713 .expect("aborted accepted bundle should exist");
4714 assert!(
4715 aborted
4716 .entity_snapshots()
4717 .get(&entity_tag)
4718 .expect("entity should remain after abort")
4719 .constraint_catalog()
4720 .activation(constraint_id)
4721 .is_none(),
4722 );
4723 assert!(
4724 aborted
4725 .source_bindings_for_tests()
4726 .constraint(entity_tag, &check_source)
4727 .is_none(),
4728 );
4729 assert!(
4730 store
4731 .with_schema(|schema| {
4732 schema.constraint_validation_job(entity_tag, constraint_id)
4733 })
4734 .expect("validation-job storage should remain readable")
4735 .is_none(),
4736 );
4737
4738 ABORT_DATA.with(|store| {
4739 ABORT_DATA_MEMORY.with(|memory| {
4740 *store.borrow_mut() = DataStore::init_journaled(memory.clone());
4741 });
4742 });
4743 ABORT_INDEX.with(|store| {
4744 ABORT_INDEX_MEMORY.with(|memory| {
4745 *store.borrow_mut() = IndexStore::init_journaled(memory.clone());
4746 });
4747 });
4748 ABORT_SCHEMA.with(|store| {
4749 ABORT_SCHEMA_MEMORY.with(|memory| {
4750 *store.borrow_mut() = SchemaStore::init_journaled(memory.clone());
4751 });
4752 });
4753 ABORT_JOURNAL.with(|store| {
4754 ABORT_JOURNAL_MEMORY.with(|memory| {
4755 *store.borrow_mut() = JournalTailStore::init(memory.clone());
4756 });
4757 });
4758 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4759 drive_startup_recovery_to_completion(&db);
4760 assert_eq!(
4761 abort_schema_application(&db, job_id, None)
4762 .expect("recovered terminal abort should replay"),
4763 progress,
4764 );
4765 assert!(
4766 store
4767 .with_schema(|schema| {
4768 schema.constraint_validation_job(entity_tag, constraint_id)
4769 })
4770 .expect("recovered validation-job storage should remain readable")
4771 .is_none(),
4772 );
4773 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4774 let startup_session =
4775 crate::db::DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &startup_root);
4776 ABORT_JOURNAL.with(|journal| {
4777 let journal = journal.borrow();
4778 assert!(
4779 journal
4780 .validate_current_tail_authority()
4781 .expect("recovered abort tail control should remain readable")
4782 .is_empty(),
4783 "recovered abort tail control should close exactly",
4784 );
4785 });
4786 assert_eq!(
4787 drive_generated_startup_recovery_page(
4788 &startup_session,
4789 &ABORT_REGISTRY,
4790 with_check.submission_key().as_str(),
4791 )
4792 .expect("an exact aborted generated submission should publish terminal startup state"),
4793 GeneratedStartupDriverStep::Terminal,
4794 );
4795 let failure = observe_generated_startup_state::<AbortCanister>(
4796 &ABORT_REGISTRY,
4797 with_check.submission_key().as_str(),
4798 )
4799 .expect_err("an aborted generated submission must not remain retryable forever");
4800 assert_eq!(
4801 failure.kind(),
4802 crate::db::StartupFailureKind::SchemaReconciliation,
4803 );
4804 assert_eq!(
4805 failure.diagnostic().error_code(),
4806 icydb_diagnostic_code::ErrorCode::RUNTIME_CONFLICT,
4807 );
4808 }
4809
4810 #[cfg(feature = "migration")]
4811 #[test]
4812 fn migration_planning_failures_retain_typed_public_classification() {
4813 use super::schema_migration_planning_error;
4814 use crate::db::schema::migration_planner::SchemaMigrationPlanningError;
4815 use icydb_diagnostic_code::{DiagnosticDetail, SchemaMigrationCode};
4816
4817 for (error, reason) in [
4818 (
4819 SchemaMigrationPlanningError::Unadopted,
4820 SchemaMigrationCode::Unadopted,
4821 ),
4822 (
4823 SchemaMigrationPlanningError::MissingMigration,
4824 SchemaMigrationCode::MissingMigration,
4825 ),
4826 (
4827 SchemaMigrationPlanningError::VersionGap,
4828 SchemaMigrationCode::VersionGap,
4829 ),
4830 (
4831 SchemaMigrationPlanningError::Downgrade,
4832 SchemaMigrationCode::Downgrade,
4833 ),
4834 (
4835 SchemaMigrationPlanningError::EmptyEntityVersionBump,
4836 SchemaMigrationCode::EmptyEntityVersionBump,
4837 ),
4838 (
4839 SchemaMigrationPlanningError::StaleAcceptedHead,
4840 SchemaMigrationCode::StaleAcceptedHead,
4841 ),
4842 (
4843 SchemaMigrationPlanningError::UnknownFromObject,
4844 SchemaMigrationCode::UnknownFromObject,
4845 ),
4846 (
4847 SchemaMigrationPlanningError::UnknownToObject,
4848 SchemaMigrationCode::UnknownToObject,
4849 ),
4850 (
4851 SchemaMigrationPlanningError::KindMismatch,
4852 SchemaMigrationCode::KindMismatch,
4853 ),
4854 (
4855 SchemaMigrationPlanningError::IdentityConflict,
4856 SchemaMigrationCode::IdentityConflict,
4857 ),
4858 (
4859 SchemaMigrationPlanningError::UnexplainedSchemaDifference,
4860 SchemaMigrationCode::UnexplainedSchemaDifference,
4861 ),
4862 (
4863 SchemaMigrationPlanningError::UnsupportedTransform,
4864 SchemaMigrationCode::UnsupportedTransform,
4865 ),
4866 (
4867 SchemaMigrationPlanningError::RekeyedCatalogInvalid,
4868 SchemaMigrationCode::CandidateMismatch,
4869 ),
4870 (
4871 SchemaMigrationPlanningError::CandidateMismatch,
4872 SchemaMigrationCode::CandidateMismatch,
4873 ),
4874 (
4875 SchemaMigrationPlanningError::CorruptLineage,
4876 SchemaMigrationCode::ProgressCorrupt,
4877 ),
4878 ] {
4879 let diagnostic = schema_migration_planning_error(error).diagnostic();
4880 assert_eq!(
4881 diagnostic.detail(),
4882 Some(&DiagnosticDetail::SchemaMigration { reason }),
4883 );
4884 assert_eq!(diagnostic.code(), reason.diagnostic_code());
4885 }
4886 }
4887
4888 #[cfg(feature = "migration")]
4889 #[test]
4890 #[expect(
4891 clippy::too_many_lines,
4892 reason = "the validation replay, staging, and unchanged-row assertions form one scenario"
4893 )]
4894 fn physical_migration_validation_is_bounded_staged_and_does_not_rewrite_rows() {
4895 use std::convert::Infallible;
4896
4897 use super::{defer_generated_schema_application_for_prepared_migration, migrate_schema};
4898 use crate::db::{
4899 data::StoreVisit,
4900 index::{IndexEntryValue, IndexId, IndexKey, IndexKeyKind},
4901 key_taxonomy::{PrimaryKeyComponent, PrimaryKeyValue},
4902 schema::{SchemaMigrationCommand, SchemaMigrationPhase},
4903 };
4904 use crate::types::EntityTag;
4905
4906 let db = Db::<MigrationCanister>::new(
4907 &MIGRATION_REGISTRY,
4908 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4909 );
4910 drive_startup_recovery_to_completion(&db);
4911 let initial_target = schema_application_target(&db).expect("initial target should issue");
4912 let store_identity = initial_target
4913 .stores()
4914 .first()
4915 .expect("migration store should exist")
4916 .identity();
4917 let initial = validation_migration_proposal(
4918 ValidationMigrationShape::Clean,
4919 false,
4920 initial_target.accepted_head().clone(),
4921 initial_target.database_identity(),
4922 store_identity,
4923 );
4924 apply_schema(&db, &initial).expect("initial schema should publish");
4925
4926 let session = DbSession::<MigrationCanister>::new(
4927 &MIGRATION_REGISTRY,
4928 &crate::db::RequestExecutionRoot::__new_runtime_root(),
4929 );
4930 for (id, value) in [(1, 7), (2, 8)] {
4931 session
4932 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4933 entity: "MigratingItem".to_string(),
4934 patch: DynamicStructuralPatch::new(vec![
4935 (
4936 "id".to_string(),
4937 DynamicWriteCell::Value(InputValue::nat64(id)),
4938 ),
4939 (
4940 "old_value".to_string(),
4941 DynamicWriteCell::Value(InputValue::int64(value)),
4942 ),
4943 ]),
4944 })
4945 .expect("predecessor row should insert");
4946 }
4947 let store = db
4948 .store_handle(MIGRATION_STORE_PATH)
4949 .expect("migration store should resolve");
4950 let row_bytes = || {
4951 store.with_data(|data| {
4952 let mut rows = Vec::new();
4953 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
4954 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
4955 Ok(StoreVisit::Continue)
4956 });
4957 result.expect("infallible row visit should complete");
4958 rows
4959 })
4960 };
4961 let before_rows = row_bytes();
4962
4963 let target = schema_application_target(&db).expect("migration target should issue");
4964 let proposal = validation_migration_proposal(
4965 ValidationMigrationShape::Clean,
4966 true,
4967 target.accepted_head().clone(),
4968 target.database_identity(),
4969 store_identity,
4970 );
4971 let plan = proposal
4972 .migration()
4973 .expect("migration plan should exist")
4974 .digest();
4975 let command = || SchemaMigrationCommand::Advance {
4976 expected_database: target.database_identity(),
4977 expected_head: target.accepted_head().clone(),
4978 expected_plan: plan,
4979 acknowledged_finding_page: None,
4980 };
4981 assert_eq!(
4982 migrate_schema(&db, &proposal, command())
4983 .unwrap_or_else(|error| {
4984 panic!(
4985 "physical migration should prepare: {:?}",
4986 error.diagnostic()
4987 )
4988 })
4989 .phase(),
4990 SchemaMigrationPhase::Prepared,
4991 );
4992 assert_eq!(
4993 migrate_schema(&db, &proposal, command())
4994 .expect("physical migration should enter validation")
4995 .phase(),
4996 SchemaMigrationPhase::Validating,
4997 );
4998 let record = super::load_schema_migration_record()
4999 .expect("migration record should remain readable")
5000 .expect("validating migration record should exist");
5001 let planned = super::recompile_active_physical_migration(&db, &proposal, &record)
5002 .expect("the exact active plan should recompile");
5003 for _ in 0..2 {
5004 let page = super::validate_migration_page(&db, &planned, record.progress())
5005 .expect("the same validation page should remain replayable");
5006 let (progress, staged, exhausted) = page.into_parts();
5007 assert!(progress.findings().is_empty());
5008 assert!(exhausted);
5009 super::stage_migration_index_entries(staged)
5010 .expect("staging before a cursor marker should be idempotent");
5011 }
5012 assert_eq!(
5013 store.with_index(IndexStore::len),
5014 2,
5015 "replaying an uncheckpointed page must retain one exact staged key per row",
5016 );
5017 let ready =
5018 migrate_schema(&db, &proposal, command()).expect("bounded validation should complete");
5019 assert_eq!(ready.phase(), SchemaMigrationPhase::ReadyToRewrite);
5020 assert_eq!(ready.rows_validated(), 2);
5021 assert!(ready.findings().is_empty());
5022 assert_eq!(row_bytes(), before_rows, "validation must not rewrite rows");
5023 assert_eq!(
5024 store.with_index(IndexStore::len),
5025 2,
5026 "the isolated candidate unique generation should be durably staged",
5027 );
5028 store.with_index_mut(|index| {
5029 for ordinal in 0..513_u64 {
5030 let component = ordinal.to_be_bytes();
5031 let key = IndexKey::new_from_components_with_primary_key_value(
5032 &IndexId::new(EntityTag::new(2), 0),
5033 IndexKeyKind::User,
5034 &[component],
5035 &PrimaryKeyValue::from(PrimaryKeyComponent::Nat64(ordinal)),
5036 )
5037 .expect("unrelated abort-scan key should build")
5038 .to_raw()
5039 .expect("unrelated abort-scan key should encode");
5040 index.insert(key, IndexEntryValue::presence());
5041 }
5042 });
5043 let abort = || SchemaMigrationCommand::Abort {
5044 expected_database: target.database_identity(),
5045 expected_head: target.accepted_head().clone(),
5046 expected_plan: plan,
5047 };
5048 let cleaning = migrate_schema(&db, &proposal, abort())
5049 .expect("the first bounded abort cleanup page should publish");
5050 assert_eq!(cleaning.phase(), SchemaMigrationPhase::ReadyToRewrite);
5051 assert_eq!(store.with_index(IndexStore::len), 513);
5052 let aborted =
5053 migrate_schema(&db, &proposal, abort()).expect("pre-rewrite migration should abort");
5054 assert_eq!(aborted.phase(), SchemaMigrationPhase::Aborted);
5055 assert_eq!(
5056 store.with_index(IndexStore::len),
5057 513,
5058 "abort must remove only planner-invisible candidate generations",
5059 );
5060 assert_eq!(
5061 row_bytes(),
5062 before_rows,
5063 "abort must retain predecessor rows"
5064 );
5065 assert!(
5066 !defer_generated_schema_application_for_prepared_migration(&db, &proposal)
5067 .expect("terminal aborted record must not block generated startup"),
5068 );
5069 }
5070
5071 #[cfg(feature = "migration")]
5072 #[test]
5073 #[expect(
5074 clippy::too_many_lines,
5075 reason = "the interrupted rewrite, recovery, final proof, and publication form one scenario"
5076 )]
5077 fn physical_migration_rewrite_recovers_and_publishes_one_complete_candidate() {
5078 use super::{
5079 defer_generated_schema_application_for_prepared_migration, migrate_schema,
5080 schema_migration_status_for_target,
5081 };
5082 use crate::db::{
5083 data::{CanonicalSlotReader, DecodedDataStoreKey, StoreVisit, StructuralSlotReader},
5084 schema::{
5085 MigrationRewriteInterruption, SchemaMigrationCommand, SchemaMigrationPhase,
5086 ensure_schema_migration_ready_for_ordinary_operations,
5087 interrupt_next_migration_rewrite_at,
5088 },
5089 };
5090 use crate::error::InternalError;
5091
5092 let db = Db::<MigrationExecutionCanister>::new(
5093 &MIGRATION_EXECUTION_REGISTRY,
5094 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5095 );
5096 drive_startup_recovery_to_completion(&db);
5097 let initial_target = schema_application_target(&db).expect("initial target should issue");
5098 let store_identity = initial_target
5099 .stores()
5100 .first()
5101 .expect("migration execution store should exist")
5102 .identity();
5103 let initial = validation_migration_proposal(
5104 ValidationMigrationShape::Clean,
5105 false,
5106 initial_target.accepted_head().clone(),
5107 initial_target.database_identity(),
5108 store_identity,
5109 );
5110 apply_schema(&db, &initial).expect("initial schema should publish");
5111 let session = DbSession::<MigrationExecutionCanister>::new(
5112 &MIGRATION_EXECUTION_REGISTRY,
5113 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5114 );
5115 for (id, value) in [(1, 7), (2, 8), (3, 9)] {
5116 session
5117 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5118 entity: "MigratingItem".to_string(),
5119 patch: DynamicStructuralPatch::new(vec![
5120 (
5121 "id".to_string(),
5122 DynamicWriteCell::Value(InputValue::nat64(id)),
5123 ),
5124 (
5125 "old_value".to_string(),
5126 DynamicWriteCell::Value(InputValue::int64(value)),
5127 ),
5128 ]),
5129 })
5130 .expect("predecessor row should insert");
5131 }
5132 let target = schema_application_target(&db).expect("migration target should issue");
5133 let proposal = validation_migration_proposal(
5134 ValidationMigrationShape::Clean,
5135 true,
5136 target.accepted_head().clone(),
5137 target.database_identity(),
5138 store_identity,
5139 );
5140 let plan = proposal
5141 .migration()
5142 .expect("migration plan should exist")
5143 .digest();
5144 let command = || SchemaMigrationCommand::Advance {
5145 expected_database: target.database_identity(),
5146 expected_head: target.accepted_head().clone(),
5147 expected_plan: plan,
5148 acknowledged_finding_page: None,
5149 };
5150 for expected in [
5151 SchemaMigrationPhase::Prepared,
5152 SchemaMigrationPhase::Validating,
5153 SchemaMigrationPhase::ReadyToRewrite,
5154 SchemaMigrationPhase::RewritingRows,
5155 ] {
5156 assert_eq!(
5157 migrate_schema(&db, &proposal, command())
5158 .expect("migration phase should advance")
5159 .phase(),
5160 expected,
5161 );
5162 }
5163
5164 for interruption in [
5165 MigrationRewriteInterruption::MarkerPersisted,
5166 MigrationRewriteInterruption::JournalPublished,
5167 MigrationRewriteInterruption::PhysicalApplied,
5168 ] {
5169 interrupt_next_migration_rewrite_at(interruption);
5170 migrate_schema(&db, &proposal, command())
5171 .expect_err("injected interruption should retain the rewrite marker");
5172
5173 forget_recovered_domain_for_tests(&db)
5174 .expect("upgrade should reset recovery ownership");
5175 drive_startup_recovery_to_completion(&db);
5176 }
5177
5178 let rebuilding = schema_migration_status_for_target(
5179 &db,
5180 &proposal,
5181 &schema_application_target(&db).expect("recovered target should issue"),
5182 )
5183 .expect("recovered status should remain readable");
5184 assert_eq!(rebuilding.phase(), SchemaMigrationPhase::RebuildingIndexes);
5185 assert_eq!(rebuilding.rows_rewritten(), 3);
5186 assert_eq!(
5187 migrate_schema(&db, &proposal, command())
5188 .expect("derived generations should complete")
5189 .phase(),
5190 SchemaMigrationPhase::FinalValidation,
5191 );
5192 assert_eq!(
5193 migrate_schema(&db, &proposal, command())
5194 .expect("final validation should complete")
5195 .phase(),
5196 SchemaMigrationPhase::Publishing,
5197 );
5198 let applied = migrate_schema(&db, &proposal, command())
5199 .expect("candidate publication should complete atomically");
5200 assert_eq!(applied.phase(), SchemaMigrationPhase::Applied);
5201 assert_eq!(applied.rows_rewritten(), 3);
5202 assert_eq!(applied.indexes_rebuilt(), 1);
5203 assert_ne!(applied.accepted_head(), target.accepted_head());
5204 let terminal_target = schema_application_target(&db).expect("terminal target should issue");
5205 let terminal_proposal = validation_migration_proposal(
5206 ValidationMigrationShape::Clean,
5207 true,
5208 terminal_target.accepted_head().clone(),
5209 terminal_target.database_identity(),
5210 store_identity,
5211 );
5212 assert!(
5213 !defer_generated_schema_application_for_prepared_migration(&db, &terminal_proposal,)
5214 .expect("terminal record must not block generated startup"),
5215 );
5216
5217 let store = db
5218 .store_handle(MIGRATION_EXECUTION_STORE_PATH)
5219 .expect("migration execution store should resolve");
5220 let runtime = db
5221 .accepted_runtime_entity_for_path("MigratingItem")
5222 .expect("published candidate entity should resolve");
5223 let selection = store
5224 .with_schema(|schema| {
5225 schema.current_accepted_catalog_selection(
5226 runtime.entity_tag(),
5227 runtime.entity_path(),
5228 runtime.store_path(),
5229 )
5230 })
5231 .expect("candidate selection should remain readable")
5232 .expect("candidate selection should exist");
5233 let contract = crate::db::data::AcceptedStructuralRowAuthority::from_catalog_selection(
5234 runtime.entity_path(),
5235 &selection,
5236 )
5237 .expect("candidate row authority should compile")
5238 .into_row_contract();
5239 let mut values = Vec::new();
5240 store
5241 .with_data(|data| {
5242 data.visit_entries(|key, row| {
5243 let decoded = DecodedDataStoreKey::try_from_raw(key)
5244 .expect("rewritten key should decode");
5245 let reader =
5246 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
5247 row, &contract,
5248 )
5249 .expect("rewritten row should use the candidate layout");
5250 reader
5251 .validate_primary_key(&decoded)
5252 .expect("rewritten row and key should remain bound");
5253 values.push(
5254 reader
5255 .required_value_by_contract(1)
5256 .expect("candidate value slot should decode"),
5257 );
5258 Ok::<StoreVisit, InternalError>(StoreVisit::Continue)
5259 })
5260 })
5261 .expect("rewritten row scan should complete");
5262 assert_eq!(
5263 values,
5264 vec![
5265 crate::value::Value::Nat64(7),
5266 crate::value::Value::Nat64(8),
5267 crate::value::Value::Nat64(9),
5268 ],
5269 );
5270 assert_eq!(store.with_index(IndexStore::len), 3);
5271 let accepted = store
5272 .with_schema(SchemaStore::current_accepted_schema_bundle)
5273 .expect("published candidate bundle should remain readable")
5274 .expect("published candidate bundle should exist");
5275 let entity_source = EntitySourceKey::try_new("MigratingItem")
5276 .expect("migration entity source should admit");
5277 let entity_tag = accepted
5278 .source_bindings_for_tests()
5279 .entity(&entity_source)
5280 .expect("candidate entity source should remain bound");
5281 let old_value =
5282 FieldSourceKey::try_new("old_value").expect("predecessor source should admit");
5283 let current_value =
5284 FieldSourceKey::try_new("value").expect("candidate source should admit");
5285 assert_eq!(
5286 accepted
5287 .source_bindings_for_tests()
5288 .field(entity_tag, &old_value),
5289 None,
5290 );
5291 assert!(
5292 accepted
5293 .source_bindings_for_tests()
5294 .field(entity_tag, ¤t_value)
5295 .is_some(),
5296 );
5297 ensure_schema_migration_ready_for_ordinary_operations()
5298 .expect("terminal publication must clear the database-wide gate");
5299 }
5300
5301 #[cfg(feature = "migration")]
5302 #[test]
5303 #[expect(
5304 clippy::too_many_lines,
5305 reason = "all four finding families share one ordered historical scan fixture"
5306 )]
5307 fn physical_migration_validation_reports_every_typed_finding_family_without_writes() {
5308 use std::convert::Infallible;
5309
5310 use super::migrate_schema;
5311 use crate::db::{
5312 data::StoreVisit,
5313 schema::{SchemaMigrationCommand, SchemaMigrationFindingKind, SchemaMigrationPhase},
5314 };
5315
5316 let db = Db::<MigrationFindingCanister>::new(
5317 &MIGRATION_FINDING_REGISTRY,
5318 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5319 );
5320 drive_startup_recovery_to_completion(&db);
5321 let initial_target = schema_application_target(&db).expect("initial target should issue");
5322 let store_identity = initial_target
5323 .stores()
5324 .first()
5325 .expect("migration finding store should exist")
5326 .identity();
5327 let initial = validation_migration_proposal(
5328 ValidationMigrationShape::AllFindingFamilies,
5329 false,
5330 initial_target.accepted_head().clone(),
5331 initial_target.database_identity(),
5332 store_identity,
5333 );
5334 apply_schema(&db, &initial).expect("initial finding schema should publish");
5335
5336 let session = DbSession::<MigrationFindingCanister>::new(
5337 &MIGRATION_FINDING_REGISTRY,
5338 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5339 );
5340 session
5341 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5342 entity: "MigrationTarget".to_string(),
5343 patch: DynamicStructuralPatch::new(vec![(
5344 "id".to_string(),
5345 DynamicWriteCell::Value(InputValue::nat64(7)),
5346 )]),
5347 })
5348 .expect("relation target should insert");
5349 for (id, value) in [(1, 9), (2, 8), (3, 7), (4, 7), (5, 300)] {
5350 session
5351 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5352 entity: "MigratingItem".to_string(),
5353 patch: DynamicStructuralPatch::new(vec![
5354 (
5355 "id".to_string(),
5356 DynamicWriteCell::Value(InputValue::nat64(id)),
5357 ),
5358 (
5359 "old_value".to_string(),
5360 DynamicWriteCell::Value(InputValue::int64(value)),
5361 ),
5362 ]),
5363 })
5364 .expect("predecessor finding row should insert");
5365 }
5366 let store = db
5367 .store_handle(MIGRATION_FINDING_STORE_PATH)
5368 .expect("migration finding store should resolve");
5369 let row_bytes = || {
5370 store.with_data(|data| {
5371 let mut rows = Vec::new();
5372 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
5373 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
5374 Ok(StoreVisit::Continue)
5375 });
5376 result.expect("infallible row visit should complete");
5377 rows
5378 })
5379 };
5380 let before_rows = row_bytes();
5381
5382 let target = schema_application_target(&db).expect("migration target should issue");
5383 let proposal = validation_migration_proposal(
5384 ValidationMigrationShape::AllFindingFamilies,
5385 true,
5386 target.accepted_head().clone(),
5387 target.database_identity(),
5388 store_identity,
5389 );
5390 let plan = proposal
5391 .migration()
5392 .expect("migration plan should exist")
5393 .digest();
5394 let command = || SchemaMigrationCommand::Advance {
5395 expected_database: target.database_identity(),
5396 expected_head: target.accepted_head().clone(),
5397 expected_plan: plan,
5398 acknowledged_finding_page: None,
5399 };
5400 assert_eq!(
5401 migrate_schema(&db, &proposal, command())
5402 .expect("finding migration should prepare")
5403 .phase(),
5404 SchemaMigrationPhase::Prepared,
5405 );
5406 assert_eq!(
5407 migrate_schema(&db, &proposal, command())
5408 .expect("finding migration should enter validation")
5409 .phase(),
5410 SchemaMigrationPhase::Validating,
5411 );
5412 let rejected =
5413 migrate_schema(&db, &proposal, command()).expect("validation should report findings");
5414 assert_eq!(rejected.phase(), SchemaMigrationPhase::Rejected);
5415 assert_eq!(rejected.rows_validated(), 5);
5416 assert_eq!(
5417 rejected
5418 .findings()
5419 .iter()
5420 .map(crate::db::schema::SchemaMigrationFinding::kind)
5421 .collect::<Vec<_>>(),
5422 vec![
5423 SchemaMigrationFindingKind::Constraint,
5424 SchemaMigrationFindingKind::Relation,
5425 SchemaMigrationFindingKind::UniqueIndex,
5426 SchemaMigrationFindingKind::Transform,
5427 ],
5428 );
5429 assert_eq!(
5430 row_bytes(),
5431 before_rows,
5432 "rejected validation must not rewrite accepted rows"
5433 );
5434 assert_eq!(
5435 store.with_index(IndexStore::len),
5436 0,
5437 "a rejected page must not publish any staged generation"
5438 );
5439 }
5440
5441 #[cfg(feature = "migration")]
5442 #[test]
5443 fn exact_migration_retry_binds_the_terminal_head_not_the_predecessor_head() {
5444 use super::exact_migration_replay_target;
5445
5446 let db = Db::<EvolutionCanister>::new(
5447 &EVOLUTION_REGISTRY,
5448 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5449 );
5450 drive_startup_recovery_to_completion(&db);
5451 let initial_target = schema_application_target(&db).expect("initial target should issue");
5452 let (proposal, _, _) = generated_check_proposal(
5453 initial_target.accepted_head().clone(),
5454 "migration-retry-initial",
5455 false,
5456 initial_target.database_identity(),
5457 initial_target
5458 .stores()
5459 .first()
5460 .expect("test store should exist")
5461 .identity(),
5462 );
5463 apply_schema(&db, &proposal).expect("initial schema should publish");
5464 let current_target = schema_application_target(&db).expect("current target should issue");
5465 assert_ne!(
5466 current_target.accepted_head(),
5467 initial_target.accepted_head(),
5468 );
5469
5470 let receipt = SchemaChangeReceipt::new(
5471 current_target.database_identity(),
5472 SchemaSubmissionKey::try_new("migration/retry")
5473 .expect("migration submission should admit"),
5474 SchemaProposalDigest::from_bytes([0x77; 32]),
5475 initial_target.accepted_head().clone(),
5476 SchemaChangeOutcome::Applied {
5477 accepted_head: current_target.accepted_head().clone(),
5478 },
5479 )
5480 .expect("terminal migration receipt should admit");
5481 let record = SchemaApplicationRecord::new(receipt, Vec::new())
5482 .expect("terminal migration record should admit");
5483
5484 assert_eq!(
5485 exact_migration_replay_target(&db, current_target.database_identity(), &record,)
5486 .expect("exact retry should resolve the terminal target"),
5487 current_target,
5488 );
5489 }
5490}