use crate::model::{
attempt_journal::{AttemptBudgetRecord, AttemptJournalRecordError},
effect_graph::{EffectGraphError, EffectGraphRecord, EffectNodeRecord, EffectNodeRequest},
execution_workflow::{
ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
ExecutionWorkflowRecord,
},
ic_snapshot_data::{
IcSnapshotDataError, IcSnapshotDataRequest, MAX_IC_SNAPSHOT_DATA_CHUNK_BYTES,
},
ic_snapshot_metadata::IcSnapshotMetadataReply,
operation_plan::{
OperationPlanError, OperationPlanRecord, OperationPlanRequest, PlanBudgetRecord,
PlannedOperationRecord, PlannedOperationRequest,
},
};
use ic_management_canister_types::SnapshotDataKind;
use thiserror::Error;
pub struct IcSnapshotDownloadPlan<'workflow, 'metadata> {
workflow: &'workflow ExecutionWorkflowRecord,
sequence: u64,
metadata: &'metadata IcSnapshotMetadataReply<'metadata>,
requests: Vec<IcSnapshotDataRequest<'metadata>>,
plan: Option<OperationPlanRecord>,
}
impl<'workflow, 'metadata> IcSnapshotDownloadPlan<'workflow, 'metadata> {
pub fn new(
workflow: &'workflow ExecutionWorkflowRecord,
sequence: u64,
metadata: &'metadata IcSnapshotMetadataReply<'metadata>,
chunk_bytes: u64,
) -> Result<Self, IcSnapshotDownloadPlanningError> {
let stage = workflow.stage(sequence)?;
if stage.target() != metadata.request().target() {
return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
}
if chunk_bytes == 0 || chunk_bytes > MAX_IC_SNAPSHOT_DATA_CHUNK_BYTES as u64 {
return Err(IcSnapshotDownloadPlanningError::InvalidChunkSize);
}
let values = metadata.metadata();
let sizes = [
values.wasm_module_size,
values.wasm_memory_size,
values.stable_memory_size,
];
let mut count = u64::try_from(values.wasm_chunk_store.len())
.map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
for size in sizes {
let region = size / chunk_bytes + u64::from(size % chunk_bytes != 0);
count = count
.checked_add(region)
.ok_or(IcSnapshotDownloadPlanningError::CountOverflow)?;
}
if count > u64::from(stage.budget().mutations()) {
return Err(IcSnapshotDownloadPlanningError::InsufficientAllowance {
required: count,
original: stage.budget().mutations(),
});
}
let capacity =
usize::try_from(count).map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
let mut requests = Vec::with_capacity(capacity);
for (region, total) in sizes.into_iter().enumerate() {
let mut offset = 0;
while offset < total {
let size = (total - offset).min(chunk_bytes);
let kind = match region {
0 => SnapshotDataKind::WasmModule { offset, size },
1 => SnapshotDataKind::WasmMemory { offset, size },
_ => SnapshotDataKind::StableMemory { offset, size },
};
requests.push(IcSnapshotDataRequest::new(metadata, kind)?);
offset += size; }
}
for chunk in &values.wasm_chunk_store {
requests.push(IcSnapshotDataRequest::new(
metadata,
SnapshotDataKind::WasmChunk {
hash: chunk.hash.clone(),
},
)?);
}
let plan = if requests.is_empty() {
None
} else {
let mut nodes = Vec::with_capacity(capacity);
let mut operations = Vec::with_capacity(capacity);
for (index, request) in requests.iter().enumerate() {
let sequence = u64::try_from(index)
.map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
nodes.push(EffectNodeRecord::new(EffectNodeRequest {
operation_sequence: sequence,
depends_on: sequence.checked_sub(1).into_iter().collect(),
})?);
operations.push(PlannedOperationRecord::new(PlannedOperationRequest {
operation_sequence: sequence,
target: stage.target().into(),
request: request.digest().hash().into(),
budget: AttemptBudgetRecord::new(1, 0)?,
})?);
}
let original = workflow.allocation();
Some(OperationPlanRecord::new(OperationPlanRequest {
context: original.context().clone(),
inventory: original.inventory().clone(),
selected_targets: vec![stage.target().into()],
graph: EffectGraphRecord::new(nodes)?,
operations,
budget: PlanBudgetRecord::new(
stage.budget().mutations(),
stage.budget().observations(),
)?,
})?)
};
Ok(Self {
workflow,
sequence,
metadata,
requests,
plan,
})
}
#[must_use]
pub fn requests(&self) -> &[IcSnapshotDataRequest<'metadata>] {
&self.requests
}
#[must_use]
pub const fn plan(&self) -> Option<&OperationPlanRecord> {
self.plan.as_ref()
}
pub fn bind(
&self,
metadata_binding: &ExecutionStageBindingRecord,
metadata_plan: &OperationPlanRecord,
predecessors: Vec<ExecutionStagePredecessorRecord>,
) -> Result<ExecutionStageBindingRecord, IcSnapshotDownloadPlanningError> {
let plan = self
.plan
.as_ref()
.ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
metadata_binding.validate(self.workflow, metadata_plan)?;
if metadata_plan.operations().len() != 1
|| metadata_plan.operations()[0].request() != self.metadata.request().digest().hash()
|| !predecessors.iter().any(|row| {
row.stage_sequence() == metadata_binding.stage_sequence()
&& row.binding() == &metadata_binding.digest()
&& row.learned_evidence() == &self.metadata.digest()
})
{
return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
}
Ok(ExecutionStageBindingRecord::new(
self.workflow,
self.sequence,
plan,
predecessors,
)?)
}
}
#[derive(Debug, Error)]
pub enum IcSnapshotDownloadPlanningError {
#[error("invalid snapshot download chunk size")]
InvalidChunkSize,
#[error("snapshot download request count overflow")]
CountOverflow,
#[error("snapshot download needs {required} updates; original allocation is {original}")]
InsufficientAllowance {
required: u64,
original: u32,
},
#[error("snapshot download metadata differs from original stage")]
MetadataMismatch,
#[error("snapshot download has no data reads to bind")]
NoDataReads,
#[error(transparent)]
Plan(#[from] OperationPlanError),
#[error(transparent)]
Graph(#[from] EffectGraphError),
#[error(transparent)]
Workflow(#[from] ExecutionWorkflowError),
#[error(transparent)]
Data(#[from] IcSnapshotDataError),
#[error(transparent)]
Budget(#[from] AttemptJournalRecordError),
}
#[cfg(test)]
mod tests;