Skip to main content

otel_bootstrap/
export_backoff.rs

1// SPDX-License-Identifier: MIT
2//! An export failure is logged once, then again only as its backoff allows.
3//!
4//! The OpenTelemetry SDK logs every failed export batch at `ERROR`: with the
5//! default five-second schedule and three signals, an unreachable collector
6//! writes a line every couple of seconds, forever, and buries every real
7//! error the service logs. [`ExportFailureBackoff`] lets the first failure of
8//! each kind through, then the next one only after a minute, then two, four,
9//! and so on up to an hour. Any other event, including the SDK's own events
10//! that are not failures, passes untouched.
11
12use std::collections::HashMap;
13use std::fmt;
14use std::sync::Mutex;
15use std::time::{Duration, Instant};
16
17use tracing::field::{Field, Visit};
18use tracing::{Event, Level, Subscriber};
19use tracing_subscriber::layer::{Context, Layer};
20
21const FIRST_BACKOFF: Duration = Duration::from_secs(60);
22const MAX_BACKOFF: Duration = Duration::from_secs(60 * 60);
23
24#[derive(Clone, Copy)]
25struct Window {
26    next_allowed: Instant,
27    backoff: Duration,
28}
29
30/// Rate-limits the SDK's own export-failure events; see the module docs.
31#[derive(Default)]
32pub struct ExportFailureBackoff {
33    windows: Mutex<HashMap<String, Window>>,
34}
35
36impl ExportFailureBackoff {
37    /// Whether an event of `kind` at `now` is let through, and records it.
38    fn admit(&self, kind: &str, now: Instant) -> bool {
39        let mut windows = self
40            .windows
41            .lock()
42            .unwrap_or_else(std::sync::PoisonError::into_inner);
43        match windows.get_mut(kind) {
44            None => {
45                windows.insert(
46                    kind.to_owned(),
47                    Window {
48                        next_allowed: now + FIRST_BACKOFF,
49                        backoff: FIRST_BACKOFF,
50                    },
51                );
52                true
53            }
54            Some(window) if now >= window.next_allowed => {
55                window.backoff = (window.backoff * 2).min(MAX_BACKOFF);
56                window.next_allowed = now + window.backoff;
57                true
58            }
59            Some(_) => false,
60        }
61    }
62}
63
64/// The SDK names each internal event in a `name` field, e.g.
65/// `BatchSpanProcessor.ExportError`.
66#[derive(Default)]
67struct EventName(Option<String>);
68
69impl Visit for EventName {
70    fn record_str(&mut self, field: &Field, value: &str) {
71        if field.name() == "name" {
72            self.0 = Some(value.to_owned());
73        }
74    }
75
76    fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) {
77        if field.name() == "name" && self.0.is_none() {
78            self.0 = Some(format!("{value:?}").trim_matches('"').to_owned());
79        }
80    }
81}
82
83/// The failure kind of an SDK export-failure event, or `None` for any other
84/// event.
85fn export_failure_kind(event: &Event<'_>) -> Option<String> {
86    let metadata = event.metadata();
87    if !metadata.target().starts_with("opentelemetry") || *metadata.level() > Level::WARN {
88        return None;
89    }
90    let mut name = EventName::default();
91    event.record(&mut name);
92    let kind = name.0.unwrap_or_else(|| metadata.name().to_owned());
93    kind.contains("Export").then_some(kind)
94}
95
96impl<S: Subscriber> Layer<S> for ExportFailureBackoff {
97    fn event_enabled(&self, event: &Event<'_>, _ctx: Context<'_, S>) -> bool {
98        match export_failure_kind(event) {
99            Some(kind) => self.admit(&kind, Instant::now()),
100            None => true,
101        }
102    }
103}
104
105#[cfg(test)]
106mod tests {
107    use super::*;
108
109    #[test]
110    fn the_first_failure_is_logged_and_repeats_wait_for_the_backoff() {
111        let backoff = ExportFailureBackoff::default();
112        let start = Instant::now();
113        assert!(backoff.admit("BatchSpanProcessor.ExportError", start));
114        assert!(!backoff.admit(
115            "BatchSpanProcessor.ExportError",
116            start + Duration::from_secs(5)
117        ));
118        assert!(!backoff.admit(
119            "BatchSpanProcessor.ExportError",
120            start + Duration::from_secs(59)
121        ));
122        assert!(backoff.admit(
123            "BatchSpanProcessor.ExportError",
124            start + Duration::from_secs(60)
125        ));
126        // The window doubles: two minutes after the second line.
127        assert!(!backoff.admit(
128            "BatchSpanProcessor.ExportError",
129            start + Duration::from_secs(170)
130        ));
131        assert!(backoff.admit(
132            "BatchSpanProcessor.ExportError",
133            start + Duration::from_secs(180)
134        ));
135    }
136
137    #[test]
138    fn each_kind_of_failure_has_its_own_window() {
139        let backoff = ExportFailureBackoff::default();
140        let start = Instant::now();
141        assert!(backoff.admit("BatchSpanProcessor.ExportError", start));
142        assert!(backoff.admit("BatchLogProcessor.ExportError", start));
143    }
144
145    #[test]
146    fn the_backoff_never_exceeds_an_hour() {
147        let backoff = ExportFailureBackoff::default();
148        let mut now = Instant::now();
149        assert!(backoff.admit("k", now));
150        for _ in 0..20 {
151            now += MAX_BACKOFF;
152            assert!(backoff.admit("k", now));
153        }
154    }
155
156    #[test]
157    fn a_name_recorded_as_debug_is_read_as_its_text() {
158        use tracing_subscriber::prelude::*;
159
160        let seen = std::sync::Arc::new(Mutex::new(0usize));
161        struct Count(std::sync::Arc<Mutex<usize>>);
162        impl<S: Subscriber> Layer<S> for Count {
163            fn on_event(&self, _event: &Event<'_>, _ctx: Context<'_, S>) {
164                *self.0.lock().unwrap() += 1;
165            }
166        }
167        let subscriber = tracing_subscriber::registry()
168            .with(ExportFailureBackoff::default())
169            .with(Count(seen.clone()));
170        tracing::subscriber::with_default(subscriber, || {
171            for _ in 0..3 {
172                tracing::warn!(target: "opentelemetry_sdk", name = ?"MetricReader.ExportError", "export failed");
173            }
174        });
175        assert_eq!(*seen.lock().unwrap(), 1);
176    }
177
178    #[test]
179    fn only_sdk_export_failures_are_rate_limited() {
180        use tracing_subscriber::prelude::*;
181
182        let seen = std::sync::Arc::new(Mutex::new(Vec::<String>::new()));
183        struct Record(std::sync::Arc<Mutex<Vec<String>>>);
184        impl<S: Subscriber> Layer<S> for Record {
185            fn on_event(&self, event: &Event<'_>, _ctx: Context<'_, S>) {
186                self.0
187                    .lock()
188                    .unwrap()
189                    .push(event.metadata().target().to_owned());
190            }
191        }
192        let subscriber = tracing_subscriber::registry()
193            .with(ExportFailureBackoff::default())
194            .with(Record(seen.clone()));
195        tracing::subscriber::with_default(subscriber, || {
196            for _ in 0..3 {
197                tracing::error!(target: "opentelemetry_sdk", name = "BatchSpanProcessor.ExportError", "export failed");
198                tracing::error!(target: "my_service", "a real error");
199            }
200        });
201        let seen = seen.lock().unwrap();
202        assert_eq!(seen.iter().filter(|t| *t == "opentelemetry_sdk").count(), 1);
203        assert_eq!(seen.iter().filter(|t| *t == "my_service").count(), 3);
204    }
205}