Skip to main content

canic_backup/restore/runner/
execute.rs

1//! Module: restore::runner::execute
2//!
3//! Responsibility: execute ready restore apply journal operations.
4//! Does not own: command rendering, apply journal validation, or restore planning.
5//! Boundary: claims operations, invokes an injected executor, and persists receipts.
6
7use super::{
8    RestoreApplyCommandOutputPair, RestoreApplyJournal, RestoreApplyOperationKind,
9    RestoreApplyOperationReceipt,
10    constants::{
11        RESTORE_RUN_COMMAND_EXIT_PREFIX, RESTORE_RUN_MISSING_UPLOADED_SNAPSHOT_ID,
12        RESTORE_RUN_OUTPUT_RECEIPT_LIMIT, RESTORE_RUN_STOPPED_PRECONDITION_FAILED,
13    },
14    io::{read_apply_journal_file, write_apply_journal_file},
15    precondition::enforce_stopped_canister_precondition,
16    status::{
17        enforce_restore_run_command_available, enforce_restore_run_executable,
18        parse_uploaded_snapshot_id, restore_command_unavailable_error,
19        restore_run_max_steps_reached, restore_run_next_action, restore_run_stopped_reason,
20    },
21    types::{
22        RestoreRunExecutedOperation, RestoreRunOperationReceipt, RestoreRunPreparedOperation,
23        RestoreRunResponse, RestoreRunResponseMode, RestoreRunStepOutcome,
24        RestoreRunnerCommandExecutor, RestoreRunnerConfig, RestoreRunnerError,
25        RestoreRunnerOutcome, RestoreStoppedPreconditionFailure,
26    },
27};
28use crate::{persistence::JournalLock, timestamp::state_updated_at};
29
30/// Execute ready restore apply journal operations through an injected command executor.
31pub fn restore_run_execute_with_executor(
32    config: &RestoreRunnerConfig,
33    executor: &mut impl RestoreRunnerCommandExecutor,
34) -> Result<RestoreRunResponse, RestoreRunnerError> {
35    let run = restore_run_execute_result_with_executor(config, executor)?;
36    if let Some(error) = run.error {
37        return Err(error);
38    }
39
40    Ok(run.response)
41}
42
43pub fn restore_run_execute_result_with_executor(
44    config: &RestoreRunnerConfig,
45    executor: &mut impl RestoreRunnerCommandExecutor,
46) -> Result<RestoreRunnerOutcome, RestoreRunnerError> {
47    let _lock = JournalLock::acquire(&config.journal)?;
48    let mut journal = read_apply_journal_file(&config.journal)?;
49    let mut executed_operations = Vec::new();
50    let mut operation_receipts = Vec::new();
51
52    loop {
53        let report = journal.report();
54        let max_steps_reached =
55            restore_run_max_steps_reached(config, executed_operations.len(), &report);
56        if report.complete || max_steps_reached {
57            return Ok(RestoreRunnerOutcome::ok(restore_run_execute_summary(
58                &journal,
59                executed_operations,
60                operation_receipts,
61                max_steps_reached,
62                config.updated_at.as_ref(),
63            )));
64        }
65
66        enforce_restore_run_executable(&journal, &report)?;
67        let prepared = restore_run_prepare_next_operation(config, &mut journal)?;
68        let sequence = prepared.sequence;
69        match restore_run_execute_prepared_operation(config, executor, &mut journal, prepared)? {
70            RestoreRunStepOutcome::Completed {
71                executed_operation,
72                operation_receipt,
73            } => {
74                executed_operations.push(executed_operation);
75                operation_receipts.push(operation_receipt);
76            }
77            RestoreRunStepOutcome::Failed {
78                executed_operation,
79                operation_receipt,
80                status,
81            } => {
82                executed_operations.push(executed_operation);
83                operation_receipts.push(operation_receipt);
84                let response = restore_run_execute_summary(
85                    &journal,
86                    executed_operations,
87                    operation_receipts,
88                    false,
89                    config.updated_at.as_ref(),
90                );
91                return Ok(RestoreRunnerOutcome {
92                    response,
93                    error: Some(RestoreRunnerError::CommandFailed { sequence, status }),
94                });
95            }
96        }
97    }
98}
99
100fn restore_run_prepare_next_operation(
101    config: &RestoreRunnerConfig,
102    journal: &mut RestoreApplyJournal,
103) -> Result<RestoreRunPreparedOperation, RestoreRunnerError> {
104    let preview = journal.next_command_preview_with_config(&config.command);
105    enforce_restore_run_command_available(&preview)?;
106
107    let operation = preview
108        .operation
109        .clone()
110        .ok_or_else(|| restore_command_unavailable_error(&preview))?;
111    let command = preview
112        .command
113        .clone()
114        .ok_or_else(|| restore_command_unavailable_error(&preview))?;
115    let sequence = operation.sequence;
116    let attempt = journal
117        .operation_receipts
118        .iter()
119        .filter(|receipt| receipt.sequence == sequence)
120        .count()
121        + 1;
122
123    enforce_apply_claim_sequence(sequence, journal)?;
124    journal
125        .mark_operation_pending_at(sequence, Some(state_updated_at(config.updated_at.as_ref())))?;
126    write_apply_journal_file(&config.journal, journal)?;
127
128    Ok(RestoreRunPreparedOperation {
129        operation,
130        command,
131        sequence,
132        attempt,
133    })
134}
135
136fn restore_run_execute_prepared_operation(
137    config: &RestoreRunnerConfig,
138    executor: &mut impl RestoreRunnerCommandExecutor,
139    journal: &mut RestoreApplyJournal,
140    prepared: RestoreRunPreparedOperation,
141) -> Result<RestoreRunStepOutcome, RestoreRunnerError> {
142    if prepared.command.requires_stopped_canister
143        && let Some(outcome) = enforce_stopped_canister_precondition(
144            config,
145            executor,
146            &prepared.operation,
147            prepared.attempt,
148            config.updated_at.as_ref(),
149        )?
150    {
151        return restore_run_commit_precondition_failure(config, journal, prepared, outcome);
152    }
153
154    let output = executor.execute(&prepared.command)?;
155    let status_label = output.status;
156    let output_pair = RestoreApplyCommandOutputPair::from_bytes(
157        &output.stdout,
158        &output.stderr,
159        RESTORE_RUN_OUTPUT_RECEIPT_LIMIT,
160    );
161
162    if output.success {
163        let is_upload_snapshot =
164            prepared.operation.operation == RestoreApplyOperationKind::UploadSnapshot;
165        let uploaded_snapshot_id = is_upload_snapshot
166            .then(|| parse_uploaded_snapshot_id(&String::from_utf8_lossy(&output.stdout)))
167            .flatten();
168        if is_upload_snapshot && uploaded_snapshot_id.is_none() {
169            return restore_run_commit_missing_uploaded_snapshot_id(
170                config,
171                journal,
172                prepared,
173                output_pair,
174            );
175        }
176
177        return restore_run_commit_command_success(
178            config,
179            journal,
180            prepared,
181            status_label,
182            output_pair,
183            uploaded_snapshot_id,
184        );
185    }
186
187    restore_run_commit_command_failure(config, journal, prepared, status_label, output_pair)
188}
189
190fn restore_run_commit_missing_uploaded_snapshot_id(
191    config: &RestoreRunnerConfig,
192    journal: &mut RestoreApplyJournal,
193    prepared: RestoreRunPreparedOperation,
194    output_pair: RestoreApplyCommandOutputPair,
195) -> Result<RestoreRunStepOutcome, RestoreRunnerError> {
196    let failed_updated_at = state_updated_at(config.updated_at.as_ref());
197    let status = RESTORE_RUN_MISSING_UPLOADED_SNAPSHOT_ID.to_string();
198    journal.mark_operation_failed_at(
199        prepared.sequence,
200        status.clone(),
201        Some(failed_updated_at.clone()),
202    )?;
203    journal.record_operation_receipt(RestoreApplyOperationReceipt::command_failed(
204        &prepared.operation,
205        prepared.command.clone(),
206        status.clone(),
207        Some(failed_updated_at.clone()),
208        output_pair,
209        prepared.attempt,
210        status.clone(),
211    ))?;
212    write_apply_journal_file(&config.journal, journal)?;
213
214    Ok(RestoreRunStepOutcome::Failed {
215        executed_operation: RestoreRunExecutedOperation::failed(
216            prepared.operation.clone(),
217            prepared.command.clone(),
218            RESTORE_RUN_MISSING_UPLOADED_SNAPSHOT_ID.to_string(),
219        ),
220        operation_receipt: RestoreRunOperationReceipt::failed(
221            prepared.operation,
222            prepared.command,
223            status.clone(),
224            Some(failed_updated_at),
225        ),
226        status,
227    })
228}
229
230fn restore_run_commit_precondition_failure(
231    config: &RestoreRunnerConfig,
232    journal: &mut RestoreApplyJournal,
233    prepared: RestoreRunPreparedOperation,
234    outcome: RestoreStoppedPreconditionFailure,
235) -> Result<RestoreRunStepOutcome, RestoreRunnerError> {
236    let failed_updated_at = state_updated_at(config.updated_at.as_ref());
237    journal.mark_operation_failed_at(
238        prepared.sequence,
239        outcome.failure_reason.clone(),
240        Some(failed_updated_at.clone()),
241    )?;
242    journal.record_operation_receipt(RestoreApplyOperationReceipt::command_failed(
243        &prepared.operation,
244        outcome.command.clone(),
245        outcome.status_label.clone(),
246        Some(failed_updated_at.clone()),
247        outcome.output,
248        prepared.attempt,
249        outcome.failure_reason,
250    ))?;
251    write_apply_journal_file(&config.journal, journal)?;
252
253    Ok(RestoreRunStepOutcome::Failed {
254        executed_operation: RestoreRunExecutedOperation::failed(
255            prepared.operation.clone(),
256            outcome.command.clone(),
257            outcome.status_label.clone(),
258        ),
259        operation_receipt: RestoreRunOperationReceipt::failed(
260            prepared.operation,
261            outcome.command,
262            outcome.status_label,
263            Some(failed_updated_at),
264        ),
265        status: RESTORE_RUN_STOPPED_PRECONDITION_FAILED.to_string(),
266    })
267}
268
269fn restore_run_commit_command_success(
270    config: &RestoreRunnerConfig,
271    journal: &mut RestoreApplyJournal,
272    prepared: RestoreRunPreparedOperation,
273    status_label: String,
274    output_pair: RestoreApplyCommandOutputPair,
275    uploaded_snapshot_id: Option<String>,
276) -> Result<RestoreRunStepOutcome, RestoreRunnerError> {
277    let completed_updated_at = state_updated_at(config.updated_at.as_ref());
278    journal.mark_operation_completed_at(prepared.sequence, Some(completed_updated_at.clone()))?;
279    journal.record_operation_receipt(RestoreApplyOperationReceipt::command_completed(
280        &prepared.operation,
281        prepared.command.clone(),
282        status_label.clone(),
283        Some(completed_updated_at.clone()),
284        output_pair,
285        prepared.attempt,
286        uploaded_snapshot_id,
287    ))?;
288    write_apply_journal_file(&config.journal, journal)?;
289
290    Ok(RestoreRunStepOutcome::Completed {
291        executed_operation: RestoreRunExecutedOperation::completed(
292            prepared.operation.clone(),
293            prepared.command.clone(),
294            status_label.clone(),
295        ),
296        operation_receipt: RestoreRunOperationReceipt::completed(
297            prepared.operation,
298            prepared.command,
299            status_label,
300            Some(completed_updated_at),
301        ),
302    })
303}
304
305fn restore_run_commit_command_failure(
306    config: &RestoreRunnerConfig,
307    journal: &mut RestoreApplyJournal,
308    prepared: RestoreRunPreparedOperation,
309    status_label: String,
310    output_pair: RestoreApplyCommandOutputPair,
311) -> Result<RestoreRunStepOutcome, RestoreRunnerError> {
312    let failed_updated_at = state_updated_at(config.updated_at.as_ref());
313    let failure_reason = format!("{RESTORE_RUN_COMMAND_EXIT_PREFIX}-{status_label}");
314    journal.mark_operation_failed_at(
315        prepared.sequence,
316        failure_reason.clone(),
317        Some(failed_updated_at.clone()),
318    )?;
319    journal.record_operation_receipt(RestoreApplyOperationReceipt::command_failed(
320        &prepared.operation,
321        prepared.command.clone(),
322        status_label.clone(),
323        Some(failed_updated_at.clone()),
324        output_pair,
325        prepared.attempt,
326        failure_reason,
327    ))?;
328    write_apply_journal_file(&config.journal, journal)?;
329
330    Ok(RestoreRunStepOutcome::Failed {
331        executed_operation: RestoreRunExecutedOperation::failed(
332            prepared.operation.clone(),
333            prepared.command.clone(),
334            status_label.clone(),
335        ),
336        operation_receipt: RestoreRunOperationReceipt::failed(
337            prepared.operation,
338            prepared.command,
339            status_label.clone(),
340            Some(failed_updated_at),
341        ),
342        status: status_label,
343    })
344}
345
346fn restore_run_execute_summary(
347    journal: &RestoreApplyJournal,
348    executed_operations: Vec<RestoreRunExecutedOperation>,
349    operation_receipts: Vec<RestoreRunOperationReceipt>,
350    max_steps_reached: bool,
351    requested_state_updated_at: Option<&String>,
352) -> RestoreRunResponse {
353    let report = journal.report();
354    let executed_operation_count = executed_operations.len();
355    let stopped_reason = restore_run_stopped_reason(&report, max_steps_reached, true);
356    let next_action = restore_run_next_action(&report);
357
358    let mut response = RestoreRunResponse::from_report(
359        journal.backup_id.clone(),
360        report,
361        RestoreRunResponseMode::execute(stopped_reason, next_action),
362    );
363    response.set_requested_state_updated_at(requested_state_updated_at);
364    response.max_steps_reached = Some(max_steps_reached);
365    response.executed_operation_count = Some(executed_operation_count);
366    response.executed_operations = executed_operations;
367    response.set_operation_receipts(operation_receipts);
368    response
369}
370
371fn enforce_apply_claim_sequence(
372    expected: usize,
373    journal: &RestoreApplyJournal,
374) -> Result<(), RestoreRunnerError> {
375    let actual = journal
376        .next_transition_operation()
377        .map(|operation| operation.sequence);
378
379    if actual == Some(expected) {
380        return Ok(());
381    }
382
383    Err(RestoreRunnerError::ClaimSequenceMismatch { expected, actual })
384}