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 fn download_snapshot<E: std::error::Error + 'static>(
52 stage: &ExecutionStageGuard<'_>,
53 download: &IcSnapshotDownloadPlan<'_, '_>,
54 mut writer: IcSnapshotArtifactWriter<'_, '_, '_>,
55 provider: &mut impl IcSnapshotTransferReadProvider,
56 mut admit: impl FnMut(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
57 mut qualify: impl FnMut(
58 &IcSnapshotTransferReadRequest<'_, '_>,
59 &IcSnapshotTransferReadResponse,
60 ) -> Result<MutationReceiptRequest, E>,
61) -> Result<ArtifactChecksumRecord, IcSnapshotDownloadExecutionError<E>> {
62 download.validate_binding(stage.binding())?;
63 let plan = download
64 .plan()
65 .ok_or(IcSnapshotDownloadPlanningError::NoDataReads)?;
66 if stage.plan() != plan
67 || download
68 .requests()
69 .iter()
70 .any(|request| request.metadata().digest() != writer.coverage().metadata().digest())
71 || writer.coverage().covered_region_bytes() != [0; 3]
72 || writer.coverage().covered_chunks() != 0
73 {
74 return Err(IcSnapshotDownloadExecutionError::OriginalMismatch);
75 }
76 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
77 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
78 if progress.attempts.mutations_used != 0 || progress.attempts.observations_used != 0 {
79 return Err(IcSnapshotDownloadExecutionError::AlreadyAttempted);
80 }
81 for (index, payload) in download.requests().iter().enumerate() {
82 let sequence =
84 u64::try_from(index).map_err(|_| IcSnapshotDownloadPlanningError::CountOverflow)?;
85 let response = read_snapshot(
86 stage,
87 sequence,
88 IcSnapshotTransferReadPayload::Data(payload),
89 provider,
90 &mut admit,
91 )?;
92 writer = append_qualified(stage, sequence, payload, writer, &response, &mut qualify)
93 .map_err(|source| IcSnapshotDownloadExecutionError::AfterReply {
94 operation_sequence: sequence,
95 source,
96 response: Box::new(response),
97 })?;
98 }
99 let progress = read_execution_progress(stage.layout()?, &plan.digest())?;
100 if progress.applied_operations != download.requests().len() {
101 return Err(IcSnapshotDownloadExecutionError::AlreadyAttempted);
102 }
103 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
104 let checksum = writer.finish()?;
105 if let Err(source) = stage.layout() {
106 return Err(IcSnapshotDownloadExecutionError::AfterPublication { source, checksum });
107 }
108 Ok(checksum)
109}
110
111fn append_qualified<'journal, 'layout, 'metadata, E: std::error::Error + 'static>(
112 stage: &ExecutionStageGuard<'_>,
113 sequence: u64,
114 payload: &crate::model::ic_snapshot_data::IcSnapshotDataRequest<'_>,
115 writer: IcSnapshotArtifactWriter<'journal, 'layout, 'metadata>,
116 response: &IcSnapshotTransferReadResponse,
117 qualify: &mut impl FnMut(
118 &IcSnapshotTransferReadRequest<'_, '_>,
119 &IcSnapshotTransferReadResponse,
120 ) -> Result<MutationReceiptRequest, E>,
121) -> Result<IcSnapshotArtifactWriter<'journal, 'layout, 'metadata>, IcSnapshotDownloadReplyError<E>>
122{
123 let plan = stage.plan();
124 let authority = plan
125 .attempt_authority(sequence)
126 .map_err(IcSnapshotTransferReadError::from)?;
127 let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
128 let request = IcSnapshotTransferReadRequest::new(
129 plan,
130 sequence,
131 journal.record()?,
132 IcSnapshotTransferReadPayload::Data(payload),
133 )?;
134 let admitted = validate_response(&request, journal.record()?, response)?;
135 let receipt =
136 qualify(&request, response).map_err(IcSnapshotDownloadReplyError::Qualification)?;
137 if receipt.outcome != MutationOutcomeRecord::Applied
138 || receipt.attempt != request.mutation_attempt()
139 || receipt.request != request.payload().digest().hash()
140 {
141 return Err(IcSnapshotDownloadReplyError::ReceiptRequired);
142 }
143 writer.validate_transfer_origin(stage.layout()?.root(), plan.digest().hash())?;
144 let IcSnapshotTransferReadReply::Data(reply) = admitted.reply() else {
145 return Err(IcSnapshotDownloadReplyError::ReceiptRequired);
146 };
147 let writer = writer.append(reply)?;
148 journal.record_mutation(receipt)?;
149 stage.layout()?;
150 Ok(writer)
151}
152
153#[derive(Debug, Error)]
155pub enum IcSnapshotDownloadExecutionError<E: std::error::Error + 'static> {
156 #[error("snapshot download originals differ")]
158 OriginalMismatch,
159 #[error("snapshot download original data stage was already attempted")]
161 AlreadyAttempted,
162 #[error(transparent)]
164 Planning(#[from] IcSnapshotDownloadPlanningError),
165 #[error(transparent)]
167 Stage(#[from] ExecutionWorkflowPersistenceError),
168 #[error(transparent)]
170 Progress(#[from] ExecutionProgressPersistenceError),
171 #[error(transparent)]
173 Read(#[from] IcSnapshotTransferReadExecutionError<E>),
174 #[error(transparent)]
176 Artifact(#[from] IcSnapshotArtifactError),
177 #[error("snapshot download reply rejected for operation {operation_sequence}: {source}")]
179 AfterReply {
180 operation_sequence: u64,
182 source: IcSnapshotDownloadReplyError<E>,
184 response: Box<IcSnapshotTransferReadResponse>,
186 },
187 #[error("snapshot download stage changed after publication: {source}")]
189 AfterPublication {
190 source: ExecutionWorkflowPersistenceError,
192 checksum: ArtifactChecksumRecord,
194 },
195}
196
197#[derive(Debug, Error)]
199pub enum IcSnapshotDownloadReplyError<E: std::error::Error + 'static> {
200 #[error("snapshot download reply qualification failed: {0}")]
202 Qualification(#[source] E),
203 #[error("snapshot download requires an explicit original Applied receipt")]
205 ReceiptRequired,
206 #[error(transparent)]
208 Stage(#[from] ExecutionWorkflowPersistenceError),
209 #[error(transparent)]
211 Journal(#[from] AttemptJournalError),
212 #[error(transparent)]
214 Request(#[from] IcSnapshotTransferReadError),
215 #[error(transparent)]
217 Association(#[from] IcSnapshotTransferReadAssociationError),
218 #[error(transparent)]
220 Artifact(#[from] IcSnapshotArtifactError),
221}
222
223#[cfg(test)]
224mod tests;