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 use super::{
3116 AcceptedSchemaPublication, AcceptedStoreHead, DirectGeneratedRowLocalProof,
3117 PendingGeneratedRowLocalConstraint, abort_schema_application,
3118 aborted_generated_row_local_candidate, accepted_head_after_candidates,
3119 application_authorities, apply_schema, continue_schema_application, derive_accepted_head,
3120 derive_schema_change_job_id, final_candidates_for_pending_row_local_constraint,
3121 generated_database_identity, include_identity_state_count, lower_existing_schema_proposal,
3122 lower_initial_schema_proposal, publish_accepted_schema_candidates_with_application_record,
3123 require_exact_empty_entity_count, schema_application_target,
3124 };
3125 use crate::{
3126 db::{
3127 DatabaseStartupState, Db, GeneratedStartupDriverStep,
3128 commit::{
3129 RecoveryProgress, continue_recovery, database_incarnation_id,
3130 forget_recovered_domain_for_tests,
3131 },
3132 data::DataStore,
3133 drive_generated_startup_recovery_page,
3134 index::IndexStore,
3135 journal::JournalTailStore,
3136 observe_generated_startup_state,
3137 registry::{
3138 StoreAllocationIdentities, StoreAllocationIdentity, StoreHandle, StoreRegistry,
3139 StoreRuntimeStorageCapabilities,
3140 },
3141 schema::{
3142 AcceptedConstraintKind, AcceptedRuleOperation, AcceptedSchemaRevisionBundle,
3143 CandidateSchemaRevision, ConstraintOrigin, ConstraintValidationJob,
3144 ExistingProposalStore, ProposalStoreTarget, SchemaApplicationRecord,
3145 SchemaApplicationRecordOp, SchemaChangeActivation, SchemaChangeJob,
3146 SchemaChangeOutcome, SchemaChangeProgressStatus, SchemaStore,
3147 cardinality_build::{
3148 CardinalityBuildAuthority, CardinalityGenerationPageOutcome,
3149 drive_cardinality_generation_page,
3150 },
3151 },
3152 },
3153 error::{ErrorClass, ErrorOrigin},
3154 testing::test_memory,
3155 traits::{CanisterKind, Path},
3156 };
3157 use ic_stable_structures::{DefaultMemoryImpl, memory_manager::VirtualMemory};
3158 use icydb_schema::{
3159 ConstraintFragment, ConstraintSourceKey, DeclaredEntityVersion, EntityFragment,
3160 EntitySourceKey, EntityStoreAssignment, ExpectedAcceptedHead, ExpectedSchemaFingerprint,
3161 FieldFragment, FieldInsertPolicy, FieldSourceKey, FieldType, NamedTypeFragment,
3162 RuleSourceKey, ScalarLiteral, ScalarType, SchemaCapability, SchemaFragment, SchemaName,
3163 SchemaProposal, SchemaSubmissionKey, SourceCheckExpr, SourceCheckInstruction,
3164 SourceRuleOperation, TargetDatabaseIdentity, TargetStoreIdentity, TargetedRuleFragment,
3165 TypeSourceKey,
3166 };
3167 use std::cell::RefCell;
3168
3169 fn drive_startup_recovery_to_completion<C: CanisterKind>(db: &Db<C>) {
3170 for _ in 0..1_024 {
3171 match continue_recovery(db).expect("test startup recovery page should succeed") {
3172 RecoveryProgress::Complete => return,
3173 RecoveryProgress::Pending => {}
3174 }
3175 }
3176 panic!("test startup recovery should complete within 1,024 bounded pages");
3177 }
3178
3179 fn drive_cardinality_to_ready(store: StoreHandle) {
3180 let journal = store
3181 .journal_tail_store()
3182 .expect("cardinality fixture store should be journaled");
3183 for _ in 0..8 {
3184 let outcome = store
3185 .with_data(|data| {
3186 store.with_index(|index| {
3187 store.with_schema_mut(|schema| {
3188 drive_cardinality_generation_page(data, index, schema, |schema| {
3189 CardinalityBuildAuthority::derive(
3190 schema,
3191 database_incarnation_id()?,
3192 store.allocation_identities(),
3193 journal.with_borrow(JournalTailStore::fold_watermark)?,
3194 )
3195 })
3196 })
3197 })
3198 })
3199 .expect("bounded cardinality generation should advance");
3200 if outcome == CardinalityGenerationPageOutcome::Quiescent {
3201 return;
3202 }
3203 }
3204 panic!("cardinality generation should become Ready within eight bounded pages");
3205 }
3206
3207 #[cfg(feature = "migration")]
3208 use crate::db::schema::SchemaChangeReceipt;
3209 use crate::{
3210 db::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell},
3211 value::InputValue,
3212 };
3213 #[cfg(feature = "migration")]
3214 use icydb_schema::{
3215 EntityMigration, IndexFragment, IndexKeyFragment, RelationDeleteAction, RelationFragment,
3216 SchemaMigrationPlan, SchemaMigrationTransform, SchemaProposalDigest,
3217 };
3218
3219 fn version_one() -> DeclaredEntityVersion {
3220 DeclaredEntityVersion::try_new(1).expect("fixture version should admit")
3221 }
3222
3223 const ABORT_STORE_PATH: &str = "schema_application_tests::AbortStore";
3224 const EVOLUTION_STORE_PATH: &str = "schema_application_tests::EvolutionStore";
3225 #[cfg(feature = "migration")]
3226 const MIGRATION_STORE_PATH: &str = "schema_application_tests::MigrationStore";
3227 #[cfg(feature = "migration")]
3228 const MIGRATION_EXECUTION_STORE_PATH: &str =
3229 "schema_application_tests::MigrationExecutionStore";
3230 #[cfg(feature = "migration")]
3231 const MIGRATION_FINDING_STORE_PATH: &str = "schema_application_tests::MigrationFindingStore";
3232
3233 #[test]
3234 fn database_identity_state_capacity_combines_store_inventories_exactly() {
3235 let below = include_identity_state_count(0, 65_535)
3236 .expect("the first store inventory should remain below the database cap");
3237 let exact = include_identity_state_count(below, 1)
3238 .expect("the combined database boundary should admit");
3239 assert_eq!(exact, 65_536);
3240
3241 let error = include_identity_state_count(exact, 1)
3242 .expect_err("the next owner in another store must reject");
3243 assert_eq!(error.class(), ErrorClass::Unsupported);
3244 assert_eq!(error.origin(), ErrorOrigin::Identity);
3245 }
3246
3247 #[test]
3248 fn generated_database_identity_cache_is_bound_to_the_incarnation() {
3249 let first_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x41);
3250 let second_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x42);
3251 let first = generated_database_identity(&ABORT_REGISTRY, first_incarnation);
3252 assert_eq!(
3253 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3254 first,
3255 );
3256 let second = generated_database_identity(&ABORT_REGISTRY, second_incarnation);
3257 assert_ne!(second, first);
3258 assert_eq!(
3259 generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3260 first,
3261 );
3262 }
3263
3264 #[test]
3265 fn exact_empty_entity_proof_distinguishes_corruption_from_non_empty_input() {
3266 let corrupt = require_exact_empty_entity_count(None)
3267 .expect_err("uninspectable cardinality must fail closed");
3268 assert_eq!(corrupt.class(), ErrorClass::Corruption);
3269
3270 let non_empty = require_exact_empty_entity_count(Some(1))
3271 .expect_err("non-empty cardinality must reject removal");
3272 assert_eq!(non_empty.class(), ErrorClass::Unsupported);
3273 assert!(require_exact_empty_entity_count(Some(0)).is_ok());
3274 }
3275
3276 thread_local! {
3277 static ABORT_DATA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(180);
3278 static ABORT_INDEX_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(181);
3279 static ABORT_SCHEMA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(182);
3280 static ABORT_JOURNAL_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(183);
3281 static ABORT_DATA: RefCell<DataStore> =
3282 ABORT_DATA_MEMORY.with(|memory| {
3283 RefCell::new(DataStore::init_journaled(memory.clone()))
3284 });
3285 static ABORT_INDEX: RefCell<IndexStore> =
3286 ABORT_INDEX_MEMORY.with(|memory| {
3287 RefCell::new(IndexStore::init_journaled(memory.clone()))
3288 });
3289 static ABORT_SCHEMA: RefCell<SchemaStore> =
3290 ABORT_SCHEMA_MEMORY.with(|memory| {
3291 RefCell::new(SchemaStore::init_journaled(memory.clone()))
3292 });
3293 static ABORT_JOURNAL: RefCell<JournalTailStore> =
3294 ABORT_JOURNAL_MEMORY.with(|memory| {
3295 RefCell::new(JournalTailStore::init(memory.clone()))
3296 });
3297 static ABORT_REGISTRY: StoreRegistry = {
3298 let mut registry = StoreRegistry::new();
3299 registry.register_journaled_store(
3300 ABORT_STORE_PATH,
3301 &ABORT_DATA,
3302 &ABORT_INDEX,
3303 &ABORT_SCHEMA,
3304 &ABORT_JOURNAL,
3305 StoreAllocationIdentities::new_journaled(
3306 StoreAllocationIdentity::new(180, "icydb.test.application_abort.data.v1"),
3307 StoreAllocationIdentity::new(181, "icydb.test.application_abort.index.v1"),
3308 StoreAllocationIdentity::new(182, "icydb.test.application_abort.schema.v1"),
3309 StoreAllocationIdentity::new(183, "icydb.test.application_abort.journal.v1"),
3310 ),
3311 StoreRuntimeStorageCapabilities::journaled(),
3312 ).expect("abort journaled store should register");
3313 registry
3314 };
3315 }
3316
3317 #[cfg(feature = "migration")]
3318 thread_local! {
3319 static MIGRATION_EXECUTION_DATA: RefCell<DataStore> =
3320 RefCell::new(DataStore::init_journaled(test_memory(210)));
3321 static MIGRATION_EXECUTION_INDEX: RefCell<IndexStore> =
3322 RefCell::new(IndexStore::init_journaled(test_memory(211)));
3323 static MIGRATION_EXECUTION_SCHEMA: RefCell<SchemaStore> =
3324 RefCell::new(SchemaStore::init_journaled(test_memory(212)));
3325 static MIGRATION_EXECUTION_JOURNAL: RefCell<JournalTailStore> =
3326 RefCell::new(JournalTailStore::init(test_memory(213)));
3327 static MIGRATION_EXECUTION_REGISTRY: StoreRegistry = {
3328 let mut registry = StoreRegistry::new();
3329 registry.register_journaled_store(
3330 MIGRATION_EXECUTION_STORE_PATH,
3331 &MIGRATION_EXECUTION_DATA,
3332 &MIGRATION_EXECUTION_INDEX,
3333 &MIGRATION_EXECUTION_SCHEMA,
3334 &MIGRATION_EXECUTION_JOURNAL,
3335 StoreAllocationIdentities::new_journaled(
3336 StoreAllocationIdentity::new(210, "icydb.test.migration_execution.data.v1"),
3337 StoreAllocationIdentity::new(211, "icydb.test.migration_execution.index.v1"),
3338 StoreAllocationIdentity::new(212, "icydb.test.migration_execution.schema.v1"),
3339 StoreAllocationIdentity::new(213, "icydb.test.migration_execution.journal.v1"),
3340 ),
3341 StoreRuntimeStorageCapabilities::journaled(),
3342 ).expect("migration execution store should register");
3343 registry
3344 };
3345 }
3346
3347 #[cfg(feature = "migration")]
3348 thread_local! {
3349 static MIGRATION_DATA: RefCell<DataStore> =
3350 RefCell::new(DataStore::init_journaled(test_memory(200)));
3351 static MIGRATION_INDEX: RefCell<IndexStore> =
3352 RefCell::new(IndexStore::init_journaled(test_memory(201)));
3353 static MIGRATION_SCHEMA: RefCell<SchemaStore> =
3354 RefCell::new(SchemaStore::init_journaled(test_memory(202)));
3355 static MIGRATION_JOURNAL: RefCell<JournalTailStore> =
3356 RefCell::new(JournalTailStore::init(test_memory(203)));
3357 static MIGRATION_REGISTRY: StoreRegistry = {
3358 let mut registry = StoreRegistry::new();
3359 registry.register_journaled_store(
3360 MIGRATION_STORE_PATH,
3361 &MIGRATION_DATA,
3362 &MIGRATION_INDEX,
3363 &MIGRATION_SCHEMA,
3364 &MIGRATION_JOURNAL,
3365 StoreAllocationIdentities::new_journaled(
3366 StoreAllocationIdentity::new(200, "icydb.test.migration_validation.data.v1"),
3367 StoreAllocationIdentity::new(201, "icydb.test.migration_validation.index.v1"),
3368 StoreAllocationIdentity::new(202, "icydb.test.migration_validation.schema.v1"),
3369 StoreAllocationIdentity::new(203, "icydb.test.migration_validation.journal.v1"),
3370 ),
3371 StoreRuntimeStorageCapabilities::journaled(),
3372 ).expect("migration validation store should register");
3373 registry
3374 };
3375 }
3376
3377 #[cfg(feature = "migration")]
3378 thread_local! {
3379 static MIGRATION_FINDING_DATA: RefCell<DataStore> =
3380 RefCell::new(DataStore::init_journaled(test_memory(206)));
3381 static MIGRATION_FINDING_INDEX: RefCell<IndexStore> =
3382 RefCell::new(IndexStore::init_journaled(test_memory(207)));
3383 static MIGRATION_FINDING_SCHEMA: RefCell<SchemaStore> =
3384 RefCell::new(SchemaStore::init_journaled(test_memory(208)));
3385 static MIGRATION_FINDING_JOURNAL: RefCell<JournalTailStore> =
3386 RefCell::new(JournalTailStore::init(test_memory(209)));
3387 static MIGRATION_FINDING_REGISTRY: StoreRegistry = {
3388 let mut registry = StoreRegistry::new();
3389 registry.register_journaled_store(
3390 MIGRATION_FINDING_STORE_PATH,
3391 &MIGRATION_FINDING_DATA,
3392 &MIGRATION_FINDING_INDEX,
3393 &MIGRATION_FINDING_SCHEMA,
3394 &MIGRATION_FINDING_JOURNAL,
3395 StoreAllocationIdentities::new_journaled(
3396 StoreAllocationIdentity::new(206, "icydb.test.migration_finding.data.v1"),
3397 StoreAllocationIdentity::new(207, "icydb.test.migration_finding.index.v1"),
3398 StoreAllocationIdentity::new(208, "icydb.test.migration_finding.schema.v1"),
3399 StoreAllocationIdentity::new(209, "icydb.test.migration_finding.journal.v1"),
3400 ),
3401 StoreRuntimeStorageCapabilities::journaled(),
3402 ).expect("migration finding store should register");
3403 registry
3404 };
3405 }
3406
3407 thread_local! {
3408 static EVOLUTION_DATA: RefCell<DataStore> =
3409 RefCell::new(DataStore::init_journaled(test_memory(192)));
3410 static EVOLUTION_INDEX: RefCell<IndexStore> =
3411 RefCell::new(IndexStore::init_journaled(test_memory(193)));
3412 static EVOLUTION_SCHEMA: RefCell<SchemaStore> =
3413 RefCell::new(SchemaStore::init_journaled(test_memory(194)));
3414 static EVOLUTION_JOURNAL: RefCell<JournalTailStore> =
3415 RefCell::new(JournalTailStore::init(test_memory(195)));
3416 static EVOLUTION_REGISTRY: StoreRegistry = {
3417 let mut registry = StoreRegistry::new();
3418 registry.register_journaled_store(
3419 EVOLUTION_STORE_PATH,
3420 &EVOLUTION_DATA,
3421 &EVOLUTION_INDEX,
3422 &EVOLUTION_SCHEMA,
3423 &EVOLUTION_JOURNAL,
3424 StoreAllocationIdentities::new_journaled(
3425 StoreAllocationIdentity::new(192, "icydb.test.rule_evolution.data.v1"),
3426 StoreAllocationIdentity::new(193, "icydb.test.rule_evolution.index.v1"),
3427 StoreAllocationIdentity::new(194, "icydb.test.rule_evolution.schema.v1"),
3428 StoreAllocationIdentity::new(195, "icydb.test.rule_evolution.journal.v1"),
3429 ),
3430 StoreRuntimeStorageCapabilities::journaled(),
3431 ).expect("rule-evolution journaled store should register");
3432 registry
3433 };
3434 }
3435
3436 struct AbortCanister;
3437
3438 impl Path for AbortCanister {
3439 const PATH: &'static str = "schema_application_tests::AbortCanister";
3440 }
3441
3442 impl CanisterKind for AbortCanister {
3443 const COMMIT_MEMORY_ID: u8 = 184;
3444 const COMMIT_STABLE_KEY: &'static str = "icydb.test.application_abort.commit.v1";
3445 const STARTUP_MEMORY_ID: u8 = 186;
3446 const STARTUP_STABLE_KEY: &'static str = "icydb.test.application_abort.startup.control.v1";
3447 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 185;
3448 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3449 "icydb.test.application_abort.integrity.v1";
3450 }
3451
3452 struct EvolutionCanister;
3453
3454 impl Path for EvolutionCanister {
3455 const PATH: &'static str = "schema_application_tests::EvolutionCanister";
3456 }
3457
3458 impl CanisterKind for EvolutionCanister {
3459 const COMMIT_MEMORY_ID: u8 = 196;
3460 const COMMIT_STABLE_KEY: &'static str = "icydb.test.rule_evolution.commit.v1";
3461 const STARTUP_MEMORY_ID: u8 = 198;
3462 const STARTUP_STABLE_KEY: &'static str = "icydb.test.rule_evolution.startup.control.v1";
3463 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 197;
3464 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3465 "icydb.test.rule_evolution.integrity.v1";
3466 }
3467
3468 #[cfg(feature = "migration")]
3469 struct MigrationCanister;
3470
3471 #[cfg(feature = "migration")]
3472 impl Path for MigrationCanister {
3473 const PATH: &'static str = "schema_application_tests::MigrationCanister";
3474 }
3475
3476 #[cfg(feature = "migration")]
3477 impl CanisterKind for MigrationCanister {
3478 const COMMIT_MEMORY_ID: u8 = 204;
3479 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_validation.commit.v1";
3480 const STARTUP_MEMORY_ID: u8 = 206;
3481 const STARTUP_STABLE_KEY: &'static str =
3482 "icydb.test.migration_validation.startup.control.v1";
3483 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 205;
3484 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3485 "icydb.test.migration_validation.integrity.v1";
3486 }
3487
3488 #[cfg(feature = "migration")]
3489 struct MigrationExecutionCanister;
3490
3491 #[cfg(feature = "migration")]
3492 impl Path for MigrationExecutionCanister {
3493 const PATH: &'static str = "schema_application_tests::MigrationExecutionCanister";
3494 }
3495
3496 #[cfg(feature = "migration")]
3497 impl CanisterKind for MigrationExecutionCanister {
3498 const COMMIT_MEMORY_ID: u8 = 214;
3499 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_execution.commit.v1";
3500 const STARTUP_MEMORY_ID: u8 = 216;
3501 const STARTUP_STABLE_KEY: &'static str =
3502 "icydb.test.migration_execution.startup.control.v1";
3503 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 215;
3504 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3505 "icydb.test.migration_execution.integrity.v1";
3506 }
3507
3508 #[cfg(feature = "migration")]
3509 struct MigrationFindingCanister;
3510
3511 #[cfg(feature = "migration")]
3512 impl Path for MigrationFindingCanister {
3513 const PATH: &'static str = "schema_application_tests::MigrationFindingCanister";
3514 }
3515
3516 #[cfg(feature = "migration")]
3517 impl CanisterKind for MigrationFindingCanister {
3518 const COMMIT_MEMORY_ID: u8 = 210;
3519 const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_finding.commit.v1";
3520 const STARTUP_MEMORY_ID: u8 = 212;
3521 const STARTUP_STABLE_KEY: &'static str = "icydb.test.migration_finding.startup.control.v1";
3522 const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 211;
3523 const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3524 "icydb.test.migration_finding.integrity.v1";
3525 }
3526
3527 fn name(value: &str) -> SchemaName {
3528 SchemaName::try_new(value).expect("test schema name should admit")
3529 }
3530
3531 fn generated_check_proposal(
3532 expected_head: ExpectedAcceptedHead,
3533 submission_key: &str,
3534 include_check: bool,
3535 database: TargetDatabaseIdentity,
3536 store: TargetStoreIdentity,
3537 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3538 let entity_source = EntitySourceKey::try_new("Item").expect("entity source should admit");
3539 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3540 let score_source = FieldSourceKey::try_new("score").expect("score source should admit");
3541 let check_source =
3542 ConstraintSourceKey::try_new("score_non_negative").expect("check source should admit");
3543 let check = SourceCheckExpr::try_new(vec![
3544 SourceCheckInstruction::Field(score_source),
3545 SourceCheckInstruction::Literal(ScalarLiteral::Int(0)),
3546 SourceCheckInstruction::GreaterThanOrEqual,
3547 ])
3548 .expect("check expression should admit");
3549 let constraints = include_check
3550 .then(|| ConstraintFragment::check(name("score_non_negative"), check))
3551 .into_iter()
3552 .collect();
3553 let entity = EntityFragment::try_new(
3554 name("Item"),
3555 version_one(),
3556 vec![
3557 FieldFragment::new(
3558 name("id"),
3559 FieldType::Scalar(ScalarType::Nat64),
3560 false,
3561 FieldInsertPolicy::Required,
3562 None,
3563 ),
3564 FieldFragment::new(
3565 name("score"),
3566 FieldType::Scalar(ScalarType::Int64),
3567 false,
3568 FieldInsertPolicy::Required,
3569 None,
3570 ),
3571 ],
3572 vec![id_source],
3573 Vec::new(),
3574 Vec::new(),
3575 constraints,
3576 )
3577 .expect("entity should admit");
3578 let proposal = SchemaProposal::try_compose(
3579 vec![SchemaCapability::ACCEPTED_CHECKS],
3580 database,
3581 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3582 expected_head,
3583 vec![
3584 SchemaFragment::try_new(vec![entity], Vec::new())
3585 .expect("schema fragment should admit"),
3586 ],
3587 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3588 Vec::new(),
3589 None,
3590 )
3591 .expect("schema proposal should compose");
3592 (proposal, entity_source, check_source)
3593 }
3594
3595 #[cfg(feature = "migration")]
3596 #[derive(Clone, Copy, Eq, PartialEq)]
3597 enum ValidationMigrationShape {
3598 Clean,
3599 AllFindingFamilies,
3600 }
3601
3602 #[cfg(feature = "migration")]
3603 #[expect(
3604 clippy::too_many_lines,
3605 reason = "the fixture keeps both predecessor and candidate source contracts adjacent"
3606 )]
3607 fn validation_migration_proposal(
3608 shape: ValidationMigrationShape,
3609 current: bool,
3610 expected_head: ExpectedAcceptedHead,
3611 database: TargetDatabaseIdentity,
3612 store: TargetStoreIdentity,
3613 ) -> SchemaProposal {
3614 let entity_source = EntitySourceKey::try_new("MigratingItem")
3615 .expect("migration entity source should admit");
3616 let old_value =
3617 FieldSourceKey::try_new("old_value").expect("predecessor field source should admit");
3618 let current_value =
3619 FieldSourceKey::try_new("value").expect("candidate field source should admit");
3620 let target_entity = EntitySourceKey::try_new("MigrationTarget")
3621 .expect("migration target source should admit");
3622 let target_id = FieldSourceKey::try_new("id").expect("target id source should admit");
3623 let constraint = SourceCheckExpr::try_new(vec![
3624 SourceCheckInstruction::Field(current_value.clone()),
3625 SourceCheckInstruction::Literal(ScalarLiteral::Nat(8)),
3626 SourceCheckInstruction::LessThanOrEqual,
3627 ])
3628 .expect("candidate check should admit");
3629 let findings = shape == ValidationMigrationShape::AllFindingFamilies;
3630 let entity = EntityFragment::try_new(
3631 name("MigratingItem"),
3632 DeclaredEntityVersion::try_new(if current { 2 } else { 1 })
3633 .expect("migration version should admit"),
3634 vec![
3635 FieldFragment::new(
3636 name("id"),
3637 FieldType::Scalar(ScalarType::Nat64),
3638 false,
3639 FieldInsertPolicy::Required,
3640 None,
3641 ),
3642 FieldFragment::new(
3643 name(if current { "value" } else { "old_value" }),
3644 FieldType::Scalar(if current {
3645 ScalarType::Nat8
3646 } else {
3647 ScalarType::Int64
3648 }),
3649 false,
3650 FieldInsertPolicy::Required,
3651 None,
3652 ),
3653 ],
3654 vec![FieldSourceKey::try_new("id").expect("id source should admit")],
3655 current
3656 .then(|| {
3657 IndexFragment::try_new(
3658 name("value_unique"),
3659 vec![IndexKeyFragment::Field(current_value.clone())],
3660 true,
3661 None,
3662 )
3663 .expect("candidate index should admit")
3664 })
3665 .into_iter()
3666 .collect(),
3667 (current && findings)
3668 .then(|| {
3669 RelationFragment::try_new(
3670 name("value_target"),
3671 vec![current_value.clone()],
3672 target_entity.clone(),
3673 vec![target_id.clone()],
3674 RelationDeleteAction::Restrict,
3675 )
3676 .expect("candidate relation should admit")
3677 })
3678 .into_iter()
3679 .collect(),
3680 (current && findings)
3681 .then(|| ConstraintFragment::check(name("value_at_most_eight"), constraint))
3682 .into_iter()
3683 .collect(),
3684 )
3685 .expect("migration entity should admit");
3686 let target = EntityFragment::try_new(
3687 name("MigrationTarget"),
3688 version_one(),
3689 vec![FieldFragment::new(
3690 name("id"),
3691 FieldType::Scalar(ScalarType::Nat8),
3692 false,
3693 FieldInsertPolicy::Required,
3694 None,
3695 )],
3696 vec![target_id],
3697 Vec::new(),
3698 Vec::new(),
3699 Vec::new(),
3700 )
3701 .expect("migration relation target should admit");
3702 let migration = current.then(|| {
3703 SchemaMigrationPlan::try_new(vec![
3704 EntityMigration::try_new(
3705 entity_source.clone(),
3706 DeclaredEntityVersion::try_new(1).expect("predecessor should admit"),
3707 None,
3708 Vec::new(),
3709 vec![SchemaMigrationTransform::CheckedCast {
3710 from: old_value.clone(),
3711 to: current_value,
3712 target: ScalarType::Nat8,
3713 }],
3714 )
3715 .expect("migration transition should admit"),
3716 ])
3717 .expect("migration plan should admit")
3718 });
3719 let mut capabilities = Vec::new();
3720 if current && findings {
3721 capabilities.extend([
3722 SchemaCapability::ACCEPTED_CHECKS,
3723 SchemaCapability::SECONDARY_INDEXES,
3724 SchemaCapability::RESTRICTIVE_RELATIONS,
3725 ]);
3726 }
3727 if migration.is_some() {
3728 capabilities.push(SchemaCapability::VERSIONED_MIGRATIONS);
3729 }
3730 let mut entities = vec![entity];
3731 let mut assignments = vec![EntityStoreAssignment::new(entity_source.clone(), store)];
3732 if findings {
3733 entities.push(target);
3734 assignments.push(EntityStoreAssignment::new(target_entity, store));
3735 }
3736 SchemaProposal::try_compose(
3737 capabilities,
3738 database,
3739 SchemaSubmissionKey::try_new(if current {
3740 "migration-validation-v2"
3741 } else {
3742 "migration-validation-v1"
3743 })
3744 .expect("submission should admit"),
3745 expected_head,
3746 vec![
3747 SchemaFragment::try_new(entities, Vec::new())
3748 .expect("migration fragment should admit"),
3749 ],
3750 assignments,
3751 current
3752 .then_some(icydb_schema::SchemaRemoval::Field {
3753 entity: entity_source,
3754 field: old_value,
3755 })
3756 .into_iter()
3757 .collect(),
3758 migration,
3759 )
3760 .expect("migration proposal should compose")
3761 }
3762
3763 fn targeted_rule_proposal(
3764 expected_head: ExpectedAcceptedHead,
3765 submission_key: &str,
3766 operation: SourceRuleOperation,
3767 database: TargetDatabaseIdentity,
3768 store: TargetStoreIdentity,
3769 ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3770 let entity_source =
3771 EntitySourceKey::try_new("Measured").expect("entity source should admit");
3772 let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3773 let value_source = FieldSourceKey::try_new("value").expect("value source should admit");
3774 let value_type = TypeSourceKey::try_new("Measure").expect("type source should admit");
3775 let rule_source = RuleSourceKey::try_new("limit").expect("rule source should admit");
3776 let constraint_source =
3777 ConstraintSourceKey::for_targeted_field_rule(&value_source, &value_type, &rule_source);
3778 let entity = EntityFragment::try_new(
3779 name("Measured"),
3780 version_one(),
3781 vec![
3782 FieldFragment::new(
3783 name("id"),
3784 FieldType::Scalar(ScalarType::Nat64),
3785 false,
3786 FieldInsertPolicy::Required,
3787 None,
3788 ),
3789 FieldFragment::new(
3790 name("value"),
3791 FieldType::Named(value_type.clone()),
3792 false,
3793 FieldInsertPolicy::Required,
3794 None,
3795 ),
3796 ],
3797 vec![id_source],
3798 Vec::new(),
3799 Vec::new(),
3800 vec![ConstraintFragment::targeted_rule(
3801 TargetedRuleFragment::new(value_source, value_type, name("limit"), operation),
3802 )],
3803 )
3804 .expect("targeted entity should admit");
3805 let proposal = SchemaProposal::try_compose(
3806 vec![SchemaCapability::ACCEPTED_CHECKS],
3807 database,
3808 SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3809 expected_head,
3810 vec![
3811 SchemaFragment::try_new(
3812 vec![entity],
3813 vec![NamedTypeFragment::newtype(
3814 name("Measure"),
3815 FieldType::Scalar(ScalarType::Nat8),
3816 )],
3817 )
3818 .expect("schema fragment should admit"),
3819 ],
3820 vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3821 Vec::new(),
3822 None,
3823 )
3824 .expect("schema proposal should compose");
3825 (proposal, entity_source, constraint_source)
3826 }
3827
3828 #[test]
3829 fn database_head_is_empty_only_when_every_store_root_is_absent() {
3830 assert_eq!(
3831 derive_accepted_head(&[("test::A", None), ("test::B", None)]),
3832 ExpectedAcceptedHead::Empty,
3833 );
3834 }
3835
3836 #[test]
3837 fn database_head_covers_store_path_revision_fingerprint_and_absence() {
3838 let first = derive_accepted_head(&[
3839 (
3840 "test::A",
3841 Some(AcceptedStoreHead {
3842 revision: 3,
3843 fingerprint: [0x11; 32],
3844 }),
3845 ),
3846 ("test::B", None),
3847 ]);
3848 let changed_fingerprint = derive_accepted_head(&[
3849 (
3850 "test::A",
3851 Some(AcceptedStoreHead {
3852 revision: 3,
3853 fingerprint: [0x12; 32],
3854 }),
3855 ),
3856 ("test::B", None),
3857 ]);
3858 let changed_absence = derive_accepted_head(&[
3859 (
3860 "test::A",
3861 Some(AcceptedStoreHead {
3862 revision: 3,
3863 fingerprint: [0x11; 32],
3864 }),
3865 ),
3866 (
3867 "test::B",
3868 Some(AcceptedStoreHead {
3869 revision: 1,
3870 fingerprint: [0x22; 32],
3871 }),
3872 ),
3873 ]);
3874
3875 assert_ne!(first, changed_fingerprint);
3876 assert_ne!(first, changed_absence);
3877 assert!(matches!(
3878 first,
3879 ExpectedAcceptedHead::Exact { revision: 3, .. }
3880 ));
3881 }
3882
3883 #[test]
3884 #[allow(
3885 clippy::too_many_lines,
3886 reason = "the end-to-end catalog assertion is clearer as one lifecycle test"
3887 )]
3888 fn generated_check_abort_retires_source_identity_and_allows_fresh_reproposal() {
3889 let database = TargetDatabaseIdentity::from_bytes([0x71; 32]);
3890 let store = TargetStoreIdentity::from_bytes([0x72; 32]);
3891 let (initial, entity_source, _) = generated_check_proposal(
3892 ExpectedAcceptedHead::Empty,
3893 "abort-initial",
3894 false,
3895 database,
3896 store,
3897 );
3898 let initial_candidate = lower_initial_schema_proposal(
3899 &initial,
3900 &[ProposalStoreTarget {
3901 path: "abort::Store",
3902 identity: store,
3903 }],
3904 )
3905 .expect("initial proposal should lower")
3906 .pop()
3907 .expect("initial proposal should produce one candidate");
3908 let (with_check, _, check_source) = generated_check_proposal(
3909 ExpectedAcceptedHead::Exact {
3910 revision: 1,
3911 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x73; 32]),
3912 },
3913 "abort-add-check",
3914 true,
3915 database,
3916 store,
3917 );
3918 let pending_candidate = lower_existing_schema_proposal(
3919 &with_check,
3920 &[ExistingProposalStore {
3921 path: "abort::Store",
3922 identity: store,
3923 bundle: initial_candidate.bundle(),
3924 }],
3925 )
3926 .expect("generated check should lower")
3927 .pop()
3928 .expect("generated check should produce one candidate");
3929 let entity_tag = pending_candidate
3930 .bundle()
3931 .source_bindings_for_tests()
3932 .entity(&entity_source)
3933 .expect("entity source should remain bound");
3934 let constraint_id = pending_candidate
3935 .bundle()
3936 .source_bindings_for_tests()
3937 .constraint(entity_tag, &check_source)
3938 .expect("generated check source should bind");
3939 let pending_snapshot = pending_candidate
3940 .bundle()
3941 .entity_snapshots()
3942 .get(&entity_tag)
3943 .expect("pending entity should exist");
3944 let activation = pending_snapshot
3945 .constraint_catalog()
3946 .activation(constraint_id)
3947 .expect("generated check should remain an activation");
3948 assert_eq!(activation.origin(), ConstraintOrigin::Generated);
3949
3950 let aborted = aborted_generated_row_local_candidate(
3951 pending_candidate.bundle(),
3952 entity_tag,
3953 constraint_id,
3954 )
3955 .expect("generated check abort should build one catalog-native candidate");
3956 let aborted_snapshot = aborted
3957 .bundle()
3958 .entity_snapshots()
3959 .get(&entity_tag)
3960 .expect("aborted entity should remain");
3961 assert!(
3962 aborted_snapshot
3963 .constraint_catalog()
3964 .activation(constraint_id)
3965 .is_none(),
3966 );
3967 assert_eq!(aborted_snapshot.row_layout(), pending_snapshot.row_layout());
3968 assert!(
3969 aborted
3970 .bundle()
3971 .source_bindings_for_tests()
3972 .constraint(entity_tag, &check_source)
3973 .is_none(),
3974 );
3975
3976 let reproposed = lower_existing_schema_proposal(
3977 &with_check,
3978 &[ExistingProposalStore {
3979 path: "abort::Store",
3980 identity: store,
3981 bundle: aborted.bundle(),
3982 }],
3983 )
3984 .expect("aborted generated check should be independently reproposable")
3985 .pop()
3986 .expect("reproposal should produce one candidate");
3987 let replacement_id = reproposed
3988 .bundle()
3989 .source_bindings_for_tests()
3990 .constraint(entity_tag, &check_source)
3991 .expect("reproposal should bind a fresh constraint identity");
3992 assert!(
3993 replacement_id > constraint_id,
3994 "aborted accepted IDs must remain retired",
3995 );
3996 }
3997
3998 #[test]
3999 fn targeted_rule_edit_abort_keeps_prior_accepted_semantics_and_source_identity() {
4000 let database = TargetDatabaseIdentity::from_bytes([0x81; 32]);
4001 let store = TargetStoreIdentity::from_bytes([0x82; 32]);
4002 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4003 ExpectedAcceptedHead::Empty,
4004 "targeted-abort-initial",
4005 SourceRuleOperation::NumericRangeInclusive {
4006 min: ScalarLiteral::Nat(0),
4007 max: ScalarLiteral::Nat(10),
4008 },
4009 database,
4010 store,
4011 );
4012 let initial_candidate = lower_initial_schema_proposal(
4013 &initial,
4014 &[ProposalStoreTarget {
4015 path: "abort::TargetedStore",
4016 identity: store,
4017 }],
4018 )
4019 .expect("initial targeted proposal should lower")
4020 .pop()
4021 .expect("initial targeted proposal should produce one candidate");
4022 let initial_bundle = initial_candidate.bundle();
4023 let entity_tag = initial_bundle
4024 .source_bindings_for_tests()
4025 .entity(&entity_source)
4026 .expect("entity source should bind");
4027 let constraint_id = initial_bundle
4028 .source_bindings_for_tests()
4029 .constraint(entity_tag, &constraint_source)
4030 .expect("targeted source should bind");
4031 let high_water = initial_bundle.entity_snapshots()[&entity_tag]
4032 .constraint_id_allocator()
4033 .high_water();
4034 let (edited, _, _) = targeted_rule_proposal(
4035 ExpectedAcceptedHead::Exact {
4036 revision: initial_bundle.revision().get(),
4037 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x83; 32]),
4038 },
4039 "targeted-abort-edit",
4040 SourceRuleOperation::NumericMaximumInclusive {
4041 value: ScalarLiteral::Nat(8),
4042 },
4043 database,
4044 store,
4045 );
4046 let staged = lower_existing_schema_proposal(
4047 &edited,
4048 &[ExistingProposalStore {
4049 path: "abort::TargetedStore",
4050 identity: store,
4051 bundle: initial_bundle,
4052 }],
4053 )
4054 .expect("targeted semantic edit should stage")
4055 .pop()
4056 .expect("targeted semantic edit should produce one candidate");
4057 let aborted =
4058 aborted_generated_row_local_candidate(staged.bundle(), entity_tag, constraint_id)
4059 .expect("targeted semantic edit should abort through catalog authority");
4060 let snapshot = &aborted.bundle().entity_snapshots()[&entity_tag];
4061
4062 assert!(
4063 snapshot
4064 .constraint_catalog()
4065 .activation(constraint_id)
4066 .is_none()
4067 );
4068 assert_eq!(snapshot.constraint_id_allocator().high_water(), high_water);
4069 assert_eq!(
4070 aborted
4071 .bundle()
4072 .source_bindings_for_tests()
4073 .constraint(entity_tag, &constraint_source),
4074 Some(constraint_id),
4075 );
4076 assert!(snapshot.constraints().iter().any(|constraint| {
4077 constraint.id() == constraint_id
4078 && matches!(
4079 constraint.kind(),
4080 AcceptedConstraintKind::TargetedRule { operation, .. }
4081 if matches!(
4082 operation.as_ref(),
4083 AcceptedRuleOperation::NumericRangeInclusive { .. }
4084 )
4085 )
4086 }));
4087 }
4088
4089 #[test]
4090 #[allow(
4091 clippy::too_many_lines,
4092 reason = "the staged publication, recovery, and promotion assertions form one lifecycle"
4093 )]
4094 fn generated_startup_driver_resumes_pending_activation_and_promotes_without_source_model() {
4095 let db = Db::<EvolutionCanister>::new(
4096 &EVOLUTION_REGISTRY,
4097 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4098 );
4099 drive_startup_recovery_to_completion(&db);
4100 let empty_target =
4101 schema_application_target(&db).expect("empty evolution target should issue");
4102 let store_identity = empty_target
4103 .stores()
4104 .first()
4105 .expect("evolution store should register")
4106 .identity();
4107 let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4108 empty_target.accepted_head().clone(),
4109 "targeted-recovery-initial",
4110 SourceRuleOperation::NumericRangeInclusive {
4111 min: ScalarLiteral::Nat(0),
4112 max: ScalarLiteral::Nat(10),
4113 },
4114 empty_target.database_identity(),
4115 store_identity,
4116 );
4117 assert!(matches!(
4118 apply_schema(&db, &initial)
4119 .expect("initial targeted proposal should publish")
4120 .outcome(),
4121 SchemaChangeOutcome::Applied { .. },
4122 ));
4123 drive_startup_recovery_to_completion(&db);
4124 drive_cardinality_to_ready(
4125 db.store_handle(EVOLUTION_STORE_PATH)
4126 .expect("evolution store should resolve"),
4127 );
4128
4129 let direct_target =
4130 schema_application_target(&db).expect("direct evolution target should issue");
4131 let (direct_edit, _, _) = targeted_rule_proposal(
4132 direct_target.accepted_head().clone(),
4133 "targeted-direct-edit",
4134 SourceRuleOperation::NumericMaximumInclusive {
4135 value: ScalarLiteral::Nat(8),
4136 },
4137 direct_target.database_identity(),
4138 store_identity,
4139 );
4140 assert!(matches!(
4141 apply_schema(&db, &direct_edit)
4142 .expect("empty-domain semantic edit should publish directly")
4143 .outcome(),
4144 SchemaChangeOutcome::Applied { .. },
4145 ));
4146 let store = db
4147 .store_handle(EVOLUTION_STORE_PATH)
4148 .expect("evolution store should resolve");
4149 let direct = store
4150 .with_schema(SchemaStore::current_accepted_schema_bundle)
4151 .expect("directly edited bundle should remain readable")
4152 .expect("directly edited bundle should exist");
4153 let entity_tag = direct
4154 .source_bindings_for_tests()
4155 .entity(&entity_source)
4156 .expect("entity source should remain bound");
4157 let constraint_id = direct
4158 .source_bindings_for_tests()
4159 .constraint(entity_tag, &constraint_source)
4160 .expect("direct edit should preserve constraint identity");
4161 assert!(
4162 direct.entity_snapshots()[&entity_tag]
4163 .constraint_catalog()
4164 .activation(constraint_id)
4165 .is_none()
4166 );
4167 assert!(
4168 direct.entity_snapshots()[&entity_tag]
4169 .constraints()
4170 .iter()
4171 .any(|constraint| {
4172 constraint.id() == constraint_id
4173 && matches!(
4174 constraint.kind(),
4175 AcceptedConstraintKind::TargetedRule { operation, .. }
4176 if matches!(
4177 operation.as_ref(),
4178 AcceptedRuleOperation::NumericMaximumInclusive { .. }
4179 )
4180 )
4181 })
4182 );
4183
4184 let target = schema_application_target(&db).expect("staged evolution target should issue");
4185 let (edited, _, _) = targeted_rule_proposal(
4186 target.accepted_head().clone(),
4187 "targeted-recovery-edit",
4188 SourceRuleOperation::MultipleOf {
4189 divisor: ScalarLiteral::Nat(2),
4190 },
4191 target.database_identity(),
4192 store_identity,
4193 );
4194 let current = store
4195 .with_schema(SchemaStore::current_accepted_schema_bundle)
4196 .expect("accepted evolution bundle should remain readable")
4197 .expect("directly edited evolution bundle should exist");
4198 let staged = lower_existing_schema_proposal(
4199 &edited,
4200 &[ExistingProposalStore {
4201 path: EVOLUTION_STORE_PATH,
4202 identity: store_identity,
4203 bundle: ¤t,
4204 }],
4205 )
4206 .expect("targeted edit should stage")
4207 .pop()
4208 .expect("targeted edit should produce one candidate");
4209 assert_eq!(
4210 staged
4211 .bundle()
4212 .source_bindings_for_tests()
4213 .constraint(entity_tag, &constraint_source),
4214 Some(constraint_id),
4215 );
4216 let proof = DirectGeneratedRowLocalProof {
4217 candidate_index: 0,
4218 store,
4219 store_path: EVOLUTION_STORE_PATH,
4220 entity_tag,
4221 entity_path: staged.bundle().entity_snapshots()[&entity_tag]
4222 .entity_path()
4223 .to_string(),
4224 constraint_id,
4225 historical_rows: 0,
4226 };
4227 let final_candidates = final_candidates_for_pending_row_local_constraint(
4228 std::slice::from_ref(&staged),
4229 &PendingGeneratedRowLocalConstraint { proof },
4230 )
4231 .expect("final semantic replacement should derive without source input");
4232 let authorities = application_authorities(&db);
4233 let candidate_head =
4234 accepted_head_after_candidates(authorities.as_slice(), &final_candidates)
4235 .expect("final candidate head should derive");
4236 let digest = edited.digest().expect("proposal digest should derive");
4237 let job_id = derive_schema_change_job_id(
4238 target.database_identity(),
4239 edited.submission_key(),
4240 digest,
4241 target.accepted_head(),
4242 )
4243 .expect("job identity should derive");
4244 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4245 target.database_identity(),
4246 edited.submission_key().clone(),
4247 digest,
4248 target.accepted_head().clone(),
4249 SchemaChangeOutcome::Pending {
4250 job: SchemaChangeJob::new(job_id),
4251 candidate_head,
4252 },
4253 )
4254 .expect("pending replacement receipt should admit");
4255 let record = SchemaApplicationRecord::new(
4256 receipt,
4257 vec![
4258 SchemaChangeActivation::new(
4259 store_identity,
4260 entity_tag.value(),
4261 constraint_id.get(),
4262 )
4263 .expect("replacement activation should admit"),
4264 ],
4265 )
4266 .expect("pending replacement record should admit");
4267 let operation =
4268 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4269 publish_accepted_schema_candidates_with_application_record(
4270 vec![AcceptedSchemaPublication::new(
4271 EVOLUTION_STORE_PATH,
4272 store,
4273 current.revision(),
4274 &staged,
4275 )],
4276 operation,
4277 )
4278 .expect("staged replacement and record should publish atomically");
4279
4280 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4281 drive_startup_recovery_to_completion(&db);
4282
4283 let recovered = store
4284 .with_schema(SchemaStore::current_accepted_schema_bundle)
4285 .expect("recovered staged bundle should decode")
4286 .expect("recovered staged bundle should exist");
4287 let recovered_snapshot = recovered.entity_snapshots()[&entity_tag].clone();
4288 let validating_catalog = recovered_snapshot
4289 .constraint_catalog()
4290 .clone()
4291 .with_validation_started(constraint_id)
4292 .expect("recovered replacement should enter validation");
4293 let mut validating_snapshots = recovered.entity_snapshots().clone();
4294 validating_snapshots.insert(
4295 entity_tag,
4296 recovered_snapshot.with_constraint_catalog(validating_catalog),
4297 );
4298 let validating_bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
4299 recovered
4300 .revision()
4301 .checked_next()
4302 .expect("validation revision should remain available"),
4303 recovered.store_path(),
4304 recovered.enum_catalog().clone(),
4305 recovered.composite_catalog().clone(),
4306 recovered.source_bindings_for_tests().clone(),
4307 validating_snapshots,
4308 )
4309 .expect("validating replacement bundle should close");
4310 let validating_candidate = CandidateSchemaRevision::new(validating_bundle)
4311 .expect("validating replacement candidate should encode");
4312 let validating_activation = validating_candidate.bundle().entity_snapshots()[&entity_tag]
4313 .constraint_catalog()
4314 .activation(constraint_id)
4315 .expect("validating replacement activation should remain present");
4316 let validation_job = ConstraintValidationJob::start(
4317 entity_tag,
4318 validating_candidate.bundle().entity_snapshots()[&entity_tag]
4319 .entity_path()
4320 .to_string(),
4321 validating_activation,
4322 None,
4323 )
4324 .expect("validating replacement job should derive from accepted state");
4325 store
4326 .with_schema(|schema| {
4327 schema.validate_live_activation_transition(validating_candidate.bundle())?;
4328 schema.validate_constraint_validation_job_closure_with_change(
4329 validating_candidate.bundle(),
4330 Some(&validation_job),
4331 None,
4332 )
4333 })
4334 .expect("validating replacement transition and job should close");
4335 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4336 let startup_session =
4337 crate::db::DbSession::<EvolutionCanister>::new(&EVOLUTION_REGISTRY, &startup_root);
4338
4339 assert_eq!(
4340 drive_generated_startup_recovery_page(
4341 &startup_session,
4342 &EVOLUTION_REGISTRY,
4343 edited.submission_key().as_str(),
4344 )
4345 .expect("generated startup should begin pending validation"),
4346 GeneratedStartupDriverStep::Recovering,
4347 );
4348 let mut terminal = false;
4349 for _ in 0..8 {
4350 match drive_generated_startup_recovery_page(
4351 &startup_session,
4352 &EVOLUTION_REGISTRY,
4353 edited.submission_key().as_str(),
4354 )
4355 .expect("generated startup should advance pending validation")
4356 {
4357 GeneratedStartupDriverStep::Recovering => {}
4358 GeneratedStartupDriverStep::Terminal => {
4359 terminal = true;
4360 break;
4361 }
4362 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4363 panic!("an exact pending receipt must resume instead of being resubmitted")
4364 }
4365 }
4366 }
4367 assert!(
4368 terminal,
4369 "empty historical domain should promote within bounded startup steps"
4370 );
4371 assert_eq!(
4372 observe_generated_startup_state::<EvolutionCanister>(
4373 &EVOLUTION_REGISTRY,
4374 edited.submission_key().as_str(),
4375 ),
4376 Ok(DatabaseStartupState::Ready),
4377 );
4378 let applied = super::exact_schema_application_receipt(
4379 &edited,
4380 edited
4381 .digest()
4382 .expect("proposal digest should remain stable"),
4383 )
4384 .expect("terminal generated receipt should remain readable")
4385 .expect("terminal generated receipt should remain present");
4386 assert!(matches!(
4387 applied.outcome(),
4388 SchemaChangeOutcome::Applied { .. }
4389 ));
4390 let promoted = store
4391 .with_schema(SchemaStore::current_accepted_schema_bundle)
4392 .expect("promoted bundle should remain readable")
4393 .expect("promoted bundle should exist");
4394 let snapshot = &promoted.entity_snapshots()[&entity_tag];
4395 assert!(
4396 snapshot
4397 .constraint_catalog()
4398 .activation(constraint_id)
4399 .is_none()
4400 );
4401 assert_eq!(
4402 promoted
4403 .source_bindings_for_tests()
4404 .constraint(entity_tag, &constraint_source),
4405 Some(constraint_id),
4406 );
4407 assert!(snapshot.constraints().iter().any(|constraint| {
4408 constraint.id() == constraint_id
4409 && matches!(
4410 constraint.kind(),
4411 AcceptedConstraintKind::TargetedRule { operation, .. }
4412 if matches!(
4413 operation.as_ref(),
4414 AcceptedRuleOperation::MultipleOf { .. }
4415 )
4416 )
4417 }));
4418 }
4419
4420 #[test]
4421 #[allow(
4422 clippy::too_many_lines,
4423 reason = "the durable pending job, startup failure, and retained finding assertions form one scenario"
4424 )]
4425 fn generated_startup_driver_persists_e223_for_a_retained_historical_finding() {
4426 let db = Db::<AbortCanister>::new(
4427 &ABORT_REGISTRY,
4428 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4429 );
4430 drive_startup_recovery_to_completion(&db);
4431 let empty_target =
4432 schema_application_target(&db).expect("empty application target should issue");
4433 let store_identity = empty_target
4434 .stores()
4435 .first()
4436 .expect("abort store should be registered")
4437 .identity();
4438 let (initial, _, _) = generated_check_proposal(
4439 empty_target.accepted_head().clone(),
4440 "startup-finding-initial",
4441 false,
4442 empty_target.database_identity(),
4443 store_identity,
4444 );
4445 apply_schema(&db, &initial).expect("initial generated schema should publish");
4446
4447 let root = crate::db::RequestExecutionRoot::__new_runtime_root();
4448 let session = DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &root);
4449 let rows = (1..=257)
4450 .map(|id| DynamicMutation::Insert {
4451 entity: "Item".to_string(),
4452 patch: DynamicStructuralPatch::new(vec![
4453 (
4454 "id".to_string(),
4455 DynamicWriteCell::Value(InputValue::nat64(id)),
4456 ),
4457 (
4458 "score".to_string(),
4459 DynamicWriteCell::Value(InputValue::int64(if id == 257 { -1 } else { 1 })),
4460 ),
4461 ]),
4462 })
4463 .collect();
4464 session
4465 .execute_trusted_dynamic_mutation_batch(rows)
4466 .expect("historical finding fixture rows should commit as one legal batch");
4467 drive_startup_recovery_to_completion(&db);
4468 drive_cardinality_to_ready(
4469 db.store_handle(ABORT_STORE_PATH)
4470 .expect("abort store should resolve"),
4471 );
4472
4473 let target =
4474 schema_application_target(&db).expect("existing application target should issue");
4475 let (with_check, _, _) = generated_check_proposal(
4476 target.accepted_head().clone(),
4477 "startup-finding-pending",
4478 true,
4479 target.database_identity(),
4480 store_identity,
4481 );
4482 let pending = apply_schema(&db, &with_check)
4483 .expect("the first clean page should admit durable continuation");
4484 let SchemaChangeOutcome::Pending { job, .. } = pending.outcome() else {
4485 panic!("a 257-row domain must exceed the 256-row direct proof page")
4486 };
4487
4488 let mut terminal = false;
4489 for _ in 0..8 {
4490 match drive_generated_startup_recovery_page(
4491 &session,
4492 &ABORT_REGISTRY,
4493 with_check.submission_key().as_str(),
4494 )
4495 .expect("generated startup should retain a typed finding failure")
4496 {
4497 GeneratedStartupDriverStep::Recovering => {}
4498 GeneratedStartupDriverStep::Terminal => {
4499 terminal = true;
4500 break;
4501 }
4502 GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4503 panic!("an exact pending receipt must not be resubmitted")
4504 }
4505 }
4506 }
4507 assert!(terminal, "the retained finding should become terminal");
4508 let failure = observe_generated_startup_state::<AbortCanister>(
4509 &ABORT_REGISTRY,
4510 with_check.submission_key().as_str(),
4511 )
4512 .expect_err("the retained finding must remain durably observable");
4513 assert_eq!(
4514 failure.kind(),
4515 crate::db::StartupFailureKind::SchemaReconciliation,
4516 );
4517 assert_eq!(
4518 failure.diagnostic().error_code(),
4519 icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_CONSTRAINT_VIOLATION,
4520 );
4521 assert_eq!(
4522 ABORT_DATA.with(|store| store.borrow().len()),
4523 257,
4524 "terminal startup publication must not change historical rows",
4525 );
4526 assert!(matches!(
4527 continue_schema_application(&db, job.id(), None)
4528 .expect("the retained finding page should replay exactly")
4529 .status(),
4530 SchemaChangeProgressStatus::Findings { findings, .. } if !findings.is_empty(),
4531 ));
4532 }
4533
4534 #[test]
4535 #[allow(
4536 clippy::too_many_lines,
4537 reason = "the journaled abort, replay, and recovery assertions form one scenario"
4538 )]
4539 fn pending_generated_check_abort_is_atomic_terminal_and_replayable() {
4540 let db = Db::<AbortCanister>::new(
4541 &ABORT_REGISTRY,
4542 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4543 );
4544 drive_startup_recovery_to_completion(&db);
4545 let empty_target =
4546 schema_application_target(&db).expect("empty application target should issue");
4547 let store_identity = empty_target
4548 .stores()
4549 .first()
4550 .expect("abort store should be registered")
4551 .identity();
4552 let (initial, entity_source, _) = generated_check_proposal(
4553 empty_target.accepted_head().clone(),
4554 "abort-runtime-initial",
4555 false,
4556 empty_target.database_identity(),
4557 store_identity,
4558 );
4559 assert!(matches!(
4560 apply_schema(&db, &initial)
4561 .expect("initial application should publish")
4562 .outcome(),
4563 SchemaChangeOutcome::Applied { .. },
4564 ));
4565
4566 let target =
4567 schema_application_target(&db).expect("existing application target should issue");
4568 let (with_check, _, check_source) = generated_check_proposal(
4569 target.accepted_head().clone(),
4570 "abort-runtime-pending",
4571 true,
4572 target.database_identity(),
4573 store_identity,
4574 );
4575 let store = db
4576 .store_handle(ABORT_STORE_PATH)
4577 .expect("abort store should resolve");
4578 let current = store
4579 .with_schema(SchemaStore::current_accepted_schema_bundle)
4580 .expect("accepted bundle should remain readable")
4581 .expect("initial accepted bundle should exist");
4582 let pending_candidate = lower_existing_schema_proposal(
4583 &with_check,
4584 &[ExistingProposalStore {
4585 path: ABORT_STORE_PATH,
4586 identity: store_identity,
4587 bundle: ¤t,
4588 }],
4589 )
4590 .expect("pending generated check should lower")
4591 .pop()
4592 .expect("pending generated check should produce one candidate");
4593 let entity_tag = pending_candidate
4594 .bundle()
4595 .source_bindings_for_tests()
4596 .entity(&entity_source)
4597 .expect("entity source should bind");
4598 let constraint_id = pending_candidate
4599 .bundle()
4600 .source_bindings_for_tests()
4601 .constraint(entity_tag, &check_source)
4602 .expect("generated check source should bind");
4603 let digest = with_check.digest().expect("proposal digest should derive");
4604 let job_id = derive_schema_change_job_id(
4605 target.database_identity(),
4606 with_check.submission_key(),
4607 digest,
4608 target.accepted_head(),
4609 )
4610 .expect("job identity should derive");
4611 let receipt = crate::db::schema::SchemaChangeReceipt::new(
4612 target.database_identity(),
4613 with_check.submission_key().clone(),
4614 digest,
4615 target.accepted_head().clone(),
4616 SchemaChangeOutcome::Pending {
4617 job: SchemaChangeJob::new(job_id),
4618 candidate_head: ExpectedAcceptedHead::Exact {
4619 revision: pending_candidate.revision().get().saturating_add(2),
4620 fingerprint: ExpectedSchemaFingerprint::from_bytes([0x76; 32]),
4621 },
4622 },
4623 )
4624 .expect("pending receipt should admit");
4625 let record = SchemaApplicationRecord::new(
4626 receipt,
4627 vec![
4628 SchemaChangeActivation::new(
4629 store_identity,
4630 entity_tag.value(),
4631 constraint_id.get(),
4632 )
4633 .expect("application activation should admit"),
4634 ],
4635 )
4636 .expect("pending application record should admit");
4637 let operation =
4638 SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4639 publish_accepted_schema_candidates_with_application_record(
4640 vec![AcceptedSchemaPublication::new(
4641 ABORT_STORE_PATH,
4642 store,
4643 current.revision(),
4644 &pending_candidate,
4645 )],
4646 operation,
4647 )
4648 .expect("pending candidate and record should publish atomically");
4649
4650 let started = continue_schema_application(&db, job_id, None)
4651 .expect("first continuation should durably start validation");
4652 assert_eq!(started.status(), &SchemaChangeProgressStatus::Started);
4653 let progress =
4654 abort_schema_application(&db, job_id, None).expect("pending application should abort");
4655 assert_eq!(progress.status(), &SchemaChangeProgressStatus::Aborted);
4656 assert!(matches!(
4657 progress.receipt().outcome(),
4658 SchemaChangeOutcome::Aborted { .. },
4659 ));
4660 let replay =
4661 abort_schema_application(&db, job_id, None).expect("terminal abort should replay");
4662 assert_eq!(replay, progress);
4663 assert_eq!(
4664 continue_schema_application(&db, job_id, None)
4665 .expect("continuation after abort should replay terminal state"),
4666 progress,
4667 );
4668
4669 let aborted = store
4670 .with_schema(SchemaStore::current_accepted_schema_bundle)
4671 .expect("accepted bundle should remain readable")
4672 .expect("aborted accepted bundle should exist");
4673 assert!(
4674 aborted
4675 .entity_snapshots()
4676 .get(&entity_tag)
4677 .expect("entity should remain after abort")
4678 .constraint_catalog()
4679 .activation(constraint_id)
4680 .is_none(),
4681 );
4682 assert!(
4683 aborted
4684 .source_bindings_for_tests()
4685 .constraint(entity_tag, &check_source)
4686 .is_none(),
4687 );
4688 assert!(
4689 store
4690 .with_schema(|schema| {
4691 schema.constraint_validation_job(entity_tag, constraint_id)
4692 })
4693 .expect("validation-job storage should remain readable")
4694 .is_none(),
4695 );
4696
4697 ABORT_DATA.with(|store| {
4698 ABORT_DATA_MEMORY.with(|memory| {
4699 *store.borrow_mut() = DataStore::init_journaled(memory.clone());
4700 });
4701 });
4702 ABORT_INDEX.with(|store| {
4703 ABORT_INDEX_MEMORY.with(|memory| {
4704 *store.borrow_mut() = IndexStore::init_journaled(memory.clone());
4705 });
4706 });
4707 ABORT_SCHEMA.with(|store| {
4708 ABORT_SCHEMA_MEMORY.with(|memory| {
4709 *store.borrow_mut() = SchemaStore::init_journaled(memory.clone());
4710 });
4711 });
4712 ABORT_JOURNAL.with(|store| {
4713 ABORT_JOURNAL_MEMORY.with(|memory| {
4714 *store.borrow_mut() = JournalTailStore::init(memory.clone());
4715 });
4716 });
4717 forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4718 drive_startup_recovery_to_completion(&db);
4719 assert_eq!(
4720 abort_schema_application(&db, job_id, None)
4721 .expect("recovered terminal abort should replay"),
4722 progress,
4723 );
4724 assert!(
4725 store
4726 .with_schema(|schema| {
4727 schema.constraint_validation_job(entity_tag, constraint_id)
4728 })
4729 .expect("recovered validation-job storage should remain readable")
4730 .is_none(),
4731 );
4732 let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4733 let startup_session =
4734 crate::db::DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &startup_root);
4735 ABORT_JOURNAL.with(|journal| {
4736 let journal = journal.borrow();
4737 assert!(
4738 journal
4739 .validate_current_tail_authority()
4740 .expect("recovered abort tail control should remain readable")
4741 .is_empty(),
4742 "recovered abort tail control should close exactly",
4743 );
4744 });
4745 assert_eq!(
4746 drive_generated_startup_recovery_page(
4747 &startup_session,
4748 &ABORT_REGISTRY,
4749 with_check.submission_key().as_str(),
4750 )
4751 .expect("an exact aborted generated submission should publish terminal startup state"),
4752 GeneratedStartupDriverStep::Terminal,
4753 );
4754 let failure = observe_generated_startup_state::<AbortCanister>(
4755 &ABORT_REGISTRY,
4756 with_check.submission_key().as_str(),
4757 )
4758 .expect_err("an aborted generated submission must not remain retryable forever");
4759 assert_eq!(
4760 failure.kind(),
4761 crate::db::StartupFailureKind::SchemaReconciliation,
4762 );
4763 assert_eq!(
4764 failure.diagnostic().error_code(),
4765 icydb_diagnostic_code::ErrorCode::RUNTIME_CONFLICT,
4766 );
4767 }
4768
4769 #[cfg(feature = "migration")]
4770 #[test]
4771 fn migration_planning_failures_retain_typed_public_classification() {
4772 use super::schema_migration_planning_error;
4773 use crate::db::schema::migration_planner::SchemaMigrationPlanningError;
4774 use icydb_diagnostic_code::{DiagnosticDetail, SchemaMigrationCode};
4775
4776 for (error, reason) in [
4777 (
4778 SchemaMigrationPlanningError::Unadopted,
4779 SchemaMigrationCode::Unadopted,
4780 ),
4781 (
4782 SchemaMigrationPlanningError::MissingMigration,
4783 SchemaMigrationCode::MissingMigration,
4784 ),
4785 (
4786 SchemaMigrationPlanningError::VersionGap,
4787 SchemaMigrationCode::VersionGap,
4788 ),
4789 (
4790 SchemaMigrationPlanningError::Downgrade,
4791 SchemaMigrationCode::Downgrade,
4792 ),
4793 (
4794 SchemaMigrationPlanningError::EmptyEntityVersionBump,
4795 SchemaMigrationCode::EmptyEntityVersionBump,
4796 ),
4797 (
4798 SchemaMigrationPlanningError::StaleAcceptedHead,
4799 SchemaMigrationCode::StaleAcceptedHead,
4800 ),
4801 (
4802 SchemaMigrationPlanningError::UnknownFromObject,
4803 SchemaMigrationCode::UnknownFromObject,
4804 ),
4805 (
4806 SchemaMigrationPlanningError::UnknownToObject,
4807 SchemaMigrationCode::UnknownToObject,
4808 ),
4809 (
4810 SchemaMigrationPlanningError::KindMismatch,
4811 SchemaMigrationCode::KindMismatch,
4812 ),
4813 (
4814 SchemaMigrationPlanningError::IdentityConflict,
4815 SchemaMigrationCode::IdentityConflict,
4816 ),
4817 (
4818 SchemaMigrationPlanningError::UnexplainedSchemaDifference,
4819 SchemaMigrationCode::UnexplainedSchemaDifference,
4820 ),
4821 (
4822 SchemaMigrationPlanningError::UnsupportedTransform,
4823 SchemaMigrationCode::UnsupportedTransform,
4824 ),
4825 (
4826 SchemaMigrationPlanningError::RekeyedCatalogInvalid,
4827 SchemaMigrationCode::CandidateMismatch,
4828 ),
4829 (
4830 SchemaMigrationPlanningError::CandidateMismatch,
4831 SchemaMigrationCode::CandidateMismatch,
4832 ),
4833 (
4834 SchemaMigrationPlanningError::CorruptLineage,
4835 SchemaMigrationCode::ProgressCorrupt,
4836 ),
4837 ] {
4838 let diagnostic = schema_migration_planning_error(error).diagnostic();
4839 assert_eq!(
4840 diagnostic.detail(),
4841 Some(&DiagnosticDetail::SchemaMigration { reason }),
4842 );
4843 assert_eq!(diagnostic.code(), reason.diagnostic_code());
4844 }
4845 }
4846
4847 #[cfg(feature = "migration")]
4848 #[test]
4849 #[expect(
4850 clippy::too_many_lines,
4851 reason = "the validation replay, staging, and unchanged-row assertions form one scenario"
4852 )]
4853 fn physical_migration_validation_is_bounded_staged_and_does_not_rewrite_rows() {
4854 use std::convert::Infallible;
4855
4856 use super::{defer_generated_schema_application_for_prepared_migration, migrate_schema};
4857 use crate::db::{
4858 data::StoreVisit,
4859 index::{IndexEntryValue, IndexId, IndexKey, IndexKeyKind},
4860 key_taxonomy::{PrimaryKeyComponent, PrimaryKeyValue},
4861 schema::{SchemaMigrationCommand, SchemaMigrationPhase},
4862 };
4863 use crate::types::EntityTag;
4864
4865 let db = Db::<MigrationCanister>::new(
4866 &MIGRATION_REGISTRY,
4867 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4868 );
4869 drive_startup_recovery_to_completion(&db);
4870 let initial_target = schema_application_target(&db).expect("initial target should issue");
4871 let store_identity = initial_target
4872 .stores()
4873 .first()
4874 .expect("migration store should exist")
4875 .identity();
4876 let initial = validation_migration_proposal(
4877 ValidationMigrationShape::Clean,
4878 false,
4879 initial_target.accepted_head().clone(),
4880 initial_target.database_identity(),
4881 store_identity,
4882 );
4883 apply_schema(&db, &initial).expect("initial schema should publish");
4884
4885 let session = DbSession::<MigrationCanister>::new(
4886 &MIGRATION_REGISTRY,
4887 &crate::db::RequestExecutionRoot::__new_runtime_root(),
4888 );
4889 for (id, value) in [(1, 7), (2, 8)] {
4890 session
4891 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4892 entity: "MigratingItem".to_string(),
4893 patch: DynamicStructuralPatch::new(vec![
4894 (
4895 "id".to_string(),
4896 DynamicWriteCell::Value(InputValue::nat64(id)),
4897 ),
4898 (
4899 "old_value".to_string(),
4900 DynamicWriteCell::Value(InputValue::int64(value)),
4901 ),
4902 ]),
4903 })
4904 .expect("predecessor row should insert");
4905 }
4906 let store = db
4907 .store_handle(MIGRATION_STORE_PATH)
4908 .expect("migration store should resolve");
4909 let row_bytes = || {
4910 store.with_data(|data| {
4911 let mut rows = Vec::new();
4912 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
4913 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
4914 Ok(StoreVisit::Continue)
4915 });
4916 result.expect("infallible row visit should complete");
4917 rows
4918 })
4919 };
4920 let before_rows = row_bytes();
4921
4922 let target = schema_application_target(&db).expect("migration target should issue");
4923 let proposal = validation_migration_proposal(
4924 ValidationMigrationShape::Clean,
4925 true,
4926 target.accepted_head().clone(),
4927 target.database_identity(),
4928 store_identity,
4929 );
4930 let plan = proposal
4931 .migration()
4932 .expect("migration plan should exist")
4933 .digest();
4934 let command = || SchemaMigrationCommand::Advance {
4935 expected_database: target.database_identity(),
4936 expected_head: target.accepted_head().clone(),
4937 expected_plan: plan,
4938 acknowledged_finding_page: None,
4939 };
4940 assert_eq!(
4941 migrate_schema(&db, &proposal, command())
4942 .unwrap_or_else(|error| {
4943 panic!(
4944 "physical migration should prepare: {:?}",
4945 error.diagnostic()
4946 )
4947 })
4948 .phase(),
4949 SchemaMigrationPhase::Prepared,
4950 );
4951 assert_eq!(
4952 migrate_schema(&db, &proposal, command())
4953 .expect("physical migration should enter validation")
4954 .phase(),
4955 SchemaMigrationPhase::Validating,
4956 );
4957 let record = super::load_schema_migration_record()
4958 .expect("migration record should remain readable")
4959 .expect("validating migration record should exist");
4960 let planned = super::recompile_active_physical_migration(&db, &proposal, &record)
4961 .expect("the exact active plan should recompile");
4962 for _ in 0..2 {
4963 let page = super::validate_migration_page(&db, &planned, record.progress())
4964 .expect("the same validation page should remain replayable");
4965 let (progress, staged, exhausted) = page.into_parts();
4966 assert!(progress.findings().is_empty());
4967 assert!(exhausted);
4968 super::stage_migration_index_entries(staged)
4969 .expect("staging before a cursor marker should be idempotent");
4970 }
4971 assert_eq!(
4972 store.with_index(IndexStore::len),
4973 2,
4974 "replaying an uncheckpointed page must retain one exact staged key per row",
4975 );
4976 let ready =
4977 migrate_schema(&db, &proposal, command()).expect("bounded validation should complete");
4978 assert_eq!(ready.phase(), SchemaMigrationPhase::ReadyToRewrite);
4979 assert_eq!(ready.rows_validated(), 2);
4980 assert!(ready.findings().is_empty());
4981 assert_eq!(row_bytes(), before_rows, "validation must not rewrite rows");
4982 assert_eq!(
4983 store.with_index(IndexStore::len),
4984 2,
4985 "the isolated candidate unique generation should be durably staged",
4986 );
4987 store.with_index_mut(|index| {
4988 for ordinal in 0..513_u64 {
4989 let component = ordinal.to_be_bytes();
4990 let key = IndexKey::new_from_components_with_primary_key_value(
4991 &IndexId::new(EntityTag::new(2), 0),
4992 IndexKeyKind::User,
4993 &[component],
4994 &PrimaryKeyValue::from(PrimaryKeyComponent::Nat64(ordinal)),
4995 )
4996 .expect("unrelated abort-scan key should build")
4997 .to_raw()
4998 .expect("unrelated abort-scan key should encode");
4999 index.insert(key, IndexEntryValue::presence());
5000 }
5001 });
5002 let abort = || SchemaMigrationCommand::Abort {
5003 expected_database: target.database_identity(),
5004 expected_head: target.accepted_head().clone(),
5005 expected_plan: plan,
5006 };
5007 let cleaning = migrate_schema(&db, &proposal, abort())
5008 .expect("the first bounded abort cleanup page should publish");
5009 assert_eq!(cleaning.phase(), SchemaMigrationPhase::ReadyToRewrite);
5010 assert_eq!(store.with_index(IndexStore::len), 513);
5011 let aborted =
5012 migrate_schema(&db, &proposal, abort()).expect("pre-rewrite migration should abort");
5013 assert_eq!(aborted.phase(), SchemaMigrationPhase::Aborted);
5014 assert_eq!(
5015 store.with_index(IndexStore::len),
5016 513,
5017 "abort must remove only planner-invisible candidate generations",
5018 );
5019 assert_eq!(
5020 row_bytes(),
5021 before_rows,
5022 "abort must retain predecessor rows"
5023 );
5024 assert!(
5025 !defer_generated_schema_application_for_prepared_migration(&db, &proposal)
5026 .expect("terminal aborted record must not block generated startup"),
5027 );
5028 }
5029
5030 #[cfg(feature = "migration")]
5031 #[test]
5032 #[expect(
5033 clippy::too_many_lines,
5034 reason = "the interrupted rewrite, recovery, final proof, and publication form one scenario"
5035 )]
5036 fn physical_migration_rewrite_recovers_and_publishes_one_complete_candidate() {
5037 use super::{
5038 defer_generated_schema_application_for_prepared_migration, migrate_schema,
5039 schema_migration_status_for_target,
5040 };
5041 use crate::db::{
5042 data::{CanonicalSlotReader, DecodedDataStoreKey, StoreVisit, StructuralSlotReader},
5043 schema::{
5044 MigrationRewriteInterruption, SchemaMigrationCommand, SchemaMigrationPhase,
5045 ensure_schema_migration_ready_for_ordinary_operations,
5046 interrupt_next_migration_rewrite_at,
5047 },
5048 };
5049 use crate::error::InternalError;
5050
5051 let db = Db::<MigrationExecutionCanister>::new(
5052 &MIGRATION_EXECUTION_REGISTRY,
5053 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5054 );
5055 drive_startup_recovery_to_completion(&db);
5056 let initial_target = schema_application_target(&db).expect("initial target should issue");
5057 let store_identity = initial_target
5058 .stores()
5059 .first()
5060 .expect("migration execution store should exist")
5061 .identity();
5062 let initial = validation_migration_proposal(
5063 ValidationMigrationShape::Clean,
5064 false,
5065 initial_target.accepted_head().clone(),
5066 initial_target.database_identity(),
5067 store_identity,
5068 );
5069 apply_schema(&db, &initial).expect("initial schema should publish");
5070 let session = DbSession::<MigrationExecutionCanister>::new(
5071 &MIGRATION_EXECUTION_REGISTRY,
5072 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5073 );
5074 for (id, value) in [(1, 7), (2, 8), (3, 9)] {
5075 session
5076 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5077 entity: "MigratingItem".to_string(),
5078 patch: DynamicStructuralPatch::new(vec![
5079 (
5080 "id".to_string(),
5081 DynamicWriteCell::Value(InputValue::nat64(id)),
5082 ),
5083 (
5084 "old_value".to_string(),
5085 DynamicWriteCell::Value(InputValue::int64(value)),
5086 ),
5087 ]),
5088 })
5089 .expect("predecessor row should insert");
5090 }
5091 let target = schema_application_target(&db).expect("migration target should issue");
5092 let proposal = validation_migration_proposal(
5093 ValidationMigrationShape::Clean,
5094 true,
5095 target.accepted_head().clone(),
5096 target.database_identity(),
5097 store_identity,
5098 );
5099 let plan = proposal
5100 .migration()
5101 .expect("migration plan should exist")
5102 .digest();
5103 let command = || SchemaMigrationCommand::Advance {
5104 expected_database: target.database_identity(),
5105 expected_head: target.accepted_head().clone(),
5106 expected_plan: plan,
5107 acknowledged_finding_page: None,
5108 };
5109 for expected in [
5110 SchemaMigrationPhase::Prepared,
5111 SchemaMigrationPhase::Validating,
5112 SchemaMigrationPhase::ReadyToRewrite,
5113 SchemaMigrationPhase::RewritingRows,
5114 ] {
5115 assert_eq!(
5116 migrate_schema(&db, &proposal, command())
5117 .expect("migration phase should advance")
5118 .phase(),
5119 expected,
5120 );
5121 }
5122
5123 for interruption in [
5124 MigrationRewriteInterruption::MarkerPersisted,
5125 MigrationRewriteInterruption::JournalPublished,
5126 MigrationRewriteInterruption::PhysicalApplied,
5127 ] {
5128 interrupt_next_migration_rewrite_at(interruption);
5129 migrate_schema(&db, &proposal, command())
5130 .expect_err("injected interruption should retain the rewrite marker");
5131
5132 forget_recovered_domain_for_tests(&db)
5133 .expect("upgrade should reset recovery ownership");
5134 drive_startup_recovery_to_completion(&db);
5135 }
5136
5137 let rebuilding = schema_migration_status_for_target(
5138 &db,
5139 &proposal,
5140 &schema_application_target(&db).expect("recovered target should issue"),
5141 )
5142 .expect("recovered status should remain readable");
5143 assert_eq!(rebuilding.phase(), SchemaMigrationPhase::RebuildingIndexes);
5144 assert_eq!(rebuilding.rows_rewritten(), 3);
5145 assert_eq!(
5146 migrate_schema(&db, &proposal, command())
5147 .expect("derived generations should complete")
5148 .phase(),
5149 SchemaMigrationPhase::FinalValidation,
5150 );
5151 assert_eq!(
5152 migrate_schema(&db, &proposal, command())
5153 .expect("final validation should complete")
5154 .phase(),
5155 SchemaMigrationPhase::Publishing,
5156 );
5157 let applied = migrate_schema(&db, &proposal, command())
5158 .expect("candidate publication should complete atomically");
5159 assert_eq!(applied.phase(), SchemaMigrationPhase::Applied);
5160 assert_eq!(applied.rows_rewritten(), 3);
5161 assert_eq!(applied.indexes_rebuilt(), 1);
5162 assert_ne!(applied.accepted_head(), target.accepted_head());
5163 let terminal_target = schema_application_target(&db).expect("terminal target should issue");
5164 let terminal_proposal = validation_migration_proposal(
5165 ValidationMigrationShape::Clean,
5166 true,
5167 terminal_target.accepted_head().clone(),
5168 terminal_target.database_identity(),
5169 store_identity,
5170 );
5171 assert!(
5172 !defer_generated_schema_application_for_prepared_migration(&db, &terminal_proposal,)
5173 .expect("terminal record must not block generated startup"),
5174 );
5175
5176 let store = db
5177 .store_handle(MIGRATION_EXECUTION_STORE_PATH)
5178 .expect("migration execution store should resolve");
5179 let runtime = db
5180 .accepted_runtime_entity_for_path("MigratingItem")
5181 .expect("published candidate entity should resolve");
5182 let selection = store
5183 .with_schema(|schema| {
5184 schema.current_accepted_catalog_selection(
5185 runtime.entity_tag(),
5186 runtime.entity_path(),
5187 runtime.store_path(),
5188 )
5189 })
5190 .expect("candidate selection should remain readable")
5191 .expect("candidate selection should exist");
5192 let contract = crate::db::data::AcceptedStructuralRowAuthority::from_catalog_selection(
5193 runtime.entity_path(),
5194 &selection,
5195 )
5196 .expect("candidate row authority should compile")
5197 .into_row_contract();
5198 let mut values = Vec::new();
5199 store
5200 .with_data(|data| {
5201 data.visit_entries(|key, row| {
5202 let decoded = DecodedDataStoreKey::try_from_raw(key)
5203 .expect("rewritten key should decode");
5204 let reader =
5205 StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
5206 row, &contract,
5207 )
5208 .expect("rewritten row should use the candidate layout");
5209 reader
5210 .validate_primary_key(&decoded)
5211 .expect("rewritten row and key should remain bound");
5212 values.push(
5213 reader
5214 .required_value_by_contract(1)
5215 .expect("candidate value slot should decode"),
5216 );
5217 Ok::<StoreVisit, InternalError>(StoreVisit::Continue)
5218 })
5219 })
5220 .expect("rewritten row scan should complete");
5221 assert_eq!(
5222 values,
5223 vec![
5224 crate::value::Value::Nat64(7),
5225 crate::value::Value::Nat64(8),
5226 crate::value::Value::Nat64(9),
5227 ],
5228 );
5229 assert_eq!(store.with_index(IndexStore::len), 3);
5230 let accepted = store
5231 .with_schema(SchemaStore::current_accepted_schema_bundle)
5232 .expect("published candidate bundle should remain readable")
5233 .expect("published candidate bundle should exist");
5234 let entity_source = EntitySourceKey::try_new("MigratingItem")
5235 .expect("migration entity source should admit");
5236 let entity_tag = accepted
5237 .source_bindings_for_tests()
5238 .entity(&entity_source)
5239 .expect("candidate entity source should remain bound");
5240 let old_value =
5241 FieldSourceKey::try_new("old_value").expect("predecessor source should admit");
5242 let current_value =
5243 FieldSourceKey::try_new("value").expect("candidate source should admit");
5244 assert_eq!(
5245 accepted
5246 .source_bindings_for_tests()
5247 .field(entity_tag, &old_value),
5248 None,
5249 );
5250 assert!(
5251 accepted
5252 .source_bindings_for_tests()
5253 .field(entity_tag, ¤t_value)
5254 .is_some(),
5255 );
5256 ensure_schema_migration_ready_for_ordinary_operations()
5257 .expect("terminal publication must clear the database-wide gate");
5258 }
5259
5260 #[cfg(feature = "migration")]
5261 #[test]
5262 #[expect(
5263 clippy::too_many_lines,
5264 reason = "all four finding families share one ordered historical scan fixture"
5265 )]
5266 fn physical_migration_validation_reports_every_typed_finding_family_without_writes() {
5267 use std::convert::Infallible;
5268
5269 use super::migrate_schema;
5270 use crate::db::{
5271 data::StoreVisit,
5272 schema::{SchemaMigrationCommand, SchemaMigrationFindingKind, SchemaMigrationPhase},
5273 };
5274
5275 let db = Db::<MigrationFindingCanister>::new(
5276 &MIGRATION_FINDING_REGISTRY,
5277 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5278 );
5279 drive_startup_recovery_to_completion(&db);
5280 let initial_target = schema_application_target(&db).expect("initial target should issue");
5281 let store_identity = initial_target
5282 .stores()
5283 .first()
5284 .expect("migration finding store should exist")
5285 .identity();
5286 let initial = validation_migration_proposal(
5287 ValidationMigrationShape::AllFindingFamilies,
5288 false,
5289 initial_target.accepted_head().clone(),
5290 initial_target.database_identity(),
5291 store_identity,
5292 );
5293 apply_schema(&db, &initial).expect("initial finding schema should publish");
5294
5295 let session = DbSession::<MigrationFindingCanister>::new(
5296 &MIGRATION_FINDING_REGISTRY,
5297 &crate::db::RequestExecutionRoot::__new_runtime_root(),
5298 );
5299 session
5300 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5301 entity: "MigrationTarget".to_string(),
5302 patch: DynamicStructuralPatch::new(vec![(
5303 "id".to_string(),
5304 DynamicWriteCell::Value(InputValue::nat64(7)),
5305 )]),
5306 })
5307 .expect("relation target should insert");
5308 for (id, value) in [(1, 9), (2, 8), (3, 7), (4, 7), (5, 300)] {
5309 session
5310 .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5311 entity: "MigratingItem".to_string(),
5312 patch: DynamicStructuralPatch::new(vec![
5313 (
5314 "id".to_string(),
5315 DynamicWriteCell::Value(InputValue::nat64(id)),
5316 ),
5317 (
5318 "old_value".to_string(),
5319 DynamicWriteCell::Value(InputValue::int64(value)),
5320 ),
5321 ]),
5322 })
5323 .expect("predecessor finding row should insert");
5324 }
5325 let store = db
5326 .store_handle(MIGRATION_FINDING_STORE_PATH)
5327 .expect("migration finding store should resolve");
5328 let row_bytes = || {
5329 store.with_data(|data| {
5330 let mut rows = Vec::new();
5331 let result: Result<(), Infallible> = data.visit_entries(|key, row| {
5332 rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
5333 Ok(StoreVisit::Continue)
5334 });
5335 result.expect("infallible row visit should complete");
5336 rows
5337 })
5338 };
5339 let before_rows = row_bytes();
5340
5341 let target = schema_application_target(&db).expect("migration target should issue");
5342 let proposal = validation_migration_proposal(
5343 ValidationMigrationShape::AllFindingFamilies,
5344 true,
5345 target.accepted_head().clone(),
5346 target.database_identity(),
5347 store_identity,
5348 );
5349 let plan = proposal
5350 .migration()
5351 .expect("migration plan should exist")
5352 .digest();
5353 let command = || SchemaMigrationCommand::Advance {
5354 expected_database: target.database_identity(),
5355 expected_head: target.accepted_head().clone(),
5356 expected_plan: plan,
5357 acknowledged_finding_page: None,
5358 };
5359 assert_eq!(
5360 migrate_schema(&db, &proposal, command())
5361 .expect("finding migration should prepare")
5362 .phase(),
5363 SchemaMigrationPhase::Prepared,
5364 );
5365 assert_eq!(
5366 migrate_schema(&db, &proposal, command())
5367 .expect("finding migration should enter validation")
5368 .phase(),
5369 SchemaMigrationPhase::Validating,
5370 );
5371 let rejected =
5372 migrate_schema(&db, &proposal, command()).expect("validation should report findings");
5373 assert_eq!(rejected.phase(), SchemaMigrationPhase::Rejected);
5374 assert_eq!(rejected.rows_validated(), 5);
5375 assert_eq!(
5376 rejected
5377 .findings()
5378 .iter()
5379 .map(crate::db::schema::SchemaMigrationFinding::kind)
5380 .collect::<Vec<_>>(),
5381 vec![
5382 SchemaMigrationFindingKind::Constraint,
5383 SchemaMigrationFindingKind::Relation,
5384 SchemaMigrationFindingKind::UniqueIndex,
5385 SchemaMigrationFindingKind::Transform,
5386 ],
5387 );
5388 assert_eq!(
5389 row_bytes(),
5390 before_rows,
5391 "rejected validation must not rewrite accepted rows"
5392 );
5393 assert_eq!(
5394 store.with_index(IndexStore::len),
5395 0,
5396 "a rejected page must not publish any staged generation"
5397 );
5398 }
5399
5400 #[cfg(feature = "migration")]
5401 #[test]
5402 fn exact_migration_retry_binds_the_terminal_head_not_the_predecessor_head() {
5403 use super::exact_migration_replay_target;
5404
5405 let db = Db::<EvolutionCanister>::new(
5406 &EVOLUTION_REGISTRY,
5407 crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5408 );
5409 drive_startup_recovery_to_completion(&db);
5410 let initial_target = schema_application_target(&db).expect("initial target should issue");
5411 let (proposal, _, _) = generated_check_proposal(
5412 initial_target.accepted_head().clone(),
5413 "migration-retry-initial",
5414 false,
5415 initial_target.database_identity(),
5416 initial_target
5417 .stores()
5418 .first()
5419 .expect("test store should exist")
5420 .identity(),
5421 );
5422 apply_schema(&db, &proposal).expect("initial schema should publish");
5423 let current_target = schema_application_target(&db).expect("current target should issue");
5424 assert_ne!(
5425 current_target.accepted_head(),
5426 initial_target.accepted_head(),
5427 );
5428
5429 let receipt = SchemaChangeReceipt::new(
5430 current_target.database_identity(),
5431 SchemaSubmissionKey::try_new("migration/retry")
5432 .expect("migration submission should admit"),
5433 SchemaProposalDigest::from_bytes([0x77; 32]),
5434 initial_target.accepted_head().clone(),
5435 SchemaChangeOutcome::Applied {
5436 accepted_head: current_target.accepted_head().clone(),
5437 },
5438 )
5439 .expect("terminal migration receipt should admit");
5440 let record = SchemaApplicationRecord::new(receipt, Vec::new())
5441 .expect("terminal migration record should admit");
5442
5443 assert_eq!(
5444 exact_migration_replay_target(&db, current_target.database_identity(), &record,)
5445 .expect("exact retry should resolve the terminal target"),
5446 current_target,
5447 );
5448 }
5449}