1mod tracker;
10
11use std::fmt::{self, Display, Formatter};
12use std::future::Future;
13
14use r402_protocol::error::FacilitatorError;
15use r402_protocol::payment::SettleResponse;
16pub use tracker::BackgroundSettlementTracker;
17
18use crate::payment_flow::{PaymentFlowName, PaymentFlowPhases};
19
20#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
25pub enum SettlementMode {
26 #[default]
28 Sequential,
29 Concurrent,
31 Background,
33}
34
35impl SettlementMode {
36 #[must_use]
38 pub const fn as_str(self) -> &'static str {
39 match self {
40 Self::Sequential => "sequential",
41 Self::Concurrent => "concurrent",
42 Self::Background => "background",
43 }
44 }
45}
46
47impl Display for SettlementMode {
48 fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
49 f.write_str(self.as_str())
50 }
51}
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
55pub enum AfterHandler {
56 WaitThenSettle,
58 JoinSettle,
60 SpawnSettle,
62 EchoReceipt,
64}
65
66#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
68#[error("incompatible settlement mode {mode} with payment flow {flow}")]
69pub struct IncompatibleSettlementMode {
70 pub mode: SettlementMode,
72 pub flow: PaymentFlowName,
74}
75
76#[derive(Debug, Clone)]
78pub struct SettlementSchedule {
79 mode: SettlementMode,
80 phases: PaymentFlowPhases,
81 after_handler: AfterHandler,
82 attach_receipt: bool,
83 tracker: Option<BackgroundSettlementTracker>,
84}
85
86impl SettlementSchedule {
87 #[must_use]
89 pub const fn after_handler(&self) -> AfterHandler {
90 self.after_handler
91 }
92
93 #[must_use]
95 pub const fn attach_receipt(&self) -> bool {
96 self.attach_receipt
97 }
98
99 #[must_use]
101 pub const fn mode(&self) -> SettlementMode {
102 self.mode
103 }
104
105 #[must_use]
107 pub const fn phases(&self) -> PaymentFlowPhases {
108 self.phases
109 }
110
111 #[must_use]
113 pub fn with_tracker(mut self, tracker: BackgroundSettlementTracker) -> Self {
114 self.tracker = Some(tracker);
115 self
116 }
117}
118
119#[derive(Debug, Clone, PartialEq, Eq)]
124#[allow(
125 clippy::large_enum_variant,
126 reason = "Settled is the sequential success path; Echo is a zero-sized marker"
127)]
128#[must_use]
129pub enum SequentialFinish {
130 Settled(SettleResponse),
132 Echo,
134}
135
136#[derive(Debug)]
138#[must_use]
139pub enum ScheduledSettlement<T, E> {
140 HandlerOkSettleOk {
142 value: T,
144 receipt: Box<SettleResponse>,
146 },
147 HandlerErrDetach {
149 error: E,
151 },
152 SettleErr {
154 value: T,
156 error: FacilitatorError,
158 },
159 Spawned {
161 value: T,
163 },
164}
165
166pub fn schedule(
173 flow: PaymentFlowPhases,
174 mode: SettlementMode,
175) -> Result<SettlementSchedule, IncompatibleSettlementMode> {
176 if mode != SettlementMode::Sequential && flow.settle_before_handler {
177 return Err(IncompatibleSettlementMode {
178 mode,
179 flow: flow_name(flow),
180 });
181 }
182 let after_handler = match mode {
183 SettlementMode::Sequential => {
184 if flow.settle_after_handler {
185 AfterHandler::WaitThenSettle
186 } else {
187 AfterHandler::EchoReceipt
188 }
189 }
190 SettlementMode::Concurrent => AfterHandler::JoinSettle,
191 SettlementMode::Background => AfterHandler::SpawnSettle,
192 };
193 Ok(SettlementSchedule {
194 mode,
195 phases: flow,
196 after_handler,
197 attach_receipt: mode != SettlementMode::Background,
198 tracker: None,
199 })
200}
201
202fn flow_name(phases: PaymentFlowPhases) -> PaymentFlowName {
203 for (name, table) in crate::PAYMENT_FLOWS {
204 if table == phases {
205 return name;
206 }
207 }
208 if phases.settle_before_handler {
209 if phases.settle_after_handler {
210 PaymentFlowName::Escrow
211 } else {
212 PaymentFlowName::Upfront
213 }
214 } else {
215 PaymentFlowName::Authorization
216 }
217}
218
219pub async fn finish<S>(
230 schedule: SettlementSchedule,
231 settle: Option<S>,
232) -> Result<SequentialFinish, FacilitatorError>
233where
234 S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send,
235{
236 match schedule.after_handler {
237 AfterHandler::WaitThenSettle => {
238 let fut = settle.ok_or_else(|| {
239 FacilitatorError::internal("WaitThenSettle requires an after-handler settle")
240 })?;
241 Ok(SequentialFinish::Settled(fut.await?))
242 }
243 AfterHandler::EchoReceipt => {
244 if settle.is_some() {
245 return Err(FacilitatorError::internal(
246 "EchoReceipt does not take an after-handler settle",
247 ));
248 }
249 Ok(SequentialFinish::Echo)
250 }
251 AfterHandler::JoinSettle | AfterHandler::SpawnSettle => {
252 Err(FacilitatorError::internal("finish is sequential-only"))
253 }
254 }
255}
256
257pub async fn run<T, E, H, S>(
261 schedule: SettlementSchedule,
262 handler: H,
263 settle: S,
264) -> ScheduledSettlement<T, E>
265where
266 T: Send,
267 E: Send,
268 H: Future<Output = Result<T, E>> + Send,
269 S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
270{
271 match schedule.after_handler {
272 AfterHandler::JoinSettle => join_settle(handler, settle).await,
273 AfterHandler::SpawnSettle => spawn_settle(schedule.tracker, handler, settle).await,
274 AfterHandler::WaitThenSettle | AfterHandler::EchoReceipt => {
275 drop(settle);
276 match handler.await {
277 Ok(value) => ScheduledSettlement::SettleErr {
278 value,
279 error: FacilitatorError::internal("run is concurrent/background"),
280 },
281 Err(error) => ScheduledSettlement::HandlerErrDetach { error },
282 }
283 }
284 }
285}
286
287async fn join_settle<T, E, H, S>(handler: H, settle: S) -> ScheduledSettlement<T, E>
288where
289 H: Future<Output = Result<T, E>> + Send,
290 S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
291{
292 let settle_handle = tokio::spawn(settle);
293 match handler.await {
294 Ok(value) => match settle_handle.await {
295 Ok(Ok(receipt)) => ScheduledSettlement::HandlerOkSettleOk {
296 value,
297 receipt: Box::new(receipt),
298 },
299 Ok(Err(error)) => ScheduledSettlement::SettleErr { value, error },
300 Err(join) => ScheduledSettlement::SettleErr {
301 value,
302 error: FacilitatorError::internal(join),
303 },
304 },
305 Err(error) => {
306 drop(settle_handle);
307 ScheduledSettlement::HandlerErrDetach { error }
308 }
309 }
310}
311
312async fn spawn_settle<T, E, H, S>(
313 tracker: Option<BackgroundSettlementTracker>,
314 handler: H,
315 settle: S,
316) -> ScheduledSettlement<T, E>
317where
318 H: Future<Output = Result<T, E>> + Send,
319 S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
320{
321 let settle_handle = tokio::spawn(settle);
322 let tracker_guard = tracker.as_ref().map(BackgroundSettlementTracker::start);
323 drop(tokio::spawn(supervise_background_settle(
324 settle_handle,
325 tracker_guard,
326 )));
327 match handler.await {
328 Ok(value) => ScheduledSettlement::Spawned { value },
329 Err(error) => ScheduledSettlement::HandlerErrDetach { error },
330 }
331}
332
333async fn supervise_background_settle(
334 handle: tokio::task::JoinHandle<Result<SettleResponse, FacilitatorError>>,
335 _tracker: Option<tracker::SettlementInFlightGuard>,
336) {
337 let outcome = handle.await;
338 log_background_settle_outcome(&outcome);
339 record_background_settle_metric(&outcome);
340}
341
342fn log_background_settle_outcome(
343 outcome: &Result<Result<SettleResponse, FacilitatorError>, tokio::task::JoinError>,
344) {
345 match outcome {
346 Ok(Ok(_)) => log_background_settle_ok(),
347 Ok(Err(err)) => log_background_settle_facilitator_err(err),
348 Err(join_err) => log_background_settle_join_err(join_err),
349 }
350}
351
352fn log_background_settle_ok() {
353 #[cfg(feature = "telemetry")]
354 tracing::debug!("background settlement completed");
355}
356
357fn log_background_settle_facilitator_err(err: &FacilitatorError) {
358 #[cfg(feature = "telemetry")]
359 tracing::error!(error = %err, "background settlement returned error");
360 #[cfg(not(feature = "telemetry"))]
361 let _ = err;
362}
363
364fn log_background_settle_join_err(join_err: &tokio::task::JoinError) {
365 #[cfg(feature = "telemetry")]
366 if join_err.is_panic() {
367 tracing::error!(error = %join_err, "background settlement task panicked");
368 } else {
369 tracing::warn!(error = %join_err, "background settlement task cancelled");
370 }
371 #[cfg(not(feature = "telemetry"))]
372 let _ = join_err;
373}
374
375fn record_background_settle_metric(
376 outcome: &Result<Result<SettleResponse, FacilitatorError>, tokio::task::JoinError>,
377) {
378 #[cfg(feature = "metrics")]
379 {
380 let result = match outcome {
381 Ok(Ok(_)) => "ok",
382 Ok(Err(_)) => "error",
383 Err(join_err) if join_err.is_panic() => "panic",
384 Err(_) => "cancelled",
385 };
386 ::metrics::counter!(
387 r402_protocol::metrics::BACKGROUND_SETTLE_TOTAL,
388 "result" => result,
389 )
390 .increment(1);
391 }
392 #[cfg(not(feature = "metrics"))]
393 let _ = outcome;
394}