1use 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
30pub 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}