Skip to main content

type_bridge_migration/
executor.rs

1//! Async migration executor.
2//!
3//! Runs an [`ExecutionPlan`] produced by the planner over a [`Database`],
4//! committing each step in its own transaction.  The executor is
5//! runtime-agnostic: the caller (`PyMigrationRunner` in Phase 3) drives it by
6//! `block_on`-ing the returned future on the shared `Arc<Runtime>`.
7//!
8//! # Per-step commit contract
9//!
10//! TypeDB forbids multiple `define` blocks per transaction.  Each
11//! [`ExecutionStep`](crate::plan::ExecutionStep) is executed in its own
12//! transaction and committed
13//! independently.  When a step fails, earlier committed steps in the same
14//! migration are **not** rolled back (there is no cross-step atomicity).
15//! This matches the current Python per-statement-commit behavior and is the
16//! documented contract for this boundary.
17
18use serde::{Deserialize, Serialize};
19use std::collections::BTreeMap;
20use type_bridge_orm::Database;
21
22use crate::backfill::{BackfillResult, execute_backfill};
23use crate::plan::{ExecutionPlan, MigrationAction, MigrationExecution, StepKind};
24use crate::state::{
25    MigrationExecutorInfo, MigrationStateStore, finished_run_record, require_legacy_writer_open,
26    require_legacy_writer_open_in_transaction, started_run_record,
27};
28
29/// Result of executing a single migration (one entry per attempted migration).
30#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
31pub struct MigrationResult {
32    /// Application or package label.
33    pub app_label: String,
34    /// Migration file stem, e.g. `0001_initial`.
35    pub name: String,
36    /// Whether this was an apply or a rollback.
37    pub action: MigrationAction,
38    /// `true` when all steps completed and committed successfully.
39    pub success: bool,
40    /// Human-readable failure reason, present only when `success` is `false`.
41    pub error: Option<String>,
42    /// Per-step backfill counts, present only when the migration contained at
43    /// least one [`StepKind::Backfill`] step.  `None` for pure-schema migrations
44    /// so the field does not appear in JSON output (no bloat — D2a).
45    #[serde(default, skip_serializing_if = "Option::is_none")]
46    pub backfill: Option<Vec<BackfillResult>>,
47}
48
49/// Execute a [`ExecutionPlan`] against `db`.
50///
51/// Rollbacks are processed first (already in reverse order from the planner),
52/// then applies.  The run halts at the first failed migration and returns all
53/// results accumulated so far including the failure.
54///
55/// Returns one [`MigrationResult`] per attempted migration, in execution order
56/// (rollbacks then applies).
57pub async fn execute_plan(db: &Database, plan: ExecutionPlan) -> Vec<MigrationResult> {
58    let mut results: Vec<MigrationResult> = Vec::new();
59
60    // Process rollbacks before applies (planner already reverse-ordered them).
61    for migration in plan.to_rollback {
62        let result = execute_migration(db, &migration).await;
63        let should_halt = !result.success;
64        results.push(result);
65        if should_halt {
66            return results;
67        }
68    }
69
70    // Process applies.
71    for migration in plan.to_apply {
72        let result = execute_migration(db, &migration).await;
73        let should_halt = !result.success;
74        results.push(result);
75        if should_halt {
76            return results;
77        }
78    }
79
80    results
81}
82
83/// Execute a plan while writing one DB-backed run-log row per attempted migration.
84pub async fn execute_plan_with_run_log<S: MigrationStateStore>(
85    db: &Database,
86    store: &S,
87    plan: ExecutionPlan,
88    checksums: &BTreeMap<(String, String), String>,
89    executor: &MigrationExecutorInfo,
90) -> crate::Result<Vec<MigrationResult>> {
91    let mut results: Vec<MigrationResult> = Vec::new();
92
93    for migration in plan.to_rollback {
94        let result = execute_logged_migration(db, store, &migration, checksums, executor).await?;
95        let should_halt = !result.success;
96        results.push(result);
97        if should_halt {
98            return Ok(results);
99        }
100    }
101
102    for migration in plan.to_apply {
103        let result = execute_logged_migration(db, store, &migration, checksums, executor).await?;
104        let should_halt = !result.success;
105        results.push(result);
106        if should_halt {
107            return Ok(results);
108        }
109    }
110
111    Ok(results)
112}
113
114async fn execute_logged_migration<S: MigrationStateStore>(
115    db: &Database,
116    store: &S,
117    migration: &MigrationExecution,
118    checksums: &BTreeMap<(String, String), String>,
119    executor: &MigrationExecutorInfo,
120) -> crate::Result<MigrationResult> {
121    // External state stores cannot share the target transaction, so reject an
122    // already-cut-over target before touching their run log.  Each target
123    // mutation repeats the check in its own transaction below.
124    require_legacy_writer_open(db).await?;
125    let checksum = checksums
126        .get(&(migration.app_label.clone(), migration.name.clone()))
127        .cloned()
128        .unwrap_or_default();
129    let run = started_run_record(migration, checksum, executor);
130    store.record_run(run.clone()).await?;
131
132    let result = execute_migration(db, migration).await;
133    let status = if result.success {
134        "succeeded"
135    } else {
136        "failed"
137    };
138    let finished = finished_run_record(run, status, result.error.clone());
139    store.record_run(finished).await?;
140    Ok(result)
141}
142
143/// Execute a single [`MigrationExecution`].
144///
145/// Returns one [`MigrationResult`] for this attempted migration.
146pub async fn execute_migration(db: &Database, migration: &MigrationExecution) -> MigrationResult {
147    // Non-reversible rollback: fail immediately without opening a transaction.
148    if migration.action == MigrationAction::Rollback && !migration.reversible {
149        return MigrationResult {
150            app_label: migration.app_label.clone(),
151            name: migration.name.clone(),
152            action: migration.action,
153            success: false,
154            error: Some(format!("{} is not reversible", migration.name)),
155            backfill: None,
156        };
157    }
158
159    if let Err(error) = require_legacy_writer_open(db).await {
160        return MigrationResult {
161            app_label: migration.app_label.clone(),
162            name: migration.name.clone(),
163            action: migration.action,
164            success: false,
165            error: Some(error.to_string()),
166            backfill: None,
167        };
168    }
169
170    let mut backfill_results: Vec<BackfillResult> = Vec::new();
171
172    for (step_index, step) in migration.steps.iter().enumerate() {
173        // Backfill steps are routed through the count-deriving path.
174        if step.kind == StepKind::Backfill && migration.action == MigrationAction::Apply {
175            match execute_backfill(db, step, step_index).await {
176                Ok(bf_result) => {
177                    backfill_results.push(bf_result);
178                    // Backfill execution committed internally; continue.
179                    continue;
180                }
181                Err(e) => {
182                    return MigrationResult {
183                        app_label: migration.app_label.clone(),
184                        name: migration.name.clone(),
185                        action: migration.action,
186                        success: false,
187                        error: Some(format!("backfill step {step_index} failed: {e}")),
188                        backfill: None,
189                    };
190                }
191            }
192        }
193
194        // Choose forward or reverse TypeQL based on the action.
195        let typeql: &str = match migration.action {
196            MigrationAction::Apply => &step.forward,
197            MigrationAction::Rollback => {
198                // reversible was true above, so every step has a reverse.
199                // Unwrap is safe here — the planner guarantees this when
200                // `migration.reversible == true`.
201                step.reverse.as_deref().unwrap_or(&step.forward)
202            }
203        };
204
205        // Version-gate annotation-bearing schema DDL before opening a
206        // transaction: pre-3.12 servers reject @doc/@meta with a syntax
207        // error; the gate produces an actionable versioned error instead.
208        if let Err(e) = db.check_schema_annotation_support(typeql) {
209            return MigrationResult {
210                app_label: migration.app_label.clone(),
211                name: migration.name.clone(),
212                action: migration.action,
213                success: false,
214                error: Some(e.to_string()),
215                backfill: None,
216            };
217        }
218
219        // Open a transaction for this step.
220        let ctx = match db.transaction_context(step.tx_type).await {
221            Ok(ctx) => ctx,
222            Err(e) => {
223                return MigrationResult {
224                    app_label: migration.app_label.clone(),
225                    name: migration.name.clone(),
226                    action: migration.action,
227                    success: false,
228                    error: Some(format!("failed to open transaction: {e}")),
229                    backfill: None,
230                };
231            }
232        };
233
234        if let Err(error) = require_legacy_writer_open_in_transaction(&ctx).await {
235            let _ = ctx.rollback().await;
236            return MigrationResult {
237                app_label: migration.app_label.clone(),
238                name: migration.name.clone(),
239                action: migration.action,
240                success: false,
241                error: Some(error.to_string()),
242                backfill: None,
243            };
244        }
245
246        // Execute the query.
247        if let Err(e) = ctx.query(typeql).await {
248            // Best-effort rollback; ignore its error.
249            let _ = ctx.rollback().await;
250            return MigrationResult {
251                app_label: migration.app_label.clone(),
252                name: migration.name.clone(),
253                action: migration.action,
254                success: false,
255                error: Some(format!("query failed: {e}")),
256                backfill: None,
257            };
258        }
259
260        // Commit the step.
261        if let Err(e) = ctx.commit().await {
262            return MigrationResult {
263                app_label: migration.app_label.clone(),
264                name: migration.name.clone(),
265                action: migration.action,
266                success: false,
267                error: Some(format!("commit failed: {e}")),
268                backfill: None,
269            };
270        }
271    }
272
273    // All steps succeeded.
274    let backfill = if backfill_results.is_empty() {
275        None
276    } else {
277        Some(backfill_results)
278    };
279    MigrationResult {
280        app_label: migration.app_label.clone(),
281        name: migration.name.clone(),
282        action: migration.action,
283        success: true,
284        error: None,
285        backfill,
286    }
287}
288
289#[cfg(test)]
290mod tests {
291    use super::*;
292    use crate::plan::ExecutionStep;
293    use crate::state::InMemoryStateStore;
294    use crate::testing::{MockEvent, MockMigrationBackend};
295    use type_bridge_orm::{Database, TxType};
296
297    // ── helpers ───────────────────────────────────────────────────────────────
298
299    fn schema_step(forward: &str, reverse: Option<&str>) -> ExecutionStep {
300        ExecutionStep {
301            tx_type: TxType::Schema,
302            kind: crate::plan::StepKind::Schema,
303            operation_kind: crate::plan::OperationKind::RunTypeql,
304            forward: forward.to_string(),
305            reverse: reverse.map(str::to_string),
306        }
307    }
308
309    fn apply_migration(name: &str, steps: Vec<ExecutionStep>) -> MigrationExecution {
310        let reversible = steps.iter().all(|s| s.reverse.is_some());
311        MigrationExecution {
312            app_label: "app".to_string(),
313            name: name.to_string(),
314            action: MigrationAction::Apply,
315            steps,
316            reversible,
317        }
318    }
319
320    fn rollback_migration(
321        name: &str,
322        steps: Vec<ExecutionStep>,
323        reversible: bool,
324    ) -> MigrationExecution {
325        MigrationExecution {
326            app_label: "app".to_string(),
327            name: name.to_string(),
328            action: MigrationAction::Rollback,
329            steps,
330            reversible,
331        }
332    }
333
334    #[tokio::test]
335    async fn cutover_rejects_before_executor_typeql_or_run_log_mutation() {
336        let (backend, log) = MockMigrationBackend::with_legacy_cutover();
337        let db = Database::with_backend(Box::new(backend), "test");
338        let migration = apply_migration(
339            "0001_blocked",
340            vec![schema_step("define entity must-not-run;", None)],
341        );
342        let plan = ExecutionPlan {
343            to_apply: vec![migration.clone()],
344            to_rollback: Vec::new(),
345        };
346
347        let results = execute_plan(&db, plan).await;
348        assert_eq!(results.len(), 1);
349        assert!(!results[0].success);
350        assert!(
351            results[0]
352                .error
353                .as_deref()
354                .is_some_and(|error| error.contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE))
355        );
356        assert_eq!(
357            *log.lock().unwrap(),
358            vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
359        );
360
361        let store = InMemoryStateStore::new();
362        let logged = execute_plan_with_run_log(
363            &db,
364            &store,
365            ExecutionPlan {
366                to_apply: vec![migration],
367                to_rollback: Vec::new(),
368            },
369            &BTreeMap::new(),
370            &MigrationExecutorInfo::default(),
371        )
372        .await
373        .expect_err("cutover must reject before the external run log");
374        assert!(
375            logged
376                .to_string()
377                .contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE)
378        );
379        assert!(store.load_runs().await.unwrap().is_empty());
380    }
381
382    // ── test: apply-only plan ─────────────────────────────────────────────────
383
384    #[tokio::test]
385    async fn apply_only_executes_steps_in_order_under_schema_tx() {
386        let (backend, log) = MockMigrationBackend::new(None);
387        let db = Database::with_backend(Box::new(backend), "test");
388
389        let plan = ExecutionPlan {
390            to_apply: vec![
391                apply_migration(
392                    "0001_initial",
393                    vec![schema_step("define attribute a, value string;", None)],
394                ),
395                apply_migration(
396                    "0002_add",
397                    vec![schema_step(
398                        "define attribute b, value string;",
399                        Some("undefine attribute b;"),
400                    )],
401                ),
402            ],
403            to_rollback: Vec::new(),
404        };
405
406        let results = execute_plan(&db, plan).await;
407
408        assert_eq!(results.len(), 2);
409        assert!(results[0].success);
410        assert_eq!(results[0].name, "0001_initial");
411        assert!(results[1].success);
412        assert_eq!(results[1].name, "0002_add");
413
414        // Each migration first performs the read-only cutover preflight, then
415        // executes its step under the intended schema transaction.
416        let events = log.lock().unwrap();
417        assert_eq!(
418            *events,
419            vec![
420                MockEvent::OpenTx(TxType::Read),
421                MockEvent::Close,
422                MockEvent::OpenTx(TxType::Schema),
423                MockEvent::Query(
424                    TxType::Schema,
425                    "define attribute a, value string;".to_string()
426                ),
427                MockEvent::Commit,
428                MockEvent::OpenTx(TxType::Read),
429                MockEvent::Close,
430                MockEvent::OpenTx(TxType::Schema),
431                MockEvent::Query(
432                    TxType::Schema,
433                    "define attribute b, value string;".to_string()
434                ),
435                MockEvent::Commit,
436            ]
437        );
438    }
439
440    #[tokio::test]
441    async fn execute_plan_with_run_log_records_started_and_finished_rows() {
442        let (backend, _log) = MockMigrationBackend::new(None);
443        let db = Database::with_backend(Box::new(backend), "test");
444        let store = InMemoryStateStore::new();
445        let mut checksums = BTreeMap::new();
446        checksums.insert(
447            ("app".to_string(), "0001_initial".to_string()),
448            "checksum-1".to_string(),
449        );
450        let executor = MigrationExecutorInfo {
451            ip: Some("127.0.0.1".to_string()),
452            mac: Some("00:11:22:33:44:55".to_string()),
453        };
454        let plan = ExecutionPlan {
455            to_apply: vec![apply_migration(
456                "0001_initial",
457                vec![schema_step("define attribute a, value string;", None)],
458            )],
459            to_rollback: Vec::new(),
460        };
461
462        let results = execute_plan_with_run_log(&db, &store, plan, &checksums, &executor)
463            .await
464            .unwrap();
465
466        assert_eq!(results.len(), 1);
467        assert!(results[0].success);
468        let runs = store.load_runs().await.unwrap();
469        assert_eq!(runs.len(), 1);
470        assert_eq!(runs[0].name, "0001_initial");
471        assert_eq!(runs[0].checksum, "checksum-1");
472        assert_eq!(runs[0].direction, "apply");
473        assert_eq!(runs[0].status, "succeeded");
474        assert!(runs[0].finished_at.is_some());
475        assert_eq!(runs[0].executor_ip.as_deref(), Some("127.0.0.1"));
476        assert_eq!(runs[0].executor_mac.as_deref(), Some("00:11:22:33:44:55"));
477    }
478
479    // ── test: rollback path ───────────────────────────────────────────────────
480
481    #[tokio::test]
482    async fn rollback_executes_reverse_typeql_under_schema_tx() {
483        let (backend, log) = MockMigrationBackend::new(None);
484        let db = Database::with_backend(Box::new(backend), "test");
485
486        let plan = ExecutionPlan {
487            to_apply: Vec::new(),
488            to_rollback: vec![rollback_migration(
489                "0002_add",
490                vec![schema_step(
491                    "define attribute b, value string;",
492                    Some("undefine attribute b;"),
493                )],
494                true,
495            )],
496        };
497
498        let results = execute_plan(&db, plan).await;
499
500        assert_eq!(results.len(), 1);
501        assert!(results[0].success);
502        assert_eq!(results[0].action, MigrationAction::Rollback);
503
504        let events = log.lock().unwrap();
505        assert_eq!(
506            *events,
507            vec![
508                MockEvent::OpenTx(TxType::Read),
509                MockEvent::Close,
510                MockEvent::OpenTx(TxType::Schema),
511                // Reverse TypeQL is used for rollback.
512                MockEvent::Query(TxType::Schema, "undefine attribute b;".to_string()),
513                MockEvent::Commit,
514            ]
515        );
516    }
517
518    // ── test: per-step-commit / no atomicity (oracle risk 1) ──────────────────
519    //
520    // A migration with 3 steps where the mock fails the 2nd query:
521    //   - Step 1: Query + Commit (durable)
522    //   - Step 2: Query fails → Rollback
523    //   - Step 3: never opened
524    //   - Run halts after the failure
525    //   - Result: success=false
526
527    #[tokio::test]
528    async fn middle_step_failure_commits_prior_steps_and_halts() {
529        // Fail the 2nd query (0-indexed: index 1).
530        let (backend, log) = MockMigrationBackend::new(Some(1));
531        let db = Database::with_backend(Box::new(backend), "test");
532
533        let plan = ExecutionPlan {
534            to_apply: vec![apply_migration(
535                "0001_three_steps",
536                vec![
537                    schema_step("define attribute a, value string;", None),
538                    schema_step("define attribute b, value string;", None),
539                    schema_step("define attribute c, value string;", None),
540                ],
541            )],
542            to_rollback: Vec::new(),
543        };
544
545        let results = execute_plan(&db, plan).await;
546
547        // Only one migration attempted; it failed.
548        assert_eq!(results.len(), 1);
549        assert!(!results[0].success);
550        assert!(results[0].error.is_some());
551
552        let events = log.lock().unwrap();
553
554        // Migration entry: read-only cutover guard, then close the snapshot.
555        assert!(matches!(events[0], MockEvent::OpenTx(TxType::Read)));
556        assert!(matches!(events[1], MockEvent::Close));
557
558        // Step 1: OpenTx → Query → Commit (durable)
559        assert!(matches!(events[2], MockEvent::OpenTx(TxType::Schema)));
560        assert!(
561            matches!(&events[3], MockEvent::Query(TxType::Schema, q) if q == "define attribute a, value string;")
562        );
563        assert!(matches!(events[4], MockEvent::Commit));
564
565        // Step 2: OpenTx → Query (fails) → Rollback
566        assert!(matches!(events[5], MockEvent::OpenTx(TxType::Schema)));
567        assert!(
568            matches!(&events[6], MockEvent::Query(TxType::Schema, q) if q == "define attribute b, value string;")
569        );
570        assert!(matches!(events[7], MockEvent::Rollback));
571
572        // Step 3 must NOT appear — run halted after step 2 failure.
573        assert_eq!(events.len(), 8, "step 3 must not be opened");
574    }
575
576    // ── test: non-reversible rollback ─────────────────────────────────────────
577
578    #[tokio::test]
579    async fn non_reversible_rollback_returns_failed_result_without_opening_tx() {
580        let (backend, log) = MockMigrationBackend::new(None);
581        let db = Database::with_backend(Box::new(backend), "test");
582
583        let plan = ExecutionPlan {
584            to_apply: Vec::new(),
585            to_rollback: vec![rollback_migration(
586                "0001_initial",
587                // Steps carry no reverse — this migration is non-reversible.
588                vec![schema_step("define attribute a, value string;", None)],
589                false,
590            )],
591        };
592
593        let results = execute_plan(&db, plan).await;
594
595        assert_eq!(results.len(), 1);
596        assert!(!results[0].success);
597        assert!(
598            results[0]
599                .error
600                .as_deref()
601                .unwrap_or("")
602                .contains("not reversible"),
603            "error message must mention 'not reversible'; got: {:?}",
604            results[0].error
605        );
606
607        // No transaction was opened.
608        let events = log.lock().unwrap();
609        assert!(
610            events.is_empty(),
611            "no events expected for non-reversible rollback; got: {events:?}"
612        );
613    }
614
615    // ── test: results are in execution order (rollbacks then applies) ─────────
616
617    #[tokio::test]
618    async fn results_are_in_execution_order_rollbacks_first() {
619        let (backend, _log) = MockMigrationBackend::new(None);
620        let db = Database::with_backend(Box::new(backend), "test");
621
622        let plan = ExecutionPlan {
623            to_apply: vec![apply_migration(
624                "0003_apply",
625                vec![schema_step("define attribute c, value string;", None)],
626            )],
627            to_rollback: vec![rollback_migration(
628                "0002_rollback",
629                vec![schema_step(
630                    "define attribute b, value string;",
631                    Some("undefine attribute b;"),
632                )],
633                true,
634            )],
635        };
636
637        let results = execute_plan(&db, plan).await;
638
639        assert_eq!(results.len(), 2);
640        // Rollback comes first in results.
641        assert_eq!(results[0].name, "0002_rollback");
642        assert_eq!(results[0].action, MigrationAction::Rollback);
643        assert_eq!(results[1].name, "0003_apply");
644        assert_eq!(results[1].action, MigrationAction::Apply);
645    }
646
647    // ── test: failure halts remaining migrations ──────────────────────────────
648
649    #[tokio::test]
650    async fn failure_in_first_migration_halts_remaining() {
651        // Fail the very first query.
652        let (backend, _log) = MockMigrationBackend::new(Some(0));
653        let db = Database::with_backend(Box::new(backend), "test");
654
655        let plan = ExecutionPlan {
656            to_apply: vec![
657                apply_migration(
658                    "0001_fail",
659                    vec![schema_step("define attribute a, value string;", None)],
660                ),
661                apply_migration(
662                    "0002_skipped",
663                    vec![schema_step("define attribute b, value string;", None)],
664                ),
665            ],
666            to_rollback: Vec::new(),
667        };
668
669        let results = execute_plan(&db, plan).await;
670
671        // Only the first migration result is returned; the second is never attempted.
672        assert_eq!(results.len(), 1);
673        assert!(!results[0].success);
674        assert_eq!(results[0].name, "0001_fail");
675    }
676}