Skip to main content

ic_backup/workflow/ic_snapshot_upload/data/
mod.rs

1//! Complete finite original data uploads with explicit integration-qualified receipts.
2
3use 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
25/// Upload every exact original data operation, then checkpoint qualified reply evidence.
26///
27/// Prepare the source-derived data plan, allocation predecessor, retained stage and
28/// complete original journals first; hold no attempt guards. Empty plans and any prior
29/// consumption reject before dispatch: this is a fresh stage, never partial upload
30/// resume. Original source and destination/command custody remain integration-owned.
31/// Fresh guarded preparation buffers one payload and rechecks its immutable binding
32/// before existing single-upload reservation, mandatory fresh `admit` and one call.
33///
34/// `qualify` must independently authenticate exact original data-write attribution and
35/// durably retain original request/reply bytes before returning its explicit Applied
36/// receipt under the selected journal lock. An empty reply never supplies a receipt.
37/// Account remote qualification separately before its calls. Record through the sole
38/// attempt owner before the next dependency can dispatch; failures stop immediately.
39///
40/// After all exact receipts, release journal locks and checkpoint the full original plan
41/// and ordered canonical request/raw-reply evidence through the existing stage owner.
42/// This is declared stage settlement, not authenticated whole-backend completeness,
43/// same-release load/start admission, terminal proof or fence/source-reference release.
44/// No hidden observation, reissue, refund, default provider or journal is added.
45/// # Errors
46/// Rejects changed originals, consumed attempts, source/preparation failures, spending,
47/// admission/provider/reply/qualification/receipt or checkpoint failure. Bounded returned
48/// replies survive post-reply rejection; spending, recorded receipts and source remain.
49pub 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        // The single-call owner rechecks the exact full source/context/request before spending.
86        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/// Data-stage failure preserves source, original spending and returned acknowledgement bytes.
134#[derive(Debug, Error)]
135pub enum IcSnapshotDataUploadExecutionError<E: std::error::Error + 'static> {
136    /// Fresh-only entrypoint rejects any original consumed mutation or observation.
137    #[error("snapshot data upload original stage was already attempted")]
138    AlreadyAttempted,
139    /// Existing exact original data-plan/binding admission.
140    #[error(transparent)]
141    Planning(#[from] IcSnapshotDataUploadPlanningError),
142    /// Existing retained workflow/stage/ancestor admission.
143    #[error(transparent)]
144    Stage(#[from] ExecutionWorkflowPersistenceError),
145    /// Complete original progress admission.
146    #[error(transparent)]
147    Progress(#[from] ExecutionProgressPersistenceError),
148    /// Fresh exact source preparation failed before the selected operation reserves.
149    #[error(transparent)]
150    Source(#[from] IcSnapshotUploadArtifactError),
151    /// One exact original upload failed; its owner retains bounded returned evidence.
152    #[error(transparent)]
153    Upload(#[from] IcSnapshotUploadExecutionError<E>),
154    /// Receipt settlement failed after a bounded returned acknowledgement.
155    #[error("snapshot data upload reply rejected for operation {operation_sequence}: {source}")]
156    AfterReply {
157        /// Exact original failed operation sequence.
158        operation_sequence: u64,
159        /// Original qualification or persistence rejection.
160        source: IcSnapshotDataUploadReplyError<E>,
161        /// Exact bounded reply, without an inferred outcome.
162        acknowledgement: Box<IcMutationAcknowledgement>,
163    },
164    /// Existing all-Applied stage checkpoint failed; receipts and occupied evidence remain.
165    #[error(transparent)]
166    Checkpoint(#[from] ExecutionStageCheckpointError),
167}
168
169/// Post-reply data-write rejection, without a second outcome/accounting owner.
170#[derive(Debug, Error)]
171pub enum IcSnapshotDataUploadReplyError<E: std::error::Error + 'static> {
172    /// Independent attribution/authentication or durable original-byte retention failed.
173    #[error("snapshot data upload qualification failed: {0}")]
174    Qualification(#[source] E),
175    /// Require the explicit exact original Applied receipt.
176    #[error("snapshot data upload requires an explicit original Applied receipt")]
177    ReceiptRequired,
178    /// Original workflow/stage/ancestor admission failed.
179    #[error(transparent)]
180    Stage(#[from] ExecutionWorkflowPersistenceError),
181    /// Existing selected journal/receipt persistence failed.
182    #[error(transparent)]
183    Journal(#[from] AttemptJournalError),
184    /// Existing original payload/current reservation binding failed.
185    #[error(transparent)]
186    Request(#[from] IcSnapshotUploadAttemptError),
187    /// Existing bounded acknowledgement association failed.
188    #[error(transparent)]
189    Association(#[from] IcSnapshotUploadAssociationError),
190}