ic_backup/ops/persistence/execution_workflow/
mod.rs1use super::{
4 BackupLayoutGuard, ExecutionSettlementPersistenceError, JournalLock, JournalLockError,
5 OperationPlanPersistenceError, PersistenceError, create_json_durable, create_operation_plan,
6 read_execution_settlement, read_json, read_operation_plan,
7};
8use crate::model::{
9 artifacts::ArtifactChecksumRecord,
10 execution_workflow::{
11 ExecutionStageBindingRecord, ExecutionWorkflowError, ExecutionWorkflowRecord,
12 MAX_EXECUTION_STAGE_BYTES, MAX_EXECUTION_WORKFLOW_BYTES,
13 },
14 operation_plan::OperationPlanRecord,
15};
16use std::{fs, path::PathBuf};
17use thiserror::Error;
18
19const WORKFLOW_FILE: &str = "execution-workflow.json";
20const BINDING_FILE: &str = "stage-binding.json";
21
22#[derive(Clone, Copy, Debug, Eq, PartialEq)]
23enum StagePreparationBarrier {
24 AfterPlanPublication,
25 AfterBindingPublication,
26}
27
28pub fn create_execution_workflow(
32 layout: &BackupLayoutGuard,
33 record: &ExecutionWorkflowRecord,
34) -> Result<(), ExecutionWorkflowPersistenceError> {
35 layout.check_root()?;
36 super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
37 let path = layout.root().join(WORKFLOW_FILE);
38 let _lock = JournalLock::acquire(&path)?;
39 create_json_durable(&path, record)?;
40 Ok(())
41}
42pub fn read_execution_workflow(
46 layout: &BackupLayoutGuard,
47 expected: &ArtifactChecksumRecord,
48) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
49 layout.check_root()?;
50 let path = layout.root().join(WORKFLOW_FILE);
51 let _lock = JournalLock::acquire(&path)?;
52 let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
53 super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
54 if &record.digest() != expected {
55 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
56 }
57 Ok(record)
58}
59
60#[derive(Debug)]
67pub struct ExecutionStageGuard<'a> {
68 workflow_layout: &'a BackupLayoutGuard,
69 stage_layout: BackupLayoutGuard,
70 binding: ExecutionStageBindingRecord,
71 plan: OperationPlanRecord,
72}
73impl<'a> ExecutionStageGuard<'a> {
74 pub fn create(
83 workflow_layout: &'a BackupLayoutGuard,
84 binding: ExecutionStageBindingRecord,
85 plan: OperationPlanRecord,
86 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
87 Self::create_with(workflow_layout, binding, plan, |_| {})
88 }
89 fn create_with(
90 workflow_layout: &'a BackupLayoutGuard,
91 binding: ExecutionStageBindingRecord,
92 plan: OperationPlanRecord,
93 mut barrier: impl FnMut(StagePreparationBarrier),
94 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
95 let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
96 binding.validate(&workflow, &plan)?;
97 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
98 super::json::check_json_size(
100 &plan,
101 crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
102 )?;
103 validate_predecessors(workflow_layout, &workflow, &binding)?;
104 let path = stage_path(workflow_layout, binding.stage_sequence());
105 create_private_directory(&path)?;
106 fs::File::open(workflow_layout.root())
107 .and_then(|directory| directory.sync_all())
108 .map_err(PersistenceError::from)?;
109 let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
110 create_operation_plan(&stage_layout, &plan)?;
111 barrier(StagePreparationBarrier::AfterPlanPublication);
112 create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
113 barrier(StagePreparationBarrier::AfterBindingPublication);
114 workflow_layout.check_root()?;
115 Ok(Self {
116 workflow_layout,
117 stage_layout,
118 binding,
119 plan,
120 })
121 }
122 pub fn open(
131 workflow_layout: &'a BackupLayoutGuard,
132 expected_workflow: &ArtifactChecksumRecord,
133 sequence: u64,
134 expected_binding: &ArtifactChecksumRecord,
135 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
136 let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
137 workflow
138 .stage(sequence)
139 .map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
140 let stage_layout = acquire_stage(workflow_layout, sequence)?;
141 let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
142 validate_predecessors(workflow_layout, &workflow, &binding)?;
143 workflow_layout.check_root()?;
144 Ok(Self {
145 workflow_layout,
146 stage_layout,
147 binding,
148 plan,
149 })
150 }
151 #[must_use]
153 pub const fn plan(&self) -> &OperationPlanRecord {
154 &self.plan
155 }
156 #[must_use]
158 pub const fn binding(&self) -> &ExecutionStageBindingRecord {
159 &self.binding
160 }
161 pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
168 let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
169 read_binding(
170 &self.stage_layout,
171 &workflow,
172 self.binding.stage_sequence(),
173 &self.binding.digest(),
174 )?;
175 validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
176 Ok(&self.stage_layout)
177 }
178}
179fn read_binding(
180 layout: &BackupLayoutGuard,
181 workflow: &ExecutionWorkflowRecord,
182 sequence: u64,
183 expected: &ArtifactChecksumRecord,
184) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
185 layout.check_root()?;
186 let binding: ExecutionStageBindingRecord =
187 read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
188 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
189 if binding.stage_sequence() != sequence || &binding.digest() != expected {
190 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
191 }
192 let plan = read_operation_plan(layout, binding.plan())?;
193 binding.validate(workflow, &plan)?;
194 Ok((binding, plan))
195}
196fn validate_predecessors(
197 layout: &BackupLayoutGuard,
198 workflow: &ExecutionWorkflowRecord,
199 binding: &ExecutionStageBindingRecord,
200) -> Result<(), ExecutionWorkflowPersistenceError> {
201 for row in binding.predecessors() {
204 let predecessor = acquire_stage(layout, row.stage_sequence())?;
205 let (original, _) =
206 read_binding(&predecessor, workflow, row.stage_sequence(), row.binding())?;
207 read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
208 }
209 Ok(())
210}
211fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
212 layout.root().join(format!("execution-stage-{sequence}"))
213}
214fn acquire_stage(
215 layout: &BackupLayoutGuard,
216 sequence: u64,
217) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
218 layout.check_root()?;
219 let path = stage_path(layout, sequence);
220 let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
221 if !metadata.is_dir() {
222 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
223 }
224 let stage = BackupLayoutGuard::acquire(&path)?;
225 if stage.root() != path {
228 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
229 }
230 layout.check_root()?;
231 Ok(stage)
232}
233fn create_private_directory(
234 path: &std::path::Path,
235) -> Result<(), ExecutionWorkflowPersistenceError> {
236 #[cfg(unix)]
237 {
238 use std::os::unix::fs::DirBuilderExt;
239 fs::DirBuilder::new()
240 .mode(0o700)
241 .create(path)
242 .map_err(PersistenceError::from)?;
243 Ok(())
244 }
245 #[cfg(not(unix))]
246 {
247 let _ = path;
248 Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
249 }
250}
251
252#[derive(Debug, Error)]
254pub enum ExecutionWorkflowPersistenceError {
255 #[error("execution workflow/stage digest mismatch")]
257 DigestMismatch,
258 #[error("unsafe execution stage directory")]
260 UnsafeStage,
261 #[error(transparent)]
263 Model(#[from] ExecutionWorkflowError),
264 #[error(transparent)]
266 Plan(#[from] OperationPlanPersistenceError),
267 #[error(transparent)]
269 Settlement(#[from] ExecutionSettlementPersistenceError),
270 #[error(transparent)]
272 Lock(#[from] JournalLockError),
273 #[error(transparent)]
275 Persistence(#[from] PersistenceError),
276}
277
278#[cfg(all(test, unix))]
279mod tests;