mod allocation;
#[cfg(unix)]
mod data;
#[cfg(unix)]
pub use data::{
IcSnapshotDataUploadExecutionError, IcSnapshotDataUploadReplyError, upload_snapshot_data,
};
pub use allocation::{
IcSnapshotAllocationExecutionError, IcSnapshotAllocationSettlementError, allocate_snapshot,
};
use crate::{
model::{
ic_mutation::IcMutationAcknowledgement,
ic_snapshot_upload::{
IcSnapshotUploadAttempt, IcSnapshotUploadAttemptError, IcSnapshotUploadRequest,
original_authority,
},
operation_plan::OperationPlanError,
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
ExecutionStageGuard, ExecutionWorkflowPersistenceError,
},
policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
ports::{ic_mutation::IcMutationProviderError, ic_snapshot_upload::IcSnapshotUploadProvider},
};
use thiserror::Error;
pub fn upload_snapshot<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
operation_sequence: u64,
payload: &IcSnapshotUploadRequest<'_>,
provider: &mut impl IcSnapshotUploadProvider,
admit: impl FnOnce(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
) -> Result<IcMutationAcknowledgement, IcSnapshotUploadExecutionError<E>> {
let plan = stage.plan();
let authority = plan.attempt_authority(operation_sequence)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
original_authority(plan, operation_sequence, journal.record()?, payload)?;
journal.reserve_planned_mutation(&plan.digest())?;
let request =
IcSnapshotUploadAttempt::new(plan, operation_sequence, journal.record()?, payload)?;
admit(&request).map_err(IcSnapshotUploadExecutionError::Admission)?;
stage.layout()?;
let acknowledgement = provider.submit_upload(&request)?;
let association = match journal.record() {
Ok(record) => record,
Err(source) => {
return Err(IcSnapshotUploadExecutionError::AfterReplyJournal {
source,
acknowledgement: Box::new(acknowledgement),
});
}
};
if let Err(source) = validate_acknowledgement(&request, association, &acknowledgement) {
return Err(IcSnapshotUploadExecutionError::Association {
source,
acknowledgement: Box::new(acknowledgement),
});
}
if let Err(source) = stage.layout() {
return Err(IcSnapshotUploadExecutionError::AfterReplyStage {
source,
acknowledgement: Box::new(acknowledgement),
});
}
Ok(acknowledgement)
}
#[derive(Debug, Error)]
pub enum IcSnapshotUploadExecutionError<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] IcSnapshotUploadAttemptError),
#[error("fresh snapshot upload admission failed: {0}")]
Admission(#[source] E),
#[error(transparent)]
Provider(#[from] IcMutationProviderError),
#[error("snapshot upload acknowledgement association failed: {source}")]
Association {
source: IcSnapshotUploadAssociationError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
#[error("snapshot upload journal changed after reply: {source}")]
AfterReplyJournal {
source: AttemptJournalError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
#[error("snapshot upload stage changed after reply: {source}")]
AfterReplyStage {
source: ExecutionWorkflowPersistenceError,
acknowledgement: Box<IcMutationAcknowledgement>,
},
}
#[cfg(all(test, unix))]
mod tests;