Skip to main content

type_bridge_schema_migration/
execution.rs

1//! Provider-neutral fenced migration journal and recovery contracts.
2
3use std::future::Future;
4use std::pin::Pin;
5
6use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
7use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileFingerprint};
8use type_bridge_contract::migration::{
9    MigrationId, MigrationManifestDigest, MigrationPlanFingerprint,
10};
11use type_bridge_contract::schema_delta::ManagedSchemaState;
12use type_bridge_contract::schema_fingerprint::{
13    ManagedDeclaredIdentityFingerprint, ManagedSemanticSchemaFingerprint,
14};
15use type_bridge_contract::schema_lowering::SchemaLoweringProfileFingerprint;
16
17use crate::{
18    VerifiedMigrationApplyManifest, VerifiedMigrationApplyPlan, VerifiedMigrationRollbackManifest,
19    VerifiedMigrationRollbackPlan, VerifiedMigrationTransactionGroup,
20};
21
22const MAX_LEASE_HOLDER_BYTES: usize = 128;
23
24/// Boxed future returned by provider-neutral execution stores.
25pub type ExecutionFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, Diagnostic>> + Send + 'a>>;
26
27/// A monotonically increasing store-issued migration fencing token.
28#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
29pub struct ExecutionFence(u64);
30
31impl ExecutionFence {
32    /// Construct a non-zero fence.
33    pub fn new(value: u64) -> Result<Self, Diagnostic> {
34        if value == 0 {
35            return Err(failure(
36                DiagnosticCategory::InvalidContract,
37                "migration_execution_zero_fence",
38                "migration execution fences must be non-zero",
39            ));
40        }
41        Ok(Self(value))
42    }
43
44    /// Return the numeric fence value.
45    pub const fn get(self) -> u64 {
46        self.0
47    }
48
49    /// Derive the next strictly greater fence without wrapping.
50    pub fn checked_successor(self) -> Result<Self, Diagnostic> {
51        let value = self.0.checked_add(1).ok_or_else(|| {
52            failure(
53                DiagnosticCategory::ResourceLimit,
54                "migration_execution_fence_exhausted",
55                "migration execution fence range is exhausted",
56            )
57        })?;
58        Self::new(value)
59    }
60}
61
62/// One durable managed migration scope.
63#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
64pub struct ExecutionScope(ManagedScopeId);
65
66impl ExecutionScope {
67    /// Bind execution to an existing managed-scope identity.
68    pub const fn new(scope: ManagedScopeId) -> Self {
69        Self(scope)
70    }
71
72    /// Return the managed-scope identity.
73    pub const fn managed_scope_id(&self) -> &ManagedScopeId {
74        &self.0
75    }
76}
77
78/// Caller-supplied lease holder identity with no machine-derived garnish.
79#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
80pub struct LeaseHolderId(String);
81
82impl LeaseHolderId {
83    /// Validate a bounded canonical holder label.
84    pub fn new(value: impl Into<String>) -> Result<Self, Diagnostic> {
85        let value = value.into();
86        if value.is_empty()
87            || value.len() > MAX_LEASE_HOLDER_BYTES
88            || !value
89                .bytes()
90                .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
91        {
92            return Err(failure(
93                DiagnosticCategory::InvalidContract,
94                "migration_execution_invalid_holder",
95                "lease holder must be bounded non-empty ASCII [A-Za-z0-9._-]",
96            ));
97        }
98        Ok(Self(value))
99    }
100
101    /// Return the canonical holder label.
102    pub fn as_str(&self) -> &str {
103        &self.0
104    }
105}
106
107/// Store-issued authority to mutate one migration scope.
108#[derive(Clone, Debug, Eq, PartialEq)]
109pub struct MigrationLease {
110    scope: ExecutionScope,
111    holder: LeaseHolderId,
112    fence: ExecutionFence,
113}
114
115impl MigrationLease {
116    /// Construct the lease returned by a store after atomic acquisition.
117    pub const fn new(scope: ExecutionScope, holder: LeaseHolderId, fence: ExecutionFence) -> Self {
118        Self {
119            scope,
120            holder,
121            fence,
122        }
123    }
124
125    /// Return the leased scope.
126    pub const fn scope(&self) -> &ExecutionScope {
127        &self.scope
128    }
129
130    /// Return the holder identity.
131    pub const fn holder(&self) -> &LeaseHolderId {
132        &self.holder
133    }
134
135    /// Return the store-issued fence.
136    pub const fn fence(&self) -> ExecutionFence {
137        self.fence
138    }
139}
140
141/// Commit certainty owned by the migration journal layer.
142///
143/// Provider adapters must map absent certainty information to [`Self::Unknown`].
144/// Only an explicit provider proof that commit could not have occurred may map
145/// to [`Self::DefinitelyAborted`].
146#[derive(Clone, Copy, Debug, Eq, PartialEq)]
147pub enum GroupCommitCertainty {
148    /// The provider proves that the group transaction did not commit.
149    DefinitelyAborted,
150    /// The provider cannot determine whether the group transaction committed.
151    Unknown,
152}
153
154impl GroupCommitCertainty {
155    /// Convert certainty into its durable journal event.
156    pub const fn journal_event(self) -> GroupJournalEventKind {
157        match self {
158            Self::DefinitelyAborted => GroupJournalEventKind::DefinitelyAborted,
159            Self::Unknown => GroupJournalEventKind::CommitOutcomeUnknown,
160        }
161    }
162}
163
164/// Durable commit-boundary event vocabulary for one transaction group.
165#[derive(Clone, Copy, Debug, Eq, PartialEq)]
166pub enum GroupJournalEventKind {
167    /// The group statements completed and the commit call is about to begin.
168    BeforeCommit,
169    /// The provider commit succeeded and exact live target semantics were observed.
170    Committed,
171    /// The commit response cannot prove whether durability occurred.
172    CommitOutcomeUnknown,
173    /// The provider proves that the transaction did not commit.
174    DefinitelyAborted,
175    /// An empty formal-only group advanced without a provider transaction.
176    FormalOnlyAdvanced,
177}
178
179/// Optional live managed-semantic evidence used during recovery.
180#[derive(Clone, Debug, Eq, PartialEq)]
181pub enum GroupRecoveryObservation {
182    /// The provider cannot supply a trustworthy managed-semantic observation.
183    Unavailable,
184    /// Exact live managed semantics observed under the current fence.
185    ManagedSemantics(ManagedSemanticSchemaFingerprint),
186}
187
188/// Fail-closed recovery decision for one positional transaction group.
189#[derive(Clone, Copy, Debug, Eq, PartialEq)]
190pub enum GroupRecoveryDecision {
191    /// The group is proven absent and may execute under the current fence.
192    ExecuteNormally,
193    /// The target is proven reached and only the journal checkpoint needs repair.
194    RepairCheckpoint,
195    /// Evidence is ambiguous or contradictory and requires verified operator action.
196    RequiresExplicitRecovery,
197}
198
199/// Decide recovery from durable event and freshly observed managed semantics.
200///
201/// Equal source and target fingerprints are intentionally uninformative for
202/// unknown commit outcomes. They never authorize replay or checkpoint repair.
203pub fn decide_group_recovery(
204    last_event: Option<GroupJournalEventKind>,
205    observation: &GroupRecoveryObservation,
206    source: &ManagedSemanticSchemaFingerprint,
207    target: &ManagedSemanticSchemaFingerprint,
208) -> GroupRecoveryDecision {
209    let distinct = source != target;
210    let observed = match observation {
211        GroupRecoveryObservation::Unavailable => ObservedRelation::Unavailable,
212        GroupRecoveryObservation::ManagedSemantics(value) if value == source && value == target => {
213            ObservedRelation::Both
214        }
215        GroupRecoveryObservation::ManagedSemantics(value) if value == source => {
216            ObservedRelation::Source
217        }
218        GroupRecoveryObservation::ManagedSemantics(value) if value == target => {
219            ObservedRelation::Target
220        }
221        GroupRecoveryObservation::ManagedSemantics(_) => ObservedRelation::Neither,
222    };
223    match (last_event, distinct, observed) {
224        (None, true, ObservedRelation::Source)
225        | (Some(GroupJournalEventKind::DefinitelyAborted), true, ObservedRelation::Source)
226        | (None, false, ObservedRelation::Both)
227        | (Some(GroupJournalEventKind::DefinitelyAborted), false, ObservedRelation::Both) => {
228            GroupRecoveryDecision::ExecuteNormally
229        }
230        (
231            Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
232            true,
233            ObservedRelation::Source,
234        ) => GroupRecoveryDecision::ExecuteNormally,
235        (
236            Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
237            true,
238            ObservedRelation::Target,
239        )
240        | (Some(GroupJournalEventKind::Committed), true, ObservedRelation::Target)
241        | (
242            Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced),
243            false,
244            ObservedRelation::Both,
245        ) => GroupRecoveryDecision::RepairCheckpoint,
246        _ => GroupRecoveryDecision::RequiresExplicitRecovery,
247    }
248}
249
250#[derive(Clone, Copy, Debug, Eq, PartialEq)]
251enum ObservedRelation {
252    Source,
253    Target,
254    Both,
255    Neither,
256    Unavailable,
257}
258
259/// Store-assigned monotonic ordering identity for journal entries.
260#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
261pub struct JournalSequence(u64);
262
263impl JournalSequence {
264    /// Construct a non-zero journal sequence.
265    pub fn new(value: u64) -> Result<Self, Diagnostic> {
266        if value == 0 {
267            return Err(failure(
268                DiagnosticCategory::InvalidContract,
269                "migration_execution_zero_sequence",
270                "journal sequence numbers must be non-zero",
271            ));
272        }
273        Ok(Self(value))
274    }
275
276    /// Return the numeric sequence value.
277    pub const fn get(self) -> u64 {
278        self.0
279    }
280}
281
282/// One record after the store atomically assigns its ordering sequence.
283#[derive(Clone, Debug, Eq, PartialEq)]
284pub struct JournalEntry<T> {
285    sequence: JournalSequence,
286    record: T,
287}
288
289impl<T> JournalEntry<T> {
290    /// Attach a sequence allocated by the authoritative store.
291    pub const fn from_store(sequence: JournalSequence, record: T) -> Self {
292        Self { sequence, record }
293    }
294
295    /// Return the store ordering sequence.
296    pub const fn sequence(&self) -> JournalSequence {
297        self.sequence
298    }
299
300    /// Return the trusted record.
301    pub const fn record(&self) -> &T {
302        &self.record
303    }
304
305    /// Consume the envelope and return its record.
306    pub fn into_record(self) -> T {
307        self.record
308    }
309}
310
311/// One trusted, still-open plan and its ordered group events.
312#[derive(Clone, Debug, Eq, PartialEq)]
313pub struct OpenPlanRecord {
314    plan: JournalEntry<PlanRecord>,
315    events: Vec<JournalEntry<GroupEventRecord>>,
316}
317
318impl OpenPlanRecord {
319    /// Rebuild store output while checking sequence, scope, fence, and manifest binding.
320    ///
321    /// The plan retains its original fence while recovery events may be written
322    /// under later fences. Event fences therefore must be monotonic and no older
323    /// than the plan fence; requiring equality would hide durable recovery work
324    /// after the lease rolls forward.
325    pub fn from_store(
326        plan: JournalEntry<PlanRecord>,
327        events: Vec<JournalEntry<GroupEventRecord>>,
328    ) -> Result<Self, Diagnostic> {
329        let mut previous = plan.sequence();
330        let mut previous_fence = plan.record().fence();
331        for event in &events {
332            let manifest_index = plan
333                .record()
334                .manifest_digests()
335                .iter()
336                .position(|digest| digest == &event.record().manifest_digest());
337            if event.sequence() <= previous
338                || event.record().scope() != plan.record().scope()
339                || event.record().fence() < previous_fence
340                || manifest_index.is_none_or(|index| {
341                    plan.record().migration_ids().get(index) != Some(event.record().migration_id())
342                })
343            {
344                return Err(failure(
345                    DiagnosticCategory::Integrity,
346                    "migration_execution_invalid_open_plan",
347                    "loaded open-plan events are not ordered and bound to the plan",
348                ));
349            }
350            previous = event.sequence();
351            previous_fence = event.record().fence();
352        }
353        Ok(Self { plan, events })
354    }
355
356    /// Return the sequenced plan record.
357    pub const fn plan(&self) -> &JournalEntry<PlanRecord> {
358        &self.plan
359    }
360
361    /// Return ordered sequenced group events.
362    pub fn events(&self) -> &[JournalEntry<GroupEventRecord>] {
363        &self.events
364    }
365}
366
367/// Identity-only journal record for one complete verified apply plan.
368#[derive(Clone, Debug, Eq, PartialEq)]
369pub struct PlanRecord {
370    scope: ExecutionScope,
371    fence: ExecutionFence,
372    source_applied: Vec<MigrationId>,
373    source_frontier: Vec<MigrationId>,
374    target_frontier: Vec<MigrationId>,
375    migration_ids: Vec<MigrationId>,
376    manifest_digests: Vec<MigrationManifestDigest>,
377    manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
378    source_declared: ManagedDeclaredIdentityFingerprint,
379    target_declared: ManagedDeclaredIdentityFingerprint,
380    source_semantics: ManagedSemanticSchemaFingerprint,
381    target_semantics: ManagedSemanticSchemaFingerprint,
382    semantic_profile: SemanticProfileFingerprint,
383    lowering_profile: SchemaLoweringProfileFingerprint,
384    observed_live_source: ManagedSemanticSchemaFingerprint,
385}
386
387impl PlanRecord {
388    /// Bind a fresh-lease ledger and live-state precondition to verified plan identities.
389    pub fn from_verified_plan(
390        lease: &MigrationLease,
391        plan: &VerifiedMigrationApplyPlan,
392        observed_applied_migrations: &[MigrationId],
393        observed_live_source: &ManagedSchemaState,
394    ) -> Result<Self, Diagnostic> {
395        let source = plan.source_state().ok_or_else(|| {
396            failure(
397                DiagnosticCategory::InvalidContract,
398                "migration_execution_empty_plan",
399                "an executable migration plan requires a source state",
400            )
401        })?;
402        let target = plan.target_state().ok_or_else(|| {
403            failure(
404                DiagnosticCategory::InvalidContract,
405                "migration_execution_empty_plan",
406                "an executable migration plan requires a target state",
407            )
408        })?;
409        let first = plan.migrations().first().ok_or_else(|| {
410            failure(
411                DiagnosticCategory::InvalidContract,
412                "migration_execution_empty_plan",
413                "an executable migration plan requires at least one manifest",
414            )
415        })?;
416        if observed_applied_migrations != plan.applied_migrations() {
417            return Err(failure(
418                DiagnosticCategory::Integrity,
419                "migration_execution_stale_applied_set",
420                "applied ledger changed after migration planning; rebuild the plan",
421            ));
422        }
423        if observed_live_source != source {
424            return Err(failure(
425                DiagnosticCategory::Integrity,
426                "migration_execution_stale_source_state",
427                "live managed state differs from the planned source; rebuild the plan",
428            ));
429        }
430        let scope = ExecutionScope::new(source.scope().id().clone());
431        if lease.scope() != &scope || target.scope() != source.scope() {
432            return Err(failure(
433                DiagnosticCategory::Integrity,
434                "migration_execution_scope_mismatch",
435                "lease, source, and target must bind the same managed scope",
436            ));
437        }
438        let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
439        let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
440        for migration in plan.migrations() {
441            if migration.manifest().managed_scope().id() != scope.managed_scope_id()
442                || migration.manifest().semantic_profile().fingerprint() != &semantic_profile
443                || migration.manifest().lowering_profile().fingerprint() != &lowering_profile
444            {
445                return Err(failure(
446                    DiagnosticCategory::Integrity,
447                    "migration_execution_plan_binding_mismatch",
448                    "planned manifests do not share exact scope and profile bindings",
449                ));
450            }
451        }
452        Ok(Self {
453            scope,
454            fence: lease.fence(),
455            source_applied: plan.applied_migrations().to_vec(),
456            source_frontier: plan.applied_frontier().to_vec(),
457            target_frontier: plan.target_frontier().to_vec(),
458            migration_ids: plan
459                .migrations()
460                .iter()
461                .map(|migration| migration.manifest().id().clone())
462                .collect(),
463            manifest_digests: plan
464                .migrations()
465                .iter()
466                .map(VerifiedMigrationApplyManifest::digest)
467                .collect(),
468            manifest_plan_fingerprints: plan
469                .migrations()
470                .iter()
471                .map(|migration| migration.manifest().plan_fingerprint().clone())
472                .collect(),
473            source_declared: source.managed_declared_identity().clone(),
474            target_declared: target.managed_declared_identity().clone(),
475            source_semantics: source.managed_semantic_schema().clone(),
476            target_semantics: target.managed_semantic_schema().clone(),
477            semantic_profile,
478            lowering_profile,
479            observed_live_source: observed_live_source.managed_semantic_schema().clone(),
480        })
481    }
482
483    /// Return the execution scope.
484    pub const fn scope(&self) -> &ExecutionScope {
485        &self.scope
486    }
487    /// Return the fence bound into this record.
488    pub const fn fence(&self) -> ExecutionFence {
489        self.fence
490    }
491    /// Return the planned source frontier.
492    pub fn source_frontier(&self) -> &[MigrationId] {
493        &self.source_frontier
494    }
495    /// Return the complete canonically ordered source applied set.
496    pub fn source_applied(&self) -> &[MigrationId] {
497        &self.source_applied
498    }
499    /// Return the planned target frontier.
500    pub fn target_frontier(&self) -> &[MigrationId] {
501        &self.target_frontier
502    }
503    /// Return ordered migration identities.
504    pub fn migration_ids(&self) -> &[MigrationId] {
505        &self.migration_ids
506    }
507    /// Return ordered canonical manifest digests.
508    pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
509        &self.manifest_digests
510    }
511    /// Return ordered manifest plan fingerprints.
512    pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
513        &self.manifest_plan_fingerprints
514    }
515    /// Return planned source managed-declared identity.
516    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
517        &self.source_declared
518    }
519    /// Return planned target managed-declared identity.
520    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
521        &self.target_declared
522    }
523    /// Return planned source managed semantics.
524    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
525        &self.source_semantics
526    }
527    /// Return planned target managed semantics.
528    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
529        &self.target_semantics
530    }
531    /// Return semantic-profile content identity.
532    pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
533        &self.semantic_profile
534    }
535    /// Return lowering-registry content identity.
536    pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
537        &self.lowering_profile
538    }
539    /// Return the pre-mutation observed source semantics.
540    pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
541        &self.observed_live_source
542    }
543}
544
545/// Identity-only journal event for one verified transaction group.
546#[derive(Clone, Debug, Eq, PartialEq)]
547pub struct GroupEventRecord {
548    scope: ExecutionScope,
549    fence: ExecutionFence,
550    manifest_digest: MigrationManifestDigest,
551    migration_id: MigrationId,
552    group_ordinal: u32,
553    first_step_index: u32,
554    schema_delta_step_index: u32,
555    end_step_index: u32,
556    kind: GroupJournalEventKind,
557    observed_target: Option<ManagedSemanticSchemaFingerprint>,
558}
559
560impl GroupEventRecord {
561    /// Derive a positional event from exact verified apply evidence.
562    pub fn new(
563        lease: &MigrationLease,
564        migration: &VerifiedMigrationApplyManifest,
565        group: &VerifiedMigrationTransactionGroup,
566        kind: GroupJournalEventKind,
567        observed_target: Option<ManagedSemanticSchemaFingerprint>,
568    ) -> Result<Self, Diagnostic> {
569        if migration.transaction_groups().get(group.ordinal()) != Some(group) {
570            return Err(failure(
571                DiagnosticCategory::Integrity,
572                "migration_execution_foreign_group",
573                "transaction group does not belong to the supplied verified manifest",
574            ));
575        }
576        let step = migration
577            .steps()
578            .get(group.schema_delta_step_index())
579            .ok_or_else(|| {
580                failure(
581                    DiagnosticCategory::Integrity,
582                    "migration_execution_group_position_mismatch",
583                    "transaction group delta position is outside the verified manifest",
584                )
585            })?;
586        let delta = step
587            .step()
588            .as_schema_delta()
589            .ok_or_else(|| {
590                failure(
591                    DiagnosticCategory::Integrity,
592                    "migration_execution_group_position_mismatch",
593                    "transaction group does not terminate in a schema delta",
594                )
595            })?
596            .delta();
597        let lowering = step.lowering().ok_or_else(|| {
598            failure(
599                DiagnosticCategory::Integrity,
600                "migration_execution_group_lowering_missing",
601                "transaction group delta has no verified lowering",
602            )
603        })?;
604        let scope = ExecutionScope::new(migration.manifest().managed_scope().id().clone());
605        if lease.scope() != &scope {
606            return Err(failure(
607                DiagnosticCategory::Integrity,
608                "migration_execution_scope_mismatch",
609                "lease scope differs from the verified migration scope",
610            ));
611        }
612        match kind {
613            GroupJournalEventKind::Committed
614                if observed_target.as_ref() == Some(delta.target().managed_semantic_schema()) => {}
615            GroupJournalEventKind::Committed => {
616                return Err(failure(
617                    DiagnosticCategory::Integrity,
618                    "migration_execution_commit_evidence_mismatch",
619                    "committed event requires the exact observed target semantics",
620                ));
621            }
622            GroupJournalEventKind::FormalOnlyAdvanced
623                if observed_target.is_none()
624                    && group.assertion_count() == 0
625                    && lowering.units().is_empty()
626                    && delta.source().managed_semantic_schema()
627                        == delta.target().managed_semantic_schema() => {}
628            GroupJournalEventKind::FormalOnlyAdvanced => {
629                return Err(failure(
630                    DiagnosticCategory::InvalidContract,
631                    "migration_execution_invalid_formal_advance",
632                    "formal-only advancement requires an assertion-free empty equal-semantic group",
633                ));
634            }
635            _ if observed_target.is_none() => {}
636            _ => {
637                return Err(failure(
638                    DiagnosticCategory::InvalidContract,
639                    "migration_execution_unexpected_observation",
640                    "only committed events may carry observed target semantics",
641                ));
642            }
643        }
644        Ok(Self {
645            scope,
646            fence: lease.fence(),
647            manifest_digest: migration.digest(),
648            migration_id: migration.manifest().id().clone(),
649            group_ordinal: position(group.ordinal())?,
650            first_step_index: position(group.first_step_index())?,
651            schema_delta_step_index: position(group.schema_delta_step_index())?,
652            end_step_index: position(group.end_step_index())?,
653            kind,
654            observed_target,
655        })
656    }
657
658    /// Return the execution scope.
659    pub const fn scope(&self) -> &ExecutionScope {
660        &self.scope
661    }
662    /// Return the event fence.
663    pub const fn fence(&self) -> ExecutionFence {
664        self.fence
665    }
666    /// Return the canonical manifest digest.
667    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
668        self.manifest_digest
669    }
670    /// Return the migration identity.
671    pub const fn migration_id(&self) -> &MigrationId {
672        &self.migration_id
673    }
674    /// Return the group ordinal.
675    pub const fn group_ordinal(&self) -> u32 {
676        self.group_ordinal
677    }
678    /// Return the first group step index.
679    pub const fn first_step_index(&self) -> u32 {
680        self.first_step_index
681    }
682    /// Return the terminal delta step index.
683    pub const fn schema_delta_step_index(&self) -> u32 {
684        self.schema_delta_step_index
685    }
686    /// Return the exclusive group step end.
687    pub const fn end_step_index(&self) -> u32 {
688        self.end_step_index
689    }
690    /// Return the event kind.
691    pub const fn kind(&self) -> GroupJournalEventKind {
692        self.kind
693    }
694    /// Return exact target evidence carried only by committed events.
695    pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
696        self.observed_target.as_ref()
697    }
698}
699
700/// Identity-only applied-ledger record for one verified manifest.
701#[derive(Clone, Debug, Eq, PartialEq)]
702pub struct AppliedRecord {
703    scope: ExecutionScope,
704    fence: ExecutionFence,
705    migration_id: MigrationId,
706    manifest_digest: MigrationManifestDigest,
707    source_declared: ManagedDeclaredIdentityFingerprint,
708    target_declared: ManagedDeclaredIdentityFingerprint,
709    source_semantics: ManagedSemanticSchemaFingerprint,
710    target_semantics: ManagedSemanticSchemaFingerprint,
711}
712
713impl AppliedRecord {
714    /// Derive an applied-ledger record from one exact verified manifest.
715    pub fn from_verified_manifest(
716        lease: &MigrationLease,
717        migration: &VerifiedMigrationApplyManifest,
718    ) -> Result<Self, Diagnostic> {
719        Self::from_verified_manifest_contract(lease, migration.manifest())
720    }
721
722    /// Derive an applied-ledger record from one verified manifest contract.
723    ///
724    /// This is the reconstruction seam for persistent stores. The manifest
725    /// remains the trust anchor: its digest is recomputed from verified bytes,
726    /// and no persisted record claim enters the trusted value.
727    pub fn from_verified_manifest_contract(
728        lease: &MigrationLease,
729        manifest: &crate::VerifiedSchemaMigrationManifest,
730    ) -> Result<Self, Diagnostic> {
731        let source = manifest.source_state();
732        let target = manifest.target_state();
733        let scope = ExecutionScope::new(source.scope().id().clone());
734        if lease.scope() != &scope || target.scope() != source.scope() {
735            return Err(failure(
736                DiagnosticCategory::Integrity,
737                "migration_execution_scope_mismatch",
738                "lease and manifest endpoints must bind the same managed scope",
739            ));
740        }
741        Ok(Self {
742            scope,
743            fence: lease.fence(),
744            migration_id: manifest.id().clone(),
745            manifest_digest: crate::verified_manifest_digest(manifest)?,
746            source_declared: source.managed_declared_identity().clone(),
747            target_declared: target.managed_declared_identity().clone(),
748            source_semantics: source.managed_semantic_schema().clone(),
749            target_semantics: target.managed_semantic_schema().clone(),
750        })
751    }
752
753    /// Return the execution scope.
754    pub const fn scope(&self) -> &ExecutionScope {
755        &self.scope
756    }
757    /// Return the record fence.
758    pub const fn fence(&self) -> ExecutionFence {
759        self.fence
760    }
761    /// Return the migration identity.
762    pub const fn migration_id(&self) -> &MigrationId {
763        &self.migration_id
764    }
765    /// Return the canonical manifest digest.
766    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
767        self.manifest_digest
768    }
769    /// Return source managed-declared identity.
770    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
771        &self.source_declared
772    }
773    /// Return target managed-declared identity.
774    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
775        &self.target_declared
776    }
777    /// Return source managed semantics.
778    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
779        &self.source_semantics
780    }
781    /// Return target managed semantics.
782    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
783        &self.target_semantics
784    }
785}
786
787/// Identity-only journal record for one complete verified rollback plan.
788///
789/// The record binds the exact reverse-topological order, manifest digests,
790/// surviving applied set, and endpoint managed states the rollback executes
791/// under. It is the rollback analogue of [`PlanRecord`].
792#[derive(Clone, Debug, Eq, PartialEq)]
793pub struct RollbackPlanRecord {
794    scope: ExecutionScope,
795    fence: ExecutionFence,
796    source_applied: Vec<MigrationId>,
797    rollback_ids: Vec<MigrationId>,
798    manifest_digests: Vec<MigrationManifestDigest>,
799    manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
800    remaining_applied: Vec<MigrationId>,
801    source_declared: ManagedDeclaredIdentityFingerprint,
802    target_declared: ManagedDeclaredIdentityFingerprint,
803    source_semantics: ManagedSemanticSchemaFingerprint,
804    target_semantics: ManagedSemanticSchemaFingerprint,
805    semantic_profile: SemanticProfileFingerprint,
806    lowering_profile: SchemaLoweringProfileFingerprint,
807    observed_live_source: ManagedSemanticSchemaFingerprint,
808}
809
810impl RollbackPlanRecord {
811    /// Bind a fresh-lease rollback ledger and live-state precondition to
812    /// verified rollback plan identities.
813    pub fn from_verified_rollback_plan(
814        lease: &MigrationLease,
815        plan: &VerifiedMigrationRollbackPlan,
816        observed_applied_migrations: &[MigrationId],
817        observed_live_source: &ManagedSchemaState,
818    ) -> Result<Self, Diagnostic> {
819        let first = plan.rollbacks().first().ok_or_else(|| {
820            failure(
821                DiagnosticCategory::InvalidContract,
822                "migration_execution_empty_plan",
823                "an executable rollback plan requires at least one manifest",
824            )
825        })?;
826        let basis: Vec<MigrationId> = plan.applied_basis().into_iter().collect();
827        if observed_applied_migrations != basis {
828            return Err(failure(
829                DiagnosticCategory::Integrity,
830                "migration_execution_stale_applied_set",
831                "applied ledger changed after rollback planning; rebuild the plan",
832            ));
833        }
834        if observed_live_source != plan.source_state() {
835            return Err(failure(
836                DiagnosticCategory::Integrity,
837                "migration_execution_stale_source_state",
838                "live managed state differs from the planned source; rebuild the plan",
839            ));
840        }
841        let scope = ExecutionScope::new(plan.source_state().scope().id().clone());
842        if lease.scope() != &scope || plan.target_state().scope() != plan.source_state().scope() {
843            return Err(failure(
844                DiagnosticCategory::Integrity,
845                "migration_execution_scope_mismatch",
846                "lease, source, and target must bind the same managed scope",
847            ));
848        }
849        let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
850        let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
851        for rollback in plan.rollbacks() {
852            if rollback.manifest().managed_scope().id() != scope.managed_scope_id()
853                || rollback.manifest().semantic_profile().fingerprint() != &semantic_profile
854                || rollback.manifest().lowering_profile().fingerprint() != &lowering_profile
855            {
856                return Err(failure(
857                    DiagnosticCategory::Integrity,
858                    "migration_execution_plan_binding_mismatch",
859                    "planned rollbacks do not share exact scope and profile bindings",
860                ));
861            }
862        }
863        Ok(Self {
864            scope,
865            fence: lease.fence(),
866            source_applied: basis,
867            rollback_ids: plan
868                .rollbacks()
869                .iter()
870                .map(|rollback| rollback.manifest().id().clone())
871                .collect(),
872            manifest_digests: plan
873                .rollbacks()
874                .iter()
875                .map(|rollback| *rollback.digest())
876                .collect(),
877            manifest_plan_fingerprints: plan
878                .rollbacks()
879                .iter()
880                .map(|rollback| rollback.manifest().plan_fingerprint().clone())
881                .collect(),
882            remaining_applied: plan.remaining_applied().to_vec(),
883            source_declared: plan.source_state().managed_declared_identity().clone(),
884            target_declared: plan.target_state().managed_declared_identity().clone(),
885            source_semantics: plan.source_state().managed_semantic_schema().clone(),
886            target_semantics: plan.target_state().managed_semantic_schema().clone(),
887            semantic_profile,
888            lowering_profile,
889            observed_live_source: observed_live_source.managed_semantic_schema().clone(),
890        })
891    }
892
893    /// Return the execution scope.
894    pub const fn scope(&self) -> &ExecutionScope {
895        &self.scope
896    }
897    /// Return the fence bound into this record.
898    pub const fn fence(&self) -> ExecutionFence {
899        self.fence
900    }
901    /// Return the complete canonically ordered pre-rollback applied set.
902    pub fn source_applied(&self) -> &[MigrationId] {
903        &self.source_applied
904    }
905    /// Return rolled-back identities in reverse-topological execution order.
906    pub fn rollback_ids(&self) -> &[MigrationId] {
907        &self.rollback_ids
908    }
909    /// Return ordered canonical manifest digests.
910    pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
911        &self.manifest_digests
912    }
913    /// Return ordered manifest plan fingerprints.
914    pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
915        &self.manifest_plan_fingerprints
916    }
917    /// Return the applied identities that survive the rollback, in order.
918    pub fn remaining_applied(&self) -> &[MigrationId] {
919        &self.remaining_applied
920    }
921    /// Return planned pre-rollback managed-declared identity.
922    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
923        &self.source_declared
924    }
925    /// Return planned restored managed-declared identity.
926    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
927        &self.target_declared
928    }
929    /// Return planned pre-rollback managed semantics.
930    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
931        &self.source_semantics
932    }
933    /// Return planned restored managed semantics.
934    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
935        &self.target_semantics
936    }
937    /// Return semantic-profile content identity.
938    pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
939        &self.semantic_profile
940    }
941    /// Return lowering-registry content identity.
942    pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
943        &self.lowering_profile
944    }
945    /// Return the pre-mutation observed source semantics.
946    pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
947        &self.observed_live_source
948    }
949}
950
951/// Identity-only journal event for one verified rollback step transaction.
952///
953/// Each rollback step executes one recorded reverse program in its own
954/// provider transaction, so the commit-boundary vocabulary and recovery
955/// decision table are shared with forward [`GroupEventRecord`] events.
956#[derive(Clone, Debug, Eq, PartialEq)]
957pub struct RollbackStepEventRecord {
958    scope: ExecutionScope,
959    fence: ExecutionFence,
960    manifest_digest: MigrationManifestDigest,
961    migration_id: MigrationId,
962    step_ordinal: u32,
963    kind: GroupJournalEventKind,
964    observed_target: Option<ManagedSemanticSchemaFingerprint>,
965}
966
967impl RollbackStepEventRecord {
968    /// Derive a positional event from exact verified rollback evidence.
969    pub fn new(
970        lease: &MigrationLease,
971        rollback: &VerifiedMigrationRollbackManifest,
972        step_index: usize,
973        kind: GroupJournalEventKind,
974        observed_target: Option<ManagedSemanticSchemaFingerprint>,
975    ) -> Result<Self, Diagnostic> {
976        let step = rollback.steps().get(step_index).ok_or_else(|| {
977            failure(
978                DiagnosticCategory::Integrity,
979                "migration_execution_rollback_step_position",
980                "rollback event position is outside the verified rollback manifest",
981            )
982        })?;
983        let reverse = rollback.reverse_delta(step)?;
984        let scope = ExecutionScope::new(rollback.manifest().managed_scope().id().clone());
985        if lease.scope() != &scope {
986            return Err(failure(
987                DiagnosticCategory::Integrity,
988                "migration_execution_scope_mismatch",
989                "lease scope differs from the verified rollback scope",
990            ));
991        }
992        match kind {
993            GroupJournalEventKind::Committed
994                if observed_target.as_ref() == Some(reverse.target().managed_semantic_schema()) => {
995            }
996            GroupJournalEventKind::Committed => {
997                return Err(failure(
998                    DiagnosticCategory::Integrity,
999                    "migration_execution_commit_evidence_mismatch",
1000                    "committed event requires the exact observed target semantics",
1001                ));
1002            }
1003            GroupJournalEventKind::FormalOnlyAdvanced
1004                if observed_target.is_none()
1005                    && step.lowering().units().is_empty()
1006                    && reverse.source().managed_semantic_schema()
1007                        == reverse.target().managed_semantic_schema() => {}
1008            GroupJournalEventKind::FormalOnlyAdvanced => {
1009                return Err(failure(
1010                    DiagnosticCategory::InvalidContract,
1011                    "migration_execution_invalid_formal_advance",
1012                    "formal-only advancement requires an empty equal-semantic reverse program",
1013                ));
1014            }
1015            _ if observed_target.is_none() => {}
1016            _ => {
1017                return Err(failure(
1018                    DiagnosticCategory::InvalidContract,
1019                    "migration_execution_unexpected_observation",
1020                    "only committed events may carry observed target semantics",
1021                ));
1022            }
1023        }
1024        Ok(Self {
1025            scope,
1026            fence: lease.fence(),
1027            manifest_digest: *rollback.digest(),
1028            migration_id: rollback.manifest().id().clone(),
1029            step_ordinal: position(step_index)?,
1030            kind,
1031            observed_target,
1032        })
1033    }
1034
1035    /// Return the execution scope.
1036    pub const fn scope(&self) -> &ExecutionScope {
1037        &self.scope
1038    }
1039    /// Return the event fence.
1040    pub const fn fence(&self) -> ExecutionFence {
1041        self.fence
1042    }
1043    /// Return the canonical digest of the manifest being rolled back.
1044    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1045        self.manifest_digest
1046    }
1047    /// Return the migration identity.
1048    pub const fn migration_id(&self) -> &MigrationId {
1049        &self.migration_id
1050    }
1051    /// Return the rollback step position in execution order.
1052    pub const fn step_ordinal(&self) -> u32 {
1053        self.step_ordinal
1054    }
1055    /// Return the event kind.
1056    pub const fn kind(&self) -> GroupJournalEventKind {
1057        self.kind
1058    }
1059    /// Return exact target evidence carried only by committed events.
1060    pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
1061        self.observed_target.as_ref()
1062    }
1063}
1064
1065/// Identity-only retirement record for one rolled-back applied migration.
1066///
1067/// Retirement is append-only history: the applied record stays in the durable
1068/// journal, and this record marks it inactive. The states enter swapped
1069/// relative to [`AppliedRecord`] because rolling back moves the managed schema
1070/// from the manifest's target back to its source.
1071#[derive(Clone, Debug, Eq, PartialEq)]
1072pub struct RolledBackRecord {
1073    scope: ExecutionScope,
1074    fence: ExecutionFence,
1075    migration_id: MigrationId,
1076    manifest_digest: MigrationManifestDigest,
1077    source_declared: ManagedDeclaredIdentityFingerprint,
1078    target_declared: ManagedDeclaredIdentityFingerprint,
1079    source_semantics: ManagedSemanticSchemaFingerprint,
1080    target_semantics: ManagedSemanticSchemaFingerprint,
1081}
1082
1083impl RolledBackRecord {
1084    /// Derive a retirement record from one exact verified rollback manifest.
1085    pub fn from_verified_rollback(
1086        lease: &MigrationLease,
1087        rollback: &VerifiedMigrationRollbackManifest,
1088    ) -> Result<Self, Diagnostic> {
1089        Self::from_verified_manifest_contract(lease, rollback.manifest())
1090    }
1091
1092    /// Derive a retirement record from one verified manifest contract.
1093    ///
1094    /// This is the reconstruction seam for persistent stores, mirroring
1095    /// [`AppliedRecord::from_verified_manifest_contract`].
1096    pub fn from_verified_manifest_contract(
1097        lease: &MigrationLease,
1098        manifest: &crate::VerifiedSchemaMigrationManifest,
1099    ) -> Result<Self, Diagnostic> {
1100        let source = manifest.target_state();
1101        let target = manifest.source_state();
1102        let scope = ExecutionScope::new(source.scope().id().clone());
1103        if lease.scope() != &scope || target.scope() != source.scope() {
1104            return Err(failure(
1105                DiagnosticCategory::Integrity,
1106                "migration_execution_scope_mismatch",
1107                "lease and manifest endpoints must bind the same managed scope",
1108            ));
1109        }
1110        Ok(Self {
1111            scope,
1112            fence: lease.fence(),
1113            migration_id: manifest.id().clone(),
1114            manifest_digest: crate::verified_manifest_digest(manifest)?,
1115            source_declared: source.managed_declared_identity().clone(),
1116            target_declared: target.managed_declared_identity().clone(),
1117            source_semantics: source.managed_semantic_schema().clone(),
1118            target_semantics: target.managed_semantic_schema().clone(),
1119        })
1120    }
1121
1122    /// Return the execution scope.
1123    pub const fn scope(&self) -> &ExecutionScope {
1124        &self.scope
1125    }
1126    /// Return the record fence.
1127    pub const fn fence(&self) -> ExecutionFence {
1128        self.fence
1129    }
1130    /// Return the retired migration identity.
1131    pub const fn migration_id(&self) -> &MigrationId {
1132        &self.migration_id
1133    }
1134    /// Return the canonical manifest digest.
1135    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1136        self.manifest_digest
1137    }
1138    /// Return pre-rollback managed-declared identity.
1139    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1140        &self.source_declared
1141    }
1142    /// Return restored managed-declared identity.
1143    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1144        &self.target_declared
1145    }
1146    /// Return pre-rollback managed semantics.
1147    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1148        &self.source_semantics
1149    }
1150    /// Return restored managed semantics.
1151    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1152        &self.target_semantics
1153    }
1154}
1155
1156/// One trusted, still-open rollback plan and its ordered step events.
1157#[derive(Clone, Debug, Eq, PartialEq)]
1158pub struct OpenRollbackPlanRecord {
1159    plan: JournalEntry<RollbackPlanRecord>,
1160    events: Vec<JournalEntry<RollbackStepEventRecord>>,
1161}
1162
1163impl OpenRollbackPlanRecord {
1164    /// Rebuild store output while checking sequence, scope, fence, and manifest binding.
1165    ///
1166    /// The same monotonic-fence rule as [`OpenPlanRecord::from_store`] applies:
1167    /// recovery events may be written under later fences than the plan itself.
1168    pub fn from_store(
1169        plan: JournalEntry<RollbackPlanRecord>,
1170        events: Vec<JournalEntry<RollbackStepEventRecord>>,
1171    ) -> Result<Self, Diagnostic> {
1172        let mut previous = plan.sequence();
1173        let mut previous_fence = plan.record().fence();
1174        for event in &events {
1175            let manifest_index = plan
1176                .record()
1177                .manifest_digests()
1178                .iter()
1179                .position(|digest| digest == &event.record().manifest_digest());
1180            if event.sequence() <= previous
1181                || event.record().scope() != plan.record().scope()
1182                || event.record().fence() < previous_fence
1183                || manifest_index.is_none_or(|index| {
1184                    plan.record().rollback_ids().get(index) != Some(event.record().migration_id())
1185                })
1186            {
1187                return Err(failure(
1188                    DiagnosticCategory::Integrity,
1189                    "migration_execution_invalid_open_plan",
1190                    "loaded open-rollback events are not ordered and bound to the plan",
1191                ));
1192            }
1193            previous = event.sequence();
1194            previous_fence = event.record().fence();
1195        }
1196        Ok(Self { plan, events })
1197    }
1198
1199    /// Return the sequenced rollback plan record.
1200    pub const fn plan(&self) -> &JournalEntry<RollbackPlanRecord> {
1201        &self.plan
1202    }
1203
1204    /// Return ordered sequenced rollback step events.
1205    pub fn events(&self) -> &[JournalEntry<RollbackStepEventRecord>] {
1206        &self.events
1207    }
1208}
1209
1210/// Store-backed exclusive lease service.
1211pub trait MigrationLeaseStore: Send + Sync {
1212    /// Atomically acquire an available scope and issue a fence strictly greater
1213    /// than every fence previously issued for that scope.
1214    ///
1215    /// Lease expiry, liveness, and takeover policy are store-defined. Every
1216    /// takeover after release, expiry, or failure must issue a strictly greater
1217    /// fence so a surviving stale holder cannot mutate the journal.
1218    fn acquire<'a>(
1219        &'a self,
1220        scope: &'a ExecutionScope,
1221        holder: &'a LeaseHolderId,
1222    ) -> ExecutionFuture<'a, MigrationLease>;
1223
1224    /// Release only the exact currently active holder and fence.
1225    fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()>;
1226}
1227
1228/// Durable applied ledger and open-plan journal.
1229///
1230/// Every write must atomically reject a lease that is not the current active
1231/// scope, holder, and fence. There is deliberately no advisory validate method.
1232pub trait MigrationExecutionJournal: Send + Sync {
1233    /// Begin one fully verified plan after stale-ledger and live-source checks.
1234    ///
1235    /// An open plan of either direction is exclusive: the store must reject a
1236    /// new apply plan while a rollback plan is open and vice versa.
1237    fn begin_plan<'a>(
1238        &'a self,
1239        lease: &'a MigrationLease,
1240        record: PlanRecord,
1241    ) -> ExecutionFuture<'a, JournalEntry<PlanRecord>>;
1242
1243    /// Append one positional commit-boundary event under the active fence.
1244    fn record_group_event<'a>(
1245        &'a self,
1246        lease: &'a MigrationLease,
1247        record: GroupEventRecord,
1248    ) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>>;
1249
1250    /// Add one exact verified manifest from the open plan to the applied ledger.
1251    ///
1252    /// The store must reject a new record whose migration identity and digest do
1253    /// not occur at the same position in the open plan. Importing an existing
1254    /// ledger is a separate concern and must not use this execution write. An
1255    /// exact retry is idempotent only under the same fence; seeing the same
1256    /// migration under a newer fence requires reloading the ledger rather than
1257    /// treating differently fenced evidence as equal.
1258    fn record_applied<'a>(
1259        &'a self,
1260        lease: &'a MigrationLease,
1261        record: AppliedRecord,
1262    ) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>>;
1263
1264    /// Load ordered active applied records under the active fence.
1265    ///
1266    /// Retired records — applied records whose migration was rolled back by a
1267    /// later [`RolledBackRecord`] — are excluded. The full append-only history
1268    /// stays durable in the store; only the active ledger is the apply basis.
1269    fn load_applied<'a>(
1270        &'a self,
1271        lease: &'a MigrationLease,
1272    ) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>>;
1273
1274    /// Load the open plan and all ordered events written at its fence or later.
1275    fn load_open_plan<'a>(
1276        &'a self,
1277        lease: &'a MigrationLease,
1278    ) -> ExecutionFuture<'a, Option<OpenPlanRecord>>;
1279
1280    /// Begin one fully verified rollback plan after stale-ledger checks.
1281    ///
1282    /// Subject to the same open-plan exclusivity as [`Self::begin_plan`].
1283    fn begin_rollback_plan<'a>(
1284        &'a self,
1285        lease: &'a MigrationLease,
1286        record: RollbackPlanRecord,
1287    ) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>>;
1288
1289    /// Append one positional rollback commit-boundary event under the active fence.
1290    fn record_rollback_step_event<'a>(
1291        &'a self,
1292        lease: &'a MigrationLease,
1293        record: RollbackStepEventRecord,
1294    ) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>>;
1295
1296    /// Retire one exact applied record from the open rollback plan.
1297    ///
1298    /// The store must reject a record whose migration identity and digest do
1299    /// not occur at the same position in the open rollback plan, or whose
1300    /// migration is not currently active in the applied ledger. Retirement is
1301    /// append-only: the applied record and this record both stay durable, and
1302    /// the migration merely leaves the active ledger. An exact retry is
1303    /// idempotent only under the same fence.
1304    fn record_rolled_back<'a>(
1305        &'a self,
1306        lease: &'a MigrationLease,
1307        record: RolledBackRecord,
1308    ) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>>;
1309
1310    /// Load ordered retirement records under the active fence.
1311    fn load_rolled_back<'a>(
1312        &'a self,
1313        lease: &'a MigrationLease,
1314    ) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>>;
1315
1316    /// Load the open rollback plan and all ordered events at its fence or later.
1317    fn load_open_rollback_plan<'a>(
1318        &'a self,
1319        lease: &'a MigrationLease,
1320    ) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>>;
1321}
1322
1323/// Filter an applied ledger down to its active records.
1324///
1325/// Each retirement record consumes the latest not-yet-retired applied record
1326/// with the same migration identity and manifest digest written before it.
1327/// A retirement that matches no applied record is corrupt history and fails
1328/// closed. Stores share this exact matching rule so every journal
1329/// implementation reports the same active basis.
1330pub fn active_applied_entries(
1331    applied: Vec<JournalEntry<AppliedRecord>>,
1332    rolled_back: &[JournalEntry<RolledBackRecord>],
1333) -> Result<Vec<JournalEntry<AppliedRecord>>, Diagnostic> {
1334    let mut retired = vec![false; applied.len()];
1335    for retirement in rolled_back {
1336        let matched = applied
1337            .iter()
1338            .enumerate()
1339            .filter(|(index, entry)| {
1340                !retired[*index]
1341                    && entry.sequence() < retirement.sequence()
1342                    && entry.record().migration_id() == retirement.record().migration_id()
1343                    && entry.record().manifest_digest() == retirement.record().manifest_digest()
1344            })
1345            .max_by_key(|(_, entry)| entry.sequence())
1346            .map(|(index, _)| index);
1347        let Some(index) = matched else {
1348            return Err(failure(
1349                DiagnosticCategory::Integrity,
1350                "migration_execution_unmatched_retirement",
1351                "retirement record matches no active applied record before it",
1352            ));
1353        };
1354        retired[index] = true;
1355    }
1356    Ok(applied
1357        .into_iter()
1358        .zip(retired)
1359        .filter_map(|(entry, retired)| (!retired).then_some(entry))
1360        .collect())
1361}
1362
1363fn position(value: usize) -> Result<u32, Diagnostic> {
1364    u32::try_from(value).map_err(|_| {
1365        failure(
1366            DiagnosticCategory::ResourceLimit,
1367            "migration_execution_position_limit",
1368            "migration transaction position exceeds the canonical u32 range",
1369        )
1370    })
1371}
1372
1373fn failure(category: DiagnosticCategory, code: &'static str, message: &'static str) -> Diagnostic {
1374    Diagnostic::new(
1375        category,
1376        DiagnosticCode::new(code).expect("static migration execution diagnostic code"),
1377        message,
1378    )
1379}
1380
1381#[cfg(test)]
1382mod tests {
1383    use std::collections::BTreeMap;
1384    use std::future::Future;
1385    use std::sync::{Arc, Mutex};
1386    use std::task::{Context, Poll, Wake, Waker};
1387
1388    use type_bridge_contract::fingerprint::SemanticProfileId;
1389    use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileBinding};
1390    use type_bridge_contract::migration::{MigrationAppLabel, MigrationName};
1391
1392    use super::*;
1393    use crate::schema_lowering_profile_binding;
1394
1395    #[derive(Default)]
1396    struct ScopeState {
1397        highest_fence: u64,
1398        active_lease: Option<MigrationLease>,
1399        next_sequence: u64,
1400        applied: Vec<JournalEntry<AppliedRecord>>,
1401        rolled_back: Vec<JournalEntry<RolledBackRecord>>,
1402        open_plan: Option<JournalEntry<PlanRecord>>,
1403        events: Vec<JournalEntry<GroupEventRecord>>,
1404        open_rollback_plan: Option<JournalEntry<RollbackPlanRecord>>,
1405        rollback_events: Vec<JournalEntry<RollbackStepEventRecord>>,
1406    }
1407
1408    #[derive(Default)]
1409    struct InMemoryStore {
1410        scopes: Mutex<BTreeMap<ExecutionScope, ScopeState>>,
1411    }
1412
1413    impl InMemoryStore {
1414        fn check_lease<'a>(
1415            state: &'a mut ScopeState,
1416            lease: &MigrationLease,
1417        ) -> Result<&'a mut ScopeState, Diagnostic> {
1418            if state.active_lease.as_ref() != Some(lease) {
1419                return Err(failure(
1420                    DiagnosticCategory::Integrity,
1421                    "migration_execution_stale_fence",
1422                    "journal write does not carry the current active lease and fence",
1423                ));
1424            }
1425            Ok(state)
1426        }
1427
1428        fn sequence(state: &mut ScopeState) -> Result<JournalSequence, Diagnostic> {
1429            state.next_sequence = state.next_sequence.checked_add(1).ok_or_else(|| {
1430                failure(
1431                    DiagnosticCategory::ResourceLimit,
1432                    "migration_execution_sequence_exhausted",
1433                    "journal sequence range is exhausted",
1434                )
1435            })?;
1436            JournalSequence::new(state.next_sequence)
1437        }
1438    }
1439
1440    impl MigrationLeaseStore for InMemoryStore {
1441        fn acquire<'a>(
1442            &'a self,
1443            scope: &'a ExecutionScope,
1444            holder: &'a LeaseHolderId,
1445        ) -> ExecutionFuture<'a, MigrationLease> {
1446            Box::pin(async move {
1447                let mut scopes = self.scopes.lock().expect("store mutex");
1448                let state = scopes.entry(scope.clone()).or_default();
1449                if state.active_lease.is_some() {
1450                    return Err(failure(
1451                        DiagnosticCategory::InvalidContract,
1452                        "migration_execution_lease_contended",
1453                        "migration scope already has an active lease",
1454                    ));
1455                }
1456                let fence = if state.highest_fence == 0 {
1457                    ExecutionFence::new(1)?
1458                } else {
1459                    ExecutionFence::new(state.highest_fence)?.checked_successor()?
1460                };
1461                state.highest_fence = fence.get();
1462                let lease = MigrationLease::new(scope.clone(), holder.clone(), fence);
1463                state.active_lease = Some(lease.clone());
1464                Ok(lease)
1465            })
1466        }
1467
1468        fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()> {
1469            Box::pin(async move {
1470                let mut scopes = self.scopes.lock().expect("store mutex");
1471                let state = scopes.get_mut(lease.scope()).ok_or_else(|| {
1472                    failure(
1473                        DiagnosticCategory::Integrity,
1474                        "migration_execution_stale_fence",
1475                        "lease scope is not active",
1476                    )
1477                })?;
1478                Self::check_lease(state, lease)?;
1479                state.active_lease = None;
1480                Ok(())
1481            })
1482        }
1483    }
1484
1485    impl MigrationExecutionJournal for InMemoryStore {
1486        fn begin_plan<'a>(
1487            &'a self,
1488            lease: &'a MigrationLease,
1489            record: PlanRecord,
1490        ) -> ExecutionFuture<'a, JournalEntry<PlanRecord>> {
1491            Box::pin(async move {
1492                let mut scopes = self.scopes.lock().expect("store mutex");
1493                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1494                Self::check_lease(state, lease)?;
1495                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1496                    return Err(stale_fence());
1497                }
1498                if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
1499                    return Err(failure(
1500                        DiagnosticCategory::InvalidContract,
1501                        "migration_execution_plan_already_open",
1502                        "migration scope already has an open plan",
1503                    ));
1504                }
1505                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1506                state.open_plan = Some(entry.clone());
1507                Ok(entry)
1508            })
1509        }
1510
1511        fn record_group_event<'a>(
1512            &'a self,
1513            lease: &'a MigrationLease,
1514            record: GroupEventRecord,
1515        ) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>> {
1516            Box::pin(async move {
1517                let mut scopes = self.scopes.lock().expect("store mutex");
1518                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1519                Self::check_lease(state, lease)?;
1520                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1521                    return Err(stale_fence());
1522                }
1523                let plan = state.open_plan.as_ref().ok_or_else(|| {
1524                    failure(
1525                        DiagnosticCategory::InvalidContract,
1526                        "migration_execution_no_open_plan",
1527                        "group event requires an open migration plan",
1528                    )
1529                })?;
1530                if !plan
1531                    .record()
1532                    .manifest_digests()
1533                    .contains(&record.manifest_digest())
1534                {
1535                    return Err(failure(
1536                        DiagnosticCategory::Integrity,
1537                        "migration_execution_foreign_event",
1538                        "group event manifest is absent from the open plan",
1539                    ));
1540                }
1541                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1542                state.events.push(entry.clone());
1543                Ok(entry)
1544            })
1545        }
1546
1547        fn record_applied<'a>(
1548            &'a self,
1549            lease: &'a MigrationLease,
1550            record: AppliedRecord,
1551        ) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>> {
1552            Box::pin(async move {
1553                let mut scopes = self.scopes.lock().expect("store mutex");
1554                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1555                Self::check_lease(state, lease)?;
1556                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1557                    return Err(stale_fence());
1558                }
1559                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
1560                if let Some(existing) = active
1561                    .iter()
1562                    .find(|entry| entry.record().migration_id() == record.migration_id())
1563                {
1564                    return if existing.record() == &record {
1565                        Ok(existing.clone())
1566                    } else {
1567                        Err(failure(
1568                            DiagnosticCategory::Integrity,
1569                            "migration_execution_applied_identity_conflict",
1570                            "applied migration identity has different evidence",
1571                        ))
1572                    };
1573                }
1574                let plan = state.open_plan.as_ref().ok_or_else(|| {
1575                    failure(
1576                        DiagnosticCategory::InvalidContract,
1577                        "migration_execution_no_open_plan",
1578                        "applied migration requires an open migration plan",
1579                    )
1580                })?;
1581                let manifest_index = plan
1582                    .record()
1583                    .migration_ids()
1584                    .iter()
1585                    .position(|id| id == record.migration_id());
1586                if manifest_index.is_none_or(|index| {
1587                    plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
1588                }) {
1589                    return Err(failure(
1590                        DiagnosticCategory::Integrity,
1591                        "migration_execution_foreign_applied_record",
1592                        "applied migration identity and digest are absent from the open plan",
1593                    ));
1594                }
1595                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1596                state.applied.push(entry.clone());
1597                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
1598                let complete = state.open_plan.as_ref().is_some_and(|plan| {
1599                    plan.record().migration_ids().iter().all(|id| {
1600                        active
1601                            .iter()
1602                            .any(|applied| applied.record().migration_id() == id)
1603                    })
1604                });
1605                if complete {
1606                    state.open_plan = None;
1607                    state.events.clear();
1608                }
1609                Ok(entry)
1610            })
1611        }
1612
1613        fn load_applied<'a>(
1614            &'a self,
1615            lease: &'a MigrationLease,
1616        ) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>> {
1617            Box::pin(async move {
1618                let mut scopes = self.scopes.lock().expect("store mutex");
1619                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1620                Self::check_lease(state, lease)?;
1621                active_applied_entries(state.applied.clone(), &state.rolled_back)
1622            })
1623        }
1624
1625        fn load_open_plan<'a>(
1626            &'a self,
1627            lease: &'a MigrationLease,
1628        ) -> ExecutionFuture<'a, Option<OpenPlanRecord>> {
1629            Box::pin(async move {
1630                let mut scopes = self.scopes.lock().expect("store mutex");
1631                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1632                Self::check_lease(state, lease)?;
1633                let Some(plan) = state.open_plan.clone() else {
1634                    return Ok(None);
1635                };
1636                Ok(Some(OpenPlanRecord::from_store(
1637                    plan,
1638                    state.events.clone(),
1639                )?))
1640            })
1641        }
1642
1643        fn begin_rollback_plan<'a>(
1644            &'a self,
1645            lease: &'a MigrationLease,
1646            record: RollbackPlanRecord,
1647        ) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>> {
1648            Box::pin(async move {
1649                let mut scopes = self.scopes.lock().expect("store mutex");
1650                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1651                Self::check_lease(state, lease)?;
1652                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1653                    return Err(stale_fence());
1654                }
1655                if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
1656                    return Err(failure(
1657                        DiagnosticCategory::InvalidContract,
1658                        "migration_execution_plan_already_open",
1659                        "migration scope already has an open plan",
1660                    ));
1661                }
1662                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1663                state.open_rollback_plan = Some(entry.clone());
1664                Ok(entry)
1665            })
1666        }
1667
1668        fn record_rollback_step_event<'a>(
1669            &'a self,
1670            lease: &'a MigrationLease,
1671            record: RollbackStepEventRecord,
1672        ) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>> {
1673            Box::pin(async move {
1674                let mut scopes = self.scopes.lock().expect("store mutex");
1675                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1676                Self::check_lease(state, lease)?;
1677                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1678                    return Err(stale_fence());
1679                }
1680                let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
1681                    failure(
1682                        DiagnosticCategory::InvalidContract,
1683                        "migration_execution_no_open_plan",
1684                        "rollback event requires an open rollback plan",
1685                    )
1686                })?;
1687                if !plan
1688                    .record()
1689                    .manifest_digests()
1690                    .contains(&record.manifest_digest())
1691                {
1692                    return Err(failure(
1693                        DiagnosticCategory::Integrity,
1694                        "migration_execution_foreign_event",
1695                        "rollback event manifest is absent from the open plan",
1696                    ));
1697                }
1698                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1699                state.rollback_events.push(entry.clone());
1700                Ok(entry)
1701            })
1702        }
1703
1704        fn record_rolled_back<'a>(
1705            &'a self,
1706            lease: &'a MigrationLease,
1707            record: RolledBackRecord,
1708        ) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>> {
1709            Box::pin(async move {
1710                let mut scopes = self.scopes.lock().expect("store mutex");
1711                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1712                Self::check_lease(state, lease)?;
1713                if record.scope() != lease.scope() || record.fence() != lease.fence() {
1714                    return Err(stale_fence());
1715                }
1716                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
1717                let is_active = active.iter().any(|entry| {
1718                    entry.record().migration_id() == record.migration_id()
1719                        && entry.record().manifest_digest() == record.manifest_digest()
1720                });
1721                if !is_active {
1722                    if let Some(existing) = state.rolled_back.iter().find(|entry| {
1723                        entry.record() == &record && entry.record().fence() == lease.fence()
1724                    }) {
1725                        return Ok(existing.clone());
1726                    }
1727                    return Err(failure(
1728                        DiagnosticCategory::Integrity,
1729                        "migration_execution_retirement_conflict",
1730                        "retirement target is not active in the applied ledger",
1731                    ));
1732                }
1733                let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
1734                    failure(
1735                        DiagnosticCategory::InvalidContract,
1736                        "migration_execution_no_open_plan",
1737                        "retirement requires an open rollback plan",
1738                    )
1739                })?;
1740                let manifest_index = plan
1741                    .record()
1742                    .rollback_ids()
1743                    .iter()
1744                    .position(|id| id == record.migration_id());
1745                if manifest_index.is_none_or(|index| {
1746                    plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
1747                }) {
1748                    return Err(failure(
1749                        DiagnosticCategory::Integrity,
1750                        "migration_execution_foreign_applied_record",
1751                        "retirement identity and digest are absent from the open plan",
1752                    ));
1753                }
1754                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
1755                state.rolled_back.push(entry.clone());
1756                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
1757                let complete = state.open_rollback_plan.as_ref().is_some_and(|plan| {
1758                    plan.record().rollback_ids().iter().all(|id| {
1759                        !active
1760                            .iter()
1761                            .any(|applied| applied.record().migration_id() == id)
1762                    })
1763                });
1764                if complete {
1765                    state.open_rollback_plan = None;
1766                    state.rollback_events.clear();
1767                }
1768                Ok(entry)
1769            })
1770        }
1771
1772        fn load_rolled_back<'a>(
1773            &'a self,
1774            lease: &'a MigrationLease,
1775        ) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>> {
1776            Box::pin(async move {
1777                let mut scopes = self.scopes.lock().expect("store mutex");
1778                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1779                Self::check_lease(state, lease)?;
1780                Ok(state.rolled_back.clone())
1781            })
1782        }
1783
1784        fn load_open_rollback_plan<'a>(
1785            &'a self,
1786            lease: &'a MigrationLease,
1787        ) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>> {
1788            Box::pin(async move {
1789                let mut scopes = self.scopes.lock().expect("store mutex");
1790                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
1791                Self::check_lease(state, lease)?;
1792                let Some(plan) = state.open_rollback_plan.clone() else {
1793                    return Ok(None);
1794                };
1795                Ok(Some(OpenRollbackPlanRecord::from_store(
1796                    plan,
1797                    state.rollback_events.clone(),
1798                )?))
1799            })
1800        }
1801    }
1802
1803    fn stale_fence() -> Diagnostic {
1804        failure(
1805            DiagnosticCategory::Integrity,
1806            "migration_execution_stale_fence",
1807            "journal write does not carry the current active lease and fence",
1808        )
1809    }
1810
1811    fn scope() -> ExecutionScope {
1812        ExecutionScope::new(ManagedScopeId::new("journal-test").expect("scope"))
1813    }
1814
1815    fn migration_id() -> MigrationId {
1816        MigrationId::from_components(
1817            MigrationAppLabel::new("example").expect("app"),
1818            MigrationName::new("0001_initial").expect("name"),
1819        )
1820    }
1821
1822    fn semantic_fingerprint(bytes: &[u8]) -> ManagedSemanticSchemaFingerprint {
1823        ManagedSemanticSchemaFingerprint::compute(
1824            SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
1825            bytes,
1826        )
1827        .expect("semantic fingerprint")
1828    }
1829
1830    fn fake_plan(lease: &MigrationLease) -> PlanRecord {
1831        let id = migration_id();
1832        let semantic = SemanticProfileBinding::resolve(
1833            SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
1834        )
1835        .expect("semantic binding");
1836        PlanRecord {
1837            scope: lease.scope().clone(),
1838            fence: lease.fence(),
1839            source_applied: Vec::new(),
1840            source_frontier: Vec::new(),
1841            target_frontier: vec![id.clone()],
1842            migration_ids: vec![id],
1843            manifest_digests: vec![MigrationManifestDigest::compute(b"manifest")],
1844            manifest_plan_fingerprints: vec![
1845                MigrationPlanFingerprint::compute(&[]).expect("plan fingerprint"),
1846            ],
1847            source_declared: ManagedDeclaredIdentityFingerprint::compute(b"source")
1848                .expect("source declared"),
1849            target_declared: ManagedDeclaredIdentityFingerprint::compute(b"target")
1850                .expect("target declared"),
1851            source_semantics: semantic_fingerprint(b"source"),
1852            target_semantics: semantic_fingerprint(b"target"),
1853            semantic_profile: semantic.fingerprint().clone(),
1854            lowering_profile: schema_lowering_profile_binding()
1855                .expect("lowering binding")
1856                .fingerprint()
1857                .clone(),
1858            observed_live_source: semantic_fingerprint(b"source"),
1859        }
1860    }
1861
1862    fn fake_event(lease: &MigrationLease, plan: &PlanRecord) -> GroupEventRecord {
1863        GroupEventRecord {
1864            scope: lease.scope().clone(),
1865            fence: lease.fence(),
1866            manifest_digest: plan.manifest_digests()[0],
1867            migration_id: plan.migration_ids()[0].clone(),
1868            group_ordinal: 0,
1869            first_step_index: 0,
1870            schema_delta_step_index: 0,
1871            end_step_index: 1,
1872            kind: GroupJournalEventKind::BeforeCommit,
1873            observed_target: None,
1874        }
1875    }
1876
1877    fn fake_applied(lease: &MigrationLease, plan: &PlanRecord) -> AppliedRecord {
1878        AppliedRecord {
1879            scope: lease.scope().clone(),
1880            fence: lease.fence(),
1881            migration_id: plan.migration_ids()[0].clone(),
1882            manifest_digest: plan.manifest_digests()[0],
1883            source_declared: plan.source_declared().clone(),
1884            target_declared: plan.target_declared().clone(),
1885            source_semantics: plan.source_semantics().clone(),
1886            target_semantics: plan.target_semantics().clone(),
1887        }
1888    }
1889
1890    #[test]
1891    fn store_assigns_monotonic_sequences_and_rejects_every_stale_fence_write() {
1892        let store = InMemoryStore::default();
1893        let scope = scope();
1894        let holder_a = LeaseHolderId::new("owner-a").expect("holder");
1895        let holder_b = LeaseHolderId::new("owner-b").expect("holder");
1896        let lease_a = block_on(store.acquire(&scope, &holder_a)).expect("lease a");
1897        assert!(block_on(store.acquire(&scope, &holder_b)).is_err());
1898        let plan_a = fake_plan(&lease_a);
1899        let plan_entry = block_on(store.begin_plan(&lease_a, plan_a.clone())).expect("plan");
1900        let event_entry =
1901            block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a)))
1902                .expect("event");
1903        assert_eq!(plan_entry.sequence().get(), 1);
1904        assert_eq!(event_entry.sequence().get(), 2);
1905        block_on(store.release(&lease_a)).expect("release a");
1906        let lease_b = block_on(store.acquire(&scope, &holder_b)).expect("lease b");
1907        assert!(lease_b.fence() > lease_a.fence());
1908        assert!(block_on(store.begin_plan(&lease_a, plan_a.clone())).is_err());
1909        assert!(
1910            block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a),)).is_err()
1911        );
1912        assert!(
1913            block_on(store.record_applied(&lease_a, fake_applied(&lease_a, &plan_a),)).is_err()
1914        );
1915        assert!(block_on(store.release(&lease_a)).is_err());
1916        let recovery_event =
1917            block_on(store.record_group_event(&lease_b, fake_event(&lease_b, &plan_a)))
1918                .expect("recovery event");
1919        assert_eq!(recovery_event.sequence().get(), 3);
1920        block_on(store.release(&lease_b)).expect("release recovery lease");
1921        let holder_c = LeaseHolderId::new("owner-c").expect("holder");
1922        let lease_c = block_on(store.acquire(&scope, &holder_c)).expect("lease c");
1923        assert!(lease_c.fence() > lease_b.fence());
1924        let open = block_on(store.load_open_plan(&lease_c))
1925            .expect("open plan")
1926            .expect("plan remains open after recovery crash");
1927        assert_eq!(open.events().len(), 2);
1928        assert_eq!(open.events()[0].record().fence(), lease_a.fence());
1929        assert_eq!(open.events()[1].record().fence(), lease_b.fence());
1930    }
1931
1932    #[test]
1933    fn applied_records_are_plan_bound_idempotent_and_close_the_completed_plan() {
1934        let store = InMemoryStore::default();
1935        let scope = scope();
1936        let holder = LeaseHolderId::new("applied-owner").expect("holder");
1937        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
1938        let plan = fake_plan(&lease);
1939        let applied = fake_applied(&lease, &plan);
1940
1941        assert!(block_on(store.record_applied(&lease, applied.clone())).is_err());
1942        block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
1943
1944        let mut foreign = applied.clone();
1945        foreign.manifest_digest = MigrationManifestDigest::compute(b"foreign");
1946        assert!(block_on(store.record_applied(&lease, foreign)).is_err());
1947
1948        let first = block_on(store.record_applied(&lease, applied.clone()))
1949            .expect("apply planned manifest");
1950        let duplicate = block_on(store.record_applied(&lease, applied))
1951            .expect("same-fence duplicate is idempotent");
1952        assert_eq!(duplicate, first);
1953        assert!(
1954            block_on(store.load_open_plan(&lease))
1955                .expect("load completed plan")
1956                .is_none()
1957        );
1958        assert_eq!(
1959            block_on(store.load_applied(&lease)).expect("load applied ledger"),
1960            vec![first],
1961        );
1962    }
1963
1964    fn fake_rollback_plan(lease: &MigrationLease, plan: &PlanRecord) -> RollbackPlanRecord {
1965        RollbackPlanRecord {
1966            scope: lease.scope().clone(),
1967            fence: lease.fence(),
1968            source_applied: plan.migration_ids().to_vec(),
1969            rollback_ids: plan.migration_ids().to_vec(),
1970            manifest_digests: plan.manifest_digests().to_vec(),
1971            manifest_plan_fingerprints: plan.manifest_plan_fingerprints().to_vec(),
1972            remaining_applied: Vec::new(),
1973            source_declared: plan.target_declared().clone(),
1974            target_declared: plan.source_declared().clone(),
1975            source_semantics: plan.target_semantics().clone(),
1976            target_semantics: plan.source_semantics().clone(),
1977            semantic_profile: plan.semantic_profile().clone(),
1978            lowering_profile: plan.lowering_profile().clone(),
1979            observed_live_source: plan.target_semantics().clone(),
1980        }
1981    }
1982
1983    fn fake_rolled_back(lease: &MigrationLease, plan: &PlanRecord) -> RolledBackRecord {
1984        RolledBackRecord {
1985            scope: lease.scope().clone(),
1986            fence: lease.fence(),
1987            migration_id: plan.migration_ids()[0].clone(),
1988            manifest_digest: plan.manifest_digests()[0],
1989            source_declared: plan.target_declared().clone(),
1990            target_declared: plan.source_declared().clone(),
1991            source_semantics: plan.target_semantics().clone(),
1992            target_semantics: plan.source_semantics().clone(),
1993        }
1994    }
1995
1996    #[test]
1997    fn retirement_is_append_only_exclusive_and_reopens_the_identity() {
1998        let store = InMemoryStore::default();
1999        let scope = scope();
2000        let holder = LeaseHolderId::new("retirement-owner").expect("holder");
2001        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
2002        let plan = fake_plan(&lease);
2003        let rollback_plan = fake_rollback_plan(&lease, &plan);
2004        let retirement = fake_rolled_back(&lease, &plan);
2005
2006        // Apply the migration through an ordinary forward plan.
2007        block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
2008        assert!(
2009            block_on(store.begin_rollback_plan(&lease, rollback_plan.clone())).is_err(),
2010            "open plans of either direction must be exclusive"
2011        );
2012        block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
2013            .expect("apply planned manifest");
2014        assert_eq!(
2015            block_on(store.load_applied(&lease)).expect("active").len(),
2016            1
2017        );
2018
2019        // Retiring outside an open rollback plan fails closed.
2020        assert!(block_on(store.record_rolled_back(&lease, retirement.clone())).is_err());
2021        block_on(store.begin_rollback_plan(&lease, rollback_plan)).expect("open rollback plan");
2022        assert!(
2023            block_on(store.begin_plan(&lease, plan.clone())).is_err(),
2024            "an open rollback plan must block a new apply plan"
2025        );
2026        let first = block_on(store.record_rolled_back(&lease, retirement.clone()))
2027            .expect("retire applied migration");
2028        let duplicate = block_on(store.record_rolled_back(&lease, retirement))
2029            .expect("same-fence duplicate retirement is idempotent");
2030        assert_eq!(duplicate, first);
2031        assert!(
2032            block_on(store.load_open_rollback_plan(&lease))
2033                .expect("closed rollback plan")
2034                .is_none()
2035        );
2036        assert!(
2037            block_on(store.load_applied(&lease))
2038                .expect("active")
2039                .is_empty()
2040        );
2041        assert_eq!(
2042            block_on(store.load_rolled_back(&lease))
2043                .expect("retired")
2044                .len(),
2045            1,
2046        );
2047
2048        // The identity is free again: a fresh forward plan re-applies it.
2049        block_on(store.begin_plan(&lease, plan.clone())).expect("reopen plan");
2050        block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
2051            .expect("re-apply retired migration");
2052        assert_eq!(
2053            block_on(store.load_applied(&lease)).expect("active").len(),
2054            1
2055        );
2056    }
2057
2058    #[test]
2059    fn unmatched_retirements_are_corrupt_history() {
2060        let store = InMemoryStore::default();
2061        let scope = scope();
2062        let holder = LeaseHolderId::new("orphan-owner").expect("holder");
2063        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
2064        let plan = fake_plan(&lease);
2065        let orphan = JournalEntry::from_store(
2066            JournalSequence::new(7).expect("sequence"),
2067            fake_rolled_back(&lease, &plan),
2068        );
2069        let error = active_applied_entries(Vec::new(), &[orphan])
2070            .expect_err("a retirement without its applied record is corrupt");
2071        assert_eq!(
2072            error.code().as_str(),
2073            "migration_execution_unmatched_retirement"
2074        );
2075    }
2076
2077    #[test]
2078    fn recovery_decision_table_is_exhaustive_for_distinct_and_equal_semantics() {
2079        let source = semantic_fingerprint(b"source");
2080        let target = semantic_fingerprint(b"target");
2081        let neither = semantic_fingerprint(b"neither");
2082        let events = [
2083            None,
2084            Some(GroupJournalEventKind::BeforeCommit),
2085            Some(GroupJournalEventKind::Committed),
2086            Some(GroupJournalEventKind::CommitOutcomeUnknown),
2087            Some(GroupJournalEventKind::DefinitelyAborted),
2088            Some(GroupJournalEventKind::FormalOnlyAdvanced),
2089        ];
2090        let observations = [
2091            GroupRecoveryObservation::ManagedSemantics(source.clone()),
2092            GroupRecoveryObservation::ManagedSemantics(target.clone()),
2093            GroupRecoveryObservation::ManagedSemantics(neither.clone()),
2094            GroupRecoveryObservation::Unavailable,
2095        ];
2096        for event in events {
2097            for observation in &observations {
2098                let expected = expected_distinct(event, observation, &source, &target);
2099                assert_eq!(
2100                    decide_group_recovery(event, observation, &source, &target),
2101                    expected,
2102                    "distinct case {event:?} {observation:?}",
2103                );
2104            }
2105        }
2106        let equal_observations = [
2107            GroupRecoveryObservation::ManagedSemantics(source.clone()),
2108            GroupRecoveryObservation::ManagedSemantics(neither),
2109            GroupRecoveryObservation::Unavailable,
2110        ];
2111        for event in events {
2112            for observation in &equal_observations {
2113                let expected = expected_equal(event, observation, &source);
2114                assert_eq!(
2115                    decide_group_recovery(event, observation, &source, &source),
2116                    expected,
2117                    "equal case {event:?} {observation:?}",
2118                );
2119            }
2120        }
2121    }
2122
2123    fn expected_distinct(
2124        event: Option<GroupJournalEventKind>,
2125        observation: &GroupRecoveryObservation,
2126        source: &ManagedSemanticSchemaFingerprint,
2127        target: &ManagedSemanticSchemaFingerprint,
2128    ) -> GroupRecoveryDecision {
2129        let source_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == source);
2130        let target_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == target);
2131        match event {
2132            None | Some(GroupJournalEventKind::DefinitelyAborted) if source_seen => {
2133                GroupRecoveryDecision::ExecuteNormally
2134            }
2135            Some(
2136                GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown,
2137            ) if source_seen => GroupRecoveryDecision::ExecuteNormally,
2138            Some(
2139                GroupJournalEventKind::BeforeCommit
2140                | GroupJournalEventKind::CommitOutcomeUnknown
2141                | GroupJournalEventKind::Committed,
2142            ) if target_seen => GroupRecoveryDecision::RepairCheckpoint,
2143            _ => GroupRecoveryDecision::RequiresExplicitRecovery,
2144        }
2145    }
2146
2147    fn expected_equal(
2148        event: Option<GroupJournalEventKind>,
2149        observation: &GroupRecoveryObservation,
2150        both: &ManagedSemanticSchemaFingerprint,
2151    ) -> GroupRecoveryDecision {
2152        let both_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == both);
2153        if !both_seen {
2154            return GroupRecoveryDecision::RequiresExplicitRecovery;
2155        }
2156        match event {
2157            None | Some(GroupJournalEventKind::DefinitelyAborted) => {
2158                GroupRecoveryDecision::ExecuteNormally
2159            }
2160            Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced) => {
2161                GroupRecoveryDecision::RepairCheckpoint
2162            }
2163            _ => GroupRecoveryDecision::RequiresExplicitRecovery,
2164        }
2165    }
2166
2167    struct NoopWake;
2168
2169    impl Wake for NoopWake {
2170        fn wake(self: Arc<Self>) {}
2171    }
2172
2173    fn block_on<F: Future>(future: F) -> F::Output {
2174        let waker = Waker::from(Arc::new(NoopWake));
2175        let mut context = Context::from_waker(&waker);
2176        let mut future = Box::pin(future);
2177        loop {
2178            match future.as_mut().poll(&mut context) {
2179                Poll::Ready(output) => return output,
2180                Poll::Pending => std::thread::yield_now(),
2181            }
2182        }
2183    }
2184}