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///
42/// All original limits/records remain unchanged. Errors consume the writer and preserve
43/// partial bytes, pending spending and every recorded receipt. Returned bounded replies
44/// survive post-read rejection. Success uses the existing fresh checksum/durable
45/// publisher, but supplies no immutable manifest/checkpoint, complete product backup,
46/// terminal proof, default provider, refund or fence/reference release. Stable bytes,
47/// complete authenticated backend transfer and application safety stay integration-owned.
48/// # Errors
49/// Rejects changed stage/plan/metadata/writer, prior consumption, spending, admission,
50/// provider or receipt failures, malformed replies and local IO/publication failure.
51pub 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        // Planning already bounded the ordinal by the original stage allowance.
83        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/// Failures preserve original spending and partial/published artifact evidence.
154#[derive(Debug, Error)]
155pub enum IcSnapshotDownloadExecutionError<E: std::error::Error + 'static> {
156    /// Exact retained writer, stage or metadata differs.
157    #[error("snapshot download originals differ")]
158    OriginalMismatch,
159    /// A prior attempt cannot be replayed by this fresh transfer entrypoint.
160    #[error("snapshot download original data stage was already attempted")]
161    AlreadyAttempted,
162    /// Existing original download plan/binding admission.
163    #[error(transparent)]
164    Planning(#[from] IcSnapshotDownloadPlanningError),
165    /// Existing retained stage/ancestor admission.
166    #[error(transparent)]
167    Stage(#[from] ExecutionWorkflowPersistenceError),
168    /// Complete original progress admission.
169    #[error(transparent)]
170    Progress(#[from] ExecutionProgressPersistenceError),
171    /// One original read failed; its owner retains any returned bounded reply.
172    #[error(transparent)]
173    Read(#[from] IcSnapshotTransferReadExecutionError<E>),
174    /// Canonical writer/checksum/durable publication failure.
175    #[error(transparent)]
176    Artifact(#[from] IcSnapshotArtifactError),
177    /// Rejection after a returned read retains the exact bounded response.
178    #[error("snapshot download reply rejected for operation {operation_sequence}: {source}")]
179    AfterReply {
180        /// Exact original data operation.
181        operation_sequence: u64,
182        /// Original rejection without undoing receipt or spending.
183        source: IcSnapshotDownloadReplyError<E>,
184        /// Exact bounded returned bytes and passive claims.
185        response: Box<IcSnapshotTransferReadResponse>,
186    },
187    /// Publication succeeded but stage re-admission failed; retain the published checksum.
188    #[error("snapshot download stage changed after publication: {source}")]
189    AfterPublication {
190        /// Retained stage/ancestor admission failure.
191        source: ExecutionWorkflowPersistenceError,
192        /// Existing durable publisher's returned checksum.
193        checksum: ArtifactChecksumRecord,
194    },
195}
196
197/// Rejection of a returned data reply while retaining its original attempt.
198#[derive(Debug, Error)]
199pub enum IcSnapshotDownloadReplyError<E: std::error::Error + 'static> {
200    /// Fresh authenticated exact attribution was not qualified by the integration.
201    #[error("snapshot download reply qualification failed: {0}")]
202    Qualification(#[source] E),
203    /// Require the explicit exact original Applied receipt, without outcome inference.
204    #[error("snapshot download requires an explicit original Applied receipt")]
205    ReceiptRequired,
206    /// Existing original stage admission.
207    #[error(transparent)]
208    Stage(#[from] ExecutionWorkflowPersistenceError),
209    /// Original selected journal/spending/receipt admission.
210    #[error(transparent)]
211    Journal(#[from] AttemptJournalError),
212    /// Original request/current reservation admission.
213    #[error(transparent)]
214    Request(#[from] IcSnapshotTransferReadError),
215    /// Existing bounded passive decoder/association.
216    #[error(transparent)]
217    Association(#[from] IcSnapshotTransferReadAssociationError),
218    /// Exact local bytes/custody admission failed.
219    #[error(transparent)]
220    Artifact(#[from] IcSnapshotArtifactError),
221}
222
223#[cfg(test)]
224mod tests;