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