1use super::{
4 AttemptJournalError, AttemptJournalGuard, BackupLayoutGuard,
5 ExecutionSettlementPersistenceError, JournalLock, JournalLockError,
6 OperationPlanPersistenceError, PersistenceError, create_json_durable, create_operation_plan,
7 read_execution_settlement, read_json, read_operation_plan,
8};
9use crate::model::{
10 artifacts::ArtifactChecksumRecord,
11 execution_workflow::{
12 ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
13 ExecutionWorkflowRecord, MAX_EXECUTION_STAGE_BYTES, MAX_EXECUTION_WORKFLOW_BYTES,
14 },
15 operation_plan::{OperationPlanError, OperationPlanRecord},
16};
17use std::{
18 collections::{BTreeMap, btree_map::Entry},
19 fs,
20 path::PathBuf,
21};
22use thiserror::Error;
23
24const WORKFLOW_FILE: &str = "execution-workflow.json";
25const BINDING_FILE: &str = "stage-binding.json";
26
27#[derive(Clone, Copy, Debug, Eq, PartialEq)]
28enum StagePreparationBarrier {
29 Plan,
30 Binding,
31 Journal(u64),
32}
33
34pub fn create_execution_workflow(
38 layout: &BackupLayoutGuard,
39 record: &ExecutionWorkflowRecord,
40) -> Result<(), ExecutionWorkflowPersistenceError> {
41 layout.check_root()?;
42 super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
43 let path = layout.root().join(WORKFLOW_FILE);
44 let _lock = JournalLock::acquire(&path)?;
45 create_json_durable(&path, record)?;
46 Ok(())
47}
48pub fn read_execution_workflow(
52 layout: &BackupLayoutGuard,
53 expected: &ArtifactChecksumRecord,
54) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
55 layout.check_root()?;
56 let path = layout.root().join(WORKFLOW_FILE);
57 let _lock = JournalLock::acquire(&path)?;
58 let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
59 super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
60 if &record.digest() != expected {
61 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
62 }
63 Ok(record)
64}
65
66#[derive(Debug)]
74pub struct ExecutionStageGuard<'a> {
75 workflow_layout: &'a BackupLayoutGuard,
76 stage_layout: BackupLayoutGuard,
77 binding: ExecutionStageBindingRecord,
78 plan: OperationPlanRecord,
79}
80impl<'a> ExecutionStageGuard<'a> {
81 pub fn prepare(
96 workflow_layout: &'a BackupLayoutGuard,
97 binding: ExecutionStageBindingRecord,
98 plan: OperationPlanRecord,
99 ) -> Result<Self, ExecutionStagePreparationError> {
100 Self::prepare_with(workflow_layout, binding, plan, |_| {})
101 }
102 fn prepare_with(
103 workflow_layout: &'a BackupLayoutGuard,
104 binding: ExecutionStageBindingRecord,
105 plan: OperationPlanRecord,
106 mut barrier: impl FnMut(StagePreparationBarrier),
107 ) -> Result<Self, ExecutionStagePreparationError> {
108 let authorities = plan.attempt_authorities()?;
109 let stage = Self::create_with(workflow_layout, binding, plan, &mut barrier)?;
110 let layout = stage.layout()?;
111 for authority in authorities {
112 let sequence = authority.binding().operation_sequence();
113 drop(AttemptJournalGuard::create(layout, authority)?);
114 barrier(StagePreparationBarrier::Journal(sequence));
115 }
116 stage.layout()?;
117 Ok(stage)
118 }
119 pub fn create(
128 workflow_layout: &'a BackupLayoutGuard,
129 binding: ExecutionStageBindingRecord,
130 plan: OperationPlanRecord,
131 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
132 Self::create_with(workflow_layout, binding, plan, |_| {})
133 }
134 fn create_with(
135 workflow_layout: &'a BackupLayoutGuard,
136 binding: ExecutionStageBindingRecord,
137 plan: OperationPlanRecord,
138 mut barrier: impl FnMut(StagePreparationBarrier),
139 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
140 let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
141 binding.validate(&workflow, &plan)?;
142 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
143 super::json::check_json_size(
145 &plan,
146 crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
147 )?;
148 validate_predecessors(workflow_layout, &workflow, &binding)?;
149 let path = stage_path(workflow_layout, binding.stage_sequence());
150 create_private_directory(&path)?;
151 fs::File::open(workflow_layout.root())
152 .and_then(|directory| directory.sync_all())
153 .map_err(PersistenceError::from)?;
154 let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
155 create_operation_plan(&stage_layout, &plan)?;
156 barrier(StagePreparationBarrier::Plan);
157 create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
158 barrier(StagePreparationBarrier::Binding);
159 workflow_layout.check_root()?;
160 Ok(Self {
161 workflow_layout,
162 stage_layout,
163 binding,
164 plan,
165 })
166 }
167 pub fn open(
176 workflow_layout: &'a BackupLayoutGuard,
177 expected_workflow: &ArtifactChecksumRecord,
178 sequence: u64,
179 expected_binding: &ArtifactChecksumRecord,
180 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
181 let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
182 workflow
183 .stage(sequence)
184 .map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
185 let stage_layout = acquire_stage(workflow_layout, sequence)?;
186 let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
187 validate_predecessors(workflow_layout, &workflow, &binding)?;
188 workflow_layout.check_root()?;
189 Ok(Self {
190 workflow_layout,
191 stage_layout,
192 binding,
193 plan,
194 })
195 }
196 #[must_use]
198 pub const fn plan(&self) -> &OperationPlanRecord {
199 &self.plan
200 }
201 #[must_use]
203 pub const fn binding(&self) -> &ExecutionStageBindingRecord {
204 &self.binding
205 }
206 pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
213 let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
214 read_binding(
215 &self.stage_layout,
216 &workflow,
217 self.binding.stage_sequence(),
218 &self.binding.digest(),
219 )?;
220 validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
221 Ok(&self.stage_layout)
222 }
223}
224
225#[derive(Debug, Error)]
227pub enum ExecutionStagePreparationError {
228 #[error(transparent)]
230 Stage(#[from] ExecutionWorkflowPersistenceError),
231 #[error(transparent)]
233 Authority(#[from] OperationPlanError),
234 #[error(transparent)]
236 Journal(#[from] AttemptJournalError),
237}
238fn read_binding(
239 layout: &BackupLayoutGuard,
240 workflow: &ExecutionWorkflowRecord,
241 sequence: u64,
242 expected: &ArtifactChecksumRecord,
243) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
244 layout.check_root()?;
245 let binding: ExecutionStageBindingRecord =
246 read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
247 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
248 if binding.stage_sequence() != sequence || &binding.digest() != expected {
249 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
250 }
251 let plan = read_operation_plan(layout, binding.plan())?;
252 binding.validate(workflow, &plan)?;
253 Ok((binding, plan))
254}
255fn validate_predecessors(
256 layout: &BackupLayoutGuard,
257 workflow: &ExecutionWorkflowRecord,
258 binding: &ExecutionStageBindingRecord,
259) -> Result<(), ExecutionWorkflowPersistenceError> {
260 let mut expected: BTreeMap<u64, ExecutionStagePredecessorRecord> = binding
264 .predecessors()
265 .iter()
266 .map(|row| (row.stage_sequence(), row.clone()))
267 .collect();
268 let mut pending: Vec<_> = expected.keys().copied().collect();
269 while let Some(sequence) = pending.pop() {
270 let row = &expected[&sequence];
271 let original = {
272 let predecessor = acquire_stage(layout, sequence)?;
275 let (original, _) = read_binding(&predecessor, workflow, sequence, row.binding())?;
276 read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
277 original
278 };
279 for ancestor in original.predecessors() {
280 match expected.entry(ancestor.stage_sequence()) {
281 Entry::Vacant(entry) => {
282 pending.push(ancestor.stage_sequence());
283 entry.insert(ancestor.clone());
284 }
285 Entry::Occupied(entry) => {
286 if entry.get().binding() != ancestor.binding()
287 || entry.get().settlement() != ancestor.settlement()
288 {
289 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
290 }
291 }
292 }
293 }
294 }
295 Ok(())
296}
297fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
298 layout.root().join(format!("execution-stage-{sequence}"))
299}
300fn acquire_stage(
301 layout: &BackupLayoutGuard,
302 sequence: u64,
303) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
304 layout.check_root()?;
305 let path = stage_path(layout, sequence);
306 let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
307 if !metadata.is_dir() {
308 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
309 }
310 let stage = BackupLayoutGuard::acquire(&path)?;
311 if stage.root() != path {
314 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
315 }
316 layout.check_root()?;
317 Ok(stage)
318}
319fn create_private_directory(
320 path: &std::path::Path,
321) -> Result<(), ExecutionWorkflowPersistenceError> {
322 #[cfg(unix)]
323 {
324 use std::os::unix::fs::DirBuilderExt;
325 fs::DirBuilder::new()
326 .mode(0o700)
327 .create(path)
328 .map_err(PersistenceError::from)?;
329 Ok(())
330 }
331 #[cfg(not(unix))]
332 {
333 let _ = path;
334 Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
335 }
336}
337
338#[derive(Debug, Error)]
340pub enum ExecutionWorkflowPersistenceError {
341 #[error("execution workflow/stage digest mismatch")]
343 DigestMismatch,
344 #[error("unsafe execution stage directory")]
346 UnsafeStage,
347 #[error(transparent)]
349 Model(#[from] ExecutionWorkflowError),
350 #[error(transparent)]
352 Plan(#[from] OperationPlanPersistenceError),
353 #[error(transparent)]
355 Settlement(#[from] ExecutionSettlementPersistenceError),
356 #[error(transparent)]
358 Lock(#[from] JournalLockError),
359 #[error(transparent)]
361 Persistence(#[from] PersistenceError),
362}
363
364#[cfg(all(test, unix))]
365mod tests;