1use 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
32pub type RecoveryFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
34
35#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
37#[serde(transparent)]
38pub struct ExecutionStepId(String);
39
40impl ExecutionStepId {
41 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
55pub struct CheckedMigrationIdentity {
56 pub app_label: String,
58 pub name: String,
60 pub checksum: String,
62 pub action: MigrationAction,
64}
65
66#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
68pub struct CheckedExecutionStep {
69 pub id: ExecutionStepId,
71 pub migration: CheckedMigrationIdentity,
73 pub step_index: usize,
75 pub operation_kind: OperationKind,
77 pub execution: ExecutionStep,
79}
80
81#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
83pub struct CheckedMigrationExecution {
84 pub identity: CheckedMigrationIdentity,
86 pub steps: Vec<CheckedExecutionStep>,
88 pub reversible: bool,
90}
91
92#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
94pub struct CheckedExecutionPlan {
95 pub to_apply: Vec<CheckedMigrationExecution>,
97 pub to_rollback: Vec<CheckedMigrationExecution>,
99}
100
101impl CheckedExecutionPlan {
102 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
111pub 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
137pub 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#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
226#[serde(tag = "kind", rename_all = "snake_case")]
227pub enum PendingProof {
228 NotCommitted,
230 IdempotentReplay {
232 strategy: String,
234 },
235}
236
237#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
239#[serde(tag = "status", rename_all = "snake_case")]
240pub enum StepRecoveryDecision {
241 Pending {
243 proof: PendingProof,
245 },
246 Applied {
248 #[serde(default, skip_serializing_if = "Option::is_none")]
250 evidence: Option<String>,
251 },
252 Indeterminate {
254 reason: String,
256 },
257}
258
259#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
261#[serde(rename_all = "snake_case")]
262pub enum StepRecoveryEventKind {
263 BeforeCommit,
265 Committed,
267 FailedBeforeCommit,
269 UnknownCommitOutcome,
271}
272
273#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
275pub struct StepRecoveryEvent {
276 pub step: CheckedExecutionStep,
278 pub kind: StepRecoveryEventKind,
280 #[serde(default, skip_serializing_if = "Option::is_none")]
282 pub message: Option<String>,
283 #[serde(default, skip_serializing_if = "Option::is_none")]
285 pub backfill: Option<BackfillResult>,
286}
287
288pub trait StepRecoveryController: Send + Sync {
295 fn classify<'a>(
297 &'a self,
298 db: &'a Database,
299 step: &'a CheckedExecutionStep,
300 ) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>>;
301
302 fn record_event<'a>(
304 &'a self,
305 event: StepRecoveryEvent,
306 ) -> RecoveryFuture<'a, crate::Result<()>>;
307}
308
309#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
311pub struct StepExecutionResult {
312 pub step: CheckedExecutionStep,
314 pub outcome: StepExecutionOutcome,
316}
317
318#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
320#[serde(tag = "status", rename_all = "snake_case")]
321pub enum StepExecutionOutcome {
322 Applied {
324 #[serde(default, skip_serializing_if = "Option::is_none")]
326 evidence: Option<String>,
327 },
328 Committed {
330 #[serde(default, skip_serializing_if = "Option::is_none")]
332 backfill: Option<BackfillResult>,
333 },
334 FailedBeforeCommit {
336 error: String,
338 },
339 Indeterminate {
341 error: String,
343 },
344}
345
346#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
348#[serde(rename_all = "snake_case")]
349pub enum RecoveryMigrationStatus {
350 Succeeded,
352 Failed,
354 Indeterminate,
356}
357
358#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
360pub struct RecoveryMigrationResult {
361 pub migration: CheckedMigrationIdentity,
363 pub status: RecoveryMigrationStatus,
365 pub steps: Vec<StepExecutionResult>,
367 #[serde(default, skip_serializing_if = "Option::is_none")]
369 pub error: Option<String>,
370}
371
372#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
374#[serde(rename_all = "snake_case")]
375pub enum RecoveryPlanStatus {
376 Succeeded,
378 Failed,
380 Indeterminate,
382}
383
384#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
386pub struct RecoveryExecutionResult {
387 pub status: RecoveryPlanStatus,
389 pub migrations: Vec<RecoveryMigrationResult>,
391}
392
393pub 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}