ic_backup/workflow/ic_snapshot_download/
mod.rs1use crate::{
4 model::{
5 artifacts::ArtifactChecksumRecord,
6 attempt_journal::{MutationOutcomeRecord, MutationReceiptRequest},
7 ic_snapshot_download::{IcSnapshotDownloadPlan, IcSnapshotDownloadPlanningError},
8 ic_snapshot_transfer_read::{
9 IcSnapshotTransferReadError, IcSnapshotTransferReadPayload,
10 IcSnapshotTransferReadRequest, IcSnapshotTransferReadResponse,
11 },
12 },
13 ops::persistence::{
14 AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
15 ExecutionStageGuard, ExecutionWorkflowPersistenceError, IcSnapshotArtifactError,
16 IcSnapshotArtifactWriter, read_execution_progress,
17 },
18 policy::ic_snapshot_transfer_read::{
19 IcSnapshotTransferReadAssociationError, IcSnapshotTransferReadReply, validate_response,
20 },
21 ports::ic_snapshot_transfer_read::IcSnapshotTransferReadProvider,
22 workflow::ic_snapshot_transfer_read::{IcSnapshotTransferReadExecutionError, read_snapshot},
23};
24use thiserror::Error;
25
26pub async fn download_snapshot<E: std::error::Error + 'static>(
56 stage: &ExecutionStageGuard<'_>,
57 download: &IcSnapshotDownloadPlan<'_, '_>,
58 mut writer: IcSnapshotArtifactWriter<'_, '_, '_>,
59 provider: &mut impl IcSnapshotTransferReadProvider,
60 mut admit: impl AsyncFnMut(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
61 mut qualify: impl AsyncFnMut(
62 &IcSnapshotTransferReadRequest<'_, '_>,
63 &IcSnapshotTransferReadResponse,
64 ) -> Result<MutationReceiptRequest, E>,
65) -> Result<ArtifactChecksumRecord, IcSnapshotDownloadExecutionError<E>> {
66 download.validate_binding(stage.binding())?;
67 let plan = download
68 .plan()
69 .ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
70 if stage.plan() != plan
71 || download
72 .requests()
73 .iter()
74 .any(|request| request.metadata().digest() != writer.coverage().metadata().digest())
75 || writer.coverage().covered_region_bytes() != [0; 3]
76 || writer.coverage().covered_chunks() != 0
77 {
78 return Err(IcSnapshotDownloadExecutionError::OriginalMismatch);
79 }
80 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
81 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
82 if progress.attempts.mutations_used != 0 || progress.attempts.observations_used != 0 {
83 return Err(IcSnapshotDownloadExecutionError::AlreadyAttempted);
84 }
85 for (index, payload) in download.requests().iter().enumerate() {
86 let sequence =
88 u64::try_from(index).map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
89 let response = read_snapshot(
90 stage,
91 sequence,
92 IcSnapshotTransferReadPayload::Data(payload),
93 provider,
94 &mut admit,
95 )
96 .await?;
97 writer = append_qualified(stage, sequence, payload, writer, &response, &mut qualify)
98 .await
99 .map_err(|source| IcSnapshotDownloadExecutionError::AfterReply {
100 operation_sequence: sequence,
101 source,
102 response: Box::new(response),
103 })?;
104 }
105 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
106 if progress.applied_operations != download.requests().len() {
107 return Err(IcSnapshotDownloadExecutionError::AlreadyAttempted);
108 }
109 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
110 let checksum = writer.finish()?;
111 if let Err(source) = stage.layout() {
112 return Err(IcSnapshotDownloadExecutionError::AfterPublication { source, checksum });
113 }
114 Ok(checksum)
115}
116
117async fn append_qualified<'journal, 'layout, 'metadata, E: std::error::Error + 'static>(
118 stage: &ExecutionStageGuard<'_>,
119 sequence: u64,
120 payload: &crate::model::ic_snapshot_data::IcSnapshotDataRequest<'_>,
121 writer: IcSnapshotArtifactWriter<'journal, 'layout, 'metadata>,
122 response: &IcSnapshotTransferReadResponse,
123 qualify: &mut impl AsyncFnMut(
124 &IcSnapshotTransferReadRequest<'_, '_>,
125 &IcSnapshotTransferReadResponse,
126 ) -> Result<MutationReceiptRequest, E>,
127) -> Result<IcSnapshotArtifactWriter<'journal, 'layout, 'metadata>, IcSnapshotDownloadReplyError<E>>
128{
129 let plan = stage.plan();
130 let authority = plan
131 .attempt_authority(sequence)
132 .map_err(IcSnapshotTransferReadError::from)?;
133 let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
134 let request = IcSnapshotTransferReadRequest::new(
135 plan,
136 sequence,
137 journal.record()?,
138 IcSnapshotTransferReadPayload::Data(payload),
139 )?;
140 let admitted = validate_response(&request, journal.record()?, response)?;
141 let receipt = qualify(&request, response)
142 .await
143 .map_err(IcSnapshotDownloadReplyError::Qualification)?;
144 if receipt.outcome != MutationOutcomeRecord::Applied
145 || receipt.attempt != request.mutation_attempt()
146 || receipt.request != request.payload().digest().hash()
147 {
148 return Err(IcSnapshotDownloadReplyError::ReceiptRequired);
149 }
150 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
151 let IcSnapshotTransferReadReply::Data(reply) = admitted.reply() else {
152 return Err(IcSnapshotDownloadReplyError::ReceiptRequired);
153 };
154 let writer = writer.append(reply)?;
155 journal.record_mutation(receipt)?;
156 stage.layout()?;
157 Ok(writer)
158}
159
160#[derive(Debug, Error)]
162pub enum IcSnapshotDownloadExecutionError<E: std::error::Error + 'static> {
163 #[error("snapshot download originals differ")]
165 OriginalMismatch,
166 #[error("snapshot download original data stage was already attempted")]
168 AlreadyAttempted,
169 #[error(transparent)]
171 Planning(#[from] IcSnapshotDownloadPlanningError),
172 #[error(transparent)]
174 Stage(#[from] ExecutionWorkflowPersistenceError),
175 #[error(transparent)]
177 Progress(#[from] ExecutionProgressPersistenceError),
178 #[error(transparent)]
180 Read(#[from] IcSnapshotTransferReadExecutionError<E>),
181 #[error(transparent)]
183 Artifact(#[from] IcSnapshotArtifactError),
184 #[error("snapshot download reply rejected for operation {operation_sequence}: {source}")]
186 AfterReply {
187 operation_sequence: u64,
189 source: IcSnapshotDownloadReplyError<E>,
191 response: Box<IcSnapshotTransferReadResponse>,
193 },
194 #[error("snapshot download stage changed after publication: {source}")]
196 AfterPublication {
197 source: ExecutionWorkflowPersistenceError,
199 checksum: ArtifactChecksumRecord,
201 },
202}
203
204#[derive(Debug, Error)]
206pub enum IcSnapshotDownloadReplyError<E: std::error::Error + 'static> {
207 #[error("snapshot download reply qualification failed: {0}")]
209 Qualification(#[source] E),
210 #[error("snapshot download requires an explicit original Applied receipt")]
212 ReceiptRequired,
213 #[error(transparent)]
215 Stage(#[from] ExecutionWorkflowPersistenceError),
216 #[error(transparent)]
218 Journal(#[from] AttemptJournalError),
219 #[error(transparent)]
221 Request(#[from] IcSnapshotTransferReadError),
222 #[error(transparent)]
224 Association(#[from] IcSnapshotTransferReadAssociationError),
225 #[error(transparent)]
227 Artifact(#[from] IcSnapshotArtifactError),
228}
229
230#[cfg(test)]
231mod tests;