Skip to main content

type_bridge_schema_migration/
execution.rs

1//! Provider-neutral fenced migration journal and recovery contracts.
2
3use std::fmt;
4use std::future::Future;
5use std::pin::Pin;
6use std::sync::Arc;
7use std::sync::atomic::{AtomicBool, Ordering};
8use std::task::{Poll, Waker};
9use std::time::Instant;
10
11use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
12use type_bridge_contract::fingerprint::Fingerprint;
13use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileFingerprint};
14use type_bridge_contract::migration::{
15    MigrationId, MigrationManifestDigest, MigrationPlanFingerprint,
16};
17use type_bridge_contract::schema_delta::ManagedSchemaState;
18use type_bridge_contract::schema_fingerprint::{
19    ManagedDeclaredIdentityFingerprint, ManagedSemanticSchemaFingerprint,
20};
21use type_bridge_contract::schema_lowering::SchemaLoweringProfileFingerprint;
22
23use crate::{
24    VerifiedMigrationApplyManifest, VerifiedMigrationApplyPlan, VerifiedMigrationRollbackManifest,
25    VerifiedMigrationRollbackPlan, VerifiedMigrationTransactionGroup,
26};
27
28const MAX_LEASE_HOLDER_BYTES: usize = 128;
29
30/// Shared maximum number of migration transaction groups in one invocation.
31pub const MAX_MIGRATION_EXECUTION_GROUPS: usize = 65_536;
32/// Shared maximum number of retained terminal backfill observations.
33pub const MAX_MIGRATION_BACKFILL_OBSERVATIONS: usize = 65_536;
34
35/// Cloneable, wakeable cancellation authority for migration operations.
36#[derive(Clone, Debug, Default)]
37pub struct MigrationCancellation {
38    inner: Arc<MigrationCancellationInner>,
39}
40
41#[derive(Debug, Default)]
42struct MigrationCancellationInner {
43    cancelled: AtomicBool,
44    waiters: std::sync::Mutex<Vec<Waker>>,
45}
46
47impl MigrationCancellation {
48    /// Request cancellation and wake every currently registered provider wait.
49    pub fn cancel(&self) {
50        if !self.inner.cancelled.swap(true, Ordering::AcqRel) {
51            let waiters = {
52                let mut waiters = self.inner.waiters.lock().expect("cancellation waiters");
53                std::mem::take(&mut *waiters)
54            };
55            for waiter in waiters {
56                waiter.wake();
57            }
58        }
59    }
60
61    /// Return whether cancellation has been requested.
62    #[must_use]
63    pub fn is_cancelled(&self) -> bool {
64        self.inner.cancelled.load(Ordering::Acquire)
65    }
66
67    /// Await cancellation without polling or spawning a worker thread.
68    pub async fn cancelled(&self) {
69        std::future::poll_fn(|context| {
70            if self.is_cancelled() {
71                return Poll::Ready(());
72            }
73            let mut waiters = self.inner.waiters.lock().expect("cancellation waiters");
74            if self.is_cancelled() {
75                return Poll::Ready(());
76            }
77            if !waiters
78                .iter()
79                .any(|waiter| waiter.will_wake(context.waker()))
80            {
81                waiters.push(context.waker().clone());
82            }
83            Poll::Pending
84        })
85        .await
86    }
87}
88
89/// Tighten-only bounded resources for one migration execution invocation.
90#[derive(Clone, Copy, Debug, Eq, PartialEq)]
91pub struct MigrationExecutionResourceLimits {
92    transaction_groups: usize,
93    backfill_observations: usize,
94}
95
96impl MigrationExecutionResourceLimits {
97    /// Construct limits clamped to the shared ceilings.
98    #[must_use]
99    pub const fn tightened(transaction_groups: usize, backfill_observations: usize) -> Self {
100        Self {
101            transaction_groups: if transaction_groups < MAX_MIGRATION_EXECUTION_GROUPS {
102                transaction_groups
103            } else {
104                MAX_MIGRATION_EXECUTION_GROUPS
105            },
106            backfill_observations: if backfill_observations < MAX_MIGRATION_BACKFILL_OBSERVATIONS {
107                backfill_observations
108            } else {
109                MAX_MIGRATION_BACKFILL_OBSERVATIONS
110            },
111        }
112    }
113
114    /// Return the maximum transaction groups admitted by this invocation.
115    pub const fn transaction_groups(self) -> usize {
116        self.transaction_groups
117    }
118
119    /// Return the maximum retained terminal backfill observations.
120    pub const fn backfill_observations(self) -> usize {
121        self.backfill_observations
122    }
123}
124
125impl Default for MigrationExecutionResourceLimits {
126    fn default() -> Self {
127        Self::tightened(
128            MAX_MIGRATION_EXECUTION_GROUPS,
129            MAX_MIGRATION_BACKFILL_OBSERVATIONS,
130        )
131    }
132}
133
134/// Immutable controls shared by apply and rollback execution.
135#[derive(Clone, Debug, Default)]
136pub struct MigrationExecutionControl {
137    cancellation: MigrationCancellation,
138    deadline: Option<Instant>,
139    resources: MigrationExecutionResourceLimits,
140}
141
142impl MigrationExecutionControl {
143    /// Bind cancellation, one absolute monotonic deadline, and tightened limits.
144    #[must_use]
145    pub const fn new(
146        cancellation: MigrationCancellation,
147        deadline: Option<Instant>,
148        resources: MigrationExecutionResourceLimits,
149    ) -> Self {
150        Self {
151            cancellation,
152            deadline,
153            resources,
154        }
155    }
156
157    /// Return the cancellation authority.
158    pub const fn cancellation(&self) -> &MigrationCancellation {
159        &self.cancellation
160    }
161
162    /// Return the absolute monotonic deadline, when bounded.
163    pub const fn deadline(&self) -> Option<Instant> {
164        self.deadline
165    }
166
167    /// Return caller-tightened resource limits.
168    pub const fn resources(&self) -> MigrationExecutionResourceLimits {
169        self.resources
170    }
171
172    /// Reject work at a safe coordinator boundary before another effect begins.
173    pub fn check(&self) -> Result<(), Diagnostic> {
174        if self.cancellation.is_cancelled() {
175            return Err(failure(
176                DiagnosticCategory::Cancelled,
177                "migration_execution_cancelled",
178                "migration execution was cancelled before the next effect",
179            ));
180        }
181        if self
182            .deadline
183            .is_some_and(|deadline| Instant::now() >= deadline)
184        {
185            return Err(failure(
186                DiagnosticCategory::ResourceLimit,
187                "migration_execution_deadline_exceeded",
188                "migration execution reached its absolute deadline before the next effect",
189            ));
190        }
191        Ok(())
192    }
193}
194
195/// Boxed future returned by provider-neutral execution stores.
196pub type ExecutionFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, Diagnostic>> + Send + 'a>>;
197
198/// A monotonically increasing store-issued migration fencing token.
199#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
200pub struct ExecutionFence(u64);
201
202impl ExecutionFence {
203    /// Construct a non-zero fence.
204    pub fn new(value: u64) -> Result<Self, Diagnostic> {
205        if value == 0 {
206            return Err(failure(
207                DiagnosticCategory::InvalidContract,
208                "migration_execution_zero_fence",
209                "migration execution fences must be non-zero",
210            ));
211        }
212        Ok(Self(value))
213    }
214
215    /// Return the numeric fence value.
216    pub const fn get(self) -> u64 {
217        self.0
218    }
219
220    /// Derive the next strictly greater fence without wrapping.
221    pub fn checked_successor(self) -> Result<Self, Diagnostic> {
222        let value = self.0.checked_add(1).ok_or_else(|| {
223            failure(
224                DiagnosticCategory::ResourceLimit,
225                "migration_execution_fence_exhausted",
226                "migration execution fence range is exhausted",
227            )
228        })?;
229        Self::new(value)
230    }
231}
232
233/// One durable managed migration scope.
234#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
235pub struct ExecutionScope(ManagedScopeId);
236
237impl ExecutionScope {
238    /// Bind execution to an existing managed-scope identity.
239    pub const fn new(scope: ManagedScopeId) -> Self {
240        Self(scope)
241    }
242
243    /// Return the managed-scope identity.
244    pub const fn managed_scope_id(&self) -> &ManagedScopeId {
245        &self.0
246    }
247}
248
249/// Caller-supplied lease holder identity with no machine-derived garnish.
250#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
251pub struct LeaseHolderId(String);
252
253impl LeaseHolderId {
254    /// Validate a bounded canonical holder label.
255    pub fn new(value: impl Into<String>) -> Result<Self, Diagnostic> {
256        let value = value.into();
257        if value.is_empty()
258            || value.len() > MAX_LEASE_HOLDER_BYTES
259            || !value
260                .bytes()
261                .all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
262        {
263            return Err(failure(
264                DiagnosticCategory::InvalidContract,
265                "migration_execution_invalid_holder",
266                "lease holder must be bounded non-empty ASCII [A-Za-z0-9._-]",
267            ));
268        }
269        Ok(Self(value))
270    }
271
272    /// Return the canonical holder label.
273    pub fn as_str(&self) -> &str {
274        &self.0
275    }
276}
277
278/// Store-issued authority to mutate one migration scope.
279#[derive(Clone, Debug)]
280pub struct MigrationLease {
281    scope: ExecutionScope,
282    holder: LeaseHolderId,
283    fence: ExecutionFence,
284    local_binding: Option<ExecutionBindingToken>,
285}
286
287/// Opaque, process-local identity shared by one provider/store execution pair.
288///
289/// The token is deliberately absent from canonical and persisted migration
290/// records. Provider-neutral stores continue to issue unbound leases through
291/// [`MigrationLease::new`]; provider adapters that require local composition
292/// integrity can mint a token and use [`MigrationLease::new_bound`].
293#[doc(hidden)]
294#[derive(Clone)]
295pub struct ExecutionBindingToken(Arc<()>);
296
297impl ExecutionBindingToken {
298    /// Mint one fresh, unforgeable process-local identity.
299    #[doc(hidden)]
300    #[must_use]
301    pub fn fresh() -> Self {
302        Self(Arc::new(()))
303    }
304
305    fn matches(&self, other: &Self) -> bool {
306        Arc::ptr_eq(&self.0, &other.0)
307    }
308}
309
310impl fmt::Debug for ExecutionBindingToken {
311    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
312        formatter.write_str("ExecutionBindingToken([OPAQUE])")
313    }
314}
315
316impl MigrationLease {
317    /// Construct an unbound lease returned by a provider-neutral store.
318    pub const fn new(scope: ExecutionScope, holder: LeaseHolderId, fence: ExecutionFence) -> Self {
319        Self {
320            scope,
321            holder,
322            fence,
323            local_binding: None,
324        }
325    }
326
327    /// Construct a lease carrying one process-local provider/store binding.
328    #[doc(hidden)]
329    #[must_use]
330    pub fn new_bound(
331        scope: ExecutionScope,
332        holder: LeaseHolderId,
333        fence: ExecutionFence,
334        local_binding: ExecutionBindingToken,
335    ) -> Self {
336        Self {
337            scope,
338            holder,
339            fence,
340            local_binding: Some(local_binding),
341        }
342    }
343
344    /// Return whether this lease carries the exact supplied local binding.
345    #[doc(hidden)]
346    #[must_use]
347    pub fn is_bound_to(&self, expected: &ExecutionBindingToken) -> bool {
348        self.local_binding
349            .as_ref()
350            .is_some_and(|actual| actual.matches(expected))
351    }
352
353    /// Return the leased scope.
354    pub const fn scope(&self) -> &ExecutionScope {
355        &self.scope
356    }
357
358    /// Return the holder identity.
359    pub const fn holder(&self) -> &LeaseHolderId {
360        &self.holder
361    }
362
363    /// Return the store-issued fence.
364    pub const fn fence(&self) -> ExecutionFence {
365        self.fence
366    }
367}
368
369impl PartialEq for MigrationLease {
370    fn eq(&self, other: &Self) -> bool {
371        self.scope == other.scope
372            && self.holder == other.holder
373            && self.fence == other.fence
374            && match (&self.local_binding, &other.local_binding) {
375                (None, None) => true,
376                (Some(left), Some(right)) => left.matches(right),
377                (None, Some(_)) | (Some(_), None) => false,
378            }
379    }
380}
381
382impl Eq for MigrationLease {}
383
384/// Commit certainty owned by the migration journal layer.
385///
386/// Provider adapters must map absent certainty information to [`Self::Unknown`].
387/// Only an explicit provider proof that commit could not have occurred may map
388/// to [`Self::DefinitelyAborted`].
389#[derive(Clone, Copy, Debug, Eq, PartialEq)]
390pub enum GroupCommitCertainty {
391    /// The provider proves that the group transaction did not commit.
392    DefinitelyAborted,
393    /// The provider cannot determine whether the group transaction committed.
394    Unknown,
395}
396
397impl GroupCommitCertainty {
398    /// Convert certainty into its durable journal event.
399    pub const fn journal_event(self) -> GroupJournalEventKind {
400        match self {
401            Self::DefinitelyAborted => GroupJournalEventKind::DefinitelyAborted,
402            Self::Unknown => GroupJournalEventKind::CommitOutcomeUnknown,
403        }
404    }
405}
406
407/// Durable commit-boundary event vocabulary for one transaction group.
408#[derive(Clone, Copy, Debug, Eq, PartialEq)]
409pub enum GroupJournalEventKind {
410    /// The group statements completed and the commit call is about to begin.
411    BeforeCommit,
412    /// The provider commit succeeded and exact live target semantics were observed.
413    Committed,
414    /// The commit response cannot prove whether durability occurred.
415    CommitOutcomeUnknown,
416    /// The provider proves that the transaction did not commit.
417    DefinitelyAborted,
418    /// An empty formal-only group advanced without a provider transaction.
419    FormalOnlyAdvanced,
420}
421
422/// Optional live managed-semantic evidence used during recovery.
423#[derive(Clone, Debug, Eq, PartialEq)]
424pub enum GroupRecoveryObservation {
425    /// The provider cannot supply a trustworthy managed-semantic observation.
426    Unavailable,
427    /// Exact live managed semantics observed under the current fence.
428    ManagedSemantics(ManagedSemanticSchemaFingerprint),
429}
430
431/// Fail-closed recovery decision for one positional transaction group.
432#[derive(Clone, Copy, Debug, Eq, PartialEq)]
433pub enum GroupRecoveryDecision {
434    /// The group is proven absent and may execute under the current fence.
435    ExecuteNormally,
436    /// The target is proven reached and only the journal checkpoint needs repair.
437    RepairCheckpoint,
438    /// Evidence is ambiguous or contradictory and requires verified operator action.
439    RequiresExplicitRecovery,
440}
441
442/// Direction of one closed backfill program at execution time.
443#[derive(Clone, Copy, Debug, Eq, PartialEq)]
444pub enum BackfillExecutionDirection {
445    /// Execute the canonical forward data program.
446    Forward,
447    /// Execute its independently verified reverse program.
448    Reverse,
449}
450
451/// Bounded aggregate counts from a completely verified backfill execution.
452#[derive(Clone, Copy, Debug, Eq, PartialEq)]
453pub struct BackfillExecutionCounts {
454    matched: u64,
455    changed: u64,
456    skipped: u64,
457    transaction_groups: u32,
458}
459
460impl BackfillExecutionCounts {
461    /// Construct internally consistent terminal counts.
462    pub fn new(
463        matched: u64,
464        changed: u64,
465        skipped: u64,
466        transaction_groups: u32,
467    ) -> Result<Self, Diagnostic> {
468        if transaction_groups == 0 {
469            return Err(failure(
470                DiagnosticCategory::InvalidContract,
471                "migration_backfill_zero_transaction_groups",
472                "terminal backfill evidence must include at least one transaction group",
473            ));
474        }
475        if changed.checked_add(skipped) != Some(matched) {
476            return Err(failure(
477                DiagnosticCategory::InvalidContract,
478                "migration_backfill_count_mismatch",
479                "backfill matched count must equal changed plus skipped counts",
480            ));
481        }
482        Ok(Self {
483            matched,
484            changed,
485            skipped,
486            transaction_groups,
487        })
488    }
489
490    /// Return the number of source rows selected by the closed plan.
491    pub const fn matched(self) -> u64 {
492        self.matched
493    }
494
495    /// Return the number of destination rows changed.
496    pub const fn changed(self) -> u64 {
497        self.changed
498    }
499
500    /// Return the number of already-equal rows skipped idempotently.
501    pub const fn skipped(self) -> u64 {
502        self.skipped
503    }
504
505    /// Return the number of committed deterministic transaction groups.
506    pub const fn transaction_groups(self) -> u32 {
507        self.transaction_groups
508    }
509}
510
511/// Exact terminal proof that one closed backfill program satisfies its postcondition.
512#[derive(Clone, Debug, Eq, PartialEq)]
513pub struct BackfillCompletionEvidence {
514    plan_fingerprint: Fingerprint,
515    direction: BackfillExecutionDirection,
516    counts: BackfillExecutionCounts,
517}
518
519impl BackfillCompletionEvidence {
520    /// Bind terminal counts to one exact canonical plan and execution direction.
521    #[must_use]
522    pub const fn new(
523        plan_fingerprint: Fingerprint,
524        direction: BackfillExecutionDirection,
525        counts: BackfillExecutionCounts,
526    ) -> Self {
527        Self {
528            plan_fingerprint,
529            direction,
530            counts,
531        }
532    }
533
534    /// Return the exact canonical backfill-plan identity.
535    pub const fn plan_fingerprint(&self) -> &Fingerprint {
536        &self.plan_fingerprint
537    }
538
539    /// Return whether the forward or checked reverse program completed.
540    pub const fn direction(&self) -> BackfillExecutionDirection {
541        self.direction
542    }
543
544    /// Return internally consistent aggregate counts.
545    pub const fn counts(&self) -> BackfillExecutionCounts {
546        self.counts
547    }
548}
549
550/// Fresh provider observation of one exact backfill postcondition.
551#[derive(Clone, Debug, Eq, PartialEq)]
552pub enum BackfillRecoveryObservation {
553    /// No trustworthy data observation is available.
554    Unavailable,
555    /// The exact postcondition is not currently satisfied.
556    Incomplete,
557    /// The exact plan and direction have a verified terminal postcondition.
558    Complete(BackfillCompletionEvidence),
559}
560
561/// Decide whether a complete backfill may execute or only repair its checkpoint.
562///
563/// Unlike schema groups, an incomplete data observation after a before-commit
564/// or unknown-commit event cannot prove that no partition committed. Automatic
565/// replay therefore remains forbidden until partition checkpoints are present.
566/// With no prior event, both incomplete and already-complete postconditions are
567/// safe to execute: the journal proves that execution has never begun, and an
568/// already-complete closed backfill is a deterministic no-op.
569pub fn decide_backfill_recovery(
570    last_event: Option<GroupJournalEventKind>,
571    observation: &BackfillRecoveryObservation,
572    expected_plan: &Fingerprint,
573    expected_direction: BackfillExecutionDirection,
574) -> GroupRecoveryDecision {
575    let complete = matches!(
576        observation,
577        BackfillRecoveryObservation::Complete(evidence)
578            if evidence.plan_fingerprint() == expected_plan
579                && evidence.direction() == expected_direction
580    );
581    match last_event {
582        None if matches!(
583            observation,
584            BackfillRecoveryObservation::Incomplete | BackfillRecoveryObservation::Complete(_)
585        ) =>
586        {
587            GroupRecoveryDecision::ExecuteNormally
588        }
589        Some(GroupJournalEventKind::DefinitelyAborted)
590            if matches!(observation, BackfillRecoveryObservation::Incomplete) =>
591        {
592            GroupRecoveryDecision::ExecuteNormally
593        }
594        Some(
595            GroupJournalEventKind::BeforeCommit
596            | GroupJournalEventKind::CommitOutcomeUnknown
597            | GroupJournalEventKind::Committed,
598        ) if complete => GroupRecoveryDecision::RepairCheckpoint,
599        _ => GroupRecoveryDecision::RequiresExplicitRecovery,
600    }
601}
602
603/// Decide recovery from durable event and freshly observed managed semantics.
604///
605/// Equal source and target fingerprints are intentionally uninformative for
606/// unknown commit outcomes. They never authorize replay or checkpoint repair.
607pub fn decide_group_recovery(
608    last_event: Option<GroupJournalEventKind>,
609    observation: &GroupRecoveryObservation,
610    source: &ManagedSemanticSchemaFingerprint,
611    target: &ManagedSemanticSchemaFingerprint,
612) -> GroupRecoveryDecision {
613    let distinct = source != target;
614    let observed = match observation {
615        GroupRecoveryObservation::Unavailable => ObservedRelation::Unavailable,
616        GroupRecoveryObservation::ManagedSemantics(value) if value == source && value == target => {
617            ObservedRelation::Both
618        }
619        GroupRecoveryObservation::ManagedSemantics(value) if value == source => {
620            ObservedRelation::Source
621        }
622        GroupRecoveryObservation::ManagedSemantics(value) if value == target => {
623            ObservedRelation::Target
624        }
625        GroupRecoveryObservation::ManagedSemantics(_) => ObservedRelation::Neither,
626    };
627    match (last_event, distinct, observed) {
628        (None, true, ObservedRelation::Source)
629        | (Some(GroupJournalEventKind::DefinitelyAborted), true, ObservedRelation::Source)
630        | (None, false, ObservedRelation::Both)
631        | (Some(GroupJournalEventKind::DefinitelyAborted), false, ObservedRelation::Both) => {
632            GroupRecoveryDecision::ExecuteNormally
633        }
634        (
635            Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
636            true,
637            ObservedRelation::Source,
638        ) => GroupRecoveryDecision::ExecuteNormally,
639        (
640            Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
641            true,
642            ObservedRelation::Target,
643        )
644        | (Some(GroupJournalEventKind::Committed), true, ObservedRelation::Target)
645        | (
646            Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced),
647            false,
648            ObservedRelation::Both,
649        ) => GroupRecoveryDecision::RepairCheckpoint,
650        _ => GroupRecoveryDecision::RequiresExplicitRecovery,
651    }
652}
653
654#[derive(Clone, Copy, Debug, Eq, PartialEq)]
655enum ObservedRelation {
656    Source,
657    Target,
658    Both,
659    Neither,
660    Unavailable,
661}
662
663/// Store-assigned monotonic ordering identity for journal entries.
664#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
665pub struct JournalSequence(u64);
666
667impl JournalSequence {
668    /// Construct a non-zero journal sequence.
669    pub fn new(value: u64) -> Result<Self, Diagnostic> {
670        if value == 0 {
671            return Err(failure(
672                DiagnosticCategory::InvalidContract,
673                "migration_execution_zero_sequence",
674                "journal sequence numbers must be non-zero",
675            ));
676        }
677        Ok(Self(value))
678    }
679
680    /// Return the numeric sequence value.
681    pub const fn get(self) -> u64 {
682        self.0
683    }
684}
685
686/// One record after the store atomically assigns its ordering sequence.
687#[derive(Clone, Debug, Eq, PartialEq)]
688pub struct JournalEntry<T> {
689    sequence: JournalSequence,
690    record: T,
691}
692
693impl<T> JournalEntry<T> {
694    /// Attach a sequence allocated by the authoritative store.
695    pub const fn from_store(sequence: JournalSequence, record: T) -> Self {
696        Self { sequence, record }
697    }
698
699    /// Return the store ordering sequence.
700    pub const fn sequence(&self) -> JournalSequence {
701        self.sequence
702    }
703
704    /// Return the trusted record.
705    pub const fn record(&self) -> &T {
706        &self.record
707    }
708
709    /// Consume the envelope and return its record.
710    pub fn into_record(self) -> T {
711        self.record
712    }
713}
714
715/// One trusted, still-open plan and its ordered group events.
716#[derive(Clone, Debug, Eq, PartialEq)]
717pub struct OpenPlanRecord {
718    plan: JournalEntry<PlanRecord>,
719    events: Vec<JournalEntry<GroupEventRecord>>,
720    backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
721}
722
723impl OpenPlanRecord {
724    /// Rebuild store output while checking sequence, scope, fence, and manifest binding.
725    ///
726    /// The plan retains its original fence while recovery events may be written
727    /// under later fences. Event fences therefore must be monotonic and no older
728    /// than the plan fence; requiring equality would hide durable recovery work
729    /// after the lease rolls forward.
730    pub fn from_store(
731        plan: JournalEntry<PlanRecord>,
732        events: Vec<JournalEntry<GroupEventRecord>>,
733    ) -> Result<Self, Diagnostic> {
734        Self::from_store_with_backfills(plan, events, Vec::new())
735    }
736
737    /// Rebuild store output including independently journaled data-step events.
738    pub fn from_store_with_backfills(
739        plan: JournalEntry<PlanRecord>,
740        events: Vec<JournalEntry<GroupEventRecord>>,
741        backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
742    ) -> Result<Self, Diagnostic> {
743        let mut previous = plan.sequence();
744        let mut previous_fence = plan.record().fence();
745        let mut ordered: Vec<(
746            JournalSequence,
747            ExecutionFence,
748            &MigrationId,
749            MigrationManifestDigest,
750        )> = events
751            .iter()
752            .map(|event| {
753                (
754                    event.sequence(),
755                    event.record().fence(),
756                    event.record().migration_id(),
757                    event.record().manifest_digest(),
758                )
759            })
760            .chain(backfill_events.iter().map(|event| {
761                (
762                    event.sequence(),
763                    event.record().fence(),
764                    event.record().migration_id(),
765                    event.record().manifest_digest(),
766                )
767            }))
768            .collect();
769        ordered.sort_by_key(|event| event.0);
770        for (sequence, fence, migration_id, manifest_digest) in ordered {
771            let manifest_index = plan
772                .record()
773                .manifest_digests()
774                .iter()
775                .position(|digest| digest == &manifest_digest);
776            if sequence <= previous
777                || fence < previous_fence
778                || manifest_index.is_none_or(|index| {
779                    plan.record().migration_ids().get(index) != Some(migration_id)
780                })
781            {
782                return Err(failure(
783                    DiagnosticCategory::Integrity,
784                    "migration_execution_invalid_open_plan",
785                    "loaded open-plan events are not ordered and bound to the plan",
786                ));
787            }
788            previous = sequence;
789            previous_fence = fence;
790        }
791        if events
792            .iter()
793            .any(|event| event.record().scope() != plan.record().scope())
794            || backfill_events
795                .iter()
796                .any(|event| event.record().scope() != plan.record().scope())
797        {
798            return Err(failure(
799                DiagnosticCategory::Integrity,
800                "migration_execution_invalid_open_plan",
801                "loaded open-plan events are not bound to the plan scope",
802            ));
803        }
804        Ok(Self {
805            plan,
806            events,
807            backfill_events,
808        })
809    }
810
811    /// Return the sequenced plan record.
812    pub const fn plan(&self) -> &JournalEntry<PlanRecord> {
813        &self.plan
814    }
815
816    /// Return ordered sequenced group events.
817    pub fn events(&self) -> &[JournalEntry<GroupEventRecord>] {
818        &self.events
819    }
820
821    /// Return ordered sequenced backfill events.
822    pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
823        &self.backfill_events
824    }
825}
826
827/// Identity-only journal record for one complete verified apply plan.
828#[derive(Clone, Debug, Eq, PartialEq)]
829pub struct PlanRecord {
830    scope: ExecutionScope,
831    fence: ExecutionFence,
832    source_applied: Vec<MigrationId>,
833    source_frontier: Vec<MigrationId>,
834    target_frontier: Vec<MigrationId>,
835    migration_ids: Vec<MigrationId>,
836    manifest_digests: Vec<MigrationManifestDigest>,
837    manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
838    source_declared: ManagedDeclaredIdentityFingerprint,
839    target_declared: ManagedDeclaredIdentityFingerprint,
840    source_semantics: ManagedSemanticSchemaFingerprint,
841    target_semantics: ManagedSemanticSchemaFingerprint,
842    semantic_profile: SemanticProfileFingerprint,
843    lowering_profile: SchemaLoweringProfileFingerprint,
844    observed_live_source: ManagedSemanticSchemaFingerprint,
845}
846
847impl PlanRecord {
848    /// Bind a fresh-lease ledger and live-state precondition to verified plan identities.
849    pub fn from_verified_plan(
850        lease: &MigrationLease,
851        plan: &VerifiedMigrationApplyPlan,
852        observed_applied_migrations: &[MigrationId],
853        observed_live_source: &ManagedSchemaState,
854    ) -> Result<Self, Diagnostic> {
855        let source = plan.source_state().ok_or_else(|| {
856            failure(
857                DiagnosticCategory::InvalidContract,
858                "migration_execution_empty_plan",
859                "an executable migration plan requires a source state",
860            )
861        })?;
862        let target = plan.target_state().ok_or_else(|| {
863            failure(
864                DiagnosticCategory::InvalidContract,
865                "migration_execution_empty_plan",
866                "an executable migration plan requires a target state",
867            )
868        })?;
869        let first = plan.migrations().first().ok_or_else(|| {
870            failure(
871                DiagnosticCategory::InvalidContract,
872                "migration_execution_empty_plan",
873                "an executable migration plan requires at least one manifest",
874            )
875        })?;
876        if observed_applied_migrations != plan.applied_migrations() {
877            return Err(failure(
878                DiagnosticCategory::Integrity,
879                "migration_execution_stale_applied_set",
880                "applied ledger changed after migration planning; rebuild the plan",
881            ));
882        }
883        if observed_live_source != source {
884            return Err(failure(
885                DiagnosticCategory::Integrity,
886                "migration_execution_stale_source_state",
887                "live managed state differs from the planned source; rebuild the plan",
888            ));
889        }
890        let scope = ExecutionScope::new(source.scope().id().clone());
891        if lease.scope() != &scope || target.scope() != source.scope() {
892            return Err(failure(
893                DiagnosticCategory::Integrity,
894                "migration_execution_scope_mismatch",
895                "lease, source, and target must bind the same managed scope",
896            ));
897        }
898        let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
899        let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
900        for migration in plan.migrations() {
901            if migration.manifest().managed_scope().id() != scope.managed_scope_id()
902                || migration.manifest().semantic_profile().fingerprint() != &semantic_profile
903                || migration.manifest().lowering_profile().fingerprint() != &lowering_profile
904            {
905                return Err(failure(
906                    DiagnosticCategory::Integrity,
907                    "migration_execution_plan_binding_mismatch",
908                    "planned manifests do not share exact scope and profile bindings",
909                ));
910            }
911        }
912        Ok(Self {
913            scope,
914            fence: lease.fence(),
915            source_applied: plan.applied_migrations().to_vec(),
916            source_frontier: plan.applied_frontier().to_vec(),
917            target_frontier: plan.target_frontier().to_vec(),
918            migration_ids: plan
919                .migrations()
920                .iter()
921                .map(|migration| migration.manifest().id().clone())
922                .collect(),
923            manifest_digests: plan
924                .migrations()
925                .iter()
926                .map(VerifiedMigrationApplyManifest::digest)
927                .collect(),
928            manifest_plan_fingerprints: plan
929                .migrations()
930                .iter()
931                .map(|migration| migration.manifest().plan_fingerprint().clone())
932                .collect(),
933            source_declared: source.managed_declared_identity().clone(),
934            target_declared: target.managed_declared_identity().clone(),
935            source_semantics: source.managed_semantic_schema().clone(),
936            target_semantics: target.managed_semantic_schema().clone(),
937            semantic_profile,
938            lowering_profile,
939            observed_live_source: observed_live_source.managed_semantic_schema().clone(),
940        })
941    }
942
943    /// Return the execution scope.
944    pub const fn scope(&self) -> &ExecutionScope {
945        &self.scope
946    }
947    /// Return the fence bound into this record.
948    pub const fn fence(&self) -> ExecutionFence {
949        self.fence
950    }
951    /// Return the planned source frontier.
952    pub fn source_frontier(&self) -> &[MigrationId] {
953        &self.source_frontier
954    }
955    /// Return the complete canonically ordered source applied set.
956    pub fn source_applied(&self) -> &[MigrationId] {
957        &self.source_applied
958    }
959    /// Return the planned target frontier.
960    pub fn target_frontier(&self) -> &[MigrationId] {
961        &self.target_frontier
962    }
963    /// Return ordered migration identities.
964    pub fn migration_ids(&self) -> &[MigrationId] {
965        &self.migration_ids
966    }
967    /// Return ordered canonical manifest digests.
968    pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
969        &self.manifest_digests
970    }
971    /// Return ordered manifest plan fingerprints.
972    pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
973        &self.manifest_plan_fingerprints
974    }
975    /// Return planned source managed-declared identity.
976    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
977        &self.source_declared
978    }
979    /// Return planned target managed-declared identity.
980    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
981        &self.target_declared
982    }
983    /// Return planned source managed semantics.
984    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
985        &self.source_semantics
986    }
987    /// Return planned target managed semantics.
988    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
989        &self.target_semantics
990    }
991    /// Return semantic-profile content identity.
992    pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
993        &self.semantic_profile
994    }
995    /// Return lowering-registry content identity.
996    pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
997        &self.lowering_profile
998    }
999    /// Return the pre-mutation observed source semantics.
1000    pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
1001        &self.observed_live_source
1002    }
1003}
1004
1005/// Identity-only journal event for one verified transaction group.
1006#[derive(Clone, Debug, Eq, PartialEq)]
1007pub struct GroupEventRecord {
1008    scope: ExecutionScope,
1009    fence: ExecutionFence,
1010    manifest_digest: MigrationManifestDigest,
1011    migration_id: MigrationId,
1012    group_ordinal: u32,
1013    first_step_index: u32,
1014    schema_delta_step_index: u32,
1015    end_step_index: u32,
1016    kind: GroupJournalEventKind,
1017    observed_target: Option<ManagedSemanticSchemaFingerprint>,
1018}
1019
1020impl GroupEventRecord {
1021    /// Derive a positional event from exact verified apply evidence.
1022    pub fn new(
1023        lease: &MigrationLease,
1024        migration: &VerifiedMigrationApplyManifest,
1025        group: &VerifiedMigrationTransactionGroup,
1026        kind: GroupJournalEventKind,
1027        observed_target: Option<ManagedSemanticSchemaFingerprint>,
1028    ) -> Result<Self, Diagnostic> {
1029        if migration.transaction_groups().get(group.ordinal()) != Some(group) {
1030            return Err(failure(
1031                DiagnosticCategory::Integrity,
1032                "migration_execution_foreign_group",
1033                "transaction group does not belong to the supplied verified manifest",
1034            ));
1035        }
1036        let step = migration
1037            .steps()
1038            .get(group.schema_delta_step_index())
1039            .ok_or_else(|| {
1040                failure(
1041                    DiagnosticCategory::Integrity,
1042                    "migration_execution_group_position_mismatch",
1043                    "transaction group delta position is outside the verified manifest",
1044                )
1045            })?;
1046        let delta = step
1047            .step()
1048            .as_schema_delta()
1049            .ok_or_else(|| {
1050                failure(
1051                    DiagnosticCategory::Integrity,
1052                    "migration_execution_group_position_mismatch",
1053                    "transaction group does not terminate in a schema delta",
1054                )
1055            })?
1056            .delta();
1057        let lowering = step.lowering().ok_or_else(|| {
1058            failure(
1059                DiagnosticCategory::Integrity,
1060                "migration_execution_group_lowering_missing",
1061                "transaction group delta has no verified lowering",
1062            )
1063        })?;
1064        let scope = ExecutionScope::new(migration.manifest().managed_scope().id().clone());
1065        if lease.scope() != &scope {
1066            return Err(failure(
1067                DiagnosticCategory::Integrity,
1068                "migration_execution_scope_mismatch",
1069                "lease scope differs from the verified migration scope",
1070            ));
1071        }
1072        match kind {
1073            GroupJournalEventKind::Committed
1074                if observed_target.as_ref() == Some(delta.target().managed_semantic_schema()) => {}
1075            GroupJournalEventKind::Committed => {
1076                return Err(failure(
1077                    DiagnosticCategory::Integrity,
1078                    "migration_execution_commit_evidence_mismatch",
1079                    "committed event requires the exact observed target semantics",
1080                ));
1081            }
1082            GroupJournalEventKind::FormalOnlyAdvanced
1083                if observed_target.is_none()
1084                    && group.assertion_count() == 0
1085                    && lowering.units().is_empty()
1086                    && delta.source().managed_semantic_schema()
1087                        == delta.target().managed_semantic_schema() => {}
1088            GroupJournalEventKind::FormalOnlyAdvanced => {
1089                return Err(failure(
1090                    DiagnosticCategory::InvalidContract,
1091                    "migration_execution_invalid_formal_advance",
1092                    "formal-only advancement requires an assertion-free empty equal-semantic group",
1093                ));
1094            }
1095            _ if observed_target.is_none() => {}
1096            _ => {
1097                return Err(failure(
1098                    DiagnosticCategory::InvalidContract,
1099                    "migration_execution_unexpected_observation",
1100                    "only committed events may carry observed target semantics",
1101                ));
1102            }
1103        }
1104        Ok(Self {
1105            scope,
1106            fence: lease.fence(),
1107            manifest_digest: migration.digest(),
1108            migration_id: migration.manifest().id().clone(),
1109            group_ordinal: position(group.ordinal())?,
1110            first_step_index: position(group.first_step_index())?,
1111            schema_delta_step_index: position(group.schema_delta_step_index())?,
1112            end_step_index: position(group.end_step_index())?,
1113            kind,
1114            observed_target,
1115        })
1116    }
1117
1118    /// Return the execution scope.
1119    pub const fn scope(&self) -> &ExecutionScope {
1120        &self.scope
1121    }
1122    /// Return the event fence.
1123    pub const fn fence(&self) -> ExecutionFence {
1124        self.fence
1125    }
1126    /// Return the canonical manifest digest.
1127    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1128        self.manifest_digest
1129    }
1130    /// Return the migration identity.
1131    pub const fn migration_id(&self) -> &MigrationId {
1132        &self.migration_id
1133    }
1134    /// Return the group ordinal.
1135    pub const fn group_ordinal(&self) -> u32 {
1136        self.group_ordinal
1137    }
1138    /// Return the first group step index.
1139    pub const fn first_step_index(&self) -> u32 {
1140        self.first_step_index
1141    }
1142    /// Return the terminal delta step index.
1143    pub const fn schema_delta_step_index(&self) -> u32 {
1144        self.schema_delta_step_index
1145    }
1146    /// Return the exclusive group step end.
1147    pub const fn end_step_index(&self) -> u32 {
1148        self.end_step_index
1149    }
1150    /// Return the event kind.
1151    pub const fn kind(&self) -> GroupJournalEventKind {
1152        self.kind
1153    }
1154    /// Return exact target evidence carried only by committed events.
1155    pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
1156        self.observed_target.as_ref()
1157    }
1158}
1159
1160/// Durable positional event for one exact forward or reverse backfill step.
1161///
1162/// Data effects cannot be inferred from unchanged schema semantics. This
1163/// record therefore carries the canonical plan identity on every boundary and
1164/// admits terminal counts only when the provider has verified the complete
1165/// postcondition.
1166#[derive(Clone, Debug, Eq, PartialEq)]
1167pub struct BackfillEventRecord {
1168    scope: ExecutionScope,
1169    fence: ExecutionFence,
1170    manifest_digest: MigrationManifestDigest,
1171    migration_id: MigrationId,
1172    operation_ordinal: u32,
1173    manifest_step_index: u32,
1174    plan_fingerprint: Fingerprint,
1175    direction: BackfillExecutionDirection,
1176    kind: GroupJournalEventKind,
1177    completion: Option<BackfillCompletionEvidence>,
1178}
1179
1180impl BackfillEventRecord {
1181    /// Derive a forward event from one exact verified apply step position.
1182    pub fn new_apply(
1183        lease: &MigrationLease,
1184        migration: &VerifiedMigrationApplyManifest,
1185        step_index: usize,
1186        kind: GroupJournalEventKind,
1187        completion: Option<BackfillCompletionEvidence>,
1188    ) -> Result<Self, Diagnostic> {
1189        if lease.scope().managed_scope_id() != migration.manifest().managed_scope().id() {
1190            return Err(failure(
1191                DiagnosticCategory::Integrity,
1192                "migration_execution_scope_mismatch",
1193                "lease scope differs from the verified migration scope",
1194            ));
1195        }
1196        if !migration.backfill_step_indices().contains(&step_index) {
1197            return Err(failure(
1198                DiagnosticCategory::Integrity,
1199                "migration_execution_backfill_step_position",
1200                "backfill event position is outside the verified apply manifest",
1201            ));
1202        }
1203        let step = migration.steps().get(step_index).ok_or_else(|| {
1204            failure(
1205                DiagnosticCategory::Integrity,
1206                "migration_execution_backfill_step_position",
1207                "backfill event position is outside the verified apply manifest",
1208            )
1209        })?;
1210        let (contract, _) = step.step().as_backfill().ok_or_else(|| {
1211            failure(
1212                DiagnosticCategory::Integrity,
1213                "migration_execution_backfill_step_position",
1214                "verified backfill position does not contain a backfill step",
1215            )
1216        })?;
1217        Self::new_checked(
1218            lease,
1219            migration.digest(),
1220            migration.manifest().id().clone(),
1221            step_index,
1222            step_index,
1223            contract.plan_fingerprint().clone(),
1224            BackfillExecutionDirection::Forward,
1225            kind,
1226            completion,
1227        )
1228    }
1229
1230    /// Derive a reverse event from one exact mixed rollback operation position.
1231    pub fn new_rollback(
1232        lease: &MigrationLease,
1233        rollback: &VerifiedMigrationRollbackManifest,
1234        operation_index: usize,
1235        kind: GroupJournalEventKind,
1236        completion: Option<BackfillCompletionEvidence>,
1237    ) -> Result<Self, Diagnostic> {
1238        if lease.scope().managed_scope_id() != rollback.manifest().managed_scope().id() {
1239            return Err(failure(
1240                DiagnosticCategory::Integrity,
1241                "migration_execution_scope_mismatch",
1242                "lease scope differs from the verified rollback scope",
1243            ));
1244        }
1245        let backfill_index = match rollback.operations().get(operation_index) {
1246            Some(crate::VerifiedMigrationRollbackOperation::Backfill(index)) => *index,
1247            _ => {
1248                return Err(failure(
1249                    DiagnosticCategory::Integrity,
1250                    "migration_execution_backfill_step_position",
1251                    "backfill event position is outside the verified rollback operations",
1252                ));
1253            }
1254        };
1255        let step = rollback
1256            .backfills()
1257            .get(backfill_index)
1258            .ok_or_else(|| {
1259                failure(
1260                    DiagnosticCategory::Integrity,
1261                    "migration_execution_backfill_step_position",
1262                    "rollback backfill index is outside the verified reverse program",
1263                )
1264            })?
1265            .forward_step();
1266        let (contract, _) = step.as_backfill().ok_or_else(|| {
1267            failure(
1268                DiagnosticCategory::Integrity,
1269                "migration_execution_backfill_step_position",
1270                "verified rollback backfill does not contain a backfill step",
1271            )
1272        })?;
1273        let manifest_step_index = rollback
1274            .manifest()
1275            .steps()
1276            .iter()
1277            .position(|candidate| candidate == step)
1278            .ok_or_else(|| {
1279                failure(
1280                    DiagnosticCategory::Integrity,
1281                    "migration_execution_backfill_step_position",
1282                    "rollback backfill is absent from its verified manifest",
1283                )
1284            })?;
1285        Self::new_checked(
1286            lease,
1287            *rollback.digest(),
1288            rollback.manifest().id().clone(),
1289            operation_index,
1290            manifest_step_index,
1291            contract.plan_fingerprint().clone(),
1292            BackfillExecutionDirection::Reverse,
1293            kind,
1294            completion,
1295        )
1296    }
1297
1298    #[allow(clippy::too_many_arguments)]
1299    fn new_checked(
1300        lease: &MigrationLease,
1301        manifest_digest: MigrationManifestDigest,
1302        migration_id: MigrationId,
1303        operation_ordinal: usize,
1304        manifest_step_index: usize,
1305        plan_fingerprint: Fingerprint,
1306        direction: BackfillExecutionDirection,
1307        kind: GroupJournalEventKind,
1308        completion: Option<BackfillCompletionEvidence>,
1309    ) -> Result<Self, Diagnostic> {
1310        match kind {
1311            GroupJournalEventKind::Committed
1312                if completion.as_ref().is_some_and(|evidence| {
1313                    evidence.plan_fingerprint() == &plan_fingerprint
1314                        && evidence.direction() == direction
1315                }) => {}
1316            GroupJournalEventKind::Committed => {
1317                return Err(failure(
1318                    DiagnosticCategory::Integrity,
1319                    "migration_execution_backfill_completion_mismatch",
1320                    "committed backfill event requires exact terminal plan evidence",
1321                ));
1322            }
1323            GroupJournalEventKind::FormalOnlyAdvanced => {
1324                return Err(failure(
1325                    DiagnosticCategory::InvalidContract,
1326                    "migration_execution_backfill_formal_advance",
1327                    "a data backfill cannot advance as a formal-only operation",
1328                ));
1329            }
1330            _ if completion.is_none() => {}
1331            _ => {
1332                return Err(failure(
1333                    DiagnosticCategory::InvalidContract,
1334                    "migration_execution_unexpected_backfill_completion",
1335                    "only a committed backfill event may carry terminal evidence",
1336                ));
1337            }
1338        }
1339        Ok(Self {
1340            scope: lease.scope().clone(),
1341            fence: lease.fence(),
1342            manifest_digest,
1343            migration_id,
1344            operation_ordinal: position(operation_ordinal)?,
1345            manifest_step_index: position(manifest_step_index)?,
1346            plan_fingerprint,
1347            direction,
1348            kind,
1349            completion,
1350        })
1351    }
1352
1353    /// Return the execution scope.
1354    pub const fn scope(&self) -> &ExecutionScope {
1355        &self.scope
1356    }
1357
1358    /// Return the event fence.
1359    pub const fn fence(&self) -> ExecutionFence {
1360        self.fence
1361    }
1362
1363    /// Return the exact manifest digest.
1364    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1365        self.manifest_digest
1366    }
1367
1368    /// Return the migration identity.
1369    pub const fn migration_id(&self) -> &MigrationId {
1370        &self.migration_id
1371    }
1372
1373    /// Return the position in forward steps or mixed reverse operations.
1374    pub const fn operation_ordinal(&self) -> u32 {
1375        self.operation_ordinal
1376    }
1377
1378    /// Return the original canonical manifest step position.
1379    pub const fn manifest_step_index(&self) -> u32 {
1380        self.manifest_step_index
1381    }
1382
1383    /// Return the exact canonical backfill plan identity.
1384    pub const fn plan_fingerprint(&self) -> &Fingerprint {
1385        &self.plan_fingerprint
1386    }
1387
1388    /// Return the execution direction.
1389    pub const fn direction(&self) -> BackfillExecutionDirection {
1390        self.direction
1391    }
1392
1393    /// Return the commit-boundary event kind.
1394    pub const fn kind(&self) -> GroupJournalEventKind {
1395        self.kind
1396    }
1397
1398    /// Return terminal evidence carried only by committed events.
1399    pub const fn completion(&self) -> Option<&BackfillCompletionEvidence> {
1400        self.completion.as_ref()
1401    }
1402}
1403
1404/// Identity-only applied-ledger record for one verified manifest.
1405#[derive(Clone, Debug, Eq, PartialEq)]
1406pub struct AppliedRecord {
1407    scope: ExecutionScope,
1408    fence: ExecutionFence,
1409    migration_id: MigrationId,
1410    manifest_digest: MigrationManifestDigest,
1411    source_declared: ManagedDeclaredIdentityFingerprint,
1412    target_declared: ManagedDeclaredIdentityFingerprint,
1413    source_semantics: ManagedSemanticSchemaFingerprint,
1414    target_semantics: ManagedSemanticSchemaFingerprint,
1415}
1416
1417impl AppliedRecord {
1418    /// Derive an applied-ledger record from one exact verified manifest.
1419    pub fn from_verified_manifest(
1420        lease: &MigrationLease,
1421        migration: &VerifiedMigrationApplyManifest,
1422    ) -> Result<Self, Diagnostic> {
1423        Self::from_verified_manifest_contract(lease, migration.manifest())
1424    }
1425
1426    /// Derive an applied-ledger record from one verified manifest contract.
1427    ///
1428    /// This is the reconstruction seam for persistent stores. The manifest
1429    /// remains the trust anchor: its digest is recomputed from verified bytes,
1430    /// and no persisted record claim enters the trusted value.
1431    pub fn from_verified_manifest_contract(
1432        lease: &MigrationLease,
1433        manifest: &crate::VerifiedSchemaMigrationManifest,
1434    ) -> Result<Self, Diagnostic> {
1435        let source = manifest.source_state();
1436        let target = manifest.target_state();
1437        let scope = ExecutionScope::new(source.scope().id().clone());
1438        if lease.scope() != &scope || target.scope() != source.scope() {
1439            return Err(failure(
1440                DiagnosticCategory::Integrity,
1441                "migration_execution_scope_mismatch",
1442                "lease and manifest endpoints must bind the same managed scope",
1443            ));
1444        }
1445        Ok(Self {
1446            scope,
1447            fence: lease.fence(),
1448            migration_id: manifest.id().clone(),
1449            manifest_digest: crate::verified_manifest_digest(manifest)?,
1450            source_declared: source.managed_declared_identity().clone(),
1451            target_declared: target.managed_declared_identity().clone(),
1452            source_semantics: source.managed_semantic_schema().clone(),
1453            target_semantics: target.managed_semantic_schema().clone(),
1454        })
1455    }
1456
1457    /// Return the execution scope.
1458    pub const fn scope(&self) -> &ExecutionScope {
1459        &self.scope
1460    }
1461    /// Return the record fence.
1462    pub const fn fence(&self) -> ExecutionFence {
1463        self.fence
1464    }
1465    /// Return the migration identity.
1466    pub const fn migration_id(&self) -> &MigrationId {
1467        &self.migration_id
1468    }
1469    /// Return the canonical manifest digest.
1470    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1471        self.manifest_digest
1472    }
1473    /// Return source managed-declared identity.
1474    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1475        &self.source_declared
1476    }
1477    /// Return target managed-declared identity.
1478    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1479        &self.target_declared
1480    }
1481    /// Return source managed semantics.
1482    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1483        &self.source_semantics
1484    }
1485    /// Return target managed semantics.
1486    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1487        &self.target_semantics
1488    }
1489}
1490
1491/// Identity-only journal record for one complete verified rollback plan.
1492///
1493/// The record binds the exact reverse-topological order, manifest digests,
1494/// surviving applied set, and endpoint managed states the rollback executes
1495/// under. It is the rollback analogue of [`PlanRecord`].
1496#[derive(Clone, Debug, Eq, PartialEq)]
1497pub struct RollbackPlanRecord {
1498    scope: ExecutionScope,
1499    fence: ExecutionFence,
1500    source_applied: Vec<MigrationId>,
1501    rollback_ids: Vec<MigrationId>,
1502    manifest_digests: Vec<MigrationManifestDigest>,
1503    manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
1504    remaining_applied: Vec<MigrationId>,
1505    source_declared: ManagedDeclaredIdentityFingerprint,
1506    target_declared: ManagedDeclaredIdentityFingerprint,
1507    source_semantics: ManagedSemanticSchemaFingerprint,
1508    target_semantics: ManagedSemanticSchemaFingerprint,
1509    semantic_profile: SemanticProfileFingerprint,
1510    lowering_profile: SchemaLoweringProfileFingerprint,
1511    observed_live_source: ManagedSemanticSchemaFingerprint,
1512}
1513
1514impl RollbackPlanRecord {
1515    /// Bind a fresh-lease rollback ledger and live-state precondition to
1516    /// verified rollback plan identities.
1517    pub fn from_verified_rollback_plan(
1518        lease: &MigrationLease,
1519        plan: &VerifiedMigrationRollbackPlan,
1520        observed_applied_migrations: &[MigrationId],
1521        observed_live_source: &ManagedSchemaState,
1522    ) -> Result<Self, Diagnostic> {
1523        let first = plan.rollbacks().first().ok_or_else(|| {
1524            failure(
1525                DiagnosticCategory::InvalidContract,
1526                "migration_execution_empty_plan",
1527                "an executable rollback plan requires at least one manifest",
1528            )
1529        })?;
1530        let basis: Vec<MigrationId> = plan.applied_basis().into_iter().collect();
1531        if observed_applied_migrations != basis {
1532            return Err(failure(
1533                DiagnosticCategory::Integrity,
1534                "migration_execution_stale_applied_set",
1535                "applied ledger changed after rollback planning; rebuild the plan",
1536            ));
1537        }
1538        if observed_live_source != plan.source_state() {
1539            return Err(failure(
1540                DiagnosticCategory::Integrity,
1541                "migration_execution_stale_source_state",
1542                "live managed state differs from the planned source; rebuild the plan",
1543            ));
1544        }
1545        let scope = ExecutionScope::new(plan.source_state().scope().id().clone());
1546        if lease.scope() != &scope || plan.target_state().scope() != plan.source_state().scope() {
1547            return Err(failure(
1548                DiagnosticCategory::Integrity,
1549                "migration_execution_scope_mismatch",
1550                "lease, source, and target must bind the same managed scope",
1551            ));
1552        }
1553        let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
1554        let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
1555        for rollback in plan.rollbacks() {
1556            if rollback.manifest().managed_scope().id() != scope.managed_scope_id()
1557                || rollback.manifest().semantic_profile().fingerprint() != &semantic_profile
1558                || rollback.manifest().lowering_profile().fingerprint() != &lowering_profile
1559            {
1560                return Err(failure(
1561                    DiagnosticCategory::Integrity,
1562                    "migration_execution_plan_binding_mismatch",
1563                    "planned rollbacks do not share exact scope and profile bindings",
1564                ));
1565            }
1566        }
1567        Ok(Self {
1568            scope,
1569            fence: lease.fence(),
1570            source_applied: basis,
1571            rollback_ids: plan
1572                .rollbacks()
1573                .iter()
1574                .map(|rollback| rollback.manifest().id().clone())
1575                .collect(),
1576            manifest_digests: plan
1577                .rollbacks()
1578                .iter()
1579                .map(|rollback| *rollback.digest())
1580                .collect(),
1581            manifest_plan_fingerprints: plan
1582                .rollbacks()
1583                .iter()
1584                .map(|rollback| rollback.manifest().plan_fingerprint().clone())
1585                .collect(),
1586            remaining_applied: plan.remaining_applied().to_vec(),
1587            source_declared: plan.source_state().managed_declared_identity().clone(),
1588            target_declared: plan.target_state().managed_declared_identity().clone(),
1589            source_semantics: plan.source_state().managed_semantic_schema().clone(),
1590            target_semantics: plan.target_state().managed_semantic_schema().clone(),
1591            semantic_profile,
1592            lowering_profile,
1593            observed_live_source: observed_live_source.managed_semantic_schema().clone(),
1594        })
1595    }
1596
1597    /// Return the execution scope.
1598    pub const fn scope(&self) -> &ExecutionScope {
1599        &self.scope
1600    }
1601    /// Return the fence bound into this record.
1602    pub const fn fence(&self) -> ExecutionFence {
1603        self.fence
1604    }
1605    /// Return the complete canonically ordered pre-rollback applied set.
1606    pub fn source_applied(&self) -> &[MigrationId] {
1607        &self.source_applied
1608    }
1609    /// Return rolled-back identities in reverse-topological execution order.
1610    pub fn rollback_ids(&self) -> &[MigrationId] {
1611        &self.rollback_ids
1612    }
1613    /// Return ordered canonical manifest digests.
1614    pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
1615        &self.manifest_digests
1616    }
1617    /// Return ordered manifest plan fingerprints.
1618    pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
1619        &self.manifest_plan_fingerprints
1620    }
1621    /// Return the applied identities that survive the rollback, in order.
1622    pub fn remaining_applied(&self) -> &[MigrationId] {
1623        &self.remaining_applied
1624    }
1625    /// Return planned pre-rollback managed-declared identity.
1626    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1627        &self.source_declared
1628    }
1629    /// Return planned restored managed-declared identity.
1630    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1631        &self.target_declared
1632    }
1633    /// Return planned pre-rollback managed semantics.
1634    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1635        &self.source_semantics
1636    }
1637    /// Return planned restored managed semantics.
1638    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1639        &self.target_semantics
1640    }
1641    /// Return semantic-profile content identity.
1642    pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
1643        &self.semantic_profile
1644    }
1645    /// Return lowering-registry content identity.
1646    pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
1647        &self.lowering_profile
1648    }
1649    /// Return the pre-mutation observed source semantics.
1650    pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
1651        &self.observed_live_source
1652    }
1653}
1654
1655/// Identity-only journal event for one verified rollback step transaction.
1656///
1657/// Each rollback step executes one recorded reverse program in its own
1658/// provider transaction, so the commit-boundary vocabulary and recovery
1659/// decision table are shared with forward [`GroupEventRecord`] events.
1660#[derive(Clone, Debug, Eq, PartialEq)]
1661pub struct RollbackStepEventRecord {
1662    scope: ExecutionScope,
1663    fence: ExecutionFence,
1664    manifest_digest: MigrationManifestDigest,
1665    migration_id: MigrationId,
1666    step_ordinal: u32,
1667    kind: GroupJournalEventKind,
1668    observed_target: Option<ManagedSemanticSchemaFingerprint>,
1669}
1670
1671impl RollbackStepEventRecord {
1672    /// Derive a positional event from exact verified rollback evidence.
1673    pub fn new(
1674        lease: &MigrationLease,
1675        rollback: &VerifiedMigrationRollbackManifest,
1676        step_index: usize,
1677        kind: GroupJournalEventKind,
1678        observed_target: Option<ManagedSemanticSchemaFingerprint>,
1679    ) -> Result<Self, Diagnostic> {
1680        let step = rollback.steps().get(step_index).ok_or_else(|| {
1681            failure(
1682                DiagnosticCategory::Integrity,
1683                "migration_execution_rollback_step_position",
1684                "rollback event position is outside the verified rollback manifest",
1685            )
1686        })?;
1687        let reverse = rollback.reverse_delta(step)?;
1688        let scope = ExecutionScope::new(rollback.manifest().managed_scope().id().clone());
1689        if lease.scope() != &scope {
1690            return Err(failure(
1691                DiagnosticCategory::Integrity,
1692                "migration_execution_scope_mismatch",
1693                "lease scope differs from the verified rollback scope",
1694            ));
1695        }
1696        match kind {
1697            GroupJournalEventKind::Committed
1698                if observed_target.as_ref() == Some(reverse.target().managed_semantic_schema()) => {
1699            }
1700            GroupJournalEventKind::Committed => {
1701                return Err(failure(
1702                    DiagnosticCategory::Integrity,
1703                    "migration_execution_commit_evidence_mismatch",
1704                    "committed event requires the exact observed target semantics",
1705                ));
1706            }
1707            GroupJournalEventKind::FormalOnlyAdvanced
1708                if observed_target.is_none()
1709                    && step.lowering().units().is_empty()
1710                    && reverse.source().managed_semantic_schema()
1711                        == reverse.target().managed_semantic_schema() => {}
1712            GroupJournalEventKind::FormalOnlyAdvanced => {
1713                return Err(failure(
1714                    DiagnosticCategory::InvalidContract,
1715                    "migration_execution_invalid_formal_advance",
1716                    "formal-only advancement requires an empty equal-semantic reverse program",
1717                ));
1718            }
1719            _ if observed_target.is_none() => {}
1720            _ => {
1721                return Err(failure(
1722                    DiagnosticCategory::InvalidContract,
1723                    "migration_execution_unexpected_observation",
1724                    "only committed events may carry observed target semantics",
1725                ));
1726            }
1727        }
1728        Ok(Self {
1729            scope,
1730            fence: lease.fence(),
1731            manifest_digest: *rollback.digest(),
1732            migration_id: rollback.manifest().id().clone(),
1733            step_ordinal: position(step_index)?,
1734            kind,
1735            observed_target,
1736        })
1737    }
1738
1739    /// Return the execution scope.
1740    pub const fn scope(&self) -> &ExecutionScope {
1741        &self.scope
1742    }
1743    /// Return the event fence.
1744    pub const fn fence(&self) -> ExecutionFence {
1745        self.fence
1746    }
1747    /// Return the canonical digest of the manifest being rolled back.
1748    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1749        self.manifest_digest
1750    }
1751    /// Return the migration identity.
1752    pub const fn migration_id(&self) -> &MigrationId {
1753        &self.migration_id
1754    }
1755    /// Return the rollback step position in execution order.
1756    pub const fn step_ordinal(&self) -> u32 {
1757        self.step_ordinal
1758    }
1759    /// Return the event kind.
1760    pub const fn kind(&self) -> GroupJournalEventKind {
1761        self.kind
1762    }
1763    /// Return exact target evidence carried only by committed events.
1764    pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
1765        self.observed_target.as_ref()
1766    }
1767}
1768
1769/// Identity-only retirement record for one rolled-back applied migration.
1770///
1771/// Retirement is append-only history: the applied record stays in the durable
1772/// journal, and this record marks it inactive. The states enter swapped
1773/// relative to [`AppliedRecord`] because rolling back moves the managed schema
1774/// from the manifest's target back to its source.
1775#[derive(Clone, Debug, Eq, PartialEq)]
1776pub struct RolledBackRecord {
1777    scope: ExecutionScope,
1778    fence: ExecutionFence,
1779    migration_id: MigrationId,
1780    manifest_digest: MigrationManifestDigest,
1781    source_declared: ManagedDeclaredIdentityFingerprint,
1782    target_declared: ManagedDeclaredIdentityFingerprint,
1783    source_semantics: ManagedSemanticSchemaFingerprint,
1784    target_semantics: ManagedSemanticSchemaFingerprint,
1785}
1786
1787impl RolledBackRecord {
1788    /// Derive a retirement record from one exact verified rollback manifest.
1789    pub fn from_verified_rollback(
1790        lease: &MigrationLease,
1791        rollback: &VerifiedMigrationRollbackManifest,
1792    ) -> Result<Self, Diagnostic> {
1793        Self::from_verified_manifest_contract(lease, rollback.manifest())
1794    }
1795
1796    /// Derive a retirement record from one verified manifest contract.
1797    ///
1798    /// This is the reconstruction seam for persistent stores, mirroring
1799    /// [`AppliedRecord::from_verified_manifest_contract`].
1800    pub fn from_verified_manifest_contract(
1801        lease: &MigrationLease,
1802        manifest: &crate::VerifiedSchemaMigrationManifest,
1803    ) -> Result<Self, Diagnostic> {
1804        let source = manifest.target_state();
1805        let target = manifest.source_state();
1806        let scope = ExecutionScope::new(source.scope().id().clone());
1807        if lease.scope() != &scope || target.scope() != source.scope() {
1808            return Err(failure(
1809                DiagnosticCategory::Integrity,
1810                "migration_execution_scope_mismatch",
1811                "lease and manifest endpoints must bind the same managed scope",
1812            ));
1813        }
1814        Ok(Self {
1815            scope,
1816            fence: lease.fence(),
1817            migration_id: manifest.id().clone(),
1818            manifest_digest: crate::verified_manifest_digest(manifest)?,
1819            source_declared: source.managed_declared_identity().clone(),
1820            target_declared: target.managed_declared_identity().clone(),
1821            source_semantics: source.managed_semantic_schema().clone(),
1822            target_semantics: target.managed_semantic_schema().clone(),
1823        })
1824    }
1825
1826    /// Return the execution scope.
1827    pub const fn scope(&self) -> &ExecutionScope {
1828        &self.scope
1829    }
1830    /// Return the record fence.
1831    pub const fn fence(&self) -> ExecutionFence {
1832        self.fence
1833    }
1834    /// Return the retired migration identity.
1835    pub const fn migration_id(&self) -> &MigrationId {
1836        &self.migration_id
1837    }
1838    /// Return the canonical manifest digest.
1839    pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1840        self.manifest_digest
1841    }
1842    /// Return pre-rollback managed-declared identity.
1843    pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1844        &self.source_declared
1845    }
1846    /// Return restored managed-declared identity.
1847    pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1848        &self.target_declared
1849    }
1850    /// Return pre-rollback managed semantics.
1851    pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1852        &self.source_semantics
1853    }
1854    /// Return restored managed semantics.
1855    pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1856        &self.target_semantics
1857    }
1858}
1859
1860/// One trusted, still-open rollback plan and its ordered step events.
1861#[derive(Clone, Debug, Eq, PartialEq)]
1862pub struct OpenRollbackPlanRecord {
1863    plan: JournalEntry<RollbackPlanRecord>,
1864    events: Vec<JournalEntry<RollbackStepEventRecord>>,
1865    backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
1866}
1867
1868impl OpenRollbackPlanRecord {
1869    /// Rebuild store output while checking sequence, scope, fence, and manifest binding.
1870    ///
1871    /// The same monotonic-fence rule as [`OpenPlanRecord::from_store`] applies:
1872    /// recovery events may be written under later fences than the plan itself.
1873    pub fn from_store(
1874        plan: JournalEntry<RollbackPlanRecord>,
1875        events: Vec<JournalEntry<RollbackStepEventRecord>>,
1876    ) -> Result<Self, Diagnostic> {
1877        Self::from_store_with_backfills(plan, events, Vec::new())
1878    }
1879
1880    /// Rebuild store output including independently journaled reverse data-step events.
1881    pub fn from_store_with_backfills(
1882        plan: JournalEntry<RollbackPlanRecord>,
1883        events: Vec<JournalEntry<RollbackStepEventRecord>>,
1884        backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
1885    ) -> Result<Self, Diagnostic> {
1886        let mut previous = plan.sequence();
1887        let mut previous_fence = plan.record().fence();
1888        let mut ordered: Vec<(
1889            JournalSequence,
1890            ExecutionFence,
1891            &MigrationId,
1892            MigrationManifestDigest,
1893        )> = events
1894            .iter()
1895            .map(|event| {
1896                (
1897                    event.sequence(),
1898                    event.record().fence(),
1899                    event.record().migration_id(),
1900                    event.record().manifest_digest(),
1901                )
1902            })
1903            .chain(backfill_events.iter().map(|event| {
1904                (
1905                    event.sequence(),
1906                    event.record().fence(),
1907                    event.record().migration_id(),
1908                    event.record().manifest_digest(),
1909                )
1910            }))
1911            .collect();
1912        ordered.sort_by_key(|event| event.0);
1913        for (sequence, fence, migration_id, manifest_digest) in ordered {
1914            let manifest_index = plan
1915                .record()
1916                .manifest_digests()
1917                .iter()
1918                .position(|digest| digest == &manifest_digest);
1919            if sequence <= previous
1920                || fence < previous_fence
1921                || manifest_index.is_none_or(|index| {
1922                    plan.record().rollback_ids().get(index) != Some(migration_id)
1923                })
1924            {
1925                return Err(failure(
1926                    DiagnosticCategory::Integrity,
1927                    "migration_execution_invalid_open_plan",
1928                    "loaded open-rollback events are not ordered and bound to the plan",
1929                ));
1930            }
1931            previous = sequence;
1932            previous_fence = fence;
1933        }
1934        if events
1935            .iter()
1936            .any(|event| event.record().scope() != plan.record().scope())
1937            || backfill_events
1938                .iter()
1939                .any(|event| event.record().scope() != plan.record().scope())
1940        {
1941            return Err(failure(
1942                DiagnosticCategory::Integrity,
1943                "migration_execution_invalid_open_plan",
1944                "loaded open-rollback events are not bound to the plan scope",
1945            ));
1946        }
1947        Ok(Self {
1948            plan,
1949            events,
1950            backfill_events,
1951        })
1952    }
1953
1954    /// Return the sequenced rollback plan record.
1955    pub const fn plan(&self) -> &JournalEntry<RollbackPlanRecord> {
1956        &self.plan
1957    }
1958
1959    /// Return ordered sequenced rollback step events.
1960    pub fn events(&self) -> &[JournalEntry<RollbackStepEventRecord>] {
1961        &self.events
1962    }
1963
1964    /// Return ordered sequenced reverse backfill events.
1965    pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
1966        &self.backfill_events
1967    }
1968}
1969
1970/// Store-backed exclusive lease service.
1971pub trait MigrationLeaseStore: Send + Sync {
1972    /// Atomically acquire an available scope and issue a fence strictly greater
1973    /// than every fence previously issued for that scope.
1974    ///
1975    /// Lease expiry, liveness, and takeover policy are store-defined. Every
1976    /// takeover after release, expiry, or failure must issue a strictly greater
1977    /// fence so a surviving stale holder cannot mutate the journal.
1978    fn acquire<'a>(
1979        &'a self,
1980        scope: &'a ExecutionScope,
1981        holder: &'a LeaseHolderId,
1982    ) -> ExecutionFuture<'a, MigrationLease>;
1983
1984    /// Release only the exact currently active holder and fence.
1985    fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()>;
1986}
1987
1988/// Durable applied ledger and open-plan journal.
1989///
1990/// Every write must atomically reject a lease that is not the current active
1991/// scope, holder, and fence. There is deliberately no advisory validate method.
1992pub trait MigrationExecutionJournal: Send + Sync {
1993    /// Begin one fully verified plan after stale-ledger and live-source checks.
1994    ///
1995    /// An open plan of either direction is exclusive: the store must reject a
1996    /// new apply plan while a rollback plan is open and vice versa.
1997    fn begin_plan<'a>(
1998        &'a self,
1999        lease: &'a MigrationLease,
2000        record: PlanRecord,
2001    ) -> ExecutionFuture<'a, JournalEntry<PlanRecord>>;
2002
2003    /// Append one positional commit-boundary event under the active fence.
2004    fn record_group_event<'a>(
2005        &'a self,
2006        lease: &'a MigrationLease,
2007        record: GroupEventRecord,
2008    ) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>>;
2009
2010    /// Append one exact forward or reverse backfill commit-boundary event.
2011    fn record_backfill_event<'a>(
2012        &'a self,
2013        lease: &'a MigrationLease,
2014        record: BackfillEventRecord,
2015    ) -> ExecutionFuture<'a, JournalEntry<BackfillEventRecord>>;
2016
2017    /// Add one exact verified manifest from the open plan to the applied ledger.
2018    ///
2019    /// The store must reject a new record whose migration identity and digest do
2020    /// not occur at the same position in the open plan. Importing an existing
2021    /// ledger is a separate concern and must not use this execution write. An
2022    /// exact retry is idempotent only under the same fence; seeing the same
2023    /// migration under a newer fence requires reloading the ledger rather than
2024    /// treating differently fenced evidence as equal.
2025    fn record_applied<'a>(
2026        &'a self,
2027        lease: &'a MigrationLease,
2028        record: AppliedRecord,
2029    ) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>>;
2030
2031    /// Load ordered active applied records under the active fence.
2032    ///
2033    /// Retired records — applied records whose migration was rolled back by a
2034    /// later [`RolledBackRecord`] — are excluded. The full append-only history
2035    /// stays durable in the store; only the active ledger is the apply basis.
2036    fn load_applied<'a>(
2037        &'a self,
2038        lease: &'a MigrationLease,
2039    ) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>>;
2040
2041    /// Load the open plan and all ordered events written at its fence or later.
2042    fn load_open_plan<'a>(
2043        &'a self,
2044        lease: &'a MigrationLease,
2045    ) -> ExecutionFuture<'a, Option<OpenPlanRecord>>;
2046
2047    /// Begin one fully verified rollback plan after stale-ledger checks.
2048    ///
2049    /// Subject to the same open-plan exclusivity as [`Self::begin_plan`].
2050    fn begin_rollback_plan<'a>(
2051        &'a self,
2052        lease: &'a MigrationLease,
2053        record: RollbackPlanRecord,
2054    ) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>>;
2055
2056    /// Append one positional rollback commit-boundary event under the active fence.
2057    fn record_rollback_step_event<'a>(
2058        &'a self,
2059        lease: &'a MigrationLease,
2060        record: RollbackStepEventRecord,
2061    ) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>>;
2062
2063    /// Retire one exact applied record from the open rollback plan.
2064    ///
2065    /// The store must reject a record whose migration identity and digest do
2066    /// not occur at the same position in the open rollback plan, or whose
2067    /// migration is not currently active in the applied ledger. Retirement is
2068    /// append-only: the applied record and this record both stay durable, and
2069    /// the migration merely leaves the active ledger. An exact retry is
2070    /// idempotent only under the same fence.
2071    fn record_rolled_back<'a>(
2072        &'a self,
2073        lease: &'a MigrationLease,
2074        record: RolledBackRecord,
2075    ) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>>;
2076
2077    /// Load ordered retirement records under the active fence.
2078    fn load_rolled_back<'a>(
2079        &'a self,
2080        lease: &'a MigrationLease,
2081    ) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>>;
2082
2083    /// Load the open rollback plan and all ordered events at its fence or later.
2084    fn load_open_rollback_plan<'a>(
2085        &'a self,
2086        lease: &'a MigrationLease,
2087    ) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>>;
2088}
2089
2090/// Filter an applied ledger down to its active records.
2091///
2092/// Each retirement record consumes the latest not-yet-retired applied record
2093/// with the same migration identity and manifest digest written before it.
2094/// A retirement that matches no applied record is corrupt history and fails
2095/// closed. Stores share this exact matching rule so every journal
2096/// implementation reports the same active basis.
2097pub fn active_applied_entries(
2098    applied: Vec<JournalEntry<AppliedRecord>>,
2099    rolled_back: &[JournalEntry<RolledBackRecord>],
2100) -> Result<Vec<JournalEntry<AppliedRecord>>, Diagnostic> {
2101    let mut retired = vec![false; applied.len()];
2102    for retirement in rolled_back {
2103        let matched = applied
2104            .iter()
2105            .enumerate()
2106            .filter(|(index, entry)| {
2107                !retired[*index]
2108                    && entry.sequence() < retirement.sequence()
2109                    && entry.record().migration_id() == retirement.record().migration_id()
2110                    && entry.record().manifest_digest() == retirement.record().manifest_digest()
2111            })
2112            .max_by_key(|(_, entry)| entry.sequence())
2113            .map(|(index, _)| index);
2114        let Some(index) = matched else {
2115            return Err(failure(
2116                DiagnosticCategory::Integrity,
2117                "migration_execution_unmatched_retirement",
2118                "retirement record matches no active applied record before it",
2119            ));
2120        };
2121        retired[index] = true;
2122    }
2123    Ok(applied
2124        .into_iter()
2125        .zip(retired)
2126        .filter_map(|(entry, retired)| (!retired).then_some(entry))
2127        .collect())
2128}
2129
2130fn position(value: usize) -> Result<u32, Diagnostic> {
2131    u32::try_from(value).map_err(|_| {
2132        failure(
2133            DiagnosticCategory::ResourceLimit,
2134            "migration_execution_position_limit",
2135            "migration transaction position exceeds the canonical u32 range",
2136        )
2137    })
2138}
2139
2140fn failure(category: DiagnosticCategory, code: &'static str, message: &'static str) -> Diagnostic {
2141    Diagnostic::new(
2142        category,
2143        DiagnosticCode::new(code).expect("static migration execution diagnostic code"),
2144        message,
2145    )
2146}
2147
2148#[cfg(test)]
2149mod tests {
2150    use std::collections::BTreeMap;
2151    use std::future::Future;
2152    use std::sync::Mutex;
2153    use std::task::{Context, Poll, Waker};
2154
2155    use type_bridge_contract::fingerprint::SemanticProfileId;
2156    use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileBinding};
2157    use type_bridge_contract::migration::{MigrationAppLabel, MigrationName};
2158
2159    use super::*;
2160
2161    #[test]
2162    fn migration_controls_cancel_and_only_tighten_shared_ceilings() {
2163        let cancellation = MigrationCancellation::default();
2164        let limits = MigrationExecutionResourceLimits::tightened(7, usize::MAX);
2165        let control = MigrationExecutionControl::new(cancellation.clone(), None, limits);
2166
2167        assert_eq!(control.resources().transaction_groups(), 7);
2168        assert_eq!(
2169            control.resources().backfill_observations(),
2170            MAX_MIGRATION_BACKFILL_OBSERVATIONS
2171        );
2172        assert!(control.check().is_ok());
2173
2174        cancellation.cancel();
2175        cancellation.cancel();
2176        let diagnostic = control.check().expect_err("cancelled control");
2177        assert_eq!(diagnostic.category(), DiagnosticCategory::Cancelled);
2178        assert_eq!(diagnostic.code().as_str(), "migration_execution_cancelled");
2179    }
2180
2181    #[test]
2182    fn migration_control_uses_one_absolute_deadline() {
2183        let control = MigrationExecutionControl::new(
2184            MigrationCancellation::default(),
2185            Some(Instant::now()),
2186            MigrationExecutionResourceLimits::default(),
2187        );
2188        let diagnostic = control.check().expect_err("expired deadline");
2189        assert_eq!(diagnostic.category(), DiagnosticCategory::ResourceLimit);
2190        assert_eq!(
2191            diagnostic.code().as_str(),
2192            "migration_execution_deadline_exceeded"
2193        );
2194    }
2195    use crate::schema_lowering_profile_binding;
2196
2197    #[derive(Default)]
2198    struct ScopeState {
2199        highest_fence: u64,
2200        active_lease: Option<MigrationLease>,
2201        next_sequence: u64,
2202        applied: Vec<JournalEntry<AppliedRecord>>,
2203        rolled_back: Vec<JournalEntry<RolledBackRecord>>,
2204        open_plan: Option<JournalEntry<PlanRecord>>,
2205        events: Vec<JournalEntry<GroupEventRecord>>,
2206        backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
2207        open_rollback_plan: Option<JournalEntry<RollbackPlanRecord>>,
2208        rollback_events: Vec<JournalEntry<RollbackStepEventRecord>>,
2209    }
2210
2211    #[derive(Default)]
2212    struct InMemoryStore {
2213        scopes: Mutex<BTreeMap<ExecutionScope, ScopeState>>,
2214    }
2215
2216    impl InMemoryStore {
2217        fn check_lease<'a>(
2218            state: &'a mut ScopeState,
2219            lease: &MigrationLease,
2220        ) -> Result<&'a mut ScopeState, Diagnostic> {
2221            if state.active_lease.as_ref() != Some(lease) {
2222                return Err(failure(
2223                    DiagnosticCategory::Integrity,
2224                    "migration_execution_stale_fence",
2225                    "journal write does not carry the current active lease and fence",
2226                ));
2227            }
2228            Ok(state)
2229        }
2230
2231        fn sequence(state: &mut ScopeState) -> Result<JournalSequence, Diagnostic> {
2232            state.next_sequence = state.next_sequence.checked_add(1).ok_or_else(|| {
2233                failure(
2234                    DiagnosticCategory::ResourceLimit,
2235                    "migration_execution_sequence_exhausted",
2236                    "journal sequence range is exhausted",
2237                )
2238            })?;
2239            JournalSequence::new(state.next_sequence)
2240        }
2241    }
2242
2243    impl MigrationLeaseStore for InMemoryStore {
2244        fn acquire<'a>(
2245            &'a self,
2246            scope: &'a ExecutionScope,
2247            holder: &'a LeaseHolderId,
2248        ) -> ExecutionFuture<'a, MigrationLease> {
2249            Box::pin(async move {
2250                let mut scopes = self.scopes.lock().expect("store mutex");
2251                let state = scopes.entry(scope.clone()).or_default();
2252                if state.active_lease.is_some() {
2253                    return Err(failure(
2254                        DiagnosticCategory::InvalidContract,
2255                        "migration_execution_lease_contended",
2256                        "migration scope already has an active lease",
2257                    ));
2258                }
2259                let fence = if state.highest_fence == 0 {
2260                    ExecutionFence::new(1)?
2261                } else {
2262                    ExecutionFence::new(state.highest_fence)?.checked_successor()?
2263                };
2264                state.highest_fence = fence.get();
2265                let lease = MigrationLease::new(scope.clone(), holder.clone(), fence);
2266                state.active_lease = Some(lease.clone());
2267                Ok(lease)
2268            })
2269        }
2270
2271        fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()> {
2272            Box::pin(async move {
2273                let mut scopes = self.scopes.lock().expect("store mutex");
2274                let state = scopes.get_mut(lease.scope()).ok_or_else(|| {
2275                    failure(
2276                        DiagnosticCategory::Integrity,
2277                        "migration_execution_stale_fence",
2278                        "lease scope is not active",
2279                    )
2280                })?;
2281                Self::check_lease(state, lease)?;
2282                state.active_lease = None;
2283                Ok(())
2284            })
2285        }
2286    }
2287
2288    impl MigrationExecutionJournal for InMemoryStore {
2289        fn begin_plan<'a>(
2290            &'a self,
2291            lease: &'a MigrationLease,
2292            record: PlanRecord,
2293        ) -> ExecutionFuture<'a, JournalEntry<PlanRecord>> {
2294            Box::pin(async move {
2295                let mut scopes = self.scopes.lock().expect("store mutex");
2296                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2297                Self::check_lease(state, lease)?;
2298                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2299                    return Err(stale_fence());
2300                }
2301                if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
2302                    return Err(failure(
2303                        DiagnosticCategory::InvalidContract,
2304                        "migration_execution_plan_already_open",
2305                        "migration scope already has an open plan",
2306                    ));
2307                }
2308                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2309                state.open_plan = Some(entry.clone());
2310                Ok(entry)
2311            })
2312        }
2313
2314        fn record_group_event<'a>(
2315            &'a self,
2316            lease: &'a MigrationLease,
2317            record: GroupEventRecord,
2318        ) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>> {
2319            Box::pin(async move {
2320                let mut scopes = self.scopes.lock().expect("store mutex");
2321                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2322                Self::check_lease(state, lease)?;
2323                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2324                    return Err(stale_fence());
2325                }
2326                let plan = state.open_plan.as_ref().ok_or_else(|| {
2327                    failure(
2328                        DiagnosticCategory::InvalidContract,
2329                        "migration_execution_no_open_plan",
2330                        "group event requires an open migration plan",
2331                    )
2332                })?;
2333                if !plan
2334                    .record()
2335                    .manifest_digests()
2336                    .contains(&record.manifest_digest())
2337                {
2338                    return Err(failure(
2339                        DiagnosticCategory::Integrity,
2340                        "migration_execution_foreign_event",
2341                        "group event manifest is absent from the open plan",
2342                    ));
2343                }
2344                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2345                state.events.push(entry.clone());
2346                Ok(entry)
2347            })
2348        }
2349
2350        fn record_backfill_event<'a>(
2351            &'a self,
2352            lease: &'a MigrationLease,
2353            record: BackfillEventRecord,
2354        ) -> ExecutionFuture<'a, JournalEntry<BackfillEventRecord>> {
2355            Box::pin(async move {
2356                let mut scopes = self.scopes.lock().expect("store mutex");
2357                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2358                Self::check_lease(state, lease)?;
2359                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2360                    return Err(stale_fence());
2361                }
2362                let member = state.open_plan.as_ref().is_some_and(|plan| {
2363                    plan.record()
2364                        .manifest_digests()
2365                        .contains(&record.manifest_digest())
2366                }) || state.open_rollback_plan.as_ref().is_some_and(|plan| {
2367                    plan.record()
2368                        .manifest_digests()
2369                        .contains(&record.manifest_digest())
2370                });
2371                if !member {
2372                    return Err(failure(
2373                        DiagnosticCategory::Integrity,
2374                        "migration_execution_foreign_event",
2375                        "backfill event manifest is absent from the open plan",
2376                    ));
2377                }
2378                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2379                state.backfill_events.push(entry.clone());
2380                Ok(entry)
2381            })
2382        }
2383
2384        fn record_applied<'a>(
2385            &'a self,
2386            lease: &'a MigrationLease,
2387            record: AppliedRecord,
2388        ) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>> {
2389            Box::pin(async move {
2390                let mut scopes = self.scopes.lock().expect("store mutex");
2391                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2392                Self::check_lease(state, lease)?;
2393                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2394                    return Err(stale_fence());
2395                }
2396                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
2397                if let Some(existing) = active
2398                    .iter()
2399                    .find(|entry| entry.record().migration_id() == record.migration_id())
2400                {
2401                    return if existing.record() == &record {
2402                        Ok(existing.clone())
2403                    } else {
2404                        Err(failure(
2405                            DiagnosticCategory::Integrity,
2406                            "migration_execution_applied_identity_conflict",
2407                            "applied migration identity has different evidence",
2408                        ))
2409                    };
2410                }
2411                let plan = state.open_plan.as_ref().ok_or_else(|| {
2412                    failure(
2413                        DiagnosticCategory::InvalidContract,
2414                        "migration_execution_no_open_plan",
2415                        "applied migration requires an open migration plan",
2416                    )
2417                })?;
2418                let manifest_index = plan
2419                    .record()
2420                    .migration_ids()
2421                    .iter()
2422                    .position(|id| id == record.migration_id());
2423                if manifest_index.is_none_or(|index| {
2424                    plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
2425                }) {
2426                    return Err(failure(
2427                        DiagnosticCategory::Integrity,
2428                        "migration_execution_foreign_applied_record",
2429                        "applied migration identity and digest are absent from the open plan",
2430                    ));
2431                }
2432                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2433                state.applied.push(entry.clone());
2434                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
2435                let complete = state.open_plan.as_ref().is_some_and(|plan| {
2436                    plan.record().migration_ids().iter().all(|id| {
2437                        active
2438                            .iter()
2439                            .any(|applied| applied.record().migration_id() == id)
2440                    })
2441                });
2442                if complete {
2443                    state.open_plan = None;
2444                    state.events.clear();
2445                    state.backfill_events.clear();
2446                }
2447                Ok(entry)
2448            })
2449        }
2450
2451        fn load_applied<'a>(
2452            &'a self,
2453            lease: &'a MigrationLease,
2454        ) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>> {
2455            Box::pin(async move {
2456                let mut scopes = self.scopes.lock().expect("store mutex");
2457                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2458                Self::check_lease(state, lease)?;
2459                active_applied_entries(state.applied.clone(), &state.rolled_back)
2460            })
2461        }
2462
2463        fn load_open_plan<'a>(
2464            &'a self,
2465            lease: &'a MigrationLease,
2466        ) -> ExecutionFuture<'a, Option<OpenPlanRecord>> {
2467            Box::pin(async move {
2468                let mut scopes = self.scopes.lock().expect("store mutex");
2469                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2470                Self::check_lease(state, lease)?;
2471                let Some(plan) = state.open_plan.clone() else {
2472                    return Ok(None);
2473                };
2474                Ok(Some(OpenPlanRecord::from_store_with_backfills(
2475                    plan,
2476                    state.events.clone(),
2477                    state
2478                        .backfill_events
2479                        .iter()
2480                        .filter(|event| {
2481                            event.record().direction() == BackfillExecutionDirection::Forward
2482                        })
2483                        .cloned()
2484                        .collect(),
2485                )?))
2486            })
2487        }
2488
2489        fn begin_rollback_plan<'a>(
2490            &'a self,
2491            lease: &'a MigrationLease,
2492            record: RollbackPlanRecord,
2493        ) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>> {
2494            Box::pin(async move {
2495                let mut scopes = self.scopes.lock().expect("store mutex");
2496                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2497                Self::check_lease(state, lease)?;
2498                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2499                    return Err(stale_fence());
2500                }
2501                if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
2502                    return Err(failure(
2503                        DiagnosticCategory::InvalidContract,
2504                        "migration_execution_plan_already_open",
2505                        "migration scope already has an open plan",
2506                    ));
2507                }
2508                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2509                state.open_rollback_plan = Some(entry.clone());
2510                Ok(entry)
2511            })
2512        }
2513
2514        fn record_rollback_step_event<'a>(
2515            &'a self,
2516            lease: &'a MigrationLease,
2517            record: RollbackStepEventRecord,
2518        ) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>> {
2519            Box::pin(async move {
2520                let mut scopes = self.scopes.lock().expect("store mutex");
2521                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2522                Self::check_lease(state, lease)?;
2523                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2524                    return Err(stale_fence());
2525                }
2526                let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
2527                    failure(
2528                        DiagnosticCategory::InvalidContract,
2529                        "migration_execution_no_open_plan",
2530                        "rollback event requires an open rollback plan",
2531                    )
2532                })?;
2533                if !plan
2534                    .record()
2535                    .manifest_digests()
2536                    .contains(&record.manifest_digest())
2537                {
2538                    return Err(failure(
2539                        DiagnosticCategory::Integrity,
2540                        "migration_execution_foreign_event",
2541                        "rollback event manifest is absent from the open plan",
2542                    ));
2543                }
2544                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2545                state.rollback_events.push(entry.clone());
2546                Ok(entry)
2547            })
2548        }
2549
2550        fn record_rolled_back<'a>(
2551            &'a self,
2552            lease: &'a MigrationLease,
2553            record: RolledBackRecord,
2554        ) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>> {
2555            Box::pin(async move {
2556                let mut scopes = self.scopes.lock().expect("store mutex");
2557                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2558                Self::check_lease(state, lease)?;
2559                if record.scope() != lease.scope() || record.fence() != lease.fence() {
2560                    return Err(stale_fence());
2561                }
2562                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
2563                let is_active = active.iter().any(|entry| {
2564                    entry.record().migration_id() == record.migration_id()
2565                        && entry.record().manifest_digest() == record.manifest_digest()
2566                });
2567                if !is_active {
2568                    if let Some(existing) = state.rolled_back.iter().find(|entry| {
2569                        entry.record() == &record && entry.record().fence() == lease.fence()
2570                    }) {
2571                        return Ok(existing.clone());
2572                    }
2573                    return Err(failure(
2574                        DiagnosticCategory::Integrity,
2575                        "migration_execution_retirement_conflict",
2576                        "retirement target is not active in the applied ledger",
2577                    ));
2578                }
2579                let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
2580                    failure(
2581                        DiagnosticCategory::InvalidContract,
2582                        "migration_execution_no_open_plan",
2583                        "retirement requires an open rollback plan",
2584                    )
2585                })?;
2586                let manifest_index = plan
2587                    .record()
2588                    .rollback_ids()
2589                    .iter()
2590                    .position(|id| id == record.migration_id());
2591                if manifest_index.is_none_or(|index| {
2592                    plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
2593                }) {
2594                    return Err(failure(
2595                        DiagnosticCategory::Integrity,
2596                        "migration_execution_foreign_applied_record",
2597                        "retirement identity and digest are absent from the open plan",
2598                    ));
2599                }
2600                let entry = JournalEntry::from_store(Self::sequence(state)?, record);
2601                state.rolled_back.push(entry.clone());
2602                let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
2603                let complete = state.open_rollback_plan.as_ref().is_some_and(|plan| {
2604                    plan.record().rollback_ids().iter().all(|id| {
2605                        !active
2606                            .iter()
2607                            .any(|applied| applied.record().migration_id() == id)
2608                    })
2609                });
2610                if complete {
2611                    state.open_rollback_plan = None;
2612                    state.rollback_events.clear();
2613                    state.backfill_events.clear();
2614                }
2615                Ok(entry)
2616            })
2617        }
2618
2619        fn load_rolled_back<'a>(
2620            &'a self,
2621            lease: &'a MigrationLease,
2622        ) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>> {
2623            Box::pin(async move {
2624                let mut scopes = self.scopes.lock().expect("store mutex");
2625                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2626                Self::check_lease(state, lease)?;
2627                Ok(state.rolled_back.clone())
2628            })
2629        }
2630
2631        fn load_open_rollback_plan<'a>(
2632            &'a self,
2633            lease: &'a MigrationLease,
2634        ) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>> {
2635            Box::pin(async move {
2636                let mut scopes = self.scopes.lock().expect("store mutex");
2637                let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
2638                Self::check_lease(state, lease)?;
2639                let Some(plan) = state.open_rollback_plan.clone() else {
2640                    return Ok(None);
2641                };
2642                Ok(Some(OpenRollbackPlanRecord::from_store_with_backfills(
2643                    plan,
2644                    state.rollback_events.clone(),
2645                    state
2646                        .backfill_events
2647                        .iter()
2648                        .filter(|event| {
2649                            event.record().direction() == BackfillExecutionDirection::Reverse
2650                        })
2651                        .cloned()
2652                        .collect(),
2653                )?))
2654            })
2655        }
2656    }
2657
2658    fn stale_fence() -> Diagnostic {
2659        failure(
2660            DiagnosticCategory::Integrity,
2661            "migration_execution_stale_fence",
2662            "journal write does not carry the current active lease and fence",
2663        )
2664    }
2665
2666    fn scope() -> ExecutionScope {
2667        ExecutionScope::new(ManagedScopeId::new("journal-test").expect("scope"))
2668    }
2669
2670    fn migration_id() -> MigrationId {
2671        MigrationId::from_components(
2672            MigrationAppLabel::new("example").expect("app"),
2673            MigrationName::new("0001_initial").expect("name"),
2674        )
2675    }
2676
2677    fn semantic_fingerprint(bytes: &[u8]) -> ManagedSemanticSchemaFingerprint {
2678        ManagedSemanticSchemaFingerprint::compute(
2679            SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
2680            bytes,
2681        )
2682        .expect("semantic fingerprint")
2683    }
2684
2685    fn fake_plan(lease: &MigrationLease) -> PlanRecord {
2686        let id = migration_id();
2687        let semantic = SemanticProfileBinding::resolve(
2688            SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
2689        )
2690        .expect("semantic binding");
2691        PlanRecord {
2692            scope: lease.scope().clone(),
2693            fence: lease.fence(),
2694            source_applied: Vec::new(),
2695            source_frontier: Vec::new(),
2696            target_frontier: vec![id.clone()],
2697            migration_ids: vec![id],
2698            manifest_digests: vec![MigrationManifestDigest::compute(b"manifest")],
2699            manifest_plan_fingerprints: vec![
2700                MigrationPlanFingerprint::compute(&[]).expect("plan fingerprint"),
2701            ],
2702            source_declared: ManagedDeclaredIdentityFingerprint::compute(b"source")
2703                .expect("source declared"),
2704            target_declared: ManagedDeclaredIdentityFingerprint::compute(b"target")
2705                .expect("target declared"),
2706            source_semantics: semantic_fingerprint(b"source"),
2707            target_semantics: semantic_fingerprint(b"target"),
2708            semantic_profile: semantic.fingerprint().clone(),
2709            lowering_profile: schema_lowering_profile_binding()
2710                .expect("lowering binding")
2711                .fingerprint()
2712                .clone(),
2713            observed_live_source: semantic_fingerprint(b"source"),
2714        }
2715    }
2716
2717    fn fake_event(lease: &MigrationLease, plan: &PlanRecord) -> GroupEventRecord {
2718        GroupEventRecord {
2719            scope: lease.scope().clone(),
2720            fence: lease.fence(),
2721            manifest_digest: plan.manifest_digests()[0],
2722            migration_id: plan.migration_ids()[0].clone(),
2723            group_ordinal: 0,
2724            first_step_index: 0,
2725            schema_delta_step_index: 0,
2726            end_step_index: 1,
2727            kind: GroupJournalEventKind::BeforeCommit,
2728            observed_target: None,
2729        }
2730    }
2731
2732    fn fake_applied(lease: &MigrationLease, plan: &PlanRecord) -> AppliedRecord {
2733        AppliedRecord {
2734            scope: lease.scope().clone(),
2735            fence: lease.fence(),
2736            migration_id: plan.migration_ids()[0].clone(),
2737            manifest_digest: plan.manifest_digests()[0],
2738            source_declared: plan.source_declared().clone(),
2739            target_declared: plan.target_declared().clone(),
2740            source_semantics: plan.source_semantics().clone(),
2741            target_semantics: plan.target_semantics().clone(),
2742        }
2743    }
2744
2745    #[test]
2746    fn bound_lease_equality_includes_process_local_binding_identity() {
2747        let scope = scope();
2748        let holder = LeaseHolderId::new("owner").expect("holder");
2749        let fence = ExecutionFence::new(1).expect("fence");
2750        let token = ExecutionBindingToken::fresh();
2751        let same = MigrationLease::new_bound(scope.clone(), holder.clone(), fence, token.clone());
2752        let clone = same.clone();
2753        let foreign = MigrationLease::new_bound(
2754            scope.clone(),
2755            holder.clone(),
2756            fence,
2757            ExecutionBindingToken::fresh(),
2758        );
2759        let unbound = MigrationLease::new(scope.clone(), holder.clone(), fence);
2760        let same_unbound = MigrationLease::new(scope, holder, fence);
2761
2762        assert_eq!(same, clone);
2763        assert_ne!(same, foreign);
2764        assert_ne!(same, unbound);
2765        assert_eq!(unbound, same_unbound);
2766    }
2767
2768    #[test]
2769    fn store_assigns_monotonic_sequences_and_rejects_every_stale_fence_write() {
2770        let store = InMemoryStore::default();
2771        let scope = scope();
2772        let holder_a = LeaseHolderId::new("owner-a").expect("holder");
2773        let holder_b = LeaseHolderId::new("owner-b").expect("holder");
2774        let lease_a = block_on(store.acquire(&scope, &holder_a)).expect("lease a");
2775        assert!(block_on(store.acquire(&scope, &holder_b)).is_err());
2776        let plan_a = fake_plan(&lease_a);
2777        let plan_entry = block_on(store.begin_plan(&lease_a, plan_a.clone())).expect("plan");
2778        let event_entry =
2779            block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a)))
2780                .expect("event");
2781        assert_eq!(plan_entry.sequence().get(), 1);
2782        assert_eq!(event_entry.sequence().get(), 2);
2783        block_on(store.release(&lease_a)).expect("release a");
2784        let lease_b = block_on(store.acquire(&scope, &holder_b)).expect("lease b");
2785        assert!(lease_b.fence() > lease_a.fence());
2786        assert!(block_on(store.begin_plan(&lease_a, plan_a.clone())).is_err());
2787        assert!(
2788            block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a),)).is_err()
2789        );
2790        assert!(
2791            block_on(store.record_applied(&lease_a, fake_applied(&lease_a, &plan_a),)).is_err()
2792        );
2793        assert!(block_on(store.release(&lease_a)).is_err());
2794        let recovery_event =
2795            block_on(store.record_group_event(&lease_b, fake_event(&lease_b, &plan_a)))
2796                .expect("recovery event");
2797        assert_eq!(recovery_event.sequence().get(), 3);
2798        block_on(store.release(&lease_b)).expect("release recovery lease");
2799        let holder_c = LeaseHolderId::new("owner-c").expect("holder");
2800        let lease_c = block_on(store.acquire(&scope, &holder_c)).expect("lease c");
2801        assert!(lease_c.fence() > lease_b.fence());
2802        let open = block_on(store.load_open_plan(&lease_c))
2803            .expect("open plan")
2804            .expect("plan remains open after recovery crash");
2805        assert_eq!(open.events().len(), 2);
2806        assert_eq!(open.events()[0].record().fence(), lease_a.fence());
2807        assert_eq!(open.events()[1].record().fence(), lease_b.fence());
2808    }
2809
2810    #[test]
2811    fn applied_records_are_plan_bound_idempotent_and_close_the_completed_plan() {
2812        let store = InMemoryStore::default();
2813        let scope = scope();
2814        let holder = LeaseHolderId::new("applied-owner").expect("holder");
2815        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
2816        let plan = fake_plan(&lease);
2817        let applied = fake_applied(&lease, &plan);
2818
2819        assert!(block_on(store.record_applied(&lease, applied.clone())).is_err());
2820        block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
2821
2822        let mut foreign = applied.clone();
2823        foreign.manifest_digest = MigrationManifestDigest::compute(b"foreign");
2824        assert!(block_on(store.record_applied(&lease, foreign)).is_err());
2825
2826        let first = block_on(store.record_applied(&lease, applied.clone()))
2827            .expect("apply planned manifest");
2828        let duplicate = block_on(store.record_applied(&lease, applied))
2829            .expect("same-fence duplicate is idempotent");
2830        assert_eq!(duplicate, first);
2831        assert!(
2832            block_on(store.load_open_plan(&lease))
2833                .expect("load completed plan")
2834                .is_none()
2835        );
2836        assert_eq!(
2837            block_on(store.load_applied(&lease)).expect("load applied ledger"),
2838            vec![first],
2839        );
2840    }
2841
2842    fn fake_rollback_plan(lease: &MigrationLease, plan: &PlanRecord) -> RollbackPlanRecord {
2843        RollbackPlanRecord {
2844            scope: lease.scope().clone(),
2845            fence: lease.fence(),
2846            source_applied: plan.migration_ids().to_vec(),
2847            rollback_ids: plan.migration_ids().to_vec(),
2848            manifest_digests: plan.manifest_digests().to_vec(),
2849            manifest_plan_fingerprints: plan.manifest_plan_fingerprints().to_vec(),
2850            remaining_applied: Vec::new(),
2851            source_declared: plan.target_declared().clone(),
2852            target_declared: plan.source_declared().clone(),
2853            source_semantics: plan.target_semantics().clone(),
2854            target_semantics: plan.source_semantics().clone(),
2855            semantic_profile: plan.semantic_profile().clone(),
2856            lowering_profile: plan.lowering_profile().clone(),
2857            observed_live_source: plan.target_semantics().clone(),
2858        }
2859    }
2860
2861    fn fake_rolled_back(lease: &MigrationLease, plan: &PlanRecord) -> RolledBackRecord {
2862        RolledBackRecord {
2863            scope: lease.scope().clone(),
2864            fence: lease.fence(),
2865            migration_id: plan.migration_ids()[0].clone(),
2866            manifest_digest: plan.manifest_digests()[0],
2867            source_declared: plan.target_declared().clone(),
2868            target_declared: plan.source_declared().clone(),
2869            source_semantics: plan.target_semantics().clone(),
2870            target_semantics: plan.source_semantics().clone(),
2871        }
2872    }
2873
2874    #[test]
2875    fn retirement_is_append_only_exclusive_and_reopens_the_identity() {
2876        let store = InMemoryStore::default();
2877        let scope = scope();
2878        let holder = LeaseHolderId::new("retirement-owner").expect("holder");
2879        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
2880        let plan = fake_plan(&lease);
2881        let rollback_plan = fake_rollback_plan(&lease, &plan);
2882        let retirement = fake_rolled_back(&lease, &plan);
2883
2884        // Apply the migration through an ordinary forward plan.
2885        block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
2886        assert!(
2887            block_on(store.begin_rollback_plan(&lease, rollback_plan.clone())).is_err(),
2888            "open plans of either direction must be exclusive"
2889        );
2890        block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
2891            .expect("apply planned manifest");
2892        assert_eq!(
2893            block_on(store.load_applied(&lease)).expect("active").len(),
2894            1
2895        );
2896
2897        // Retiring outside an open rollback plan fails closed.
2898        assert!(block_on(store.record_rolled_back(&lease, retirement.clone())).is_err());
2899        block_on(store.begin_rollback_plan(&lease, rollback_plan)).expect("open rollback plan");
2900        assert!(
2901            block_on(store.begin_plan(&lease, plan.clone())).is_err(),
2902            "an open rollback plan must block a new apply plan"
2903        );
2904        let first = block_on(store.record_rolled_back(&lease, retirement.clone()))
2905            .expect("retire applied migration");
2906        let duplicate = block_on(store.record_rolled_back(&lease, retirement))
2907            .expect("same-fence duplicate retirement is idempotent");
2908        assert_eq!(duplicate, first);
2909        assert!(
2910            block_on(store.load_open_rollback_plan(&lease))
2911                .expect("closed rollback plan")
2912                .is_none()
2913        );
2914        assert!(
2915            block_on(store.load_applied(&lease))
2916                .expect("active")
2917                .is_empty()
2918        );
2919        assert_eq!(
2920            block_on(store.load_rolled_back(&lease))
2921                .expect("retired")
2922                .len(),
2923            1,
2924        );
2925
2926        // The identity is free again: a fresh forward plan re-applies it.
2927        block_on(store.begin_plan(&lease, plan.clone())).expect("reopen plan");
2928        block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
2929            .expect("re-apply retired migration");
2930        assert_eq!(
2931            block_on(store.load_applied(&lease)).expect("active").len(),
2932            1
2933        );
2934    }
2935
2936    #[test]
2937    fn unmatched_retirements_are_corrupt_history() {
2938        let store = InMemoryStore::default();
2939        let scope = scope();
2940        let holder = LeaseHolderId::new("orphan-owner").expect("holder");
2941        let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
2942        let plan = fake_plan(&lease);
2943        let orphan = JournalEntry::from_store(
2944            JournalSequence::new(7).expect("sequence"),
2945            fake_rolled_back(&lease, &plan),
2946        );
2947        let error = active_applied_entries(Vec::new(), &[orphan])
2948            .expect_err("a retirement without its applied record is corrupt");
2949        assert_eq!(
2950            error.code().as_str(),
2951            "migration_execution_unmatched_retirement"
2952        );
2953    }
2954
2955    #[test]
2956    fn recovery_decision_table_is_exhaustive_for_distinct_and_equal_semantics() {
2957        let source = semantic_fingerprint(b"source");
2958        let target = semantic_fingerprint(b"target");
2959        let neither = semantic_fingerprint(b"neither");
2960        let events = [
2961            None,
2962            Some(GroupJournalEventKind::BeforeCommit),
2963            Some(GroupJournalEventKind::Committed),
2964            Some(GroupJournalEventKind::CommitOutcomeUnknown),
2965            Some(GroupJournalEventKind::DefinitelyAborted),
2966            Some(GroupJournalEventKind::FormalOnlyAdvanced),
2967        ];
2968        let observations = [
2969            GroupRecoveryObservation::ManagedSemantics(source.clone()),
2970            GroupRecoveryObservation::ManagedSemantics(target.clone()),
2971            GroupRecoveryObservation::ManagedSemantics(neither.clone()),
2972            GroupRecoveryObservation::Unavailable,
2973        ];
2974        for event in events {
2975            for observation in &observations {
2976                let expected = expected_distinct(event, observation, &source, &target);
2977                assert_eq!(
2978                    decide_group_recovery(event, observation, &source, &target),
2979                    expected,
2980                    "distinct case {event:?} {observation:?}",
2981                );
2982            }
2983        }
2984        let equal_observations = [
2985            GroupRecoveryObservation::ManagedSemantics(source.clone()),
2986            GroupRecoveryObservation::ManagedSemantics(neither),
2987            GroupRecoveryObservation::Unavailable,
2988        ];
2989        for event in events {
2990            for observation in &equal_observations {
2991                let expected = expected_equal(event, observation, &source);
2992                assert_eq!(
2993                    decide_group_recovery(event, observation, &source, &source),
2994                    expected,
2995                    "equal case {event:?} {observation:?}",
2996                );
2997            }
2998        }
2999    }
3000
3001    #[test]
3002    fn backfill_counts_and_recovery_fail_closed_without_exact_completion() {
3003        assert_eq!(
3004            BackfillExecutionCounts::new(5, 3, 1, 1)
3005                .expect_err("inconsistent counts reject")
3006                .code()
3007                .as_str(),
3008            "migration_backfill_count_mismatch"
3009        );
3010        assert_eq!(
3011            BackfillExecutionCounts::new(0, 0, 0, 0)
3012                .expect_err("zero transaction groups reject")
3013                .code()
3014                .as_str(),
3015            "migration_backfill_zero_transaction_groups"
3016        );
3017
3018        let plan = semantic_fingerprint(b"backfill-plan")
3019            .as_fingerprint()
3020            .clone();
3021        let foreign = semantic_fingerprint(b"foreign-backfill-plan")
3022            .as_fingerprint()
3023            .clone();
3024        let counts = BackfillExecutionCounts::new(5, 3, 2, 2).expect("valid counts");
3025        let complete = BackfillRecoveryObservation::Complete(BackfillCompletionEvidence::new(
3026            plan.clone(),
3027            BackfillExecutionDirection::Forward,
3028            counts,
3029        ));
3030
3031        assert_eq!(
3032            decide_backfill_recovery(
3033                None,
3034                &BackfillRecoveryObservation::Incomplete,
3035                &plan,
3036                BackfillExecutionDirection::Forward,
3037            ),
3038            GroupRecoveryDecision::ExecuteNormally
3039        );
3040        assert_eq!(
3041            decide_backfill_recovery(None, &complete, &plan, BackfillExecutionDirection::Forward,),
3042            GroupRecoveryDecision::ExecuteNormally
3043        );
3044        assert_eq!(
3045            decide_backfill_recovery(
3046                Some(GroupJournalEventKind::BeforeCommit),
3047                &BackfillRecoveryObservation::Incomplete,
3048                &plan,
3049                BackfillExecutionDirection::Forward,
3050            ),
3051            GroupRecoveryDecision::RequiresExplicitRecovery
3052        );
3053        assert_eq!(
3054            decide_backfill_recovery(
3055                Some(GroupJournalEventKind::CommitOutcomeUnknown),
3056                &complete,
3057                &plan,
3058                BackfillExecutionDirection::Forward,
3059            ),
3060            GroupRecoveryDecision::RepairCheckpoint
3061        );
3062        assert_eq!(
3063            decide_backfill_recovery(
3064                Some(GroupJournalEventKind::Committed),
3065                &complete,
3066                &foreign,
3067                BackfillExecutionDirection::Forward,
3068            ),
3069            GroupRecoveryDecision::RequiresExplicitRecovery
3070        );
3071        assert_eq!(
3072            decide_backfill_recovery(
3073                Some(GroupJournalEventKind::Committed),
3074                &complete,
3075                &plan,
3076                BackfillExecutionDirection::Reverse,
3077            ),
3078            GroupRecoveryDecision::RequiresExplicitRecovery
3079        );
3080    }
3081
3082    fn expected_distinct(
3083        event: Option<GroupJournalEventKind>,
3084        observation: &GroupRecoveryObservation,
3085        source: &ManagedSemanticSchemaFingerprint,
3086        target: &ManagedSemanticSchemaFingerprint,
3087    ) -> GroupRecoveryDecision {
3088        let source_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == source);
3089        let target_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == target);
3090        match event {
3091            None | Some(GroupJournalEventKind::DefinitelyAborted) if source_seen => {
3092                GroupRecoveryDecision::ExecuteNormally
3093            }
3094            Some(
3095                GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown,
3096            ) if source_seen => GroupRecoveryDecision::ExecuteNormally,
3097            Some(
3098                GroupJournalEventKind::BeforeCommit
3099                | GroupJournalEventKind::CommitOutcomeUnknown
3100                | GroupJournalEventKind::Committed,
3101            ) if target_seen => GroupRecoveryDecision::RepairCheckpoint,
3102            _ => GroupRecoveryDecision::RequiresExplicitRecovery,
3103        }
3104    }
3105
3106    fn expected_equal(
3107        event: Option<GroupJournalEventKind>,
3108        observation: &GroupRecoveryObservation,
3109        both: &ManagedSemanticSchemaFingerprint,
3110    ) -> GroupRecoveryDecision {
3111        let both_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == both);
3112        if !both_seen {
3113            return GroupRecoveryDecision::RequiresExplicitRecovery;
3114        }
3115        match event {
3116            None | Some(GroupJournalEventKind::DefinitelyAborted) => {
3117                GroupRecoveryDecision::ExecuteNormally
3118            }
3119            Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced) => {
3120                GroupRecoveryDecision::RepairCheckpoint
3121            }
3122            _ => GroupRecoveryDecision::RequiresExplicitRecovery,
3123        }
3124    }
3125
3126    fn block_on<F: Future>(future: F) -> F::Output {
3127        let mut context = Context::from_waker(Waker::noop());
3128        let mut future = Box::pin(future);
3129        loop {
3130            match future.as_mut().poll(&mut context) {
3131                Poll::Ready(output) => return output,
3132                Poll::Pending => std::thread::yield_now(),
3133            }
3134        }
3135    }
3136}