1use 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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
29pub struct MigrationResult {
30 pub app_label: String,
32 pub name: String,
34 pub action: MigrationAction,
36 pub success: bool,
38 pub error: Option<String>,
40 #[serde(default, skip_serializing_if = "Option::is_none")]
44 pub backfill: Option<Vec<BackfillResult>>,
45}
46
47pub async fn execute_plan(db: &Database, plan: ExecutionPlan) -> Vec<MigrationResult> {
56 let mut results: Vec<MigrationResult> = Vec::new();
57
58 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 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
81pub 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
137pub async fn execute_migration(db: &Database, migration: &MigrationExecution) -> MigrationResult {
141 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 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 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 let typeql: &str = match migration.action {
179 MigrationAction::Apply => &step.forward,
180 MigrationAction::Rollback => {
181 step.reverse.as_deref().unwrap_or(&step.forward)
185 }
186 };
187
188 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 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 if let Err(e) = ctx.query(typeql).await {
219 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 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 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 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 #[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 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 #[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 MockEvent::Query(TxType::Schema, "undefine attribute b;".to_string()),
430 MockEvent::Commit,
431 ]
432 );
433 }
434
435 #[tokio::test]
445 async fn middle_step_failure_commits_prior_steps_and_halts() {
446 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 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 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 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 assert_eq!(events.len(), 6, "step 3 must not be opened");
487 }
488
489 #[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 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 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 #[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 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 #[tokio::test]
563 async fn failure_in_first_migration_halts_remaining() {
564 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 assert_eq!(results.len(), 1);
586 assert!(!results[0].success);
587 assert_eq!(results[0].name, "0001_fail");
588 }
589}