use super::{IcSnapshotUploadExecutionError, upload_snapshot};
use crate::{
model::{
artifacts::ArtifactChecksumRecord,
attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
execution_workflow::ExecutionStagePredecessorRecord,
ic_mutation::IcMutationAcknowledgement,
ic_snapshot_upload::{
IcSnapshotDataUploadPlan, IcSnapshotDataUploadPlanningError, IcSnapshotUploadAttempt,
IcSnapshotUploadAttemptError, IcSnapshotUploadRequest,
},
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, DownloadJournalGuard,
ExecutionProgressPersistenceError, ExecutionStageCheckpointError, ExecutionStageGuard,
ExecutionWorkflowPersistenceError, IcSnapshotUploadArtifactError, read_execution_progress,
},
policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
ports::ic_snapshot_upload::IcSnapshotUploadProvider,
};
use thiserror::Error;
pub fn upload_snapshot_data<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
upload: &IcSnapshotDataUploadPlan<'_, '_, '_>,
source: &DownloadJournalGuard<'_>,
snapshot: &str,
provider: &mut impl IcSnapshotUploadProvider,
mut admit: impl FnMut(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
mut qualify: impl FnMut(
&IcSnapshotUploadAttempt<'_, '_>,
&IcMutationAcknowledgement,
) -> Result<MutationReceiptRequest, E>,
) -> Result<ExecutionStagePredecessorRecord, IcSnapshotDataUploadExecutionError<E>> {
upload.validate_binding(stage.binding())?;
let plan = upload
.plan()
.ok_or(IcSnapshotDataUploadPlanningError::NoDataWrites)?;
if stage.plan() != plan {
return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch.into());
}
let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
if progress.attempts.mutations_used != 0 || progress.attempts.observations_used != 0 {
return Err(IcSnapshotDataUploadExecutionError::AlreadyAttempted);
}
let mut evidence = b"ic-backup/ic-snapshot-data-upload/v1\0".to_vec();
evidence.extend_from_slice(plan.digest().hash().as_bytes());
for (index, kind) in upload.kinds().iter().enumerate() {
let sequence =
u64::try_from(index).map_err(|_| IcSnapshotDataUploadPlanningError::CountOverflow)?;
let metadata = upload.metadata();
let payload = source.prepare_ic_snapshot_upload_data(
metadata.source_plan(),
snapshot,
metadata,
upload.destination(),
kind.clone(),
)?;
let acknowledgement = upload_snapshot(stage, sequence, &payload, provider, &mut admit)?;
let digest = record_qualified(stage, sequence, &payload, &acknowledgement, &mut qualify)
.map_err(|source| IcSnapshotDataUploadExecutionError::AfterReply {
operation_sequence: sequence,
source,
acknowledgement: Box::new(acknowledgement),
})?;
evidence.extend_from_slice(digest.hash().as_bytes());
}
let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
if progress.applied_operations != upload.kinds().len() {
return Err(IcSnapshotDataUploadExecutionError::AlreadyAttempted);
}
Ok(stage.checkpoint(ArtifactChecksumRecord::from_bytes(&evidence))?)
}
fn record_qualified<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
sequence: u64,
payload: &IcSnapshotUploadRequest<'_>,
acknowledgement: &IcMutationAcknowledgement,
qualify: &mut impl FnMut(
&IcSnapshotUploadAttempt<'_, '_>,
&IcMutationAcknowledgement,
) -> Result<MutationReceiptRequest, E>,
) -> Result<ArtifactChecksumRecord, IcSnapshotDataUploadReplyError<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(IcSnapshotDataUploadReplyError::Qualification)?;
if receipt.outcome != MutationOutcomeRecord::Applied
|| receipt.attempt != request.mutation_attempt()
|| receipt.request != payload.binding_digest().hash()
{
return Err(IcSnapshotDataUploadReplyError::ReceiptRequired);
}
stage.layout()?;
journal.record_mutation(receipt)?;
stage.layout()?;
Ok(admitted.reply().digest())
}
#[derive(Debug, Error)]
pub enum IcSnapshotDataUploadExecutionError<E: std::error::Error + 'static> {
#[error("snapshot data upload original stage was already attempted")]
AlreadyAttempted,
#[error(transparent)]
Planning(#[from] IcSnapshotDataUploadPlanningError),
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Progress(#[from] ExecutionProgressPersistenceError),
#[error(transparent)]
Source(#[from] IcSnapshotUploadArtifactError),
#[error(transparent)]
Upload(#[from] IcSnapshotUploadExecutionError<E>),
#[error("snapshot data upload reply rejected for operation {operation_sequence}: {source}")]
AfterReply {
operation_sequence: u64,
source: IcSnapshotDataUploadReplyError<E>,
acknowledgement: Box<IcMutationAcknowledgement>,
},
#[error(transparent)]
Checkpoint(#[from] ExecutionStageCheckpointError),
}
#[derive(Debug, Error)]
pub enum IcSnapshotDataUploadReplyError<E: std::error::Error + 'static> {
#[error("snapshot data upload qualification failed: {0}")]
Qualification(#[source] E),
#[error("snapshot data upload 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),
}