use thiserror::Error;
use crate::hook::{SettlementFailureReason, SettlementSkipReason};
use crate::retry::RetryPolicy;
pub const MAX_SETTLEMENT_WORKER_ID_BYTES: usize = 128;
pub const MAX_SETTLEMENT_LEASE_MS: u64 = 86_400_000;
pub const MAX_SETTLEMENT_CLAIM_BATCH: usize = 1024;
#[derive(Debug, Error, Clone, Copy, PartialEq, Eq)]
pub enum SettlementClaimValidationError {
#[error("settlement worker id length is out of bounds")]
WorkerIdLength,
#[error("settlement lease duration is out of bounds")]
LeaseDuration,
#[error("settlement claim batch size is out of bounds")]
ClaimBatch,
}
pub fn validate_settlement_claim(
worker_id: &str,
lease_ms: u64,
limit: usize,
) -> Result<(), SettlementClaimValidationError> {
if worker_id.is_empty() || worker_id.len() > MAX_SETTLEMENT_WORKER_ID_BYTES {
return Err(SettlementClaimValidationError::WorkerIdLength);
}
if lease_ms == 0 || lease_ms > MAX_SETTLEMENT_LEASE_MS {
return Err(SettlementClaimValidationError::LeaseDuration);
}
if limit == 0 || limit > MAX_SETTLEMENT_CLAIM_BATCH {
return Err(SettlementClaimValidationError::ClaimBatch);
}
Ok(())
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct SettlementStoreBinding([u8; 32]);
impl SettlementStoreBinding {
#[must_use]
pub const fn from_digest(digest: [u8; 32]) -> Self {
Self(digest)
}
#[must_use]
pub const fn as_bytes(&self) -> &[u8; 32] {
&self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum SettlementRoutingInput {
Accepted,
Skipped { reason: SettlementSkipReason },
Retryable { reason: SettlementFailureReason },
Permanent { reason: SettlementFailureReason },
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SettlementAttemptClaim {
pub receipt_id: String,
pub finalized_at: u64,
pub attempts: u32,
pub row_version: u64,
pub lease_owner: String,
pub lease_token: String,
pub lease_until_ms: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SettlementRoute {
NoAction,
RetryScheduled {
attempt: u32,
next_visible_at_ms: u64,
},
DeadLettered {
attempts: u32,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SettlementRouteErrorClass {
Backend,
Conflict,
InvalidRecord,
}
#[derive(Debug, Error, Clone, PartialEq, Eq)]
pub enum SettlementRouteError {
#[error("settlement routing backend failure: {detail}")]
Backend { detail: String },
#[error("settlement routing conflict: {detail}")]
Conflict { detail: String },
#[error("invalid settlement routing record: {detail}")]
InvalidRecord { detail: String },
}
impl SettlementRouteError {
#[must_use]
pub const fn class(&self) -> SettlementRouteErrorClass {
match self {
Self::Backend { .. } => SettlementRouteErrorClass::Backend,
Self::Conflict { .. } => SettlementRouteErrorClass::Conflict,
Self::InvalidRecord { .. } => SettlementRouteErrorClass::InvalidRecord,
}
}
}
pub trait SettlementOutcomeStore: Send + Sync {
fn settlement_store_binding(&self) -> SettlementStoreBinding;
fn claim_receipt(
&self,
receipt_id: &str,
worker_id: &str,
now_ms: u64,
lease_ms: u64,
) -> Result<Option<SettlementAttemptClaim>, SettlementRouteError>;
fn claim_due(
&self,
worker_id: &str,
now_ms: u64,
lease_ms: u64,
limit: usize,
) -> Result<Vec<SettlementAttemptClaim>, SettlementRouteError>;
fn record_claimed_outcome(
&self,
claim: &SettlementAttemptClaim,
outcome: &SettlementRoutingInput,
policy: RetryPolicy,
observed_at_ms: u64,
) -> Result<SettlementRoute, SettlementRouteError>;
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn attempt_claim_carries_the_persisted_cas_state() {
let claim = SettlementAttemptClaim {
receipt_id: "receipt-1".to_string(),
finalized_at: 41,
attempts: 2,
row_version: 3,
lease_owner: "worker-1".to_string(),
lease_token: "token-1".to_string(),
lease_until_ms: 50,
};
assert_eq!(claim.finalized_at, 41);
assert_eq!(claim.attempts, 2);
assert_eq!(claim.lease_owner, "worker-1");
}
#[test]
fn claim_parameters_are_bounded() {
assert!(validate_settlement_claim("worker", 1, 1).is_ok());
assert!(validate_settlement_claim(
&"w".repeat(MAX_SETTLEMENT_WORKER_ID_BYTES),
MAX_SETTLEMENT_LEASE_MS,
MAX_SETTLEMENT_CLAIM_BATCH,
)
.is_ok());
assert!(validate_settlement_claim(&"\u{e9}".repeat(64), 1, 1).is_ok());
let invalid = [
(
validate_settlement_claim("", 1, 1),
SettlementClaimValidationError::WorkerIdLength,
),
(
validate_settlement_claim(&"w".repeat(MAX_SETTLEMENT_WORKER_ID_BYTES + 1), 1, 1),
SettlementClaimValidationError::WorkerIdLength,
),
(
validate_settlement_claim(&"\u{e9}".repeat(65), 1, 1),
SettlementClaimValidationError::WorkerIdLength,
),
(
validate_settlement_claim("worker", 0, 1),
SettlementClaimValidationError::LeaseDuration,
),
(
validate_settlement_claim("worker", MAX_SETTLEMENT_LEASE_MS + 1, 1),
SettlementClaimValidationError::LeaseDuration,
),
(
validate_settlement_claim("worker", 1, 0),
SettlementClaimValidationError::ClaimBatch,
),
(
validate_settlement_claim("worker", 1, MAX_SETTLEMENT_CLAIM_BATCH + 1),
SettlementClaimValidationError::ClaimBatch,
),
];
for (result, expected) in invalid {
assert_eq!(result, Err(expected));
}
}
#[test]
fn outcome_store_error_classes_are_bounded() {
let cases = [
(
SettlementRouteError::Backend {
detail: "db unavailable".to_string(),
},
SettlementRouteErrorClass::Backend,
),
(
SettlementRouteError::Conflict {
detail: "stale lease".to_string(),
},
SettlementRouteErrorClass::Conflict,
),
(
SettlementRouteError::InvalidRecord {
detail: "negative counter".to_string(),
},
SettlementRouteErrorClass::InvalidRecord,
),
];
for (error, expected) in cases {
assert_eq!(error.class(), expected);
}
}
#[test]
fn outcome_store_trait_is_object_safe() {
fn accepts_trait_object(_store: &dyn SettlementOutcomeStore) {}
let _ = accepts_trait_object;
}
}