polyc_runtime/
observability.rs1use std::io;
22use std::sync::OnceLock;
23
24use opentelemetry::KeyValue;
25use opentelemetry::global;
26use opentelemetry::trace::TracerProvider as _;
27use opentelemetry_otlp::SpanExporter;
28use opentelemetry_sdk::Resource;
29use opentelemetry_sdk::propagation::TraceContextPropagator;
30use opentelemetry_sdk::trace::SdkTracerProvider;
31use tokio::sync::broadcast;
32use tracing::Subscriber;
33use tracing_subscriber::fmt::format::{FormatEvent, Writer};
34use tracing_subscriber::fmt::{FmtContext, MakeWriter};
35use tracing_subscriber::registry::LookupSpan;
36use tracing_subscriber::{EnvFilter, fmt, prelude::*};
37
38const LOG_BROADCAST_CAPACITY: usize = 512;
42
43static LOG_BROADCAST: OnceLock<broadcast::Sender<String>> = OnceLock::new();
48
49#[must_use]
56pub fn subscribe_logs() -> Option<broadcast::Receiver<String>> {
57 LOG_BROADCAST.get().map(broadcast::Sender::subscribe)
58}
59
60#[derive(Clone)]
63struct BroadcastTee {
64 tx: broadcast::Sender<String>,
65}
66
67impl<'a> MakeWriter<'a> for BroadcastTee {
68 type Writer = TeeLine;
69
70 fn make_writer(&'a self) -> Self::Writer {
71 TeeLine {
72 buf: Vec::new(),
73 tx: self.tx.clone(),
74 }
75 }
76}
77
78struct TeeLine {
82 buf: Vec<u8>,
83 tx: broadcast::Sender<String>,
84}
85
86impl io::Write for TeeLine {
87 fn write(&mut self, data: &[u8]) -> io::Result<usize> {
88 io::stderr().write_all(data)?;
90 self.buf.extend_from_slice(data);
91 Ok(data.len())
92 }
93
94 fn flush(&mut self) -> io::Result<()> {
95 io::stderr().flush()
96 }
97}
98
99impl Drop for TeeLine {
100 fn drop(&mut self) {
101 if self.buf.is_empty() {
102 return;
103 }
104 let line = String::from_utf8_lossy(&self.buf).trim_end().to_owned();
106 if !line.is_empty() {
107 let _ = self.tx.send(line);
110 }
111 }
112}
113
114pub struct ShutdownGuard {
118 provider: Option<SdkTracerProvider>,
119}
120
121impl Drop for ShutdownGuard {
122 fn drop(&mut self) {
123 if let Some(provider) = self.provider.take() {
124 if let Err(err) = provider.shutdown() {
126 eprintln!("opentelemetry shutdown failed: {err}");
127 }
128 }
129 }
130}
131
132#[must_use]
137pub fn init(service_name: &str) -> ShutdownGuard {
138 let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
139 let json = std::env::var("RUST_LOG_FORMAT").is_ok_and(|v| v.eq_ignore_ascii_case("json"));
140
141 let (tracer, provider) = build_otel_tracer(service_name);
142 let otel_layer = tracer.map(|t| tracing_opentelemetry::layer().with_tracer(t));
143
144 global::set_text_map_propagator(TraceContextPropagator::new());
152
153 let (log_tx, _seed) = broadcast::channel(LOG_BROADCAST_CAPACITY);
158 let _ = LOG_BROADCAST.set(log_tx.clone());
159 let tee = BroadcastTee { tx: log_tx };
160
161 let registry = tracing_subscriber::registry().with(filter).with(otel_layer);
162 if json {
163 let fmt_layer = fmt::layer()
164 .json()
165 .flatten_event(true)
166 .map_event_format(CloudLoggingSeverity::wrap)
167 .with_writer(tee);
168 registry.with(fmt_layer).init();
169 } else {
170 let fmt_layer = fmt::layer().with_writer(tee);
171 registry.with(fmt_layer).init();
172 }
173
174 install_panic_hook();
175 ShutdownGuard { provider }
176}
177
178fn build_otel_tracer(
189 service_name: &str,
190) -> (
191 Option<opentelemetry_sdk::trace::Tracer>,
192 Option<SdkTracerProvider>,
193) {
194 if std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT").is_err() {
195 return (None, None);
196 }
197 let resolved_name = std::env::var("OTEL_SERVICE_NAME")
198 .ok()
199 .unwrap_or_else(|| service_name.to_owned());
200 let exporter = match SpanExporter::builder().with_tonic().build() {
201 Ok(exp) => exp,
202 Err(err) => {
203 eprintln!("opentelemetry OTLP exporter build failed: {err}; continuing without traces");
204 return (None, None);
205 }
206 };
207 let resource = Resource::builder()
208 .with_attribute(KeyValue::new("service.name", resolved_name.clone()))
209 .build();
210 let provider = SdkTracerProvider::builder()
211 .with_batch_exporter(exporter)
212 .with_resource(resource)
213 .build();
214 let tracer = provider.tracer(resolved_name);
215 global::set_tracer_provider(provider.clone());
216 (Some(tracer), Some(provider))
217}
218
219fn install_panic_hook() {
222 let previous = std::panic::take_hook();
223 std::panic::set_hook(Box::new(move |info| {
224 let location = info
225 .location()
226 .map_or_else(|| "unknown".to_owned(), ToString::to_string);
227 tracing::error!(panic = %info, location = %location, "process panicked");
228 previous(info);
229 }));
230}
231
232struct CloudLoggingSeverity<F> {
237 inner: F,
238}
239
240impl<F> CloudLoggingSeverity<F> {
241 const fn wrap(inner: F) -> Self {
242 Self { inner }
243 }
244}
245
246impl<S, N, F> FormatEvent<S, N> for CloudLoggingSeverity<F>
247where
248 S: Subscriber + for<'a> LookupSpan<'a>,
249 N: for<'a> tracing_subscriber::fmt::FormatFields<'a> + 'static,
250 F: FormatEvent<S, N>,
251{
252 fn format_event(
253 &self,
254 ctx: &FmtContext<'_, S, N>,
255 mut writer: Writer<'_>,
256 event: &tracing::Event<'_>,
257 ) -> std::fmt::Result {
258 let mut formatted = String::new();
259 self.inner
260 .format_event(ctx, Writer::new(&mut formatted), event)?;
261 let line = insert_severity(&formatted, *event.metadata().level());
262 writer.write_str(&line)
263 }
264}
265
266fn cloud_logging_severity(level: tracing::Level) -> &'static str {
268 match level.as_str() {
269 "ERROR" => "ERROR",
270 "WARN" => "WARNING",
271 "INFO" => "INFO",
272 _ => "DEBUG",
273 }
274}
275
276fn insert_severity(formatted: &str, level: tracing::Level) -> String {
281 let Some(rest) = formatted.strip_prefix('{') else {
282 return formatted.to_owned();
283 };
284 let severity = cloud_logging_severity(level);
285 let mut out = String::with_capacity(formatted.len() + severity.len() + 14);
286 out.push('{');
287 out.push_str("\"severity\":\"");
288 out.push_str(severity);
289 out.push_str("\",");
290 out.push_str(rest);
291 out
292}
293
294#[cfg(test)]
295mod tests {
296 use super::{BroadcastTee, CloudLoggingSeverity, insert_severity};
297
298 use tracing_subscriber::fmt;
299 use tracing_subscriber::prelude::*;
300
301 fn json_lines(emit: impl FnOnce()) -> Vec<String> {
302 let (tx, mut rx) = tokio::sync::broadcast::channel(8);
303 let layer = fmt::layer()
304 .json()
305 .flatten_event(true)
306 .map_event_format(CloudLoggingSeverity::wrap)
307 .with_writer(BroadcastTee { tx });
308 tracing::subscriber::with_default(tracing_subscriber::registry().with(layer), emit);
309 let mut lines = Vec::new();
310 while let Ok(line) = rx.try_recv() {
311 lines.push(line);
312 }
313 lines
314 }
315
316 fn top_level_severity(line: &str) -> String {
317 let value: serde_json::Value = serde_json::from_str(line).expect("json line");
318 let object = value.as_object().expect("object");
319 let severities: Vec<_> = object
320 .iter()
321 .filter(|(key, _)| *key == "severity")
322 .map(|(_, value)| value.as_str().expect("severity string").to_owned())
323 .collect();
324 assert_eq!(
325 severities.len(),
326 1,
327 "exactly one top-level severity in {line}"
328 );
329 severities.into_iter().next().expect("one severity")
330 }
331
332 #[test]
333 fn json_lines_carry_cloud_logging_severity() {
334 let lines = json_lines(|| {
335 tracing::info!(answer = 7, "ready");
336 tracing::warn!(answer = 7, "careful");
337 tracing::error!(answer = 7, "failed");
338 });
339 assert_eq!(lines.len(), 3);
340 assert_eq!(top_level_severity(&lines[0]), "INFO");
341 assert_eq!(top_level_severity(&lines[1]), "WARNING");
342 assert_eq!(top_level_severity(&lines[2]), "ERROR");
343 for line in &lines {
344 let value: serde_json::Value = serde_json::from_str(line).expect("json line");
345 assert!(value.get("level").is_some(), "keeps level: {line}");
346 assert_eq!(
347 value.get("answer").and_then(serde_json::Value::as_i64),
348 Some(7)
349 );
350 assert!(value.get("message").is_some(), "keeps message: {line}");
351 }
352 }
353
354 #[test]
355 fn field_text_does_not_add_a_second_severity() {
356 let note = "has { and \"severity\":\"FAKE\"";
357 let body = "body with { and \"severity\":\"FAKE\"";
358 let lines = json_lines(|| {
359 tracing::warn!(note, "{body}");
360 });
361 assert_eq!(lines.len(), 1);
362 let value: serde_json::Value = serde_json::from_str(&lines[0]).expect("json line");
363 assert_eq!(top_level_severity(&lines[0]), "WARNING");
364 assert_eq!(
365 value.get("level").and_then(serde_json::Value::as_str),
366 Some("WARN")
367 );
368 let note = value
369 .get("note")
370 .and_then(serde_json::Value::as_str)
371 .expect("note");
372 let message = value
373 .get("message")
374 .and_then(serde_json::Value::as_str)
375 .expect("message");
376 assert!(note.contains('{'), "{note}");
377 assert!(note.contains("\"severity\""), "{note}");
378 assert!(message.contains('{'), "{message}");
379 assert!(message.contains("\"severity\""), "{message}");
380 }
381
382 #[test]
383 fn pretty_lines_do_not_gain_a_severity_field() {
384 let (tx, mut rx) = tokio::sync::broadcast::channel(4);
385 let layer = fmt::layer().with_writer(BroadcastTee { tx });
386 tracing::subscriber::with_default(tracing_subscriber::registry().with(layer), || {
387 tracing::info!(answer = 7, "ready");
388 });
389 let line = rx.try_recv().expect("pretty line");
390 assert!(!line.contains("\"severity\""), "{line}");
391 assert!(line.contains("ready"), "{line}");
392 assert!(!line.starts_with('{'), "{line}");
393 }
394
395 #[test]
396 fn insert_severity_keeps_a_line_that_does_not_start_with_an_object() {
397 assert_eq!(
398 insert_severity("not json", tracing::Level::INFO),
399 "not json"
400 );
401 }
402}