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    use super::{
3116        AcceptedSchemaPublication, AcceptedStoreHead, DirectGeneratedRowLocalProof,
3117        PendingGeneratedRowLocalConstraint, abort_schema_application,
3118        aborted_generated_row_local_candidate, accepted_head_after_candidates,
3119        application_authorities, apply_schema, continue_schema_application, derive_accepted_head,
3120        derive_schema_change_job_id, final_candidates_for_pending_row_local_constraint,
3121        generated_database_identity, include_identity_state_count, lower_existing_schema_proposal,
3122        lower_initial_schema_proposal, publish_accepted_schema_candidates_with_application_record,
3123        require_exact_empty_entity_count, schema_application_target,
3124    };
3125    use crate::{
3126        db::{
3127            DatabaseStartupState, Db, GeneratedStartupDriverStep,
3128            commit::{
3129                RecoveryProgress, continue_recovery, database_incarnation_id,
3130                forget_recovered_domain_for_tests,
3131            },
3132            data::DataStore,
3133            drive_generated_startup_recovery_page,
3134            index::IndexStore,
3135            journal::JournalTailStore,
3136            observe_generated_startup_state,
3137            registry::{
3138                StoreAllocationIdentities, StoreAllocationIdentity, StoreHandle, StoreRegistry,
3139                StoreRuntimeStorageCapabilities,
3140            },
3141            schema::{
3142                AcceptedConstraintKind, AcceptedRuleOperation, AcceptedSchemaRevisionBundle,
3143                CandidateSchemaRevision, ConstraintOrigin, ConstraintValidationJob,
3144                ExistingProposalStore, ProposalStoreTarget, SchemaApplicationRecord,
3145                SchemaApplicationRecordOp, SchemaChangeActivation, SchemaChangeJob,
3146                SchemaChangeOutcome, SchemaChangeProgressStatus, SchemaStore,
3147                cardinality_build::{
3148                    CardinalityBuildAuthority, CardinalityGenerationPageOutcome,
3149                    drive_cardinality_generation_page,
3150                },
3151            },
3152        },
3153        error::{ErrorClass, ErrorOrigin},
3154        testing::test_memory,
3155        traits::{CanisterKind, Path},
3156    };
3157    use ic_stable_structures::{DefaultMemoryImpl, memory_manager::VirtualMemory};
3158    use icydb_schema::{
3159        ConstraintFragment, ConstraintSourceKey, DeclaredEntityVersion, EntityFragment,
3160        EntitySourceKey, EntityStoreAssignment, ExpectedAcceptedHead, ExpectedSchemaFingerprint,
3161        FieldFragment, FieldInsertPolicy, FieldSourceKey, FieldType, NamedTypeFragment,
3162        RuleSourceKey, ScalarLiteral, ScalarType, SchemaCapability, SchemaFragment, SchemaName,
3163        SchemaProposal, SchemaSubmissionKey, SourceCheckExpr, SourceCheckInstruction,
3164        SourceRuleOperation, TargetDatabaseIdentity, TargetStoreIdentity, TargetedRuleFragment,
3165        TypeSourceKey,
3166    };
3167    use std::cell::RefCell;
3168
3169    fn drive_startup_recovery_to_completion<C: CanisterKind>(db: &Db<C>) {
3170        for _ in 0..1_024 {
3171            match continue_recovery(db).expect("test startup recovery page should succeed") {
3172                RecoveryProgress::Complete => return,
3173                RecoveryProgress::Pending => {}
3174            }
3175        }
3176        panic!("test startup recovery should complete within 1,024 bounded pages");
3177    }
3178
3179    fn drive_cardinality_to_ready(store: StoreHandle) {
3180        let journal = store
3181            .journal_tail_store()
3182            .expect("cardinality fixture store should be journaled");
3183        for _ in 0..8 {
3184            let outcome = store
3185                .with_data(|data| {
3186                    store.with_index(|index| {
3187                        store.with_schema_mut(|schema| {
3188                            drive_cardinality_generation_page(data, index, schema, |schema| {
3189                                CardinalityBuildAuthority::derive(
3190                                    schema,
3191                                    database_incarnation_id()?,
3192                                    store.allocation_identities(),
3193                                    journal.with_borrow(JournalTailStore::fold_watermark)?,
3194                                )
3195                            })
3196                        })
3197                    })
3198                })
3199                .expect("bounded cardinality generation should advance");
3200            if outcome == CardinalityGenerationPageOutcome::Quiescent {
3201                return;
3202            }
3203        }
3204        panic!("cardinality generation should become Ready within eight bounded pages");
3205    }
3206
3207    #[cfg(feature = "migration")]
3208    use crate::db::schema::SchemaChangeReceipt;
3209    use crate::{
3210        db::{DbSession, DynamicMutation, DynamicStructuralPatch, DynamicWriteCell},
3211        value::InputValue,
3212    };
3213    #[cfg(feature = "migration")]
3214    use icydb_schema::{
3215        EntityMigration, IndexFragment, IndexKeyFragment, RelationDeleteAction, RelationFragment,
3216        SchemaMigrationPlan, SchemaMigrationTransform, SchemaProposalDigest,
3217    };
3218
3219    fn version_one() -> DeclaredEntityVersion {
3220        DeclaredEntityVersion::try_new(1).expect("fixture version should admit")
3221    }
3222
3223    const ABORT_STORE_PATH: &str = "schema_application_tests::AbortStore";
3224    const EVOLUTION_STORE_PATH: &str = "schema_application_tests::EvolutionStore";
3225    #[cfg(feature = "migration")]
3226    const MIGRATION_STORE_PATH: &str = "schema_application_tests::MigrationStore";
3227    #[cfg(feature = "migration")]
3228    const MIGRATION_EXECUTION_STORE_PATH: &str =
3229        "schema_application_tests::MigrationExecutionStore";
3230    #[cfg(feature = "migration")]
3231    const MIGRATION_FINDING_STORE_PATH: &str = "schema_application_tests::MigrationFindingStore";
3232
3233    #[test]
3234    fn database_identity_state_capacity_combines_store_inventories_exactly() {
3235        let below = include_identity_state_count(0, 65_535)
3236            .expect("the first store inventory should remain below the database cap");
3237        let exact = include_identity_state_count(below, 1)
3238            .expect("the combined database boundary should admit");
3239        assert_eq!(exact, 65_536);
3240
3241        let error = include_identity_state_count(exact, 1)
3242            .expect_err("the next owner in another store must reject");
3243        assert_eq!(error.class(), ErrorClass::Unsupported);
3244        assert_eq!(error.origin(), ErrorOrigin::Identity);
3245    }
3246
3247    #[test]
3248    fn generated_database_identity_cache_is_bound_to_the_incarnation() {
3249        let first_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x41);
3250        let second_incarnation = crate::db::DatabaseIncarnationId::for_tests(0x42);
3251        let first = generated_database_identity(&ABORT_REGISTRY, first_incarnation);
3252        assert_eq!(
3253            generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3254            first,
3255        );
3256        let second = generated_database_identity(&ABORT_REGISTRY, second_incarnation);
3257        assert_ne!(second, first);
3258        assert_eq!(
3259            generated_database_identity(&ABORT_REGISTRY, first_incarnation),
3260            first,
3261        );
3262    }
3263
3264    #[test]
3265    fn exact_empty_entity_proof_distinguishes_corruption_from_non_empty_input() {
3266        let corrupt = require_exact_empty_entity_count(None)
3267            .expect_err("uninspectable cardinality must fail closed");
3268        assert_eq!(corrupt.class(), ErrorClass::Corruption);
3269
3270        let non_empty = require_exact_empty_entity_count(Some(1))
3271            .expect_err("non-empty cardinality must reject removal");
3272        assert_eq!(non_empty.class(), ErrorClass::Unsupported);
3273        assert!(require_exact_empty_entity_count(Some(0)).is_ok());
3274    }
3275
3276    thread_local! {
3277        static ABORT_DATA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(180);
3278        static ABORT_INDEX_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(181);
3279        static ABORT_SCHEMA_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(182);
3280        static ABORT_JOURNAL_MEMORY: VirtualMemory<DefaultMemoryImpl> = test_memory(183);
3281        static ABORT_DATA: RefCell<DataStore> =
3282            ABORT_DATA_MEMORY.with(|memory| {
3283                RefCell::new(DataStore::init_journaled(memory.clone()))
3284            });
3285        static ABORT_INDEX: RefCell<IndexStore> =
3286            ABORT_INDEX_MEMORY.with(|memory| {
3287                RefCell::new(IndexStore::init_journaled(memory.clone()))
3288            });
3289        static ABORT_SCHEMA: RefCell<SchemaStore> =
3290            ABORT_SCHEMA_MEMORY.with(|memory| {
3291                RefCell::new(SchemaStore::init_journaled(memory.clone()))
3292            });
3293        static ABORT_JOURNAL: RefCell<JournalTailStore> =
3294            ABORT_JOURNAL_MEMORY.with(|memory| {
3295                RefCell::new(JournalTailStore::init(memory.clone()))
3296            });
3297        static ABORT_REGISTRY: StoreRegistry = {
3298            let mut registry = StoreRegistry::new();
3299            registry.register_journaled_store(
3300                ABORT_STORE_PATH,
3301                &ABORT_DATA,
3302                &ABORT_INDEX,
3303                &ABORT_SCHEMA,
3304                &ABORT_JOURNAL,
3305                StoreAllocationIdentities::new_journaled(
3306                    StoreAllocationIdentity::new(180, "icydb.test.application_abort.data.v1"),
3307                    StoreAllocationIdentity::new(181, "icydb.test.application_abort.index.v1"),
3308                    StoreAllocationIdentity::new(182, "icydb.test.application_abort.schema.v1"),
3309                    StoreAllocationIdentity::new(183, "icydb.test.application_abort.journal.v1"),
3310                ),
3311                StoreRuntimeStorageCapabilities::journaled(),
3312            ).expect("abort journaled store should register");
3313            registry
3314        };
3315    }
3316
3317    #[cfg(feature = "migration")]
3318    thread_local! {
3319        static MIGRATION_EXECUTION_DATA: RefCell<DataStore> =
3320            RefCell::new(DataStore::init_journaled(test_memory(210)));
3321        static MIGRATION_EXECUTION_INDEX: RefCell<IndexStore> =
3322            RefCell::new(IndexStore::init_journaled(test_memory(211)));
3323        static MIGRATION_EXECUTION_SCHEMA: RefCell<SchemaStore> =
3324            RefCell::new(SchemaStore::init_journaled(test_memory(212)));
3325        static MIGRATION_EXECUTION_JOURNAL: RefCell<JournalTailStore> =
3326            RefCell::new(JournalTailStore::init(test_memory(213)));
3327        static MIGRATION_EXECUTION_REGISTRY: StoreRegistry = {
3328            let mut registry = StoreRegistry::new();
3329            registry.register_journaled_store(
3330                MIGRATION_EXECUTION_STORE_PATH,
3331                &MIGRATION_EXECUTION_DATA,
3332                &MIGRATION_EXECUTION_INDEX,
3333                &MIGRATION_EXECUTION_SCHEMA,
3334                &MIGRATION_EXECUTION_JOURNAL,
3335                StoreAllocationIdentities::new_journaled(
3336                    StoreAllocationIdentity::new(210, "icydb.test.migration_execution.data.v1"),
3337                    StoreAllocationIdentity::new(211, "icydb.test.migration_execution.index.v1"),
3338                    StoreAllocationIdentity::new(212, "icydb.test.migration_execution.schema.v1"),
3339                    StoreAllocationIdentity::new(213, "icydb.test.migration_execution.journal.v1"),
3340                ),
3341                StoreRuntimeStorageCapabilities::journaled(),
3342            ).expect("migration execution store should register");
3343            registry
3344        };
3345    }
3346
3347    #[cfg(feature = "migration")]
3348    thread_local! {
3349        static MIGRATION_DATA: RefCell<DataStore> =
3350            RefCell::new(DataStore::init_journaled(test_memory(200)));
3351        static MIGRATION_INDEX: RefCell<IndexStore> =
3352            RefCell::new(IndexStore::init_journaled(test_memory(201)));
3353        static MIGRATION_SCHEMA: RefCell<SchemaStore> =
3354            RefCell::new(SchemaStore::init_journaled(test_memory(202)));
3355        static MIGRATION_JOURNAL: RefCell<JournalTailStore> =
3356            RefCell::new(JournalTailStore::init(test_memory(203)));
3357        static MIGRATION_REGISTRY: StoreRegistry = {
3358            let mut registry = StoreRegistry::new();
3359            registry.register_journaled_store(
3360                MIGRATION_STORE_PATH,
3361                &MIGRATION_DATA,
3362                &MIGRATION_INDEX,
3363                &MIGRATION_SCHEMA,
3364                &MIGRATION_JOURNAL,
3365                StoreAllocationIdentities::new_journaled(
3366                    StoreAllocationIdentity::new(200, "icydb.test.migration_validation.data.v1"),
3367                    StoreAllocationIdentity::new(201, "icydb.test.migration_validation.index.v1"),
3368                    StoreAllocationIdentity::new(202, "icydb.test.migration_validation.schema.v1"),
3369                    StoreAllocationIdentity::new(203, "icydb.test.migration_validation.journal.v1"),
3370                ),
3371                StoreRuntimeStorageCapabilities::journaled(),
3372            ).expect("migration validation store should register");
3373            registry
3374        };
3375    }
3376
3377    #[cfg(feature = "migration")]
3378    thread_local! {
3379        static MIGRATION_FINDING_DATA: RefCell<DataStore> =
3380            RefCell::new(DataStore::init_journaled(test_memory(206)));
3381        static MIGRATION_FINDING_INDEX: RefCell<IndexStore> =
3382            RefCell::new(IndexStore::init_journaled(test_memory(207)));
3383        static MIGRATION_FINDING_SCHEMA: RefCell<SchemaStore> =
3384            RefCell::new(SchemaStore::init_journaled(test_memory(208)));
3385        static MIGRATION_FINDING_JOURNAL: RefCell<JournalTailStore> =
3386            RefCell::new(JournalTailStore::init(test_memory(209)));
3387        static MIGRATION_FINDING_REGISTRY: StoreRegistry = {
3388            let mut registry = StoreRegistry::new();
3389            registry.register_journaled_store(
3390                MIGRATION_FINDING_STORE_PATH,
3391                &MIGRATION_FINDING_DATA,
3392                &MIGRATION_FINDING_INDEX,
3393                &MIGRATION_FINDING_SCHEMA,
3394                &MIGRATION_FINDING_JOURNAL,
3395                StoreAllocationIdentities::new_journaled(
3396                    StoreAllocationIdentity::new(206, "icydb.test.migration_finding.data.v1"),
3397                    StoreAllocationIdentity::new(207, "icydb.test.migration_finding.index.v1"),
3398                    StoreAllocationIdentity::new(208, "icydb.test.migration_finding.schema.v1"),
3399                    StoreAllocationIdentity::new(209, "icydb.test.migration_finding.journal.v1"),
3400                ),
3401                StoreRuntimeStorageCapabilities::journaled(),
3402            ).expect("migration finding store should register");
3403            registry
3404        };
3405    }
3406
3407    thread_local! {
3408        static EVOLUTION_DATA: RefCell<DataStore> =
3409            RefCell::new(DataStore::init_journaled(test_memory(192)));
3410        static EVOLUTION_INDEX: RefCell<IndexStore> =
3411            RefCell::new(IndexStore::init_journaled(test_memory(193)));
3412        static EVOLUTION_SCHEMA: RefCell<SchemaStore> =
3413            RefCell::new(SchemaStore::init_journaled(test_memory(194)));
3414        static EVOLUTION_JOURNAL: RefCell<JournalTailStore> =
3415            RefCell::new(JournalTailStore::init(test_memory(195)));
3416        static EVOLUTION_REGISTRY: StoreRegistry = {
3417            let mut registry = StoreRegistry::new();
3418            registry.register_journaled_store(
3419                EVOLUTION_STORE_PATH,
3420                &EVOLUTION_DATA,
3421                &EVOLUTION_INDEX,
3422                &EVOLUTION_SCHEMA,
3423                &EVOLUTION_JOURNAL,
3424                StoreAllocationIdentities::new_journaled(
3425                    StoreAllocationIdentity::new(192, "icydb.test.rule_evolution.data.v1"),
3426                    StoreAllocationIdentity::new(193, "icydb.test.rule_evolution.index.v1"),
3427                    StoreAllocationIdentity::new(194, "icydb.test.rule_evolution.schema.v1"),
3428                    StoreAllocationIdentity::new(195, "icydb.test.rule_evolution.journal.v1"),
3429                ),
3430                StoreRuntimeStorageCapabilities::journaled(),
3431            ).expect("rule-evolution journaled store should register");
3432            registry
3433        };
3434    }
3435
3436    struct AbortCanister;
3437
3438    impl Path for AbortCanister {
3439        const PATH: &'static str = "schema_application_tests::AbortCanister";
3440    }
3441
3442    impl CanisterKind for AbortCanister {
3443        const COMMIT_MEMORY_ID: u8 = 184;
3444        const COMMIT_STABLE_KEY: &'static str = "icydb.test.application_abort.commit.v1";
3445        const STARTUP_MEMORY_ID: u8 = 186;
3446        const STARTUP_STABLE_KEY: &'static str = "icydb.test.application_abort.startup.control.v1";
3447        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 185;
3448        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3449            "icydb.test.application_abort.integrity.v1";
3450    }
3451
3452    struct EvolutionCanister;
3453
3454    impl Path for EvolutionCanister {
3455        const PATH: &'static str = "schema_application_tests::EvolutionCanister";
3456    }
3457
3458    impl CanisterKind for EvolutionCanister {
3459        const COMMIT_MEMORY_ID: u8 = 196;
3460        const COMMIT_STABLE_KEY: &'static str = "icydb.test.rule_evolution.commit.v1";
3461        const STARTUP_MEMORY_ID: u8 = 198;
3462        const STARTUP_STABLE_KEY: &'static str = "icydb.test.rule_evolution.startup.control.v1";
3463        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 197;
3464        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3465            "icydb.test.rule_evolution.integrity.v1";
3466    }
3467
3468    #[cfg(feature = "migration")]
3469    struct MigrationCanister;
3470
3471    #[cfg(feature = "migration")]
3472    impl Path for MigrationCanister {
3473        const PATH: &'static str = "schema_application_tests::MigrationCanister";
3474    }
3475
3476    #[cfg(feature = "migration")]
3477    impl CanisterKind for MigrationCanister {
3478        const COMMIT_MEMORY_ID: u8 = 204;
3479        const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_validation.commit.v1";
3480        const STARTUP_MEMORY_ID: u8 = 206;
3481        const STARTUP_STABLE_KEY: &'static str =
3482            "icydb.test.migration_validation.startup.control.v1";
3483        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 205;
3484        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3485            "icydb.test.migration_validation.integrity.v1";
3486    }
3487
3488    #[cfg(feature = "migration")]
3489    struct MigrationExecutionCanister;
3490
3491    #[cfg(feature = "migration")]
3492    impl Path for MigrationExecutionCanister {
3493        const PATH: &'static str = "schema_application_tests::MigrationExecutionCanister";
3494    }
3495
3496    #[cfg(feature = "migration")]
3497    impl CanisterKind for MigrationExecutionCanister {
3498        const COMMIT_MEMORY_ID: u8 = 214;
3499        const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_execution.commit.v1";
3500        const STARTUP_MEMORY_ID: u8 = 216;
3501        const STARTUP_STABLE_KEY: &'static str =
3502            "icydb.test.migration_execution.startup.control.v1";
3503        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 215;
3504        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3505            "icydb.test.migration_execution.integrity.v1";
3506    }
3507
3508    #[cfg(feature = "migration")]
3509    struct MigrationFindingCanister;
3510
3511    #[cfg(feature = "migration")]
3512    impl Path for MigrationFindingCanister {
3513        const PATH: &'static str = "schema_application_tests::MigrationFindingCanister";
3514    }
3515
3516    #[cfg(feature = "migration")]
3517    impl CanisterKind for MigrationFindingCanister {
3518        const COMMIT_MEMORY_ID: u8 = 210;
3519        const COMMIT_STABLE_KEY: &'static str = "icydb.test.migration_finding.commit.v1";
3520        const STARTUP_MEMORY_ID: u8 = 212;
3521        const STARTUP_STABLE_KEY: &'static str = "icydb.test.migration_finding.startup.control.v1";
3522        const INTEGRITY_PROGRESS_MEMORY_ID: u8 = 211;
3523        const INTEGRITY_PROGRESS_STABLE_KEY: &'static str =
3524            "icydb.test.migration_finding.integrity.v1";
3525    }
3526
3527    fn name(value: &str) -> SchemaName {
3528        SchemaName::try_new(value).expect("test schema name should admit")
3529    }
3530
3531    fn generated_check_proposal(
3532        expected_head: ExpectedAcceptedHead,
3533        submission_key: &str,
3534        include_check: bool,
3535        database: TargetDatabaseIdentity,
3536        store: TargetStoreIdentity,
3537    ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3538        let entity_source = EntitySourceKey::try_new("Item").expect("entity source should admit");
3539        let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3540        let score_source = FieldSourceKey::try_new("score").expect("score source should admit");
3541        let check_source =
3542            ConstraintSourceKey::try_new("score_non_negative").expect("check source should admit");
3543        let check = SourceCheckExpr::try_new(vec![
3544            SourceCheckInstruction::Field(score_source),
3545            SourceCheckInstruction::Literal(ScalarLiteral::Int(0)),
3546            SourceCheckInstruction::GreaterThanOrEqual,
3547        ])
3548        .expect("check expression should admit");
3549        let constraints = include_check
3550            .then(|| ConstraintFragment::check(name("score_non_negative"), check))
3551            .into_iter()
3552            .collect();
3553        let entity = EntityFragment::try_new(
3554            name("Item"),
3555            version_one(),
3556            vec![
3557                FieldFragment::new(
3558                    name("id"),
3559                    FieldType::Scalar(ScalarType::Nat64),
3560                    false,
3561                    FieldInsertPolicy::Required,
3562                    None,
3563                ),
3564                FieldFragment::new(
3565                    name("score"),
3566                    FieldType::Scalar(ScalarType::Int64),
3567                    false,
3568                    FieldInsertPolicy::Required,
3569                    None,
3570                ),
3571            ],
3572            vec![id_source],
3573            Vec::new(),
3574            Vec::new(),
3575            constraints,
3576        )
3577        .expect("entity should admit");
3578        let proposal = SchemaProposal::try_compose(
3579            vec![SchemaCapability::ACCEPTED_CHECKS],
3580            database,
3581            SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3582            expected_head,
3583            vec![
3584                SchemaFragment::try_new(vec![entity], Vec::new())
3585                    .expect("schema fragment should admit"),
3586            ],
3587            vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3588            Vec::new(),
3589            None,
3590        )
3591        .expect("schema proposal should compose");
3592        (proposal, entity_source, check_source)
3593    }
3594
3595    #[cfg(feature = "migration")]
3596    #[derive(Clone, Copy, Eq, PartialEq)]
3597    enum ValidationMigrationShape {
3598        Clean,
3599        AllFindingFamilies,
3600    }
3601
3602    #[cfg(feature = "migration")]
3603    #[expect(
3604        clippy::too_many_lines,
3605        reason = "the fixture keeps both predecessor and candidate source contracts adjacent"
3606    )]
3607    fn validation_migration_proposal(
3608        shape: ValidationMigrationShape,
3609        current: bool,
3610        expected_head: ExpectedAcceptedHead,
3611        database: TargetDatabaseIdentity,
3612        store: TargetStoreIdentity,
3613    ) -> SchemaProposal {
3614        let entity_source = EntitySourceKey::try_new("MigratingItem")
3615            .expect("migration entity source should admit");
3616        let old_value =
3617            FieldSourceKey::try_new("old_value").expect("predecessor field source should admit");
3618        let current_value =
3619            FieldSourceKey::try_new("value").expect("candidate field source should admit");
3620        let target_entity = EntitySourceKey::try_new("MigrationTarget")
3621            .expect("migration target source should admit");
3622        let target_id = FieldSourceKey::try_new("id").expect("target id source should admit");
3623        let constraint = SourceCheckExpr::try_new(vec![
3624            SourceCheckInstruction::Field(current_value.clone()),
3625            SourceCheckInstruction::Literal(ScalarLiteral::Nat(8)),
3626            SourceCheckInstruction::LessThanOrEqual,
3627        ])
3628        .expect("candidate check should admit");
3629        let findings = shape == ValidationMigrationShape::AllFindingFamilies;
3630        let entity = EntityFragment::try_new(
3631            name("MigratingItem"),
3632            DeclaredEntityVersion::try_new(if current { 2 } else { 1 })
3633                .expect("migration version should admit"),
3634            vec![
3635                FieldFragment::new(
3636                    name("id"),
3637                    FieldType::Scalar(ScalarType::Nat64),
3638                    false,
3639                    FieldInsertPolicy::Required,
3640                    None,
3641                ),
3642                FieldFragment::new(
3643                    name(if current { "value" } else { "old_value" }),
3644                    FieldType::Scalar(if current {
3645                        ScalarType::Nat8
3646                    } else {
3647                        ScalarType::Int64
3648                    }),
3649                    false,
3650                    FieldInsertPolicy::Required,
3651                    None,
3652                ),
3653            ],
3654            vec![FieldSourceKey::try_new("id").expect("id source should admit")],
3655            current
3656                .then(|| {
3657                    IndexFragment::try_new(
3658                        name("value_unique"),
3659                        vec![IndexKeyFragment::Field(current_value.clone())],
3660                        true,
3661                        None,
3662                    )
3663                    .expect("candidate index should admit")
3664                })
3665                .into_iter()
3666                .collect(),
3667            (current && findings)
3668                .then(|| {
3669                    RelationFragment::try_new(
3670                        name("value_target"),
3671                        vec![current_value.clone()],
3672                        target_entity.clone(),
3673                        vec![target_id.clone()],
3674                        RelationDeleteAction::Restrict,
3675                    )
3676                    .expect("candidate relation should admit")
3677                })
3678                .into_iter()
3679                .collect(),
3680            (current && findings)
3681                .then(|| ConstraintFragment::check(name("value_at_most_eight"), constraint))
3682                .into_iter()
3683                .collect(),
3684        )
3685        .expect("migration entity should admit");
3686        let target = EntityFragment::try_new(
3687            name("MigrationTarget"),
3688            version_one(),
3689            vec![FieldFragment::new(
3690                name("id"),
3691                FieldType::Scalar(ScalarType::Nat8),
3692                false,
3693                FieldInsertPolicy::Required,
3694                None,
3695            )],
3696            vec![target_id],
3697            Vec::new(),
3698            Vec::new(),
3699            Vec::new(),
3700        )
3701        .expect("migration relation target should admit");
3702        let migration = current.then(|| {
3703            SchemaMigrationPlan::try_new(vec![
3704                EntityMigration::try_new(
3705                    entity_source.clone(),
3706                    DeclaredEntityVersion::try_new(1).expect("predecessor should admit"),
3707                    None,
3708                    Vec::new(),
3709                    vec![SchemaMigrationTransform::CheckedCast {
3710                        from: old_value.clone(),
3711                        to: current_value,
3712                        target: ScalarType::Nat8,
3713                    }],
3714                )
3715                .expect("migration transition should admit"),
3716            ])
3717            .expect("migration plan should admit")
3718        });
3719        let mut capabilities = Vec::new();
3720        if current && findings {
3721            capabilities.extend([
3722                SchemaCapability::ACCEPTED_CHECKS,
3723                SchemaCapability::SECONDARY_INDEXES,
3724                SchemaCapability::RESTRICTIVE_RELATIONS,
3725            ]);
3726        }
3727        if migration.is_some() {
3728            capabilities.push(SchemaCapability::VERSIONED_MIGRATIONS);
3729        }
3730        let mut entities = vec![entity];
3731        let mut assignments = vec![EntityStoreAssignment::new(entity_source.clone(), store)];
3732        if findings {
3733            entities.push(target);
3734            assignments.push(EntityStoreAssignment::new(target_entity, store));
3735        }
3736        SchemaProposal::try_compose(
3737            capabilities,
3738            database,
3739            SchemaSubmissionKey::try_new(if current {
3740                "migration-validation-v2"
3741            } else {
3742                "migration-validation-v1"
3743            })
3744            .expect("submission should admit"),
3745            expected_head,
3746            vec![
3747                SchemaFragment::try_new(entities, Vec::new())
3748                    .expect("migration fragment should admit"),
3749            ],
3750            assignments,
3751            current
3752                .then_some(icydb_schema::SchemaRemoval::Field {
3753                    entity: entity_source,
3754                    field: old_value,
3755                })
3756                .into_iter()
3757                .collect(),
3758            migration,
3759        )
3760        .expect("migration proposal should compose")
3761    }
3762
3763    fn targeted_rule_proposal(
3764        expected_head: ExpectedAcceptedHead,
3765        submission_key: &str,
3766        operation: SourceRuleOperation,
3767        database: TargetDatabaseIdentity,
3768        store: TargetStoreIdentity,
3769    ) -> (SchemaProposal, EntitySourceKey, ConstraintSourceKey) {
3770        let entity_source =
3771            EntitySourceKey::try_new("Measured").expect("entity source should admit");
3772        let id_source = FieldSourceKey::try_new("id").expect("id source should admit");
3773        let value_source = FieldSourceKey::try_new("value").expect("value source should admit");
3774        let value_type = TypeSourceKey::try_new("Measure").expect("type source should admit");
3775        let rule_source = RuleSourceKey::try_new("limit").expect("rule source should admit");
3776        let constraint_source =
3777            ConstraintSourceKey::for_targeted_field_rule(&value_source, &value_type, &rule_source);
3778        let entity = EntityFragment::try_new(
3779            name("Measured"),
3780            version_one(),
3781            vec![
3782                FieldFragment::new(
3783                    name("id"),
3784                    FieldType::Scalar(ScalarType::Nat64),
3785                    false,
3786                    FieldInsertPolicy::Required,
3787                    None,
3788                ),
3789                FieldFragment::new(
3790                    name("value"),
3791                    FieldType::Named(value_type.clone()),
3792                    false,
3793                    FieldInsertPolicy::Required,
3794                    None,
3795                ),
3796            ],
3797            vec![id_source],
3798            Vec::new(),
3799            Vec::new(),
3800            vec![ConstraintFragment::targeted_rule(
3801                TargetedRuleFragment::new(value_source, value_type, name("limit"), operation),
3802            )],
3803        )
3804        .expect("targeted entity should admit");
3805        let proposal = SchemaProposal::try_compose(
3806            vec![SchemaCapability::ACCEPTED_CHECKS],
3807            database,
3808            SchemaSubmissionKey::try_new(submission_key).expect("submission key should admit"),
3809            expected_head,
3810            vec![
3811                SchemaFragment::try_new(
3812                    vec![entity],
3813                    vec![NamedTypeFragment::newtype(
3814                        name("Measure"),
3815                        FieldType::Scalar(ScalarType::Nat8),
3816                    )],
3817                )
3818                .expect("schema fragment should admit"),
3819            ],
3820            vec![EntityStoreAssignment::new(entity_source.clone(), store)],
3821            Vec::new(),
3822            None,
3823        )
3824        .expect("schema proposal should compose");
3825        (proposal, entity_source, constraint_source)
3826    }
3827
3828    #[test]
3829    fn database_head_is_empty_only_when_every_store_root_is_absent() {
3830        assert_eq!(
3831            derive_accepted_head(&[("test::A", None), ("test::B", None)]),
3832            ExpectedAcceptedHead::Empty,
3833        );
3834    }
3835
3836    #[test]
3837    fn database_head_covers_store_path_revision_fingerprint_and_absence() {
3838        let first = derive_accepted_head(&[
3839            (
3840                "test::A",
3841                Some(AcceptedStoreHead {
3842                    revision: 3,
3843                    fingerprint: [0x11; 32],
3844                }),
3845            ),
3846            ("test::B", None),
3847        ]);
3848        let changed_fingerprint = derive_accepted_head(&[
3849            (
3850                "test::A",
3851                Some(AcceptedStoreHead {
3852                    revision: 3,
3853                    fingerprint: [0x12; 32],
3854                }),
3855            ),
3856            ("test::B", None),
3857        ]);
3858        let changed_absence = derive_accepted_head(&[
3859            (
3860                "test::A",
3861                Some(AcceptedStoreHead {
3862                    revision: 3,
3863                    fingerprint: [0x11; 32],
3864                }),
3865            ),
3866            (
3867                "test::B",
3868                Some(AcceptedStoreHead {
3869                    revision: 1,
3870                    fingerprint: [0x22; 32],
3871                }),
3872            ),
3873        ]);
3874
3875        assert_ne!(first, changed_fingerprint);
3876        assert_ne!(first, changed_absence);
3877        assert!(matches!(
3878            first,
3879            ExpectedAcceptedHead::Exact { revision: 3, .. }
3880        ));
3881    }
3882
3883    #[test]
3884    #[allow(
3885        clippy::too_many_lines,
3886        reason = "the end-to-end catalog assertion is clearer as one lifecycle test"
3887    )]
3888    fn generated_check_abort_retires_source_identity_and_allows_fresh_reproposal() {
3889        let database = TargetDatabaseIdentity::from_bytes([0x71; 32]);
3890        let store = TargetStoreIdentity::from_bytes([0x72; 32]);
3891        let (initial, entity_source, _) = generated_check_proposal(
3892            ExpectedAcceptedHead::Empty,
3893            "abort-initial",
3894            false,
3895            database,
3896            store,
3897        );
3898        let initial_candidate = lower_initial_schema_proposal(
3899            &initial,
3900            &[ProposalStoreTarget {
3901                path: "abort::Store",
3902                identity: store,
3903            }],
3904        )
3905        .expect("initial proposal should lower")
3906        .pop()
3907        .expect("initial proposal should produce one candidate");
3908        let (with_check, _, check_source) = generated_check_proposal(
3909            ExpectedAcceptedHead::Exact {
3910                revision: 1,
3911                fingerprint: ExpectedSchemaFingerprint::from_bytes([0x73; 32]),
3912            },
3913            "abort-add-check",
3914            true,
3915            database,
3916            store,
3917        );
3918        let pending_candidate = lower_existing_schema_proposal(
3919            &with_check,
3920            &[ExistingProposalStore {
3921                path: "abort::Store",
3922                identity: store,
3923                bundle: initial_candidate.bundle(),
3924            }],
3925        )
3926        .expect("generated check should lower")
3927        .pop()
3928        .expect("generated check should produce one candidate");
3929        let entity_tag = pending_candidate
3930            .bundle()
3931            .source_bindings_for_tests()
3932            .entity(&entity_source)
3933            .expect("entity source should remain bound");
3934        let constraint_id = pending_candidate
3935            .bundle()
3936            .source_bindings_for_tests()
3937            .constraint(entity_tag, &check_source)
3938            .expect("generated check source should bind");
3939        let pending_snapshot = pending_candidate
3940            .bundle()
3941            .entity_snapshots()
3942            .get(&entity_tag)
3943            .expect("pending entity should exist");
3944        let activation = pending_snapshot
3945            .constraint_catalog()
3946            .activation(constraint_id)
3947            .expect("generated check should remain an activation");
3948        assert_eq!(activation.origin(), ConstraintOrigin::Generated);
3949
3950        let aborted = aborted_generated_row_local_candidate(
3951            pending_candidate.bundle(),
3952            entity_tag,
3953            constraint_id,
3954        )
3955        .expect("generated check abort should build one catalog-native candidate");
3956        let aborted_snapshot = aborted
3957            .bundle()
3958            .entity_snapshots()
3959            .get(&entity_tag)
3960            .expect("aborted entity should remain");
3961        assert!(
3962            aborted_snapshot
3963                .constraint_catalog()
3964                .activation(constraint_id)
3965                .is_none(),
3966        );
3967        assert_eq!(aborted_snapshot.row_layout(), pending_snapshot.row_layout());
3968        assert!(
3969            aborted
3970                .bundle()
3971                .source_bindings_for_tests()
3972                .constraint(entity_tag, &check_source)
3973                .is_none(),
3974        );
3975
3976        let reproposed = lower_existing_schema_proposal(
3977            &with_check,
3978            &[ExistingProposalStore {
3979                path: "abort::Store",
3980                identity: store,
3981                bundle: aborted.bundle(),
3982            }],
3983        )
3984        .expect("aborted generated check should be independently reproposable")
3985        .pop()
3986        .expect("reproposal should produce one candidate");
3987        let replacement_id = reproposed
3988            .bundle()
3989            .source_bindings_for_tests()
3990            .constraint(entity_tag, &check_source)
3991            .expect("reproposal should bind a fresh constraint identity");
3992        assert!(
3993            replacement_id > constraint_id,
3994            "aborted accepted IDs must remain retired",
3995        );
3996    }
3997
3998    #[test]
3999    fn targeted_rule_edit_abort_keeps_prior_accepted_semantics_and_source_identity() {
4000        let database = TargetDatabaseIdentity::from_bytes([0x81; 32]);
4001        let store = TargetStoreIdentity::from_bytes([0x82; 32]);
4002        let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4003            ExpectedAcceptedHead::Empty,
4004            "targeted-abort-initial",
4005            SourceRuleOperation::NumericRangeInclusive {
4006                min: ScalarLiteral::Nat(0),
4007                max: ScalarLiteral::Nat(10),
4008            },
4009            database,
4010            store,
4011        );
4012        let initial_candidate = lower_initial_schema_proposal(
4013            &initial,
4014            &[ProposalStoreTarget {
4015                path: "abort::TargetedStore",
4016                identity: store,
4017            }],
4018        )
4019        .expect("initial targeted proposal should lower")
4020        .pop()
4021        .expect("initial targeted proposal should produce one candidate");
4022        let initial_bundle = initial_candidate.bundle();
4023        let entity_tag = initial_bundle
4024            .source_bindings_for_tests()
4025            .entity(&entity_source)
4026            .expect("entity source should bind");
4027        let constraint_id = initial_bundle
4028            .source_bindings_for_tests()
4029            .constraint(entity_tag, &constraint_source)
4030            .expect("targeted source should bind");
4031        let high_water = initial_bundle.entity_snapshots()[&entity_tag]
4032            .constraint_id_allocator()
4033            .high_water();
4034        let (edited, _, _) = targeted_rule_proposal(
4035            ExpectedAcceptedHead::Exact {
4036                revision: initial_bundle.revision().get(),
4037                fingerprint: ExpectedSchemaFingerprint::from_bytes([0x83; 32]),
4038            },
4039            "targeted-abort-edit",
4040            SourceRuleOperation::NumericMaximumInclusive {
4041                value: ScalarLiteral::Nat(8),
4042            },
4043            database,
4044            store,
4045        );
4046        let staged = lower_existing_schema_proposal(
4047            &edited,
4048            &[ExistingProposalStore {
4049                path: "abort::TargetedStore",
4050                identity: store,
4051                bundle: initial_bundle,
4052            }],
4053        )
4054        .expect("targeted semantic edit should stage")
4055        .pop()
4056        .expect("targeted semantic edit should produce one candidate");
4057        let aborted =
4058            aborted_generated_row_local_candidate(staged.bundle(), entity_tag, constraint_id)
4059                .expect("targeted semantic edit should abort through catalog authority");
4060        let snapshot = &aborted.bundle().entity_snapshots()[&entity_tag];
4061
4062        assert!(
4063            snapshot
4064                .constraint_catalog()
4065                .activation(constraint_id)
4066                .is_none()
4067        );
4068        assert_eq!(snapshot.constraint_id_allocator().high_water(), high_water);
4069        assert_eq!(
4070            aborted
4071                .bundle()
4072                .source_bindings_for_tests()
4073                .constraint(entity_tag, &constraint_source),
4074            Some(constraint_id),
4075        );
4076        assert!(snapshot.constraints().iter().any(|constraint| {
4077            constraint.id() == constraint_id
4078                && matches!(
4079                    constraint.kind(),
4080                    AcceptedConstraintKind::TargetedRule { operation, .. }
4081                        if matches!(
4082                            operation.as_ref(),
4083                            AcceptedRuleOperation::NumericRangeInclusive { .. }
4084                        )
4085                )
4086        }));
4087    }
4088
4089    #[test]
4090    #[allow(
4091        clippy::too_many_lines,
4092        reason = "the staged publication, recovery, and promotion assertions form one lifecycle"
4093    )]
4094    fn generated_startup_driver_resumes_pending_activation_and_promotes_without_source_model() {
4095        let db = Db::<EvolutionCanister>::new(
4096            &EVOLUTION_REGISTRY,
4097            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4098        );
4099        drive_startup_recovery_to_completion(&db);
4100        let empty_target =
4101            schema_application_target(&db).expect("empty evolution target should issue");
4102        let store_identity = empty_target
4103            .stores()
4104            .first()
4105            .expect("evolution store should register")
4106            .identity();
4107        let (initial, entity_source, constraint_source) = targeted_rule_proposal(
4108            empty_target.accepted_head().clone(),
4109            "targeted-recovery-initial",
4110            SourceRuleOperation::NumericRangeInclusive {
4111                min: ScalarLiteral::Nat(0),
4112                max: ScalarLiteral::Nat(10),
4113            },
4114            empty_target.database_identity(),
4115            store_identity,
4116        );
4117        assert!(matches!(
4118            apply_schema(&db, &initial)
4119                .expect("initial targeted proposal should publish")
4120                .outcome(),
4121            SchemaChangeOutcome::Applied { .. },
4122        ));
4123        drive_startup_recovery_to_completion(&db);
4124        drive_cardinality_to_ready(
4125            db.store_handle(EVOLUTION_STORE_PATH)
4126                .expect("evolution store should resolve"),
4127        );
4128
4129        let direct_target =
4130            schema_application_target(&db).expect("direct evolution target should issue");
4131        let (direct_edit, _, _) = targeted_rule_proposal(
4132            direct_target.accepted_head().clone(),
4133            "targeted-direct-edit",
4134            SourceRuleOperation::NumericMaximumInclusive {
4135                value: ScalarLiteral::Nat(8),
4136            },
4137            direct_target.database_identity(),
4138            store_identity,
4139        );
4140        assert!(matches!(
4141            apply_schema(&db, &direct_edit)
4142                .expect("empty-domain semantic edit should publish directly")
4143                .outcome(),
4144            SchemaChangeOutcome::Applied { .. },
4145        ));
4146        let store = db
4147            .store_handle(EVOLUTION_STORE_PATH)
4148            .expect("evolution store should resolve");
4149        let direct = store
4150            .with_schema(SchemaStore::current_accepted_schema_bundle)
4151            .expect("directly edited bundle should remain readable")
4152            .expect("directly edited bundle should exist");
4153        let entity_tag = direct
4154            .source_bindings_for_tests()
4155            .entity(&entity_source)
4156            .expect("entity source should remain bound");
4157        let constraint_id = direct
4158            .source_bindings_for_tests()
4159            .constraint(entity_tag, &constraint_source)
4160            .expect("direct edit should preserve constraint identity");
4161        assert!(
4162            direct.entity_snapshots()[&entity_tag]
4163                .constraint_catalog()
4164                .activation(constraint_id)
4165                .is_none()
4166        );
4167        assert!(
4168            direct.entity_snapshots()[&entity_tag]
4169                .constraints()
4170                .iter()
4171                .any(|constraint| {
4172                    constraint.id() == constraint_id
4173                        && matches!(
4174                            constraint.kind(),
4175                            AcceptedConstraintKind::TargetedRule { operation, .. }
4176                                if matches!(
4177                                    operation.as_ref(),
4178                                    AcceptedRuleOperation::NumericMaximumInclusive { .. }
4179                                )
4180                        )
4181                })
4182        );
4183
4184        let target = schema_application_target(&db).expect("staged evolution target should issue");
4185        let (edited, _, _) = targeted_rule_proposal(
4186            target.accepted_head().clone(),
4187            "targeted-recovery-edit",
4188            SourceRuleOperation::MultipleOf {
4189                divisor: ScalarLiteral::Nat(2),
4190            },
4191            target.database_identity(),
4192            store_identity,
4193        );
4194        let current = store
4195            .with_schema(SchemaStore::current_accepted_schema_bundle)
4196            .expect("accepted evolution bundle should remain readable")
4197            .expect("directly edited evolution bundle should exist");
4198        let staged = lower_existing_schema_proposal(
4199            &edited,
4200            &[ExistingProposalStore {
4201                path: EVOLUTION_STORE_PATH,
4202                identity: store_identity,
4203                bundle: &current,
4204            }],
4205        )
4206        .expect("targeted edit should stage")
4207        .pop()
4208        .expect("targeted edit should produce one candidate");
4209        assert_eq!(
4210            staged
4211                .bundle()
4212                .source_bindings_for_tests()
4213                .constraint(entity_tag, &constraint_source),
4214            Some(constraint_id),
4215        );
4216        let proof = DirectGeneratedRowLocalProof {
4217            candidate_index: 0,
4218            store,
4219            store_path: EVOLUTION_STORE_PATH,
4220            entity_tag,
4221            entity_path: staged.bundle().entity_snapshots()[&entity_tag]
4222                .entity_path()
4223                .to_string(),
4224            constraint_id,
4225            historical_rows: 0,
4226        };
4227        let final_candidates = final_candidates_for_pending_row_local_constraint(
4228            std::slice::from_ref(&staged),
4229            &PendingGeneratedRowLocalConstraint { proof },
4230        )
4231        .expect("final semantic replacement should derive without source input");
4232        let authorities = application_authorities(&db);
4233        let candidate_head =
4234            accepted_head_after_candidates(authorities.as_slice(), &final_candidates)
4235                .expect("final candidate head should derive");
4236        let digest = edited.digest().expect("proposal digest should derive");
4237        let job_id = derive_schema_change_job_id(
4238            target.database_identity(),
4239            edited.submission_key(),
4240            digest,
4241            target.accepted_head(),
4242        )
4243        .expect("job identity should derive");
4244        let receipt = crate::db::schema::SchemaChangeReceipt::new(
4245            target.database_identity(),
4246            edited.submission_key().clone(),
4247            digest,
4248            target.accepted_head().clone(),
4249            SchemaChangeOutcome::Pending {
4250                job: SchemaChangeJob::new(job_id),
4251                candidate_head,
4252            },
4253        )
4254        .expect("pending replacement receipt should admit");
4255        let record = SchemaApplicationRecord::new(
4256            receipt,
4257            vec![
4258                SchemaChangeActivation::new(
4259                    store_identity,
4260                    entity_tag.value(),
4261                    constraint_id.get(),
4262                )
4263                .expect("replacement activation should admit"),
4264            ],
4265        )
4266        .expect("pending replacement record should admit");
4267        let operation =
4268            SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4269        publish_accepted_schema_candidates_with_application_record(
4270            vec![AcceptedSchemaPublication::new(
4271                EVOLUTION_STORE_PATH,
4272                store,
4273                current.revision(),
4274                &staged,
4275            )],
4276            operation,
4277        )
4278        .expect("staged replacement and record should publish atomically");
4279
4280        forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4281        drive_startup_recovery_to_completion(&db);
4282
4283        let recovered = store
4284            .with_schema(SchemaStore::current_accepted_schema_bundle)
4285            .expect("recovered staged bundle should decode")
4286            .expect("recovered staged bundle should exist");
4287        let recovered_snapshot = recovered.entity_snapshots()[&entity_tag].clone();
4288        let validating_catalog = recovered_snapshot
4289            .constraint_catalog()
4290            .clone()
4291            .with_validation_started(constraint_id)
4292            .expect("recovered replacement should enter validation");
4293        let mut validating_snapshots = recovered.entity_snapshots().clone();
4294        validating_snapshots.insert(
4295            entity_tag,
4296            recovered_snapshot.with_constraint_catalog(validating_catalog),
4297        );
4298        let validating_bundle = AcceptedSchemaRevisionBundle::new_with_source_bindings(
4299            recovered
4300                .revision()
4301                .checked_next()
4302                .expect("validation revision should remain available"),
4303            recovered.store_path(),
4304            recovered.enum_catalog().clone(),
4305            recovered.composite_catalog().clone(),
4306            recovered.source_bindings_for_tests().clone(),
4307            validating_snapshots,
4308        )
4309        .expect("validating replacement bundle should close");
4310        let validating_candidate = CandidateSchemaRevision::new(validating_bundle)
4311            .expect("validating replacement candidate should encode");
4312        let validating_activation = validating_candidate.bundle().entity_snapshots()[&entity_tag]
4313            .constraint_catalog()
4314            .activation(constraint_id)
4315            .expect("validating replacement activation should remain present");
4316        let validation_job = ConstraintValidationJob::start(
4317            entity_tag,
4318            validating_candidate.bundle().entity_snapshots()[&entity_tag]
4319                .entity_path()
4320                .to_string(),
4321            validating_activation,
4322            None,
4323        )
4324        .expect("validating replacement job should derive from accepted state");
4325        store
4326            .with_schema(|schema| {
4327                schema.validate_live_activation_transition(validating_candidate.bundle())?;
4328                schema.validate_constraint_validation_job_closure_with_change(
4329                    validating_candidate.bundle(),
4330                    Some(&validation_job),
4331                    None,
4332                )
4333            })
4334            .expect("validating replacement transition and job should close");
4335        let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4336        let startup_session =
4337            crate::db::DbSession::<EvolutionCanister>::new(&EVOLUTION_REGISTRY, &startup_root);
4338
4339        assert_eq!(
4340            drive_generated_startup_recovery_page(
4341                &startup_session,
4342                &EVOLUTION_REGISTRY,
4343                edited.submission_key().as_str(),
4344            )
4345            .expect("generated startup should begin pending validation"),
4346            GeneratedStartupDriverStep::Recovering,
4347        );
4348        let mut terminal = false;
4349        for _ in 0..8 {
4350            match drive_generated_startup_recovery_page(
4351                &startup_session,
4352                &EVOLUTION_REGISTRY,
4353                edited.submission_key().as_str(),
4354            )
4355            .expect("generated startup should advance pending validation")
4356            {
4357                GeneratedStartupDriverStep::Recovering => {}
4358                GeneratedStartupDriverStep::Terminal => {
4359                    terminal = true;
4360                    break;
4361                }
4362                GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4363                    panic!("an exact pending receipt must resume instead of being resubmitted")
4364                }
4365            }
4366        }
4367        assert!(
4368            terminal,
4369            "empty historical domain should promote within bounded startup steps"
4370        );
4371        assert_eq!(
4372            observe_generated_startup_state::<EvolutionCanister>(
4373                &EVOLUTION_REGISTRY,
4374                edited.submission_key().as_str(),
4375            ),
4376            Ok(DatabaseStartupState::Ready),
4377        );
4378        let applied = super::exact_schema_application_receipt(
4379            &edited,
4380            edited
4381                .digest()
4382                .expect("proposal digest should remain stable"),
4383        )
4384        .expect("terminal generated receipt should remain readable")
4385        .expect("terminal generated receipt should remain present");
4386        assert!(matches!(
4387            applied.outcome(),
4388            SchemaChangeOutcome::Applied { .. }
4389        ));
4390        let promoted = store
4391            .with_schema(SchemaStore::current_accepted_schema_bundle)
4392            .expect("promoted bundle should remain readable")
4393            .expect("promoted bundle should exist");
4394        let snapshot = &promoted.entity_snapshots()[&entity_tag];
4395        assert!(
4396            snapshot
4397                .constraint_catalog()
4398                .activation(constraint_id)
4399                .is_none()
4400        );
4401        assert_eq!(
4402            promoted
4403                .source_bindings_for_tests()
4404                .constraint(entity_tag, &constraint_source),
4405            Some(constraint_id),
4406        );
4407        assert!(snapshot.constraints().iter().any(|constraint| {
4408            constraint.id() == constraint_id
4409                && matches!(
4410                    constraint.kind(),
4411                    AcceptedConstraintKind::TargetedRule { operation, .. }
4412                        if matches!(
4413                            operation.as_ref(),
4414                            AcceptedRuleOperation::MultipleOf { .. }
4415                        )
4416                )
4417        }));
4418    }
4419
4420    #[test]
4421    #[allow(
4422        clippy::too_many_lines,
4423        reason = "the durable pending job, startup failure, and retained finding assertions form one scenario"
4424    )]
4425    fn generated_startup_driver_persists_e223_for_a_retained_historical_finding() {
4426        let db = Db::<AbortCanister>::new(
4427            &ABORT_REGISTRY,
4428            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4429        );
4430        drive_startup_recovery_to_completion(&db);
4431        let empty_target =
4432            schema_application_target(&db).expect("empty application target should issue");
4433        let store_identity = empty_target
4434            .stores()
4435            .first()
4436            .expect("abort store should be registered")
4437            .identity();
4438        let (initial, _, _) = generated_check_proposal(
4439            empty_target.accepted_head().clone(),
4440            "startup-finding-initial",
4441            false,
4442            empty_target.database_identity(),
4443            store_identity,
4444        );
4445        apply_schema(&db, &initial).expect("initial generated schema should publish");
4446
4447        let root = crate::db::RequestExecutionRoot::__new_runtime_root();
4448        let session = DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &root);
4449        let rows = (1..=257)
4450            .map(|id| DynamicMutation::Insert {
4451                entity: "Item".to_string(),
4452                patch: DynamicStructuralPatch::new(vec![
4453                    (
4454                        "id".to_string(),
4455                        DynamicWriteCell::Value(InputValue::nat64(id)),
4456                    ),
4457                    (
4458                        "score".to_string(),
4459                        DynamicWriteCell::Value(InputValue::int64(if id == 257 { -1 } else { 1 })),
4460                    ),
4461                ]),
4462            })
4463            .collect();
4464        session
4465            .execute_trusted_dynamic_mutation_batch(rows)
4466            .expect("historical finding fixture rows should commit as one legal batch");
4467        drive_startup_recovery_to_completion(&db);
4468        drive_cardinality_to_ready(
4469            db.store_handle(ABORT_STORE_PATH)
4470                .expect("abort store should resolve"),
4471        );
4472
4473        let target =
4474            schema_application_target(&db).expect("existing application target should issue");
4475        let (with_check, _, _) = generated_check_proposal(
4476            target.accepted_head().clone(),
4477            "startup-finding-pending",
4478            true,
4479            target.database_identity(),
4480            store_identity,
4481        );
4482        let pending = apply_schema(&db, &with_check)
4483            .expect("the first clean page should admit durable continuation");
4484        let SchemaChangeOutcome::Pending { job, .. } = pending.outcome() else {
4485            panic!("a 257-row domain must exceed the 256-row direct proof page")
4486        };
4487
4488        let mut terminal = false;
4489        for _ in 0..8 {
4490            match drive_generated_startup_recovery_page(
4491                &session,
4492                &ABORT_REGISTRY,
4493                with_check.submission_key().as_str(),
4494            )
4495            .expect("generated startup should retain a typed finding failure")
4496            {
4497                GeneratedStartupDriverStep::Recovering => {}
4498                GeneratedStartupDriverStep::Terminal => {
4499                    terminal = true;
4500                    break;
4501                }
4502                GeneratedStartupDriverStep::ApplyGeneratedSchema => {
4503                    panic!("an exact pending receipt must not be resubmitted")
4504                }
4505            }
4506        }
4507        assert!(terminal, "the retained finding should become terminal");
4508        let failure = observe_generated_startup_state::<AbortCanister>(
4509            &ABORT_REGISTRY,
4510            with_check.submission_key().as_str(),
4511        )
4512        .expect_err("the retained finding must remain durably observable");
4513        assert_eq!(
4514            failure.kind(),
4515            crate::db::StartupFailureKind::SchemaReconciliation,
4516        );
4517        assert_eq!(
4518            failure.diagnostic().error_code(),
4519            icydb_diagnostic_code::ErrorCode::RUNTIME_BOUNDARY_CONSTRAINT_VIOLATION,
4520        );
4521        assert_eq!(
4522            ABORT_DATA.with(|store| store.borrow().len()),
4523            257,
4524            "terminal startup publication must not change historical rows",
4525        );
4526        assert!(matches!(
4527            continue_schema_application(&db, job.id(), None)
4528                .expect("the retained finding page should replay exactly")
4529                .status(),
4530            SchemaChangeProgressStatus::Findings { findings, .. } if !findings.is_empty(),
4531        ));
4532    }
4533
4534    #[test]
4535    #[allow(
4536        clippy::too_many_lines,
4537        reason = "the journaled abort, replay, and recovery assertions form one scenario"
4538    )]
4539    fn pending_generated_check_abort_is_atomic_terminal_and_replayable() {
4540        let db = Db::<AbortCanister>::new(
4541            &ABORT_REGISTRY,
4542            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4543        );
4544        drive_startup_recovery_to_completion(&db);
4545        let empty_target =
4546            schema_application_target(&db).expect("empty application target should issue");
4547        let store_identity = empty_target
4548            .stores()
4549            .first()
4550            .expect("abort store should be registered")
4551            .identity();
4552        let (initial, entity_source, _) = generated_check_proposal(
4553            empty_target.accepted_head().clone(),
4554            "abort-runtime-initial",
4555            false,
4556            empty_target.database_identity(),
4557            store_identity,
4558        );
4559        assert!(matches!(
4560            apply_schema(&db, &initial)
4561                .expect("initial application should publish")
4562                .outcome(),
4563            SchemaChangeOutcome::Applied { .. },
4564        ));
4565
4566        let target =
4567            schema_application_target(&db).expect("existing application target should issue");
4568        let (with_check, _, check_source) = generated_check_proposal(
4569            target.accepted_head().clone(),
4570            "abort-runtime-pending",
4571            true,
4572            target.database_identity(),
4573            store_identity,
4574        );
4575        let store = db
4576            .store_handle(ABORT_STORE_PATH)
4577            .expect("abort store should resolve");
4578        let current = store
4579            .with_schema(SchemaStore::current_accepted_schema_bundle)
4580            .expect("accepted bundle should remain readable")
4581            .expect("initial accepted bundle should exist");
4582        let pending_candidate = lower_existing_schema_proposal(
4583            &with_check,
4584            &[ExistingProposalStore {
4585                path: ABORT_STORE_PATH,
4586                identity: store_identity,
4587                bundle: &current,
4588            }],
4589        )
4590        .expect("pending generated check should lower")
4591        .pop()
4592        .expect("pending generated check should produce one candidate");
4593        let entity_tag = pending_candidate
4594            .bundle()
4595            .source_bindings_for_tests()
4596            .entity(&entity_source)
4597            .expect("entity source should bind");
4598        let constraint_id = pending_candidate
4599            .bundle()
4600            .source_bindings_for_tests()
4601            .constraint(entity_tag, &check_source)
4602            .expect("generated check source should bind");
4603        let digest = with_check.digest().expect("proposal digest should derive");
4604        let job_id = derive_schema_change_job_id(
4605            target.database_identity(),
4606            with_check.submission_key(),
4607            digest,
4608            target.accepted_head(),
4609        )
4610        .expect("job identity should derive");
4611        let receipt = crate::db::schema::SchemaChangeReceipt::new(
4612            target.database_identity(),
4613            with_check.submission_key().clone(),
4614            digest,
4615            target.accepted_head().clone(),
4616            SchemaChangeOutcome::Pending {
4617                job: SchemaChangeJob::new(job_id),
4618                candidate_head: ExpectedAcceptedHead::Exact {
4619                    revision: pending_candidate.revision().get().saturating_add(2),
4620                    fingerprint: ExpectedSchemaFingerprint::from_bytes([0x76; 32]),
4621                },
4622            },
4623        )
4624        .expect("pending receipt should admit");
4625        let record = SchemaApplicationRecord::new(
4626            receipt,
4627            vec![
4628                SchemaChangeActivation::new(
4629                    store_identity,
4630                    entity_tag.value(),
4631                    constraint_id.get(),
4632                )
4633                .expect("application activation should admit"),
4634            ],
4635        )
4636        .expect("pending application record should admit");
4637        let operation =
4638            SchemaApplicationRecordOp::insert(&record).expect("pending insert should prepare");
4639        publish_accepted_schema_candidates_with_application_record(
4640            vec![AcceptedSchemaPublication::new(
4641                ABORT_STORE_PATH,
4642                store,
4643                current.revision(),
4644                &pending_candidate,
4645            )],
4646            operation,
4647        )
4648        .expect("pending candidate and record should publish atomically");
4649
4650        let started = continue_schema_application(&db, job_id, None)
4651            .expect("first continuation should durably start validation");
4652        assert_eq!(started.status(), &SchemaChangeProgressStatus::Started);
4653        let progress =
4654            abort_schema_application(&db, job_id, None).expect("pending application should abort");
4655        assert_eq!(progress.status(), &SchemaChangeProgressStatus::Aborted);
4656        assert!(matches!(
4657            progress.receipt().outcome(),
4658            SchemaChangeOutcome::Aborted { .. },
4659        ));
4660        let replay =
4661            abort_schema_application(&db, job_id, None).expect("terminal abort should replay");
4662        assert_eq!(replay, progress);
4663        assert_eq!(
4664            continue_schema_application(&db, job_id, None)
4665                .expect("continuation after abort should replay terminal state"),
4666            progress,
4667        );
4668
4669        let aborted = store
4670            .with_schema(SchemaStore::current_accepted_schema_bundle)
4671            .expect("accepted bundle should remain readable")
4672            .expect("aborted accepted bundle should exist");
4673        assert!(
4674            aborted
4675                .entity_snapshots()
4676                .get(&entity_tag)
4677                .expect("entity should remain after abort")
4678                .constraint_catalog()
4679                .activation(constraint_id)
4680                .is_none(),
4681        );
4682        assert!(
4683            aborted
4684                .source_bindings_for_tests()
4685                .constraint(entity_tag, &check_source)
4686                .is_none(),
4687        );
4688        assert!(
4689            store
4690                .with_schema(|schema| {
4691                    schema.constraint_validation_job(entity_tag, constraint_id)
4692                })
4693                .expect("validation-job storage should remain readable")
4694                .is_none(),
4695        );
4696
4697        ABORT_DATA.with(|store| {
4698            ABORT_DATA_MEMORY.with(|memory| {
4699                *store.borrow_mut() = DataStore::init_journaled(memory.clone());
4700            });
4701        });
4702        ABORT_INDEX.with(|store| {
4703            ABORT_INDEX_MEMORY.with(|memory| {
4704                *store.borrow_mut() = IndexStore::init_journaled(memory.clone());
4705            });
4706        });
4707        ABORT_SCHEMA.with(|store| {
4708            ABORT_SCHEMA_MEMORY.with(|memory| {
4709                *store.borrow_mut() = SchemaStore::init_journaled(memory.clone());
4710            });
4711        });
4712        ABORT_JOURNAL.with(|store| {
4713            ABORT_JOURNAL_MEMORY.with(|memory| {
4714                *store.borrow_mut() = JournalTailStore::init(memory.clone());
4715            });
4716        });
4717        forget_recovered_domain_for_tests(&db).expect("upgrade should reset recovery ownership");
4718        drive_startup_recovery_to_completion(&db);
4719        assert_eq!(
4720            abort_schema_application(&db, job_id, None)
4721                .expect("recovered terminal abort should replay"),
4722            progress,
4723        );
4724        assert!(
4725            store
4726                .with_schema(|schema| {
4727                    schema.constraint_validation_job(entity_tag, constraint_id)
4728                })
4729                .expect("recovered validation-job storage should remain readable")
4730                .is_none(),
4731        );
4732        let startup_root = crate::db::RequestExecutionRoot::__new_runtime_root();
4733        let startup_session =
4734            crate::db::DbSession::<AbortCanister>::new(&ABORT_REGISTRY, &startup_root);
4735        ABORT_JOURNAL.with(|journal| {
4736            let journal = journal.borrow();
4737            assert!(
4738                journal
4739                    .validate_current_tail_authority()
4740                    .expect("recovered abort tail control should remain readable")
4741                    .is_empty(),
4742                "recovered abort tail control should close exactly",
4743            );
4744        });
4745        assert_eq!(
4746            drive_generated_startup_recovery_page(
4747                &startup_session,
4748                &ABORT_REGISTRY,
4749                with_check.submission_key().as_str(),
4750            )
4751            .expect("an exact aborted generated submission should publish terminal startup state"),
4752            GeneratedStartupDriverStep::Terminal,
4753        );
4754        let failure = observe_generated_startup_state::<AbortCanister>(
4755            &ABORT_REGISTRY,
4756            with_check.submission_key().as_str(),
4757        )
4758        .expect_err("an aborted generated submission must not remain retryable forever");
4759        assert_eq!(
4760            failure.kind(),
4761            crate::db::StartupFailureKind::SchemaReconciliation,
4762        );
4763        assert_eq!(
4764            failure.diagnostic().error_code(),
4765            icydb_diagnostic_code::ErrorCode::RUNTIME_CONFLICT,
4766        );
4767    }
4768
4769    #[cfg(feature = "migration")]
4770    #[test]
4771    fn migration_planning_failures_retain_typed_public_classification() {
4772        use super::schema_migration_planning_error;
4773        use crate::db::schema::migration_planner::SchemaMigrationPlanningError;
4774        use icydb_diagnostic_code::{DiagnosticDetail, SchemaMigrationCode};
4775
4776        for (error, reason) in [
4777            (
4778                SchemaMigrationPlanningError::Unadopted,
4779                SchemaMigrationCode::Unadopted,
4780            ),
4781            (
4782                SchemaMigrationPlanningError::MissingMigration,
4783                SchemaMigrationCode::MissingMigration,
4784            ),
4785            (
4786                SchemaMigrationPlanningError::VersionGap,
4787                SchemaMigrationCode::VersionGap,
4788            ),
4789            (
4790                SchemaMigrationPlanningError::Downgrade,
4791                SchemaMigrationCode::Downgrade,
4792            ),
4793            (
4794                SchemaMigrationPlanningError::EmptyEntityVersionBump,
4795                SchemaMigrationCode::EmptyEntityVersionBump,
4796            ),
4797            (
4798                SchemaMigrationPlanningError::StaleAcceptedHead,
4799                SchemaMigrationCode::StaleAcceptedHead,
4800            ),
4801            (
4802                SchemaMigrationPlanningError::UnknownFromObject,
4803                SchemaMigrationCode::UnknownFromObject,
4804            ),
4805            (
4806                SchemaMigrationPlanningError::UnknownToObject,
4807                SchemaMigrationCode::UnknownToObject,
4808            ),
4809            (
4810                SchemaMigrationPlanningError::KindMismatch,
4811                SchemaMigrationCode::KindMismatch,
4812            ),
4813            (
4814                SchemaMigrationPlanningError::IdentityConflict,
4815                SchemaMigrationCode::IdentityConflict,
4816            ),
4817            (
4818                SchemaMigrationPlanningError::UnexplainedSchemaDifference,
4819                SchemaMigrationCode::UnexplainedSchemaDifference,
4820            ),
4821            (
4822                SchemaMigrationPlanningError::UnsupportedTransform,
4823                SchemaMigrationCode::UnsupportedTransform,
4824            ),
4825            (
4826                SchemaMigrationPlanningError::RekeyedCatalogInvalid,
4827                SchemaMigrationCode::CandidateMismatch,
4828            ),
4829            (
4830                SchemaMigrationPlanningError::CandidateMismatch,
4831                SchemaMigrationCode::CandidateMismatch,
4832            ),
4833            (
4834                SchemaMigrationPlanningError::CorruptLineage,
4835                SchemaMigrationCode::ProgressCorrupt,
4836            ),
4837        ] {
4838            let diagnostic = schema_migration_planning_error(error).diagnostic();
4839            assert_eq!(
4840                diagnostic.detail(),
4841                Some(&DiagnosticDetail::SchemaMigration { reason }),
4842            );
4843            assert_eq!(diagnostic.code(), reason.diagnostic_code());
4844        }
4845    }
4846
4847    #[cfg(feature = "migration")]
4848    #[test]
4849    #[expect(
4850        clippy::too_many_lines,
4851        reason = "the validation replay, staging, and unchanged-row assertions form one scenario"
4852    )]
4853    fn physical_migration_validation_is_bounded_staged_and_does_not_rewrite_rows() {
4854        use std::convert::Infallible;
4855
4856        use super::{defer_generated_schema_application_for_prepared_migration, migrate_schema};
4857        use crate::db::{
4858            data::StoreVisit,
4859            index::{IndexEntryValue, IndexId, IndexKey, IndexKeyKind},
4860            key_taxonomy::{PrimaryKeyComponent, PrimaryKeyValue},
4861            schema::{SchemaMigrationCommand, SchemaMigrationPhase},
4862        };
4863        use crate::types::EntityTag;
4864
4865        let db = Db::<MigrationCanister>::new(
4866            &MIGRATION_REGISTRY,
4867            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
4868        );
4869        drive_startup_recovery_to_completion(&db);
4870        let initial_target = schema_application_target(&db).expect("initial target should issue");
4871        let store_identity = initial_target
4872            .stores()
4873            .first()
4874            .expect("migration store should exist")
4875            .identity();
4876        let initial = validation_migration_proposal(
4877            ValidationMigrationShape::Clean,
4878            false,
4879            initial_target.accepted_head().clone(),
4880            initial_target.database_identity(),
4881            store_identity,
4882        );
4883        apply_schema(&db, &initial).expect("initial schema should publish");
4884
4885        let session = DbSession::<MigrationCanister>::new(
4886            &MIGRATION_REGISTRY,
4887            &crate::db::RequestExecutionRoot::__new_runtime_root(),
4888        );
4889        for (id, value) in [(1, 7), (2, 8)] {
4890            session
4891                .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
4892                    entity: "MigratingItem".to_string(),
4893                    patch: DynamicStructuralPatch::new(vec![
4894                        (
4895                            "id".to_string(),
4896                            DynamicWriteCell::Value(InputValue::nat64(id)),
4897                        ),
4898                        (
4899                            "old_value".to_string(),
4900                            DynamicWriteCell::Value(InputValue::int64(value)),
4901                        ),
4902                    ]),
4903                })
4904                .expect("predecessor row should insert");
4905        }
4906        let store = db
4907            .store_handle(MIGRATION_STORE_PATH)
4908            .expect("migration store should resolve");
4909        let row_bytes = || {
4910            store.with_data(|data| {
4911                let mut rows = Vec::new();
4912                let result: Result<(), Infallible> = data.visit_entries(|key, row| {
4913                    rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
4914                    Ok(StoreVisit::Continue)
4915                });
4916                result.expect("infallible row visit should complete");
4917                rows
4918            })
4919        };
4920        let before_rows = row_bytes();
4921
4922        let target = schema_application_target(&db).expect("migration target should issue");
4923        let proposal = validation_migration_proposal(
4924            ValidationMigrationShape::Clean,
4925            true,
4926            target.accepted_head().clone(),
4927            target.database_identity(),
4928            store_identity,
4929        );
4930        let plan = proposal
4931            .migration()
4932            .expect("migration plan should exist")
4933            .digest();
4934        let command = || SchemaMigrationCommand::Advance {
4935            expected_database: target.database_identity(),
4936            expected_head: target.accepted_head().clone(),
4937            expected_plan: plan,
4938            acknowledged_finding_page: None,
4939        };
4940        assert_eq!(
4941            migrate_schema(&db, &proposal, command())
4942                .unwrap_or_else(|error| {
4943                    panic!(
4944                        "physical migration should prepare: {:?}",
4945                        error.diagnostic()
4946                    )
4947                })
4948                .phase(),
4949            SchemaMigrationPhase::Prepared,
4950        );
4951        assert_eq!(
4952            migrate_schema(&db, &proposal, command())
4953                .expect("physical migration should enter validation")
4954                .phase(),
4955            SchemaMigrationPhase::Validating,
4956        );
4957        let record = super::load_schema_migration_record()
4958            .expect("migration record should remain readable")
4959            .expect("validating migration record should exist");
4960        let planned = super::recompile_active_physical_migration(&db, &proposal, &record)
4961            .expect("the exact active plan should recompile");
4962        for _ in 0..2 {
4963            let page = super::validate_migration_page(&db, &planned, record.progress())
4964                .expect("the same validation page should remain replayable");
4965            let (progress, staged, exhausted) = page.into_parts();
4966            assert!(progress.findings().is_empty());
4967            assert!(exhausted);
4968            super::stage_migration_index_entries(staged)
4969                .expect("staging before a cursor marker should be idempotent");
4970        }
4971        assert_eq!(
4972            store.with_index(IndexStore::len),
4973            2,
4974            "replaying an uncheckpointed page must retain one exact staged key per row",
4975        );
4976        let ready =
4977            migrate_schema(&db, &proposal, command()).expect("bounded validation should complete");
4978        assert_eq!(ready.phase(), SchemaMigrationPhase::ReadyToRewrite);
4979        assert_eq!(ready.rows_validated(), 2);
4980        assert!(ready.findings().is_empty());
4981        assert_eq!(row_bytes(), before_rows, "validation must not rewrite rows");
4982        assert_eq!(
4983            store.with_index(IndexStore::len),
4984            2,
4985            "the isolated candidate unique generation should be durably staged",
4986        );
4987        store.with_index_mut(|index| {
4988            for ordinal in 0..513_u64 {
4989                let component = ordinal.to_be_bytes();
4990                let key = IndexKey::new_from_components_with_primary_key_value(
4991                    &IndexId::new(EntityTag::new(2), 0),
4992                    IndexKeyKind::User,
4993                    &[component],
4994                    &PrimaryKeyValue::from(PrimaryKeyComponent::Nat64(ordinal)),
4995                )
4996                .expect("unrelated abort-scan key should build")
4997                .to_raw()
4998                .expect("unrelated abort-scan key should encode");
4999                index.insert(key, IndexEntryValue::presence());
5000            }
5001        });
5002        let abort = || SchemaMigrationCommand::Abort {
5003            expected_database: target.database_identity(),
5004            expected_head: target.accepted_head().clone(),
5005            expected_plan: plan,
5006        };
5007        let cleaning = migrate_schema(&db, &proposal, abort())
5008            .expect("the first bounded abort cleanup page should publish");
5009        assert_eq!(cleaning.phase(), SchemaMigrationPhase::ReadyToRewrite);
5010        assert_eq!(store.with_index(IndexStore::len), 513);
5011        let aborted =
5012            migrate_schema(&db, &proposal, abort()).expect("pre-rewrite migration should abort");
5013        assert_eq!(aborted.phase(), SchemaMigrationPhase::Aborted);
5014        assert_eq!(
5015            store.with_index(IndexStore::len),
5016            513,
5017            "abort must remove only planner-invisible candidate generations",
5018        );
5019        assert_eq!(
5020            row_bytes(),
5021            before_rows,
5022            "abort must retain predecessor rows"
5023        );
5024        assert!(
5025            !defer_generated_schema_application_for_prepared_migration(&db, &proposal)
5026                .expect("terminal aborted record must not block generated startup"),
5027        );
5028    }
5029
5030    #[cfg(feature = "migration")]
5031    #[test]
5032    #[expect(
5033        clippy::too_many_lines,
5034        reason = "the interrupted rewrite, recovery, final proof, and publication form one scenario"
5035    )]
5036    fn physical_migration_rewrite_recovers_and_publishes_one_complete_candidate() {
5037        use super::{
5038            defer_generated_schema_application_for_prepared_migration, migrate_schema,
5039            schema_migration_status_for_target,
5040        };
5041        use crate::db::{
5042            data::{CanonicalSlotReader, DecodedDataStoreKey, StoreVisit, StructuralSlotReader},
5043            schema::{
5044                MigrationRewriteInterruption, SchemaMigrationCommand, SchemaMigrationPhase,
5045                ensure_schema_migration_ready_for_ordinary_operations,
5046                interrupt_next_migration_rewrite_at,
5047            },
5048        };
5049        use crate::error::InternalError;
5050
5051        let db = Db::<MigrationExecutionCanister>::new(
5052            &MIGRATION_EXECUTION_REGISTRY,
5053            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5054        );
5055        drive_startup_recovery_to_completion(&db);
5056        let initial_target = schema_application_target(&db).expect("initial target should issue");
5057        let store_identity = initial_target
5058            .stores()
5059            .first()
5060            .expect("migration execution store should exist")
5061            .identity();
5062        let initial = validation_migration_proposal(
5063            ValidationMigrationShape::Clean,
5064            false,
5065            initial_target.accepted_head().clone(),
5066            initial_target.database_identity(),
5067            store_identity,
5068        );
5069        apply_schema(&db, &initial).expect("initial schema should publish");
5070        let session = DbSession::<MigrationExecutionCanister>::new(
5071            &MIGRATION_EXECUTION_REGISTRY,
5072            &crate::db::RequestExecutionRoot::__new_runtime_root(),
5073        );
5074        for (id, value) in [(1, 7), (2, 8), (3, 9)] {
5075            session
5076                .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5077                    entity: "MigratingItem".to_string(),
5078                    patch: DynamicStructuralPatch::new(vec![
5079                        (
5080                            "id".to_string(),
5081                            DynamicWriteCell::Value(InputValue::nat64(id)),
5082                        ),
5083                        (
5084                            "old_value".to_string(),
5085                            DynamicWriteCell::Value(InputValue::int64(value)),
5086                        ),
5087                    ]),
5088                })
5089                .expect("predecessor row should insert");
5090        }
5091        let target = schema_application_target(&db).expect("migration target should issue");
5092        let proposal = validation_migration_proposal(
5093            ValidationMigrationShape::Clean,
5094            true,
5095            target.accepted_head().clone(),
5096            target.database_identity(),
5097            store_identity,
5098        );
5099        let plan = proposal
5100            .migration()
5101            .expect("migration plan should exist")
5102            .digest();
5103        let command = || SchemaMigrationCommand::Advance {
5104            expected_database: target.database_identity(),
5105            expected_head: target.accepted_head().clone(),
5106            expected_plan: plan,
5107            acknowledged_finding_page: None,
5108        };
5109        for expected in [
5110            SchemaMigrationPhase::Prepared,
5111            SchemaMigrationPhase::Validating,
5112            SchemaMigrationPhase::ReadyToRewrite,
5113            SchemaMigrationPhase::RewritingRows,
5114        ] {
5115            assert_eq!(
5116                migrate_schema(&db, &proposal, command())
5117                    .expect("migration phase should advance")
5118                    .phase(),
5119                expected,
5120            );
5121        }
5122
5123        for interruption in [
5124            MigrationRewriteInterruption::MarkerPersisted,
5125            MigrationRewriteInterruption::JournalPublished,
5126            MigrationRewriteInterruption::PhysicalApplied,
5127        ] {
5128            interrupt_next_migration_rewrite_at(interruption);
5129            migrate_schema(&db, &proposal, command())
5130                .expect_err("injected interruption should retain the rewrite marker");
5131
5132            forget_recovered_domain_for_tests(&db)
5133                .expect("upgrade should reset recovery ownership");
5134            drive_startup_recovery_to_completion(&db);
5135        }
5136
5137        let rebuilding = schema_migration_status_for_target(
5138            &db,
5139            &proposal,
5140            &schema_application_target(&db).expect("recovered target should issue"),
5141        )
5142        .expect("recovered status should remain readable");
5143        assert_eq!(rebuilding.phase(), SchemaMigrationPhase::RebuildingIndexes);
5144        assert_eq!(rebuilding.rows_rewritten(), 3);
5145        assert_eq!(
5146            migrate_schema(&db, &proposal, command())
5147                .expect("derived generations should complete")
5148                .phase(),
5149            SchemaMigrationPhase::FinalValidation,
5150        );
5151        assert_eq!(
5152            migrate_schema(&db, &proposal, command())
5153                .expect("final validation should complete")
5154                .phase(),
5155            SchemaMigrationPhase::Publishing,
5156        );
5157        let applied = migrate_schema(&db, &proposal, command())
5158            .expect("candidate publication should complete atomically");
5159        assert_eq!(applied.phase(), SchemaMigrationPhase::Applied);
5160        assert_eq!(applied.rows_rewritten(), 3);
5161        assert_eq!(applied.indexes_rebuilt(), 1);
5162        assert_ne!(applied.accepted_head(), target.accepted_head());
5163        let terminal_target = schema_application_target(&db).expect("terminal target should issue");
5164        let terminal_proposal = validation_migration_proposal(
5165            ValidationMigrationShape::Clean,
5166            true,
5167            terminal_target.accepted_head().clone(),
5168            terminal_target.database_identity(),
5169            store_identity,
5170        );
5171        assert!(
5172            !defer_generated_schema_application_for_prepared_migration(&db, &terminal_proposal,)
5173                .expect("terminal record must not block generated startup"),
5174        );
5175
5176        let store = db
5177            .store_handle(MIGRATION_EXECUTION_STORE_PATH)
5178            .expect("migration execution store should resolve");
5179        let runtime = db
5180            .accepted_runtime_entity_for_path("MigratingItem")
5181            .expect("published candidate entity should resolve");
5182        let selection = store
5183            .with_schema(|schema| {
5184                schema.current_accepted_catalog_selection(
5185                    runtime.entity_tag(),
5186                    runtime.entity_path(),
5187                    runtime.store_path(),
5188                )
5189            })
5190            .expect("candidate selection should remain readable")
5191            .expect("candidate selection should exist");
5192        let contract = crate::db::data::AcceptedStructuralRowAuthority::from_catalog_selection(
5193            runtime.entity_path(),
5194            &selection,
5195        )
5196        .expect("candidate row authority should compile")
5197        .into_row_contract();
5198        let mut values = Vec::new();
5199        store
5200            .with_data(|data| {
5201                data.visit_entries(|key, row| {
5202                    let decoded = DecodedDataStoreKey::try_from_raw(key)
5203                        .expect("rewritten key should decode");
5204                    let reader =
5205                        StructuralSlotReader::from_raw_row_with_validated_borrowed_contract(
5206                            row, &contract,
5207                        )
5208                        .expect("rewritten row should use the candidate layout");
5209                    reader
5210                        .validate_primary_key(&decoded)
5211                        .expect("rewritten row and key should remain bound");
5212                    values.push(
5213                        reader
5214                            .required_value_by_contract(1)
5215                            .expect("candidate value slot should decode"),
5216                    );
5217                    Ok::<StoreVisit, InternalError>(StoreVisit::Continue)
5218                })
5219            })
5220            .expect("rewritten row scan should complete");
5221        assert_eq!(
5222            values,
5223            vec![
5224                crate::value::Value::Nat64(7),
5225                crate::value::Value::Nat64(8),
5226                crate::value::Value::Nat64(9),
5227            ],
5228        );
5229        assert_eq!(store.with_index(IndexStore::len), 3);
5230        let accepted = store
5231            .with_schema(SchemaStore::current_accepted_schema_bundle)
5232            .expect("published candidate bundle should remain readable")
5233            .expect("published candidate bundle should exist");
5234        let entity_source = EntitySourceKey::try_new("MigratingItem")
5235            .expect("migration entity source should admit");
5236        let entity_tag = accepted
5237            .source_bindings_for_tests()
5238            .entity(&entity_source)
5239            .expect("candidate entity source should remain bound");
5240        let old_value =
5241            FieldSourceKey::try_new("old_value").expect("predecessor source should admit");
5242        let current_value =
5243            FieldSourceKey::try_new("value").expect("candidate source should admit");
5244        assert_eq!(
5245            accepted
5246                .source_bindings_for_tests()
5247                .field(entity_tag, &old_value),
5248            None,
5249        );
5250        assert!(
5251            accepted
5252                .source_bindings_for_tests()
5253                .field(entity_tag, &current_value)
5254                .is_some(),
5255        );
5256        ensure_schema_migration_ready_for_ordinary_operations()
5257            .expect("terminal publication must clear the database-wide gate");
5258    }
5259
5260    #[cfg(feature = "migration")]
5261    #[test]
5262    #[expect(
5263        clippy::too_many_lines,
5264        reason = "all four finding families share one ordered historical scan fixture"
5265    )]
5266    fn physical_migration_validation_reports_every_typed_finding_family_without_writes() {
5267        use std::convert::Infallible;
5268
5269        use super::migrate_schema;
5270        use crate::db::{
5271            data::StoreVisit,
5272            schema::{SchemaMigrationCommand, SchemaMigrationFindingKind, SchemaMigrationPhase},
5273        };
5274
5275        let db = Db::<MigrationFindingCanister>::new(
5276            &MIGRATION_FINDING_REGISTRY,
5277            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5278        );
5279        drive_startup_recovery_to_completion(&db);
5280        let initial_target = schema_application_target(&db).expect("initial target should issue");
5281        let store_identity = initial_target
5282            .stores()
5283            .first()
5284            .expect("migration finding store should exist")
5285            .identity();
5286        let initial = validation_migration_proposal(
5287            ValidationMigrationShape::AllFindingFamilies,
5288            false,
5289            initial_target.accepted_head().clone(),
5290            initial_target.database_identity(),
5291            store_identity,
5292        );
5293        apply_schema(&db, &initial).expect("initial finding schema should publish");
5294
5295        let session = DbSession::<MigrationFindingCanister>::new(
5296            &MIGRATION_FINDING_REGISTRY,
5297            &crate::db::RequestExecutionRoot::__new_runtime_root(),
5298        );
5299        session
5300            .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5301                entity: "MigrationTarget".to_string(),
5302                patch: DynamicStructuralPatch::new(vec![(
5303                    "id".to_string(),
5304                    DynamicWriteCell::Value(InputValue::nat64(7)),
5305                )]),
5306            })
5307            .expect("relation target should insert");
5308        for (id, value) in [(1, 9), (2, 8), (3, 7), (4, 7), (5, 300)] {
5309            session
5310                .execute_trusted_dynamic_mutation(&DynamicMutation::Insert {
5311                    entity: "MigratingItem".to_string(),
5312                    patch: DynamicStructuralPatch::new(vec![
5313                        (
5314                            "id".to_string(),
5315                            DynamicWriteCell::Value(InputValue::nat64(id)),
5316                        ),
5317                        (
5318                            "old_value".to_string(),
5319                            DynamicWriteCell::Value(InputValue::int64(value)),
5320                        ),
5321                    ]),
5322                })
5323                .expect("predecessor finding row should insert");
5324        }
5325        let store = db
5326            .store_handle(MIGRATION_FINDING_STORE_PATH)
5327            .expect("migration finding store should resolve");
5328        let row_bytes = || {
5329            store.with_data(|data| {
5330                let mut rows = Vec::new();
5331                let result: Result<(), Infallible> = data.visit_entries(|key, row| {
5332                    rows.push((key.as_bytes().to_vec(), row.as_bytes().to_vec()));
5333                    Ok(StoreVisit::Continue)
5334                });
5335                result.expect("infallible row visit should complete");
5336                rows
5337            })
5338        };
5339        let before_rows = row_bytes();
5340
5341        let target = schema_application_target(&db).expect("migration target should issue");
5342        let proposal = validation_migration_proposal(
5343            ValidationMigrationShape::AllFindingFamilies,
5344            true,
5345            target.accepted_head().clone(),
5346            target.database_identity(),
5347            store_identity,
5348        );
5349        let plan = proposal
5350            .migration()
5351            .expect("migration plan should exist")
5352            .digest();
5353        let command = || SchemaMigrationCommand::Advance {
5354            expected_database: target.database_identity(),
5355            expected_head: target.accepted_head().clone(),
5356            expected_plan: plan,
5357            acknowledged_finding_page: None,
5358        };
5359        assert_eq!(
5360            migrate_schema(&db, &proposal, command())
5361                .expect("finding migration should prepare")
5362                .phase(),
5363            SchemaMigrationPhase::Prepared,
5364        );
5365        assert_eq!(
5366            migrate_schema(&db, &proposal, command())
5367                .expect("finding migration should enter validation")
5368                .phase(),
5369            SchemaMigrationPhase::Validating,
5370        );
5371        let rejected =
5372            migrate_schema(&db, &proposal, command()).expect("validation should report findings");
5373        assert_eq!(rejected.phase(), SchemaMigrationPhase::Rejected);
5374        assert_eq!(rejected.rows_validated(), 5);
5375        assert_eq!(
5376            rejected
5377                .findings()
5378                .iter()
5379                .map(crate::db::schema::SchemaMigrationFinding::kind)
5380                .collect::<Vec<_>>(),
5381            vec![
5382                SchemaMigrationFindingKind::Constraint,
5383                SchemaMigrationFindingKind::Relation,
5384                SchemaMigrationFindingKind::UniqueIndex,
5385                SchemaMigrationFindingKind::Transform,
5386            ],
5387        );
5388        assert_eq!(
5389            row_bytes(),
5390            before_rows,
5391            "rejected validation must not rewrite accepted rows"
5392        );
5393        assert_eq!(
5394            store.with_index(IndexStore::len),
5395            0,
5396            "a rejected page must not publish any staged generation"
5397        );
5398    }
5399
5400    #[cfg(feature = "migration")]
5401    #[test]
5402    fn exact_migration_retry_binds_the_terminal_head_not_the_predecessor_head() {
5403        use super::exact_migration_replay_target;
5404
5405        let db = Db::<EvolutionCanister>::new(
5406            &EVOLUTION_REGISTRY,
5407            crate::db::RequestExecutionRoot::__new_runtime_root().scope(),
5408        );
5409        drive_startup_recovery_to_completion(&db);
5410        let initial_target = schema_application_target(&db).expect("initial target should issue");
5411        let (proposal, _, _) = generated_check_proposal(
5412            initial_target.accepted_head().clone(),
5413            "migration-retry-initial",
5414            false,
5415            initial_target.database_identity(),
5416            initial_target
5417                .stores()
5418                .first()
5419                .expect("test store should exist")
5420                .identity(),
5421        );
5422        apply_schema(&db, &proposal).expect("initial schema should publish");
5423        let current_target = schema_application_target(&db).expect("current target should issue");
5424        assert_ne!(
5425            current_target.accepted_head(),
5426            initial_target.accepted_head(),
5427        );
5428
5429        let receipt = SchemaChangeReceipt::new(
5430            current_target.database_identity(),
5431            SchemaSubmissionKey::try_new("migration/retry")
5432                .expect("migration submission should admit"),
5433            SchemaProposalDigest::from_bytes([0x77; 32]),
5434            initial_target.accepted_head().clone(),
5435            SchemaChangeOutcome::Applied {
5436                accepted_head: current_target.accepted_head().clone(),
5437            },
5438        )
5439        .expect("terminal migration receipt should admit");
5440        let record = SchemaApplicationRecord::new(receipt, Vec::new())
5441            .expect("terminal migration record should admit");
5442
5443        assert_eq!(
5444            exact_migration_replay_target(&db, current_target.database_identity(), &record,)
5445                .expect("exact retry should resolve the terminal target"),
5446            current_target,
5447        );
5448    }
5449}