Skip to main content

ic_backup/model/execution_workflow/
mod.rs

1//! Original stage allocations and exact learned-plan bindings; no second spending ledger.
2
3use crate::model::{
4    artifacts::ArtifactChecksumRecord,
5    effect_graph::MAX_EFFECT_DEPENDENCIES,
6    operation_plan::{OperationPlanError, OperationPlanRecord, PlannedOperationRecord},
7};
8use serde::{Deserialize, Deserializer, Serialize, de};
9use std::fmt;
10use thiserror::Error;
11
12/// Maximum raw and canonical workflow declaration bytes.
13pub const MAX_EXECUTION_WORKFLOW_BYTES: u64 = 1024 * 1024;
14/// Maximum raw and canonical learned-stage binding bytes.
15pub const MAX_EXECUTION_STAGE_BYTES: u64 = 512 * 1024;
16
17/// Strict v1 original stage catalog, admitted by the existing plan structural owner.
18///
19/// Each allocation operation denotes one stage: its request hash commits the
20/// integration-owned input/purpose contract, its target is the exact stage target,
21/// and its budget is the original ceiling for that stage's later exact child plan.
22/// The allocation is never an executable operation plan or an attempt-journal
23/// authority. Only bound child plans own attempt journals. Unassigned aggregate
24/// headroom cannot be transferred to a stage. This distinct record does not change
25/// the meaning or schema of ordinary v1 operation plans.
26#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
27#[serde(try_from = "WorkflowFields")]
28pub struct ExecutionWorkflowRecord {
29    version: u16,
30    allocation: OperationPlanRecord,
31}
32#[derive(Deserialize)]
33#[serde(deny_unknown_fields)]
34struct WorkflowFields {
35    version: u16,
36    allocation: OperationPlanRecord,
37}
38impl TryFrom<WorkflowFields> for ExecutionWorkflowRecord {
39    type Error = ExecutionWorkflowError;
40    fn try_from(fields: WorkflowFields) -> Result<Self, Self::Error> {
41        if fields.version != 1 {
42            return Err(ExecutionWorkflowError::UnsupportedVersion(fields.version));
43        }
44        Ok(Self::new(fields.allocation))
45    }
46}
47impl ExecutionWorkflowRecord {
48    pub(crate) const fn allocation(&self) -> &OperationPlanRecord {
49        &self.allocation
50    }
51    /// Retain a canonical allocation without creating journals or dispatch authority.
52    #[must_use]
53    pub const fn new(allocation: OperationPlanRecord) -> Self {
54        Self {
55            version: 1,
56            allocation,
57        }
58    }
59    /// Read one original stage's purpose commitment, target and attempt ceilings.
60    /// # Errors
61    /// Rejects a stage not present in the original catalog.
62    pub fn stage(&self, sequence: u64) -> Result<&PlannedOperationRecord, OperationPlanError> {
63        self.allocation.operation(sequence)
64    }
65    /// Read the canonical original stage catalog; rows are declarations, not journals.
66    #[must_use]
67    pub fn stages(&self) -> &[PlannedOperationRecord] {
68        self.allocation.operations()
69    }
70    /// Bind the full original allocation under a distinct NUL-terminated v1 domain.
71    #[must_use]
72    pub fn digest(&self) -> ArtifactChecksumRecord {
73        let mut bytes = b"ic-backup/execution-workflow/v1\0".to_vec();
74        bytes.extend_from_slice(self.allocation.digest().hash().as_bytes());
75        ArtifactChecksumRecord::from_bytes(&bytes)
76    }
77}
78
79/// Exact preceding stage binding, chronological settlement and learned-input commitment.
80///
81/// The evidence digest is integration-owned. It proves neither authenticity nor
82/// that learned snapshot IDs/dimensions were derived from those settled attempts.
83#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
84#[serde(deny_unknown_fields)]
85pub struct ExecutionStagePredecessorRecord {
86    stage_sequence: u64,
87    binding: ArtifactChecksumRecord,
88    settlement: ArtifactChecksumRecord,
89    learned_evidence: ArtifactChecksumRecord,
90}
91impl ExecutionStagePredecessorRecord {
92    /// Declare exact predecessor identities without inferring an effect outcome.
93    #[must_use]
94    pub const fn new(
95        stage_sequence: u64,
96        binding: ArtifactChecksumRecord,
97        settlement: ArtifactChecksumRecord,
98        learned_evidence: ArtifactChecksumRecord,
99    ) -> Self {
100        Self {
101            stage_sequence,
102            binding,
103            settlement,
104            learned_evidence,
105        }
106    }
107    /// Read the original predecessor stage identity.
108    #[must_use]
109    pub const fn stage_sequence(&self) -> u64 {
110        self.stage_sequence
111    }
112    /// Read the exact immutable predecessor binding identity.
113    #[must_use]
114    pub const fn binding(&self) -> &ArtifactChecksumRecord {
115        &self.binding
116    }
117    /// Read the exact retained chronological settlement identity.
118    #[must_use]
119    pub const fn settlement(&self) -> &ArtifactChecksumRecord {
120        &self.settlement
121    }
122    /// Read the opaque learned-input commitment, without qualifying its meaning.
123    #[must_use]
124    pub const fn learned_evidence(&self) -> &ArtifactChecksumRecord {
125        &self.learned_evidence
126    }
127}
128
129/// Strict v1 immutable association of one allocated stage with one exact ordinary plan.
130///
131/// The canonical child plan remains the sole request/authority owner. This record
132/// copies no counters, receipts, completion flags or payloads. Local persistence
133/// additionally requires every declared predecessor's complete Applied journal
134/// settlement; that still grants no authentic capture, transfer or dispatch permit.
135#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
136#[serde(try_from = "StageFields")]
137pub struct ExecutionStageBindingRecord {
138    version: u16,
139    workflow: ArtifactChecksumRecord,
140    stage_sequence: u64,
141    plan: ArtifactChecksumRecord,
142    predecessors: Vec<ExecutionStagePredecessorRecord>,
143}
144#[derive(Deserialize)]
145#[serde(deny_unknown_fields)]
146struct StageFields {
147    version: u16,
148    workflow: ArtifactChecksumRecord,
149    stage_sequence: u64,
150    plan: ArtifactChecksumRecord,
151    #[serde(deserialize_with = "bounded_predecessors")]
152    predecessors: Vec<ExecutionStagePredecessorRecord>,
153}
154impl TryFrom<StageFields> for ExecutionStageBindingRecord {
155    type Error = ExecutionWorkflowError;
156    fn try_from(mut fields: StageFields) -> Result<Self, Self::Error> {
157        if fields.version != 1 {
158            return Err(ExecutionWorkflowError::UnsupportedVersion(fields.version));
159        }
160        canonical_predecessors(&mut fields.predecessors)?;
161        Ok(Self {
162            version: 1,
163            workflow: fields.workflow,
164            stage_sequence: fields.stage_sequence,
165            plan: fields.plan,
166            predecessors: fields.predecessors,
167        })
168    }
169}
170impl ExecutionStageBindingRecord {
171    /// Bind exact learned requests within one original target and allocation.
172    ///
173    /// Child aggregate ceilings must equal the stage's original ceilings; existing
174    /// plan admission bounds assigned operations below those ceilings. Zero-budget
175    /// child operations reject, so the original 65,536 workflow attempt ceiling
176    /// also bounds total bound child operations. No allowance is borrowed/reset.
177    /// # Errors
178    /// Rejects changed context/inventory/target/limits, unknown stages and incomplete dependencies.
179    pub fn new(
180        workflow: &ExecutionWorkflowRecord,
181        stage_sequence: u64,
182        plan: &OperationPlanRecord,
183        mut predecessors: Vec<ExecutionStagePredecessorRecord>,
184    ) -> Result<Self, ExecutionWorkflowError> {
185        let stage = workflow
186            .stage(stage_sequence)
187            .map_err(|_| ExecutionWorkflowError::UnknownStage(stage_sequence))?;
188        if plan.context() != workflow.allocation.context()
189            || plan.inventory() != workflow.allocation.inventory()
190            || plan.selected_targets() != [stage.target()]
191        {
192            return Err(ExecutionWorkflowError::ChildIdentityMismatch);
193        }
194        if plan.budget().mutations() != stage.budget().mutations()
195            || plan.budget().observations() != stage.budget().observations()
196        {
197            return Err(ExecutionWorkflowError::ChildBudgetMismatch);
198        }
199        if plan.operations().iter().any(|operation| {
200            operation.budget().mutations() == 0 && operation.budget().observations() == 0
201        }) {
202            return Err(ExecutionWorkflowError::UnallocatedOperation);
203        }
204        canonical_predecessors(&mut predecessors)?;
205        let dependencies = workflow
206            .allocation
207            .graph()
208            .node(stage_sequence)
209            .map_err(|_| ExecutionWorkflowError::PredecessorMismatch)?
210            .depends_on();
211        if !predecessors
212            .iter()
213            .map(ExecutionStagePredecessorRecord::stage_sequence)
214            .eq(dependencies.iter().copied())
215        {
216            return Err(ExecutionWorkflowError::PredecessorMismatch);
217        }
218        Ok(Self {
219            version: 1,
220            workflow: workflow.digest(),
221            stage_sequence,
222            plan: plan.digest(),
223            predecessors,
224        })
225    }
226    /// Re-admit decoded bindings against full original workflow and child plan.
227    /// # Errors
228    /// Rejects original identity, plan or allocation/dependency changes.
229    pub fn validate(
230        &self,
231        workflow: &ExecutionWorkflowRecord,
232        plan: &OperationPlanRecord,
233    ) -> Result<(), ExecutionWorkflowError> {
234        if self.workflow != workflow.digest() || self.plan != plan.digest() {
235            return Err(ExecutionWorkflowError::BindingMismatch);
236        }
237        Self::new(
238            workflow,
239            self.stage_sequence,
240            plan,
241            self.predecessors.clone(),
242        )?;
243        Ok(())
244    }
245    /// Read full original workflow identity.
246    #[must_use]
247    pub const fn workflow(&self) -> &ArtifactChecksumRecord {
248        &self.workflow
249    }
250    /// Read the fixed original stage identity.
251    #[must_use]
252    pub const fn stage_sequence(&self) -> u64 {
253        self.stage_sequence
254    }
255    /// Read the exact canonical child plan identity.
256    #[must_use]
257    pub const fn plan(&self) -> &ArtifactChecksumRecord {
258        &self.plan
259    }
260    /// Read canonical exact predecessors and learned-input commitments.
261    #[must_use]
262    pub fn predecessors(&self) -> &[ExecutionStagePredecessorRecord] {
263        &self.predecessors
264    }
265    /// Hash the v1 NUL domain, workflow/plan ASCII hashes, big-endian u64 stage/count,
266    /// then canonical predecessor u64 identities and binding/settlement/evidence hashes.
267    #[must_use]
268    pub fn digest(&self) -> ArtifactChecksumRecord {
269        let mut bytes = b"ic-backup/execution-stage/v1\0".to_vec();
270        bytes.extend_from_slice(self.workflow.hash().as_bytes());
271        bytes.extend_from_slice(&self.stage_sequence.to_be_bytes());
272        bytes.extend_from_slice(self.plan.hash().as_bytes());
273        bytes.extend_from_slice(&(self.predecessors.len() as u64).to_be_bytes());
274        for row in &self.predecessors {
275            bytes.extend_from_slice(&row.stage_sequence.to_be_bytes());
276            for digest in [&row.binding, &row.settlement, &row.learned_evidence] {
277                bytes.extend_from_slice(digest.hash().as_bytes());
278            }
279        }
280        ArtifactChecksumRecord::from_bytes(&bytes)
281    }
282}
283fn canonical_predecessors(
284    rows: &mut [ExecutionStagePredecessorRecord],
285) -> Result<(), ExecutionWorkflowError> {
286    if rows.len() > MAX_EFFECT_DEPENDENCIES {
287        return Err(ExecutionWorkflowError::PredecessorMismatch);
288    }
289    rows.sort_by_key(ExecutionStagePredecessorRecord::stage_sequence);
290    if rows
291        .windows(2)
292        .any(|pair| pair[0].stage_sequence == pair[1].stage_sequence)
293    {
294        return Err(ExecutionWorkflowError::PredecessorMismatch);
295    }
296    Ok(())
297}
298fn bounded_predecessors<'de, D: Deserializer<'de>>(
299    deserializer: D,
300) -> Result<Vec<ExecutionStagePredecessorRecord>, D::Error> {
301    struct Rows;
302    impl<'de> de::Visitor<'de> for Rows {
303        type Value = Vec<ExecutionStagePredecessorRecord>;
304        fn expecting(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
305            f.write_str("at most 1,024 exact stage predecessors")
306        }
307        fn visit_seq<A: de::SeqAccess<'de>>(
308            self,
309            mut sequence: A,
310        ) -> Result<Self::Value, A::Error> {
311            let mut rows = Vec::new();
312            while rows.len() < MAX_EFFECT_DEPENDENCIES {
313                match sequence.next_element()? {
314                    Some(row) => rows.push(row),
315                    None => return Ok(rows),
316                }
317            }
318            if sequence.next_element::<de::IgnoredAny>()?.is_some() {
319                return Err(de::Error::custom(
320                    ExecutionWorkflowError::PredecessorMismatch,
321                ));
322            }
323            Ok(rows)
324        }
325    }
326    deserializer.deserialize_seq(Rows)
327}
328/// Typed rejection without changing original records or spending.
329#[derive(Debug, Error, Eq, PartialEq)]
330pub enum ExecutionWorkflowError {
331    /// Only the distinct maintained v1 generation is admitted.
332    #[error("unsupported execution workflow/stage version {0}")]
333    UnsupportedVersion(u16),
334    /// Full context, inventory or exact single allocated target changed.
335    #[error("child plan differs from original stage identity")]
336    ChildIdentityMismatch,
337    /// Child aggregate ceilings differ from original assigned stage limits.
338    #[error("child plan differs from original stage budget")]
339    ChildBudgetMismatch,
340    /// A child operation has no assigned allowance.
341    #[error("child plan contains an unallocated operation")]
342    UnallocatedOperation,
343    /// Predecessors are excessive, duplicated, missing or outside the original graph.
344    #[error("stage predecessors differ from original dependencies")]
345    PredecessorMismatch,
346    /// Retained full workflow or child plan digest differs.
347    #[error("execution stage binding mismatch")]
348    BindingMismatch,
349    /// The original catalog has no stage with this identity.
350    #[error("unknown execution stage {0}")]
351    UnknownStage(u64),
352}
353
354#[cfg(test)]
355pub(crate) mod tests;