mod tracker;
use std::fmt::{self, Display, Formatter};
use std::future::Future;
use r402_protocol::error::FacilitatorError;
use r402_protocol::payment::SettleResponse;
pub use tracker::BackgroundSettlementTracker;
use crate::payment_flow::{PaymentFlowName, PaymentFlowPhases};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Hash)]
pub enum SettlementMode {
#[default]
Sequential,
Concurrent,
Background,
}
impl SettlementMode {
#[must_use]
pub const fn as_str(self) -> &'static str {
match self {
Self::Sequential => "sequential",
Self::Concurrent => "concurrent",
Self::Background => "background",
}
}
}
impl Display for SettlementMode {
fn fmt(&self, f: &mut Formatter<'_>) -> fmt::Result {
f.write_str(self.as_str())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum AfterHandler {
WaitThenSettle,
JoinSettle,
SpawnSettle,
EchoReceipt,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, thiserror::Error)]
#[error("incompatible settlement mode {mode} with payment flow {flow}")]
pub struct IncompatibleSettlementMode {
pub mode: SettlementMode,
pub flow: PaymentFlowName,
}
#[derive(Debug, Clone)]
pub struct SettlementSchedule {
mode: SettlementMode,
phases: PaymentFlowPhases,
after_handler: AfterHandler,
attach_receipt: bool,
tracker: Option<BackgroundSettlementTracker>,
}
impl SettlementSchedule {
#[must_use]
pub const fn after_handler(&self) -> AfterHandler {
self.after_handler
}
#[must_use]
pub const fn attach_receipt(&self) -> bool {
self.attach_receipt
}
#[must_use]
pub const fn mode(&self) -> SettlementMode {
self.mode
}
#[must_use]
pub const fn phases(&self) -> PaymentFlowPhases {
self.phases
}
#[must_use]
pub fn with_tracker(mut self, tracker: BackgroundSettlementTracker) -> Self {
self.tracker = Some(tracker);
self
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[allow(
clippy::large_enum_variant,
reason = "Settled is the sequential success path; Echo is a zero-sized marker"
)]
#[must_use]
pub enum SequentialFinish {
Settled(SettleResponse),
Echo,
}
#[derive(Debug)]
#[must_use]
pub enum ScheduledSettlement<T, E> {
HandlerOkSettleOk {
value: T,
receipt: Box<SettleResponse>,
},
HandlerErrDetach {
error: E,
},
SettleErr {
value: T,
error: FacilitatorError,
},
Spawned {
value: T,
},
}
pub fn schedule(
flow: PaymentFlowPhases,
mode: SettlementMode,
) -> Result<SettlementSchedule, IncompatibleSettlementMode> {
if mode != SettlementMode::Sequential && flow.settle_before_handler {
return Err(IncompatibleSettlementMode {
mode,
flow: flow_name(flow),
});
}
let after_handler = match mode {
SettlementMode::Sequential => {
if flow.settle_after_handler {
AfterHandler::WaitThenSettle
} else {
AfterHandler::EchoReceipt
}
}
SettlementMode::Concurrent => AfterHandler::JoinSettle,
SettlementMode::Background => AfterHandler::SpawnSettle,
};
Ok(SettlementSchedule {
mode,
phases: flow,
after_handler,
attach_receipt: mode != SettlementMode::Background,
tracker: None,
})
}
fn flow_name(phases: PaymentFlowPhases) -> PaymentFlowName {
for (name, table) in crate::PAYMENT_FLOWS {
if table == phases {
return name;
}
}
if phases.settle_before_handler {
if phases.settle_after_handler {
PaymentFlowName::Escrow
} else {
PaymentFlowName::Upfront
}
} else {
PaymentFlowName::Authorization
}
}
pub async fn finish<S>(
schedule: SettlementSchedule,
settle: Option<S>,
) -> Result<SequentialFinish, FacilitatorError>
where
S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send,
{
match schedule.after_handler {
AfterHandler::WaitThenSettle => {
let fut = settle.ok_or_else(|| {
FacilitatorError::internal("WaitThenSettle requires an after-handler settle")
})?;
Ok(SequentialFinish::Settled(fut.await?))
}
AfterHandler::EchoReceipt => {
if settle.is_some() {
return Err(FacilitatorError::internal(
"EchoReceipt does not take an after-handler settle",
));
}
Ok(SequentialFinish::Echo)
}
AfterHandler::JoinSettle | AfterHandler::SpawnSettle => {
Err(FacilitatorError::internal("finish is sequential-only"))
}
}
}
pub async fn run<T, E, H, S>(
schedule: SettlementSchedule,
handler: H,
settle: S,
) -> ScheduledSettlement<T, E>
where
T: Send,
E: Send,
H: Future<Output = Result<T, E>> + Send,
S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
{
match schedule.after_handler {
AfterHandler::JoinSettle => join_settle(handler, settle).await,
AfterHandler::SpawnSettle => spawn_settle(schedule.tracker, handler, settle).await,
AfterHandler::WaitThenSettle | AfterHandler::EchoReceipt => {
drop(settle);
match handler.await {
Ok(value) => ScheduledSettlement::SettleErr {
value,
error: FacilitatorError::internal("run is concurrent/background"),
},
Err(error) => ScheduledSettlement::HandlerErrDetach { error },
}
}
}
}
async fn join_settle<T, E, H, S>(handler: H, settle: S) -> ScheduledSettlement<T, E>
where
H: Future<Output = Result<T, E>> + Send,
S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
{
let settle_handle = tokio::spawn(settle);
match handler.await {
Ok(value) => match settle_handle.await {
Ok(Ok(receipt)) => ScheduledSettlement::HandlerOkSettleOk {
value,
receipt: Box::new(receipt),
},
Ok(Err(error)) => ScheduledSettlement::SettleErr { value, error },
Err(join) => ScheduledSettlement::SettleErr {
value,
error: FacilitatorError::internal(join),
},
},
Err(error) => {
drop(settle_handle);
ScheduledSettlement::HandlerErrDetach { error }
}
}
}
async fn spawn_settle<T, E, H, S>(
tracker: Option<BackgroundSettlementTracker>,
handler: H,
settle: S,
) -> ScheduledSettlement<T, E>
where
H: Future<Output = Result<T, E>> + Send,
S: Future<Output = Result<SettleResponse, FacilitatorError>> + Send + 'static,
{
let settle_handle = tokio::spawn(settle);
let tracker_guard = tracker.as_ref().map(BackgroundSettlementTracker::start);
drop(tokio::spawn(supervise_background_settle(
settle_handle,
tracker_guard,
)));
match handler.await {
Ok(value) => ScheduledSettlement::Spawned { value },
Err(error) => ScheduledSettlement::HandlerErrDetach { error },
}
}
async fn supervise_background_settle(
handle: tokio::task::JoinHandle<Result<SettleResponse, FacilitatorError>>,
_tracker: Option<tracker::SettlementInFlightGuard>,
) {
let outcome = handle.await;
log_background_settle_outcome(&outcome);
record_background_settle_metric(&outcome);
}
fn log_background_settle_outcome(
outcome: &Result<Result<SettleResponse, FacilitatorError>, tokio::task::JoinError>,
) {
match outcome {
Ok(Ok(_)) => log_background_settle_ok(),
Ok(Err(err)) => log_background_settle_facilitator_err(err),
Err(join_err) => log_background_settle_join_err(join_err),
}
}
fn log_background_settle_ok() {
#[cfg(feature = "telemetry")]
tracing::debug!("background settlement completed");
}
fn log_background_settle_facilitator_err(err: &FacilitatorError) {
#[cfg(feature = "telemetry")]
tracing::error!(error = %err, "background settlement returned error");
#[cfg(not(feature = "telemetry"))]
let _ = err;
}
fn log_background_settle_join_err(join_err: &tokio::task::JoinError) {
#[cfg(feature = "telemetry")]
if join_err.is_panic() {
tracing::error!(error = %join_err, "background settlement task panicked");
} else {
tracing::warn!(error = %join_err, "background settlement task cancelled");
}
#[cfg(not(feature = "telemetry"))]
let _ = join_err;
}
fn record_background_settle_metric(
outcome: &Result<Result<SettleResponse, FacilitatorError>, tokio::task::JoinError>,
) {
#[cfg(feature = "metrics")]
{
let result = match outcome {
Ok(Ok(_)) => "ok",
Ok(Err(_)) => "error",
Err(join_err) if join_err.is_panic() => "panic",
Err(_) => "cancelled",
};
::metrics::counter!(
r402_protocol::metrics::BACKGROUND_SETTLE_TOTAL,
"result" => result,
)
.increment(1);
}
#[cfg(not(feature = "metrics"))]
let _ = outcome;
}