1use super::{
4 AttemptJournalError, AttemptJournalGuard, BackupLayoutGuard, ExecutionProgressPersistenceError,
5 ExecutionSettlementCheckpointError, ExecutionSettlementPersistenceError, JournalLock,
6 JournalLockError, OperationPlanPersistenceError, PersistenceError,
7 checkpoint_execution_settlement, create_json_durable, create_operation_plan,
8 read_execution_progress, read_execution_settlement, read_json, read_operation_plan,
9};
10use crate::model::{
11 artifacts::ArtifactChecksumRecord,
12 execution_workflow::{
13 ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
14 ExecutionWorkflowRecord, MAX_EXECUTION_STAGE_BYTES, MAX_EXECUTION_WORKFLOW_BYTES,
15 },
16 operation_plan::{OperationPlanError, OperationPlanRecord},
17};
18use crate::policy::execution_progress::ExecutionProgressView;
19use std::{
20 collections::{BTreeMap, btree_map::Entry},
21 fs,
22 path::PathBuf,
23};
24use thiserror::Error;
25
26const WORKFLOW_FILE: &str = "execution-workflow.json";
27const BINDING_FILE: &str = "stage-binding.json";
28
29#[derive(Clone, Copy, Debug, Eq, PartialEq)]
30enum StagePreparationBarrier {
31 Plan,
32 Binding,
33 Journal(u64),
34}
35
36pub fn create_execution_workflow(
40 layout: &BackupLayoutGuard,
41 record: &ExecutionWorkflowRecord,
42) -> Result<(), ExecutionWorkflowPersistenceError> {
43 layout.check_root()?;
44 super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
45 let path = layout.root().join(WORKFLOW_FILE);
46 let _lock = JournalLock::acquire(&path)?;
47 create_json_durable(&path, record)?;
48 Ok(())
49}
50pub fn read_execution_workflow(
54 layout: &BackupLayoutGuard,
55 expected: &ArtifactChecksumRecord,
56) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
57 layout.check_root()?;
58 let path = layout.root().join(WORKFLOW_FILE);
59 let _lock = JournalLock::acquire(&path)?;
60 let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
61 super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
62 if &record.digest() != expected {
63 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
64 }
65 Ok(record)
66}
67
68#[derive(Debug)]
76pub struct ExecutionStageGuard<'a> {
77 workflow_layout: &'a BackupLayoutGuard,
78 stage_layout: BackupLayoutGuard,
79 binding: ExecutionStageBindingRecord,
80 plan: OperationPlanRecord,
81}
82impl<'a> ExecutionStageGuard<'a> {
83 pub fn checkpoint(
96 &self,
97 learned_evidence: ArtifactChecksumRecord,
98 ) -> Result<ExecutionStagePredecessorRecord, ExecutionStageCheckpointError> {
99 self.checkpoint_with(learned_evidence, || {})
100 }
101 fn checkpoint_with(
102 &self,
103 learned_evidence: ArtifactChecksumRecord,
104 after_publication: impl FnOnce(),
105 ) -> Result<ExecutionStagePredecessorRecord, ExecutionStageCheckpointError> {
106 let settlement = checkpoint_execution_settlement(self.layout()?, &self.plan.digest())?;
107 after_publication();
108 self.layout()?;
109 Ok(ExecutionStagePredecessorRecord::new(
110 self.binding.stage_sequence(),
111 self.binding.digest(),
112 settlement.digest(),
113 learned_evidence,
114 ))
115 }
116 pub fn resume(
129 workflow_layout: &'a BackupLayoutGuard,
130 expected_workflow: &ArtifactChecksumRecord,
131 sequence: u64,
132 expected_binding: &ArtifactChecksumRecord,
133 ) -> Result<(Self, ExecutionProgressView), ExecutionStageResumeError> {
134 let stage = Self::open(
135 workflow_layout,
136 expected_workflow,
137 sequence,
138 expected_binding,
139 )?;
140 let progress = read_execution_progress(&stage.stage_layout, &stage.plan.digest())?;
141 stage.layout()?;
142 Ok((stage, progress))
143 }
144 pub fn prepare(
159 workflow_layout: &'a BackupLayoutGuard,
160 binding: ExecutionStageBindingRecord,
161 plan: OperationPlanRecord,
162 ) -> Result<Self, ExecutionStagePreparationError> {
163 Self::prepare_with(workflow_layout, binding, plan, |_| {})
164 }
165 fn prepare_with(
166 workflow_layout: &'a BackupLayoutGuard,
167 binding: ExecutionStageBindingRecord,
168 plan: OperationPlanRecord,
169 mut barrier: impl FnMut(StagePreparationBarrier),
170 ) -> Result<Self, ExecutionStagePreparationError> {
171 let authorities = plan.attempt_authorities()?;
172 let stage = Self::create_with(workflow_layout, binding, plan, &mut barrier)?;
173 let layout = stage.layout()?;
174 for authority in authorities {
175 let sequence = authority.binding().operation_sequence();
176 drop(AttemptJournalGuard::create(layout, authority)?);
177 barrier(StagePreparationBarrier::Journal(sequence));
178 }
179 stage.layout()?;
180 Ok(stage)
181 }
182 pub fn create(
191 workflow_layout: &'a BackupLayoutGuard,
192 binding: ExecutionStageBindingRecord,
193 plan: OperationPlanRecord,
194 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
195 Self::create_with(workflow_layout, binding, plan, |_| {})
196 }
197 fn create_with(
198 workflow_layout: &'a BackupLayoutGuard,
199 binding: ExecutionStageBindingRecord,
200 plan: OperationPlanRecord,
201 mut barrier: impl FnMut(StagePreparationBarrier),
202 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
203 let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
204 binding.validate(&workflow, &plan)?;
205 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
206 super::json::check_json_size(
208 &plan,
209 crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
210 )?;
211 validate_predecessors(workflow_layout, &workflow, &binding)?;
212 let path = stage_path(workflow_layout, binding.stage_sequence());
213 create_private_directory(&path)?;
214 fs::File::open(workflow_layout.root())
215 .and_then(|directory| directory.sync_all())
216 .map_err(PersistenceError::from)?;
217 let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
218 create_operation_plan(&stage_layout, &plan)?;
219 barrier(StagePreparationBarrier::Plan);
220 create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
221 barrier(StagePreparationBarrier::Binding);
222 workflow_layout.check_root()?;
223 Ok(Self {
224 workflow_layout,
225 stage_layout,
226 binding,
227 plan,
228 })
229 }
230 pub fn open(
239 workflow_layout: &'a BackupLayoutGuard,
240 expected_workflow: &ArtifactChecksumRecord,
241 sequence: u64,
242 expected_binding: &ArtifactChecksumRecord,
243 ) -> Result<Self, ExecutionWorkflowPersistenceError> {
244 let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
245 workflow
246 .stage(sequence)
247 .map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
248 let stage_layout = acquire_stage(workflow_layout, sequence)?;
249 let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
250 validate_predecessors(workflow_layout, &workflow, &binding)?;
251 workflow_layout.check_root()?;
252 Ok(Self {
253 workflow_layout,
254 stage_layout,
255 binding,
256 plan,
257 })
258 }
259 #[must_use]
261 pub const fn plan(&self) -> &OperationPlanRecord {
262 &self.plan
263 }
264 #[must_use]
266 pub const fn binding(&self) -> &ExecutionStageBindingRecord {
267 &self.binding
268 }
269 pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
276 let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
277 read_binding(
278 &self.stage_layout,
279 &workflow,
280 self.binding.stage_sequence(),
281 &self.binding.digest(),
282 )?;
283 validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
284 Ok(&self.stage_layout)
285 }
286}
287
288#[derive(Debug, Error)]
290pub enum ExecutionStageCheckpointError {
291 #[error(transparent)]
293 Stage(#[from] ExecutionWorkflowPersistenceError),
294 #[error(transparent)]
296 Settlement(#[from] ExecutionSettlementCheckpointError),
297}
298
299#[derive(Debug, Error)]
301pub enum ExecutionStagePreparationError {
302 #[error(transparent)]
304 Stage(#[from] ExecutionWorkflowPersistenceError),
305 #[error(transparent)]
307 Authority(#[from] OperationPlanError),
308 #[error(transparent)]
310 Journal(#[from] AttemptJournalError),
311}
312
313#[derive(Debug, Error)]
315pub enum ExecutionStageResumeError {
316 #[error(transparent)]
318 Stage(#[from] ExecutionWorkflowPersistenceError),
319 #[error(transparent)]
321 Progress(#[from] ExecutionProgressPersistenceError),
322}
323fn read_binding(
324 layout: &BackupLayoutGuard,
325 workflow: &ExecutionWorkflowRecord,
326 sequence: u64,
327 expected: &ArtifactChecksumRecord,
328) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
329 layout.check_root()?;
330 let binding: ExecutionStageBindingRecord =
331 read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
332 super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
333 if binding.stage_sequence() != sequence || &binding.digest() != expected {
334 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
335 }
336 let plan = read_operation_plan(layout, binding.plan())?;
337 binding.validate(workflow, &plan)?;
338 Ok((binding, plan))
339}
340fn validate_predecessors(
341 layout: &BackupLayoutGuard,
342 workflow: &ExecutionWorkflowRecord,
343 binding: &ExecutionStageBindingRecord,
344) -> Result<(), ExecutionWorkflowPersistenceError> {
345 let mut expected: BTreeMap<u64, ExecutionStagePredecessorRecord> = binding
349 .predecessors()
350 .iter()
351 .map(|row| (row.stage_sequence(), row.clone()))
352 .collect();
353 let mut pending: Vec<_> = expected.keys().copied().collect();
354 while let Some(sequence) = pending.pop() {
355 let row = &expected[&sequence];
356 let original = {
357 let predecessor = acquire_stage(layout, sequence)?;
360 let (original, _) = read_binding(&predecessor, workflow, sequence, row.binding())?;
361 read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
362 original
363 };
364 for ancestor in original.predecessors() {
365 match expected.entry(ancestor.stage_sequence()) {
366 Entry::Vacant(entry) => {
367 pending.push(ancestor.stage_sequence());
368 entry.insert(ancestor.clone());
369 }
370 Entry::Occupied(entry) => {
371 if entry.get().binding() != ancestor.binding()
372 || entry.get().settlement() != ancestor.settlement()
373 {
374 return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
375 }
376 }
377 }
378 }
379 }
380 Ok(())
381}
382fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
383 layout.root().join(format!("execution-stage-{sequence}"))
384}
385fn acquire_stage(
386 layout: &BackupLayoutGuard,
387 sequence: u64,
388) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
389 layout.check_root()?;
390 let path = stage_path(layout, sequence);
391 let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
392 if !metadata.is_dir() {
393 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
394 }
395 let stage = BackupLayoutGuard::acquire(&path)?;
396 if stage.root() != path {
399 return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
400 }
401 layout.check_root()?;
402 Ok(stage)
403}
404fn create_private_directory(
405 path: &std::path::Path,
406) -> Result<(), ExecutionWorkflowPersistenceError> {
407 #[cfg(unix)]
408 {
409 use std::os::unix::fs::DirBuilderExt;
410 fs::DirBuilder::new()
411 .mode(0o700)
412 .create(path)
413 .map_err(PersistenceError::from)?;
414 Ok(())
415 }
416 #[cfg(not(unix))]
417 {
418 let _ = path;
419 Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
420 }
421}
422
423#[derive(Debug, Error)]
425pub enum ExecutionWorkflowPersistenceError {
426 #[error("execution workflow/stage digest mismatch")]
428 DigestMismatch,
429 #[error("unsafe execution stage directory")]
431 UnsafeStage,
432 #[error(transparent)]
434 Model(#[from] ExecutionWorkflowError),
435 #[error(transparent)]
437 Plan(#[from] OperationPlanPersistenceError),
438 #[error(transparent)]
440 Settlement(#[from] ExecutionSettlementPersistenceError),
441 #[error(transparent)]
443 Lock(#[from] JournalLockError),
444 #[error(transparent)]
446 Persistence(#[from] PersistenceError),
447}
448
449#[cfg(all(test, unix))]
450mod tests;