1use 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
30pub const MAX_MIGRATION_EXECUTION_GROUPS: usize = 65_536;
32pub const MAX_MIGRATION_BACKFILL_OBSERVATIONS: usize = 65_536;
34
35#[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 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 #[must_use]
63 pub fn is_cancelled(&self) -> bool {
64 self.inner.cancelled.load(Ordering::Acquire)
65 }
66
67 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
91pub struct MigrationExecutionResourceLimits {
92 transaction_groups: usize,
93 backfill_observations: usize,
94}
95
96impl MigrationExecutionResourceLimits {
97 #[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 pub const fn transaction_groups(self) -> usize {
116 self.transaction_groups
117 }
118
119 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#[derive(Clone, Debug, Default)]
136pub struct MigrationExecutionControl {
137 cancellation: MigrationCancellation,
138 deadline: Option<Instant>,
139 resources: MigrationExecutionResourceLimits,
140}
141
142impl MigrationExecutionControl {
143 #[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 pub const fn cancellation(&self) -> &MigrationCancellation {
159 &self.cancellation
160 }
161
162 pub const fn deadline(&self) -> Option<Instant> {
164 self.deadline
165 }
166
167 pub const fn resources(&self) -> MigrationExecutionResourceLimits {
169 self.resources
170 }
171
172 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
195pub type ExecutionFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, Diagnostic>> + Send + 'a>>;
197
198#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
200pub struct ExecutionFence(u64);
201
202impl ExecutionFence {
203 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 pub const fn get(self) -> u64 {
217 self.0
218 }
219
220 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#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
235pub struct ExecutionScope(ManagedScopeId);
236
237impl ExecutionScope {
238 pub const fn new(scope: ManagedScopeId) -> Self {
240 Self(scope)
241 }
242
243 pub const fn managed_scope_id(&self) -> &ManagedScopeId {
245 &self.0
246 }
247}
248
249#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
251pub struct LeaseHolderId(String);
252
253impl LeaseHolderId {
254 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 pub fn as_str(&self) -> &str {
274 &self.0
275 }
276}
277
278#[derive(Clone, Debug)]
280pub struct MigrationLease {
281 scope: ExecutionScope,
282 holder: LeaseHolderId,
283 fence: ExecutionFence,
284 local_binding: Option<ExecutionBindingToken>,
285}
286
287#[doc(hidden)]
294#[derive(Clone)]
295pub struct ExecutionBindingToken(Arc<()>);
296
297impl ExecutionBindingToken {
298 #[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 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 #[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 #[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 pub const fn scope(&self) -> &ExecutionScope {
355 &self.scope
356 }
357
358 pub const fn holder(&self) -> &LeaseHolderId {
360 &self.holder
361 }
362
363 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
390pub enum GroupCommitCertainty {
391 DefinitelyAborted,
393 Unknown,
395}
396
397impl GroupCommitCertainty {
398 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#[derive(Clone, Copy, Debug, Eq, PartialEq)]
409pub enum GroupJournalEventKind {
410 BeforeCommit,
412 Committed,
414 CommitOutcomeUnknown,
416 DefinitelyAborted,
418 FormalOnlyAdvanced,
420}
421
422#[derive(Clone, Debug, Eq, PartialEq)]
424pub enum GroupRecoveryObservation {
425 Unavailable,
427 ManagedSemantics(ManagedSemanticSchemaFingerprint),
429}
430
431#[derive(Clone, Copy, Debug, Eq, PartialEq)]
433pub enum GroupRecoveryDecision {
434 ExecuteNormally,
436 RepairCheckpoint,
438 RequiresExplicitRecovery,
440}
441
442#[derive(Clone, Copy, Debug, Eq, PartialEq)]
444pub enum BackfillExecutionDirection {
445 Forward,
447 Reverse,
449}
450
451#[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 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 pub const fn matched(self) -> u64 {
492 self.matched
493 }
494
495 pub const fn changed(self) -> u64 {
497 self.changed
498 }
499
500 pub const fn skipped(self) -> u64 {
502 self.skipped
503 }
504
505 pub const fn transaction_groups(self) -> u32 {
507 self.transaction_groups
508 }
509}
510
511#[derive(Clone, Debug, Eq, PartialEq)]
513pub struct BackfillCompletionEvidence {
514 plan_fingerprint: Fingerprint,
515 direction: BackfillExecutionDirection,
516 counts: BackfillExecutionCounts,
517}
518
519impl BackfillCompletionEvidence {
520 #[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 pub const fn plan_fingerprint(&self) -> &Fingerprint {
536 &self.plan_fingerprint
537 }
538
539 pub const fn direction(&self) -> BackfillExecutionDirection {
541 self.direction
542 }
543
544 pub const fn counts(&self) -> BackfillExecutionCounts {
546 self.counts
547 }
548}
549
550#[derive(Clone, Debug, Eq, PartialEq)]
552pub enum BackfillRecoveryObservation {
553 Unavailable,
555 Incomplete,
557 Complete(BackfillCompletionEvidence),
559}
560
561pub 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
603pub 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#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
665pub struct JournalSequence(u64);
666
667impl JournalSequence {
668 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 pub const fn get(self) -> u64 {
682 self.0
683 }
684}
685
686#[derive(Clone, Debug, Eq, PartialEq)]
688pub struct JournalEntry<T> {
689 sequence: JournalSequence,
690 record: T,
691}
692
693impl<T> JournalEntry<T> {
694 pub const fn from_store(sequence: JournalSequence, record: T) -> Self {
696 Self { sequence, record }
697 }
698
699 pub const fn sequence(&self) -> JournalSequence {
701 self.sequence
702 }
703
704 pub const fn record(&self) -> &T {
706 &self.record
707 }
708
709 pub fn into_record(self) -> T {
711 self.record
712 }
713}
714
715#[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 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 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 pub const fn plan(&self) -> &JournalEntry<PlanRecord> {
813 &self.plan
814 }
815
816 pub fn events(&self) -> &[JournalEntry<GroupEventRecord>] {
818 &self.events
819 }
820
821 pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
823 &self.backfill_events
824 }
825}
826
827#[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 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 pub const fn scope(&self) -> &ExecutionScope {
945 &self.scope
946 }
947 pub const fn fence(&self) -> ExecutionFence {
949 self.fence
950 }
951 pub fn source_frontier(&self) -> &[MigrationId] {
953 &self.source_frontier
954 }
955 pub fn source_applied(&self) -> &[MigrationId] {
957 &self.source_applied
958 }
959 pub fn target_frontier(&self) -> &[MigrationId] {
961 &self.target_frontier
962 }
963 pub fn migration_ids(&self) -> &[MigrationId] {
965 &self.migration_ids
966 }
967 pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
969 &self.manifest_digests
970 }
971 pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
973 &self.manifest_plan_fingerprints
974 }
975 pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
977 &self.source_declared
978 }
979 pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
981 &self.target_declared
982 }
983 pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
985 &self.source_semantics
986 }
987 pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
989 &self.target_semantics
990 }
991 pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
993 &self.semantic_profile
994 }
995 pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
997 &self.lowering_profile
998 }
999 pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
1001 &self.observed_live_source
1002 }
1003}
1004
1005#[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 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 pub const fn scope(&self) -> &ExecutionScope {
1120 &self.scope
1121 }
1122 pub const fn fence(&self) -> ExecutionFence {
1124 self.fence
1125 }
1126 pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1128 self.manifest_digest
1129 }
1130 pub const fn migration_id(&self) -> &MigrationId {
1132 &self.migration_id
1133 }
1134 pub const fn group_ordinal(&self) -> u32 {
1136 self.group_ordinal
1137 }
1138 pub const fn first_step_index(&self) -> u32 {
1140 self.first_step_index
1141 }
1142 pub const fn schema_delta_step_index(&self) -> u32 {
1144 self.schema_delta_step_index
1145 }
1146 pub const fn end_step_index(&self) -> u32 {
1148 self.end_step_index
1149 }
1150 pub const fn kind(&self) -> GroupJournalEventKind {
1152 self.kind
1153 }
1154 pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
1156 self.observed_target.as_ref()
1157 }
1158}
1159
1160#[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 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 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 pub const fn scope(&self) -> &ExecutionScope {
1355 &self.scope
1356 }
1357
1358 pub const fn fence(&self) -> ExecutionFence {
1360 self.fence
1361 }
1362
1363 pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1365 self.manifest_digest
1366 }
1367
1368 pub const fn migration_id(&self) -> &MigrationId {
1370 &self.migration_id
1371 }
1372
1373 pub const fn operation_ordinal(&self) -> u32 {
1375 self.operation_ordinal
1376 }
1377
1378 pub const fn manifest_step_index(&self) -> u32 {
1380 self.manifest_step_index
1381 }
1382
1383 pub const fn plan_fingerprint(&self) -> &Fingerprint {
1385 &self.plan_fingerprint
1386 }
1387
1388 pub const fn direction(&self) -> BackfillExecutionDirection {
1390 self.direction
1391 }
1392
1393 pub const fn kind(&self) -> GroupJournalEventKind {
1395 self.kind
1396 }
1397
1398 pub const fn completion(&self) -> Option<&BackfillCompletionEvidence> {
1400 self.completion.as_ref()
1401 }
1402}
1403
1404#[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 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 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 pub const fn scope(&self) -> &ExecutionScope {
1459 &self.scope
1460 }
1461 pub const fn fence(&self) -> ExecutionFence {
1463 self.fence
1464 }
1465 pub const fn migration_id(&self) -> &MigrationId {
1467 &self.migration_id
1468 }
1469 pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1471 self.manifest_digest
1472 }
1473 pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1475 &self.source_declared
1476 }
1477 pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1479 &self.target_declared
1480 }
1481 pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1483 &self.source_semantics
1484 }
1485 pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1487 &self.target_semantics
1488 }
1489}
1490
1491#[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 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 pub const fn scope(&self) -> &ExecutionScope {
1599 &self.scope
1600 }
1601 pub const fn fence(&self) -> ExecutionFence {
1603 self.fence
1604 }
1605 pub fn source_applied(&self) -> &[MigrationId] {
1607 &self.source_applied
1608 }
1609 pub fn rollback_ids(&self) -> &[MigrationId] {
1611 &self.rollback_ids
1612 }
1613 pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
1615 &self.manifest_digests
1616 }
1617 pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
1619 &self.manifest_plan_fingerprints
1620 }
1621 pub fn remaining_applied(&self) -> &[MigrationId] {
1623 &self.remaining_applied
1624 }
1625 pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1627 &self.source_declared
1628 }
1629 pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1631 &self.target_declared
1632 }
1633 pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1635 &self.source_semantics
1636 }
1637 pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1639 &self.target_semantics
1640 }
1641 pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
1643 &self.semantic_profile
1644 }
1645 pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
1647 &self.lowering_profile
1648 }
1649 pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
1651 &self.observed_live_source
1652 }
1653}
1654
1655#[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 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 pub const fn scope(&self) -> &ExecutionScope {
1741 &self.scope
1742 }
1743 pub const fn fence(&self) -> ExecutionFence {
1745 self.fence
1746 }
1747 pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1749 self.manifest_digest
1750 }
1751 pub const fn migration_id(&self) -> &MigrationId {
1753 &self.migration_id
1754 }
1755 pub const fn step_ordinal(&self) -> u32 {
1757 self.step_ordinal
1758 }
1759 pub const fn kind(&self) -> GroupJournalEventKind {
1761 self.kind
1762 }
1763 pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
1765 self.observed_target.as_ref()
1766 }
1767}
1768
1769#[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 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 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 pub const fn scope(&self) -> &ExecutionScope {
1828 &self.scope
1829 }
1830 pub const fn fence(&self) -> ExecutionFence {
1832 self.fence
1833 }
1834 pub const fn migration_id(&self) -> &MigrationId {
1836 &self.migration_id
1837 }
1838 pub const fn manifest_digest(&self) -> MigrationManifestDigest {
1840 self.manifest_digest
1841 }
1842 pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1844 &self.source_declared
1845 }
1846 pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
1848 &self.target_declared
1849 }
1850 pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1852 &self.source_semantics
1853 }
1854 pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
1856 &self.target_semantics
1857 }
1858}
1859
1860#[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 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 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 pub const fn plan(&self) -> &JournalEntry<RollbackPlanRecord> {
1956 &self.plan
1957 }
1958
1959 pub fn events(&self) -> &[JournalEntry<RollbackStepEventRecord>] {
1961 &self.events
1962 }
1963
1964 pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
1966 &self.backfill_events
1967 }
1968}
1969
1970pub trait MigrationLeaseStore: Send + Sync {
1972 fn acquire<'a>(
1979 &'a self,
1980 scope: &'a ExecutionScope,
1981 holder: &'a LeaseHolderId,
1982 ) -> ExecutionFuture<'a, MigrationLease>;
1983
1984 fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()>;
1986}
1987
1988pub trait MigrationExecutionJournal: Send + Sync {
1993 fn begin_plan<'a>(
1998 &'a self,
1999 lease: &'a MigrationLease,
2000 record: PlanRecord,
2001 ) -> ExecutionFuture<'a, JournalEntry<PlanRecord>>;
2002
2003 fn record_group_event<'a>(
2005 &'a self,
2006 lease: &'a MigrationLease,
2007 record: GroupEventRecord,
2008 ) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>>;
2009
2010 fn record_backfill_event<'a>(
2012 &'a self,
2013 lease: &'a MigrationLease,
2014 record: BackfillEventRecord,
2015 ) -> ExecutionFuture<'a, JournalEntry<BackfillEventRecord>>;
2016
2017 fn record_applied<'a>(
2026 &'a self,
2027 lease: &'a MigrationLease,
2028 record: AppliedRecord,
2029 ) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>>;
2030
2031 fn load_applied<'a>(
2037 &'a self,
2038 lease: &'a MigrationLease,
2039 ) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>>;
2040
2041 fn load_open_plan<'a>(
2043 &'a self,
2044 lease: &'a MigrationLease,
2045 ) -> ExecutionFuture<'a, Option<OpenPlanRecord>>;
2046
2047 fn begin_rollback_plan<'a>(
2051 &'a self,
2052 lease: &'a MigrationLease,
2053 record: RollbackPlanRecord,
2054 ) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>>;
2055
2056 fn record_rollback_step_event<'a>(
2058 &'a self,
2059 lease: &'a MigrationLease,
2060 record: RollbackStepEventRecord,
2061 ) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>>;
2062
2063 fn record_rolled_back<'a>(
2072 &'a self,
2073 lease: &'a MigrationLease,
2074 record: RolledBackRecord,
2075 ) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>>;
2076
2077 fn load_rolled_back<'a>(
2079 &'a self,
2080 lease: &'a MigrationLease,
2081 ) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>>;
2082
2083 fn load_open_rollback_plan<'a>(
2085 &'a self,
2086 lease: &'a MigrationLease,
2087 ) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>>;
2088}
2089
2090pub 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 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 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 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}