use std::time::Duration;
use thiserror::Error;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Error)]
pub enum TransactionPolicyError {
#[error("retry limit must be between {min} and {max}")]
InvalidRetryLimit {
min: u32,
max: u32,
value: u32,
},
#[error("transaction timeout must be between {min_ms}ms and {max_ms}ms")]
InvalidTimeout {
min_ms: u64,
max_ms: u64,
value_ms: u64,
},
}
const MIN_RETRY_LIMIT: u32 = 1;
const MAX_RETRY_LIMIT: u32 = 16;
const MIN_TIMEOUT_MS: u64 = 1;
const MAX_TIMEOUT_MS: u64 = 60_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ReadTxnPolicy {
retry_limit: u32,
time_out: Duration,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WriteTxnPolicy {
retry_limit: u32,
time_out: Duration,
}
impl ReadTxnPolicy {
pub fn try_new(retry_limit: u32, time_out: Duration) -> Result<Self, TransactionPolicyError> {
validate_policy_limits(retry_limit, time_out)?;
Ok(Self {
retry_limit,
time_out,
})
}
pub const fn default() -> Self {
Self {
retry_limit: 2,
time_out: Duration::from_millis(500),
}
}
pub const fn retry_limit(&self) -> u32 {
self.retry_limit
}
pub const fn time_out(&self) -> Duration {
self.time_out
}
pub const fn to_transact_option(self) -> foundationdb::TransactOption {
foundationdb::TransactOption {
retry_limit: Some(self.retry_limit),
time_out: Some(self.time_out),
is_idempotent: true,
}
}
}
impl WriteTxnPolicy {
pub fn try_new(retry_limit: u32, time_out: Duration) -> Result<Self, TransactionPolicyError> {
validate_policy_limits(retry_limit, time_out)?;
Ok(Self {
retry_limit,
time_out,
})
}
pub const fn default() -> Self {
Self {
retry_limit: 3,
time_out: Duration::from_secs(5),
}
}
pub const fn retry_limit(&self) -> u32 {
self.retry_limit
}
pub const fn time_out(&self) -> Duration {
self.time_out
}
pub const fn to_transact_option(self) -> foundationdb::TransactOption {
foundationdb::TransactOption {
retry_limit: Some(self.retry_limit),
time_out: Some(self.time_out),
is_idempotent: false,
}
}
}
pub fn idempotent_read_option() -> foundationdb::TransactOption {
ReadTxnPolicy::default().to_transact_option()
}
pub fn mutation_option() -> foundationdb::TransactOption {
WriteTxnPolicy::default().to_transact_option()
}
fn validate_policy_limits(
retry_limit: u32,
time_out: Duration,
) -> Result<(), TransactionPolicyError> {
if !(MIN_RETRY_LIMIT..=MAX_RETRY_LIMIT).contains(&retry_limit) {
return Err(TransactionPolicyError::InvalidRetryLimit {
min: MIN_RETRY_LIMIT,
max: MAX_RETRY_LIMIT,
value: retry_limit,
});
}
let time_out_ms = u64::try_from(time_out.as_millis()).map_err(|_| {
TransactionPolicyError::InvalidTimeout {
min_ms: MIN_TIMEOUT_MS,
max_ms: MAX_TIMEOUT_MS,
value_ms: MAX_TIMEOUT_MS.saturating_add(1),
}
})?;
if !(MIN_TIMEOUT_MS..=MAX_TIMEOUT_MS).contains(&time_out_ms) {
return Err(TransactionPolicyError::InvalidTimeout {
min_ms: MIN_TIMEOUT_MS,
max_ms: MAX_TIMEOUT_MS,
value_ms: time_out_ms,
});
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::{ReadTxnPolicy, TransactionPolicyError, WriteTxnPolicy};
use std::time::Duration;
#[test]
fn read_default_policy_is_fast_and_idempotent() {
let policy = ReadTxnPolicy::default();
let options = policy.to_transact_option();
assert!(options.is_idempotent);
assert_eq!(policy.retry_limit(), 2);
assert_eq!(policy.time_out(), Duration::from_millis(500));
assert_eq!(options.retry_limit, Some(2));
assert_eq!(options.time_out, Some(Duration::from_millis(500)));
}
#[test]
fn write_default_policy_is_retryable_and_non_idempotent() {
let policy = WriteTxnPolicy::default();
let options = policy.to_transact_option();
assert!(!options.is_idempotent);
assert_eq!(policy.retry_limit(), 3);
assert_eq!(policy.time_out(), Duration::from_secs(5));
assert_eq!(options.retry_limit, Some(3));
assert_eq!(options.time_out, Some(Duration::from_secs(5)));
}
#[test]
fn policy_validation_rejects_zero_retries() {
assert!(matches!(
ReadTxnPolicy::try_new(0, Duration::from_millis(10)),
Err(TransactionPolicyError::InvalidRetryLimit { .. })
));
assert!(matches!(
WriteTxnPolicy::try_new(0, Duration::from_millis(10)),
Err(TransactionPolicyError::InvalidRetryLimit { .. })
));
}
#[test]
fn policy_validation_rejects_too_long_timeout() {
assert!(matches!(
ReadTxnPolicy::try_new(1, Duration::from_millis(100_000)),
Err(TransactionPolicyError::InvalidTimeout { .. })
));
}
}