ic_backup/workflow/ic_snapshot_upload/data/
mod.rs1use super::{IcSnapshotUploadExecutionError, upload_snapshot};
4use crate::{
5 model::{
6 artifacts::ArtifactChecksumRecord,
7 attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
8 execution_workflow::ExecutionStagePredecessorRecord,
9 ic_mutation::IcMutationAcknowledgement,
10 ic_snapshot_upload::{
11 IcSnapshotDataUploadPlan, IcSnapshotDataUploadPlanningError, IcSnapshotUploadAttempt,
12 IcSnapshotUploadAttemptError, IcSnapshotUploadRequest,
13 },
14 },
15 ops::persistence::{
16 AttemptJournalError, AttemptJournalGuard, DownloadJournalGuard,
17 ExecutionProgressPersistenceError, ExecutionStageCheckpointError, ExecutionStageGuard,
18 ExecutionWorkflowPersistenceError, IcSnapshotUploadArtifactError, read_execution_progress,
19 },
20 policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
21 ports::ic_snapshot_upload::IcSnapshotUploadProvider,
22};
23use thiserror::Error;
24
25pub fn upload_snapshot_data<E: std::error::Error + 'static>(
50 stage: &ExecutionStageGuard<'_>,
51 upload: &IcSnapshotDataUploadPlan<'_, '_, '_>,
52 source: &DownloadJournalGuard<'_>,
53 snapshot: &str,
54 provider: &mut impl IcSnapshotUploadProvider,
55 mut admit: impl FnMut(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
56 mut qualify: impl FnMut(
57 &IcSnapshotUploadAttempt<'_, '_>,
58 &IcMutationAcknowledgement,
59 ) -> Result<MutationReceiptRequest, E>,
60) -> Result<ExecutionStagePredecessorRecord, IcSnapshotDataUploadExecutionError<E>> {
61 upload.validate_binding(stage.binding())?;
62 let plan = upload
63 .plan()
64 .ok_or(IcSnapshotDataUploadPlanningError::NoDataWrites)?;
65 if stage.plan() != plan {
66 return Err(IcSnapshotDataUploadPlanningError::OriginalMismatch.into());
67 }
68 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
69 if progress.attempts.mutations_used != 0 || progress.attempts.observations_used != 0 {
70 return Err(IcSnapshotDataUploadExecutionError::AlreadyAttempted);
71 }
72 let mut evidence = b"ic-backup/ic-snapshot-data-upload/v1\0".to_vec();
73 evidence.extend_from_slice(plan.digest().hash().as_bytes());
74 for (index, kind) in upload.kinds().iter().enumerate() {
75 let sequence =
76 u64::try_from(index).map_err(|_| IcSnapshotDataUploadPlanningError::CountOverflow)?;
77 let metadata = upload.metadata();
78 let payload = source.prepare_ic_snapshot_upload_data(
79 metadata.source_plan(),
80 snapshot,
81 metadata,
82 upload.destination(),
83 kind.clone(),
84 )?;
85 let acknowledgement = upload_snapshot(stage, sequence, &payload, provider, &mut admit)?;
87 let digest = record_qualified(stage, sequence, &payload, &acknowledgement, &mut qualify)
88 .map_err(|source| IcSnapshotDataUploadExecutionError::AfterReply {
89 operation_sequence: sequence,
90 source,
91 acknowledgement: Box::new(acknowledgement),
92 })?;
93 evidence.extend_from_slice(digest.hash().as_bytes());
94 }
95 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
96 if progress.applied_operations != upload.kinds().len() {
97 return Err(IcSnapshotDataUploadExecutionError::AlreadyAttempted);
98 }
99 Ok(stage.checkpoint(ArtifactChecksumRecord::from_bytes(&evidence))?)
100}
101
102fn record_qualified<E: std::error::Error + 'static>(
103 stage: &ExecutionStageGuard<'_>,
104 sequence: u64,
105 payload: &IcSnapshotUploadRequest<'_>,
106 acknowledgement: &IcMutationAcknowledgement,
107 qualify: &mut impl FnMut(
108 &IcSnapshotUploadAttempt<'_, '_>,
109 &IcMutationAcknowledgement,
110 ) -> Result<MutationReceiptRequest, E>,
111) -> Result<ArtifactChecksumRecord, IcSnapshotDataUploadReplyError<E>> {
112 let plan = stage.plan();
113 let authority = plan
114 .attempt_authority(sequence)
115 .map_err(IcSnapshotUploadAttemptError::from)?;
116 let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
117 let request = IcSnapshotUploadAttempt::new(plan, sequence, journal.record()?, payload)?;
118 let admitted = validate_acknowledgement(&request, journal.record()?, acknowledgement)?;
119 let receipt = qualify(&request, acknowledgement)
120 .map_err(IcSnapshotDataUploadReplyError::Qualification)?;
121 if receipt.outcome != MutationOutcomeRecord::Applied
122 || receipt.attempt != request.mutation_attempt()
123 || receipt.request != payload.binding_digest().hash()
124 {
125 return Err(IcSnapshotDataUploadReplyError::ReceiptRequired);
126 }
127 stage.layout()?;
128 journal.record_mutation(receipt)?;
129 stage.layout()?;
130 Ok(admitted.reply().digest())
131}
132
133#[derive(Debug, Error)]
135pub enum IcSnapshotDataUploadExecutionError<E: std::error::Error + 'static> {
136 #[error("snapshot data upload original stage was already attempted")]
138 AlreadyAttempted,
139 #[error(transparent)]
141 Planning(#[from] IcSnapshotDataUploadPlanningError),
142 #[error(transparent)]
144 Stage(#[from] ExecutionWorkflowPersistenceError),
145 #[error(transparent)]
147 Progress(#[from] ExecutionProgressPersistenceError),
148 #[error(transparent)]
150 Source(#[from] IcSnapshotUploadArtifactError),
151 #[error(transparent)]
153 Upload(#[from] IcSnapshotUploadExecutionError<E>),
154 #[error("snapshot data upload reply rejected for operation {operation_sequence}: {source}")]
156 AfterReply {
157 operation_sequence: u64,
159 source: IcSnapshotDataUploadReplyError<E>,
161 acknowledgement: Box<IcMutationAcknowledgement>,
163 },
164 #[error(transparent)]
166 Checkpoint(#[from] ExecutionStageCheckpointError),
167}
168
169#[derive(Debug, Error)]
171pub enum IcSnapshotDataUploadReplyError<E: std::error::Error + 'static> {
172 #[error("snapshot data upload qualification failed: {0}")]
174 Qualification(#[source] E),
175 #[error("snapshot data upload requires an explicit original Applied receipt")]
177 ReceiptRequired,
178 #[error(transparent)]
180 Stage(#[from] ExecutionWorkflowPersistenceError),
181 #[error(transparent)]
183 Journal(#[from] AttemptJournalError),
184 #[error(transparent)]
186 Request(#[from] IcSnapshotUploadAttemptError),
187 #[error(transparent)]
189 Association(#[from] IcSnapshotUploadAssociationError),
190}