Skip to main content

type_bridge_migration/
recovery.rs

1//! Fail-closed per-step execution for externally owned migration ledgers.
2//!
3//! Recovery execution is deliberately separate from TypeBridge's migration
4//! state store. A caller first prepares a [`CheckedExecutionPlan`], persists
5//! or inspects its complete ordered step sequence, then supplies a
6//! [`StepRecoveryController`] that classifies every step before execution.
7//!
8//! The controller callbacks are not atomic with TypeDB commits. A durable
9//! `BeforeCommit` event narrows the uncertainty window, but a process can still
10//! disappear after TypeDB commits and before `Committed` is recorded. On the
11//! next attempt the controller must reconcile that step or classify it as
12//! indeterminate; the executor never infers safe replay from a callback alone.
13
14use std::collections::BTreeMap;
15use std::fmt;
16use std::future::Future;
17use std::pin::Pin;
18
19use serde::{Deserialize, Serialize};
20use sha2::{Digest, Sha256};
21use type_bridge_orm::Database;
22
23use crate::backfill::{BackfillResult, prepare_backfill};
24use crate::error::MigrationError;
25use crate::graph::AppliedMigrationRecord;
26use crate::plan::{
27    ExecutionPlan, ExecutionStep, MigrationAction, MigrationExecution, OperationKind, StepKind,
28    plan,
29};
30use crate::spec::MigrationGraph;
31use crate::state::{require_legacy_writer_open, require_legacy_writer_open_in_transaction};
32
33/// Boxed future returned by [`StepRecoveryController`] callbacks.
34pub type RecoveryFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
35
36/// Stable deterministic identity for one checked execution step.
37#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
38#[serde(transparent)]
39pub struct ExecutionStepId(String);
40
41impl ExecutionStepId {
42    /// Borrow the versioned hexadecimal identity.
43    pub fn as_str(&self) -> &str {
44        &self.0
45    }
46}
47
48impl fmt::Display for ExecutionStepId {
49    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
50        formatter.write_str(&self.0)
51    }
52}
53
54/// Artifact identity shared by all checked steps in one migration execution.
55#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
56pub struct CheckedMigrationIdentity {
57    /// Application or migration package label.
58    pub app_label: String,
59    /// Migration file stem.
60    pub name: String,
61    /// Checked artifact checksum.
62    pub checksum: String,
63    /// Apply or rollback direction.
64    pub action: MigrationAction,
65}
66
67/// One checked execution step exposed before mutation begins.
68#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
69pub struct CheckedExecutionStep {
70    /// Deterministic identity bound to the artifact and ordered position.
71    pub id: ExecutionStepId,
72    /// Checked migration identity.
73    pub migration: CheckedMigrationIdentity,
74    /// Zero-based position within the migration's lowered step sequence.
75    pub step_index: usize,
76    /// Artifact operation kind that produced the step.
77    pub operation_kind: OperationKind,
78    /// Transaction and TypeQL payload executed for this step.
79    pub execution: ExecutionStep,
80}
81
82/// One checked migration execution.
83#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
84pub struct CheckedMigrationExecution {
85    /// Checked artifact and direction identity.
86    pub identity: CheckedMigrationIdentity,
87    /// Ordered checked steps.
88    pub steps: Vec<CheckedExecutionStep>,
89    /// Whether every step has a reverse operation.
90    pub reversible: bool,
91}
92
93/// Checked recovery plan whose complete step sequence is inspectable.
94#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
95pub struct CheckedExecutionPlan {
96    /// Migrations scheduled for apply.
97    pub to_apply: Vec<CheckedMigrationExecution>,
98    /// Migrations scheduled for rollback.
99    pub to_rollback: Vec<CheckedMigrationExecution>,
100}
101
102impl CheckedExecutionPlan {
103    /// Iterate all steps in execution order: rollbacks, then applies.
104    pub fn ordered_steps(&self) -> impl Iterator<Item = &CheckedExecutionStep> {
105        self.to_rollback
106            .iter()
107            .chain(self.to_apply.iter())
108            .flat_map(|migration| migration.steps.iter())
109    }
110}
111
112/// Validate, plan, and bind a migration graph for fail-closed recovery.
113///
114/// This is the preferred entry point for embedders. It performs the normal
115/// graph and checksum-drift gates, then derives step identities directly from
116/// the checksums carried by that same graph.
117pub fn plan_recovery(
118    graph: &MigrationGraph,
119    applied: &[AppliedMigrationRecord],
120    target: Option<&str>,
121) -> crate::Result<CheckedExecutionPlan> {
122    let execution_plan = plan(graph, applied, target)?;
123    let checksums = graph
124        .migrations
125        .iter()
126        .filter_map(|migration| {
127            migration.checksum.as_ref().map(|checksum| {
128                (
129                    (migration.app_label.clone(), migration.name.clone()),
130                    checksum.clone(),
131                )
132            })
133        })
134        .collect();
135    prepare_recovery_plan(&execution_plan, &checksums)
136}
137
138/// Bind an execution plan to checked artifact checksums and stable step IDs.
139///
140/// This function is pure and opens no database or migration-state store. Every
141/// migration in `plan` must have a non-empty checksum entry keyed by
142/// `(app_label, name)`.
143pub fn prepare_recovery_plan(
144    plan: &ExecutionPlan,
145    checksums: &BTreeMap<(String, String), String>,
146) -> crate::Result<CheckedExecutionPlan> {
147    Ok(CheckedExecutionPlan {
148        to_apply: prepare_migrations(&plan.to_apply, checksums)?,
149        to_rollback: prepare_migrations(&plan.to_rollback, checksums)?,
150    })
151}
152
153fn prepare_migrations(
154    migrations: &[MigrationExecution],
155    checksums: &BTreeMap<(String, String), String>,
156) -> crate::Result<Vec<CheckedMigrationExecution>> {
157    migrations
158        .iter()
159        .map(|migration| {
160            let checksum = checksums
161                .get(&(migration.app_label.clone(), migration.name.clone()))
162                .filter(|checksum| !checksum.is_empty())
163                .ok_or_else(|| MigrationError::MissingRecoveryChecksum {
164                    app_label: migration.app_label.clone(),
165                    name: migration.name.clone(),
166                })?
167                .clone();
168            let identity = CheckedMigrationIdentity {
169                app_label: migration.app_label.clone(),
170                name: migration.name.clone(),
171                checksum,
172                action: migration.action,
173            };
174            let steps = migration
175                .steps
176                .iter()
177                .enumerate()
178                .map(|(step_index, execution)| CheckedExecutionStep {
179                    id: step_id(&identity, step_index, execution.operation_kind),
180                    migration: identity.clone(),
181                    step_index,
182                    operation_kind: execution.operation_kind,
183                    execution: execution.clone(),
184                })
185                .collect();
186            Ok(CheckedMigrationExecution {
187                identity,
188                steps,
189                reversible: migration.reversible,
190            })
191        })
192        .collect()
193}
194
195fn step_id(
196    migration: &CheckedMigrationIdentity,
197    step_index: usize,
198    operation_kind: OperationKind,
199) -> ExecutionStepId {
200    let mut digest = Sha256::new();
201    hash_field(&mut digest, b"type-bridge-execution-step-v1");
202    hash_field(&mut digest, migration.checksum.as_bytes());
203    hash_field(&mut digest, migration.app_label.as_bytes());
204    hash_field(&mut digest, migration.name.as_bytes());
205    hash_field(
206        &mut digest,
207        match migration.action {
208            MigrationAction::Apply => b"apply",
209            MigrationAction::Rollback => b"rollback",
210        },
211    );
212    hash_field(
213        &mut digest,
214        &u64::try_from(step_index).unwrap_or(u64::MAX).to_be_bytes(),
215    );
216    hash_field(&mut digest, operation_kind.as_str().as_bytes());
217    ExecutionStepId(format!("tb-step-v1:{:x}", digest.finalize()))
218}
219
220fn hash_field(digest: &mut Sha256, value: &[u8]) {
221    digest.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
222    digest.update(value);
223}
224
225/// Evidence supporting execution of a step currently classified as pending.
226#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227#[serde(tag = "kind", rename_all = "snake_case")]
228pub enum PendingProof {
229    /// External state proves that the step has not committed.
230    NotCommitted,
231    /// Replay is safe under a named operation-specific idempotency contract.
232    IdempotentReplay {
233        /// Stable name of the caller-owned idempotency strategy.
234        strategy: String,
235    },
236}
237
238/// External recovery decision made before a checked step can execute.
239#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
240#[serde(tag = "status", rename_all = "snake_case")]
241pub enum StepRecoveryDecision {
242    /// The step may execute under the supplied proof.
243    Pending {
244        /// Why execution or replay is safe.
245        proof: PendingProof,
246    },
247    /// The step is proven committed and must be skipped.
248    Applied {
249        /// Optional caller-owned reconciliation evidence description.
250        #[serde(default, skip_serializing_if = "Option::is_none")]
251        evidence: Option<String>,
252    },
253    /// The outcome cannot be proven; automatic execution must stop.
254    Indeterminate {
255        /// Human-readable reconciliation blocker.
256        reason: String,
257    },
258}
259
260/// Per-step event emitted around the TypeDB commit boundary.
261#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
262#[serde(rename_all = "snake_case")]
263pub enum StepRecoveryEventKind {
264    /// The query has succeeded and the transaction is ready to commit.
265    BeforeCommit,
266    /// TypeDB returned a successful commit response.
267    Committed,
268    /// The step failed before commit was called.
269    FailedBeforeCommit,
270    /// Commit returned an error, so its durable outcome is not assumed.
271    UnknownCommitOutcome,
272}
273
274/// Typed event delivered to an external recovery controller.
275#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
276pub struct StepRecoveryEvent {
277    /// Checked step associated with the event.
278    pub step: CheckedExecutionStep,
279    /// Commit-boundary event kind.
280    pub kind: StepRecoveryEventKind,
281    /// Failure or diagnostic detail when relevant.
282    #[serde(default, skip_serializing_if = "Option::is_none")]
283    pub message: Option<String>,
284    /// Derived counts for a prepared or committed backfill.
285    #[serde(default, skip_serializing_if = "Option::is_none")]
286    pub backfill: Option<BackfillResult>,
287}
288
289/// External-ledger classification and durable event boundary.
290///
291/// `classify` is invoked before every checked step. It may inspect `db` to
292/// reconcile schema state or apply an operation-specific strategy. Returning
293/// [`StepRecoveryDecision::Indeterminate`] is the required fail-closed answer
294/// whenever the outcome cannot be proven.
295pub trait StepRecoveryController: Send + Sync {
296    /// Classify a checked step before any transaction for it is opened.
297    fn classify<'a>(
298        &'a self,
299        db: &'a Database,
300        step: &'a CheckedExecutionStep,
301    ) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>>;
302
303    /// Persist or otherwise handle a typed commit-boundary event.
304    fn record_event<'a>(
305        &'a self,
306        event: StepRecoveryEvent,
307    ) -> RecoveryFuture<'a, crate::Result<()>>;
308}
309
310/// Typed result of one checked step.
311#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
312pub struct StepExecutionResult {
313    /// Checked step and stable identity.
314    pub step: CheckedExecutionStep,
315    /// Proven execution outcome.
316    pub outcome: StepExecutionOutcome,
317}
318
319/// Proven result category for one recovery step.
320#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
321#[serde(tag = "status", rename_all = "snake_case")]
322pub enum StepExecutionOutcome {
323    /// Reconciliation proved the step applied, so it was skipped.
324    Applied {
325        /// Optional caller-owned evidence description.
326        #[serde(default, skip_serializing_if = "Option::is_none")]
327        evidence: Option<String>,
328    },
329    /// The step committed and its committed event was recorded successfully.
330    Committed {
331        /// Derived counts for a backfill step.
332        #[serde(default, skip_serializing_if = "Option::is_none")]
333        backfill: Option<BackfillResult>,
334    },
335    /// The step is proven not to have reached commit.
336    FailedBeforeCommit {
337        /// Failure detail.
338        error: String,
339    },
340    /// The durable outcome or recovery evidence cannot be proven.
341    Indeterminate {
342        /// Reconciliation blocker or ambiguous commit detail.
343        error: String,
344    },
345}
346
347/// Terminal status for a checked migration attempt.
348#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
349#[serde(rename_all = "snake_case")]
350pub enum RecoveryMigrationStatus {
351    /// Every step was skipped as applied or committed successfully.
352    Succeeded,
353    /// A step failed before commit.
354    Failed,
355    /// A step was or became indeterminate.
356    Indeterminate,
357}
358
359/// Result of one checked migration attempt.
360#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
361pub struct RecoveryMigrationResult {
362    /// Checked migration identity.
363    pub migration: CheckedMigrationIdentity,
364    /// Terminal migration status.
365    pub status: RecoveryMigrationStatus,
366    /// Results for all classified steps up to the terminal step.
367    pub steps: Vec<StepExecutionResult>,
368    /// Migration-level failure summary.
369    #[serde(default, skip_serializing_if = "Option::is_none")]
370    pub error: Option<String>,
371}
372
373/// Terminal status for a recovery plan attempt.
374#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
375#[serde(rename_all = "snake_case")]
376pub enum RecoveryPlanStatus {
377    /// Every scheduled migration succeeded.
378    Succeeded,
379    /// Execution stopped on a known pre-commit failure.
380    Failed,
381    /// Execution stopped because an outcome could not be proven.
382    Indeterminate,
383}
384
385/// Results accumulated by fail-closed recovery execution.
386#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
387pub struct RecoveryExecutionResult {
388    /// Terminal plan status.
389    pub status: RecoveryPlanStatus,
390    /// Attempted migration results in execution order.
391    pub migrations: Vec<RecoveryMigrationResult>,
392}
393
394/// Execute a checked plan under an external fail-closed recovery controller.
395///
396/// Rollbacks run first, then applies. Execution stops at the first failed or
397/// indeterminate migration. No TypeBridge migration-state store is created.
398pub async fn execute_recovery_plan<C: StepRecoveryController + ?Sized>(
399    db: &Database,
400    plan: &CheckedExecutionPlan,
401    controller: &C,
402) -> RecoveryExecutionResult {
403    let mut migrations = Vec::new();
404    for migration in plan.to_rollback.iter().chain(plan.to_apply.iter()) {
405        let result = execute_recovery_migration(db, migration, controller).await;
406        let terminal = result.status;
407        migrations.push(result);
408        match terminal {
409            RecoveryMigrationStatus::Succeeded => {}
410            RecoveryMigrationStatus::Failed => {
411                return RecoveryExecutionResult {
412                    status: RecoveryPlanStatus::Failed,
413                    migrations,
414                };
415            }
416            RecoveryMigrationStatus::Indeterminate => {
417                return RecoveryExecutionResult {
418                    status: RecoveryPlanStatus::Indeterminate,
419                    migrations,
420                };
421            }
422        }
423    }
424    RecoveryExecutionResult {
425        status: RecoveryPlanStatus::Succeeded,
426        migrations,
427    }
428}
429
430async fn execute_recovery_migration<C: StepRecoveryController + ?Sized>(
431    db: &Database,
432    migration: &CheckedMigrationExecution,
433    controller: &C,
434) -> RecoveryMigrationResult {
435    if migration.identity.action == MigrationAction::Rollback && !migration.reversible {
436        return RecoveryMigrationResult {
437            migration: migration.identity.clone(),
438            status: RecoveryMigrationStatus::Failed,
439            steps: Vec::new(),
440            error: Some(format!("{} is not reversible", migration.identity.name)),
441        };
442    }
443    if let Err(error) = require_legacy_writer_open(db).await {
444        return RecoveryMigrationResult {
445            migration: migration.identity.clone(),
446            status: RecoveryMigrationStatus::Failed,
447            steps: Vec::new(),
448            error: Some(error.to_string()),
449        };
450    }
451
452    let mut results = Vec::new();
453    for step in &migration.steps {
454        let decision = match controller.classify(db, step).await {
455            Ok(decision) => decision,
456            Err(error) => {
457                let message = format!("failed to classify step {}: {error}", step.id);
458                results.push(step_result(
459                    step,
460                    StepExecutionOutcome::Indeterminate {
461                        error: message.clone(),
462                    },
463                ));
464                return migration_result(
465                    migration,
466                    RecoveryMigrationStatus::Indeterminate,
467                    results,
468                    message,
469                );
470            }
471        };
472
473        match decision {
474            StepRecoveryDecision::Applied { evidence } => {
475                results.push(step_result(
476                    step,
477                    StepExecutionOutcome::Applied { evidence },
478                ));
479            }
480            StepRecoveryDecision::Indeterminate { reason } => {
481                results.push(step_result(
482                    step,
483                    StepExecutionOutcome::Indeterminate {
484                        error: reason.clone(),
485                    },
486                ));
487                return migration_result(
488                    migration,
489                    RecoveryMigrationStatus::Indeterminate,
490                    results,
491                    reason,
492                );
493            }
494            StepRecoveryDecision::Pending { proof: _ } => {
495                let outcome = execute_pending_step(db, step, controller).await;
496                let terminal = match &outcome {
497                    StepExecutionOutcome::Applied { .. }
498                    | StepExecutionOutcome::Committed { .. } => None,
499                    StepExecutionOutcome::FailedBeforeCommit { error } => {
500                        Some((RecoveryMigrationStatus::Failed, error.clone()))
501                    }
502                    StepExecutionOutcome::Indeterminate { error } => {
503                        Some((RecoveryMigrationStatus::Indeterminate, error.clone()))
504                    }
505                };
506                results.push(step_result(step, outcome));
507                if let Some((status, error)) = terminal {
508                    return migration_result(migration, status, results, error);
509                }
510            }
511        }
512    }
513
514    RecoveryMigrationResult {
515        migration: migration.identity.clone(),
516        status: RecoveryMigrationStatus::Succeeded,
517        steps: results,
518        error: None,
519    }
520}
521
522async fn execute_pending_step<C: StepRecoveryController + ?Sized>(
523    db: &Database,
524    step: &CheckedExecutionStep,
525    controller: &C,
526) -> StepExecutionOutcome {
527    if step.execution.kind == StepKind::Backfill && step.migration.action == MigrationAction::Apply
528    {
529        return execute_pending_backfill(db, step, controller).await;
530    }
531
532    let typeql = match step.migration.action {
533        MigrationAction::Apply => step.execution.forward.as_str(),
534        MigrationAction::Rollback => match step.execution.reverse.as_deref() {
535            Some(reverse) => reverse,
536            None => {
537                return failed_before_commit(
538                    controller,
539                    step,
540                    "rollback step has no reverse TypeQL".to_string(),
541                    None,
542                )
543                .await;
544            }
545        },
546    };
547
548    if let Err(error) = db.check_schema_annotation_support(typeql) {
549        return failed_before_commit(controller, step, error.to_string(), None).await;
550    }
551    let transaction = match db.transaction_context(step.execution.tx_type).await {
552        Ok(transaction) => transaction,
553        Err(error) => {
554            return failed_before_commit(
555                controller,
556                step,
557                format!("failed to open transaction: {error}"),
558                None,
559            )
560            .await;
561        }
562    };
563    if let Err(error) = require_legacy_writer_open_in_transaction(&transaction).await {
564        let _ = transaction.rollback().await;
565        return StepExecutionOutcome::FailedBeforeCommit {
566            error: error.to_string(),
567        };
568    }
569    if let Err(error) = transaction.query(typeql).await {
570        let _ = transaction.rollback().await;
571        return failed_before_commit(controller, step, format!("query failed: {error}"), None)
572            .await;
573    }
574
575    if let Err(error) = controller
576        .record_event(event(step, StepRecoveryEventKind::BeforeCommit, None, None))
577        .await
578    {
579        let _ = transaction.rollback().await;
580        return failed_before_commit(
581            controller,
582            step,
583            format!("before-commit event was not recorded: {error}"),
584            None,
585        )
586        .await;
587    }
588
589    if let Err(error) = transaction.commit().await {
590        let message = format!("commit outcome is unknown: {error}");
591        return unknown_commit(controller, step, message, None).await;
592    }
593    committed(controller, step, None).await
594}
595
596async fn execute_pending_backfill<C: StepRecoveryController + ?Sized>(
597    db: &Database,
598    step: &CheckedExecutionStep,
599    controller: &C,
600) -> StepExecutionOutcome {
601    let prepared = match prepare_backfill(db, &step.execution, step.step_index).await {
602        Ok(prepared) => prepared,
603        Err(error) => {
604            return failed_before_commit(controller, step, error.to_string(), None).await;
605        }
606    };
607    let backfill = prepared.result.clone();
608    if let Err(error) = controller
609        .record_event(event(
610            step,
611            StepRecoveryEventKind::BeforeCommit,
612            None,
613            Some(backfill.clone()),
614        ))
615        .await
616    {
617        let _ = prepared.transaction.rollback().await;
618        return failed_before_commit(
619            controller,
620            step,
621            format!("before-commit event was not recorded: {error}"),
622            Some(backfill),
623        )
624        .await;
625    }
626    if let Err(error) = prepared.transaction.commit().await {
627        let message = format!("backfill commit outcome is unknown: {error}");
628        return unknown_commit(controller, step, message, Some(backfill)).await;
629    }
630    committed(controller, step, Some(backfill)).await
631}
632
633async fn committed<C: StepRecoveryController + ?Sized>(
634    controller: &C,
635    step: &CheckedExecutionStep,
636    backfill: Option<BackfillResult>,
637) -> StepExecutionOutcome {
638    if let Err(error) = controller
639        .record_event(event(
640            step,
641            StepRecoveryEventKind::Committed,
642            None,
643            backfill.clone(),
644        ))
645        .await
646    {
647        return StepExecutionOutcome::Indeterminate {
648            error: format!(
649                "TypeDB committed step {}, but its committed event was not durably recorded: {error}",
650                step.id
651            ),
652        };
653    }
654    StepExecutionOutcome::Committed { backfill }
655}
656
657async fn failed_before_commit<C: StepRecoveryController + ?Sized>(
658    controller: &C,
659    step: &CheckedExecutionStep,
660    mut message: String,
661    backfill: Option<BackfillResult>,
662) -> StepExecutionOutcome {
663    if let Err(event_error) = controller
664        .record_event(event(
665            step,
666            StepRecoveryEventKind::FailedBeforeCommit,
667            Some(message.clone()),
668            backfill,
669        ))
670        .await
671    {
672        message.push_str(&format!(
673            "; failed to record failed-before-commit event: {event_error}"
674        ));
675    }
676    StepExecutionOutcome::FailedBeforeCommit { error: message }
677}
678
679async fn unknown_commit<C: StepRecoveryController + ?Sized>(
680    controller: &C,
681    step: &CheckedExecutionStep,
682    mut message: String,
683    backfill: Option<BackfillResult>,
684) -> StepExecutionOutcome {
685    if let Err(event_error) = controller
686        .record_event(event(
687            step,
688            StepRecoveryEventKind::UnknownCommitOutcome,
689            Some(message.clone()),
690            backfill,
691        ))
692        .await
693    {
694        message.push_str(&format!(
695            "; failed to record unknown-commit event: {event_error}"
696        ));
697    }
698    StepExecutionOutcome::Indeterminate { error: message }
699}
700
701fn event(
702    step: &CheckedExecutionStep,
703    kind: StepRecoveryEventKind,
704    message: Option<String>,
705    backfill: Option<BackfillResult>,
706) -> StepRecoveryEvent {
707    StepRecoveryEvent {
708        step: step.clone(),
709        kind,
710        message,
711        backfill,
712    }
713}
714
715fn step_result(step: &CheckedExecutionStep, outcome: StepExecutionOutcome) -> StepExecutionResult {
716    StepExecutionResult {
717        step: step.clone(),
718        outcome,
719    }
720}
721
722fn migration_result(
723    migration: &CheckedMigrationExecution,
724    status: RecoveryMigrationStatus,
725    steps: Vec<StepExecutionResult>,
726    error: String,
727) -> RecoveryMigrationResult {
728    RecoveryMigrationResult {
729        migration: migration.identity.clone(),
730        status,
731        steps,
732        error: Some(error),
733    }
734}
735
736#[cfg(test)]
737mod tests {
738    use std::collections::{BTreeMap, VecDeque};
739    use std::sync::Mutex;
740
741    use serde_json::json;
742    use type_bridge_orm::session::backend::QueryResult;
743    use type_bridge_orm::{Database, TxType};
744
745    use super::*;
746    use crate::spec::{MigrationGraph, MigrationSpec, OperationSpec};
747    use crate::testing::{MockEvent, MockMigrationBackend};
748
749    struct TestController {
750        decisions: Mutex<VecDeque<StepRecoveryDecision>>,
751        events: Mutex<Vec<StepRecoveryEvent>>,
752        fail_event: Option<StepRecoveryEventKind>,
753    }
754
755    impl TestController {
756        fn new(decisions: Vec<StepRecoveryDecision>) -> Self {
757            Self {
758                decisions: Mutex::new(decisions.into()),
759                events: Mutex::new(Vec::new()),
760                fail_event: None,
761            }
762        }
763
764        fn failing_event(
765            decisions: Vec<StepRecoveryDecision>,
766            kind: StepRecoveryEventKind,
767        ) -> Self {
768            Self {
769                decisions: Mutex::new(decisions.into()),
770                events: Mutex::new(Vec::new()),
771                fail_event: Some(kind),
772            }
773        }
774
775        fn event_kinds(&self) -> Vec<StepRecoveryEventKind> {
776            self.events
777                .lock()
778                .unwrap()
779                .iter()
780                .map(|event| event.kind)
781                .collect()
782        }
783    }
784
785    impl StepRecoveryController for TestController {
786        fn classify<'a>(
787            &'a self,
788            _db: &'a Database,
789            _step: &'a CheckedExecutionStep,
790        ) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>> {
791            let decision = self.decisions.lock().unwrap().pop_front();
792            Box::pin(async move {
793                decision.ok_or_else(|| MigrationError::Recovery {
794                    message: "test controller has no decision".to_string(),
795                })
796            })
797        }
798
799        fn record_event<'a>(
800            &'a self,
801            event: StepRecoveryEvent,
802        ) -> RecoveryFuture<'a, crate::Result<()>> {
803            let should_fail = self.fail_event == Some(event.kind);
804            self.events.lock().unwrap().push(event);
805            Box::pin(async move {
806                if should_fail {
807                    Err(MigrationError::Recovery {
808                        message: "injected event persistence failure".to_string(),
809                    })
810                } else {
811                    Ok(())
812                }
813            })
814        }
815    }
816
817    fn pending() -> StepRecoveryDecision {
818        StepRecoveryDecision::Pending {
819            proof: PendingProof::NotCommitted,
820        }
821    }
822
823    fn applied(evidence: &str) -> StepRecoveryDecision {
824        StepRecoveryDecision::Applied {
825            evidence: Some(evidence.to_string()),
826        }
827    }
828
829    fn execution_step(
830        tx_type: TxType,
831        kind: StepKind,
832        operation_kind: OperationKind,
833        forward: &str,
834    ) -> ExecutionStep {
835        ExecutionStep {
836            tx_type,
837            kind,
838            operation_kind,
839            forward: forward.to_string(),
840            reverse: Some(format!("reverse {forward}")),
841        }
842    }
843
844    fn schema_step(forward: &str) -> ExecutionStep {
845        execution_step(
846            TxType::Schema,
847            StepKind::Schema,
848            OperationKind::AddAttribute,
849            forward,
850        )
851    }
852
853    fn write_step(forward: &str) -> ExecutionStep {
854        execution_step(
855            TxType::Write,
856            StepKind::Write,
857            OperationKind::RunTypeql,
858            forward,
859        )
860    }
861
862    fn backfill_step() -> ExecutionStep {
863        execution_step(
864            TxType::Write,
865            StepKind::Backfill,
866            OperationKind::CopyAttribute,
867            "match\n  $x isa person, has old-name $v;\n  not { $x has new-name $d; };\ninsert\n  $x has new-name == $v;",
868        )
869    }
870
871    fn checked_plan(steps: Vec<ExecutionStep>) -> CheckedExecutionPlan {
872        let plan = ExecutionPlan {
873            to_apply: vec![MigrationExecution {
874                app_label: "app".to_string(),
875                name: "0001_initial".to_string(),
876                action: MigrationAction::Apply,
877                reversible: steps.iter().all(|step| step.reverse.is_some()),
878                steps,
879            }],
880            to_rollback: Vec::new(),
881        };
882        let checksums = BTreeMap::from([(
883            ("app".to_string(), "0001_initial".to_string()),
884            "checksum-1".to_string(),
885        )]);
886        prepare_recovery_plan(&plan, &checksums).unwrap()
887    }
888
889    #[test]
890    fn checked_plan_exposes_stable_checksum_bound_step_sequence() {
891        let steps = vec![
892            schema_step("define attribute a, value string;"),
893            write_step("insert $p isa person;"),
894        ];
895        let first = checked_plan(steps.clone());
896        let second = checked_plan(steps);
897        let first_ids: Vec<_> = first.ordered_steps().map(|step| step.id.clone()).collect();
898        let second_ids: Vec<_> = second.ordered_steps().map(|step| step.id.clone()).collect();
899
900        assert_eq!(first_ids, second_ids);
901        assert_ne!(first_ids[0], first_ids[1]);
902        assert_eq!(
903            first_ids[0].as_str(),
904            "tb-step-v1:0a308e166ea81d161a8c92cbf01157935fbf53b55111df3c75f5129a183b74f9"
905        );
906        assert_eq!(first_ids[0].as_str().len(), 75);
907        assert_eq!(
908            serde_json::to_value(&first).unwrap()["to_apply"][0]["steps"]
909                .as_array()
910                .unwrap()
911                .len(),
912            2
913        );
914
915        let mut different_checksums = BTreeMap::from([(
916            ("app".to_string(), "0001_initial".to_string()),
917            "checksum-2".to_string(),
918        )]);
919        let raw = ExecutionPlan {
920            to_apply: vec![MigrationExecution {
921                app_label: "app".to_string(),
922                name: "0001_initial".to_string(),
923                action: MigrationAction::Apply,
924                steps: vec![schema_step("define attribute a, value string;")],
925                reversible: true,
926            }],
927            to_rollback: Vec::new(),
928        };
929        let changed = prepare_recovery_plan(&raw, &different_checksums).unwrap();
930        assert_ne!(first_ids[0], changed.ordered_steps().next().unwrap().id);
931        different_checksums.clear();
932        assert!(matches!(
933            prepare_recovery_plan(&raw, &different_checksums),
934            Err(MigrationError::MissingRecoveryChecksum { .. })
935        ));
936    }
937
938    #[test]
939    fn plan_recovery_binds_ids_to_the_validated_graph_artifact() {
940        let graph = MigrationGraph {
941            migrations: vec![MigrationSpec {
942                app_label: "app".to_string(),
943                name: "0001_initial".to_string(),
944                dependencies: Vec::new(),
945                operations: vec![OperationSpec::RunTypeql {
946                    forward: "define attribute a, value string;".to_string(),
947                    reverse: None,
948                }],
949                checksum: Some("artifact-checksum".to_string()),
950                source_sha256: None,
951                reversible: false,
952            }],
953        };
954
955        let checked = plan_recovery(&graph, &[], None).unwrap();
956        let step = checked.ordered_steps().next().unwrap();
957
958        assert_eq!(step.migration.checksum, "artifact-checksum");
959        assert_eq!(step.operation_kind, OperationKind::RunTypeql);
960        assert_eq!(checked.ordered_steps().count(), 1);
961    }
962
963    #[tokio::test]
964    async fn safe_resume_skips_applied_schema_and_executes_proven_pending_data() {
965        let plan = checked_plan(vec![
966            schema_step("define attribute a, value string;"),
967            write_step("insert $p isa person;"),
968        ]);
969        let controller = TestController::new(vec![applied("schema introspection"), pending()]);
970        let (backend, log) = MockMigrationBackend::new(None);
971        let db = Database::with_backend(Box::new(backend), "test");
972
973        let result = execute_recovery_plan(&db, &plan, &controller).await;
974
975        assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
976        assert!(matches!(
977            result.migrations[0].steps[0].outcome,
978            StepExecutionOutcome::Applied { .. }
979        ));
980        assert!(matches!(
981            result.migrations[0].steps[1].outcome,
982            StepExecutionOutcome::Committed { .. }
983        ));
984        assert_eq!(
985            controller.event_kinds(),
986            vec![
987                StepRecoveryEventKind::BeforeCommit,
988                StepRecoveryEventKind::Committed
989            ]
990        );
991        assert_eq!(
992            *log.lock().unwrap(),
993            vec![
994                MockEvent::OpenTx(TxType::Read),
995                MockEvent::Close,
996                MockEvent::OpenTx(TxType::Write),
997                MockEvent::Query(TxType::Write, "insert $p isa person;".to_string()),
998                MockEvent::Commit,
999            ]
1000        );
1001    }
1002
1003    #[tokio::test]
1004    async fn cutover_rejects_recovery_before_controller_or_typeql() {
1005        let plan = checked_plan(vec![write_step("insert $p isa person;")]);
1006        let controller = TestController::new(vec![pending()]);
1007        let (backend, log) = MockMigrationBackend::with_legacy_cutover();
1008        let db = Database::with_backend(Box::new(backend), "test");
1009
1010        let result = execute_recovery_plan(&db, &plan, &controller).await;
1011
1012        assert_eq!(result.status, RecoveryPlanStatus::Failed);
1013        assert!(result.migrations[0].steps.is_empty());
1014        assert!(
1015            result.migrations[0]
1016                .error
1017                .as_deref()
1018                .is_some_and(|error| error.contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE))
1019        );
1020        assert!(controller.event_kinds().is_empty());
1021        assert_eq!(
1022            *log.lock().unwrap(),
1023            vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1024        );
1025    }
1026
1027    #[tokio::test]
1028    async fn indeterminate_run_typeql_blocks_without_replay() {
1029        let plan = checked_plan(vec![write_step("insert $p isa person;")]);
1030        let controller = TestController::new(vec![StepRecoveryDecision::Indeterminate {
1031            reason: "prior before-commit receipt has no matching outcome".to_string(),
1032        }]);
1033        let (backend, log) = MockMigrationBackend::new(None);
1034        let db = Database::with_backend(Box::new(backend), "test");
1035
1036        let result = execute_recovery_plan(&db, &plan, &controller).await;
1037
1038        assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1039        assert!(matches!(
1040            result.migrations[0].steps[0].outcome,
1041            StepExecutionOutcome::Indeterminate { .. }
1042        ));
1043        assert_eq!(
1044            *log.lock().unwrap(),
1045            vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1046        );
1047        assert!(controller.event_kinds().is_empty());
1048    }
1049
1050    #[tokio::test]
1051    async fn failure_after_earlier_commit_is_typed_failed_before_commit() {
1052        let plan = checked_plan(vec![
1053            schema_step("define attribute a, value string;"),
1054            schema_step("define attribute b, value string;"),
1055        ]);
1056        let controller = TestController::new(vec![pending(), pending()]);
1057        let (backend, _log) = MockMigrationBackend::new(Some(1));
1058        let db = Database::with_backend(Box::new(backend), "test");
1059
1060        let result = execute_recovery_plan(&db, &plan, &controller).await;
1061
1062        assert_eq!(result.status, RecoveryPlanStatus::Failed);
1063        assert!(matches!(
1064            result.migrations[0].steps[0].outcome,
1065            StepExecutionOutcome::Committed { .. }
1066        ));
1067        assert!(matches!(
1068            result.migrations[0].steps[1].outcome,
1069            StepExecutionOutcome::FailedBeforeCommit { .. }
1070        ));
1071        assert_eq!(
1072            controller.event_kinds(),
1073            vec![
1074                StepRecoveryEventKind::BeforeCommit,
1075                StepRecoveryEventKind::Committed,
1076                StepRecoveryEventKind::FailedBeforeCommit,
1077            ]
1078        );
1079    }
1080
1081    #[tokio::test]
1082    async fn definitely_aborted_commit_response_remains_indeterminate_in_v1_recovery() {
1083        let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
1084        let controller = TestController::new(vec![pending()]);
1085        let (backend, _log) = MockMigrationBackend::with_definitely_aborted_commit_failure(0);
1086        let db = Database::with_backend(Box::new(backend), "test");
1087
1088        let result = execute_recovery_plan(&db, &plan, &controller).await;
1089
1090        assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1091        assert!(matches!(
1092            result.migrations[0].steps[0].outcome,
1093            StepExecutionOutcome::Indeterminate { .. }
1094        ));
1095        assert_eq!(
1096            controller.event_kinds(),
1097            vec![
1098                StepRecoveryEventKind::BeforeCommit,
1099                StepRecoveryEventKind::UnknownCommitOutcome,
1100            ]
1101        );
1102    }
1103
1104    #[tokio::test]
1105    async fn ambiguous_commit_response_is_indeterminate() {
1106        let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
1107        let controller = TestController::new(vec![pending()]);
1108        let (backend, _log) = MockMigrationBackend::with_commit_failure(0);
1109        let db = Database::with_backend(Box::new(backend), "test");
1110
1111        let result = execute_recovery_plan(&db, &plan, &controller).await;
1112
1113        assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1114        assert!(matches!(
1115            result.migrations[0].steps[0].outcome,
1116            StepExecutionOutcome::Indeterminate { .. }
1117        ));
1118        assert_eq!(
1119            controller.event_kinds(),
1120            vec![
1121                StepRecoveryEventKind::BeforeCommit,
1122                StepRecoveryEventKind::UnknownCommitOutcome,
1123            ]
1124        );
1125    }
1126
1127    #[tokio::test]
1128    async fn lost_committed_event_delivery_is_indeterminate_and_halts() {
1129        let plan = checked_plan(vec![
1130            schema_step("define attribute a, value string;"),
1131            schema_step("define attribute b, value string;"),
1132        ]);
1133        let controller = TestController::failing_event(
1134            vec![pending(), pending()],
1135            StepRecoveryEventKind::Committed,
1136        );
1137        let (backend, log) = MockMigrationBackend::new(None);
1138        let db = Database::with_backend(Box::new(backend), "test");
1139
1140        let result = execute_recovery_plan(&db, &plan, &controller).await;
1141
1142        assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1143        assert_eq!(result.migrations[0].steps.len(), 1);
1144        assert!(matches!(
1145            result.migrations[0].steps[0].outcome,
1146            StepExecutionOutcome::Indeterminate { .. }
1147        ));
1148        assert_eq!(
1149            log.lock()
1150                .unwrap()
1151                .iter()
1152                .filter(|event| matches!(event, MockEvent::Commit))
1153                .count(),
1154            1
1155        );
1156    }
1157
1158    #[tokio::test]
1159    async fn final_external_applied_record_retry_skips_all_proven_steps() {
1160        let plan = checked_plan(vec![
1161            schema_step("define attribute a, value string;"),
1162            write_step("insert $p isa person;"),
1163        ]);
1164        let controller = TestController::new(vec![
1165            applied("durable step receipt"),
1166            applied("operation-specific reconciliation"),
1167        ]);
1168        let (backend, log) = MockMigrationBackend::new(None);
1169        let db = Database::with_backend(Box::new(backend), "test");
1170
1171        let result = execute_recovery_plan(&db, &plan, &controller).await;
1172
1173        assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
1174        assert!(
1175            result.migrations[0]
1176                .steps
1177                .iter()
1178                .all(|step| matches!(step.outcome, StepExecutionOutcome::Applied { .. }))
1179        );
1180        assert_eq!(
1181            *log.lock().unwrap(),
1182            vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1183        );
1184    }
1185
1186    #[tokio::test]
1187    async fn backfill_emits_counts_on_both_commit_boundary_events() {
1188        let plan = checked_plan(vec![backfill_step()]);
1189        let controller = TestController::new(vec![StepRecoveryDecision::Pending {
1190            proof: PendingProof::IdempotentReplay {
1191                strategy: "copy-if-absent".to_string(),
1192            },
1193        }]);
1194        let responses = vec![
1195            QueryResult::Rows(vec![json!({"c": 4})]),
1196            QueryResult::Rows(vec![json!({"c": 7})]),
1197            QueryResult::Ok,
1198        ];
1199        let (backend, _log) = MockMigrationBackend::with_responses(responses);
1200        let db = Database::with_backend(Box::new(backend), "test");
1201
1202        let result = execute_recovery_plan(&db, &plan, &controller).await;
1203
1204        assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
1205        let events = controller.events.lock().unwrap();
1206        assert_eq!(events.len(), 2);
1207        for event in events.iter() {
1208            let counts = event.backfill.as_ref().unwrap();
1209            assert_eq!(counts.inserted, 4);
1210            assert_eq!(counts.matched, 7);
1211            assert_eq!(counts.skipped, 3);
1212        }
1213        assert!(matches!(
1214            result.migrations[0].steps[0].outcome,
1215            StepExecutionOutcome::Committed { backfill: Some(_) }
1216        ));
1217    }
1218}