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    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
35/// Durably create the full original workflow allocation without replacement.
36/// # Errors
37/// Rejects occupied/unsafe paths, bounds, ownership and publication failures.
38pub 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}
49/// Read exact original workflow allocation; missing evidence never supplies a new budget.
50/// # Errors
51/// Rejects missing/changed/unsafe/oversized records and ownership failures.
52pub 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/// One original stage's fixed layout, retained under both workflow and child exclusion.
68///
69/// Existing attempt journals remain the only spending owner. `create` supplies no
70/// journals; `prepare` creates the complete original set. Reopening never creates,
71/// resets or reconstructs journals. Integrations
72/// retain exact learned inputs, authenticate their relation to settled originals,
73/// and qualify fresh authority, command custody and proof of no previous dispatch.
74#[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    /// Resume an exact retained stage with its complete original journal progress.
83    ///
84    /// Reuses record-only `open` and canonical complete-journal admission, then
85    /// rechecks the original stage and ancestor histories before returning. Hold
86    /// no attempt guards. Reads only retained local records, without creating
87    /// journals, reading artifact trees or calling providers. Pending and exhausted
88    /// spending is retained.
89    /// The returned progress is a sequential local projection, not atomic custody,
90    /// fresh permission, dispatch, authenticated receipt or terminal/release proof.
91    /// # Errors
92    /// Rejects missing/changed/unsafe originals, journal contention and invalid
93    /// retained accounting/causality without repair or replenishing allowance.
94    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    /// Durably prepare a new stage and its complete original attempt-journal set.
111    ///
112    /// Reuses create-only stage admission and the existing journal owner, taking
113    /// one journal lock at a time. Every original authority is derived before
114    /// stage allocation. No reservation, receipt, backend call or fresh dispatch
115    /// authority follows. Hold no other attempt guards during preparation.
116    ///
117    /// Partial preparation remains occupied and is never repaired by retrying.
118    /// Reopen only retained originals and require complete execution admission;
119    /// missing journals never mean unused allowance. A lost successful response
120    /// can reopen the complete exact stage without creating another journal.
121    /// # Errors
122    /// Rejects invalid originals, occupied stages, journal contention and failed
123    /// durable publication, preserving all partial records and original limits.
124    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    /// Create a fixed private stage directory, exact child plan and immutable binding.
149    ///
150    /// Every original predecessor and ancestor needs its exact retained binding, plan and
151    /// complete chronological Applied settlement. The directory is create-only:
152    /// occupied or interrupted preparation is retained and cannot be recreated
153    /// with another plan. No journal or backend effect occurs before return.
154    /// # Errors
155    /// Rejects changed allocation, unavailable prerequisites, occupied stages and IO failures.
156    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        // Validate both documents before allocating any stage paths.
173        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    /// Reopen an exact fixed stage without allocation, repair or journal creation.
197    ///
198    /// Lost successful creation responses reconcile here. A partially prepared
199    /// directory without both immutable records rejects and remains retained.
200    /// Missing pending journals cannot be treated as unused allowance by the
201    /// existing execution-progress or settlement owners.
202    /// # Errors
203    /// Rejects missing/changed original records, prerequisite history or ownership failures.
204    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    /// Read the exact bound original child plan; no fresh permission is implied.
226    #[must_use]
227    pub const fn plan(&self) -> &OperationPlanRecord {
228        &self.plan
229    }
230    /// Read the immutable original stage and predecessor commitments.
231    #[must_use]
232    pub const fn binding(&self) -> &ExecutionStageBindingRecord {
233        &self.binding
234    }
235    /// Re-admit both original records and expose the existing journal/layout owner.
236    ///
237    /// Callers must use complete original journal admission for resume. Never
238    /// create a missing journal on resume or substitute a freshly derived plan.
239    /// # Errors
240    /// Rejects changed roots, retained records or any ancestor settlement history.
241    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/// Failed create-only preparation; retained records never grant repair or extra spending.
255#[derive(Debug, Error)]
256pub enum ExecutionStagePreparationError {
257    /// Original stage records, predecessor settlement or layout were rejected.
258    #[error(transparent)]
259    Stage(#[from] ExecutionWorkflowPersistenceError),
260    /// Complete original child authority derivation failed before allocation.
261    #[error(transparent)]
262    Authority(#[from] OperationPlanError),
263    /// Existing attempt-journal locking or durable create-only publication failed.
264    #[error(transparent)]
265    Journal(#[from] AttemptJournalError),
266}
267
268/// Complete retained stage-resume refusal; original evidence and spending remain unchanged.
269#[derive(Debug, Error)]
270pub enum ExecutionStageResumeError {
271    /// Original stage, ancestor history or layout cannot be admitted.
272    #[error(transparent)]
273    Stage(#[from] ExecutionWorkflowPersistenceError),
274    /// Complete original child journals, accounting or causality cannot be admitted.
275    #[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    // Traverse the original acyclic catalog iteratively, admitting each ancestor
301    // once. Keep at most one row per catalog node, not one per dependency path.
302    // Learned evidence may differ between edges; binding/settlement identity cannot.
303    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            // Release this layout and every journal lock before admitting another
313            // ancestor, including an ancestor shared by multiple branches.
314            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    // Unlike an explicitly operator-selected root, a derived stage symlink is
352    // never adopted. Noncooperating namespace stability remains caller-owned.
353    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/// Typed immutable workflow/stage admission failure; no recovery evidence is removed.
379#[derive(Debug, Error)]
380pub enum ExecutionWorkflowPersistenceError {
381    /// Workflow/stage identity differs from the retained expected digest.
382    #[error("execution workflow/stage digest mismatch")]
383    DigestMismatch,
384    /// A derived stage directory is a symlink or another unsafe entry.
385    #[error("unsafe execution stage directory")]
386    UnsafeStage,
387    /// Original model allocation or learned binding was rejected.
388    #[error(transparent)]
389    Model(#[from] ExecutionWorkflowError),
390    /// Existing immutable operation-plan admission failed.
391    #[error(transparent)]
392    Plan(#[from] OperationPlanPersistenceError),
393    /// Existing complete chronological Applied settlement admission failed.
394    #[error(transparent)]
395    Settlement(#[from] ExecutionSettlementPersistenceError),
396    /// Cooperating layout or journal exclusion failed.
397    #[error(transparent)]
398    Lock(#[from] JournalLockError),
399    /// Existing bounded JSON/durable filesystem owner failed.
400    #[error(transparent)]
401    Persistence(#[from] PersistenceError),
402}
403
404#[cfg(all(test, unix))]
405mod tests;