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