Skip to main content

icydb_core/db/schema/
application.rs

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