- 1
//! Cross-run circuit breaker for provider health. Only retryable failures - 2
//! (429/529/network) trip it; auth/config errors are the caller's problem. - 3
//! While open, steps fail fast with the remaining cooldown instead of - 4
//! burning their retry budget against a dead provider. - 5
- 6
use std::collections::HashMap; - 7
use std::sync::Mutex; - 8
use std::time::{Duration, Instant}; - 9
- 10
#[derive(Debug, Clone)] - 11
pub struct CircuitBreakerConfig { - 12
/// Consecutive retryable failures before the circuit opens. - 13
pub threshold: u32, - 14
/// How long an open circuit stays open before allowing a probe. - 15
pub cooldown: Duration, - 16
} - 17
- 18
impl Default for CircuitBreakerConfig { - 19
fn default() -> Self { - 20
CircuitBreakerConfig { - 21
threshold: 5, - 22
cooldown: Duration::from_secs(60), - 23
} - 24
} - 25
} - 26
- 27
#[derive(Debug, Default)] - 28
struct State { - 29
consecutive_failures: u32, - 30
opened_at: Option<Instant>, - 31
probe_in_flight: bool, - 32
} - 33
- 34
#[derive(Debug)] - 35
pub struct CircuitBreaker { - 36
config: CircuitBreakerConfig, - 37
state: Mutex<HashMap<String, State>>, - 38
} - 39
- 40
#[derive(Debug, thiserror::Error)] - 41
#[error( - 42
"circuit open for provider (cooling down {remaining_secs}s after {failures} consecutive failures)" - 43
)] - 44
pub struct CircuitOpen { - 45
pub remaining_secs: u64, - 46
pub failures: u32, - 47
} - 48
- 49
impl CircuitBreaker { - 50
pub fn new(config: CircuitBreakerConfig) -> Self { - 51
CircuitBreaker { - 52
config, - 53
state: Mutex::new(HashMap::new()), - 54
} - 55
} - 56
- 57
/// Err when the circuit is open and the cooldown has not elapsed. - 58
pub fn check(&self) -> Result<(), CircuitOpen> { - 59
self.check_key("") - 60
} - 61
- 62
/// Check one provider/key circuit. Provider names are used by the agent - 63
/// as the stable health domain; an empty key means the process-global - 64
/// helper semantics for callers that do not have a route identity. - 65
pub fn check_key(&self, key: &str) -> Result<(), CircuitOpen> { - 66
let mut states = self.lock(); - 67
let st = states.entry(key.to_string()).or_default(); - 68
if st.probe_in_flight { - 69
return Err(CircuitOpen { - 70
remaining_secs: 1, - 71
failures: st.consecutive_failures, - 72
}); - 73
} - 74
if let Some(opened_at) = st.opened_at { - 75
let elapsed = opened_at.elapsed(); - 76
if elapsed < self.config.cooldown { - 77
return Err(CircuitOpen { - 78
remaining_secs: (self.config.cooldown - elapsed).as_secs().max(1), - 79
failures: st.consecutive_failures, - 80
}); - 81
} - 82
// Cooldown elapsed: reserve exactly one half-open probe. Other - 83
// callers remain fail-closed until that probe settles. - 84
st.opened_at = None; - 85
st.probe_in_flight = true; - 86
} - 87
Ok(()) - 88
} - 89
- 90
pub fn record_success(&self) { - 91
self.record_success_key(""); - 92
} - 93
- 94
pub fn record_success_key(&self, key: &str) { - 95
let mut states = self.lock(); - 96
let st = states.entry(key.to_string()).or_default(); - 97
st.consecutive_failures = 0; - 98
st.opened_at = None; - 99
st.probe_in_flight = false; - 100
} - 101
- 102
/// Only retryable failures should call this. - 103
pub fn record_failure(&self) { - 104
self.record_failure_key(""); - 105
} - 106
- 107
pub fn record_failure_key(&self, key: &str) { - 108
let mut states = self.lock(); - 109
let st = states.entry(key.to_string()).or_default(); - 110
st.consecutive_failures = st.consecutive_failures.saturating_add(1); - 111
if st.consecutive_failures >= self.config.threshold && st.opened_at.is_none() { - 112
st.opened_at = Some(Instant::now()); - 113
} - 114
st.probe_in_flight = false; - 115
} - 116
- 117
fn lock(&self) -> std::sync::MutexGuard<'_, HashMap<String, State>> { - 118
self.state - 119
.lock() - 120
.unwrap_or_else(std::sync::PoisonError::into_inner) - 121
} - 122
} - 123
Indexing the workspace…
Vakyartha documentation is discovering safe artifacts, anchors, and source references.