Skip to main content

ic_backup/workflow/ic_snapshot_transfer_read/
mod.rs

1//! One freshly reserved metadata/data update; no automatic receipt or retry.
2
3use crate::{
4    model::{
5        ic_snapshot_transfer_read::{
6            IcSnapshotTransferReadError, IcSnapshotTransferReadPayload,
7            IcSnapshotTransferReadRequest, IcSnapshotTransferReadResponse,
8        },
9        operation_plan::OperationPlanError,
10    },
11    ops::persistence::{
12        AttemptJournalError, AttemptJournalGuard, ExecutionProgressPersistenceError,
13        ExecutionStageGuard, ExecutionWorkflowPersistenceError,
14    },
15    policy::ic_snapshot_transfer_read::{
16        IcSnapshotTransferReadAssociationError, validate_response,
17    },
18    ports::{
19        ic_observation::IcObservationProviderError,
20        ic_snapshot_transfer_read::IcSnapshotTransferReadProvider,
21    },
22};
23use thiserror::Error;
24
25/// Durably reserve one original planned read, freshly admit it and invoke its provider once.
26///
27/// Hold no attempt guards on entry. The selected original journal remains locked
28/// through admission, dispatch and passive response association. Missing originals,
29/// pending attempts and unfulfilled prerequisites reject before provider invocation.
30/// This only accepts a newly reserved attempt; reconstructed pending requests are
31/// never dispatched. Every failure after reservation retains its consumed allowance.
32/// Admission and provider submission are awaited under that same selected guard.
33/// Dropping the future releases the guard without refunding or granting reentry.
34///
35/// `admit` must qualify actual fresh method-specific access, original metadata/raw-ID
36/// custody, current application requirements and exclusive never-dispatched command
37/// custody for this exact request. There is no default admission. Any remote preflight
38/// calls require their own prior accounting. The provider retains its existing
39/// authenticated single-update contract, without retries or hidden observations.
40///
41/// The returned bounded reply has only structural association. Integrations must
42/// independently authenticate it and explicitly use the existing journal receipt
43/// owner to record a qualified outcome. Even success leaves the attempt pending.
44/// No artifact, complete-transfer, fence or terminal/reference-release proof follows.
45/// # Errors
46/// Rejects changed originals, payloads, spending, fresh admission, provider failure
47/// or malformed/mismatched replies. Replies returned before a later rejection are
48/// retained in the error. No failure grants a repeat call, refund or cleanup.
49pub async fn read_snapshot<E: std::error::Error + 'static>(
50    stage: &ExecutionStageGuard<'_>,
51    operation_sequence: u64,
52    payload: IcSnapshotTransferReadPayload<'_, '_>,
53    provider: &mut impl IcSnapshotTransferReadProvider,
54    admit: impl AsyncFnOnce(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
55) -> Result<IcSnapshotTransferReadResponse, IcSnapshotTransferReadExecutionError<E>> {
56    let plan = stage.plan();
57    let authority = plan.attempt_authority(operation_sequence)?;
58    payload.validate_binding(&authority)?;
59    let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
60    journal.reserve_planned_mutation(&plan.digest())?;
61    let request =
62        IcSnapshotTransferReadRequest::new(plan, operation_sequence, journal.record()?, payload)?;
63    admit(&request)
64        .await
65        .map_err(IcSnapshotTransferReadExecutionError::Admission)?;
66    stage.layout()?;
67    let response = provider.read_snapshot(&request, journal.record()?).await?;
68    let association = match journal.record() {
69        Ok(record) => record,
70        Err(source) => {
71            return Err(IcSnapshotTransferReadExecutionError::AfterReplyJournal {
72                source,
73                response: Box::new(response),
74            });
75        }
76    };
77    if let Err(source) = validate_response(&request, association, &response) {
78        return Err(IcSnapshotTransferReadExecutionError::Association {
79            source,
80            response: Box::new(response),
81        });
82    }
83    if let Err(source) = stage.layout() {
84        return Err(IcSnapshotTransferReadExecutionError::AfterReplyStage {
85            source,
86            response: Box::new(response),
87        });
88    }
89    Ok(response)
90}
91
92/// Typed step failures retain original spending; returned raw replies remain available.
93#[derive(Debug, Error)]
94pub enum IcSnapshotTransferReadExecutionError<E: std::error::Error + 'static> {
95    /// The original operation cannot be derived.
96    #[error(transparent)]
97    Plan(#[from] OperationPlanError),
98    /// Retained workflow, stage or ancestors changed before dispatch.
99    #[error(transparent)]
100    Stage(#[from] ExecutionWorkflowPersistenceError),
101    /// Original journal admission or durable reservation failed.
102    #[error(transparent)]
103    Journal(#[from] AttemptJournalError),
104    /// Complete original progress, prerequisites or reservation failed.
105    #[error(transparent)]
106    Progress(#[from] ExecutionProgressPersistenceError),
107    /// Exact original payload/reservation admission failed.
108    #[error(transparent)]
109    Request(#[from] IcSnapshotTransferReadError),
110    /// Integration-owned fresh admission failed after durable reservation.
111    #[error("fresh snapshot read admission failed: {0}")]
112    Admission(#[source] E),
113    /// The single provider call failed; its reservation stays pending.
114    #[error(transparent)]
115    Provider(#[from] IcObservationProviderError),
116    /// Passive association rejected a bounded returned reply.
117    #[error("snapshot read response association failed: {source}")]
118    Association {
119        /// Canonical structural rejection.
120        source: IcSnapshotTransferReadAssociationError,
121        /// Exact bounded returned response, without authentication or outcome.
122        response: Box<IcSnapshotTransferReadResponse>,
123    },
124    /// Journal re-admission failed after a reply; retain the reply.
125    #[error("snapshot read journal changed after reply: {source}")]
126    AfterReplyJournal {
127        /// Original journal rejection.
128        source: AttemptJournalError,
129        /// Exact bounded returned response.
130        response: Box<IcSnapshotTransferReadResponse>,
131    },
132    /// Stage/ancestor re-admission failed after a reply; retain the reply.
133    #[error("snapshot read stage changed after reply: {source}")]
134    AfterReplyStage {
135        /// Original stage or ancestor rejection.
136        source: ExecutionWorkflowPersistenceError,
137        /// Exact bounded returned response.
138        response: Box<IcSnapshotTransferReadResponse>,
139    },
140}
141
142#[cfg(all(test, unix))]
143mod tests;