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