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