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