use crate::{
model::{
ic_snapshot_transfer_read::{
IcSnapshotTransferReadError, IcSnapshotTransferReadPayload,
IcSnapshotTransferReadRequest, IcSnapshotTransferReadResponse,
},
operation_plan::OperationPlanError,
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
ExecutionStageGuard, ExecutionWorkflowPersistenceError,
},
policy::ic_snapshot_transfer_read::{
IcSnapshotTransferReadAssociationError, validate_response,
},
ports::{
ic_observation::IcObservationProviderError,
ic_snapshot_transfer_read::IcSnapshotTransferReadProvider,
},
};
use thiserror::Error;
pub async fn read_snapshot<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
operation_sequence: u64,
payload: IcSnapshotTransferReadPayload<'_, '_>,
provider: &mut impl IcSnapshotTransferReadProvider,
admit: impl AsyncFnOnce(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
) -> Result<IcSnapshotTransferReadResponse, IcSnapshotTransferReadExecutionError<E>> {
let plan = stage.plan();
let authority = plan.attempt_authority(operation_sequence)?;
payload.validate_binding(&authority)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
journal.reserve_planned_mutation(&plan.digest())?;
let request =
IcSnapshotTransferReadRequest::new(plan, operation_sequence, journal.record()?, payload)?;
admit(&request)
.await
.map_err(IcSnapshotTransferReadExecutionError::Admission)?;
stage.layout()?;
let response = provider.read_snapshot(&request, journal.record()?).await?;
let association = match journal.record() {
Ok(record) => record,
Err(source) => {
return Err(IcSnapshotTransferReadExecutionError::AfterReplyJournal {
source,
response: Box::new(response),
});
}
};
if let Err(source) = validate_response(&request, association, &response) {
return Err(IcSnapshotTransferReadExecutionError::Association {
source,
response: Box::new(response),
});
}
if let Err(source) = stage.layout() {
return Err(IcSnapshotTransferReadExecutionError::AfterReplyStage {
source,
response: Box::new(response),
});
}
Ok(response)
}
#[derive(Debug, Error)]
pub enum IcSnapshotTransferReadExecutionError<E: std::error::Error + 'static> {
#[error(transparent)]
Plan(#[from] OperationPlanError),
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Progress(#[from] ExecutionProgressPersistenceError),
#[error(transparent)]
Request(#[from] IcSnapshotTransferReadError),
#[error("fresh snapshot read admission failed: {0}")]
Admission(#[source] E),
#[error(transparent)]
Provider(#[from] IcObservationProviderError),
#[error("snapshot read response association failed: {source}")]
Association {
source: IcSnapshotTransferReadAssociationError,
response: Box<IcSnapshotTransferReadResponse>,
},
#[error("snapshot read journal changed after reply: {source}")]
AfterReplyJournal {
source: AttemptJournalError,
response: Box<IcSnapshotTransferReadResponse>,
},
#[error("snapshot read stage changed after reply: {source}")]
AfterReplyStage {
source: ExecutionWorkflowPersistenceError,
response: Box<IcSnapshotTransferReadResponse>,
},
}
#[cfg(all(test, unix))]
mod tests;