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 fn read_snapshot<E: std::error::Error + 'static>(
48 stage: &ExecutionStageGuard<'_>,
49 operation_sequence: u64,
50 payload: IcSnapshotTransferReadPayload<'_, '_>,
51 provider: &mut impl IcSnapshotTransferReadProvider,
52 admit: impl FnOnce(&IcSnapshotTransferReadRequest<'_, '_>) -> Result<(), E>,
53) -> Result<IcSnapshotTransferReadResponse, IcSnapshotTransferReadExecutionError<E>> {
54 let plan = stage.plan();
55 let authority = plan.attempt_authority(operation_sequence)?;
56 payload.validate_binding(&authority)?;
57 let mut journal = AttemptJournalGuard::open(stage.layout()?, &authority)?;
58 journal.reserve_planned_mutation(&plan.digest())?;
59 let request =
60 IcSnapshotTransferReadRequest::new(plan, operation_sequence, journal.record()?, payload)?;
61 admit(&request).map_err(IcSnapshotTransferReadExecutionError::Admission)?;
62 stage.layout()?;
63 let response = provider.read_snapshot(&request)?;
64 let association = match journal.record() {
65 Ok(record) => record,
66 Err(source) => {
67 return Err(IcSnapshotTransferReadExecutionError::AfterReplyJournal {
68 source,
69 response: Box::new(response),
70 });
71 }
72 };
73 if let Err(source) = validate_response(&request, association, &response) {
74 return Err(IcSnapshotTransferReadExecutionError::Association {
75 source,
76 response: Box::new(response),
77 });
78 }
79 if let Err(source) = stage.layout() {
80 return Err(IcSnapshotTransferReadExecutionError::AfterReplyStage {
81 source,
82 response: Box::new(response),
83 });
84 }
85 Ok(response)
86}
87
88#[derive(Debug, Error)]
90pub enum IcSnapshotTransferReadExecutionError<E: std::error::Error + 'static> {
91 #[error(transparent)]
93 Plan(#[from] OperationPlanError),
94 #[error(transparent)]
96 Stage(#[from] ExecutionWorkflowPersistenceError),
97 #[error(transparent)]
99 Journal(#[from] AttemptJournalError),
100 #[error(transparent)]
102 Progress(#[from] ExecutionProgressPersistenceError),
103 #[error(transparent)]
105 Request(#[from] IcSnapshotTransferReadError),
106 #[error("fresh snapshot read admission failed: {0}")]
108 Admission(#[source] E),
109 #[error(transparent)]
111 Provider(#[from] IcObservationProviderError),
112 #[error("snapshot read response association failed: {source}")]
114 Association {
115 source: IcSnapshotTransferReadAssociationError,
117 response: Box<IcSnapshotTransferReadResponse>,
119 },
120 #[error("snapshot read journal changed after reply: {source}")]
122 AfterReplyJournal {
123 source: AttemptJournalError,
125 response: Box<IcSnapshotTransferReadResponse>,
127 },
128 #[error("snapshot read stage changed after reply: {source}")]
130 AfterReplyStage {
131 source: ExecutionWorkflowPersistenceError,
133 response: Box<IcSnapshotTransferReadResponse>,
135 },
136}
137
138#[cfg(all(test, unix))]
139mod tests;