1mod 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 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 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 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 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 #[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 #[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 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 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 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 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 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 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 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 #[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}