ic_backup/workflow/ic_snapshot_upload/allocation/
mod.rs1use super::{IcSnapshotUploadExecutionError, upload_snapshot};
4use crate::{
5 model::{
6 attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
7 execution_workflow::ExecutionStagePredecessorRecord,
8 ic_mutation::IcMutationAcknowledgement,
9 ic_snapshot_upload::{
10 IcSnapshotUploadAttempt, IcSnapshotUploadAttemptError, IcSnapshotUploadKind,
11 IcSnapshotUploadRequest,
12 },
13 },
14 ops::persistence::{
15 AttemptJournalError, AttemptJournalGuard, ExecutionStageCheckpointError,
16 ExecutionStageGuard, ExecutionWorkflowPersistenceError,
17 },
18 policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
19 ports::ic_snapshot_upload::IcSnapshotUploadProvider,
20};
21use thiserror::Error;
22
23pub fn allocate_snapshot<E: std::error::Error + 'static>(
48 stage: &ExecutionStageGuard<'_>,
49 operation_sequence: u64,
50 payload: &IcSnapshotUploadRequest<'_>,
51 provider: &mut impl IcSnapshotUploadProvider,
52 admit: impl FnOnce(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
53 qualify: impl FnOnce(
54 &IcSnapshotUploadAttempt<'_, '_>,
55 &IcMutationAcknowledgement,
56 ) -> Result<MutationReceiptRequest, E>,
57) -> Result<
58 (IcMutationAcknowledgement, ExecutionStagePredecessorRecord),
59 IcSnapshotAllocationExecutionError<E>,
60> {
61 let operations = stage.plan().operations();
62 if !matches!(payload.kind(), IcSnapshotUploadKind::Metadata)
63 || operations.len() != 1
64 || operations[0].operation_sequence() != operation_sequence
65 || operations[0].request() != payload.binding_digest().hash()
66 {
67 return Err(IcSnapshotAllocationExecutionError::OriginalMismatch);
68 }
69 let acknowledgement = upload_snapshot(stage, operation_sequence, payload, provider, admit)?;
70 match settle_allocation(
71 stage,
72 operation_sequence,
73 payload,
74 &acknowledgement,
75 qualify,
76 ) {
77 Ok(predecessor) => Ok((acknowledgement, predecessor)),
78 Err(source) => Err(IcSnapshotAllocationExecutionError::AfterReply {
79 source,
80 acknowledgement: Box::new(acknowledgement),
81 }),
82 }
83}
84
85fn settle_allocation<E: std::error::Error + 'static>(
86 stage: &ExecutionStageGuard<'_>,
87 sequence: u64,
88 payload: &IcSnapshotUploadRequest<'_>,
89 acknowledgement: &IcMutationAcknowledgement,
90 qualify: impl FnOnce(
91 &IcSnapshotUploadAttempt<'_, '_>,
92 &IcMutationAcknowledgement,
93 ) -> Result<MutationReceiptRequest, E>,
94) -> Result<ExecutionStagePredecessorRecord, IcSnapshotAllocationSettlementError<E>> {
95 let plan = stage.plan();
96 let authority = plan
97 .attempt_authority(sequence)
98 .map_err(IcSnapshotUploadAttemptError::from)?;
99 let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
100 let request = IcSnapshotUploadAttempt::new(plan, sequence, journal.record()?, payload)?;
101 let admitted = validate_acknowledgement(&request, journal.record()?, acknowledgement)?;
102 let receipt = qualify(&request, acknowledgement)
103 .map_err(IcSnapshotAllocationSettlementError::Qualification)?;
104 if receipt.outcome != MutationOutcomeRecord::Applied
105 || receipt.attempt != request.mutation_attempt()
106 || receipt.request != payload.binding_digest().hash()
107 {
108 return Err(IcSnapshotAllocationSettlementError::ReceiptRequired);
109 }
110 stage.layout()?;
111 journal.record_mutation(receipt)?;
112 drop(journal);
113 Ok(stage.checkpoint(admitted.reply().digest())?)
114}
115
116#[derive(Debug, Error)]
118pub enum IcSnapshotAllocationExecutionError<E: std::error::Error + 'static> {
119 #[error("snapshot allocation requires an exact singleton original metadata stage")]
121 OriginalMismatch,
122 #[error(transparent)]
124 Upload(#[from] IcSnapshotUploadExecutionError<E>),
125 #[error("snapshot allocation reply settlement failed: {source}")]
127 AfterReply {
128 source: IcSnapshotAllocationSettlementError<E>,
130 acknowledgement: Box<IcMutationAcknowledgement>,
132 },
133}
134
135#[derive(Debug, Error)]
137pub enum IcSnapshotAllocationSettlementError<E: std::error::Error + 'static> {
138 #[error("snapshot allocation qualification failed: {0}")]
140 Qualification(#[source] E),
141 #[error("snapshot allocation requires an explicit original Applied receipt")]
143 ReceiptRequired,
144 #[error(transparent)]
146 Stage(#[from] ExecutionWorkflowPersistenceError),
147 #[error(transparent)]
149 Journal(#[from] AttemptJournalError),
150 #[error(transparent)]
152 Request(#[from] IcSnapshotUploadAttemptError),
153 #[error(transparent)]
155 Association(#[from] IcSnapshotUploadAssociationError),
156 #[error(transparent)]
158 Checkpoint(#[from] ExecutionStageCheckpointError),
159}
160
161#[cfg(all(test, unix))]
162mod tests;