1use super::{
4 AttemptJournalError, AttemptJournalGuard, BackupLayoutGuard, ExecutionProgressPersistenceError,
5 ExecutionSettlementPersistenceError, JournalLock, JournalLockError,
6 OperationPlanPersistenceError, PersistenceError, create_json_durable, create_operation_plan,
7 read_execution_progress, 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 crate::policy::execution_progress::ExecutionProgressView;
18use std::{
19 collections::{BTreeMap, btree_map::Entry},
20 fs,
21 path::PathBuf,
22};
23use thiserror::Error;
24
25const WORKFLOW_FILE: &str = "execution-workflow.json";
26const BINDING_FILE: &str = "stage-binding.json";
27
28#[derive(Clone, Copy, Debug, Eq, PartialEq)]
29enum StagePreparationBarrier {
30 Plan,
31 Binding,
32 Journal(u64),
33}
34
35pub fn create_execution_workflow(
39 layout: &BackupLayoutGuard,
40 record: &ExecutionWorkflowRecord,
41) -> Result<(), ExecutionWorkflowPersistenceError> {
42 layout.check_root()?;
43 super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
44 let path = layout.root().join(WORKFLOW_FILE);
45 let _lock = JournalLock::acquire(&path)?;
46 create_json_durable(&path, record)?;
47 Ok(())
48}
49pub fn read_execution_workflow(
53 layout: &BackupLayoutGuard,
54 expected: &ArtifactChecksumRecord,
55) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
56 layout.check_root()?;
57 let path = layout.root().join(WORKFLOW_FILE);
58 let _lock = JournalLock::acquire(&path)?;
59 let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
60 super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
61 if &record.digest() != expected {
62 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
63 }
64 Ok(record)
65}
66
67#[derive(Debug)]
75pub struct ExecutionStageGuard<'a> {
76 workflow_layout: &'a BackupLayoutGuard,
77 stage_layout: BackupLayoutGuard,
78 binding: ExecutionStageBindingRecord,
79 plan: OperationPlanRecord,
80}
81impl<'a> ExecutionStageGuard<'a> {
82 pub fn resume(
95 workflow_layout: &'a BackupLayoutGuard,
96 expected_workflow: &ArtifactChecksumRecord,
97 sequence: u64,
98 expected_binding: &ArtifactChecksumRecord,
99 ) -> Result<(Self, ExecutionProgressView), ExecutionStageResumeError> {
100 let stage = Self::open(
101 workflow_layout,
102 expected_workflow,
103 sequence,
104 expected_binding,
105 )?;
106 let progress = read_execution_progress(&stage.stage_layout, &stage.plan.digest())?;
107 stage.layout()?;
108 Ok((stage, progress))
109 }
110 pub fn prepare(
125 workflow_layout: &'a BackupLayoutGuard,
126 binding: ExecutionStageBindingRecord,
127 plan: OperationPlanRecord,
128 ) -> Result<Self, ExecutionStagePreparationError> {
129 Self::prepare_with(workflow_layout, binding, plan, |_| {})
130 }
131 fn prepare_with(
132 workflow_layout: &'a BackupLayoutGuard,
133 binding: ExecutionStageBindingRecord,
134 plan: OperationPlanRecord,
135 mut barrier: impl FnMut(StagePreparationBarrier),
136 ) -> Result<Self, ExecutionStagePreparationError> {
137 let authorities = plan.attempt_authorities()?;
138 let stage = Self::create_with(workflow_layout, binding, plan, &mut barrier)?;
139 let layout = stage.layout()?;
140 for authority in authorities {
141 let sequence = authority.binding().operation_sequence();
142 drop(AttemptJournalGuard::create(layout, authority)?);
143 barrier(StagePreparationBarrier::Journal(sequence));
144 }
145 stage.layout()?;
146 Ok(stage)
147 }
148 pub fn create(
157 workflow_layout: &'a BackupLayoutGuard,
158 binding: ExecutionStageBindingRecord,
159 plan: OperationPlanRecord,
160 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
161 Self::create_with(workflow_layout, binding, plan, |_| {})
162 }
163 fn create_with(
164 workflow_layout: &'a BackupLayoutGuard,
165 binding: ExecutionStageBindingRecord,
166 plan: OperationPlanRecord,
167 mut barrier: impl FnMut(StagePreparationBarrier),
168 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
169 let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
170 binding.validate(&workflow, &plan)?;
171 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
172 super::json::check_json_size(
174 &plan,
175 crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
176 )?;
177 validate_predecessors(workflow_layout, &workflow, &binding)?;
178 let path = stage_path(workflow_layout, binding.stage_sequence());
179 create_private_directory(&path)?;
180 fs::File::open(workflow_layout.root())
181 .and_then(|directory| directory.sync_all())
182 .map_err(PersistenceError::from)?;
183 let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
184 create_operation_plan(&stage_layout, &plan)?;
185 barrier(StagePreparationBarrier::Plan);
186 create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
187 barrier(StagePreparationBarrier::Binding);
188 workflow_layout.check_root()?;
189 Ok(Self {
190 workflow_layout,
191 stage_layout,
192 binding,
193 plan,
194 })
195 }
196 pub fn open(
205 workflow_layout: &'a BackupLayoutGuard,
206 expected_workflow: &ArtifactChecksumRecord,
207 sequence: u64,
208 expected_binding: &ArtifactChecksumRecord,
209 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
210 let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
211 workflow
212 .stage(sequence)
213 .map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
214 let stage_layout = acquire_stage(workflow_layout, sequence)?;
215 let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
216 validate_predecessors(workflow_layout, &workflow, &binding)?;
217 workflow_layout.check_root()?;
218 Ok(Self {
219 workflow_layout,
220 stage_layout,
221 binding,
222 plan,
223 })
224 }
225 #[must_use]
227 pub const fn plan(&self) -> &OperationPlanRecord {
228 &self.plan
229 }
230 #[must_use]
232 pub const fn binding(&self) -> &ExecutionStageBindingRecord {
233 &self.binding
234 }
235 pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
242 let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
243 read_binding(
244 &self.stage_layout,
245 &workflow,
246 self.binding.stage_sequence(),
247 &self.binding.digest(),
248 )?;
249 validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
250 Ok(&self.stage_layout)
251 }
252}
253
254#[derive(Debug, Error)]
256pub enum ExecutionStagePreparationError {
257 #[error(transparent)]
259 Stage(#[from] ExecutionWorkflowPersistenceError),
260 #[error(transparent)]
262 Authority(#[from] OperationPlanError),
263 #[error(transparent)]
265 Journal(#[from] AttemptJournalError),
266}
267
268#[derive(Debug, Error)]
270pub enum ExecutionStageResumeError {
271 #[error(transparent)]
273 Stage(#[from] ExecutionWorkflowPersistenceError),
274 #[error(transparent)]
276 Progress(#[from] ExecutionProgressPersistenceError),
277}
278fn read_binding(
279 layout: &BackupLayoutGuard,
280 workflow: &ExecutionWorkflowRecord,
281 sequence: u64,
282 expected: &ArtifactChecksumRecord,
283) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
284 layout.check_root()?;
285 let binding: ExecutionStageBindingRecord =
286 read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
287 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
288 if binding.stage_sequence() != sequence || &binding.digest() != expected {
289 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
290 }
291 let plan = read_operation_plan(layout, binding.plan())?;
292 binding.validate(workflow, &plan)?;
293 Ok((binding, plan))
294}
295fn validate_predecessors(
296 layout: &BackupLayoutGuard,
297 workflow: &ExecutionWorkflowRecord,
298 binding: &ExecutionStageBindingRecord,
299) -> Result<(), ExecutionWorkflowPersistenceError> {
300 let mut expected: BTreeMap<u64, ExecutionStagePredecessorRecord> = binding
304 .predecessors()
305 .iter()
306 .map(|row| (row.stage_sequence(), row.clone()))
307 .collect();
308 let mut pending: Vec<_> = expected.keys().copied().collect();
309 while let Some(sequence) = pending.pop() {
310 let row = &expected[&sequence];
311 let original = {
312 let predecessor = acquire_stage(layout, sequence)?;
315 let (original, _) = read_binding(&predecessor, workflow, sequence, row.binding())?;
316 read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
317 original
318 };
319 for ancestor in original.predecessors() {
320 match expected.entry(ancestor.stage_sequence()) {
321 Entry::Vacant(entry) => {
322 pending.push(ancestor.stage_sequence());
323 entry.insert(ancestor.clone());
324 }
325 Entry::Occupied(entry) => {
326 if entry.get().binding() != ancestor.binding()
327 || entry.get().settlement() != ancestor.settlement()
328 {
329 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
330 }
331 }
332 }
333 }
334 }
335 Ok(())
336}
337fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
338 layout.root().join(format!("execution-stage-{sequence}"))
339}
340fn acquire_stage(
341 layout: &BackupLayoutGuard,
342 sequence: u64,
343) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
344 layout.check_root()?;
345 let path = stage_path(layout, sequence);
346 let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
347 if !metadata.is_dir() {
348 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
349 }
350 let stage = BackupLayoutGuard::acquire(&path)?;
351 if stage.root() != path {
354 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
355 }
356 layout.check_root()?;
357 Ok(stage)
358}
359fn create_private_directory(
360 path: &std::path::Path,
361) -> Result<(), ExecutionWorkflowPersistenceError> {
362 #[cfg(unix)]
363 {
364 use std::os::unix::fs::DirBuilderExt;
365 fs::DirBuilder::new()
366 .mode(0o700)
367 .create(path)
368 .map_err(PersistenceError::from)?;
369 Ok(())
370 }
371 #[cfg(not(unix))]
372 {
373 let _ = path;
374 Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
375 }
376}
377
378#[derive(Debug, Error)]
380pub enum ExecutionWorkflowPersistenceError {
381 #[error("execution workflow/stage digest mismatch")]
383 DigestMismatch,
384 #[error("unsafe execution stage directory")]
386 UnsafeStage,
387 #[error(transparent)]
389 Model(#[from] ExecutionWorkflowError),
390 #[error(transparent)]
392 Plan(#[from] OperationPlanPersistenceError),
393 #[error(transparent)]
395 Settlement(#[from] ExecutionSettlementPersistenceError),
396 #[error(transparent)]
398 Lock(#[from] JournalLockError),
399 #[error(transparent)]
401 Persistence(#[from] PersistenceError),
402}
403
404#[cfg(all(test, unix))]
405mod tests;