use std::collections::BTreeMap;
use std::future::Future;
use std::sync::{Mutex, OnceLock};
use std::time::{Duration, Instant};
use crate::errors::TelemetryError;
use crate::health::{increment_retries, record_export_failure, record_export_latency};
use crate::sampling::Signal;
pub(crate) const CIRCUIT_BREAKER_THRESHOLD: u32 = 3;
pub(crate) const CIRCUIT_COOLDOWN: Duration = Duration::from_secs(30);
pub(crate) const MAX_EXPORT_ATTEMPTS: u32 = crate::config::MAX_EXPORTER_RETRIES as u32 + 1;
fn capped_attempts(retries: u32) -> u32 {
retries.saturating_add(1).min(MAX_EXPORT_ATTEMPTS)
}
#[derive(Clone, Debug, PartialEq)]
pub struct ExporterPolicy {
pub retries: u32,
pub backoff_seconds: f64,
pub timeout_seconds: f64,
pub fail_open: bool,
pub allow_blocking_in_event_loop: bool,
}
impl Default for ExporterPolicy {
fn default() -> Self {
Self {
retries: 0,
backoff_seconds: 0.0,
timeout_seconds: 10.0,
fail_open: true,
allow_blocking_in_event_loop: false,
}
}
}
#[derive(Clone, Debug, Default)]
struct CircuitState {
consecutive_timeouts: u32,
tripped_at: Option<Instant>,
open_count: u32,
half_open_probing: bool,
}
static POLICIES: OnceLock<Mutex<BTreeMap<Signal, ExporterPolicy>>> = OnceLock::new();
static CIRCUITS: OnceLock<Mutex<BTreeMap<Signal, CircuitState>>> = OnceLock::new();
fn default_policies_mutex() -> Mutex<BTreeMap<Signal, ExporterPolicy>> {
Mutex::new(BTreeMap::from([
(Signal::Logs, ExporterPolicy::default()),
(Signal::Traces, ExporterPolicy::default()),
(Signal::Metrics, ExporterPolicy::default()),
]))
}
fn policies() -> &'static Mutex<BTreeMap<Signal, ExporterPolicy>> {
POLICIES.get_or_init(default_policies_mutex)
}
fn default_circuits_mutex() -> Mutex<BTreeMap<Signal, CircuitState>> {
Mutex::new(BTreeMap::from([
(Signal::Logs, CircuitState::default()),
(Signal::Traces, CircuitState::default()),
(Signal::Metrics, CircuitState::default()),
]))
}
fn circuits() -> &'static Mutex<BTreeMap<Signal, CircuitState>> {
CIRCUITS.get_or_init(default_circuits_mutex)
}
fn backoff_duration(backoff_seconds: f64, has_tokio_reactor: bool) -> Option<Duration> {
if backoff_seconds <= 0.0 || !has_tokio_reactor {
None
} else {
Some(Duration::from_secs_f64(backoff_seconds))
}
}
async fn wait_before_retry(
signal: Signal,
attempt: u32,
backoff_seconds: f64,
has_tokio_reactor: bool,
) {
if attempt == 0 {
return;
}
if let Some(backoff) = backoff_duration(backoff_seconds, has_tokio_reactor) {
tokio::time::sleep(backoff).await;
}
increment_retries(signal, 1);
}
pub fn set_exporter_policy(
signal: Signal,
policy: ExporterPolicy,
) -> Result<ExporterPolicy, TelemetryError> {
crate::_lock::lock(policies()).insert(signal, policy.clone());
Ok(policy)
}
pub fn get_exporter_policy(signal: Signal) -> Result<ExporterPolicy, TelemetryError> {
let policy_lock = crate::_lock::lock(policies());
match policy_lock.get(&signal).cloned() {
Some(policy) => Ok(policy),
None => Err(TelemetryError::new("unknown signal")),
}
}
pub fn get_circuit_state(signal: Signal) -> Result<(String, u32, f64), TelemetryError> {
let circuits = crate::_lock::lock(circuits());
let state = match circuits.get(&signal).cloned() {
Some(state) => state,
None => return Err(TelemetryError::new("unknown signal")),
};
Ok(describe_circuit_state(&state))
}
pub(crate) async fn run_with_resilience_inner<F, Fut, T, E>(
signal: Signal,
policy: &ExporterPolicy,
operation: F,
timeout_err: impl Fn(Duration) -> E,
is_sdk_timeout: impl Fn(&E) -> bool,
circuit_open_err: impl Fn() -> E,
) -> Result<Option<T>, E>
where
F: Fn() -> Fut,
Fut: Future<Output = Result<T, E>>,
{
let timeout = Duration::from_secs_f64(policy.timeout_seconds.max(0.0));
let has_tokio_reactor = tokio::runtime::Handle::try_current().is_ok();
let timeout_active = if timeout.is_zero() {
false
} else {
has_tokio_reactor
};
let should_probe = if timeout_active {
_check_and_start_probe_for_wrappers(signal)
} else {
false
};
if should_probe {
return if policy.fail_open {
Ok(None)
} else {
Err(circuit_open_err())
};
}
let max_attempts = capped_attempts(policy.retries);
let mut last_err: Option<E> = None;
for attempt in 0..max_attempts {
wait_before_retry(signal, attempt, policy.backoff_seconds, has_tokio_reactor).await;
let started = Instant::now();
let (result, wrapper_timeout) = if !timeout_active {
(operation().await, false)
} else {
match tokio::time::timeout(timeout, operation()).await {
Ok(inner) => (inner, false),
Err(_) => (Err(timeout_err(timeout)), true),
}
};
match result {
Ok(value) => {
record_export_latency(signal, started.elapsed().as_secs_f64() * 1000.0);
_record_circuit_success_for_wrappers(signal);
return Ok(Some(value));
}
Err(err) => {
record_export_failure(signal);
let is_timeout = if wrapper_timeout {
true
} else {
is_sdk_timeout(&err)
};
_record_circuit_failure_for_wrappers(signal, is_timeout);
last_err = Some(err);
}
}
}
if policy.fail_open {
Ok(None)
} else {
Err(last_err.expect("retry loop ran at least once"))
}
}
fn cooldown_remaining(tripped_at: Option<Instant>) -> f64 {
match tripped_at {
Some(instant) => CIRCUIT_COOLDOWN
.saturating_sub(instant.elapsed())
.as_secs_f64(),
None => 0.0,
}
}
fn describe_circuit_state(state: &CircuitState) -> (String, u32, f64) {
if state.half_open_probing {
return ("half-open".to_string(), state.open_count, 0.0);
}
if state.consecutive_timeouts >= CIRCUIT_BREAKER_THRESHOLD {
let remaining = cooldown_remaining(state.tripped_at);
if remaining > 0.0 {
return ("open".to_string(), state.open_count, remaining);
}
return ("half-open".to_string(), state.open_count, 0.0);
}
("closed".to_string(), state.open_count, 0.0)
}
pub async fn run_with_resilience<F, Fut, T>(
signal: Signal,
operation: F,
) -> Result<Option<T>, TelemetryError>
where
F: Fn() -> Fut,
Fut: Future<Output = Result<T, TelemetryError>>,
{
let policy = get_exporter_policy(signal)?;
run_with_resilience_inner(
signal,
&policy,
operation,
|_| TelemetryError::new("operation timed out"),
|_| false,
|| TelemetryError::new("circuit breaker open"),
)
.await
}
pub(crate) fn _record_circuit_failure_for_wrappers(signal: Signal, is_timeout: bool) {
let mut circuit_lock = crate::_lock::lock(circuits());
let Some(state) = circuit_lock.get_mut(&signal) else {
return;
};
if state.half_open_probing {
state.half_open_probing = false;
state.open_count += 1;
state.tripped_at = Some(Instant::now());
return;
}
if !is_timeout {
state.consecutive_timeouts = 0;
return;
}
state.consecutive_timeouts += 1;
if state.consecutive_timeouts >= CIRCUIT_BREAKER_THRESHOLD {
state.open_count += 1;
state.tripped_at = Some(Instant::now());
}
}
pub(crate) fn _record_circuit_success_for_wrappers(signal: Signal) {
let mut circuit_lock = crate::_lock::lock(circuits());
let Some(state) = circuit_lock.get_mut(&signal) else {
return;
};
if state.half_open_probing {
state.half_open_probing = false;
}
state.consecutive_timeouts = 0;
}
fn circuit_cooldown_is_active(elapsed: Duration) -> bool {
elapsed < CIRCUIT_COOLDOWN
}
pub(crate) fn _check_and_start_probe_for_wrappers(signal: Signal) -> bool {
let mut circuit_lock = crate::_lock::lock(circuits());
let Some(state) = circuit_lock.get_mut(&signal) else {
return false;
};
if state.consecutive_timeouts < CIRCUIT_BREAKER_THRESHOLD {
return false;
}
let cooldown_active = state
.tripped_at
.map(|instant| circuit_cooldown_is_active(instant.elapsed()))
.unwrap_or(false);
if cooldown_active {
return true; }
if state.half_open_probing {
return true; }
state.half_open_probing = true;
false
}
pub fn _reset_resilience_for_tests() {
*crate::_lock::lock(policies()) = BTreeMap::from([
(Signal::Logs, ExporterPolicy::default()),
(Signal::Traces, ExporterPolicy::default()),
(Signal::Metrics, ExporterPolicy::default()),
]);
*crate::_lock::lock(circuits()) = BTreeMap::from([
(Signal::Logs, CircuitState::default()),
(Signal::Traces, CircuitState::default()),
(Signal::Metrics, CircuitState::default()),
]);
}
pub fn _clear_resilience_state_for_tests() {
crate::_lock::lock(policies()).clear();
crate::_lock::lock(circuits()).clear();
}
#[cfg(test)]
#[path = "resilience_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "resilience_inner_callback_tests.rs"]
mod inner_callback_tests;
#[cfg(test)]
#[path = "resilience_state_tests.rs"]
mod state_tests;