ic_backup/model/ic_snapshot_download/
mod.rs1use crate::model::{
4 attempt_journal::{AttemptBudgetRecord, AttemptJournalRecordError},
5 effect_graph::{EffectGraphError, EffectGraphRecord, EffectNodeRecord, EffectNodeRequest},
6 execution_workflow::{
7 ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
8 ExecutionWorkflowRecord,
9 },
10 ic_snapshot_data::{
11 ExtentPlanningError, IcSnapshotDataError, IcSnapshotDataRequest, planned_extents,
12 },
13 ic_snapshot_metadata::IcSnapshotMetadataReply,
14 operation_plan::{
15 OperationPlanError, OperationPlanRecord, OperationPlanRequest, PlanBudgetRecord,
16 PlannedOperationRecord, PlannedOperationRequest,
17 },
18};
19use thiserror::Error;
20
21pub struct IcSnapshotDownloadPlan<'workflow, 'metadata> {
32 workflow: &'workflow ExecutionWorkflowRecord,
33 sequence: u64,
34 metadata: &'metadata IcSnapshotMetadataReply<'metadata>,
35 requests: Vec<IcSnapshotDataRequest<'metadata>>,
36 plan: Option<OperationPlanRecord>,
37}
38impl<'workflow, 'metadata> IcSnapshotDownloadPlan<'workflow, 'metadata> {
39 pub fn new(
48 workflow: &'workflow ExecutionWorkflowRecord,
49 sequence: u64,
50 metadata: &'metadata IcSnapshotMetadataReply<'metadata>,
51 chunk_bytes: u64,
52 ) -> Result<Self, IcSnapshotDownloadPlanningError> {
53 let stage = workflow.stage(sequence)?;
54 if stage.target() != metadata.request().target() {
55 return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
56 }
57 let kinds = planned_extents(metadata, chunk_bytes, stage.budget().mutations())
58 .map_err(IcSnapshotDownloadPlanningError::from)?;
59 let capacity = kinds.len();
60 let requests = kinds
61 .into_iter()
62 .map(|kind| IcSnapshotDataRequest::new(metadata, kind))
63 .collect::<Result<Vec<_>, _>>()?;
64 let plan = if requests.is_empty() {
65 None
66 } else {
67 let mut nodes = Vec::with_capacity(capacity);
68 let mut operations = Vec::with_capacity(capacity);
69 for (index, request) in requests.iter().enumerate() {
70 let sequence = u64::try_from(index)
71 .map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
72 nodes.push(EffectNodeRecord::new(EffectNodeRequest {
73 operation_sequence: sequence,
74 depends_on: sequence.checked_sub(1).into_iter().collect(),
75 })?);
76 operations.push(PlannedOperationRecord::new(PlannedOperationRequest {
77 operation_sequence: sequence,
78 target: stage.target().into(),
79 request: request.digest().hash().into(),
80 budget: AttemptBudgetRecord::new(1, 0)?,
81 })?);
82 }
83 let original = workflow.allocation();
84 Some(OperationPlanRecord::new(OperationPlanRequest {
85 context: original.context().clone(),
86 inventory: original.inventory().clone(),
87 selected_targets: vec![stage.target().into()],
88 graph: EffectGraphRecord::new(nodes)?,
89 operations,
90 budget: PlanBudgetRecord::new(
91 stage.budget().mutations(),
92 stage.budget().observations(),
93 )?,
94 })?)
95 };
96 Ok(Self {
97 workflow,
98 sequence,
99 metadata,
100 requests,
101 plan,
102 })
103 }
104 #[must_use]
106 pub fn requests(&self) -> &[IcSnapshotDataRequest<'metadata>] {
107 &self.requests
108 }
109 #[must_use]
112 pub const fn plan(&self) -> Option<&OperationPlanRecord> {
113 self.plan.as_ref()
114 }
115
116 pub(crate) fn validate_binding(
119 &self,
120 binding: &ExecutionStageBindingRecord,
121 ) -> Result<(), IcSnapshotDownloadPlanningError> {
122 let plan = self
123 .plan
124 .as_ref()
125 .ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
126 binding.validate(self.workflow, plan)?;
127 if binding.stage_sequence() != self.sequence
128 || !binding
129 .predecessors()
130 .iter()
131 .any(|row| row.learned_evidence() == &self.metadata.digest())
132 {
133 return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
134 }
135 Ok(())
136 }
137 pub fn bind(
147 &self,
148 metadata_binding: &ExecutionStageBindingRecord,
149 metadata_plan: &OperationPlanRecord,
150 predecessors: Vec<ExecutionStagePredecessorRecord>,
151 ) -> Result<ExecutionStageBindingRecord, IcSnapshotDownloadPlanningError> {
152 let plan = self
153 .plan
154 .as_ref()
155 .ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
156 metadata_binding.validate(self.workflow, metadata_plan)?;
157 if metadata_plan.operations().len() != 1
158 || metadata_plan.operations()[0].request() != self.metadata.request().digest().hash()
159 || !predecessors.iter().any(|row| {
160 row.stage_sequence() == metadata_binding.stage_sequence()
161 && row.binding() == &metadata_binding.digest()
162 && row.learned_evidence() == &self.metadata.digest()
163 })
164 {
165 return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
166 }
167 Ok(ExecutionStageBindingRecord::new(
168 self.workflow,
169 self.sequence,
170 plan,
171 predecessors,
172 )?)
173 }
174}
175
176#[derive(Debug, Error)]
178pub enum IcSnapshotDownloadPlanningError {
179 #[error("invalid snapshot download chunk size")]
181 InvalidChunkSize,
182 #[error("snapshot download request count overflow")]
184 CountOverflow,
185 #[error("snapshot download needs {required} updates; original allocation is {original}")]
187 InsufficientAllowance {
188 required: u64,
190 original: u32,
192 },
193 #[error("snapshot download metadata differs from original stage")]
195 MetadataMismatch,
196 #[error("snapshot download has no data reads to bind")]
198 NoDataReads,
199 #[error(transparent)]
201 Plan(#[from] OperationPlanError),
202 #[error(transparent)]
204 Graph(#[from] EffectGraphError),
205 #[error(transparent)]
207 Workflow(#[from] ExecutionWorkflowError),
208 #[error(transparent)]
210 Data(#[from] IcSnapshotDataError),
211 #[error(transparent)]
213 Budget(#[from] AttemptJournalRecordError),
214}
215
216#[cfg(test)]
217pub(crate) mod tests;
218
219impl From<ExtentPlanningError> for IcSnapshotDownloadPlanningError {
220 fn from(error: ExtentPlanningError) -> Self {
221 match error {
222 ExtentPlanningError::InvalidChunkSize => Self::InvalidChunkSize,
223 ExtentPlanningError::CountOverflow => Self::CountOverflow,
224 ExtentPlanningError::InsufficientAllowance { required, original } => {
225 Self::InsufficientAllowance { required, original }
226 }
227 }
228 }
229}