use crate::model::{
attempt_journal::{AttemptBudgetRecord, AttemptJournalRecordError},
effect_graph::{EffectGraphError, EffectGraphRecord, EffectNodeRecord, EffectNodeRequest},
execution_workflow::{
ExecutionStageBindingRecord, ExecutionStagePredecessorRecord, ExecutionWorkflowError,
ExecutionWorkflowRecord,
},
ic_snapshot_data::{
ExtentPlanningError, IcSnapshotDataError, IcSnapshotDataRequest, planned_extents,
},
ic_snapshot_metadata::IcSnapshotMetadataReply,
operation_plan::{
OperationPlanError, OperationPlanRecord, OperationPlanRequest, PlanBudgetRecord,
PlannedOperationRecord, PlannedOperationRequest,
},
};
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);
}
let kinds = planned_extents(metadata, chunk_bytes, stage.budget().mutations())
.map_err(IcSnapshotDownloadPlanningError::from)?;
let capacity = kinds.len();
let requests = kinds
.into_iter()
.map(|kind| IcSnapshotDataRequest::new(metadata, kind))
.collect::<Result<Vec<_>, _>>()?;
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(crate) fn validate_binding(
&self,
binding: &ExecutionStageBindingRecord,
) -> Result<(), IcSnapshotDownloadPlanningError> {
let plan = self
.plan
.as_ref()
.ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
binding.validate(self.workflow, plan)?;
if binding.stage_sequence() != self.sequence
|| !binding
.predecessors()
.iter()
.any(|row| row.learned_evidence() == &self.metadata.digest())
{
return Err(IcSnapshotDownloadPlanningError::MetadataMismatch);
}
Ok(())
}
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)]
pub(crate) mod tests;
impl From<ExtentPlanningError> for IcSnapshotDownloadPlanningError {
fn from(error: ExtentPlanningError) -> Self {
match error {
ExtentPlanningError::InvalidChunkSize => Self::InvalidChunkSize,
ExtentPlanningError::CountOverflow => Self::CountOverflow,
ExtentPlanningError::InsufficientAllowance { required, original } => {
Self::InsufficientAllowance { required, original }
}
}
}
}