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;
31use crate::state::{require_legacy_writer_open, require_legacy_writer_open_in_transaction};
32
33pub type RecoveryFuture<'a, T> = Pin<Box<dyn Future<Output = T> + Send + 'a>>;
35
36#[derive(Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, Serialize, Deserialize)]
38#[serde(transparent)]
39pub struct ExecutionStepId(String);
40
41impl ExecutionStepId {
42 pub fn as_str(&self) -> &str {
44 &self.0
45 }
46}
47
48impl fmt::Display for ExecutionStepId {
49 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
50 formatter.write_str(&self.0)
51 }
52}
53
54#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
56pub struct CheckedMigrationIdentity {
57 pub app_label: String,
59 pub name: String,
61 pub checksum: String,
63 pub action: MigrationAction,
65}
66
67#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
69pub struct CheckedExecutionStep {
70 pub id: ExecutionStepId,
72 pub migration: CheckedMigrationIdentity,
74 pub step_index: usize,
76 pub operation_kind: OperationKind,
78 pub execution: ExecutionStep,
80}
81
82#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
84pub struct CheckedMigrationExecution {
85 pub identity: CheckedMigrationIdentity,
87 pub steps: Vec<CheckedExecutionStep>,
89 pub reversible: bool,
91}
92
93#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
95pub struct CheckedExecutionPlan {
96 pub to_apply: Vec<CheckedMigrationExecution>,
98 pub to_rollback: Vec<CheckedMigrationExecution>,
100}
101
102impl CheckedExecutionPlan {
103 pub fn ordered_steps(&self) -> impl Iterator<Item = &CheckedExecutionStep> {
105 self.to_rollback
106 .iter()
107 .chain(self.to_apply.iter())
108 .flat_map(|migration| migration.steps.iter())
109 }
110}
111
112pub fn plan_recovery(
118 graph: &MigrationGraph,
119 applied: &[AppliedMigrationRecord],
120 target: Option<&str>,
121) -> crate::Result<CheckedExecutionPlan> {
122 let execution_plan = plan(graph, applied, target)?;
123 let checksums = graph
124 .migrations
125 .iter()
126 .filter_map(|migration| {
127 migration.checksum.as_ref().map(|checksum| {
128 (
129 (migration.app_label.clone(), migration.name.clone()),
130 checksum.clone(),
131 )
132 })
133 })
134 .collect();
135 prepare_recovery_plan(&execution_plan, &checksums)
136}
137
138pub fn prepare_recovery_plan(
144 plan: &ExecutionPlan,
145 checksums: &BTreeMap<(String, String), String>,
146) -> crate::Result<CheckedExecutionPlan> {
147 Ok(CheckedExecutionPlan {
148 to_apply: prepare_migrations(&plan.to_apply, checksums)?,
149 to_rollback: prepare_migrations(&plan.to_rollback, checksums)?,
150 })
151}
152
153fn prepare_migrations(
154 migrations: &[MigrationExecution],
155 checksums: &BTreeMap<(String, String), String>,
156) -> crate::Result<Vec<CheckedMigrationExecution>> {
157 migrations
158 .iter()
159 .map(|migration| {
160 let checksum = checksums
161 .get(&(migration.app_label.clone(), migration.name.clone()))
162 .filter(|checksum| !checksum.is_empty())
163 .ok_or_else(|| MigrationError::MissingRecoveryChecksum {
164 app_label: migration.app_label.clone(),
165 name: migration.name.clone(),
166 })?
167 .clone();
168 let identity = CheckedMigrationIdentity {
169 app_label: migration.app_label.clone(),
170 name: migration.name.clone(),
171 checksum,
172 action: migration.action,
173 };
174 let steps = migration
175 .steps
176 .iter()
177 .enumerate()
178 .map(|(step_index, execution)| CheckedExecutionStep {
179 id: step_id(&identity, step_index, execution.operation_kind),
180 migration: identity.clone(),
181 step_index,
182 operation_kind: execution.operation_kind,
183 execution: execution.clone(),
184 })
185 .collect();
186 Ok(CheckedMigrationExecution {
187 identity,
188 steps,
189 reversible: migration.reversible,
190 })
191 })
192 .collect()
193}
194
195fn step_id(
196 migration: &CheckedMigrationIdentity,
197 step_index: usize,
198 operation_kind: OperationKind,
199) -> ExecutionStepId {
200 let mut digest = Sha256::new();
201 hash_field(&mut digest, b"type-bridge-execution-step-v1");
202 hash_field(&mut digest, migration.checksum.as_bytes());
203 hash_field(&mut digest, migration.app_label.as_bytes());
204 hash_field(&mut digest, migration.name.as_bytes());
205 hash_field(
206 &mut digest,
207 match migration.action {
208 MigrationAction::Apply => b"apply",
209 MigrationAction::Rollback => b"rollback",
210 },
211 );
212 hash_field(
213 &mut digest,
214 &u64::try_from(step_index).unwrap_or(u64::MAX).to_be_bytes(),
215 );
216 hash_field(&mut digest, operation_kind.as_str().as_bytes());
217 ExecutionStepId(format!("tb-step-v1:{:x}", digest.finalize()))
218}
219
220fn hash_field(digest: &mut Sha256, value: &[u8]) {
221 digest.update(u64::try_from(value.len()).unwrap_or(u64::MAX).to_be_bytes());
222 digest.update(value);
223}
224
225#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
227#[serde(tag = "kind", rename_all = "snake_case")]
228pub enum PendingProof {
229 NotCommitted,
231 IdempotentReplay {
233 strategy: String,
235 },
236}
237
238#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
240#[serde(tag = "status", rename_all = "snake_case")]
241pub enum StepRecoveryDecision {
242 Pending {
244 proof: PendingProof,
246 },
247 Applied {
249 #[serde(default, skip_serializing_if = "Option::is_none")]
251 evidence: Option<String>,
252 },
253 Indeterminate {
255 reason: String,
257 },
258}
259
260#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
262#[serde(rename_all = "snake_case")]
263pub enum StepRecoveryEventKind {
264 BeforeCommit,
266 Committed,
268 FailedBeforeCommit,
270 UnknownCommitOutcome,
272}
273
274#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
276pub struct StepRecoveryEvent {
277 pub step: CheckedExecutionStep,
279 pub kind: StepRecoveryEventKind,
281 #[serde(default, skip_serializing_if = "Option::is_none")]
283 pub message: Option<String>,
284 #[serde(default, skip_serializing_if = "Option::is_none")]
286 pub backfill: Option<BackfillResult>,
287}
288
289pub trait StepRecoveryController: Send + Sync {
296 fn classify<'a>(
298 &'a self,
299 db: &'a Database,
300 step: &'a CheckedExecutionStep,
301 ) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>>;
302
303 fn record_event<'a>(
305 &'a self,
306 event: StepRecoveryEvent,
307 ) -> RecoveryFuture<'a, crate::Result<()>>;
308}
309
310#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
312pub struct StepExecutionResult {
313 pub step: CheckedExecutionStep,
315 pub outcome: StepExecutionOutcome,
317}
318
319#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
321#[serde(tag = "status", rename_all = "snake_case")]
322pub enum StepExecutionOutcome {
323 Applied {
325 #[serde(default, skip_serializing_if = "Option::is_none")]
327 evidence: Option<String>,
328 },
329 Committed {
331 #[serde(default, skip_serializing_if = "Option::is_none")]
333 backfill: Option<BackfillResult>,
334 },
335 FailedBeforeCommit {
337 error: String,
339 },
340 Indeterminate {
342 error: String,
344 },
345}
346
347#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
349#[serde(rename_all = "snake_case")]
350pub enum RecoveryMigrationStatus {
351 Succeeded,
353 Failed,
355 Indeterminate,
357}
358
359#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
361pub struct RecoveryMigrationResult {
362 pub migration: CheckedMigrationIdentity,
364 pub status: RecoveryMigrationStatus,
366 pub steps: Vec<StepExecutionResult>,
368 #[serde(default, skip_serializing_if = "Option::is_none")]
370 pub error: Option<String>,
371}
372
373#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
375#[serde(rename_all = "snake_case")]
376pub enum RecoveryPlanStatus {
377 Succeeded,
379 Failed,
381 Indeterminate,
383}
384
385#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
387pub struct RecoveryExecutionResult {
388 pub status: RecoveryPlanStatus,
390 pub migrations: Vec<RecoveryMigrationResult>,
392}
393
394pub async fn execute_recovery_plan<C: StepRecoveryController + ?Sized>(
399 db: &Database,
400 plan: &CheckedExecutionPlan,
401 controller: &C,
402) -> RecoveryExecutionResult {
403 let mut migrations = Vec::new();
404 for migration in plan.to_rollback.iter().chain(plan.to_apply.iter()) {
405 let result = execute_recovery_migration(db, migration, controller).await;
406 let terminal = result.status;
407 migrations.push(result);
408 match terminal {
409 RecoveryMigrationStatus::Succeeded => {}
410 RecoveryMigrationStatus::Failed => {
411 return RecoveryExecutionResult {
412 status: RecoveryPlanStatus::Failed,
413 migrations,
414 };
415 }
416 RecoveryMigrationStatus::Indeterminate => {
417 return RecoveryExecutionResult {
418 status: RecoveryPlanStatus::Indeterminate,
419 migrations,
420 };
421 }
422 }
423 }
424 RecoveryExecutionResult {
425 status: RecoveryPlanStatus::Succeeded,
426 migrations,
427 }
428}
429
430async fn execute_recovery_migration<C: StepRecoveryController + ?Sized>(
431 db: &Database,
432 migration: &CheckedMigrationExecution,
433 controller: &C,
434) -> RecoveryMigrationResult {
435 if migration.identity.action == MigrationAction::Rollback && !migration.reversible {
436 return RecoveryMigrationResult {
437 migration: migration.identity.clone(),
438 status: RecoveryMigrationStatus::Failed,
439 steps: Vec::new(),
440 error: Some(format!("{} is not reversible", migration.identity.name)),
441 };
442 }
443 if let Err(error) = require_legacy_writer_open(db).await {
444 return RecoveryMigrationResult {
445 migration: migration.identity.clone(),
446 status: RecoveryMigrationStatus::Failed,
447 steps: Vec::new(),
448 error: Some(error.to_string()),
449 };
450 }
451
452 let mut results = Vec::new();
453 for step in &migration.steps {
454 let decision = match controller.classify(db, step).await {
455 Ok(decision) => decision,
456 Err(error) => {
457 let message = format!("failed to classify step {}: {error}", step.id);
458 results.push(step_result(
459 step,
460 StepExecutionOutcome::Indeterminate {
461 error: message.clone(),
462 },
463 ));
464 return migration_result(
465 migration,
466 RecoveryMigrationStatus::Indeterminate,
467 results,
468 message,
469 );
470 }
471 };
472
473 match decision {
474 StepRecoveryDecision::Applied { evidence } => {
475 results.push(step_result(
476 step,
477 StepExecutionOutcome::Applied { evidence },
478 ));
479 }
480 StepRecoveryDecision::Indeterminate { reason } => {
481 results.push(step_result(
482 step,
483 StepExecutionOutcome::Indeterminate {
484 error: reason.clone(),
485 },
486 ));
487 return migration_result(
488 migration,
489 RecoveryMigrationStatus::Indeterminate,
490 results,
491 reason,
492 );
493 }
494 StepRecoveryDecision::Pending { proof: _ } => {
495 let outcome = execute_pending_step(db, step, controller).await;
496 let terminal = match &outcome {
497 StepExecutionOutcome::Applied { .. }
498 | StepExecutionOutcome::Committed { .. } => None,
499 StepExecutionOutcome::FailedBeforeCommit { error } => {
500 Some((RecoveryMigrationStatus::Failed, error.clone()))
501 }
502 StepExecutionOutcome::Indeterminate { error } => {
503 Some((RecoveryMigrationStatus::Indeterminate, error.clone()))
504 }
505 };
506 results.push(step_result(step, outcome));
507 if let Some((status, error)) = terminal {
508 return migration_result(migration, status, results, error);
509 }
510 }
511 }
512 }
513
514 RecoveryMigrationResult {
515 migration: migration.identity.clone(),
516 status: RecoveryMigrationStatus::Succeeded,
517 steps: results,
518 error: None,
519 }
520}
521
522async fn execute_pending_step<C: StepRecoveryController + ?Sized>(
523 db: &Database,
524 step: &CheckedExecutionStep,
525 controller: &C,
526) -> StepExecutionOutcome {
527 if step.execution.kind == StepKind::Backfill && step.migration.action == MigrationAction::Apply
528 {
529 return execute_pending_backfill(db, step, controller).await;
530 }
531
532 let typeql = match step.migration.action {
533 MigrationAction::Apply => step.execution.forward.as_str(),
534 MigrationAction::Rollback => match step.execution.reverse.as_deref() {
535 Some(reverse) => reverse,
536 None => {
537 return failed_before_commit(
538 controller,
539 step,
540 "rollback step has no reverse TypeQL".to_string(),
541 None,
542 )
543 .await;
544 }
545 },
546 };
547
548 if let Err(error) = db.check_schema_annotation_support(typeql) {
549 return failed_before_commit(controller, step, error.to_string(), None).await;
550 }
551 let transaction = match db.transaction_context(step.execution.tx_type).await {
552 Ok(transaction) => transaction,
553 Err(error) => {
554 return failed_before_commit(
555 controller,
556 step,
557 format!("failed to open transaction: {error}"),
558 None,
559 )
560 .await;
561 }
562 };
563 if let Err(error) = require_legacy_writer_open_in_transaction(&transaction).await {
564 let _ = transaction.rollback().await;
565 return StepExecutionOutcome::FailedBeforeCommit {
566 error: error.to_string(),
567 };
568 }
569 if let Err(error) = transaction.query(typeql).await {
570 let _ = transaction.rollback().await;
571 return failed_before_commit(controller, step, format!("query failed: {error}"), None)
572 .await;
573 }
574
575 if let Err(error) = controller
576 .record_event(event(step, StepRecoveryEventKind::BeforeCommit, None, None))
577 .await
578 {
579 let _ = transaction.rollback().await;
580 return failed_before_commit(
581 controller,
582 step,
583 format!("before-commit event was not recorded: {error}"),
584 None,
585 )
586 .await;
587 }
588
589 if let Err(error) = transaction.commit().await {
590 let message = format!("commit outcome is unknown: {error}");
591 return unknown_commit(controller, step, message, None).await;
592 }
593 committed(controller, step, None).await
594}
595
596async fn execute_pending_backfill<C: StepRecoveryController + ?Sized>(
597 db: &Database,
598 step: &CheckedExecutionStep,
599 controller: &C,
600) -> StepExecutionOutcome {
601 let prepared = match prepare_backfill(db, &step.execution, step.step_index).await {
602 Ok(prepared) => prepared,
603 Err(error) => {
604 return failed_before_commit(controller, step, error.to_string(), None).await;
605 }
606 };
607 let backfill = prepared.result.clone();
608 if let Err(error) = controller
609 .record_event(event(
610 step,
611 StepRecoveryEventKind::BeforeCommit,
612 None,
613 Some(backfill.clone()),
614 ))
615 .await
616 {
617 let _ = prepared.transaction.rollback().await;
618 return failed_before_commit(
619 controller,
620 step,
621 format!("before-commit event was not recorded: {error}"),
622 Some(backfill),
623 )
624 .await;
625 }
626 if let Err(error) = prepared.transaction.commit().await {
627 let message = format!("backfill commit outcome is unknown: {error}");
628 return unknown_commit(controller, step, message, Some(backfill)).await;
629 }
630 committed(controller, step, Some(backfill)).await
631}
632
633async fn committed<C: StepRecoveryController + ?Sized>(
634 controller: &C,
635 step: &CheckedExecutionStep,
636 backfill: Option<BackfillResult>,
637) -> StepExecutionOutcome {
638 if let Err(error) = controller
639 .record_event(event(
640 step,
641 StepRecoveryEventKind::Committed,
642 None,
643 backfill.clone(),
644 ))
645 .await
646 {
647 return StepExecutionOutcome::Indeterminate {
648 error: format!(
649 "TypeDB committed step {}, but its committed event was not durably recorded: {error}",
650 step.id
651 ),
652 };
653 }
654 StepExecutionOutcome::Committed { backfill }
655}
656
657async fn failed_before_commit<C: StepRecoveryController + ?Sized>(
658 controller: &C,
659 step: &CheckedExecutionStep,
660 mut message: String,
661 backfill: Option<BackfillResult>,
662) -> StepExecutionOutcome {
663 if let Err(event_error) = controller
664 .record_event(event(
665 step,
666 StepRecoveryEventKind::FailedBeforeCommit,
667 Some(message.clone()),
668 backfill,
669 ))
670 .await
671 {
672 message.push_str(&format!(
673 "; failed to record failed-before-commit event: {event_error}"
674 ));
675 }
676 StepExecutionOutcome::FailedBeforeCommit { error: message }
677}
678
679async fn unknown_commit<C: StepRecoveryController + ?Sized>(
680 controller: &C,
681 step: &CheckedExecutionStep,
682 mut message: String,
683 backfill: Option<BackfillResult>,
684) -> StepExecutionOutcome {
685 if let Err(event_error) = controller
686 .record_event(event(
687 step,
688 StepRecoveryEventKind::UnknownCommitOutcome,
689 Some(message.clone()),
690 backfill,
691 ))
692 .await
693 {
694 message.push_str(&format!(
695 "; failed to record unknown-commit event: {event_error}"
696 ));
697 }
698 StepExecutionOutcome::Indeterminate { error: message }
699}
700
701fn event(
702 step: &CheckedExecutionStep,
703 kind: StepRecoveryEventKind,
704 message: Option<String>,
705 backfill: Option<BackfillResult>,
706) -> StepRecoveryEvent {
707 StepRecoveryEvent {
708 step: step.clone(),
709 kind,
710 message,
711 backfill,
712 }
713}
714
715fn step_result(step: &CheckedExecutionStep, outcome: StepExecutionOutcome) -> StepExecutionResult {
716 StepExecutionResult {
717 step: step.clone(),
718 outcome,
719 }
720}
721
722fn migration_result(
723 migration: &CheckedMigrationExecution,
724 status: RecoveryMigrationStatus,
725 steps: Vec<StepExecutionResult>,
726 error: String,
727) -> RecoveryMigrationResult {
728 RecoveryMigrationResult {
729 migration: migration.identity.clone(),
730 status,
731 steps,
732 error: Some(error),
733 }
734}
735
736#[cfg(test)]
737mod tests {
738 use std::collections::{BTreeMap, VecDeque};
739 use std::sync::Mutex;
740
741 use serde_json::json;
742 use type_bridge_orm::session::backend::QueryResult;
743 use type_bridge_orm::{Database, TxType};
744
745 use super::*;
746 use crate::spec::{MigrationGraph, MigrationSpec, OperationSpec};
747 use crate::testing::{MockEvent, MockMigrationBackend};
748
749 struct TestController {
750 decisions: Mutex<VecDeque<StepRecoveryDecision>>,
751 events: Mutex<Vec<StepRecoveryEvent>>,
752 fail_event: Option<StepRecoveryEventKind>,
753 }
754
755 impl TestController {
756 fn new(decisions: Vec<StepRecoveryDecision>) -> Self {
757 Self {
758 decisions: Mutex::new(decisions.into()),
759 events: Mutex::new(Vec::new()),
760 fail_event: None,
761 }
762 }
763
764 fn failing_event(
765 decisions: Vec<StepRecoveryDecision>,
766 kind: StepRecoveryEventKind,
767 ) -> Self {
768 Self {
769 decisions: Mutex::new(decisions.into()),
770 events: Mutex::new(Vec::new()),
771 fail_event: Some(kind),
772 }
773 }
774
775 fn event_kinds(&self) -> Vec<StepRecoveryEventKind> {
776 self.events
777 .lock()
778 .unwrap()
779 .iter()
780 .map(|event| event.kind)
781 .collect()
782 }
783 }
784
785 impl StepRecoveryController for TestController {
786 fn classify<'a>(
787 &'a self,
788 _db: &'a Database,
789 _step: &'a CheckedExecutionStep,
790 ) -> RecoveryFuture<'a, crate::Result<StepRecoveryDecision>> {
791 let decision = self.decisions.lock().unwrap().pop_front();
792 Box::pin(async move {
793 decision.ok_or_else(|| MigrationError::Recovery {
794 message: "test controller has no decision".to_string(),
795 })
796 })
797 }
798
799 fn record_event<'a>(
800 &'a self,
801 event: StepRecoveryEvent,
802 ) -> RecoveryFuture<'a, crate::Result<()>> {
803 let should_fail = self.fail_event == Some(event.kind);
804 self.events.lock().unwrap().push(event);
805 Box::pin(async move {
806 if should_fail {
807 Err(MigrationError::Recovery {
808 message: "injected event persistence failure".to_string(),
809 })
810 } else {
811 Ok(())
812 }
813 })
814 }
815 }
816
817 fn pending() -> StepRecoveryDecision {
818 StepRecoveryDecision::Pending {
819 proof: PendingProof::NotCommitted,
820 }
821 }
822
823 fn applied(evidence: &str) -> StepRecoveryDecision {
824 StepRecoveryDecision::Applied {
825 evidence: Some(evidence.to_string()),
826 }
827 }
828
829 fn execution_step(
830 tx_type: TxType,
831 kind: StepKind,
832 operation_kind: OperationKind,
833 forward: &str,
834 ) -> ExecutionStep {
835 ExecutionStep {
836 tx_type,
837 kind,
838 operation_kind,
839 forward: forward.to_string(),
840 reverse: Some(format!("reverse {forward}")),
841 }
842 }
843
844 fn schema_step(forward: &str) -> ExecutionStep {
845 execution_step(
846 TxType::Schema,
847 StepKind::Schema,
848 OperationKind::AddAttribute,
849 forward,
850 )
851 }
852
853 fn write_step(forward: &str) -> ExecutionStep {
854 execution_step(
855 TxType::Write,
856 StepKind::Write,
857 OperationKind::RunTypeql,
858 forward,
859 )
860 }
861
862 fn backfill_step() -> ExecutionStep {
863 execution_step(
864 TxType::Write,
865 StepKind::Backfill,
866 OperationKind::CopyAttribute,
867 "match\n $x isa person, has old-name $v;\n not { $x has new-name $d; };\ninsert\n $x has new-name == $v;",
868 )
869 }
870
871 fn checked_plan(steps: Vec<ExecutionStep>) -> CheckedExecutionPlan {
872 let plan = ExecutionPlan {
873 to_apply: vec![MigrationExecution {
874 app_label: "app".to_string(),
875 name: "0001_initial".to_string(),
876 action: MigrationAction::Apply,
877 reversible: steps.iter().all(|step| step.reverse.is_some()),
878 steps,
879 }],
880 to_rollback: Vec::new(),
881 };
882 let checksums = BTreeMap::from([(
883 ("app".to_string(), "0001_initial".to_string()),
884 "checksum-1".to_string(),
885 )]);
886 prepare_recovery_plan(&plan, &checksums).unwrap()
887 }
888
889 #[test]
890 fn checked_plan_exposes_stable_checksum_bound_step_sequence() {
891 let steps = vec![
892 schema_step("define attribute a, value string;"),
893 write_step("insert $p isa person;"),
894 ];
895 let first = checked_plan(steps.clone());
896 let second = checked_plan(steps);
897 let first_ids: Vec<_> = first.ordered_steps().map(|step| step.id.clone()).collect();
898 let second_ids: Vec<_> = second.ordered_steps().map(|step| step.id.clone()).collect();
899
900 assert_eq!(first_ids, second_ids);
901 assert_ne!(first_ids[0], first_ids[1]);
902 assert_eq!(
903 first_ids[0].as_str(),
904 "tb-step-v1:0a308e166ea81d161a8c92cbf01157935fbf53b55111df3c75f5129a183b74f9"
905 );
906 assert_eq!(first_ids[0].as_str().len(), 75);
907 assert_eq!(
908 serde_json::to_value(&first).unwrap()["to_apply"][0]["steps"]
909 .as_array()
910 .unwrap()
911 .len(),
912 2
913 );
914
915 let mut different_checksums = BTreeMap::from([(
916 ("app".to_string(), "0001_initial".to_string()),
917 "checksum-2".to_string(),
918 )]);
919 let raw = ExecutionPlan {
920 to_apply: vec![MigrationExecution {
921 app_label: "app".to_string(),
922 name: "0001_initial".to_string(),
923 action: MigrationAction::Apply,
924 steps: vec![schema_step("define attribute a, value string;")],
925 reversible: true,
926 }],
927 to_rollback: Vec::new(),
928 };
929 let changed = prepare_recovery_plan(&raw, &different_checksums).unwrap();
930 assert_ne!(first_ids[0], changed.ordered_steps().next().unwrap().id);
931 different_checksums.clear();
932 assert!(matches!(
933 prepare_recovery_plan(&raw, &different_checksums),
934 Err(MigrationError::MissingRecoveryChecksum { .. })
935 ));
936 }
937
938 #[test]
939 fn plan_recovery_binds_ids_to_the_validated_graph_artifact() {
940 let graph = MigrationGraph {
941 migrations: vec![MigrationSpec {
942 app_label: "app".to_string(),
943 name: "0001_initial".to_string(),
944 dependencies: Vec::new(),
945 operations: vec![OperationSpec::RunTypeql {
946 forward: "define attribute a, value string;".to_string(),
947 reverse: None,
948 }],
949 checksum: Some("artifact-checksum".to_string()),
950 source_sha256: None,
951 reversible: false,
952 }],
953 };
954
955 let checked = plan_recovery(&graph, &[], None).unwrap();
956 let step = checked.ordered_steps().next().unwrap();
957
958 assert_eq!(step.migration.checksum, "artifact-checksum");
959 assert_eq!(step.operation_kind, OperationKind::RunTypeql);
960 assert_eq!(checked.ordered_steps().count(), 1);
961 }
962
963 #[tokio::test]
964 async fn safe_resume_skips_applied_schema_and_executes_proven_pending_data() {
965 let plan = checked_plan(vec![
966 schema_step("define attribute a, value string;"),
967 write_step("insert $p isa person;"),
968 ]);
969 let controller = TestController::new(vec![applied("schema introspection"), pending()]);
970 let (backend, log) = MockMigrationBackend::new(None);
971 let db = Database::with_backend(Box::new(backend), "test");
972
973 let result = execute_recovery_plan(&db, &plan, &controller).await;
974
975 assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
976 assert!(matches!(
977 result.migrations[0].steps[0].outcome,
978 StepExecutionOutcome::Applied { .. }
979 ));
980 assert!(matches!(
981 result.migrations[0].steps[1].outcome,
982 StepExecutionOutcome::Committed { .. }
983 ));
984 assert_eq!(
985 controller.event_kinds(),
986 vec![
987 StepRecoveryEventKind::BeforeCommit,
988 StepRecoveryEventKind::Committed
989 ]
990 );
991 assert_eq!(
992 *log.lock().unwrap(),
993 vec![
994 MockEvent::OpenTx(TxType::Read),
995 MockEvent::Close,
996 MockEvent::OpenTx(TxType::Write),
997 MockEvent::Query(TxType::Write, "insert $p isa person;".to_string()),
998 MockEvent::Commit,
999 ]
1000 );
1001 }
1002
1003 #[tokio::test]
1004 async fn cutover_rejects_recovery_before_controller_or_typeql() {
1005 let plan = checked_plan(vec![write_step("insert $p isa person;")]);
1006 let controller = TestController::new(vec![pending()]);
1007 let (backend, log) = MockMigrationBackend::with_legacy_cutover();
1008 let db = Database::with_backend(Box::new(backend), "test");
1009
1010 let result = execute_recovery_plan(&db, &plan, &controller).await;
1011
1012 assert_eq!(result.status, RecoveryPlanStatus::Failed);
1013 assert!(result.migrations[0].steps.is_empty());
1014 assert!(
1015 result.migrations[0]
1016 .error
1017 .as_deref()
1018 .is_some_and(|error| error.contains(crate::LEGACY_WRITER_CUTOVER_MESSAGE))
1019 );
1020 assert!(controller.event_kinds().is_empty());
1021 assert_eq!(
1022 *log.lock().unwrap(),
1023 vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1024 );
1025 }
1026
1027 #[tokio::test]
1028 async fn indeterminate_run_typeql_blocks_without_replay() {
1029 let plan = checked_plan(vec![write_step("insert $p isa person;")]);
1030 let controller = TestController::new(vec![StepRecoveryDecision::Indeterminate {
1031 reason: "prior before-commit receipt has no matching outcome".to_string(),
1032 }]);
1033 let (backend, log) = MockMigrationBackend::new(None);
1034 let db = Database::with_backend(Box::new(backend), "test");
1035
1036 let result = execute_recovery_plan(&db, &plan, &controller).await;
1037
1038 assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1039 assert!(matches!(
1040 result.migrations[0].steps[0].outcome,
1041 StepExecutionOutcome::Indeterminate { .. }
1042 ));
1043 assert_eq!(
1044 *log.lock().unwrap(),
1045 vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1046 );
1047 assert!(controller.event_kinds().is_empty());
1048 }
1049
1050 #[tokio::test]
1051 async fn failure_after_earlier_commit_is_typed_failed_before_commit() {
1052 let plan = checked_plan(vec![
1053 schema_step("define attribute a, value string;"),
1054 schema_step("define attribute b, value string;"),
1055 ]);
1056 let controller = TestController::new(vec![pending(), pending()]);
1057 let (backend, _log) = MockMigrationBackend::new(Some(1));
1058 let db = Database::with_backend(Box::new(backend), "test");
1059
1060 let result = execute_recovery_plan(&db, &plan, &controller).await;
1061
1062 assert_eq!(result.status, RecoveryPlanStatus::Failed);
1063 assert!(matches!(
1064 result.migrations[0].steps[0].outcome,
1065 StepExecutionOutcome::Committed { .. }
1066 ));
1067 assert!(matches!(
1068 result.migrations[0].steps[1].outcome,
1069 StepExecutionOutcome::FailedBeforeCommit { .. }
1070 ));
1071 assert_eq!(
1072 controller.event_kinds(),
1073 vec![
1074 StepRecoveryEventKind::BeforeCommit,
1075 StepRecoveryEventKind::Committed,
1076 StepRecoveryEventKind::FailedBeforeCommit,
1077 ]
1078 );
1079 }
1080
1081 #[tokio::test]
1082 async fn definitely_aborted_commit_response_remains_indeterminate_in_v1_recovery() {
1083 let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
1084 let controller = TestController::new(vec![pending()]);
1085 let (backend, _log) = MockMigrationBackend::with_definitely_aborted_commit_failure(0);
1086 let db = Database::with_backend(Box::new(backend), "test");
1087
1088 let result = execute_recovery_plan(&db, &plan, &controller).await;
1089
1090 assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1091 assert!(matches!(
1092 result.migrations[0].steps[0].outcome,
1093 StepExecutionOutcome::Indeterminate { .. }
1094 ));
1095 assert_eq!(
1096 controller.event_kinds(),
1097 vec![
1098 StepRecoveryEventKind::BeforeCommit,
1099 StepRecoveryEventKind::UnknownCommitOutcome,
1100 ]
1101 );
1102 }
1103
1104 #[tokio::test]
1105 async fn ambiguous_commit_response_is_indeterminate() {
1106 let plan = checked_plan(vec![schema_step("define attribute a, value string;")]);
1107 let controller = TestController::new(vec![pending()]);
1108 let (backend, _log) = MockMigrationBackend::with_commit_failure(0);
1109 let db = Database::with_backend(Box::new(backend), "test");
1110
1111 let result = execute_recovery_plan(&db, &plan, &controller).await;
1112
1113 assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1114 assert!(matches!(
1115 result.migrations[0].steps[0].outcome,
1116 StepExecutionOutcome::Indeterminate { .. }
1117 ));
1118 assert_eq!(
1119 controller.event_kinds(),
1120 vec![
1121 StepRecoveryEventKind::BeforeCommit,
1122 StepRecoveryEventKind::UnknownCommitOutcome,
1123 ]
1124 );
1125 }
1126
1127 #[tokio::test]
1128 async fn lost_committed_event_delivery_is_indeterminate_and_halts() {
1129 let plan = checked_plan(vec![
1130 schema_step("define attribute a, value string;"),
1131 schema_step("define attribute b, value string;"),
1132 ]);
1133 let controller = TestController::failing_event(
1134 vec![pending(), pending()],
1135 StepRecoveryEventKind::Committed,
1136 );
1137 let (backend, log) = MockMigrationBackend::new(None);
1138 let db = Database::with_backend(Box::new(backend), "test");
1139
1140 let result = execute_recovery_plan(&db, &plan, &controller).await;
1141
1142 assert_eq!(result.status, RecoveryPlanStatus::Indeterminate);
1143 assert_eq!(result.migrations[0].steps.len(), 1);
1144 assert!(matches!(
1145 result.migrations[0].steps[0].outcome,
1146 StepExecutionOutcome::Indeterminate { .. }
1147 ));
1148 assert_eq!(
1149 log.lock()
1150 .unwrap()
1151 .iter()
1152 .filter(|event| matches!(event, MockEvent::Commit))
1153 .count(),
1154 1
1155 );
1156 }
1157
1158 #[tokio::test]
1159 async fn final_external_applied_record_retry_skips_all_proven_steps() {
1160 let plan = checked_plan(vec![
1161 schema_step("define attribute a, value string;"),
1162 write_step("insert $p isa person;"),
1163 ]);
1164 let controller = TestController::new(vec![
1165 applied("durable step receipt"),
1166 applied("operation-specific reconciliation"),
1167 ]);
1168 let (backend, log) = MockMigrationBackend::new(None);
1169 let db = Database::with_backend(Box::new(backend), "test");
1170
1171 let result = execute_recovery_plan(&db, &plan, &controller).await;
1172
1173 assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
1174 assert!(
1175 result.migrations[0]
1176 .steps
1177 .iter()
1178 .all(|step| matches!(step.outcome, StepExecutionOutcome::Applied { .. }))
1179 );
1180 assert_eq!(
1181 *log.lock().unwrap(),
1182 vec![MockEvent::OpenTx(TxType::Read), MockEvent::Close]
1183 );
1184 }
1185
1186 #[tokio::test]
1187 async fn backfill_emits_counts_on_both_commit_boundary_events() {
1188 let plan = checked_plan(vec![backfill_step()]);
1189 let controller = TestController::new(vec![StepRecoveryDecision::Pending {
1190 proof: PendingProof::IdempotentReplay {
1191 strategy: "copy-if-absent".to_string(),
1192 },
1193 }]);
1194 let responses = vec![
1195 QueryResult::Rows(vec![json!({"c": 4})]),
1196 QueryResult::Rows(vec![json!({"c": 7})]),
1197 QueryResult::Ok,
1198 ];
1199 let (backend, _log) = MockMigrationBackend::with_responses(responses);
1200 let db = Database::with_backend(Box::new(backend), "test");
1201
1202 let result = execute_recovery_plan(&db, &plan, &controller).await;
1203
1204 assert_eq!(result.status, RecoveryPlanStatus::Succeeded);
1205 let events = controller.events.lock().unwrap();
1206 assert_eq!(events.len(), 2);
1207 for event in events.iter() {
1208 let counts = event.backfill.as_ref().unwrap();
1209 assert_eq!(counts.inserted, 4);
1210 assert_eq!(counts.matched, 7);
1211 assert_eq!(counts.skipped, 3);
1212 }
1213 assert!(matches!(
1214 result.migrations[0].steps[0].outcome,
1215 StepExecutionOutcome::Committed { backfill: Some(_) }
1216 ));
1217 }
1218}