use std::collections::BTreeMap;
use std::fmt;
use std::future::Future;
use std::pin::Pin;
use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use type_bridge_orm::Database;
use crate::backfill::{BackfillResult, prepare_backfill};
use crate::error::MigrationError;
use crate::graph::AppliedMigrationRecord;
use crate::plan::{
ExecutionPlan, ExecutionStep, MigrationAction, MigrationExecution, OperationKind, StepKind,
plan,
};
use crate::spec::MigrationGraph;
use crate::state::{require_legacy_writer_open, require_legacy_writer_open_in_transaction};
pub type RecoveryFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
#[serde(transparent)]
pub struct ExecutionStepId(String);
impl ExecutionStepId {
pub fn as_str(&self) -> &str {
&self.0
}
}
impl fmt::Display for ExecutionStepId {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str(&self.0)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct CheckedMigrationIdentity {
pub app_label: String,
pub name: String,
pub checksum: String,
pub action: MigrationAction,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CheckedExecutionStep {
pub id: ExecutionStepId,
pub migration: CheckedMigrationIdentity,
pub step_index: usize,
pub operation_kind: OperationKind,
pub execution: ExecutionStep,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CheckedMigrationExecution {
pub identity: CheckedMigrationIdentity,
pub steps: Vec<CheckedExecutionStep>,
pub reversible: bool,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CheckedExecutionPlan {
pub to_apply: Vec<CheckedMigrationExecution>,
pub to_rollback: Vec<CheckedMigrationExecution>,
}
impl CheckedExecutionPlan {
pub fn ordered_steps(&self) -> impl Iterator<Item = &CheckedExecutionStep> {
self.to_rollback
.iter()
.chain(self.to_apply.iter())
.flat_map(|migration| migration.steps.iter())
}
}
pub fn plan_recovery(
graph: &MigrationGraph,
applied: &[AppliedMigrationRecord],
target: Option<&str>,
) -> crate::Result<CheckedExecutionPlan> {
let execution_plan = plan(graph, applied, target)?;
let checksums = graph
.migrations
.iter()
.filter_map(|migration| {
migration.checksum.as_ref().map(|checksum| {
(
(migration.app_label.clone(), migration.name.clone()),
checksum.clone(),
)
})
})
.collect();
prepare_recovery_plan(&execution_plan, &checksums)
}
pub fn prepare_recovery_plan(
plan: &ExecutionPlan,
checksums: &BTreeMap<(String, String), String>,
) -> crate::Result<CheckedExecutionPlan> {
Ok(CheckedExecutionPlan {
to_apply: prepare_migrations(&plan.to_apply, checksums)?,
to_rollback: prepare_migrations(&plan.to_rollback, checksums)?,
})
}
fn prepare_migrations(
migrations: &[MigrationExecution],
checksums: &BTreeMap<(String, String), String>,
) -> crate::Result<Vec<CheckedMigrationExecution>> {
migrations
.iter()
.map(|migration| {
let checksum = checksums
.get(&(migration.app_label.clone(), migration.name.clone()))
.filter(|checksum| !checksum.is_empty())
.ok_or_else(|| MigrationError::MissingRecoveryChecksum {
app_label: migration.app_label.clone(),
name: migration.name.clone(),
})?
.clone();
let identity = CheckedMigrationIdentity {
app_label: migration.app_label.clone(),
name: migration.name.clone(),
checksum,
action: migration.action,
};
let steps = migration
.steps
.iter()
.enumerate()
.map(|(step_index, execution)| CheckedExecutionStep {
id: step_id(&identity, step_index, execution.operation_kind),
migration: identity.clone(),
step_index,
operation_kind: execution.operation_kind,
execution: execution.clone(),
})
.collect();
Ok(CheckedMigrationExecution {
identity,
steps,
reversible: migration.reversible,
})
})
.collect()
}
fn step_id(
migration: &CheckedMigrationIdentity,
step_index: usize,
operation_kind: OperationKind,
) -> ExecutionStepId {
let mut digest = Sha256::new();
hash_field(&mut digest, b"type-bridge-execution-step-v1");
hash_field(&mut digest, migration.checksum.as_bytes());
hash_field(&mut digest, migration.app_label.as_bytes());
hash_field(&mut digest, migration.name.as_bytes());
hash_field(
&mut digest,
match migration.action {
MigrationAction::Apply => b"apply",
MigrationAction::Rollback => b"rollback",
},
);
hash_field(
&mut digest,
&u64::try_from(step_index).unwrap_or(u64::MAX).to_be_bytes(),
);
hash_field(&mut digest, operation_kind.as_str().as_bytes());
ExecutionStepId(format!("tb-step-v1:{:x}", digest.finalize()))
}
fn hash_field(digest: &mut Sha256, value: &[u8]) {
digest.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
digest.update(value);
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "kind", rename_all = "snake_case")]
pub enum PendingProof {
NotCommitted,
IdempotentReplay {
strategy: String,
},
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum StepRecoveryDecision {
Pending {
proof: PendingProof,
},
Applied {
#[serde(default, skip_serializing_if = "Option::is_none")]
evidence: Option<String>,
},
Indeterminate {
reason: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum StepRecoveryEventKind {
BeforeCommit,
Committed,
FailedBeforeCommit,
UnknownCommitOutcome,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StepRecoveryEvent {
pub step: CheckedExecutionStep,
pub kind: StepRecoveryEventKind,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub message: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub backfill: Option<BackfillResult>,
}
pub trait StepRecoveryController: Send + Sync {
fn classify<'a>(
&'a self,
db: &'a Database,
step: &'a CheckedExecutionStep,
) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>>;
fn record_event<'a>(
&'a self,
event: StepRecoveryEvent,
) -> RecoveryFuture<'a, crate::Result<()>>;
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StepExecutionResult {
pub step: CheckedExecutionStep,
pub outcome: StepExecutionOutcome,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum StepExecutionOutcome {
Applied {
#[serde(default, skip_serializing_if = "Option::is_none")]
evidence: Option<String>,
},
Committed {
#[serde(default, skip_serializing_if = "Option::is_none")]
backfill: Option<BackfillResult>,
},
FailedBeforeCommit {
error: String,
},
Indeterminate {
error: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RecoveryMigrationStatus {
Succeeded,
Failed,
Indeterminate,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RecoveryMigrationResult {
pub migration: CheckedMigrationIdentity,
pub status: RecoveryMigrationStatus,
pub steps: Vec<StepExecutionResult>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum RecoveryPlanStatus {
Succeeded,
Failed,
Indeterminate,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RecoveryExecutionResult {
pub status: RecoveryPlanStatus,
pub migrations: Vec<RecoveryMigrationResult>,
}
pub async fn execute_recovery_plan<C: StepRecoveryController + ?Sized>(
db: &Database,
plan: &CheckedExecutionPlan,
controller: &C,
) -> RecoveryExecutionResult {
let mut migrations = Vec::new();
for migration in plan.to_rollback.iter().chain(plan.to_apply.iter()) {
let result = execute_recovery_migration(db, migration, controller).await;
let terminal = result.status;
migrations.push(result);
match terminal {
RecoveryMigrationStatus::Succeeded => {}
RecoveryMigrationStatus::Failed => {
return RecoveryExecutionResult {
status: RecoveryPlanStatus::Failed,
migrations,
};
}
RecoveryMigrationStatus::Indeterminate => {
return RecoveryExecutionResult {
status: RecoveryPlanStatus::Indeterminate,
migrations,
};
}
}
}
RecoveryExecutionResult {
status: RecoveryPlanStatus::Succeeded,
migrations,
}
}
async fn execute_recovery_migration<C: StepRecoveryController + ?Sized>(
db: &Database,
migration: &CheckedMigrationExecution,
controller: &C,
) -> RecoveryMigrationResult {
if migration.identity.action == MigrationAction::Rollback && !migration.reversible {
return RecoveryMigrationResult {
migration: migration.identity.clone(),
status: RecoveryMigrationStatus::Failed,
steps: Vec::new(),
error: Some(format!("{} is not reversible", migration.identity.name)),
};
}
if let Err(error) = require_legacy_writer_open(db).await {
return RecoveryMigrationResult {
migration: migration.identity.clone(),
status: RecoveryMigrationStatus::Failed,
steps: Vec::new(),
error: Some(error.to_string()),
};
}
let mut results = Vec::new();
for step in &migration.steps {
let decision = match controller.classify(db, step).await {
Ok(decision) => decision,
Err(error) => {
let message = format!("failed to classify step {}: {error}", step.id);
results.push(step_result(
step,
StepExecutionOutcome::Indeterminate {
error: message.clone(),
},
));
return migration_result(
migration,
RecoveryMigrationStatus::Indeterminate,
results,
message,
);
}
};
match decision {
StepRecoveryDecision::Applied { evidence } => {
results.push(step_result(
step,
StepExecutionOutcome::Applied { evidence },
));
}
StepRecoveryDecision::Indeterminate { reason } => {
results.push(step_result(
step,
StepExecutionOutcome::Indeterminate {
error: reason.clone(),
},
));
return migration_result(
migration,
RecoveryMigrationStatus::Indeterminate,
results,
reason,
);
}
StepRecoveryDecision::Pending { proof: _ } => {
let outcome = execute_pending_step(db, step, controller).await;
let terminal = match &outcome {
StepExecutionOutcome::Applied { .. }
| StepExecutionOutcome::Committed { .. } => None,
StepExecutionOutcome::FailedBeforeCommit { error } => {
Some((RecoveryMigrationStatus::Failed, error.clone()))
}
StepExecutionOutcome::Indeterminate { error } => {
Some((RecoveryMigrationStatus::Indeterminate, error.clone()))
}
};
results.push(step_result(step, outcome));
if let Some((status, error)) = terminal {
return migration_result(migration, status, results, error);
}
}
}
}
RecoveryMigrationResult {
migration: migration.identity.clone(),
status: RecoveryMigrationStatus::Succeeded,
steps: results,
error: None,
}
}
async fn execute_pending_step<C: StepRecoveryController + ?Sized>(
db: &Database,
step: &CheckedExecutionStep,
controller: &C,
) -> StepExecutionOutcome {
if step.execution.kind == StepKind::Backfill && step.migration.action == MigrationAction::Apply
{
return execute_pending_backfill(db, step, controller).await;
}
let typeql = match step.migration.action {
MigrationAction::Apply => step.execution.forward.as_str(),
MigrationAction::Rollback => match step.execution.reverse.as_deref() {
Some(reverse) => reverse,
None => {
return failed_before_commit(
controller,
step,
"rollback step has no reverse TypeQL".to_string(),
None,
)
.await;
}
},
};
if let Err(error) = db.check_schema_annotation_support(typeql) {
return failed_before_commit(controller, step, error.to_string(), None).await;
}
let transaction = match db.transaction_context(step.execution.tx_type).await {
Ok(transaction) => transaction,
Err(error) => {
return failed_before_commit(
controller,
step,
format!("failed to open transaction: {error}"),
None,
)
.await;
}
};
if let Err(error) = require_legacy_writer_open_in_transaction(&transaction).await {
let _ = transaction.rollback().await;
return StepExecutionOutcome::FailedBeforeCommit {
error: error.to_string(),
};
}
if let Err(error) = transaction.query(typeql).await {
let _ = transaction.rollback().await;
return failed_before_commit(controller, step, format!("query failed: {error}"), None)
.await;
}
if let Err(error) = controller
.record_event(event(step, StepRecoveryEventKind::BeforeCommit, None, None))
.await
{
let _ = transaction.rollback().await;
return failed_before_commit(
controller,
step,
format!("before-commit event was not recorded: {error}"),
None,
)
.await;
}
if let Err(error) = transaction.commit().await {
let message = format!("commit outcome is unknown: {error}");
return unknown_commit(controller, step, message, None).await;
}
committed(controller, step, None).await
}
async fn execute_pending_backfill<C: StepRecoveryController + ?Sized>(
db: &Database,
step: &CheckedExecutionStep,
controller: &C,
) -> StepExecutionOutcome {
let prepared = match prepare_backfill(db, &step.execution, step.step_index).await {
Ok(prepared) => prepared,
Err(error) => {
return failed_before_commit(controller, step, error.to_string(), None).await;
}
};
let backfill = prepared.result.clone();
if let Err(error) = controller
.record_event(event(
step,
StepRecoveryEventKind::BeforeCommit,
None,
Some(backfill.clone()),
))
.await
{
let _ = prepared.transaction.rollback().await;
return failed_before_commit(
controller,
step,
format!("before-commit event was not recorded: {error}"),
Some(backfill),
)
.await;
}
if let Err(error) = prepared.transaction.commit().await {
let message = format!("backfill commit outcome is unknown: {error}");
return unknown_commit(controller, step, message, Some(backfill)).await;
}
committed(controller, step, Some(backfill)).await
}
async fn committed<C: StepRecoveryController + ?Sized>(
controller: &C,
step: &CheckedExecutionStep,
backfill: Option<BackfillResult>,
) -> StepExecutionOutcome {
if let Err(error) = controller
.record_event(event(
step,
StepRecoveryEventKind::Committed,
None,
backfill.clone(),
))
.await
{
return StepExecutionOutcome::Indeterminate {
error: format!(
"TypeDB committed step {}, but its committed event was not durably recorded: {error}",
step.id
),
};
}
StepExecutionOutcome::Committed { backfill }
}
async fn failed_before_commit<C: StepRecoveryController + ?Sized>(
controller: &C,
step: &CheckedExecutionStep,
mut message: String,
backfill: Option<BackfillResult>,
) -> StepExecutionOutcome {
if let Err(event_error) = controller
.record_event(event(
step,
StepRecoveryEventKind::FailedBeforeCommit,
Some(message.clone()),
backfill,
))
.await
{
message.push_str(&format!(
"; failed to record failed-before-commit event: {event_error}"
));
}
StepExecutionOutcome::FailedBeforeCommit { error: message }
}
async fn unknown_commit<C: StepRecoveryController + ?Sized>(
controller: &C,
step: &CheckedExecutionStep,
mut message: String,
backfill: Option<BackfillResult>,
) -> StepExecutionOutcome {
if let Err(event_error) = controller
.record_event(event(
step,
StepRecoveryEventKind::UnknownCommitOutcome,
Some(message.clone()),
backfill,
))
.await
{
message.push_str(&format!(
"; failed to record unknown-commit event: {event_error}"
));
}
StepExecutionOutcome::Indeterminate { error: message }
}
fn event(
step: &CheckedExecutionStep,
kind: StepRecoveryEventKind,
message: Option<String>,
backfill: Option<BackfillResult>,
) -> StepRecoveryEvent {
StepRecoveryEvent {
step: step.clone(),
kind,
message,
backfill,
}
}
fn step_result(step: &CheckedExecutionStep, outcome: StepExecutionOutcome) -> StepExecutionResult {
StepExecutionResult {
step: step.clone(),
outcome,
}
}
fn migration_result(
migration: &CheckedMigrationExecution,
status: RecoveryMigrationStatus,
steps: Vec<StepExecutionResult>,
error: String,
) -> RecoveryMigrationResult {
RecoveryMigrationResult {
migration: migration.identity.clone(),
status,
steps,
error: Some(error),
}
}
#[cfg(test)]
mod tests {
use std::collections::{BTreeMap, VecDeque};
use std::sync::Mutex;
use serde_json::json;
use type_bridge_orm::session::backend::QueryResult;
use type_bridge_orm::{Database, TxType};
use super::*;
use crate::spec::{MigrationGraph, MigrationSpec, OperationSpec};
use crate::testing::{MockEvent, MockMigrationBackend};
struct TestController {
decisions: Mutex<VecDeque<StepRecoveryDecision>>,
events: Mutex<Vec<StepRecoveryEvent>>,
fail_event: Option<StepRecoveryEventKind>,
}
impl TestController {
fn new(decisions: Vec<StepRecoveryDecision>) -> Self {
Self {
decisions: Mutex::new(decisions.into()),
events: Mutex::new(Vec::new()),
fail_event: None,
}
}
fn failing_event(
decisions: Vec<StepRecoveryDecision>,
kind: StepRecoveryEventKind,
) -> Self {
Self {
decisions: Mutex::new(decisions.into()),
events: Mutex::new(Vec::new()),
fail_event: Some(kind),
}
}
fn event_kinds(&self) -> Vec<StepRecoveryEventKind> {
self.events
.lock()
.unwrap()
.iter()
.map(|event| event.kind)
.collect()
}
}
impl StepRecoveryController for TestController {
fn classify<'a>(
&'a self,
_db: &'a Database,
_step: &'a CheckedExecutionStep,
) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>> {
let decision = self.decisions.lock().unwrap().pop_front();
Box::pin(async move {
decision.ok_or_else(|| MigrationError::Recovery {
message: "test controller has no decision".to_string(),
})
})
}
fn record_event<'a>(
&'a self,
event: StepRecoveryEvent,
) -> RecoveryFuture<'a, crate::Result<()>> {
let should_fail = self.fail_event == Some(event.kind);
self.events.lock().unwrap().push(event);
Box::pin(async move {
if should_fail {
Err(MigrationError::Recovery {
message: "injected event persistence failure".to_string(),
})
} else {
Ok(())
}
})
}
}
fn pending() -> StepRecoveryDecision {
StepRecoveryDecision::Pending {
proof: PendingProof::NotCommitted,
}
}
fn applied(evidence: &str) -> StepRecoveryDecision {
StepRecoveryDecision::Applied {
evidence: Some(evidence.to_string()),
}
}
fn execution_step(
tx_type: TxType,
kind: StepKind,
operation_kind: OperationKind,
forward: &str,
) -> ExecutionStep {
ExecutionStep {
tx_type,
kind,
operation_kind,
forward: forward.to_string(),
reverse: Some(format!("reverse {forward}")),
}
}
fn schema_step(forward: &str) -> ExecutionStep {
execution_step(
TxType::Schema,
StepKind::Schema,
OperationKind::AddAttribute,
forward,
)
}
fn write_step(forward: &str) -> ExecutionStep {
execution_step(
TxType::Write,
StepKind::Write,
OperationKind::RunTypeql,
forward,
)
}
fn backfill_step() -> ExecutionStep {
execution_step(
TxType::Write,
StepKind::Backfill,
OperationKind::CopyAttribute,
"match\n $x isa person, has old-name $v;\n not { $x has new-name $d; };\ninsert\n $x has new-name == $v;",
)
}
fn checked_plan(steps: Vec<ExecutionStep>) -> CheckedExecutionPlan {
let plan = ExecutionPlan {
to_apply: vec![MigrationExecution {
app_label: "app".to_string(),
name: "0001_initial".to_string(),
action: MigrationAction::Apply,
reversible: steps.iter().all(|step| step.reverse.is_some()),
steps,
}],
to_rollback: Vec::new(),
};
let checksums = BTreeMap::from([(
("app".to_string(), "0001_initial".to_string()),
"checksum-1".to_string(),
)]);
prepare_recovery_plan(&plan, &checksums).unwrap()
}
#[test]
fn checked_plan_exposes_stable_checksum_bound_step_sequence() {
let steps = vec![
schema_step("define attribute a, value string;"),
write_step("insert $p isa person;"),
];
let first = checked_plan(steps.clone());
let second = checked_plan(steps);
let first_ids: Vec<_> = first.ordered_steps().map(|step| step.id.clone()).collect();
let second_ids: Vec<_> = second.ordered_steps().map(|step| step.id.clone()).collect();
assert_eq!(first_ids, second_ids);
assert_ne!(first_ids[0], first_ids[1]);
assert_eq!(
first_ids[0].as_str(),
"tb-step-v1:0a308e166ea81d161a8c92cbf01157935fbf53b55111df3c75f5129a183b74f9"
);
assert_eq!(first_ids[0].as_str().len(), 75);
assert_eq!(
serde_json::to_value(&first).unwrap()["to_apply"][0]["steps"]
.as_array()
.unwrap()
.len(),
2
);
let mut different_checksums = BTreeMap::from([(
("app".to_string(), "0001_initial".to_string()),
"checksum-2".to_string(),
)]);
let raw = ExecutionPlan {
to_apply: vec![MigrationExecution {
app_label: "app".to_string(),
name: "0001_initial".to_string(),
action: MigrationAction::Apply,
steps: vec![schema_step("define attribute a, value string;")],
reversible: true,
}],
to_rollback: Vec::new(),
};
let changed = prepare_recovery_plan(&raw, &different_checksums).unwrap();
assert_ne!(first_ids[0], changed.ordered_steps().next().unwrap().id);
different_checksums.clear();
assert!(matches!(
prepare_recovery_plan(&raw, &different_checksums),
Err(MigrationError::MissingRecoveryChecksum { .. })
));
}
#[test]
fn plan_recovery_binds_ids_to_the_validated_graph_artifact() {
let graph = MigrationGraph {
migrations: vec![MigrationSpec {
app_label: "app".to_string(),
name: "0001_initial".to_string(),
dependencies: Vec::new(),
operations: vec![OperationSpec::RunTypeql {
forward: "define attribute a, value string;".to_string(),
reverse: None,
}],
checksum: Some("artifact-checksum".to_string()),
source_sha256: None,
reversible: false,
}],
};
let checked = plan_recovery(&graph, &[], None).unwrap();
let step = checked.ordered_steps().next().unwrap();
assert_eq!(step.migration.checksum, "artifact-checksum");
assert_eq!(step.operation_kind, OperationKind::RunTypeql);
assert_eq!(checked.ordered_steps().count(), 1);
}
#[tokio::test]
async fn safe_resume_skips_applied_schema_and_executes_proven_pending_data() {
let plan = checked_plan(vec![
schema_step("define attribute a, value string;"),
write_step("insert $p isa person;"),
]);
let controller = TestController::new(vec![applied("schema introspection"), pending()]);
let (backend, log) = MockMigrationBackend::new(None);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Applied { .. }
));
assert!(matches!(
result.migrations[0].steps[1].outcome,
StepExecutionOutcome::Committed { .. }
));
assert_eq!(
controller.event_kinds(),
vec![
StepRecoveryEventKind::BeforeCommit,
StepRecoveryEventKind::Committed
]
);
assert_eq!(
*log.lock().unwrap(),
vec![
MockEvent::OpenTx(TxType::Read),
MockEvent::Close,
MockEvent::OpenTx(TxType::Write),
MockEvent::Query(TxType::Write, "insert $p isa person;".to_string()),
MockEvent::Commit,
]
);
}
#[tokio::test]
async fn cutover_rejects_recovery_before_controller_or_typeql() {
let plan = checked_plan(vec![write_step("insert $p isa person;")]);
let controller = TestController::new(vec![pending()]);
let (backend, log) = MockMigrationBackend::with_legacy_cutover();
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Failed);
assert!(result.migrations[0].steps.is_empty());
assert!(
result.migrations[0]
.error
.as_deref()
.is_some_and(|error| error.contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE))
);
assert!(controller.event_kinds().is_empty());
assert_eq!(
*log.lock().unwrap(),
vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
);
}
#[tokio::test]
async fn indeterminate_run_typeql_blocks_without_replay() {
let plan = checked_plan(vec![write_step("insert $p isa person;")]);
let controller = TestController::new(vec![StepRecoveryDecision::Indeterminate {
reason: "prior before-commit receipt has no matching outcome".to_string(),
}]);
let (backend, log) = MockMigrationBackend::new(None);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Indeterminate { .. }
));
assert_eq!(
*log.lock().unwrap(),
vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
);
assert!(controller.event_kinds().is_empty());
}
#[tokio::test]
async fn failure_after_earlier_commit_is_typed_failed_before_commit() {
let plan = checked_plan(vec![
schema_step("define attribute a, value string;"),
schema_step("define attribute b, value string;"),
]);
let controller = TestController::new(vec![pending(), pending()]);
let (backend, _log) = MockMigrationBackend::new(Some(1));
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Failed);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Committed { .. }
));
assert!(matches!(
result.migrations[0].steps[1].outcome,
StepExecutionOutcome::FailedBeforeCommit { .. }
));
assert_eq!(
controller.event_kinds(),
vec![
StepRecoveryEventKind::BeforeCommit,
StepRecoveryEventKind::Committed,
StepRecoveryEventKind::FailedBeforeCommit,
]
);
}
#[tokio::test]
async fn definitely_aborted_commit_response_remains_indeterminate_in_v1_recovery() {
let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
let controller = TestController::new(vec![pending()]);
let (backend, _log) = MockMigrationBackend::with_definitely_aborted_commit_failure(0);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Indeterminate { .. }
));
assert_eq!(
controller.event_kinds(),
vec![
StepRecoveryEventKind::BeforeCommit,
StepRecoveryEventKind::UnknownCommitOutcome,
]
);
}
#[tokio::test]
async fn ambiguous_commit_response_is_indeterminate() {
let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
let controller = TestController::new(vec![pending()]);
let (backend, _log) = MockMigrationBackend::with_commit_failure(0);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Indeterminate { .. }
));
assert_eq!(
controller.event_kinds(),
vec![
StepRecoveryEventKind::BeforeCommit,
StepRecoveryEventKind::UnknownCommitOutcome,
]
);
}
#[tokio::test]
async fn lost_committed_event_delivery_is_indeterminate_and_halts() {
let plan = checked_plan(vec![
schema_step("define attribute a, value string;"),
schema_step("define attribute b, value string;"),
]);
let controller = TestController::failing_event(
vec![pending(), pending()],
StepRecoveryEventKind::Committed,
);
let (backend, log) = MockMigrationBackend::new(None);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
assert_eq!(result.migrations[0].steps.len(), 1);
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Indeterminate { .. }
));
assert_eq!(
log.lock()
.unwrap()
.iter()
.filter(|event| matches!(event, MockEvent::Commit))
.count(),
1
);
}
#[tokio::test]
async fn final_external_applied_record_retry_skips_all_proven_steps() {
let plan = checked_plan(vec![
schema_step("define attribute a, value string;"),
write_step("insert $p isa person;"),
]);
let controller = TestController::new(vec![
applied("durable step receipt"),
applied("operation-specific reconciliation"),
]);
let (backend, log) = MockMigrationBackend::new(None);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
assert!(
result.migrations[0]
.steps
.iter()
.all(|step| matches!(step.outcome, StepExecutionOutcome::Applied { .. }))
);
assert_eq!(
*log.lock().unwrap(),
vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
);
}
#[tokio::test]
async fn backfill_emits_counts_on_both_commit_boundary_events() {
let plan = checked_plan(vec![backfill_step()]);
let controller = TestController::new(vec![StepRecoveryDecision::Pending {
proof: PendingProof::IdempotentReplay {
strategy: "copy-if-absent".to_string(),
},
}]);
let responses = vec![
QueryResult::Rows(vec![json!({"c": 4})]),
QueryResult::Rows(vec![json!({"c": 7})]),
QueryResult::Ok,
];
let (backend, _log) = MockMigrationBackend::with_responses(responses);
let db = Database::with_backend(Box::new(backend), "test");
let result = execute_recovery_plan(&db, &plan, &controller).await;
assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
let events = controller.events.lock().unwrap();
assert_eq!(events.len(), 2);
for event in events.iter() {
let counts = event.backfill.as_ref().unwrap();
assert_eq!(counts.inserted, 4);
assert_eq!(counts.matched, 7);
assert_eq!(counts.skipped, 3);
}
assert!(matches!(
result.migrations[0].steps[0].outcome,
StepExecutionOutcome::Committed { backfill: Some(_) }
));
}
}