Skip to main content

icydb_core/db/schema/
application.rs

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