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