Skip to main content

ic_backup/ops/persistence/execution_workflow/
mod.rs

1//! Fixed immutable stage admission using the existing layout, plan and settlement owners.
2
3use 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
36/// Durably create the full original workflow allocation without replacement.
37/// # Errors
38/// Rejects occupied/unsafe paths, bounds, ownership and publication failures.
39pub 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}
50/// Read exact original workflow allocation; missing evidence never supplies a new budget.
51/// # Errors
52/// Rejects missing/changed/unsafe/oversized records and ownership failures.
53pub 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/// One original stage's fixed layout, retained under both workflow and child exclusion.
69///
70/// Existing attempt journals remain the only spending owner. `create` supplies no
71/// journals; `prepare` creates the complete original set. Reopening never creates,
72/// resets or reconstructs journals. Integrations
73/// retain exact learned inputs, authenticate their relation to settled originals,
74/// and qualify fresh authority, command custody and proof of no previous dispatch.
75#[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    /// Publish the original all-Applied checkpoint and derive its exact predecessor commitment.
84    ///
85    /// Re-admits the workflow, stage and complete ancestor histories before and after
86    /// delegating to the canonical original-journal checkpoint owner. Hold no attempt
87    /// guards. Failed post-publication admission retains the checkpoint and all originals;
88    /// an occupied checkpoint is never replaced. Identity-bound local replay remains
89    /// owned by `read_execution_settlement`.
90    ///
91    /// The supplied learned evidence stays integration-owned: this authenticates no
92    /// learned input, receipt or effect, and grants no dispatch or terminal/release proof.
93    /// # Errors
94    /// Rejects changed/missing/held originals, unsettled journals and publication failures.
95    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    /// Resume an exact retained stage with its complete original journal progress.
117    ///
118    /// Reuses record-only `open` and canonical complete-journal admission, then
119    /// rechecks the original stage and ancestor histories before returning. Hold
120    /// no attempt guards. Reads only retained local records, without creating
121    /// journals, reading artifact trees or calling providers. Pending and exhausted
122    /// spending is retained.
123    /// The returned progress is a sequential local projection, not atomic custody,
124    /// fresh permission, dispatch, authenticated receipt or terminal/release proof.
125    /// # Errors
126    /// Rejects missing/changed/unsafe originals, journal contention and invalid
127    /// retained accounting/causality without repair or replenishing allowance.
128    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    /// Durably prepare a new stage and its complete original attempt-journal set.
145    ///
146    /// Reuses create-only stage admission and the existing journal owner, taking
147    /// one journal lock at a time. Every original authority is derived before
148    /// stage allocation. No reservation, receipt, backend call or fresh dispatch
149    /// authority follows. Hold no other attempt guards during preparation.
150    ///
151    /// Partial preparation remains occupied and is never repaired by retrying.
152    /// Reopen only retained originals and require complete execution admission;
153    /// missing journals never mean unused allowance. A lost successful response
154    /// can reopen the complete exact stage without creating another journal.
155    /// # Errors
156    /// Rejects invalid originals, occupied stages, journal contention and failed
157    /// durable publication, preserving all partial records and original limits.
158    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    /// Create a fixed private stage directory, exact child plan and immutable binding.
183    ///
184    /// Every original predecessor and ancestor needs its exact retained binding, plan and
185    /// complete chronological Applied settlement. The directory is create-only:
186    /// occupied or interrupted preparation is retained and cannot be recreated
187    /// with another plan. No journal or backend effect occurs before return.
188    /// # Errors
189    /// Rejects changed allocation, unavailable prerequisites, occupied stages and IO failures.
190    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        // Validate both documents before allocating any stage paths.
207        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    /// Reopen an exact fixed stage without allocation, repair or journal creation.
231    ///
232    /// Lost successful creation responses reconcile here. A partially prepared
233    /// directory without both immutable records rejects and remains retained.
234    /// Missing pending journals cannot be treated as unused allowance by the
235    /// existing execution-progress or settlement owners.
236    /// # Errors
237    /// Rejects missing/changed original records, prerequisite history or ownership failures.
238    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    /// Read the exact bound original child plan; no fresh permission is implied.
260    #[must_use]
261    pub const fn plan(&self) -> &OperationPlanRecord {
262        &self.plan
263    }
264    /// Read the immutable original stage and predecessor commitments.
265    #[must_use]
266    pub const fn binding(&self) -> &ExecutionStageBindingRecord {
267        &self.binding
268    }
269    /// Re-admit both original records and expose the existing journal/layout owner.
270    ///
271    /// Callers must use complete original journal admission for resume. Never
272    /// create a missing journal on resume or substitute a freshly derived plan.
273    /// # Errors
274    /// Rejects changed roots, retained records or any ancestor settlement history.
275    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/// Stage-bound checkpoint refusal; original spending and published evidence remain retained.
289#[derive(Debug, Error)]
290pub enum ExecutionStageCheckpointError {
291    /// Original workflow/stage or complete ancestor history cannot be admitted.
292    #[error(transparent)]
293    Stage(#[from] ExecutionWorkflowPersistenceError),
294    /// Complete original journal checkpoint derivation/publication failed.
295    #[error(transparent)]
296    Settlement(#[from] ExecutionSettlementCheckpointError),
297}
298
299/// Failed create-only preparation; retained records never grant repair or extra spending.
300#[derive(Debug, Error)]
301pub enum ExecutionStagePreparationError {
302    /// Original stage records, predecessor settlement or layout were rejected.
303    #[error(transparent)]
304    Stage(#[from] ExecutionWorkflowPersistenceError),
305    /// Complete original child authority derivation failed before allocation.
306    #[error(transparent)]
307    Authority(#[from] OperationPlanError),
308    /// Existing attempt-journal locking or durable create-only publication failed.
309    #[error(transparent)]
310    Journal(#[from] AttemptJournalError),
311}
312
313/// Complete retained stage-resume refusal; original evidence and spending remain unchanged.
314#[derive(Debug, Error)]
315pub enum ExecutionStageResumeError {
316    /// Original stage, ancestor history or layout cannot be admitted.
317    #[error(transparent)]
318    Stage(#[from] ExecutionWorkflowPersistenceError),
319    /// Complete original child journals, accounting or causality cannot be admitted.
320    #[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    // Traverse the original acyclic catalog iteratively, admitting each ancestor
346    // once. Keep at most one row per catalog node, not one per dependency path.
347    // Learned evidence may differ between edges; binding/settlement identity cannot.
348    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            // Release this layout and every journal lock before admitting another
358            // ancestor, including an ancestor shared by multiple branches.
359            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    // Unlike an explicitly operator-selected root, a derived stage symlink is
397    // never adopted. Noncooperating namespace stability remains caller-owned.
398    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/// Typed immutable workflow/stage admission failure; no recovery evidence is removed.
424#[derive(Debug, Error)]
425pub enum ExecutionWorkflowPersistenceError {
426    /// Workflow/stage identity differs from the retained expected digest.
427    #[error("execution workflow/stage digest mismatch")]
428    DigestMismatch,
429    /// A derived stage directory is a symlink or another unsafe entry.
430    #[error("unsafe execution stage directory")]
431    UnsafeStage,
432    /// Original model allocation or learned binding was rejected.
433    #[error(transparent)]
434    Model(#[from] ExecutionWorkflowError),
435    /// Existing immutable operation-plan admission failed.
436    #[error(transparent)]
437    Plan(#[from] OperationPlanPersistenceError),
438    /// Existing complete chronological Applied settlement admission failed.
439    #[error(transparent)]
440    Settlement(#[from] ExecutionSettlementPersistenceError),
441    /// Cooperating layout or journal exclusion failed.
442    #[error(transparent)]
443    Lock(#[from] JournalLockError),
444    /// Existing bounded JSON/durable filesystem owner failed.
445    #[error(transparent)]
446    Persistence(#[from] PersistenceError),
447}
448
449#[cfg(all(test, unix))]
450mod tests;