Skip to main content

made_app/workers/
recoverable_ceremony_worker.rs

1use std::sync::Arc;
2
3use async_trait::async_trait;
4use made_core::error::DomainError;
5use made_core::value_objects::ExecutionReceiptLinkKind;
6
7use super::{
8    CompleteExecutionReceiptInput, CompleteExecutionReceiptUseCase, ExecuteCeremonyOperationInput,
9    ExecuteCeremonyOperationOutcome, ExecuteCeremonyOperationUseCase, ExecutionRecoveryItem,
10    RecoverExecutionIntentOutcome, RecoverExecutionIntentUseCase, RecoverableCeremonyWorkerOutcome,
11    RecoverableCeremonyWorkerPort,
12};
13
14/// Runs an accepted claim through receipt persistence and fenced completion.
15pub struct RecoverableCeremonyWorker {
16    execute_operation: Arc<ExecuteCeremonyOperationUseCase>,
17    recover_intent: Arc<RecoverExecutionIntentUseCase>,
18    complete_receipt: Arc<CompleteExecutionReceiptUseCase>,
19}
20
21impl std::fmt::Debug for RecoverableCeremonyWorker {
22    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
23        formatter
24            .debug_struct("RecoverableCeremonyWorker")
25            .finish_non_exhaustive()
26    }
27}
28
29impl RecoverableCeremonyWorker {
30    #[must_use]
31    pub const fn new(
32        execute_operation: Arc<ExecuteCeremonyOperationUseCase>,
33        recover_intent: Arc<RecoverExecutionIntentUseCase>,
34        complete_receipt: Arc<CompleteExecutionReceiptUseCase>,
35    ) -> Self {
36        Self {
37            execute_operation,
38            recover_intent,
39            complete_receipt,
40        }
41    }
42
43    pub async fn execute_claim(
44        &self,
45        input: ExecuteCeremonyOperationInput,
46    ) -> Result<RecoverableCeremonyWorkerOutcome, DomainError> {
47        let ceremony_id = input.handler_request.instance_id().clone();
48        let step_id = input.handler_request.step_id().clone();
49        let claim_fence = input.claim_fence.clone();
50        let actor_kind = input.actor_kind;
51        let receipt = match self.execute_operation.execute(input).await? {
52            ExecuteCeremonyOperationOutcome::Receipt(receipt) => *receipt,
53            ExecuteCeremonyOperationOutcome::ReconciliationRequired(operation_id) => {
54                return Ok(RecoverableCeremonyWorkerOutcome::ReconciliationRequired(
55                    operation_id,
56                ));
57            }
58        };
59        let instance = self
60            .complete_receipt
61            .execute(CompleteExecutionReceiptInput {
62                ceremony_id,
63                step_id,
64                operation_id: receipt.operation_id().clone(),
65                claim_fence,
66                link_kind: ExecutionReceiptLinkKind::Direct,
67                actor_kind,
68            })
69            .await?;
70        Ok(RecoverableCeremonyWorkerOutcome::Completed {
71            receipt: Box::new(receipt),
72            instance: Box::new(instance),
73        })
74    }
75
76    pub async fn recover(
77        &self,
78        item: ExecutionRecoveryItem,
79    ) -> Result<RecoverableCeremonyWorkerOutcome, DomainError> {
80        let (operation, intents, receipt, current_claim_fence) = item.into_parts();
81        let intent = current_claim_fence
82            .as_ref()
83            .and_then(|fence| intents.iter().find(|intent| intent.claim_fence() == fence))
84            .or_else(|| intents.iter().min_by_key(|intent| intent.recorded_at()))
85            .ok_or(DomainError::InvariantViolated {
86                reason: "execution recovery item has no durable intent",
87            })?;
88        let receipt = match receipt {
89            Some(receipt) => receipt,
90            None => match self.recover_intent.execute(intent).await? {
91                RecoverExecutionIntentOutcome::Receipt(receipt) => *receipt,
92                RecoverExecutionIntentOutcome::ReconciliationRequired(operation_id) => {
93                    return Ok(RecoverableCeremonyWorkerOutcome::ReconciliationRequired(
94                        operation_id,
95                    ));
96                }
97            },
98        };
99        let applied_fence =
100            current_claim_fence.unwrap_or_else(|| receipt.producer_claim_fence().clone());
101        let link_kind = if &applied_fence == receipt.producer_claim_fence() {
102            ExecutionReceiptLinkKind::Direct
103        } else {
104            ExecutionReceiptLinkKind::Adopted
105        };
106        let actor_kind = intents
107            .iter()
108            .find(|candidate| candidate.claim_fence() == &applied_fence)
109            .unwrap_or(intent)
110            .actor_kind();
111        let instance = self
112            .complete_receipt
113            .execute(CompleteExecutionReceiptInput {
114                ceremony_id: operation.ceremony_id().clone(),
115                step_id: operation.step_id().clone(),
116                operation_id: operation.operation_id().clone(),
117                claim_fence: applied_fence,
118                link_kind,
119                actor_kind,
120            })
121            .await?;
122        Ok(RecoverableCeremonyWorkerOutcome::Completed {
123            receipt: Box::new(receipt),
124            instance: Box::new(instance),
125        })
126    }
127}
128
129#[async_trait]
130impl RecoverableCeremonyWorkerPort for RecoverableCeremonyWorker {
131    async fn execute_claim(
132        &self,
133        input: ExecuteCeremonyOperationInput,
134    ) -> Result<RecoverableCeremonyWorkerOutcome, DomainError> {
135        RecoverableCeremonyWorker::execute_claim(self, input).await
136    }
137
138    async fn recover(
139        &self,
140        item: ExecutionRecoveryItem,
141    ) -> Result<RecoverableCeremonyWorkerOutcome, DomainError> {
142        RecoverableCeremonyWorker::recover(self, item).await
143    }
144}