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