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    BackupLayoutGuard, ExecutionSettlementPersistenceError, JournalLock, JournalLockError,
5    OperationPlanPersistenceError, PersistenceError, create_json_durable, create_operation_plan,
6    read_execution_settlement, read_json, read_operation_plan,
7};
8use crate::model::{
9    artifacts::ArtifactChecksumRecord,
10    execution_workflow::{
11        ExecutionStageBindingRecord, ExecutionWorkflowError, ExecutionWorkflowRecord,
12        MAX_EXECUTION_STAGE_BYTES, MAX_EXECUTION_WORKFLOW_BYTES,
13    },
14    operation_plan::OperationPlanRecord,
15};
16use std::{fs, path::PathBuf};
17use thiserror::Error;
18
19const WORKFLOW_FILE: &str = "execution-workflow.json";
20const BINDING_FILE: &str = "stage-binding.json";
21
22#[derive(Clone, Copy, Debug, Eq, PartialEq)]
23enum StagePreparationBarrier {
24    AfterPlanPublication,
25    AfterBindingPublication,
26}
27
28/// Durably create the full original workflow allocation without replacement.
29/// # Errors
30/// Rejects occupied/unsafe paths, bounds, ownership and publication failures.
31pub fn create_execution_workflow(
32    layout: &BackupLayoutGuard,
33    record: &ExecutionWorkflowRecord,
34) -> Result<(), ExecutionWorkflowPersistenceError> {
35    layout.check_root()?;
36    super::json::check_json_size(record, MAX_EXECUTION_WORKFLOW_BYTES)?;
37    let path = layout.root().join(WORKFLOW_FILE);
38    let _lock = JournalLock::acquire(&path)?;
39    create_json_durable(&path, record)?;
40    Ok(())
41}
42/// Read exact original workflow allocation; missing evidence never supplies a new budget.
43/// # Errors
44/// Rejects missing/changed/unsafe/oversized records and ownership failures.
45pub fn read_execution_workflow(
46    layout: &BackupLayoutGuard,
47    expected: &ArtifactChecksumRecord,
48) -> Result<ExecutionWorkflowRecord, ExecutionWorkflowPersistenceError> {
49    layout.check_root()?;
50    let path = layout.root().join(WORKFLOW_FILE);
51    let _lock = JournalLock::acquire(&path)?;
52    let record: ExecutionWorkflowRecord = read_json(&path, MAX_EXECUTION_WORKFLOW_BYTES)?;
53    super::json::check_json_size(&record, MAX_EXECUTION_WORKFLOW_BYTES)?;
54    if &record.digest() != expected {
55        return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
56    }
57    Ok(record)
58}
59
60/// One original stage's fixed layout, retained under both workflow and child exclusion.
61///
62/// Existing attempt journals remain the only spending owner. Creation supplies no
63/// journals and reopening never creates, resets or reconstructs them. Integrations
64/// retain exact learned inputs, authenticate their relation to settled originals,
65/// and qualify fresh authority, command custody and proof of no previous dispatch.
66#[derive(Debug)]
67pub struct ExecutionStageGuard<'a> {
68    workflow_layout: &'a BackupLayoutGuard,
69    stage_layout: BackupLayoutGuard,
70    binding: ExecutionStageBindingRecord,
71    plan: OperationPlanRecord,
72}
73impl<'a> ExecutionStageGuard<'a> {
74    /// Create a fixed private stage directory, exact child plan and immutable binding.
75    ///
76    /// Every original predecessor must have its exact retained binding, plan and
77    /// complete chronological Applied settlement. The directory is create-only:
78    /// occupied or interrupted preparation is retained and cannot be recreated
79    /// with another plan. No journal or backend effect occurs before return.
80    /// # Errors
81    /// Rejects changed allocation, unavailable prerequisites, occupied stages and IO failures.
82    pub fn create(
83        workflow_layout: &'a BackupLayoutGuard,
84        binding: ExecutionStageBindingRecord,
85        plan: OperationPlanRecord,
86    ) -> Result<Self, ExecutionWorkflowPersistenceError> {
87        Self::create_with(workflow_layout, binding, plan, |_| {})
88    }
89    fn create_with(
90        workflow_layout: &'a BackupLayoutGuard,
91        binding: ExecutionStageBindingRecord,
92        plan: OperationPlanRecord,
93        mut barrier: impl FnMut(StagePreparationBarrier),
94    ) -> Result<Self, ExecutionWorkflowPersistenceError> {
95        let workflow = read_execution_workflow(workflow_layout, binding.workflow())?;
96        binding.validate(&workflow, &plan)?;
97        super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
98        // Validate both documents before allocating any stage paths.
99        super::json::check_json_size(
100            &plan,
101            crate::model::operation_plan::MAX_OPERATION_PLAN_BYTES,
102        )?;
103        validate_predecessors(workflow_layout, &workflow, &binding)?;
104        let path = stage_path(workflow_layout, binding.stage_sequence());
105        create_private_directory(&path)?;
106        fs::File::open(workflow_layout.root())
107            .and_then(|directory| directory.sync_all())
108            .map_err(PersistenceError::from)?;
109        let stage_layout = acquire_stage(workflow_layout, binding.stage_sequence())?;
110        create_operation_plan(&stage_layout, &plan)?;
111        barrier(StagePreparationBarrier::AfterPlanPublication);
112        create_json_durable(&stage_layout.root().join(BINDING_FILE), &binding)?;
113        barrier(StagePreparationBarrier::AfterBindingPublication);
114        workflow_layout.check_root()?;
115        Ok(Self {
116            workflow_layout,
117            stage_layout,
118            binding,
119            plan,
120        })
121    }
122    /// Reopen an exact fixed stage without allocation, repair or journal creation.
123    ///
124    /// Lost successful creation responses reconcile here. A partially prepared
125    /// directory without both immutable records rejects and remains retained.
126    /// Missing pending journals cannot be treated as unused allowance by the
127    /// existing execution-progress or settlement owners.
128    /// # Errors
129    /// Rejects missing/changed original records, prerequisite history or ownership failures.
130    pub fn open(
131        workflow_layout: &'a BackupLayoutGuard,
132        expected_workflow: &ArtifactChecksumRecord,
133        sequence: u64,
134        expected_binding: &ArtifactChecksumRecord,
135    ) -> Result<Self, ExecutionWorkflowPersistenceError> {
136        let workflow = read_execution_workflow(workflow_layout, expected_workflow)?;
137        workflow
138            .stage(sequence)
139            .map_err(|_| ExecutionWorkflowError::UnknownStage(sequence))?;
140        let stage_layout = acquire_stage(workflow_layout, sequence)?;
141        let (binding, plan) = read_binding(&stage_layout, &workflow, sequence, expected_binding)?;
142        validate_predecessors(workflow_layout, &workflow, &binding)?;
143        workflow_layout.check_root()?;
144        Ok(Self {
145            workflow_layout,
146            stage_layout,
147            binding,
148            plan,
149        })
150    }
151    /// Read the exact bound original child plan; no fresh permission is implied.
152    #[must_use]
153    pub const fn plan(&self) -> &OperationPlanRecord {
154        &self.plan
155    }
156    /// Read the immutable original stage and predecessor commitments.
157    #[must_use]
158    pub const fn binding(&self) -> &ExecutionStageBindingRecord {
159        &self.binding
160    }
161    /// Re-admit both original records and expose the existing journal/layout owner.
162    ///
163    /// Callers must use complete original journal admission for resume. Never
164    /// create a missing journal on resume or substitute a freshly derived plan.
165    /// # Errors
166    /// Rejects changed roots, retained records or predecessor settlement histories.
167    pub fn layout(&self) -> Result<&BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
168        let workflow = read_execution_workflow(self.workflow_layout, self.binding.workflow())?;
169        read_binding(
170            &self.stage_layout,
171            &workflow,
172            self.binding.stage_sequence(),
173            &self.binding.digest(),
174        )?;
175        validate_predecessors(self.workflow_layout, &workflow, &self.binding)?;
176        Ok(&self.stage_layout)
177    }
178}
179fn read_binding(
180    layout: &BackupLayoutGuard,
181    workflow: &ExecutionWorkflowRecord,
182    sequence: u64,
183    expected: &ArtifactChecksumRecord,
184) -> Result<(ExecutionStageBindingRecord, OperationPlanRecord), ExecutionWorkflowPersistenceError> {
185    layout.check_root()?;
186    let binding: ExecutionStageBindingRecord =
187        read_json(&layout.root().join(BINDING_FILE), MAX_EXECUTION_STAGE_BYTES)?;
188    super::json::check_json_size(&binding, MAX_EXECUTION_STAGE_BYTES)?;
189    if binding.stage_sequence() != sequence || &binding.digest() != expected {
190        return Err(ExecutionWorkflowPersistenceError::DigestMismatch);
191    }
192    let plan = read_operation_plan(layout, binding.plan())?;
193    binding.validate(workflow, &plan)?;
194    Ok((binding, plan))
195}
196fn validate_predecessors(
197    layout: &BackupLayoutGuard,
198    workflow: &ExecutionWorkflowRecord,
199    binding: &ExecutionStageBindingRecord,
200) -> Result<(), ExecutionWorkflowPersistenceError> {
201    // Sequential existing layout/journal admission avoids holding multiple
202    // predecessor locks or inventing a second outcome/accounting projection.
203    for row in binding.predecessors() {
204        let predecessor = acquire_stage(layout, row.stage_sequence())?;
205        let (original, _) =
206            read_binding(&predecessor, workflow, row.stage_sequence(), row.binding())?;
207        read_execution_settlement(&predecessor, original.plan(), row.settlement())?;
208    }
209    Ok(())
210}
211fn stage_path(layout: &BackupLayoutGuard, sequence: u64) -> PathBuf {
212    layout.root().join(format!("execution-stage-{sequence}"))
213}
214fn acquire_stage(
215    layout: &BackupLayoutGuard,
216    sequence: u64,
217) -> Result<BackupLayoutGuard, ExecutionWorkflowPersistenceError> {
218    layout.check_root()?;
219    let path = stage_path(layout, sequence);
220    let metadata = fs::symlink_metadata(&path).map_err(PersistenceError::from)?;
221    if !metadata.is_dir() {
222        return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
223    }
224    let stage = BackupLayoutGuard::acquire(&path)?;
225    // Unlike an explicitly operator-selected root, a derived stage symlink is
226    // never adopted. Noncooperating namespace stability remains caller-owned.
227    if stage.root() != path {
228        return Err(ExecutionWorkflowPersistenceError::UnsafeStage);
229    }
230    layout.check_root()?;
231    Ok(stage)
232}
233fn create_private_directory(
234    path: &std::path::Path,
235) -> Result<(), ExecutionWorkflowPersistenceError> {
236    #[cfg(unix)]
237    {
238        use std::os::unix::fs::DirBuilderExt;
239        fs::DirBuilder::new()
240            .mode(0o700)
241            .create(path)
242            .map_err(PersistenceError::from)?;
243        Ok(())
244    }
245    #[cfg(not(unix))]
246    {
247        let _ = path;
248        Err(PersistenceError::Io(std::io::Error::from(std::io::ErrorKind::Unsupported)).into())
249    }
250}
251
252/// Typed immutable workflow/stage admission failure; no recovery evidence is removed.
253#[derive(Debug, Error)]
254pub enum ExecutionWorkflowPersistenceError {
255    /// Workflow/stage identity differs from the retained expected digest.
256    #[error("execution workflow/stage digest mismatch")]
257    DigestMismatch,
258    /// A derived stage directory is a symlink or another unsafe entry.
259    #[error("unsafe execution stage directory")]
260    UnsafeStage,
261    /// Original model allocation or learned binding was rejected.
262    #[error(transparent)]
263    Model(#[from] ExecutionWorkflowError),
264    /// Existing immutable operation-plan admission failed.
265    #[error(transparent)]
266    Plan(#[from] OperationPlanPersistenceError),
267    /// Existing complete chronological Applied settlement admission failed.
268    #[error(transparent)]
269    Settlement(#[from] ExecutionSettlementPersistenceError),
270    /// Cooperating layout or journal exclusion failed.
271    #[error(transparent)]
272    Lock(#[from] JournalLockError),
273    /// Existing bounded JSON/durable filesystem owner failed.
274    #[error(transparent)]
275    Persistence(#[from] PersistenceError),
276}
277
278#[cfg(all(test, unix))]
279mod tests;