1use std::collections::BTreeMap;
7use std::future::Future;
8use std::sync::{Mutex, OnceLock};
9use std::time::{Duration, Instant};
10
11use crate::errors::TelemetryError;
12use crate::health::{increment_retries, record_export_failure, record_export_latency};
13use crate::sampling::Signal;
14
15pub(crate) const CIRCUIT_BREAKER_THRESHOLD: u32 = 3;
16pub(crate) const CIRCUIT_COOLDOWN: Duration = Duration::from_secs(30);
17pub(crate) const MAX_EXPORT_ATTEMPTS: u32 = crate::config::MAX_EXPORTER_RETRIES as u32 + 1;
20
21fn capped_attempts(retries: u32) -> u32 {
28 retries.saturating_add(1).min(MAX_EXPORT_ATTEMPTS)
29}
30
31#[derive(Clone, Debug, PartialEq)]
32pub struct ExporterPolicy {
33 pub retries: u32,
34 pub backoff_seconds: f64,
35 pub timeout_seconds: f64,
36 pub fail_open: bool,
37 pub allow_blocking_in_event_loop: bool,
38}
39
40impl Default for ExporterPolicy {
41 fn default() -> Self {
42 Self {
43 retries: 0,
44 backoff_seconds: 0.0,
45 timeout_seconds: 10.0,
46 fail_open: true,
47 allow_blocking_in_event_loop: false,
48 }
49 }
50}
51
52#[derive(Clone, Debug, Default)]
53struct CircuitState {
54 consecutive_timeouts: u32,
55 tripped_at: Option<Instant>,
56 open_count: u32,
57 half_open_probing: bool,
59}
60
61static POLICIES: OnceLock<Mutex<BTreeMap<Signal, ExporterPolicy>>> = OnceLock::new();
62static CIRCUITS: OnceLock<Mutex<BTreeMap<Signal, CircuitState>>> = OnceLock::new();
63
64fn default_policies_mutex() -> Mutex<BTreeMap<Signal, ExporterPolicy>> {
65 Mutex::new(BTreeMap::from([
66 (Signal::Logs, ExporterPolicy::default()),
67 (Signal::Traces, ExporterPolicy::default()),
68 (Signal::Metrics, ExporterPolicy::default()),
69 ]))
70}
71
72fn policies() -> &'static Mutex<BTreeMap<Signal, ExporterPolicy>> {
73 POLICIES.get_or_init(default_policies_mutex)
74}
75
76fn default_circuits_mutex() -> Mutex<BTreeMap<Signal, CircuitState>> {
77 Mutex::new(BTreeMap::from([
78 (Signal::Logs, CircuitState::default()),
79 (Signal::Traces, CircuitState::default()),
80 (Signal::Metrics, CircuitState::default()),
81 ]))
82}
83
84fn circuits() -> &'static Mutex<BTreeMap<Signal, CircuitState>> {
85 CIRCUITS.get_or_init(default_circuits_mutex)
86}
87
88fn backoff_duration(backoff_seconds: f64, has_tokio_reactor: bool) -> Option<Duration> {
89 if backoff_seconds <= 0.0 || !has_tokio_reactor {
90 None
91 } else {
92 Some(Duration::from_secs_f64(backoff_seconds))
93 }
94}
95
96async fn wait_before_retry(
97 signal: Signal,
98 attempt: u32,
99 backoff_seconds: f64,
100 has_tokio_reactor: bool,
101) {
102 if attempt == 0 {
103 return;
104 }
105 if let Some(backoff) = backoff_duration(backoff_seconds, has_tokio_reactor) {
106 tokio::time::sleep(backoff).await;
107 }
108 increment_retries(signal, 1);
109}
110
111pub fn set_exporter_policy(
112 signal: Signal,
113 policy: ExporterPolicy,
114) -> Result<ExporterPolicy, TelemetryError> {
115 crate::_lock::lock(policies()).insert(signal, policy.clone());
116 Ok(policy)
117}
118
119pub fn get_exporter_policy(signal: Signal) -> Result<ExporterPolicy, TelemetryError> {
120 let policy_lock = crate::_lock::lock(policies());
121 match policy_lock.get(&signal).cloned() {
122 Some(policy) => Ok(policy),
123 None => Err(TelemetryError::new("unknown signal")),
124 }
125}
126
127pub fn get_circuit_state(signal: Signal) -> Result<(String, u32, f64), TelemetryError> {
128 let circuits = crate::_lock::lock(circuits());
129 let state = match circuits.get(&signal).cloned() {
130 Some(state) => state,
131 None => return Err(TelemetryError::new("unknown signal")),
132 };
133 Ok(describe_circuit_state(&state))
134}
135
136pub(crate) async fn run_with_resilience_inner<F, Fut, T, E>(
152 signal: Signal,
153 policy: &ExporterPolicy,
154 operation: F,
155 timeout_err: impl Fn(Duration) -> E,
156 is_sdk_timeout: impl Fn(&E) -> bool,
157 circuit_open_err: impl Fn() -> E,
158) -> Result<Option<T>, E>
159where
160 F: Fn() -> Fut,
161 Fut: Future<Output = Result<T, E>>,
162{
163 let timeout = Duration::from_secs_f64(policy.timeout_seconds.max(0.0));
164 let has_tokio_reactor = tokio::runtime::Handle::try_current().is_ok();
171 let timeout_active = if timeout.is_zero() {
172 false
173 } else {
174 has_tokio_reactor
175 };
176 let should_probe = if timeout_active {
181 _check_and_start_probe_for_wrappers(signal)
182 } else {
183 false
184 };
185 if should_probe {
186 return if policy.fail_open {
187 Ok(None)
188 } else {
189 Err(circuit_open_err())
190 };
191 }
192
193 let max_attempts = capped_attempts(policy.retries);
194 let mut last_err: Option<E> = None;
195 for attempt in 0..max_attempts {
196 wait_before_retry(signal, attempt, policy.backoff_seconds, has_tokio_reactor).await;
197
198 let started = Instant::now();
199 let (result, wrapper_timeout) = if !timeout_active {
200 (operation().await, false)
201 } else {
202 match tokio::time::timeout(timeout, operation()).await {
203 Ok(inner) => (inner, false),
204 Err(_) => (Err(timeout_err(timeout)), true),
205 }
206 };
207
208 match result {
209 Ok(value) => {
210 record_export_latency(signal, started.elapsed().as_secs_f64() * 1000.0);
211 _record_circuit_success_for_wrappers(signal);
212 return Ok(Some(value));
213 }
214 Err(err) => {
215 record_export_failure(signal);
216 let is_timeout = if wrapper_timeout {
220 true
221 } else {
222 is_sdk_timeout(&err)
223 };
224 _record_circuit_failure_for_wrappers(signal, is_timeout);
225 last_err = Some(err);
226 }
227 }
228 }
229
230 if policy.fail_open {
231 Ok(None)
232 } else {
233 Err(last_err.expect("retry loop ran at least once"))
236 }
237}
238
239fn cooldown_remaining(tripped_at: Option<Instant>) -> f64 {
240 match tripped_at {
241 Some(instant) => CIRCUIT_COOLDOWN
242 .saturating_sub(instant.elapsed())
243 .as_secs_f64(),
244 None => 0.0,
245 }
246}
247
248fn describe_circuit_state(state: &CircuitState) -> (String, u32, f64) {
249 if state.half_open_probing {
250 return ("half-open".to_string(), state.open_count, 0.0);
251 }
252 if state.consecutive_timeouts >= CIRCUIT_BREAKER_THRESHOLD {
253 let remaining = cooldown_remaining(state.tripped_at);
254 if remaining > 0.0 {
255 return ("open".to_string(), state.open_count, remaining);
256 }
257 return ("half-open".to_string(), state.open_count, 0.0);
258 }
259 ("closed".to_string(), state.open_count, 0.0)
260}
261
262pub async fn run_with_resilience<F, Fut, T>(
263 signal: Signal,
264 operation: F,
265) -> Result<Option<T>, TelemetryError>
266where
267 F: Fn() -> Fut,
268 Fut: Future<Output = Result<T, TelemetryError>>,
269{
270 let policy = get_exporter_policy(signal)?;
271 run_with_resilience_inner(
272 signal,
273 &policy,
274 operation,
275 |_| TelemetryError::new("operation timed out"),
276 |_| false,
277 || TelemetryError::new("circuit breaker open"),
278 )
279 .await
280}
281
282pub(crate) fn _record_circuit_failure_for_wrappers(signal: Signal, is_timeout: bool) {
291 let mut circuit_lock = crate::_lock::lock(circuits());
292 let Some(state) = circuit_lock.get_mut(&signal) else {
293 return;
294 };
295 if state.half_open_probing {
296 state.half_open_probing = false;
297 state.open_count += 1;
298 state.tripped_at = Some(Instant::now());
299 return;
300 }
301 if !is_timeout {
302 state.consecutive_timeouts = 0;
303 return;
304 }
305 state.consecutive_timeouts += 1;
306 if state.consecutive_timeouts >= CIRCUIT_BREAKER_THRESHOLD {
307 state.open_count += 1;
308 state.tripped_at = Some(Instant::now());
309 }
310}
311
312pub(crate) fn _record_circuit_success_for_wrappers(signal: Signal) {
316 let mut circuit_lock = crate::_lock::lock(circuits());
317 let Some(state) = circuit_lock.get_mut(&signal) else {
318 return;
319 };
320 if state.half_open_probing {
321 state.half_open_probing = false;
322 }
323 state.consecutive_timeouts = 0;
324}
325
326fn circuit_cooldown_is_active(elapsed: Duration) -> bool {
327 elapsed < CIRCUIT_COOLDOWN
328}
329
330pub(crate) fn _check_and_start_probe_for_wrappers(signal: Signal) -> bool {
336 let mut circuit_lock = crate::_lock::lock(circuits());
337 let Some(state) = circuit_lock.get_mut(&signal) else {
338 return false;
339 };
340 if state.consecutive_timeouts < CIRCUIT_BREAKER_THRESHOLD {
341 return false;
342 }
343 let cooldown_active = state
344 .tripped_at
345 .map(|instant| circuit_cooldown_is_active(instant.elapsed()))
346 .unwrap_or(false);
347 if cooldown_active {
348 return true; }
350 if state.half_open_probing {
351 return true; }
353 state.half_open_probing = true;
355 false
356}
357
358pub fn _reset_resilience_for_tests() {
359 *crate::_lock::lock(policies()) = BTreeMap::from([
360 (Signal::Logs, ExporterPolicy::default()),
361 (Signal::Traces, ExporterPolicy::default()),
362 (Signal::Metrics, ExporterPolicy::default()),
363 ]);
364 *crate::_lock::lock(circuits()) = BTreeMap::from([
365 (Signal::Logs, CircuitState::default()),
366 (Signal::Traces, CircuitState::default()),
367 (Signal::Metrics, CircuitState::default()),
368 ]);
369}
370
371pub fn _clear_resilience_state_for_tests() {
372 crate::_lock::lock(policies()).clear();
373 crate::_lock::lock(circuits()).clear();
374}
375
376#[cfg(test)]
377#[path = "resilience_tests.rs"]
378mod tests;
379
380#[cfg(test)]
381#[path = "resilience_inner_callback_tests.rs"]
382mod inner_callback_tests;
383
384#[cfg(test)]
385#[path = "resilience_state_tests.rs"]
386mod state_tests;