ic_backup/model/execution_workflow/
mod.rs1use 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
12pub const MAX_EXECUTION_WORKFLOW_BYTES: u64 = 1024 * 1024;
14pub const MAX_EXECUTION_STAGE_BYTES: u64 = 512 * 1024;
16
17#[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 #[must_use]
53 pub const fn new(allocation: OperationPlanRecord) -> Self {
54 Self {
55 version: 1,
56 allocation,
57 }
58 }
59 pub fn stage(&self, sequence: u64) -> Result<&PlannedOperationRecord, OperationPlanError> {
63 self.allocation.operation(sequence)
64 }
65 #[must_use]
67 pub fn stages(&self) -> &[PlannedOperationRecord] {
68 self.allocation.operations()
69 }
70 #[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#[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 #[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 #[must_use]
109 pub const fn stage_sequence(&self) -> u64 {
110 self.stage_sequence
111 }
112 #[must_use]
114 pub const fn binding(&self) -> &ArtifactChecksumRecord {
115 &self.binding
116 }
117 #[must_use]
119 pub const fn settlement(&self) -> &ArtifactChecksumRecord {
120 &self.settlement
121 }
122 #[must_use]
124 pub const fn learned_evidence(&self) -> &ArtifactChecksumRecord {
125 &self.learned_evidence
126 }
127}
128
129#[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 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 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 #[must_use]
247 pub const fn workflow(&self) -> &ArtifactChecksumRecord {
248 &self.workflow
249 }
250 #[must_use]
252 pub const fn stage_sequence(&self) -> u64 {
253 self.stage_sequence
254 }
255 #[must_use]
257 pub const fn plan(&self) -> &ArtifactChecksumRecord {
258 &self.plan
259 }
260 #[must_use]
262 pub fn predecessors(&self) -> &[ExecutionStagePredecessorRecord] {
263 &self.predecessors
264 }
265 #[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#[derive(Debug, Error, Eq, PartialEq)]
330pub enum ExecutionWorkflowError {
331 #[error("unsupported execution workflow/stage version {0}")]
333 UnsupportedVersion(u16),
334 #[error("child plan differs from original stage identity")]
336 ChildIdentityMismatch,
337 #[error("child plan differs from original stage budget")]
339 ChildBudgetMismatch,
340 #[error("child plan contains an unallocated operation")]
342 UnallocatedOperation,
343 #[error("stage predecessors differ from original dependencies")]
345 PredecessorMismatch,
346 #[error("execution stage binding mismatch")]
348 BindingMismatch,
349 #[error("unknown execution stage {0}")]
351 UnknownStage(u64),
352}
353
354#[cfg(test)]
355pub(crate) mod tests;