Skip to main content

ic_backup/workflow/ic_snapshot_upload/
mod.rs

1//! One fresh original source-bound metadata/data upload; no automatic receipt or retry.
2
3mod allocation;
4#[cfg(unix)]
5mod data;
6#[cfg(unix)]
7pub use data::{
8    IcSnapshotDataUploadExecutionError, IcSnapshotDataUploadReplyError, upload_snapshot_data,
9};
10
11pub use allocation::{
12    IcSnapshotAllocationExecutionError, IcSnapshotAllocationSettlementError, allocate_snapshot,
13};
14
15use crate::{
16    model::{
17        ic_mutation::IcMutationAcknowledgement,
18        ic_snapshot_upload::{
19            IcSnapshotUploadAttempt, IcSnapshotUploadAttemptError, IcSnapshotUploadRequest,
20            original_authority,
21        },
22        operation_plan::OperationPlanError,
23    },
24    ops::persistence::{
25        AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
26        ExecutionStageGuard, ExecutionWorkflowPersistenceError,
27    },
28    policy::ic_snapshot_upload::{IcSnapshotUploadAssociationError, validate_acknowledgement},
29    ports::{ic_mutation::IcMutationProviderError, ic_snapshot_upload::IcSnapshotUploadProvider},
30};
31use thiserror::Error;
32
33/// Durably reserve one exact original upload, freshly admit it and submit once.
34///
35/// Prepare the retained stage and every original journal before entry; hold no attempt
36/// guards. Guarded source preparation and exact metadata/tree/payload retention precede
37/// the upload plan. Data bytes and their independently attributed destination must be
38/// bound before that data stage is published, within its original workflow allocation.
39/// This checks the canonical full source/upload context and binding before spending,
40/// then delegates reservation/prerequisites to the existing complete-plan owner.
41/// Missing, pending, Applied or exhausted originals never reach provider invocation.
42///
43/// The mandatory fallible `admit` callback qualifies authentic complete source bytes,
44/// actual fresh controllers and application/restore requirements, exact metadata/data
45/// destination attribution, stable source/payload custody and exclusive never-dispatched
46/// command custody. It runs under the selected journal lock; remote preflight requires
47/// its own prior accounting. Local checks grant no byte fence or allocation authority.
48/// The existing provider sends only the exact accounted replicated update, with no
49/// batching, retries, hidden observations or implicit funding. No default is installed.
50/// Await admission and submission under the selected original guard. Cancellation
51/// releases exclusion without changing spending or permitting reentry. The provider
52/// receives the currently guarded original journal; the core supplies no executor.
53///
54/// Re-admit original stage/ancestor evidence before and after dispatch and check the
55/// bounded acknowledgement through the canonical upload policy/decoder. Success
56/// remains pending; integrations independently authenticate attribution and retain
57/// exact replies before explicitly recording receipts with the existing journal owner.
58/// A returned ID/empty reply alone grants no allocation outcome, learned data authority,
59/// complete transfer, load/start admission, terminal or fence/source-reference release.
60/// # Errors
61/// Rejects changed originals/payload/context, spending, fresh admission, provider or
62/// reply failures. Post-reservation failures retain consumption; post-reply rejections
63/// retain the bounded acknowledgement. No failure grants redispatch/refund/cleanup.
64pub async fn upload_snapshot<E: std::error::Error + 'static>(
65    stage: &ExecutionStageGuard<'_>,
66    operation_sequence: u64,
67    payload: &IcSnapshotUploadRequest<'_>,
68    provider: &mut impl IcSnapshotUploadProvider,
69    admit: impl AsyncFnOnce(&IcSnapshotUploadAttempt<'_, '_>) -> Result<(), E>,
70) -> Result<IcMutationAcknowledgement, IcSnapshotUploadExecutionError<E>> {
71    let plan = stage.plan();
72    let authority = plan.attempt_authority(operation_sequence)?;
73    let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
74    original_authority(plan, operation_sequence, journal.record()?, payload)?;
75    journal.reserve_planned_mutation(&plan.digest())?;
76    let request =
77        IcSnapshotUploadAttempt::new(plan, operation_sequence, journal.record()?, payload)?;
78    admit(&request)
79        .await
80        .map_err(IcSnapshotUploadExecutionError::Admission)?;
81    stage.layout()?;
82    let acknowledgement = provider.submit_upload(&request, journal.record()?).await?;
83    let association = match journal.record() {
84        Ok(record) => record,
85        Err(source) => {
86            return Err(IcSnapshotUploadExecutionError::AfterReplyJournal {
87                source,
88                acknowledgement: Box::new(acknowledgement),
89            });
90        }
91    };
92    if let Err(source) = validate_acknowledgement(&request, association, &acknowledgement) {
93        return Err(IcSnapshotUploadExecutionError::Association {
94            source,
95            acknowledgement: Box::new(acknowledgement),
96        });
97    }
98    if let Err(source) = stage.layout() {
99        return Err(IcSnapshotUploadExecutionError::AfterReplyStage {
100            source,
101            acknowledgement: Box::new(acknowledgement),
102        });
103    }
104    Ok(acknowledgement)
105}
106
107/// Upload-step failures retain original spending and any bounded returned acknowledgement.
108#[derive(Debug, Error)]
109pub enum IcSnapshotUploadExecutionError<E: std::error::Error + 'static> {
110    /// Original operation authority cannot be derived.
111    #[error(transparent)]
112    Plan(#[from] OperationPlanError),
113    /// Workflow, stage or original ancestor evidence failed admission.
114    #[error(transparent)]
115    Stage(#[from] ExecutionWorkflowPersistenceError),
116    /// Original selected journal admission failed.
117    #[error(transparent)]
118    Journal(#[from] AttemptJournalError),
119    /// Complete original progress, prerequisites or durable reservation failed.
120    #[error(transparent)]
121    Progress(#[from] ExecutionProgressPersistenceError),
122    /// Canonical original mutation binding or current reservation differs.
123    #[error(transparent)]
124    Request(#[from] IcSnapshotUploadAttemptError),
125    /// Fresh integration-owned admission rejected after reservation.
126    #[error("fresh snapshot upload admission failed: {0}")]
127    Admission(#[source] E),
128    /// The single provider call failed; spending remains pending.
129    #[error(transparent)]
130    Provider(#[from] IcMutationProviderError),
131    /// Passive bounded association rejected a returned acknowledgement.
132    #[error("snapshot upload acknowledgement association failed: {source}")]
133    Association {
134        /// Existing canonical mutation association rejection.
135        source: IcSnapshotUploadAssociationError,
136        /// Exact returned acknowledgement, without authenticated outcome.
137        acknowledgement: Box<IcMutationAcknowledgement>,
138    },
139    /// Selected journal re-admission failed after a reply.
140    #[error("snapshot upload journal changed after reply: {source}")]
141    AfterReplyJournal {
142        /// Original journal rejection.
143        source: AttemptJournalError,
144        /// Exact bounded returned acknowledgement.
145        acknowledgement: Box<IcMutationAcknowledgement>,
146    },
147    /// Stage or ancestor re-admission failed after a reply.
148    #[error("snapshot upload stage changed after reply: {source}")]
149    AfterReplyStage {
150        /// Original stage or ancestor rejection.
151        source: ExecutionWorkflowPersistenceError,
152        /// Exact bounded returned acknowledgement.
153        acknowledgement: Box<IcMutationAcknowledgement>,
154    },
155}
156
157#[cfg(all(test, unix))]
158mod tests;