otel_bootstrap/
export_backoff.rs1use 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#[derive(Default)]
32pub struct ExportFailureBackoff {
33 windows: Mutex<HashMap<String, Window>>,
34}
35
36impl ExportFailureBackoff {
37 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#[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
83fn 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 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}