Skip to main content

ic_backup/workflow/ic_snapshot_download/
mod.rs

1//! Complete original planned data transfer with explicit integration-qualified receipts.
2
3use 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
26/// Stream every original data request into an existing private writer and publish it.
27///
28/// The metadata read and its independently qualified receipt/checkpoint must already
29/// exist. Prepare the metadata-derived data stage and its complete original journals,
30/// retain exact token/raw-ID association, then create the writer with exact raw metadata.
31/// Hold no attempt guards on entry. Read-free plans and any consumed original attempt
32/// reject; this is a fresh transfer, never partial-read resume or artifact recovery.
33///
34/// For every request, `admit` qualifies fresh access, metadata/command custody and
35/// application requirements under the existing single-read coordinator. `qualify`
36/// must independently authenticate exact original response attribution and return an
37/// explicit original Applied receipt. Wire shape alone cannot qualify that receipt.
38/// It runs under the selected original journal lock. Append exact decoded bytes before
39/// the sole journal owner records that receipt; only then may the next dependent read
40/// run. Failed qualification, append or persistence stops without another provider call.
41/// Admission, submission and qualification are awaited under their selected journal
42/// guards. Retain each reply durably before cancellable qualification work. Dropping
43/// this future consumes the writer, retaining partial bytes, earlier receipts and
44/// the current pending reservation; it grants no resume or repeat-call authority.
45///
46/// All original limits/records remain unchanged. Errors consume the writer and preserve
47/// partial bytes, pending spending and every recorded receipt. Returned bounded replies
48/// survive post-read rejection. Success uses the existing fresh checksum/durable
49/// publisher, but supplies no immutable manifest/checkpoint, complete product backup,
50/// terminal proof, default provider, refund or fence/reference release. Stable bytes,
51/// complete authenticated backend transfer and application safety stay integration-owned.
52/// # Errors
53/// Rejects changed stage/plan/metadata/writer, prior consumption, spending, admission,
54/// provider or receipt failures, malformed replies and local IO/publication failure.
55pub 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        // Planning already bounded the ordinal by the original stage allowance.
87        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/// Failures preserve original spending and partial/published artifact evidence.
161#[derive(Debug, Error)]
162pub enum IcSnapshotDownloadExecutionError<E: std::error::Error + 'static> {
163    /// Exact retained writer, stage or metadata differs.
164    #[error("snapshot download originals differ")]
165    OriginalMismatch,
166    /// A prior attempt cannot be replayed by this fresh transfer entrypoint.
167    #[error("snapshot download original data stage was already attempted")]
168    AlreadyAttempted,
169    /// Existing original download plan/binding admission.
170    #[error(transparent)]
171    Planning(#[from] IcSnapshotDownloadPlanningError),
172    /// Existing retained stage/ancestor admission.
173    #[error(transparent)]
174    Stage(#[from] ExecutionWorkflowPersistenceError),
175    /// Complete original progress admission.
176    #[error(transparent)]
177    Progress(#[from] ExecutionProgressPersistenceError),
178    /// One original read failed; its owner retains any returned bounded reply.
179    #[error(transparent)]
180    Read(#[from] IcSnapshotTransferReadExecutionError<E>),
181    /// Canonical writer/checksum/durable publication failure.
182    #[error(transparent)]
183    Artifact(#[from] IcSnapshotArtifactError),
184    /// Rejection after a returned read retains the exact bounded response.
185    #[error("snapshot download reply rejected for operation {operation_sequence}: {source}")]
186    AfterReply {
187        /// Exact original data operation.
188        operation_sequence: u64,
189        /// Original rejection without undoing receipt or spending.
190        source: IcSnapshotDownloadReplyError<E>,
191        /// Exact bounded returned bytes and passive claims.
192        response: Box<IcSnapshotTransferReadResponse>,
193    },
194    /// Publication succeeded but stage re-admission failed; retain the published checksum.
195    #[error("snapshot download stage changed after publication: {source}")]
196    AfterPublication {
197        /// Retained stage/ancestor admission failure.
198        source: ExecutionWorkflowPersistenceError,
199        /// Existing durable publisher's returned checksum.
200        checksum: ArtifactChecksumRecord,
201    },
202}
203
204/// Rejection of a returned data reply while retaining its original attempt.
205#[derive(Debug, Error)]
206pub enum IcSnapshotDownloadReplyError<E: std::error::Error + 'static> {
207    /// Fresh authenticated exact attribution was not qualified by the integration.
208    #[error("snapshot download reply qualification failed: {0}")]
209    Qualification(#[source] E),
210    /// Require the explicit exact original Applied receipt, without outcome inference.
211    #[error("snapshot download requires an explicit original Applied receipt")]
212    ReceiptRequired,
213    /// Existing original stage admission.
214    #[error(transparent)]
215    Stage(#[from] ExecutionWorkflowPersistenceError),
216    /// Original selected journal/spending/receipt admission.
217    #[error(transparent)]
218    Journal(#[from] AttemptJournalError),
219    /// Original request/current reservation admission.
220    #[error(transparent)]
221    Request(#[from] IcSnapshotTransferReadError),
222    /// Existing bounded passive decoder/association.
223    #[error(transparent)]
224    Association(#[from] IcSnapshotTransferReadAssociationError),
225    /// Exact local bytes/custody admission failed.
226    #[error(transparent)]
227    Artifact(#[from] IcSnapshotArtifactError),
228}
229
230#[cfg(test)]
231mod tests;