ic_backup/workflow/ic_snapshot_transfer_read/
mod.rs1use 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
25pub 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#[derive(Debug, Error)]
94pub enum IcSnapshotTransferReadExecutionError<E: std::error::Error + 'static> {
95 #[error(transparent)]
97 Plan(#[from] OperationPlanError),
98 #[error(transparent)]
100 Stage(#[from] ExecutionWorkflowPersistenceError),
101 #[error(transparent)]
103 Journal(#[from] AttemptJournalError),
104 #[error(transparent)]
106 Progress(#[from] ExecutionProgressPersistenceError),
107 #[error(transparent)]
109 Request(#[from] IcSnapshotTransferReadError),
110 #[error("fresh snapshot read admission failed: {0}")]
112 Admission(#[source] E),
113 #[error(transparent)]
115 Provider(#[from] IcObservationProviderError),
116 #[error("snapshot read response association failed: {source}")]
118 Association {
119 source: IcSnapshotTransferReadAssociationError,
121 response: Box<IcSnapshotTransferReadResponse>,
123 },
124 #[error("snapshot read journal changed after reply: {source}")]
126 AfterReplyJournal {
127 source: AttemptJournalError,
129 response: Box<IcSnapshotTransferReadResponse>,
131 },
132 #[error("snapshot read stage changed after reply: {source}")]
134 AfterReplyStage {
135 source: ExecutionWorkflowPersistenceError,
137 response: Box<IcSnapshotTransferReadResponse>,
139 },
140}
141
142#[cfg(all(test, unix))]
143mod tests;