made_app/workers/
recoverable_ceremony_worker.rs1use 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
14pub 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}