use super::{IcSnapshotUploadExecutionError, upload_snapshot};
use crate::{
model::{
attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
execution_workflow::ExecutionStagePredecessorRecord,
ic_mutation::IcMutationAcknowledgement,
ic_snapshot_upload::{
IcSnapshotUploadAttempt, IcSnapshotUploadAttemptError, IcSnapshotUploadKind,
IcSnapshotUploadRequest,
},
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, ExecutionStageCheckpointError,
ExecutionStageGuard, ExecutionWorkflowPersistenceError,
},
policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
ports::ic_snapshot_upload::IcSnapshotUploadProvider,
};
use thiserror::Error;
pub fn allocate_snapshot<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
operation_sequence: u64,
payload: &IcSnapshotUploadRequest<'_>,
provider: &mut impl IcSnapshotUploadProvider,
admit: impl FnOnce(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
qualify: impl FnOnce(
&IcSnapshotUploadAttempt<'_, '_>,
&IcMutationAcknowledgement,
) -> Result<MutationReceiptRequest, E>,
) -> Result<
(IcMutationAcknowledgement, ExecutionStagePredecessorRecord),
IcSnapshotAllocationExecutionError<E>,
> {
let operations = stage.plan().operations();
if !matches!(payload.kind(), IcSnapshotUploadKind::Metadata)
|| operations.len() != 1
|| operations[0].operation_sequence() != operation_sequence
|| operations[0].request() != payload.binding_digest().hash()
{
return Err(IcSnapshotAllocationExecutionError::OriginalMismatch);
}
let acknowledgement = upload_snapshot(stage, operation_sequence, payload, provider, admit)?;
match settle_allocation(
stage,
operation_sequence,
payload,
&acknowledgement,
qualify,
) {
Ok(predecessor) => Ok((acknowledgement, predecessor)),
Err(source) => Err(IcSnapshotAllocationExecutionError::AfterReply {
source,
acknowledgement: Box::new(acknowledgement),
}),
}
}
fn settle_allocation<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
sequence: u64,
payload: &IcSnapshotUploadRequest<'_>,
acknowledgement: &IcMutationAcknowledgement,
qualify: impl FnOnce(
&IcSnapshotUploadAttempt<'_, '_>,
&IcMutationAcknowledgement,
) -> Result<MutationReceiptRequest, E>,
) -> Result<ExecutionStagePredecessorRecord, IcSnapshotAllocationSettlementError<E>> {
let plan = stage.plan();
let authority = plan
.attempt_authority(sequence)
.map_err(IcSnapshotUploadAttemptError::from)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
let request = IcSnapshotUploadAttempt::new(plan, sequence, journal.record()?, payload)?;
let admitted = validate_acknowledgement(&request, journal.record()?, acknowledgement)?;
let receipt = qualify(&request, acknowledgement)
.map_err(IcSnapshotAllocationSettlementError::Qualification)?;
if receipt.outcome != MutationOutcomeRecord::Applied
|| receipt.attempt != request.mutation_attempt()
|| receipt.request != payload.binding_digest().hash()
{
return Err(IcSnapshotAllocationSettlementError::ReceiptRequired);
}
stage.layout()?;
journal.record_mutation(receipt)?;
drop(journal);
Ok(stage.checkpoint(admitted.reply().digest())?)
}
#[derive(Debug, Error)]
pub enum IcSnapshotAllocationExecutionError<E: std::error::Error + 'static> {
#[error("snapshot allocation requires an exact singleton original metadata stage")]
OriginalMismatch,
#[error(transparent)]
Upload(#[from] IcSnapshotUploadExecutionError<E>),
#[error("snapshot allocation reply settlement failed: {source}")]
AfterReply {
source: IcSnapshotAllocationSettlementError<E>,
acknowledgement: Box<IcMutationAcknowledgement>,
},
}
#[derive(Debug, Error)]
pub enum IcSnapshotAllocationSettlementError<E: std::error::Error + 'static> {
#[error("snapshot allocation qualification failed: {0}")]
Qualification(#[source] E),
#[error("snapshot allocation requires an explicit original Applied receipt")]
ReceiptRequired,
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Request(#[from] IcSnapshotUploadAttemptError),
#[error(transparent)]
Association(#[from] IcSnapshotUploadAssociationError),
#[error(transparent)]
Checkpoint(#[from] ExecutionStageCheckpointError),
}
#[cfg(all(test, unix))]
mod tests;