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