ic_backup/model/ic_snapshot_upload/plan/
mod.rs1use super::{
4 IcSnapshotUploadError, IcSnapshotUploadReply, IcSnapshotUploadReplyKind,
5 IcSnapshotUploadRequest,
6};
7use crate::model::{
8 artifacts::ArtifactChecksumRecord,
9 attempt_journal::{AttemptBudgetRecord, AttemptJournalRecordError},
10 effect_graph::{EffectGraphError, EffectGraphRecord, EffectNodeRecord, EffectNodeRequest},
11 execution_workflow::{
12 ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
13 ExecutionWorkflowRecord,
14 },
15 ic_snapshot_data::{ExtentPlanningError, planned_extents},
16 operation_plan::{
17 OperationPlanError, OperationPlanRecord, OperationPlanRequest, PlanBudgetRecord,
18 PlannedOperationRecord, PlannedOperationRequest,
19 },
20};
21use ic_management_canister_types::SnapshotDataKind;
22use thiserror::Error;
23
24pub struct IcSnapshotDataUploadPlan<'workflow, 'request, 'source> {
32 workflow: &'workflow ExecutionWorkflowRecord,
33 sequence: u64,
34 metadata: &'request IcSnapshotUploadRequest<'source>,
35 destination: Vec<u8>,
36 allocation_evidence: ArtifactChecksumRecord,
37 kinds: Vec<SnapshotDataKind>,
38 plan: Option<OperationPlanRecord>,
39}
40impl<'workflow, 'request, 'source> IcSnapshotDataUploadPlan<'workflow, 'request, 'source> {
41 pub(crate) fn derive_kinds(
42 workflow: &ExecutionWorkflowRecord,
43 sequence: u64,
44 allocation: &IcSnapshotUploadReply<'_, '_>,
45 chunk_bytes: u64,
46 ) -> Result<Vec<SnapshotDataKind>, IcSnapshotDataUploadPlanningError> {
47 let metadata = allocation.request();
48 let stage = workflow.stage(sequence)?;
49 let IcSnapshotUploadReplyKind::Metadata { snapshot_id } = allocation.kind() else {
50 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
51 };
52 metadata.validate_data_destination(snapshot_id)?;
53 if stage.target() != metadata.target()
54 || !super::attempt::source_context_matches(workflow.allocation(), metadata)
55 {
56 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
57 }
58 Ok(planned_extents(
59 metadata.source(),
60 chunk_bytes,
61 stage.budget().mutations(),
62 )?)
63 }
64
65 pub(crate) fn from_bindings(
66 workflow: &'workflow ExecutionWorkflowRecord,
67 sequence: u64,
68 allocation: &IcSnapshotUploadReply<'request, 'source>,
69 chunk_bytes: u64,
70 bindings: &[ArtifactChecksumRecord],
71 ) -> Result<Self, IcSnapshotDataUploadPlanningError> {
72 let kinds = Self::derive_kinds(workflow, sequence, allocation, chunk_bytes)?;
73 if bindings.len() != kinds.len() {
74 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
75 }
76 let IcSnapshotUploadReplyKind::Metadata { snapshot_id } = allocation.kind() else {
77 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
78 };
79 let stage = workflow.stage(sequence)?;
80 let plan = if kinds.is_empty() {
81 None
82 } else {
83 let mut nodes = Vec::with_capacity(kinds.len());
84 let mut operations = Vec::with_capacity(kinds.len());
85 for (index, request) in bindings.iter().enumerate() {
86 let operation_sequence = u64::try_from(index)
87 .map_err(|_| IcSnapshotDataUploadPlanningError::CountOverflow)?;
88 nodes.push(EffectNodeRecord::new(EffectNodeRequest {
89 operation_sequence,
90 depends_on: operation_sequence.checked_sub(1).into_iter().collect(),
91 })?);
92 operations.push(PlannedOperationRecord::new(PlannedOperationRequest {
93 operation_sequence,
94 target: stage.target().into(),
95 request: request.hash().into(),
96 budget: AttemptBudgetRecord::new(1, 0)?,
97 })?);
98 }
99 let original = workflow.allocation();
100 Some(OperationPlanRecord::new(OperationPlanRequest {
101 context: original.context().clone(),
102 inventory: original.inventory().clone(),
103 selected_targets: vec![stage.target().into()],
104 graph: EffectGraphRecord::new(nodes)?,
105 operations,
106 budget: PlanBudgetRecord::new(
107 stage.budget().mutations(),
108 stage.budget().observations(),
109 )?,
110 })?)
111 };
112 Ok(Self {
113 workflow,
114 sequence,
115 metadata: allocation.request(),
116 destination: snapshot_id.clone(),
117 allocation_evidence: allocation.digest(),
118 kinds,
119 plan,
120 })
121 }
122
123 #[must_use]
125 pub fn kinds(&self) -> &[SnapshotDataKind] {
126 &self.kinds
127 }
128 #[must_use]
130 pub const fn metadata(&self) -> &'request IcSnapshotUploadRequest<'source> {
131 self.metadata
132 }
133 #[must_use]
135 pub fn destination(&self) -> &[u8] {
136 &self.destination
137 }
138 #[must_use]
140 pub const fn plan(&self) -> Option<&OperationPlanRecord> {
141 self.plan.as_ref()
142 }
143
144 pub fn bind(
152 &self,
153 allocation_binding: &ExecutionStageBindingRecord,
154 allocation_plan: &OperationPlanRecord,
155 predecessors: Vec<ExecutionStagePredecessorRecord>,
156 ) -> Result<ExecutionStageBindingRecord, IcSnapshotDataUploadPlanningError> {
157 let plan = self
158 .plan
159 .as_ref()
160 .ok_or(IcSnapshotDataUploadPlanningError::NoDataWrites)?;
161 allocation_binding.validate(self.workflow, allocation_plan)?;
162 if allocation_plan.operations().len() != 1
163 || allocation_plan.operations()[0].request() != self.metadata.binding_digest().hash()
164 || !predecessors.iter().any(|row| {
165 row.stage_sequence() == allocation_binding.stage_sequence()
166 && row.binding() == &allocation_binding.digest()
167 && row.learned_evidence() == &self.allocation_evidence
168 })
169 {
170 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
171 }
172 Ok(ExecutionStageBindingRecord::new(
173 self.workflow,
174 self.sequence,
175 plan,
176 predecessors,
177 )?)
178 }
179 pub(crate) fn validate_binding(
180 &self,
181 binding: &ExecutionStageBindingRecord,
182 ) -> Result<(), IcSnapshotDataUploadPlanningError> {
183 let plan = self
184 .plan
185 .as_ref()
186 .ok_or(IcSnapshotDataUploadPlanningError::NoDataWrites)?;
187 binding.validate(self.workflow, plan)?;
188 if binding.stage_sequence() != self.sequence
189 || !binding
190 .predecessors()
191 .iter()
192 .any(|row| row.learned_evidence() == &self.allocation_evidence)
193 {
194 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch);
195 }
196 Ok(())
197 }
198}
199
200#[derive(Debug, Error)]
202pub enum IcSnapshotDataUploadPlanningError {
203 #[error("invalid snapshot data upload chunk size")]
205 InvalidChunkSize,
206 #[error("snapshot data upload count overflow")]
208 CountOverflow,
209 #[error("snapshot data upload needs {required} updates; original allocation is {original}")]
211 InsufficientAllowance {
212 required: u64,
214 original: u32,
216 },
217 #[error("snapshot data upload originals differ")]
219 OriginalMismatch,
220 #[error("snapshot data upload has no data writes to bind")]
222 NoDataWrites,
223 #[error(transparent)]
225 Plan(#[from] OperationPlanError),
226 #[error(transparent)]
228 Graph(#[from] EffectGraphError),
229 #[error(transparent)]
231 Workflow(#[from] ExecutionWorkflowError),
232 #[error(transparent)]
234 Upload(#[from] IcSnapshotUploadError),
235 #[error(transparent)]
237 Budget(#[from] AttemptJournalRecordError),
238}
239impl From<ExtentPlanningError> for IcSnapshotDataUploadPlanningError {
240 fn from(error: ExtentPlanningError) -> Self {
241 match error {
242 ExtentPlanningError::InvalidChunkSize => Self::InvalidChunkSize,
243 ExtentPlanningError::CountOverflow => Self::CountOverflow,
244 ExtentPlanningError::InsufficientAllowance { required, original } => {
245 Self::InsufficientAllowance { required, original }
246 }
247 }
248 }
249}
250
251#[cfg(test)]
252mod tests;