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