Skip to main content

canic_backup/execution/
mod.rs

1//! Module: execution
2//!
3//! Responsibility: build and advance backup execution journals.
4//! Does not own: backup plan construction, artifact IO, or manifest storage.
5//! Boundary: tracks runner progress from validated plans through receipts.
6
7mod operation;
8mod receipt;
9#[cfg(test)]
10mod tests;
11mod types;
12mod validation;
13
14pub use types::*;
15
16use crate::plan::{BackupExecutionPreflightReceipts, BackupOperationKind, BackupPlan};
17use validation::{
18    operation_kind_is_mutating, operation_kind_is_preflight, validate_nonempty,
19    validate_operation_sequences,
20};
21
22const BACKUP_EXECUTION_JOURNAL_VERSION: u16 = 1;
23const PREFLIGHT_NOT_ACCEPTED: &str = "preflight-not-accepted";
24
25impl BackupExecutionJournal {
26    /// Build an execution journal from a validated backup plan.
27    pub fn from_plan(plan: &BackupPlan) -> Result<Self, BackupExecutionJournalError> {
28        plan.validate()
29            .map_err(|error| BackupExecutionJournalError::InvalidPlan(error.to_string()))?;
30        let operations = plan
31            .phases
32            .iter()
33            .map(BackupExecutionJournalOperation::from_plan_operation)
34            .collect::<Vec<_>>();
35        let mut journal = Self {
36            journal_version: BACKUP_EXECUTION_JOURNAL_VERSION,
37            plan_id: plan.plan_id.clone(),
38            run_id: plan.run_id.clone(),
39            preflight_id: None,
40            preflight_accepted: false,
41            restart_required: false,
42            operations,
43            operation_receipts: Vec::new(),
44        };
45        journal.refresh_blocked_operations();
46        journal.validate()?;
47        Ok(journal)
48    }
49
50    /// Validate journal structure and operation receipts.
51    pub fn validate(&self) -> Result<(), BackupExecutionJournalError> {
52        if self.journal_version != BACKUP_EXECUTION_JOURNAL_VERSION {
53            return Err(BackupExecutionJournalError::UnsupportedVersion(
54                self.journal_version,
55            ));
56        }
57        validate_nonempty("plan_id", &self.plan_id)?;
58        validate_nonempty("run_id", &self.run_id)?;
59        if let Some(preflight_id) = &self.preflight_id {
60            validate_nonempty("preflight_id", preflight_id)?;
61        } else if self.preflight_accepted {
62            return Err(BackupExecutionJournalError::AcceptedPreflightMissingId);
63        }
64        if self.restart_required != self.derived_restart_required() {
65            return Err(BackupExecutionJournalError::RestartRequiredMismatch);
66        }
67        validate_operation_sequences(&self.operations)?;
68        for operation in &self.operations {
69            operation.validate()?;
70            if !self.preflight_accepted && operation_kind_is_mutating(&operation.kind) {
71                match operation.state {
72                    BackupExecutionOperationState::Blocked => {}
73                    BackupExecutionOperationState::Ready
74                    | BackupExecutionOperationState::Pending
75                    | BackupExecutionOperationState::Completed
76                    | BackupExecutionOperationState::Failed
77                    | BackupExecutionOperationState::Skipped => {
78                        return Err(BackupExecutionJournalError::MutationReadyBeforePreflight {
79                            sequence: operation.sequence,
80                        });
81                    }
82                }
83            }
84        }
85        for receipt in &self.operation_receipts {
86            receipt.validate_against(self)?;
87        }
88        Ok(())
89    }
90
91    /// Mark all preflight operations completed and unblock mutating operations.
92    pub fn accept_preflight_bundle_at(
93        &mut self,
94        preflight_id: String,
95        updated_at: Option<String>,
96    ) -> Result<(), BackupExecutionJournalError> {
97        validate_nonempty("preflight_id", &preflight_id)?;
98        validate_nonempty("updated_at", updated_at.as_deref().unwrap_or_default())?;
99        if let Some(existing) = &self.preflight_id
100            && existing != &preflight_id
101        {
102            return Err(BackupExecutionJournalError::PreflightAlreadyAccepted {
103                existing: existing.clone(),
104                attempted: preflight_id,
105            });
106        }
107
108        self.preflight_id = Some(preflight_id);
109        self.preflight_accepted = true;
110        for operation in &mut self.operations {
111            if operation_kind_is_preflight(&operation.kind) {
112                operation.state = BackupExecutionOperationState::Completed;
113                operation.state_updated_at.clone_from(&updated_at);
114                operation.blocking_reasons.clear();
115            } else if operation.state == BackupExecutionOperationState::Blocked {
116                operation.state = BackupExecutionOperationState::Ready;
117                operation.blocking_reasons.clear();
118            }
119        }
120        self.refresh_restart_required();
121        self.validate()
122    }
123
124    /// Accept a typed preflight receipt bundle and unblock mutating operations.
125    pub fn accept_preflight_receipts_at(
126        &mut self,
127        receipts: &BackupExecutionPreflightReceipts,
128        updated_at: Option<String>,
129    ) -> Result<(), BackupExecutionJournalError> {
130        validate_nonempty("preflight_receipts.plan_id", &receipts.plan_id)?;
131        if receipts.plan_id != self.plan_id {
132            return Err(BackupExecutionJournalError::PreflightPlanMismatch {
133                expected: self.plan_id.clone(),
134                actual: receipts.plan_id.clone(),
135            });
136        }
137        self.accept_preflight_bundle_at(receipts.preflight_id.clone(), updated_at)
138    }
139
140    /// Return the next operation that should control runner progress.
141    #[must_use]
142    pub fn next_ready_operation(&self) -> Option<&BackupExecutionJournalOperation> {
143        self.operations
144            .iter()
145            .filter(|operation| {
146                matches!(
147                    operation.state,
148                    BackupExecutionOperationState::Ready
149                        | BackupExecutionOperationState::Pending
150                        | BackupExecutionOperationState::Failed
151                )
152            })
153            .min_by_key(|operation| operation.sequence)
154    }
155
156    /// Return the next start needed to restore availability after a backup failure.
157    #[must_use]
158    pub(crate) fn next_failure_containment_start(
159        &self,
160    ) -> Option<&BackupExecutionJournalOperation> {
161        self.operations
162            .iter()
163            .filter(|operation| {
164                operation.kind == BackupOperationKind::Start
165                    && matches!(
166                        operation.state,
167                        BackupExecutionOperationState::Ready
168                            | BackupExecutionOperationState::Pending
169                            | BackupExecutionOperationState::Failed
170                    )
171                    && operation.target_canister_id.as_ref().is_some_and(|target| {
172                        self.operations.iter().any(|candidate| {
173                            candidate.kind == BackupOperationKind::Stop
174                                && candidate.target_canister_id.as_ref() == Some(target)
175                                && candidate.state == BackupExecutionOperationState::Completed
176                        })
177                    })
178            })
179            .min_by_key(|operation| operation.sequence)
180    }
181
182    /// Claim one start that is allowed to bypass a failed primary operation.
183    pub(crate) fn mark_failure_containment_start_pending_at(
184        &mut self,
185        sequence: usize,
186        updated_at: Option<String>,
187    ) -> Result<(), BackupExecutionJournalError> {
188        validate_nonempty("updated_at", updated_at.as_deref().unwrap_or_default())?;
189        let eligible = self
190            .next_failure_containment_start()
191            .is_some_and(|operation| operation.sequence == sequence);
192        if !eligible {
193            return Err(BackupExecutionJournalError::OperationNotFailureContainmentStart(sequence));
194        }
195        let index = self.operation_index(sequence)?;
196        if !matches!(
197            self.operations[index].state,
198            BackupExecutionOperationState::Ready | BackupExecutionOperationState::Failed
199        ) {
200            return Err(BackupExecutionJournalError::InvalidOperationTransition {
201                sequence,
202                from: self.operations[index].state.clone(),
203                to: BackupExecutionOperationState::Pending,
204            });
205        }
206
207        let previous_operation = self.operations[index].clone();
208        let previous_restart_required = self.restart_required;
209        self.operations[index].state = BackupExecutionOperationState::Pending;
210        self.operations[index].state_updated_at = updated_at;
211        self.operations[index].blocking_reasons.clear();
212        self.refresh_restart_required();
213        if let Err(error) = self.validate() {
214            self.operations[index] = previous_operation;
215            self.restart_required = previous_restart_required;
216            return Err(error);
217        }
218        Ok(())
219    }
220
221    /// Rearm the paired stop/start after availability was restored before snapshot completion.
222    pub(crate) fn rearm_after_failure_containment(
223        &mut self,
224        start_sequence: usize,
225        updated_at: Option<String>,
226    ) -> Result<(), BackupExecutionJournalError> {
227        validate_nonempty("updated_at", updated_at.as_deref().unwrap_or_default())?;
228        let start_index = self.operation_index(start_sequence)?;
229        let start = &self.operations[start_index];
230        if start.kind != BackupOperationKind::Start
231            || start.state != BackupExecutionOperationState::Completed
232        {
233            return Err(
234                BackupExecutionJournalError::OperationNotFailureContainmentStart(start_sequence),
235            );
236        }
237        let target = start.target_canister_id.clone();
238        let stop_index = self
239            .operations
240            .iter()
241            .position(|operation| {
242                operation.kind == BackupOperationKind::Stop
243                    && operation.target_canister_id == target
244                    && operation.state == BackupExecutionOperationState::Completed
245            })
246            .ok_or(
247                BackupExecutionJournalError::OperationNotFailureContainmentStart(start_sequence),
248            )?;
249        let previous_stop = self.operations[stop_index].clone();
250        let previous_start = self.operations[start_index].clone();
251        let previous_restart_required = self.restart_required;
252        for index in [stop_index, start_index] {
253            self.operations[index].state = BackupExecutionOperationState::Ready;
254            self.operations[index]
255                .state_updated_at
256                .clone_from(&updated_at);
257            self.operations[index].blocking_reasons.clear();
258        }
259        self.refresh_restart_required();
260        if let Err(error) = self.validate() {
261            self.operations[stop_index] = previous_stop;
262            self.operations[start_index] = previous_start;
263            self.restart_required = previous_restart_required;
264            return Err(error);
265        }
266        Ok(())
267    }
268
269    /// Mark the next transitionable operation pending.
270    pub fn mark_next_operation_pending_at(
271        &mut self,
272        updated_at: Option<String>,
273    ) -> Result<(), BackupExecutionJournalError> {
274        let sequence = self
275            .next_ready_operation()
276            .ok_or(BackupExecutionJournalError::NoTransitionableOperation)?
277            .sequence;
278        self.mark_operation_pending_at(sequence, updated_at)
279    }
280
281    /// Mark one operation pending.
282    pub fn mark_operation_pending_at(
283        &mut self,
284        sequence: usize,
285        updated_at: Option<String>,
286    ) -> Result<(), BackupExecutionJournalError> {
287        self.mark_operation_pending_with_snapshot_inventory_at(sequence, updated_at, None)
288    }
289
290    /// Mark one snapshot-create operation pending with its exact pre-effect inventory.
291    pub fn mark_snapshot_create_pending_at(
292        &mut self,
293        sequence: usize,
294        updated_at: Option<String>,
295        snapshot_ids_before: Vec<String>,
296    ) -> Result<(), BackupExecutionJournalError> {
297        self.mark_operation_pending_with_snapshot_inventory_at(
298            sequence,
299            updated_at,
300            Some(snapshot_ids_before),
301        )
302    }
303
304    fn mark_operation_pending_with_snapshot_inventory_at(
305        &mut self,
306        sequence: usize,
307        updated_at: Option<String>,
308        snapshot_ids_before: Option<Vec<String>>,
309    ) -> Result<(), BackupExecutionJournalError> {
310        validate_nonempty("updated_at", updated_at.as_deref().unwrap_or_default())?;
311        let expected = self
312            .next_ready_operation()
313            .ok_or(BackupExecutionJournalError::NoTransitionableOperation)?
314            .sequence;
315        if sequence != expected {
316            return Err(BackupExecutionJournalError::OutOfOrderOperationTransition {
317                requested: sequence,
318                next: expected,
319            });
320        }
321        let index = self.operation_index(sequence)?;
322        let operation = &self.operations[index];
323        if operation_kind_is_mutating(&operation.kind) && !self.preflight_accepted {
324            return Err(BackupExecutionJournalError::MutationBeforePreflightAccepted { sequence });
325        }
326        if !matches!(
327            operation.state,
328            BackupExecutionOperationState::Ready | BackupExecutionOperationState::Failed
329        ) {
330            return Err(BackupExecutionJournalError::InvalidOperationTransition {
331                sequence,
332                from: operation.state.clone(),
333                to: BackupExecutionOperationState::Pending,
334            });
335        }
336
337        let previous_operation = self.operations[index].clone();
338        let previous_restart_required = self.restart_required;
339        let operation = &mut self.operations[index];
340        operation.state = BackupExecutionOperationState::Pending;
341        operation.state_updated_at = updated_at;
342        operation.snapshot_ids_before = snapshot_ids_before;
343        operation.blocking_reasons.clear();
344        self.refresh_restart_required();
345        if let Err(error) = self.validate() {
346            self.operations[index] = previous_operation;
347            self.restart_required = previous_restart_required;
348            return Err(error);
349        }
350        Ok(())
351    }
352
353    /// Record one operation receipt and transition the matching operation.
354    pub fn record_operation_receipt(
355        &mut self,
356        receipt: BackupExecutionOperationReceipt,
357    ) -> Result<(), BackupExecutionJournalError> {
358        receipt.validate_against(self)?;
359        let index = self.operation_index(receipt.sequence)?;
360        let operation = &self.operations[index];
361        if operation.state != BackupExecutionOperationState::Pending {
362            return Err(
363                BackupExecutionJournalError::ReceiptWithoutPendingOperation {
364                    sequence: receipt.sequence,
365                },
366            );
367        }
368
369        let next_state = match receipt.outcome {
370            BackupExecutionOperationReceiptOutcome::Completed => {
371                BackupExecutionOperationState::Completed
372            }
373            BackupExecutionOperationReceiptOutcome::Failed => BackupExecutionOperationState::Failed,
374            BackupExecutionOperationReceiptOutcome::Skipped => {
375                BackupExecutionOperationState::Skipped
376            }
377        };
378        let failure_reason = receipt.failure_reason.clone();
379        let previous_operation = self.operations[index].clone();
380        let previous_restart_required = self.restart_required;
381        self.operation_receipts.push(receipt);
382
383        let operation = &mut self.operations[index];
384        operation.state = next_state;
385        operation.state_updated_at = self
386            .operation_receipts
387            .last()
388            .and_then(|receipt| receipt.updated_at.clone());
389        operation.blocking_reasons = failure_reason.into_iter().collect();
390        self.refresh_restart_required();
391        if let Err(error) = self.validate() {
392            self.operation_receipts.pop();
393            self.operations[index] = previous_operation;
394            self.restart_required = previous_restart_required;
395            return Err(error);
396        }
397        Ok(())
398    }
399
400    /// Move a failed operation back to ready for retry.
401    pub fn retry_failed_operation_at(
402        &mut self,
403        sequence: usize,
404        updated_at: Option<String>,
405    ) -> Result<(), BackupExecutionJournalError> {
406        validate_nonempty("updated_at", updated_at.as_deref().unwrap_or_default())?;
407        let index = self.operation_index(sequence)?;
408        if self.operations[index].state != BackupExecutionOperationState::Failed {
409            return Err(BackupExecutionJournalError::OperationNotFailed(sequence));
410        }
411        self.operations[index].state = BackupExecutionOperationState::Ready;
412        self.operations[index].state_updated_at = updated_at;
413        self.operations[index].blocking_reasons.clear();
414        self.refresh_restart_required();
415        self.validate()
416    }
417
418    /// Build a compact resumability summary.
419    #[must_use]
420    pub fn resume_summary(&self) -> BackupExecutionResumeSummary {
421        let mut summary = BackupExecutionResumeSummary {
422            plan_id: self.plan_id.clone(),
423            run_id: self.run_id.clone(),
424            preflight_id: self.preflight_id.clone(),
425            preflight_accepted: self.preflight_accepted,
426            restart_required: self.restart_required,
427            total_operations: self.operations.len(),
428            ready_operations: 0,
429            pending_operations: 0,
430            blocked_operations: 0,
431            completed_operations: 0,
432            failed_operations: 0,
433            skipped_operations: 0,
434            next_operation: self.next_ready_operation().cloned(),
435        };
436        for operation in &self.operations {
437            match operation.state {
438                BackupExecutionOperationState::Ready => summary.ready_operations += 1,
439                BackupExecutionOperationState::Pending => summary.pending_operations += 1,
440                BackupExecutionOperationState::Blocked => summary.blocked_operations += 1,
441                BackupExecutionOperationState::Completed => summary.completed_operations += 1,
442                BackupExecutionOperationState::Failed => summary.failed_operations += 1,
443                BackupExecutionOperationState::Skipped => summary.skipped_operations += 1,
444            }
445        }
446        summary
447    }
448
449    fn operation_index(&self, sequence: usize) -> Result<usize, BackupExecutionJournalError> {
450        self.operations
451            .iter()
452            .position(|operation| operation.sequence == sequence)
453            .ok_or(BackupExecutionJournalError::OperationNotFound(sequence))
454    }
455
456    fn refresh_blocked_operations(&mut self) {
457        if self.preflight_accepted {
458            return;
459        }
460        for operation in &mut self.operations {
461            if operation_kind_is_mutating(&operation.kind) {
462                operation.state = BackupExecutionOperationState::Blocked;
463                operation.blocking_reasons = vec![PREFLIGHT_NOT_ACCEPTED.to_string()];
464            }
465        }
466    }
467
468    fn refresh_restart_required(&mut self) {
469        self.restart_required = self.derived_restart_required();
470    }
471
472    fn derived_restart_required(&self) -> bool {
473        self.operations.iter().any(|stop| {
474            stop.kind == BackupOperationKind::Stop
475                && stop.state == BackupExecutionOperationState::Completed
476                && self.operations.iter().any(|start| {
477                    start.kind == BackupOperationKind::Start
478                        && start.target_canister_id == stop.target_canister_id
479                        && !matches!(
480                            start.state,
481                            BackupExecutionOperationState::Completed
482                                | BackupExecutionOperationState::Skipped
483                        )
484                })
485        })
486    }
487}