use crate::{
model::{
attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
execution_workflow::ExecutionStagePredecessorRecord,
ic_mutation::{IcMutationAcknowledgement, IcMutationRequest, IcMutationRequestError},
operation_plan::OperationPlanRecord,
restore_safety::{RestoreSafetyObservation, RestoreSafetyRequest},
},
ops::persistence::{
AttemptJournalError, AttemptJournalGuard, BackupLayoutGuard,
ExecutionProgressPersistenceError, ExecutionStageCheckpointError, ExecutionStageGuard,
ExecutionWorkflowPersistenceError, RestoreSafetyPersistenceError,
read_restore_safety_requirement,
},
policy::{
ic_mutation::{IcMutationAssociationError, IcMutationReplyView, validate_acknowledgement},
restore_safety::{RestoreSafetyError, validate},
},
ports::ic_mutation::{IcMutationProvider, IcMutationProviderError},
};
use thiserror::Error;
pub async fn restore_snapshot<E: std::error::Error + 'static>(
stage: &ExecutionStageGuard<'_>,
source_layout: &BackupLayoutGuard,
source_plan: &OperationPlanRecord,
safety: &RestoreSafetyRequest<'_>,
provider: &mut impl IcMutationProvider,
admit: impl AsyncFnOnce(
&IcMutationRequest<'_>,
&RestoreSafetyRequest<'_>,
) -> Result<RestoreSafetyObservation, E>,
qualify: impl AsyncFnOnce(
&IcMutationRequest<'_>,
&IcMutationAcknowledgement,
) -> Result<MutationReceiptRequest, E>,
) -> Result<(IcMutationAcknowledgement, ExecutionStagePredecessorRecord), IcRestoreExecutionError<E>>
{
let plan = stage.plan();
let sequence = safety.binding().operation_sequence();
if plan.operations().len() != 1 || plan.operations()[0].operation_sequence() != sequence {
return Err(IcRestoreExecutionError::OriginalMismatch);
}
let retain = || {
read_restore_safety_requirement(
stage.layout()?,
source_layout,
plan,
source_plan,
&safety.requirement().digest(),
)?;
Ok::<_, IcRestoreReplyError<E>>(())
};
retain()?;
let authority = plan
.attempt_authority(sequence)
.map_err(IcMutationRequestError::from)?;
safety
.wire()
.validate_mutation_binding(authority.binding())
.map_err(IcMutationRequestError::from)?;
let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
journal.reserve_planned_mutation(&plan.digest())?;
let request = IcMutationRequest::new(plan, sequence, journal.record()?, safety.wire())?;
let observation = admit(&request, safety)
.await
.map_err(IcRestoreExecutionError::Admission)?;
validate(safety, &observation)?;
retain()?;
let acknowledgement = provider
.submit_mutation(&request, journal.record()?)
.await?;
let settle = async || {
let admitted = validate_acknowledgement(&request, journal.record()?, &acknowledgement)?;
retain()?;
let receipt = qualify(&request, &acknowledgement)
.await
.map_err(IcRestoreReplyError::Qualification)?;
if receipt.outcome != MutationOutcomeRecord::Applied
|| receipt.attempt != request.mutation_attempt()
|| receipt.request != safety.wire().digest().hash()
{
return Err(IcRestoreReplyError::ReceiptRequired);
}
retain()?;
journal.record_mutation(receipt)?;
retain()?;
let IcMutationReplyView::Lifecycle(reply) = admitted.reply() else {
return Err(IcRestoreReplyError::ReceiptRequired);
};
Ok(reply.digest())
};
let evidence = match settle().await {
Ok(evidence) => evidence,
Err(source) => {
return Err(IcRestoreExecutionError::AfterReply {
source,
acknowledgement: Box::new(acknowledgement),
});
}
};
drop(journal);
match stage.checkpoint(evidence) {
Ok(predecessor) => Ok((acknowledgement, predecessor)),
Err(source) => Err(IcRestoreExecutionError::AfterReply {
source: IcRestoreReplyError::Checkpoint(source),
acknowledgement: Box::new(acknowledgement),
}),
}
}
#[derive(Debug, Error)]
pub enum IcRestoreExecutionError<E: std::error::Error + 'static> {
#[error("restore requires an exact singleton original load/start stage")]
OriginalMismatch,
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Progress(#[from] ExecutionProgressPersistenceError),
#[error(transparent)]
Request(#[from] IcMutationRequestError),
#[error(transparent)]
Retention(#[from] IcRestoreReplyError<E>),
#[error("fresh restore admission failed: {0}")]
Admission(#[source] E),
#[error(transparent)]
Safety(#[from] RestoreSafetyError),
#[error(transparent)]
Provider(#[from] IcMutationProviderError),
#[error("restore reply settlement failed: {source}")]
AfterReply {
source: IcRestoreReplyError<E>,
acknowledgement: Box<IcMutationAcknowledgement>,
},
}
#[derive(Debug, Error)]
pub enum IcRestoreReplyError<E: std::error::Error + 'static> {
#[error(transparent)]
Stage(#[from] ExecutionWorkflowPersistenceError),
#[error(transparent)]
Retention(#[from] RestoreSafetyPersistenceError),
#[error(transparent)]
Journal(#[from] AttemptJournalError),
#[error(transparent)]
Association(#[from] IcMutationAssociationError),
#[error("restore reply qualification failed: {0}")]
Qualification(#[source] E),
#[error("restore requires an explicit original Applied receipt")]
ReceiptRequired,
#[error(transparent)]
Checkpoint(#[from] ExecutionStageCheckpointError),
}
#[cfg(test)]
mod tests;