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