use crate::{
model::{
attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
execution_workflow::ExecutionStagePredecessorRecord,
ic_snapshot_metadata::IcSnapshotMetadataRequest,
ic_snapshot_transfer_read::{
IcSnapshotTransferReadError, IcSnapshotTransferReadPayload,
IcSnapshotTransferReadRequest, IcSnapshotTransferReadResponse,
},
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, ExecutionStageCheckpointError,
ExecutionStageGuard, ExecutionWorkflowPersistenceError,
},
policy::ic_snapshot_transfer_read::{
IcSnapshotTransferReadAssociationError, IcSnapshotTransferReadReply, validate_response,
},
ports::ic_snapshot_transfer_read::IcSnapshotTransferReadProvider,
workflow::ic_snapshot_transfer_read::{IcSnapshotTransferReadExecutionError, read_snapshot},
};
use thiserror::Error;
pub async fn read_snapshot_metadata<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
operation_sequence: u64,
payload: &IcSnapshotMetadataRequest,
provider: &mut impl IcSnapshotTransferReadProvider,
admit: impl AsyncFnOnce(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
qualify: impl AsyncFnOnce(
&IcSnapshotTransferReadRequest<'_, '_>,
&IcSnapshotTransferReadResponse,
) -> Result<MutationReceiptRequest, E>,
) -> Result<
(
IcSnapshotTransferReadResponse,
ExecutionStagePredecessorRecord,
),
IcSnapshotMetadataExecutionError<E>,
> {
let operations = stage.plan().operations();
if operations.len() != 1
|| operations[0].operation_sequence() != operation_sequence
|| operations[0].request() != payload.digest().hash()
{
return Err(IcSnapshotMetadataExecutionError::OriginalMismatch);
}
let response = read_snapshot(
stage,
operation_sequence,
IcSnapshotTransferReadPayload::Metadata(payload),
provider,
admit,
)
.await?;
match settle_metadata(stage, operation_sequence, payload, &response, qualify).await {
Ok(predecessor) => Ok((response, predecessor)),
Err(source) => Err(IcSnapshotMetadataExecutionError::AfterReply {
source,
response: Box::new(response),
}),
}
}
async fn settle_metadata<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
sequence: u64,
payload: &IcSnapshotMetadataRequest,
response: &IcSnapshotTransferReadResponse,
qualify: impl AsyncFnOnce(
&IcSnapshotTransferReadRequest<'_, '_>,
&IcSnapshotTransferReadResponse,
) -> Result<MutationReceiptRequest, E>,
) -> Result<ExecutionStagePredecessorRecord, IcSnapshotMetadataSettlementError<E>> {
let plan = stage.plan();
let authority = plan
.attempt_authority(sequence)
.map_err(IcSnapshotTransferReadError::from)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
let request = IcSnapshotTransferReadRequest::new(
plan,
sequence,
journal.record()?,
IcSnapshotTransferReadPayload::Metadata(payload),
)?;
let admitted = validate_response(&request, journal.record()?, response)?;
let receipt = qualify(&request, response)
.await
.map_err(IcSnapshotMetadataSettlementError::Qualification)?;
if receipt.outcome != MutationOutcomeRecord::Applied
|| receipt.attempt != request.mutation_attempt()
|| receipt.request != payload.digest().hash()
{
return Err(IcSnapshotMetadataSettlementError::ReceiptRequired);
}
let IcSnapshotTransferReadReply::Metadata(metadata) = admitted.reply() else {
return Err(IcSnapshotMetadataSettlementError::ReceiptRequired);
};
stage.layout()?;
journal.record_mutation(receipt)?;
drop(journal);
Ok(stage.checkpoint(metadata.digest())?)
}
#[derive(Debug, Error)]
pub enum IcSnapshotMetadataExecutionError<E: std::error::Error + 'static> {
#[error("snapshot metadata requires an exact singleton original stage")]
OriginalMismatch,
#[error(transparent)]
Read(#[from] IcSnapshotTransferReadExecutionError<E>),
#[error("snapshot metadata reply settlement failed: {source}")]
AfterReply {
source: IcSnapshotMetadataSettlementError<E>,
response: Box<IcSnapshotTransferReadResponse>,
},
}
#[derive(Debug, Error)]
pub enum IcSnapshotMetadataSettlementError<E: std::error::Error + 'static> {
#[error("snapshot metadata qualification failed: {0}")]
Qualification(#[source] E),
#[error("snapshot metadata requires an explicit original Applied receipt")]
ReceiptRequired,
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Request(#[from] IcSnapshotTransferReadError),
#[error(transparent)]
Association(#[from] IcSnapshotTransferReadAssociationError),
#[error(transparent)]
Checkpoint(#[from] ExecutionStageCheckpointError),
}
#[cfg(all(test, unix))]
mod tests;