use crate::{
model::{
ic_mutation::{IcMutationAcknowledgement, IcMutationRequest, IcMutationRequestError},
ic_request::{IcManagementMethodRecord, IcManagementRequestRecord},
operation_plan::OperationPlanError,
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
ExecutionStageGuard, ExecutionWorkflowPersistenceError,
},
policy::ic_mutation::{IcMutationAssociationError, validate_acknowledgement},
ports::ic_mutation::{IcMutationProvider, IcMutationProviderError},
};
use thiserror::Error;
pub fn capture_snapshot<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
operation_sequence: u64,
payload: &IcManagementRequestRecord,
provider: &mut impl IcMutationProvider,
admit: impl FnOnce(&IcMutationRequest<'_>) -> Result<(), E>,
) -> Result<IcMutationAcknowledgement, IcSnapshotCaptureExecutionError<E>> {
if payload.method() != IcManagementMethodRecord::TakeCanisterSnapshot {
return Err(IcSnapshotCaptureExecutionError::CaptureRequired);
}
let plan = stage.plan();
let authority = plan.attempt_authority(operation_sequence)?;
payload
.validate_mutation_binding(authority.binding())
.map_err(IcMutationRequestError::from)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
journal.reserve_planned_mutation(&plan.digest())?;
let request = IcMutationRequest::new(plan, operation_sequence, journal.record()?, payload)?;
admit(&request).map_err(IcSnapshotCaptureExecutionError::Admission)?;
stage.layout()?;
let acknowledgement = provider.submit_mutation(&request)?;
let association = match journal.record() {
Ok(record) => record,
Err(source) => {
return Err(IcSnapshotCaptureExecutionError::AfterReplyJournal {
source,
acknowledgement: Box::new(acknowledgement),
});
}
};
if let Err(source) = validate_acknowledgement(&request, association, &acknowledgement) {
return Err(IcSnapshotCaptureExecutionError::Association {
source,
acknowledgement: Box::new(acknowledgement),
});
}
if let Err(source) = stage.layout() {
return Err(IcSnapshotCaptureExecutionError::AfterReplyStage {
source,
acknowledgement: Box::new(acknowledgement),
});
}
Ok(acknowledgement)
}
#[derive(Debug, Error)]
pub enum IcSnapshotCaptureExecutionError<E: std::error::Error + 'static> {
#[error("snapshot capture requires take_canister_snapshot")]
CaptureRequired,
#[error(transparent)]
Plan(#[from] OperationPlanError),
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Progress(#[from] ExecutionProgressPersistenceError),
#[error(transparent)]
Request(#[from] IcMutationRequestError),
#[error("fresh snapshot capture admission failed: {0}")]
Admission(#[source] E),
#[error(transparent)]
Provider(#[from] IcMutationProviderError),
#[error("snapshot capture acknowledgement association failed: {source}")]
Association {
source: IcMutationAssociationError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
#[error("snapshot capture journal changed after reply: {source}")]
AfterReplyJournal {
source: AttemptJournalError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
#[error("snapshot capture stage changed after reply: {source}")]
AfterReplyStage {
source: ExecutionWorkflowPersistenceError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
}
#[cfg(all(test, unix))]
mod tests;