use std::time::Duration;
use super::env_parsing::{
parse_duration_millis_from_env, parse_from_env, parse_optional_bool_from_env, ValidationBounds,
};
#[non_exhaustive]
#[derive(Clone, Debug)]
pub struct PartitionFailoverOptions {
circuit_breaker_enabled: bool,
circuit_breaker_enabled_override: Option<bool>,
read_failure_threshold: u32,
write_failure_threshold: u32,
counter_reset_window: Duration,
partition_unavailability_duration: Duration,
failback_sweep_interval: Duration,
consecutive_hedge_win_threshold: u32,
}
impl Default for PartitionFailoverOptions {
fn default() -> Self {
Self {
circuit_breaker_enabled: true, circuit_breaker_enabled_override: None,
read_failure_threshold: 10,
write_failure_threshold: 5,
counter_reset_window: Duration::from_millis(300_000),
partition_unavailability_duration: Duration::from_millis(5_000),
failback_sweep_interval: Duration::from_millis(300_000),
consecutive_hedge_win_threshold: 5,
}
}
}
impl PartitionFailoverOptions {
pub fn builder() -> PartitionFailoverOptionsBuilder {
PartitionFailoverOptionsBuilder::new()
}
pub fn circuit_breaker_enabled(&self) -> bool {
self.circuit_breaker_enabled
}
pub(crate) fn circuit_breaker_enabled_override(&self) -> Option<bool> {
self.circuit_breaker_enabled_override
}
pub fn read_failure_threshold(&self) -> u32 {
self.read_failure_threshold
}
pub fn write_failure_threshold(&self) -> u32 {
self.write_failure_threshold
}
pub fn counter_reset_window(&self) -> Duration {
self.counter_reset_window
}
pub fn partition_unavailability_duration(&self) -> Duration {
self.partition_unavailability_duration
}
pub fn failback_sweep_interval(&self) -> Duration {
self.failback_sweep_interval
}
pub fn consecutive_hedge_win_threshold(&self) -> u32 {
self.consecutive_hedge_win_threshold
}
}
#[non_exhaustive]
#[derive(Clone, Debug, Default)]
pub struct PartitionFailoverOptionsBuilder {
circuit_breaker_enabled: Option<bool>,
circuit_breaker_enabled_override: Option<bool>,
read_failure_threshold: Option<u32>,
write_failure_threshold: Option<u32>,
counter_reset_window: Option<Duration>,
partition_unavailability_duration: Option<Duration>,
failback_sweep_interval: Option<Duration>,
consecutive_hedge_win_threshold: Option<u32>,
}
impl PartitionFailoverOptionsBuilder {
pub fn new() -> Self {
Self::default()
}
pub fn with_circuit_breaker_enabled(mut self, value: bool) -> Self {
self.circuit_breaker_enabled = Some(value);
self
}
#[cfg(test)]
pub(crate) fn with_circuit_breaker_enabled_override(mut self, value: bool) -> Self {
self.circuit_breaker_enabled_override = Some(value);
self
}
pub fn with_read_failure_threshold(mut self, value: u32) -> Self {
self.read_failure_threshold = Some(value);
self
}
pub fn with_write_failure_threshold(mut self, value: u32) -> Self {
self.write_failure_threshold = Some(value);
self
}
pub fn with_counter_reset_window(mut self, value: Duration) -> Self {
self.counter_reset_window = Some(value);
self
}
pub fn with_partition_unavailability_duration(mut self, value: Duration) -> Self {
self.partition_unavailability_duration = Some(value);
self
}
pub fn with_failback_sweep_interval(mut self, value: Duration) -> Self {
self.failback_sweep_interval = Some(value);
self
}
pub fn with_consecutive_hedge_win_threshold(mut self, value: u32) -> Self {
self.consecutive_hedge_win_threshold = Some(value);
self
}
pub fn build(self) -> crate::error::Result<PartitionFailoverOptions> {
self.build_from_env(&|k| std::env::var(k).ok())
}
pub(crate) fn build_from_env(
self,
get_env: &dyn Fn(&str) -> Option<String>,
) -> crate::error::Result<PartitionFailoverOptions> {
let defaults = PartitionFailoverOptions::default();
let circuit_breaker_enabled = parse_from_env(
self.circuit_breaker_enabled,
"AZURE_COSMOS_PPCB_ENABLED",
defaults.circuit_breaker_enabled,
ValidationBounds::none(),
get_env,
)?;
let circuit_breaker_enabled_override = parse_optional_bool_from_env(
self.circuit_breaker_enabled_override,
"AZURE_COSMOS_PPCB_ENABLED_OVERRIDE",
get_env,
);
let read_failure_threshold = parse_from_env(
self.read_failure_threshold,
"AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD",
defaults.read_failure_threshold,
ValidationBounds::min(1),
get_env,
)?;
let write_failure_threshold = parse_from_env(
self.write_failure_threshold,
"AZURE_COSMOS_PPCB_WRITE_FAILURE_THRESHOLD",
defaults.write_failure_threshold,
ValidationBounds::min(1),
get_env,
)?;
let counter_reset_window = parse_duration_millis_from_env(
self.counter_reset_window,
"AZURE_COSMOS_PPCB_COUNTER_RESET_WINDOW_MS",
defaults.counter_reset_window.as_millis() as u64,
1_000,
u64::MAX,
get_env,
)?;
let partition_unavailability_duration = parse_duration_millis_from_env(
self.partition_unavailability_duration,
"AZURE_COSMOS_PPCB_PARTITION_UNAVAILABILITY_DURATION_MS",
defaults.partition_unavailability_duration.as_millis() as u64,
1_000,
u64::MAX,
get_env,
)?;
let failback_sweep_interval = parse_duration_millis_from_env(
self.failback_sweep_interval,
"AZURE_COSMOS_PPCB_FAILBACK_SWEEP_INTERVAL_MS",
defaults.failback_sweep_interval.as_millis() as u64,
1_000,
u64::MAX,
get_env,
)?;
let consecutive_hedge_win_threshold = parse_from_env(
self.consecutive_hedge_win_threshold,
"AZURE_COSMOS_PPCB_CONSECUTIVE_HEDGE_WIN_THRESHOLD",
defaults.consecutive_hedge_win_threshold,
ValidationBounds::min(1),
get_env,
)?;
Ok(PartitionFailoverOptions {
circuit_breaker_enabled,
circuit_breaker_enabled_override,
read_failure_threshold,
write_failure_threshold,
counter_reset_window,
partition_unavailability_duration,
failback_sweep_interval,
consecutive_hedge_win_threshold,
})
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builder_defaults_match_documented_values() {
let options = PartitionFailoverOptionsBuilder::new()
.build_from_env(&|_| None)
.unwrap();
assert!(options.circuit_breaker_enabled());
assert_eq!(options.read_failure_threshold(), 10);
assert_eq!(options.write_failure_threshold(), 5);
assert_eq!(options.counter_reset_window(), Duration::from_secs(5 * 60));
assert_eq!(
options.partition_unavailability_duration(),
Duration::from_secs(5)
);
assert_eq!(options.failback_sweep_interval(), Duration::from_secs(300));
assert_eq!(options.consecutive_hedge_win_threshold(), 5);
}
#[test]
fn builder_round_trips_custom_values() {
let options = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled(true)
.with_read_failure_threshold(20)
.with_write_failure_threshold(7)
.with_counter_reset_window(Duration::from_secs(60))
.with_partition_unavailability_duration(Duration::from_secs(30))
.with_failback_sweep_interval(Duration::from_secs(120))
.with_consecutive_hedge_win_threshold(3)
.build_from_env(&|_| None)
.unwrap();
assert!(options.circuit_breaker_enabled());
assert_eq!(options.read_failure_threshold(), 20);
assert_eq!(options.write_failure_threshold(), 7);
assert_eq!(options.counter_reset_window(), Duration::from_secs(60));
assert_eq!(
options.partition_unavailability_duration(),
Duration::from_secs(30)
);
assert_eq!(options.failback_sweep_interval(), Duration::from_secs(120));
assert_eq!(options.consecutive_hedge_win_threshold(), 3);
}
#[test]
fn circuit_breaker_override_unset_by_default() {
let options = PartitionFailoverOptionsBuilder::new()
.build_from_env(&|_| None)
.unwrap();
assert_eq!(options.circuit_breaker_enabled_override(), None);
}
#[test]
fn circuit_breaker_override_builder_value_round_trips() {
let off = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled_override(false)
.build_from_env(&|_| None)
.unwrap();
assert_eq!(off.circuit_breaker_enabled_override(), Some(false));
let on = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled_override(true)
.build_from_env(&|_| None)
.unwrap();
assert_eq!(on.circuit_breaker_enabled_override(), Some(true));
}
#[test]
fn read_failure_threshold_zero_rejected() {
let err = PartitionFailoverOptionsBuilder::new()
.with_read_failure_threshold(0)
.build_from_env(&|_| None)
.unwrap_err()
.to_string();
assert!(
err.contains("read_failure_threshold must be at least 1"),
"unexpected error: {err}",
);
}
#[test]
fn counter_reset_window_below_min_rejected() {
let err = PartitionFailoverOptionsBuilder::new()
.with_counter_reset_window(Duration::from_millis(500))
.build_from_env(&|_| None)
.unwrap_err()
.to_string();
assert!(
err.contains("counter_reset_window_ms must be at least 1000ms"),
"unexpected error: {err}",
);
}
#[test]
fn partition_unavailability_duration_below_min_rejected() {
let err = PartitionFailoverOptionsBuilder::new()
.with_partition_unavailability_duration(Duration::from_millis(50))
.build_from_env(&|_| None)
.unwrap_err()
.to_string();
assert!(
err.contains("partition_unavailability_duration_ms must be at least 1000ms"),
"unexpected error: {err}",
);
}
#[test]
fn failback_sweep_interval_below_min_rejected() {
let err = PartitionFailoverOptionsBuilder::new()
.with_failback_sweep_interval(Duration::from_millis(100))
.build_from_env(&|_| None)
.unwrap_err()
.to_string();
assert!(
err.contains("failback_sweep_interval_ms must be at least 1000ms"),
"unexpected error: {err}",
);
}
}
#[cfg(test)]
mod env_matrix_tests {
use super::*;
use std::collections::HashMap;
fn env_of(pairs: &[(&str, &str)]) -> impl Fn(&str) -> Option<String> {
let map: HashMap<String, String> = pairs
.iter()
.map(|(k, v)| ((*k).to_string(), (*v).to_string()))
.collect();
move |k: &str| map.get(k).cloned()
}
fn empty_env() -> impl Fn(&str) -> Option<String> {
|_: &str| None
}
#[test]
fn enabled_defaults_true_when_nothing_set() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&empty_env())
.unwrap();
assert!(o.circuit_breaker_enabled());
}
#[test]
fn enabled_env_false_honored_when_builder_unset() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED", "false")]))
.unwrap();
assert!(!o.circuit_breaker_enabled());
}
#[test]
fn enabled_builder_true_wins_over_env_false() {
let o = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled(true)
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED", "false")]))
.unwrap();
assert!(o.circuit_breaker_enabled());
}
#[test]
fn enabled_builder_false_wins_over_env_true() {
let o = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled(false)
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED", "true")]))
.unwrap();
assert!(!o.circuit_breaker_enabled());
}
#[test]
fn enabled_env_uses_strict_bool_parsing() {
for garbage in ["1", "yes", "on", "TRUE", "False", " false ", "enabled", ""] {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED", garbage)]))
.unwrap();
assert!(
o.circuit_breaker_enabled(),
"{garbage:?} should be ignored and fall back to default true",
);
}
}
#[test]
fn override_unset_by_default() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&empty_env())
.unwrap();
assert_eq!(o.circuit_breaker_enabled_override(), None);
}
#[test]
fn override_env_lenient_true_spellings() {
for raw in ["true", "1", "yes", "on", "TRUE", "On", " yes "] {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", raw)]))
.unwrap();
assert_eq!(
o.circuit_breaker_enabled_override(),
Some(true),
"{raw:?} should parse as Some(true)",
);
}
}
#[test]
fn override_env_lenient_false_spellings() {
for raw in ["false", "0", "no", "off", "FALSE", "Off", " no "] {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", raw)]))
.unwrap();
assert_eq!(
o.circuit_breaker_enabled_override(),
Some(false),
"{raw:?} should parse as Some(false)",
);
}
}
#[test]
fn override_env_garbage_ignored() {
for raw in ["maybe", "2", "", "tru", "noo"] {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", raw)]))
.unwrap();
assert_eq!(
o.circuit_breaker_enabled_override(),
None,
"{raw:?} should be ignored (treated as unset)",
);
}
}
#[test]
fn override_is_independent_of_base_enable() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[
("AZURE_COSMOS_PPCB_ENABLED", "false"),
("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", "true"),
]))
.unwrap();
assert!(!o.circuit_breaker_enabled());
assert_eq!(o.circuit_breaker_enabled_override(), Some(true));
}
#[test]
fn override_builder_wins_over_env() {
let o = PartitionFailoverOptionsBuilder::new()
.with_circuit_breaker_enabled_override(false)
.build_from_env(&env_of(&[("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", "true")]))
.unwrap();
assert_eq!(o.circuit_breaker_enabled_override(), Some(false));
}
#[test]
fn thresholds_env_values_honored() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[
("AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD", "20"),
("AZURE_COSMOS_PPCB_WRITE_FAILURE_THRESHOLD", "8"),
("AZURE_COSMOS_PPCB_CONSECUTIVE_HEDGE_WIN_THRESHOLD", "3"),
]))
.unwrap();
assert_eq!(o.read_failure_threshold(), 20);
assert_eq!(o.write_failure_threshold(), 8);
assert_eq!(o.consecutive_hedge_win_threshold(), 3);
}
#[test]
fn thresholds_builder_wins_over_env() {
let o = PartitionFailoverOptionsBuilder::new()
.with_read_failure_threshold(99)
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD",
"20",
)]))
.unwrap();
assert_eq!(o.read_failure_threshold(), 99);
}
#[test]
fn thresholds_default_when_env_unparseable() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD",
"not-a-number",
)]))
.unwrap();
assert_eq!(o.read_failure_threshold(), 10);
}
#[test]
fn threshold_env_zero_is_out_of_bounds_error() {
let err = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD",
"0",
)]))
.unwrap_err()
.to_string();
assert!(
err.contains("read_failure_threshold must be at least 1"),
"unexpected error: {err}",
);
}
#[test]
fn threshold_builder_zero_is_out_of_bounds_even_with_valid_env() {
let err = PartitionFailoverOptionsBuilder::new()
.with_write_failure_threshold(0)
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_WRITE_FAILURE_THRESHOLD",
"5",
)]))
.unwrap_err()
.to_string();
assert!(
err.contains("write_failure_threshold must be at least 1"),
"unexpected error: {err}",
);
}
#[test]
fn durations_env_values_honored() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[
("AZURE_COSMOS_PPCB_COUNTER_RESET_WINDOW_MS", "60000"),
(
"AZURE_COSMOS_PPCB_PARTITION_UNAVAILABILITY_DURATION_MS",
"30000",
),
("AZURE_COSMOS_PPCB_FAILBACK_SWEEP_INTERVAL_MS", "120000"),
]))
.unwrap();
assert_eq!(o.counter_reset_window(), Duration::from_millis(60_000));
assert_eq!(
o.partition_unavailability_duration(),
Duration::from_millis(30_000)
);
assert_eq!(o.failback_sweep_interval(), Duration::from_millis(120_000));
}
#[test]
fn duration_builder_wins_over_env() {
let o = PartitionFailoverOptionsBuilder::new()
.with_counter_reset_window(Duration::from_secs(90))
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_COUNTER_RESET_WINDOW_MS",
"60000",
)]))
.unwrap();
assert_eq!(o.counter_reset_window(), Duration::from_secs(90));
}
#[test]
fn duration_env_below_min_is_error() {
let err = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_COUNTER_RESET_WINDOW_MS",
"500",
)]))
.unwrap_err()
.to_string();
assert!(
err.contains("counter_reset_window_ms must be at least 1000ms"),
"unexpected error: {err}",
);
}
#[test]
fn duration_env_unparseable_falls_back_to_default() {
let o = PartitionFailoverOptionsBuilder::new()
.build_from_env(&env_of(&[(
"AZURE_COSMOS_PPCB_FAILBACK_SWEEP_INTERVAL_MS",
"soon",
)]))
.unwrap();
assert_eq!(o.failback_sweep_interval(), Duration::from_millis(300_000));
}
#[test]
fn kitchen_sink_mixed_builder_env_and_defaults() {
let o = PartitionFailoverOptionsBuilder::new()
.with_read_failure_threshold(33)
.with_circuit_breaker_enabled_override(false)
.build_from_env(&env_of(&[
("AZURE_COSMOS_PPCB_ENABLED", "false"),
("AZURE_COSMOS_PPCB_ENABLED_OVERRIDE", "true"), ("AZURE_COSMOS_PPCB_READ_FAILURE_THRESHOLD", "11"), ("AZURE_COSMOS_PPCB_WRITE_FAILURE_THRESHOLD", "9"),
("AZURE_COSMOS_PPCB_CONSECUTIVE_HEDGE_WIN_THRESHOLD", "xx"), ("AZURE_COSMOS_PPCB_COUNTER_RESET_WINDOW_MS", "45000"),
(
"AZURE_COSMOS_PPCB_PARTITION_UNAVAILABILITY_DURATION_MS",
"later",
), ]))
.unwrap();
assert!(!o.circuit_breaker_enabled()); assert_eq!(o.circuit_breaker_enabled_override(), Some(false)); assert_eq!(o.read_failure_threshold(), 33); assert_eq!(o.write_failure_threshold(), 9); assert_eq!(o.consecutive_hedge_win_threshold(), 5); assert_eq!(o.counter_reset_window(), Duration::from_millis(45_000)); assert_eq!(
o.partition_unavailability_duration(),
Duration::from_secs(5)
); assert_eq!(o.failback_sweep_interval(), Duration::from_secs(300)); }
}
#[cfg(test)]
mod real_env_tests {
use super::*;
use crate::options::env_parsing::test_env::{with_scoped_env, PPCB_ENV_VARS};
#[test]
fn real_env_enable_false_disables_ppcb() {
with_scoped_env(
PPCB_ENV_VARS,
&[("AZURE_COSMOS_PPCB_ENABLED", "false")],
|| {
let o = PartitionFailoverOptionsBuilder::new().build().unwrap();
assert!(!o.circuit_breaker_enabled());
},
);
}
#[test]
fn real_env_empty_uses_default_enabled() {
with_scoped_env(PPCB_ENV_VARS, &[], || {
let o = PartitionFailoverOptionsBuilder::new().build().unwrap();
assert!(o.circuit_breaker_enabled());
assert_eq!(o.circuit_breaker_enabled_override(), None);
});
}
}