use std::fmt;
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Poll, Waker};
use std::time::Instant;
use type_bridge_contract::diagnostic::{Diagnostic, DiagnosticCategory, DiagnosticCode};
use type_bridge_contract::fingerprint::Fingerprint;
use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileFingerprint};
use type_bridge_contract::migration::{
MigrationId, MigrationManifestDigest, MigrationPlanFingerprint,
};
use type_bridge_contract::schema_delta::ManagedSchemaState;
use type_bridge_contract::schema_fingerprint::{
ManagedDeclaredIdentityFingerprint, ManagedSemanticSchemaFingerprint,
};
use type_bridge_contract::schema_lowering::SchemaLoweringProfileFingerprint;
use crate::{
VerifiedMigrationApplyManifest, VerifiedMigrationApplyPlan, VerifiedMigrationRollbackManifest,
VerifiedMigrationRollbackPlan, VerifiedMigrationTransactionGroup,
};
const MAX_LEASE_HOLDER_BYTES: usize = 128;
pub const MAX_MIGRATION_EXECUTION_GROUPS: usize = 65_536;
pub const MAX_MIGRATION_BACKFILL_OBSERVATIONS: usize = 65_536;
#[derive(Clone, Debug, Default)]
pub struct MigrationCancellation {
inner: Arc<MigrationCancellationInner>,
}
#[derive(Debug, Default)]
struct MigrationCancellationInner {
cancelled: AtomicBool,
waiters: std::sync::Mutex<Vec<Waker>>,
}
impl MigrationCancellation {
pub fn cancel(&self) {
if !self.inner.cancelled.swap(true, Ordering::AcqRel) {
let waiters = {
let mut waiters = self.inner.waiters.lock().expect("cancellation waiters");
std::mem::take(&mut *waiters)
};
for waiter in waiters {
waiter.wake();
}
}
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
self.inner.cancelled.load(Ordering::Acquire)
}
pub async fn cancelled(&self) {
std::future::poll_fn(|context| {
if self.is_cancelled() {
return Poll::Ready(());
}
let mut waiters = self.inner.waiters.lock().expect("cancellation waiters");
if self.is_cancelled() {
return Poll::Ready(());
}
if !waiters
.iter()
.any(|waiter| waiter.will_wake(context.waker()))
{
waiters.push(context.waker().clone());
}
Poll::Pending
})
.await
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct MigrationExecutionResourceLimits {
transaction_groups: usize,
backfill_observations: usize,
}
impl MigrationExecutionResourceLimits {
#[must_use]
pub const fn tightened(transaction_groups: usize, backfill_observations: usize) -> Self {
Self {
transaction_groups: if transaction_groups < MAX_MIGRATION_EXECUTION_GROUPS {
transaction_groups
} else {
MAX_MIGRATION_EXECUTION_GROUPS
},
backfill_observations: if backfill_observations < MAX_MIGRATION_BACKFILL_OBSERVATIONS {
backfill_observations
} else {
MAX_MIGRATION_BACKFILL_OBSERVATIONS
},
}
}
pub const fn transaction_groups(self) -> usize {
self.transaction_groups
}
pub const fn backfill_observations(self) -> usize {
self.backfill_observations
}
}
impl Default for MigrationExecutionResourceLimits {
fn default() -> Self {
Self::tightened(
MAX_MIGRATION_EXECUTION_GROUPS,
MAX_MIGRATION_BACKFILL_OBSERVATIONS,
)
}
}
#[derive(Clone, Debug, Default)]
pub struct MigrationExecutionControl {
cancellation: MigrationCancellation,
deadline: Option<Instant>,
resources: MigrationExecutionResourceLimits,
}
impl MigrationExecutionControl {
#[must_use]
pub const fn new(
cancellation: MigrationCancellation,
deadline: Option<Instant>,
resources: MigrationExecutionResourceLimits,
) -> Self {
Self {
cancellation,
deadline,
resources,
}
}
pub const fn cancellation(&self) -> &MigrationCancellation {
&self.cancellation
}
pub const fn deadline(&self) -> Option<Instant> {
self.deadline
}
pub const fn resources(&self) -> MigrationExecutionResourceLimits {
self.resources
}
pub fn check(&self) -> Result<(), Diagnostic> {
if self.cancellation.is_cancelled() {
return Err(failure(
DiagnosticCategory::Cancelled,
"migration_execution_cancelled",
"migration execution was cancelled before the next effect",
));
}
if self
.deadline
.is_some_and(|deadline| Instant::now() >= deadline)
{
return Err(failure(
DiagnosticCategory::ResourceLimit,
"migration_execution_deadline_exceeded",
"migration execution reached its absolute deadline before the next effect",
));
}
Ok(())
}
}
pub type ExecutionFuture<'a, T> = Pin<Box<dyn Future<Output = Result<T, Diagnostic>> + Send + 'a>>;
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct ExecutionFence(u64);
impl ExecutionFence {
pub fn new(value: u64) -> Result<Self, Diagnostic> {
if value == 0 {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_zero_fence",
"migration execution fences must be non-zero",
));
}
Ok(Self(value))
}
pub const fn get(self) -> u64 {
self.0
}
pub fn checked_successor(self) -> Result<Self, Diagnostic> {
let value = self.0.checked_add(1).ok_or_else(|| {
failure(
DiagnosticCategory::ResourceLimit,
"migration_execution_fence_exhausted",
"migration execution fence range is exhausted",
)
})?;
Self::new(value)
}
}
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct ExecutionScope(ManagedScopeId);
impl ExecutionScope {
pub const fn new(scope: ManagedScopeId) -> Self {
Self(scope)
}
pub const fn managed_scope_id(&self) -> &ManagedScopeId {
&self.0
}
}
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct LeaseHolderId(String);
impl LeaseHolderId {
pub fn new(value: impl Into<String>) -> Result<Self, Diagnostic> {
let value = value.into();
if value.is_empty()
|| value.len() > MAX_LEASE_HOLDER_BYTES
|| !value
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'.' | b'_' | b'-'))
{
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_invalid_holder",
"lease holder must be bounded non-empty ASCII [A-Za-z0-9._-]",
));
}
Ok(Self(value))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug)]
pub struct MigrationLease {
scope: ExecutionScope,
holder: LeaseHolderId,
fence: ExecutionFence,
local_binding: Option<ExecutionBindingToken>,
}
#[doc(hidden)]
#[derive(Clone)]
pub struct ExecutionBindingToken(Arc<()>);
impl ExecutionBindingToken {
#[doc(hidden)]
#[must_use]
pub fn fresh() -> Self {
Self(Arc::new(()))
}
fn matches(&self, other: &Self) -> bool {
Arc::ptr_eq(&self.0, &other.0)
}
}
impl fmt::Debug for ExecutionBindingToken {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("ExecutionBindingToken([OPAQUE])")
}
}
impl MigrationLease {
pub const fn new(scope: ExecutionScope, holder: LeaseHolderId, fence: ExecutionFence) -> Self {
Self {
scope,
holder,
fence,
local_binding: None,
}
}
#[doc(hidden)]
#[must_use]
pub fn new_bound(
scope: ExecutionScope,
holder: LeaseHolderId,
fence: ExecutionFence,
local_binding: ExecutionBindingToken,
) -> Self {
Self {
scope,
holder,
fence,
local_binding: Some(local_binding),
}
}
#[doc(hidden)]
#[must_use]
pub fn is_bound_to(&self, expected: &ExecutionBindingToken) -> bool {
self.local_binding
.as_ref()
.is_some_and(|actual| actual.matches(expected))
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn holder(&self) -> &LeaseHolderId {
&self.holder
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
}
impl PartialEq for MigrationLease {
fn eq(&self, other: &Self) -> bool {
self.scope == other.scope
&& self.holder == other.holder
&& self.fence == other.fence
&& match (&self.local_binding, &other.local_binding) {
(None, None) => true,
(Some(left), Some(right)) => left.matches(right),
(None, Some(_)) | (Some(_), None) => false,
}
}
}
impl Eq for MigrationLease {}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum GroupCommitCertainty {
DefinitelyAborted,
Unknown,
}
impl GroupCommitCertainty {
pub const fn journal_event(self) -> GroupJournalEventKind {
match self {
Self::DefinitelyAborted => GroupJournalEventKind::DefinitelyAborted,
Self::Unknown => GroupJournalEventKind::CommitOutcomeUnknown,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum GroupJournalEventKind {
BeforeCommit,
Committed,
CommitOutcomeUnknown,
DefinitelyAborted,
FormalOnlyAdvanced,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum GroupRecoveryObservation {
Unavailable,
ManagedSemantics(ManagedSemanticSchemaFingerprint),
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum GroupRecoveryDecision {
ExecuteNormally,
RepairCheckpoint,
RequiresExplicitRecovery,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum BackfillExecutionDirection {
Forward,
Reverse,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct BackfillExecutionCounts {
matched: u64,
changed: u64,
skipped: u64,
transaction_groups: u32,
}
impl BackfillExecutionCounts {
pub fn new(
matched: u64,
changed: u64,
skipped: u64,
transaction_groups: u32,
) -> Result<Self, Diagnostic> {
if transaction_groups == 0 {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_backfill_zero_transaction_groups",
"terminal backfill evidence must include at least one transaction group",
));
}
if changed.checked_add(skipped) != Some(matched) {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_backfill_count_mismatch",
"backfill matched count must equal changed plus skipped counts",
));
}
Ok(Self {
matched,
changed,
skipped,
transaction_groups,
})
}
pub const fn matched(self) -> u64 {
self.matched
}
pub const fn changed(self) -> u64 {
self.changed
}
pub const fn skipped(self) -> u64 {
self.skipped
}
pub const fn transaction_groups(self) -> u32 {
self.transaction_groups
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct BackfillCompletionEvidence {
plan_fingerprint: Fingerprint,
direction: BackfillExecutionDirection,
counts: BackfillExecutionCounts,
}
impl BackfillCompletionEvidence {
#[must_use]
pub const fn new(
plan_fingerprint: Fingerprint,
direction: BackfillExecutionDirection,
counts: BackfillExecutionCounts,
) -> Self {
Self {
plan_fingerprint,
direction,
counts,
}
}
pub const fn plan_fingerprint(&self) -> &Fingerprint {
&self.plan_fingerprint
}
pub const fn direction(&self) -> BackfillExecutionDirection {
self.direction
}
pub const fn counts(&self) -> BackfillExecutionCounts {
self.counts
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub enum BackfillRecoveryObservation {
Unavailable,
Incomplete,
Complete(BackfillCompletionEvidence),
}
pub fn decide_backfill_recovery(
last_event: Option<GroupJournalEventKind>,
observation: &BackfillRecoveryObservation,
expected_plan: &Fingerprint,
expected_direction: BackfillExecutionDirection,
) -> GroupRecoveryDecision {
let complete = matches!(
observation,
BackfillRecoveryObservation::Complete(evidence)
if evidence.plan_fingerprint() == expected_plan
&& evidence.direction() == expected_direction
);
match last_event {
None if matches!(
observation,
BackfillRecoveryObservation::Incomplete | BackfillRecoveryObservation::Complete(_)
) =>
{
GroupRecoveryDecision::ExecuteNormally
}
Some(GroupJournalEventKind::DefinitelyAborted)
if matches!(observation, BackfillRecoveryObservation::Incomplete) =>
{
GroupRecoveryDecision::ExecuteNormally
}
Some(
GroupJournalEventKind::BeforeCommit
| GroupJournalEventKind::CommitOutcomeUnknown
| GroupJournalEventKind::Committed,
) if complete => GroupRecoveryDecision::RepairCheckpoint,
_ => GroupRecoveryDecision::RequiresExplicitRecovery,
}
}
pub fn decide_group_recovery(
last_event: Option<GroupJournalEventKind>,
observation: &GroupRecoveryObservation,
source: &ManagedSemanticSchemaFingerprint,
target: &ManagedSemanticSchemaFingerprint,
) -> GroupRecoveryDecision {
let distinct = source != target;
let observed = match observation {
GroupRecoveryObservation::Unavailable => ObservedRelation::Unavailable,
GroupRecoveryObservation::ManagedSemantics(value) if value == source && value == target => {
ObservedRelation::Both
}
GroupRecoveryObservation::ManagedSemantics(value) if value == source => {
ObservedRelation::Source
}
GroupRecoveryObservation::ManagedSemantics(value) if value == target => {
ObservedRelation::Target
}
GroupRecoveryObservation::ManagedSemantics(_) => ObservedRelation::Neither,
};
match (last_event, distinct, observed) {
(None, true, ObservedRelation::Source)
| (Some(GroupJournalEventKind::DefinitelyAborted), true, ObservedRelation::Source)
| (None, false, ObservedRelation::Both)
| (Some(GroupJournalEventKind::DefinitelyAborted), false, ObservedRelation::Both) => {
GroupRecoveryDecision::ExecuteNormally
}
(
Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
true,
ObservedRelation::Source,
) => GroupRecoveryDecision::ExecuteNormally,
(
Some(GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown),
true,
ObservedRelation::Target,
)
| (Some(GroupJournalEventKind::Committed), true, ObservedRelation::Target)
| (
Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced),
false,
ObservedRelation::Both,
) => GroupRecoveryDecision::RepairCheckpoint,
_ => GroupRecoveryDecision::RequiresExplicitRecovery,
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ObservedRelation {
Source,
Target,
Both,
Neither,
Unavailable,
}
#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct JournalSequence(u64);
impl JournalSequence {
pub fn new(value: u64) -> Result<Self, Diagnostic> {
if value == 0 {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_zero_sequence",
"journal sequence numbers must be non-zero",
));
}
Ok(Self(value))
}
pub const fn get(self) -> u64 {
self.0
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct JournalEntry<T> {
sequence: JournalSequence,
record: T,
}
impl<T> JournalEntry<T> {
pub const fn from_store(sequence: JournalSequence, record: T) -> Self {
Self { sequence, record }
}
pub const fn sequence(&self) -> JournalSequence {
self.sequence
}
pub const fn record(&self) -> &T {
&self.record
}
pub fn into_record(self) -> T {
self.record
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OpenPlanRecord {
plan: JournalEntry<PlanRecord>,
events: Vec<JournalEntry<GroupEventRecord>>,
backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
}
impl OpenPlanRecord {
pub fn from_store(
plan: JournalEntry<PlanRecord>,
events: Vec<JournalEntry<GroupEventRecord>>,
) -> Result<Self, Diagnostic> {
Self::from_store_with_backfills(plan, events, Vec::new())
}
pub fn from_store_with_backfills(
plan: JournalEntry<PlanRecord>,
events: Vec<JournalEntry<GroupEventRecord>>,
backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
) -> Result<Self, Diagnostic> {
let mut previous = plan.sequence();
let mut previous_fence = plan.record().fence();
let mut ordered: Vec<(
JournalSequence,
ExecutionFence,
&MigrationId,
MigrationManifestDigest,
)> = events
.iter()
.map(|event| {
(
event.sequence(),
event.record().fence(),
event.record().migration_id(),
event.record().manifest_digest(),
)
})
.chain(backfill_events.iter().map(|event| {
(
event.sequence(),
event.record().fence(),
event.record().migration_id(),
event.record().manifest_digest(),
)
}))
.collect();
ordered.sort_by_key(|event| event.0);
for (sequence, fence, migration_id, manifest_digest) in ordered {
let manifest_index = plan
.record()
.manifest_digests()
.iter()
.position(|digest| digest == &manifest_digest);
if sequence <= previous
|| fence < previous_fence
|| manifest_index.is_none_or(|index| {
plan.record().migration_ids().get(index) != Some(migration_id)
})
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_invalid_open_plan",
"loaded open-plan events are not ordered and bound to the plan",
));
}
previous = sequence;
previous_fence = fence;
}
if events
.iter()
.any(|event| event.record().scope() != plan.record().scope())
|| backfill_events
.iter()
.any(|event| event.record().scope() != plan.record().scope())
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_invalid_open_plan",
"loaded open-plan events are not bound to the plan scope",
));
}
Ok(Self {
plan,
events,
backfill_events,
})
}
pub const fn plan(&self) -> &JournalEntry<PlanRecord> {
&self.plan
}
pub fn events(&self) -> &[JournalEntry<GroupEventRecord>] {
&self.events
}
pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
&self.backfill_events
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct PlanRecord {
scope: ExecutionScope,
fence: ExecutionFence,
source_applied: Vec<MigrationId>,
source_frontier: Vec<MigrationId>,
target_frontier: Vec<MigrationId>,
migration_ids: Vec<MigrationId>,
manifest_digests: Vec<MigrationManifestDigest>,
manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
source_declared: ManagedDeclaredIdentityFingerprint,
target_declared: ManagedDeclaredIdentityFingerprint,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
semantic_profile: SemanticProfileFingerprint,
lowering_profile: SchemaLoweringProfileFingerprint,
observed_live_source: ManagedSemanticSchemaFingerprint,
}
impl PlanRecord {
pub fn from_verified_plan(
lease: &MigrationLease,
plan: &VerifiedMigrationApplyPlan,
observed_applied_migrations: &[MigrationId],
observed_live_source: &ManagedSchemaState,
) -> Result<Self, Diagnostic> {
let source = plan.source_state().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_empty_plan",
"an executable migration plan requires a source state",
)
})?;
let target = plan.target_state().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_empty_plan",
"an executable migration plan requires a target state",
)
})?;
let first = plan.migrations().first().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_empty_plan",
"an executable migration plan requires at least one manifest",
)
})?;
if observed_applied_migrations != plan.applied_migrations() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_applied_set",
"applied ledger changed after migration planning; rebuild the plan",
));
}
if observed_live_source != source {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_source_state",
"live managed state differs from the planned source; rebuild the plan",
));
}
let scope = ExecutionScope::new(source.scope().id().clone());
if lease.scope() != &scope || target.scope() != source.scope() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease, source, and target must bind the same managed scope",
));
}
let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
for migration in plan.migrations() {
if migration.manifest().managed_scope().id() != scope.managed_scope_id()
|| migration.manifest().semantic_profile().fingerprint() != &semantic_profile
|| migration.manifest().lowering_profile().fingerprint() != &lowering_profile
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_plan_binding_mismatch",
"planned manifests do not share exact scope and profile bindings",
));
}
}
Ok(Self {
scope,
fence: lease.fence(),
source_applied: plan.applied_migrations().to_vec(),
source_frontier: plan.applied_frontier().to_vec(),
target_frontier: plan.target_frontier().to_vec(),
migration_ids: plan
.migrations()
.iter()
.map(|migration| migration.manifest().id().clone())
.collect(),
manifest_digests: plan
.migrations()
.iter()
.map(VerifiedMigrationApplyManifest::digest)
.collect(),
manifest_plan_fingerprints: plan
.migrations()
.iter()
.map(|migration| migration.manifest().plan_fingerprint().clone())
.collect(),
source_declared: source.managed_declared_identity().clone(),
target_declared: target.managed_declared_identity().clone(),
source_semantics: source.managed_semantic_schema().clone(),
target_semantics: target.managed_semantic_schema().clone(),
semantic_profile,
lowering_profile,
observed_live_source: observed_live_source.managed_semantic_schema().clone(),
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub fn source_frontier(&self) -> &[MigrationId] {
&self.source_frontier
}
pub fn source_applied(&self) -> &[MigrationId] {
&self.source_applied
}
pub fn target_frontier(&self) -> &[MigrationId] {
&self.target_frontier
}
pub fn migration_ids(&self) -> &[MigrationId] {
&self.migration_ids
}
pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
&self.manifest_digests
}
pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
&self.manifest_plan_fingerprints
}
pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.source_declared
}
pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.target_declared
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
&self.semantic_profile
}
pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
&self.lowering_profile
}
pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
&self.observed_live_source
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct GroupEventRecord {
scope: ExecutionScope,
fence: ExecutionFence,
manifest_digest: MigrationManifestDigest,
migration_id: MigrationId,
group_ordinal: u32,
first_step_index: u32,
schema_delta_step_index: u32,
end_step_index: u32,
kind: GroupJournalEventKind,
observed_target: Option<ManagedSemanticSchemaFingerprint>,
}
impl GroupEventRecord {
pub fn new(
lease: &MigrationLease,
migration: &VerifiedMigrationApplyManifest,
group: &VerifiedMigrationTransactionGroup,
kind: GroupJournalEventKind,
observed_target: Option<ManagedSemanticSchemaFingerprint>,
) -> Result<Self, Diagnostic> {
if migration.transaction_groups().get(group.ordinal()) != Some(group) {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_group",
"transaction group does not belong to the supplied verified manifest",
));
}
let step = migration
.steps()
.get(group.schema_delta_step_index())
.ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_group_position_mismatch",
"transaction group delta position is outside the verified manifest",
)
})?;
let delta = step
.step()
.as_schema_delta()
.ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_group_position_mismatch",
"transaction group does not terminate in a schema delta",
)
})?
.delta();
let lowering = step.lowering().ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_group_lowering_missing",
"transaction group delta has no verified lowering",
)
})?;
let scope = ExecutionScope::new(migration.manifest().managed_scope().id().clone());
if lease.scope() != &scope {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease scope differs from the verified migration scope",
));
}
match kind {
GroupJournalEventKind::Committed
if observed_target.as_ref() == Some(delta.target().managed_semantic_schema()) => {}
GroupJournalEventKind::Committed => {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_commit_evidence_mismatch",
"committed event requires the exact observed target semantics",
));
}
GroupJournalEventKind::FormalOnlyAdvanced
if observed_target.is_none()
&& group.assertion_count() == 0
&& lowering.units().is_empty()
&& delta.source().managed_semantic_schema()
== delta.target().managed_semantic_schema() => {}
GroupJournalEventKind::FormalOnlyAdvanced => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_invalid_formal_advance",
"formal-only advancement requires an assertion-free empty equal-semantic group",
));
}
_ if observed_target.is_none() => {}
_ => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_unexpected_observation",
"only committed events may carry observed target semantics",
));
}
}
Ok(Self {
scope,
fence: lease.fence(),
manifest_digest: migration.digest(),
migration_id: migration.manifest().id().clone(),
group_ordinal: position(group.ordinal())?,
first_step_index: position(group.first_step_index())?,
schema_delta_step_index: position(group.schema_delta_step_index())?,
end_step_index: position(group.end_step_index())?,
kind,
observed_target,
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub const fn manifest_digest(&self) -> MigrationManifestDigest {
self.manifest_digest
}
pub const fn migration_id(&self) -> &MigrationId {
&self.migration_id
}
pub const fn group_ordinal(&self) -> u32 {
self.group_ordinal
}
pub const fn first_step_index(&self) -> u32 {
self.first_step_index
}
pub const fn schema_delta_step_index(&self) -> u32 {
self.schema_delta_step_index
}
pub const fn end_step_index(&self) -> u32 {
self.end_step_index
}
pub const fn kind(&self) -> GroupJournalEventKind {
self.kind
}
pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
self.observed_target.as_ref()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct BackfillEventRecord {
scope: ExecutionScope,
fence: ExecutionFence,
manifest_digest: MigrationManifestDigest,
migration_id: MigrationId,
operation_ordinal: u32,
manifest_step_index: u32,
plan_fingerprint: Fingerprint,
direction: BackfillExecutionDirection,
kind: GroupJournalEventKind,
completion: Option<BackfillCompletionEvidence>,
}
impl BackfillEventRecord {
pub fn new_apply(
lease: &MigrationLease,
migration: &VerifiedMigrationApplyManifest,
step_index: usize,
kind: GroupJournalEventKind,
completion: Option<BackfillCompletionEvidence>,
) -> Result<Self, Diagnostic> {
if lease.scope().managed_scope_id() != migration.manifest().managed_scope().id() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease scope differs from the verified migration scope",
));
}
if !migration.backfill_step_indices().contains(&step_index) {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"backfill event position is outside the verified apply manifest",
));
}
let step = migration.steps().get(step_index).ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"backfill event position is outside the verified apply manifest",
)
})?;
let (contract, _) = step.step().as_backfill().ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"verified backfill position does not contain a backfill step",
)
})?;
Self::new_checked(
lease,
migration.digest(),
migration.manifest().id().clone(),
step_index,
step_index,
contract.plan_fingerprint().clone(),
BackfillExecutionDirection::Forward,
kind,
completion,
)
}
pub fn new_rollback(
lease: &MigrationLease,
rollback: &VerifiedMigrationRollbackManifest,
operation_index: usize,
kind: GroupJournalEventKind,
completion: Option<BackfillCompletionEvidence>,
) -> Result<Self, Diagnostic> {
if lease.scope().managed_scope_id() != rollback.manifest().managed_scope().id() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease scope differs from the verified rollback scope",
));
}
let backfill_index = match rollback.operations().get(operation_index) {
Some(crate::VerifiedMigrationRollbackOperation::Backfill(index)) => *index,
_ => {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"backfill event position is outside the verified rollback operations",
));
}
};
let step = rollback
.backfills()
.get(backfill_index)
.ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"rollback backfill index is outside the verified reverse program",
)
})?
.forward_step();
let (contract, _) = step.as_backfill().ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"verified rollback backfill does not contain a backfill step",
)
})?;
let manifest_step_index = rollback
.manifest()
.steps()
.iter()
.position(|candidate| candidate == step)
.ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_step_position",
"rollback backfill is absent from its verified manifest",
)
})?;
Self::new_checked(
lease,
*rollback.digest(),
rollback.manifest().id().clone(),
operation_index,
manifest_step_index,
contract.plan_fingerprint().clone(),
BackfillExecutionDirection::Reverse,
kind,
completion,
)
}
#[allow(clippy::too_many_arguments)]
fn new_checked(
lease: &MigrationLease,
manifest_digest: MigrationManifestDigest,
migration_id: MigrationId,
operation_ordinal: usize,
manifest_step_index: usize,
plan_fingerprint: Fingerprint,
direction: BackfillExecutionDirection,
kind: GroupJournalEventKind,
completion: Option<BackfillCompletionEvidence>,
) -> Result<Self, Diagnostic> {
match kind {
GroupJournalEventKind::Committed
if completion.as_ref().is_some_and(|evidence| {
evidence.plan_fingerprint() == &plan_fingerprint
&& evidence.direction() == direction
}) => {}
GroupJournalEventKind::Committed => {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_backfill_completion_mismatch",
"committed backfill event requires exact terminal plan evidence",
));
}
GroupJournalEventKind::FormalOnlyAdvanced => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_backfill_formal_advance",
"a data backfill cannot advance as a formal-only operation",
));
}
_ if completion.is_none() => {}
_ => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_unexpected_backfill_completion",
"only a committed backfill event may carry terminal evidence",
));
}
}
Ok(Self {
scope: lease.scope().clone(),
fence: lease.fence(),
manifest_digest,
migration_id,
operation_ordinal: position(operation_ordinal)?,
manifest_step_index: position(manifest_step_index)?,
plan_fingerprint,
direction,
kind,
completion,
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub const fn manifest_digest(&self) -> MigrationManifestDigest {
self.manifest_digest
}
pub const fn migration_id(&self) -> &MigrationId {
&self.migration_id
}
pub const fn operation_ordinal(&self) -> u32 {
self.operation_ordinal
}
pub const fn manifest_step_index(&self) -> u32 {
self.manifest_step_index
}
pub const fn plan_fingerprint(&self) -> &Fingerprint {
&self.plan_fingerprint
}
pub const fn direction(&self) -> BackfillExecutionDirection {
self.direction
}
pub const fn kind(&self) -> GroupJournalEventKind {
self.kind
}
pub const fn completion(&self) -> Option<&BackfillCompletionEvidence> {
self.completion.as_ref()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct AppliedRecord {
scope: ExecutionScope,
fence: ExecutionFence,
migration_id: MigrationId,
manifest_digest: MigrationManifestDigest,
source_declared: ManagedDeclaredIdentityFingerprint,
target_declared: ManagedDeclaredIdentityFingerprint,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
}
impl AppliedRecord {
pub fn from_verified_manifest(
lease: &MigrationLease,
migration: &VerifiedMigrationApplyManifest,
) -> Result<Self, Diagnostic> {
Self::from_verified_manifest_contract(lease, migration.manifest())
}
pub fn from_verified_manifest_contract(
lease: &MigrationLease,
manifest: &crate::VerifiedSchemaMigrationManifest,
) -> Result<Self, Diagnostic> {
let source = manifest.source_state();
let target = manifest.target_state();
let scope = ExecutionScope::new(source.scope().id().clone());
if lease.scope() != &scope || target.scope() != source.scope() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease and manifest endpoints must bind the same managed scope",
));
}
Ok(Self {
scope,
fence: lease.fence(),
migration_id: manifest.id().clone(),
manifest_digest: crate::verified_manifest_digest(manifest)?,
source_declared: source.managed_declared_identity().clone(),
target_declared: target.managed_declared_identity().clone(),
source_semantics: source.managed_semantic_schema().clone(),
target_semantics: target.managed_semantic_schema().clone(),
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub const fn migration_id(&self) -> &MigrationId {
&self.migration_id
}
pub const fn manifest_digest(&self) -> MigrationManifestDigest {
self.manifest_digest
}
pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.source_declared
}
pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.target_declared
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RollbackPlanRecord {
scope: ExecutionScope,
fence: ExecutionFence,
source_applied: Vec<MigrationId>,
rollback_ids: Vec<MigrationId>,
manifest_digests: Vec<MigrationManifestDigest>,
manifest_plan_fingerprints: Vec<MigrationPlanFingerprint>,
remaining_applied: Vec<MigrationId>,
source_declared: ManagedDeclaredIdentityFingerprint,
target_declared: ManagedDeclaredIdentityFingerprint,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
semantic_profile: SemanticProfileFingerprint,
lowering_profile: SchemaLoweringProfileFingerprint,
observed_live_source: ManagedSemanticSchemaFingerprint,
}
impl RollbackPlanRecord {
pub fn from_verified_rollback_plan(
lease: &MigrationLease,
plan: &VerifiedMigrationRollbackPlan,
observed_applied_migrations: &[MigrationId],
observed_live_source: &ManagedSchemaState,
) -> Result<Self, Diagnostic> {
let first = plan.rollbacks().first().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_empty_plan",
"an executable rollback plan requires at least one manifest",
)
})?;
let basis: Vec<MigrationId> = plan.applied_basis().into_iter().collect();
if observed_applied_migrations != basis {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_applied_set",
"applied ledger changed after rollback planning; rebuild the plan",
));
}
if observed_live_source != plan.source_state() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_source_state",
"live managed state differs from the planned source; rebuild the plan",
));
}
let scope = ExecutionScope::new(plan.source_state().scope().id().clone());
if lease.scope() != &scope || plan.target_state().scope() != plan.source_state().scope() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease, source, and target must bind the same managed scope",
));
}
let semantic_profile = first.manifest().semantic_profile().fingerprint().clone();
let lowering_profile = first.manifest().lowering_profile().fingerprint().clone();
for rollback in plan.rollbacks() {
if rollback.manifest().managed_scope().id() != scope.managed_scope_id()
|| rollback.manifest().semantic_profile().fingerprint() != &semantic_profile
|| rollback.manifest().lowering_profile().fingerprint() != &lowering_profile
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_plan_binding_mismatch",
"planned rollbacks do not share exact scope and profile bindings",
));
}
}
Ok(Self {
scope,
fence: lease.fence(),
source_applied: basis,
rollback_ids: plan
.rollbacks()
.iter()
.map(|rollback| rollback.manifest().id().clone())
.collect(),
manifest_digests: plan
.rollbacks()
.iter()
.map(|rollback| *rollback.digest())
.collect(),
manifest_plan_fingerprints: plan
.rollbacks()
.iter()
.map(|rollback| rollback.manifest().plan_fingerprint().clone())
.collect(),
remaining_applied: plan.remaining_applied().to_vec(),
source_declared: plan.source_state().managed_declared_identity().clone(),
target_declared: plan.target_state().managed_declared_identity().clone(),
source_semantics: plan.source_state().managed_semantic_schema().clone(),
target_semantics: plan.target_state().managed_semantic_schema().clone(),
semantic_profile,
lowering_profile,
observed_live_source: observed_live_source.managed_semantic_schema().clone(),
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub fn source_applied(&self) -> &[MigrationId] {
&self.source_applied
}
pub fn rollback_ids(&self) -> &[MigrationId] {
&self.rollback_ids
}
pub fn manifest_digests(&self) -> &[MigrationManifestDigest] {
&self.manifest_digests
}
pub fn manifest_plan_fingerprints(&self) -> &[MigrationPlanFingerprint] {
&self.manifest_plan_fingerprints
}
pub fn remaining_applied(&self) -> &[MigrationId] {
&self.remaining_applied
}
pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.source_declared
}
pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.target_declared
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
pub const fn semantic_profile(&self) -> &SemanticProfileFingerprint {
&self.semantic_profile
}
pub const fn lowering_profile(&self) -> &SchemaLoweringProfileFingerprint {
&self.lowering_profile
}
pub const fn observed_live_source(&self) -> &ManagedSemanticSchemaFingerprint {
&self.observed_live_source
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RollbackStepEventRecord {
scope: ExecutionScope,
fence: ExecutionFence,
manifest_digest: MigrationManifestDigest,
migration_id: MigrationId,
step_ordinal: u32,
kind: GroupJournalEventKind,
observed_target: Option<ManagedSemanticSchemaFingerprint>,
}
impl RollbackStepEventRecord {
pub fn new(
lease: &MigrationLease,
rollback: &VerifiedMigrationRollbackManifest,
step_index: usize,
kind: GroupJournalEventKind,
observed_target: Option<ManagedSemanticSchemaFingerprint>,
) -> Result<Self, Diagnostic> {
let step = rollback.steps().get(step_index).ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_rollback_step_position",
"rollback event position is outside the verified rollback manifest",
)
})?;
let reverse = rollback.reverse_delta(step)?;
let scope = ExecutionScope::new(rollback.manifest().managed_scope().id().clone());
if lease.scope() != &scope {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease scope differs from the verified rollback scope",
));
}
match kind {
GroupJournalEventKind::Committed
if observed_target.as_ref() == Some(reverse.target().managed_semantic_schema()) => {
}
GroupJournalEventKind::Committed => {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_commit_evidence_mismatch",
"committed event requires the exact observed target semantics",
));
}
GroupJournalEventKind::FormalOnlyAdvanced
if observed_target.is_none()
&& step.lowering().units().is_empty()
&& reverse.source().managed_semantic_schema()
== reverse.target().managed_semantic_schema() => {}
GroupJournalEventKind::FormalOnlyAdvanced => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_invalid_formal_advance",
"formal-only advancement requires an empty equal-semantic reverse program",
));
}
_ if observed_target.is_none() => {}
_ => {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_unexpected_observation",
"only committed events may carry observed target semantics",
));
}
}
Ok(Self {
scope,
fence: lease.fence(),
manifest_digest: *rollback.digest(),
migration_id: rollback.manifest().id().clone(),
step_ordinal: position(step_index)?,
kind,
observed_target,
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub const fn manifest_digest(&self) -> MigrationManifestDigest {
self.manifest_digest
}
pub const fn migration_id(&self) -> &MigrationId {
&self.migration_id
}
pub const fn step_ordinal(&self) -> u32 {
self.step_ordinal
}
pub const fn kind(&self) -> GroupJournalEventKind {
self.kind
}
pub const fn observed_target(&self) -> Option<&ManagedSemanticSchemaFingerprint> {
self.observed_target.as_ref()
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct RolledBackRecord {
scope: ExecutionScope,
fence: ExecutionFence,
migration_id: MigrationId,
manifest_digest: MigrationManifestDigest,
source_declared: ManagedDeclaredIdentityFingerprint,
target_declared: ManagedDeclaredIdentityFingerprint,
source_semantics: ManagedSemanticSchemaFingerprint,
target_semantics: ManagedSemanticSchemaFingerprint,
}
impl RolledBackRecord {
pub fn from_verified_rollback(
lease: &MigrationLease,
rollback: &VerifiedMigrationRollbackManifest,
) -> Result<Self, Diagnostic> {
Self::from_verified_manifest_contract(lease, rollback.manifest())
}
pub fn from_verified_manifest_contract(
lease: &MigrationLease,
manifest: &crate::VerifiedSchemaMigrationManifest,
) -> Result<Self, Diagnostic> {
let source = manifest.target_state();
let target = manifest.source_state();
let scope = ExecutionScope::new(source.scope().id().clone());
if lease.scope() != &scope || target.scope() != source.scope() {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_scope_mismatch",
"lease and manifest endpoints must bind the same managed scope",
));
}
Ok(Self {
scope,
fence: lease.fence(),
migration_id: manifest.id().clone(),
manifest_digest: crate::verified_manifest_digest(manifest)?,
source_declared: source.managed_declared_identity().clone(),
target_declared: target.managed_declared_identity().clone(),
source_semantics: source.managed_semantic_schema().clone(),
target_semantics: target.managed_semantic_schema().clone(),
})
}
pub const fn scope(&self) -> &ExecutionScope {
&self.scope
}
pub const fn fence(&self) -> ExecutionFence {
self.fence
}
pub const fn migration_id(&self) -> &MigrationId {
&self.migration_id
}
pub const fn manifest_digest(&self) -> MigrationManifestDigest {
self.manifest_digest
}
pub const fn source_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.source_declared
}
pub const fn target_declared(&self) -> &ManagedDeclaredIdentityFingerprint {
&self.target_declared
}
pub const fn source_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.source_semantics
}
pub const fn target_semantics(&self) -> &ManagedSemanticSchemaFingerprint {
&self.target_semantics
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct OpenRollbackPlanRecord {
plan: JournalEntry<RollbackPlanRecord>,
events: Vec<JournalEntry<RollbackStepEventRecord>>,
backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
}
impl OpenRollbackPlanRecord {
pub fn from_store(
plan: JournalEntry<RollbackPlanRecord>,
events: Vec<JournalEntry<RollbackStepEventRecord>>,
) -> Result<Self, Diagnostic> {
Self::from_store_with_backfills(plan, events, Vec::new())
}
pub fn from_store_with_backfills(
plan: JournalEntry<RollbackPlanRecord>,
events: Vec<JournalEntry<RollbackStepEventRecord>>,
backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
) -> Result<Self, Diagnostic> {
let mut previous = plan.sequence();
let mut previous_fence = plan.record().fence();
let mut ordered: Vec<(
JournalSequence,
ExecutionFence,
&MigrationId,
MigrationManifestDigest,
)> = events
.iter()
.map(|event| {
(
event.sequence(),
event.record().fence(),
event.record().migration_id(),
event.record().manifest_digest(),
)
})
.chain(backfill_events.iter().map(|event| {
(
event.sequence(),
event.record().fence(),
event.record().migration_id(),
event.record().manifest_digest(),
)
}))
.collect();
ordered.sort_by_key(|event| event.0);
for (sequence, fence, migration_id, manifest_digest) in ordered {
let manifest_index = plan
.record()
.manifest_digests()
.iter()
.position(|digest| digest == &manifest_digest);
if sequence <= previous
|| fence < previous_fence
|| manifest_index.is_none_or(|index| {
plan.record().rollback_ids().get(index) != Some(migration_id)
})
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_invalid_open_plan",
"loaded open-rollback events are not ordered and bound to the plan",
));
}
previous = sequence;
previous_fence = fence;
}
if events
.iter()
.any(|event| event.record().scope() != plan.record().scope())
|| backfill_events
.iter()
.any(|event| event.record().scope() != plan.record().scope())
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_invalid_open_plan",
"loaded open-rollback events are not bound to the plan scope",
));
}
Ok(Self {
plan,
events,
backfill_events,
})
}
pub const fn plan(&self) -> &JournalEntry<RollbackPlanRecord> {
&self.plan
}
pub fn events(&self) -> &[JournalEntry<RollbackStepEventRecord>] {
&self.events
}
pub fn backfill_events(&self) -> &[JournalEntry<BackfillEventRecord>] {
&self.backfill_events
}
}
pub trait MigrationLeaseStore: Send + Sync {
fn acquire<'a>(
&'a self,
scope: &'a ExecutionScope,
holder: &'a LeaseHolderId,
) -> ExecutionFuture<'a, MigrationLease>;
fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()>;
}
pub trait MigrationExecutionJournal: Send + Sync {
fn begin_plan<'a>(
&'a self,
lease: &'a MigrationLease,
record: PlanRecord,
) -> ExecutionFuture<'a, JournalEntry<PlanRecord>>;
fn record_group_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: GroupEventRecord,
) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>>;
fn record_backfill_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: BackfillEventRecord,
) -> ExecutionFuture<'a, JournalEntry<BackfillEventRecord>>;
fn record_applied<'a>(
&'a self,
lease: &'a MigrationLease,
record: AppliedRecord,
) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>>;
fn load_applied<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>>;
fn load_open_plan<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Option<OpenPlanRecord>>;
fn begin_rollback_plan<'a>(
&'a self,
lease: &'a MigrationLease,
record: RollbackPlanRecord,
) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>>;
fn record_rollback_step_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: RollbackStepEventRecord,
) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>>;
fn record_rolled_back<'a>(
&'a self,
lease: &'a MigrationLease,
record: RolledBackRecord,
) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>>;
fn load_rolled_back<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>>;
fn load_open_rollback_plan<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>>;
}
pub fn active_applied_entries(
applied: Vec<JournalEntry<AppliedRecord>>,
rolled_back: &[JournalEntry<RolledBackRecord>],
) -> Result<Vec<JournalEntry<AppliedRecord>>, Diagnostic> {
let mut retired = vec![false; applied.len()];
for retirement in rolled_back {
let matched = applied
.iter()
.enumerate()
.filter(|(index, entry)| {
!retired[*index]
&& entry.sequence() < retirement.sequence()
&& entry.record().migration_id() == retirement.record().migration_id()
&& entry.record().manifest_digest() == retirement.record().manifest_digest()
})
.max_by_key(|(_, entry)| entry.sequence())
.map(|(index, _)| index);
let Some(index) = matched else {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_unmatched_retirement",
"retirement record matches no active applied record before it",
));
};
retired[index] = true;
}
Ok(applied
.into_iter()
.zip(retired)
.filter_map(|(entry, retired)| (!retired).then_some(entry))
.collect())
}
fn position(value: usize) -> Result<u32, Diagnostic> {
u32::try_from(value).map_err(|_| {
failure(
DiagnosticCategory::ResourceLimit,
"migration_execution_position_limit",
"migration transaction position exceeds the canonical u32 range",
)
})
}
fn failure(category: DiagnosticCategory, code: &'static str, message: &'static str) -> Diagnostic {
Diagnostic::new(
category,
DiagnosticCode::new(code).expect("static migration execution diagnostic code"),
message,
)
}
#[cfg(test)]
mod tests {
use std::collections::BTreeMap;
use std::future::Future;
use std::sync::Mutex;
use std::task::{Context, Poll, Waker};
use type_bridge_contract::fingerprint::SemanticProfileId;
use type_bridge_contract::managed_scope::{ManagedScopeId, SemanticProfileBinding};
use type_bridge_contract::migration::{MigrationAppLabel, MigrationName};
use super::*;
#[test]
fn migration_controls_cancel_and_only_tighten_shared_ceilings() {
let cancellation = MigrationCancellation::default();
let limits = MigrationExecutionResourceLimits::tightened(7, usize::MAX);
let control = MigrationExecutionControl::new(cancellation.clone(), None, limits);
assert_eq!(control.resources().transaction_groups(), 7);
assert_eq!(
control.resources().backfill_observations(),
MAX_MIGRATION_BACKFILL_OBSERVATIONS
);
assert!(control.check().is_ok());
cancellation.cancel();
cancellation.cancel();
let diagnostic = control.check().expect_err("cancelled control");
assert_eq!(diagnostic.category(), DiagnosticCategory::Cancelled);
assert_eq!(diagnostic.code().as_str(), "migration_execution_cancelled");
}
#[test]
fn migration_control_uses_one_absolute_deadline() {
let control = MigrationExecutionControl::new(
MigrationCancellation::default(),
Some(Instant::now()),
MigrationExecutionResourceLimits::default(),
);
let diagnostic = control.check().expect_err("expired deadline");
assert_eq!(diagnostic.category(), DiagnosticCategory::ResourceLimit);
assert_eq!(
diagnostic.code().as_str(),
"migration_execution_deadline_exceeded"
);
}
use crate::schema_lowering_profile_binding;
#[derive(Default)]
struct ScopeState {
highest_fence: u64,
active_lease: Option<MigrationLease>,
next_sequence: u64,
applied: Vec<JournalEntry<AppliedRecord>>,
rolled_back: Vec<JournalEntry<RolledBackRecord>>,
open_plan: Option<JournalEntry<PlanRecord>>,
events: Vec<JournalEntry<GroupEventRecord>>,
backfill_events: Vec<JournalEntry<BackfillEventRecord>>,
open_rollback_plan: Option<JournalEntry<RollbackPlanRecord>>,
rollback_events: Vec<JournalEntry<RollbackStepEventRecord>>,
}
#[derive(Default)]
struct InMemoryStore {
scopes: Mutex<BTreeMap<ExecutionScope, ScopeState>>,
}
impl InMemoryStore {
fn check_lease<'a>(
state: &'a mut ScopeState,
lease: &MigrationLease,
) -> Result<&'a mut ScopeState, Diagnostic> {
if state.active_lease.as_ref() != Some(lease) {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_fence",
"journal write does not carry the current active lease and fence",
));
}
Ok(state)
}
fn sequence(state: &mut ScopeState) -> Result<JournalSequence, Diagnostic> {
state.next_sequence = state.next_sequence.checked_add(1).ok_or_else(|| {
failure(
DiagnosticCategory::ResourceLimit,
"migration_execution_sequence_exhausted",
"journal sequence range is exhausted",
)
})?;
JournalSequence::new(state.next_sequence)
}
}
impl MigrationLeaseStore for InMemoryStore {
fn acquire<'a>(
&'a self,
scope: &'a ExecutionScope,
holder: &'a LeaseHolderId,
) -> ExecutionFuture<'a, MigrationLease> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.entry(scope.clone()).or_default();
if state.active_lease.is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_lease_contended",
"migration scope already has an active lease",
));
}
let fence = if state.highest_fence == 0 {
ExecutionFence::new(1)?
} else {
ExecutionFence::new(state.highest_fence)?.checked_successor()?
};
state.highest_fence = fence.get();
let lease = MigrationLease::new(scope.clone(), holder.clone(), fence);
state.active_lease = Some(lease.clone());
Ok(lease)
})
}
fn release<'a>(&'a self, lease: &'a MigrationLease) -> ExecutionFuture<'a, ()> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(|| {
failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_fence",
"lease scope is not active",
)
})?;
Self::check_lease(state, lease)?;
state.active_lease = None;
Ok(())
})
}
}
impl MigrationExecutionJournal for InMemoryStore {
fn begin_plan<'a>(
&'a self,
lease: &'a MigrationLease,
record: PlanRecord,
) -> ExecutionFuture<'a, JournalEntry<PlanRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_plan_already_open",
"migration scope already has an open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.open_plan = Some(entry.clone());
Ok(entry)
})
}
fn record_group_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: GroupEventRecord,
) -> ExecutionFuture<'a, JournalEntry<GroupEventRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
let plan = state.open_plan.as_ref().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_no_open_plan",
"group event requires an open migration plan",
)
})?;
if !plan
.record()
.manifest_digests()
.contains(&record.manifest_digest())
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_event",
"group event manifest is absent from the open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.events.push(entry.clone());
Ok(entry)
})
}
fn record_backfill_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: BackfillEventRecord,
) -> ExecutionFuture<'a, JournalEntry<BackfillEventRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
let member = state.open_plan.as_ref().is_some_and(|plan| {
plan.record()
.manifest_digests()
.contains(&record.manifest_digest())
}) || state.open_rollback_plan.as_ref().is_some_and(|plan| {
plan.record()
.manifest_digests()
.contains(&record.manifest_digest())
});
if !member {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_event",
"backfill event manifest is absent from the open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.backfill_events.push(entry.clone());
Ok(entry)
})
}
fn record_applied<'a>(
&'a self,
lease: &'a MigrationLease,
record: AppliedRecord,
) -> ExecutionFuture<'a, JournalEntry<AppliedRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
if let Some(existing) = active
.iter()
.find(|entry| entry.record().migration_id() == record.migration_id())
{
return if existing.record() == &record {
Ok(existing.clone())
} else {
Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_applied_identity_conflict",
"applied migration identity has different evidence",
))
};
}
let plan = state.open_plan.as_ref().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_no_open_plan",
"applied migration requires an open migration plan",
)
})?;
let manifest_index = plan
.record()
.migration_ids()
.iter()
.position(|id| id == record.migration_id());
if manifest_index.is_none_or(|index| {
plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
}) {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_applied_record",
"applied migration identity and digest are absent from the open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.applied.push(entry.clone());
let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
let complete = state.open_plan.as_ref().is_some_and(|plan| {
plan.record().migration_ids().iter().all(|id| {
active
.iter()
.any(|applied| applied.record().migration_id() == id)
})
});
if complete {
state.open_plan = None;
state.events.clear();
state.backfill_events.clear();
}
Ok(entry)
})
}
fn load_applied<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Vec<JournalEntry<AppliedRecord>>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
active_applied_entries(state.applied.clone(), &state.rolled_back)
})
}
fn load_open_plan<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Option<OpenPlanRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
let Some(plan) = state.open_plan.clone() else {
return Ok(None);
};
Ok(Some(OpenPlanRecord::from_store_with_backfills(
plan,
state.events.clone(),
state
.backfill_events
.iter()
.filter(|event| {
event.record().direction() == BackfillExecutionDirection::Forward
})
.cloned()
.collect(),
)?))
})
}
fn begin_rollback_plan<'a>(
&'a self,
lease: &'a MigrationLease,
record: RollbackPlanRecord,
) -> ExecutionFuture<'a, JournalEntry<RollbackPlanRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
if state.open_plan.is_some() || state.open_rollback_plan.is_some() {
return Err(failure(
DiagnosticCategory::InvalidContract,
"migration_execution_plan_already_open",
"migration scope already has an open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.open_rollback_plan = Some(entry.clone());
Ok(entry)
})
}
fn record_rollback_step_event<'a>(
&'a self,
lease: &'a MigrationLease,
record: RollbackStepEventRecord,
) -> ExecutionFuture<'a, JournalEntry<RollbackStepEventRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_no_open_plan",
"rollback event requires an open rollback plan",
)
})?;
if !plan
.record()
.manifest_digests()
.contains(&record.manifest_digest())
{
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_event",
"rollback event manifest is absent from the open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.rollback_events.push(entry.clone());
Ok(entry)
})
}
fn record_rolled_back<'a>(
&'a self,
lease: &'a MigrationLease,
record: RolledBackRecord,
) -> ExecutionFuture<'a, JournalEntry<RolledBackRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
if record.scope() != lease.scope() || record.fence() != lease.fence() {
return Err(stale_fence());
}
let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
let is_active = active.iter().any(|entry| {
entry.record().migration_id() == record.migration_id()
&& entry.record().manifest_digest() == record.manifest_digest()
});
if !is_active {
if let Some(existing) = state.rolled_back.iter().find(|entry| {
entry.record() == &record && entry.record().fence() == lease.fence()
}) {
return Ok(existing.clone());
}
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_retirement_conflict",
"retirement target is not active in the applied ledger",
));
}
let plan = state.open_rollback_plan.as_ref().ok_or_else(|| {
failure(
DiagnosticCategory::InvalidContract,
"migration_execution_no_open_plan",
"retirement requires an open rollback plan",
)
})?;
let manifest_index = plan
.record()
.rollback_ids()
.iter()
.position(|id| id == record.migration_id());
if manifest_index.is_none_or(|index| {
plan.record().manifest_digests().get(index) != Some(&record.manifest_digest())
}) {
return Err(failure(
DiagnosticCategory::Integrity,
"migration_execution_foreign_applied_record",
"retirement identity and digest are absent from the open plan",
));
}
let entry = JournalEntry::from_store(Self::sequence(state)?, record);
state.rolled_back.push(entry.clone());
let active = active_applied_entries(state.applied.clone(), &state.rolled_back)?;
let complete = state.open_rollback_plan.as_ref().is_some_and(|plan| {
plan.record().rollback_ids().iter().all(|id| {
!active
.iter()
.any(|applied| applied.record().migration_id() == id)
})
});
if complete {
state.open_rollback_plan = None;
state.rollback_events.clear();
state.backfill_events.clear();
}
Ok(entry)
})
}
fn load_rolled_back<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Vec<JournalEntry<RolledBackRecord>>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
Ok(state.rolled_back.clone())
})
}
fn load_open_rollback_plan<'a>(
&'a self,
lease: &'a MigrationLease,
) -> ExecutionFuture<'a, Option<OpenRollbackPlanRecord>> {
Box::pin(async move {
let mut scopes = self.scopes.lock().expect("store mutex");
let state = scopes.get_mut(lease.scope()).ok_or_else(stale_fence)?;
Self::check_lease(state, lease)?;
let Some(plan) = state.open_rollback_plan.clone() else {
return Ok(None);
};
Ok(Some(OpenRollbackPlanRecord::from_store_with_backfills(
plan,
state.rollback_events.clone(),
state
.backfill_events
.iter()
.filter(|event| {
event.record().direction() == BackfillExecutionDirection::Reverse
})
.cloned()
.collect(),
)?))
})
}
}
fn stale_fence() -> Diagnostic {
failure(
DiagnosticCategory::Integrity,
"migration_execution_stale_fence",
"journal write does not carry the current active lease and fence",
)
}
fn scope() -> ExecutionScope {
ExecutionScope::new(ManagedScopeId::new("journal-test").expect("scope"))
}
fn migration_id() -> MigrationId {
MigrationId::from_components(
MigrationAppLabel::new("example").expect("app"),
MigrationName::new("0001_initial").expect("name"),
)
}
fn semantic_fingerprint(bytes: &[u8]) -> ManagedSemanticSchemaFingerprint {
ManagedSemanticSchemaFingerprint::compute(
SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
bytes,
)
.expect("semantic fingerprint")
}
fn fake_plan(lease: &MigrationLease) -> PlanRecord {
let id = migration_id();
let semantic = SemanticProfileBinding::resolve(
SemanticProfileId::new("typedb-3.12.1/v1").expect("profile"),
)
.expect("semantic binding");
PlanRecord {
scope: lease.scope().clone(),
fence: lease.fence(),
source_applied: Vec::new(),
source_frontier: Vec::new(),
target_frontier: vec![id.clone()],
migration_ids: vec![id],
manifest_digests: vec![MigrationManifestDigest::compute(b"manifest")],
manifest_plan_fingerprints: vec![
MigrationPlanFingerprint::compute(&[]).expect("plan fingerprint"),
],
source_declared: ManagedDeclaredIdentityFingerprint::compute(b"source")
.expect("source declared"),
target_declared: ManagedDeclaredIdentityFingerprint::compute(b"target")
.expect("target declared"),
source_semantics: semantic_fingerprint(b"source"),
target_semantics: semantic_fingerprint(b"target"),
semantic_profile: semantic.fingerprint().clone(),
lowering_profile: schema_lowering_profile_binding()
.expect("lowering binding")
.fingerprint()
.clone(),
observed_live_source: semantic_fingerprint(b"source"),
}
}
fn fake_event(lease: &MigrationLease, plan: &PlanRecord) -> GroupEventRecord {
GroupEventRecord {
scope: lease.scope().clone(),
fence: lease.fence(),
manifest_digest: plan.manifest_digests()[0],
migration_id: plan.migration_ids()[0].clone(),
group_ordinal: 0,
first_step_index: 0,
schema_delta_step_index: 0,
end_step_index: 1,
kind: GroupJournalEventKind::BeforeCommit,
observed_target: None,
}
}
fn fake_applied(lease: &MigrationLease, plan: &PlanRecord) -> AppliedRecord {
AppliedRecord {
scope: lease.scope().clone(),
fence: lease.fence(),
migration_id: plan.migration_ids()[0].clone(),
manifest_digest: plan.manifest_digests()[0],
source_declared: plan.source_declared().clone(),
target_declared: plan.target_declared().clone(),
source_semantics: plan.source_semantics().clone(),
target_semantics: plan.target_semantics().clone(),
}
}
#[test]
fn bound_lease_equality_includes_process_local_binding_identity() {
let scope = scope();
let holder = LeaseHolderId::new("owner").expect("holder");
let fence = ExecutionFence::new(1).expect("fence");
let token = ExecutionBindingToken::fresh();
let same = MigrationLease::new_bound(scope.clone(), holder.clone(), fence, token.clone());
let clone = same.clone();
let foreign = MigrationLease::new_bound(
scope.clone(),
holder.clone(),
fence,
ExecutionBindingToken::fresh(),
);
let unbound = MigrationLease::new(scope.clone(), holder.clone(), fence);
let same_unbound = MigrationLease::new(scope, holder, fence);
assert_eq!(same, clone);
assert_ne!(same, foreign);
assert_ne!(same, unbound);
assert_eq!(unbound, same_unbound);
}
#[test]
fn store_assigns_monotonic_sequences_and_rejects_every_stale_fence_write() {
let store = InMemoryStore::default();
let scope = scope();
let holder_a = LeaseHolderId::new("owner-a").expect("holder");
let holder_b = LeaseHolderId::new("owner-b").expect("holder");
let lease_a = block_on(store.acquire(&scope, &holder_a)).expect("lease a");
assert!(block_on(store.acquire(&scope, &holder_b)).is_err());
let plan_a = fake_plan(&lease_a);
let plan_entry = block_on(store.begin_plan(&lease_a, plan_a.clone())).expect("plan");
let event_entry =
block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a)))
.expect("event");
assert_eq!(plan_entry.sequence().get(), 1);
assert_eq!(event_entry.sequence().get(), 2);
block_on(store.release(&lease_a)).expect("release a");
let lease_b = block_on(store.acquire(&scope, &holder_b)).expect("lease b");
assert!(lease_b.fence() > lease_a.fence());
assert!(block_on(store.begin_plan(&lease_a, plan_a.clone())).is_err());
assert!(
block_on(store.record_group_event(&lease_a, fake_event(&lease_a, &plan_a),)).is_err()
);
assert!(
block_on(store.record_applied(&lease_a, fake_applied(&lease_a, &plan_a),)).is_err()
);
assert!(block_on(store.release(&lease_a)).is_err());
let recovery_event =
block_on(store.record_group_event(&lease_b, fake_event(&lease_b, &plan_a)))
.expect("recovery event");
assert_eq!(recovery_event.sequence().get(), 3);
block_on(store.release(&lease_b)).expect("release recovery lease");
let holder_c = LeaseHolderId::new("owner-c").expect("holder");
let lease_c = block_on(store.acquire(&scope, &holder_c)).expect("lease c");
assert!(lease_c.fence() > lease_b.fence());
let open = block_on(store.load_open_plan(&lease_c))
.expect("open plan")
.expect("plan remains open after recovery crash");
assert_eq!(open.events().len(), 2);
assert_eq!(open.events()[0].record().fence(), lease_a.fence());
assert_eq!(open.events()[1].record().fence(), lease_b.fence());
}
#[test]
fn applied_records_are_plan_bound_idempotent_and_close_the_completed_plan() {
let store = InMemoryStore::default();
let scope = scope();
let holder = LeaseHolderId::new("applied-owner").expect("holder");
let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
let plan = fake_plan(&lease);
let applied = fake_applied(&lease, &plan);
assert!(block_on(store.record_applied(&lease, applied.clone())).is_err());
block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
let mut foreign = applied.clone();
foreign.manifest_digest = MigrationManifestDigest::compute(b"foreign");
assert!(block_on(store.record_applied(&lease, foreign)).is_err());
let first = block_on(store.record_applied(&lease, applied.clone()))
.expect("apply planned manifest");
let duplicate = block_on(store.record_applied(&lease, applied))
.expect("same-fence duplicate is idempotent");
assert_eq!(duplicate, first);
assert!(
block_on(store.load_open_plan(&lease))
.expect("load completed plan")
.is_none()
);
assert_eq!(
block_on(store.load_applied(&lease)).expect("load applied ledger"),
vec![first],
);
}
fn fake_rollback_plan(lease: &MigrationLease, plan: &PlanRecord) -> RollbackPlanRecord {
RollbackPlanRecord {
scope: lease.scope().clone(),
fence: lease.fence(),
source_applied: plan.migration_ids().to_vec(),
rollback_ids: plan.migration_ids().to_vec(),
manifest_digests: plan.manifest_digests().to_vec(),
manifest_plan_fingerprints: plan.manifest_plan_fingerprints().to_vec(),
remaining_applied: Vec::new(),
source_declared: plan.target_declared().clone(),
target_declared: plan.source_declared().clone(),
source_semantics: plan.target_semantics().clone(),
target_semantics: plan.source_semantics().clone(),
semantic_profile: plan.semantic_profile().clone(),
lowering_profile: plan.lowering_profile().clone(),
observed_live_source: plan.target_semantics().clone(),
}
}
fn fake_rolled_back(lease: &MigrationLease, plan: &PlanRecord) -> RolledBackRecord {
RolledBackRecord {
scope: lease.scope().clone(),
fence: lease.fence(),
migration_id: plan.migration_ids()[0].clone(),
manifest_digest: plan.manifest_digests()[0],
source_declared: plan.target_declared().clone(),
target_declared: plan.source_declared().clone(),
source_semantics: plan.target_semantics().clone(),
target_semantics: plan.source_semantics().clone(),
}
}
#[test]
fn retirement_is_append_only_exclusive_and_reopens_the_identity() {
let store = InMemoryStore::default();
let scope = scope();
let holder = LeaseHolderId::new("retirement-owner").expect("holder");
let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
let plan = fake_plan(&lease);
let rollback_plan = fake_rollback_plan(&lease, &plan);
let retirement = fake_rolled_back(&lease, &plan);
block_on(store.begin_plan(&lease, plan.clone())).expect("open plan");
assert!(
block_on(store.begin_rollback_plan(&lease, rollback_plan.clone())).is_err(),
"open plans of either direction must be exclusive"
);
block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
.expect("apply planned manifest");
assert_eq!(
block_on(store.load_applied(&lease)).expect("active").len(),
1
);
assert!(block_on(store.record_rolled_back(&lease, retirement.clone())).is_err());
block_on(store.begin_rollback_plan(&lease, rollback_plan)).expect("open rollback plan");
assert!(
block_on(store.begin_plan(&lease, plan.clone())).is_err(),
"an open rollback plan must block a new apply plan"
);
let first = block_on(store.record_rolled_back(&lease, retirement.clone()))
.expect("retire applied migration");
let duplicate = block_on(store.record_rolled_back(&lease, retirement))
.expect("same-fence duplicate retirement is idempotent");
assert_eq!(duplicate, first);
assert!(
block_on(store.load_open_rollback_plan(&lease))
.expect("closed rollback plan")
.is_none()
);
assert!(
block_on(store.load_applied(&lease))
.expect("active")
.is_empty()
);
assert_eq!(
block_on(store.load_rolled_back(&lease))
.expect("retired")
.len(),
1,
);
block_on(store.begin_plan(&lease, plan.clone())).expect("reopen plan");
block_on(store.record_applied(&lease, fake_applied(&lease, &plan)))
.expect("re-apply retired migration");
assert_eq!(
block_on(store.load_applied(&lease)).expect("active").len(),
1
);
}
#[test]
fn unmatched_retirements_are_corrupt_history() {
let store = InMemoryStore::default();
let scope = scope();
let holder = LeaseHolderId::new("orphan-owner").expect("holder");
let lease = block_on(store.acquire(&scope, &holder)).expect("lease");
let plan = fake_plan(&lease);
let orphan = JournalEntry::from_store(
JournalSequence::new(7).expect("sequence"),
fake_rolled_back(&lease, &plan),
);
let error = active_applied_entries(Vec::new(), &[orphan])
.expect_err("a retirement without its applied record is corrupt");
assert_eq!(
error.code().as_str(),
"migration_execution_unmatched_retirement"
);
}
#[test]
fn recovery_decision_table_is_exhaustive_for_distinct_and_equal_semantics() {
let source = semantic_fingerprint(b"source");
let target = semantic_fingerprint(b"target");
let neither = semantic_fingerprint(b"neither");
let events = [
None,
Some(GroupJournalEventKind::BeforeCommit),
Some(GroupJournalEventKind::Committed),
Some(GroupJournalEventKind::CommitOutcomeUnknown),
Some(GroupJournalEventKind::DefinitelyAborted),
Some(GroupJournalEventKind::FormalOnlyAdvanced),
];
let observations = [
GroupRecoveryObservation::ManagedSemantics(source.clone()),
GroupRecoveryObservation::ManagedSemantics(target.clone()),
GroupRecoveryObservation::ManagedSemantics(neither.clone()),
GroupRecoveryObservation::Unavailable,
];
for event in events {
for observation in &observations {
let expected = expected_distinct(event, observation, &source, &target);
assert_eq!(
decide_group_recovery(event, observation, &source, &target),
expected,
"distinct case {event:?} {observation:?}",
);
}
}
let equal_observations = [
GroupRecoveryObservation::ManagedSemantics(source.clone()),
GroupRecoveryObservation::ManagedSemantics(neither),
GroupRecoveryObservation::Unavailable,
];
for event in events {
for observation in &equal_observations {
let expected = expected_equal(event, observation, &source);
assert_eq!(
decide_group_recovery(event, observation, &source, &source),
expected,
"equal case {event:?} {observation:?}",
);
}
}
}
#[test]
fn backfill_counts_and_recovery_fail_closed_without_exact_completion() {
assert_eq!(
BackfillExecutionCounts::new(5, 3, 1, 1)
.expect_err("inconsistent counts reject")
.code()
.as_str(),
"migration_backfill_count_mismatch"
);
assert_eq!(
BackfillExecutionCounts::new(0, 0, 0, 0)
.expect_err("zero transaction groups reject")
.code()
.as_str(),
"migration_backfill_zero_transaction_groups"
);
let plan = semantic_fingerprint(b"backfill-plan")
.as_fingerprint()
.clone();
let foreign = semantic_fingerprint(b"foreign-backfill-plan")
.as_fingerprint()
.clone();
let counts = BackfillExecutionCounts::new(5, 3, 2, 2).expect("valid counts");
let complete = BackfillRecoveryObservation::Complete(BackfillCompletionEvidence::new(
plan.clone(),
BackfillExecutionDirection::Forward,
counts,
));
assert_eq!(
decide_backfill_recovery(
None,
&BackfillRecoveryObservation::Incomplete,
&plan,
BackfillExecutionDirection::Forward,
),
GroupRecoveryDecision::ExecuteNormally
);
assert_eq!(
decide_backfill_recovery(None, &complete, &plan, BackfillExecutionDirection::Forward,),
GroupRecoveryDecision::ExecuteNormally
);
assert_eq!(
decide_backfill_recovery(
Some(GroupJournalEventKind::BeforeCommit),
&BackfillRecoveryObservation::Incomplete,
&plan,
BackfillExecutionDirection::Forward,
),
GroupRecoveryDecision::RequiresExplicitRecovery
);
assert_eq!(
decide_backfill_recovery(
Some(GroupJournalEventKind::CommitOutcomeUnknown),
&complete,
&plan,
BackfillExecutionDirection::Forward,
),
GroupRecoveryDecision::RepairCheckpoint
);
assert_eq!(
decide_backfill_recovery(
Some(GroupJournalEventKind::Committed),
&complete,
&foreign,
BackfillExecutionDirection::Forward,
),
GroupRecoveryDecision::RequiresExplicitRecovery
);
assert_eq!(
decide_backfill_recovery(
Some(GroupJournalEventKind::Committed),
&complete,
&plan,
BackfillExecutionDirection::Reverse,
),
GroupRecoveryDecision::RequiresExplicitRecovery
);
}
fn expected_distinct(
event: Option<GroupJournalEventKind>,
observation: &GroupRecoveryObservation,
source: &ManagedSemanticSchemaFingerprint,
target: &ManagedSemanticSchemaFingerprint,
) -> GroupRecoveryDecision {
let source_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == source);
let target_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == target);
match event {
None | Some(GroupJournalEventKind::DefinitelyAborted) if source_seen => {
GroupRecoveryDecision::ExecuteNormally
}
Some(
GroupJournalEventKind::BeforeCommit | GroupJournalEventKind::CommitOutcomeUnknown,
) if source_seen => GroupRecoveryDecision::ExecuteNormally,
Some(
GroupJournalEventKind::BeforeCommit
| GroupJournalEventKind::CommitOutcomeUnknown
| GroupJournalEventKind::Committed,
) if target_seen => GroupRecoveryDecision::RepairCheckpoint,
_ => GroupRecoveryDecision::RequiresExplicitRecovery,
}
}
fn expected_equal(
event: Option<GroupJournalEventKind>,
observation: &GroupRecoveryObservation,
both: &ManagedSemanticSchemaFingerprint,
) -> GroupRecoveryDecision {
let both_seen = matches!(observation, GroupRecoveryObservation::ManagedSemantics(value) if value == both);
if !both_seen {
return GroupRecoveryDecision::RequiresExplicitRecovery;
}
match event {
None | Some(GroupJournalEventKind::DefinitelyAborted) => {
GroupRecoveryDecision::ExecuteNormally
}
Some(GroupJournalEventKind::Committed | GroupJournalEventKind::FormalOnlyAdvanced) => {
GroupRecoveryDecision::RepairCheckpoint
}
_ => GroupRecoveryDecision::RequiresExplicitRecovery,
}
}
fn block_on<F: Future>(future: F) -> F::Output {
let mut context = Context::from_waker(Waker::noop());
let mut future = Box::pin(future);
loop {
match future.as_mut().poll(&mut context) {
Poll::Ready(output) => return output,
Poll::Pending => std::thread::yield_now(),
}
}
}
}