Skip to main content

polyc_runtime/
observability.rs

1//! Logging + traces + panic discipline (PRD §11).
2//!
3//! All logs go to **stderr** (stdout is reserved for program data). Format
4//! is selected by `RUST_LOG_FORMAT` (`json` in prod, pretty otherwise);
5//! verbosity by `RUST_LOG` via `EnvFilter`.
6//!
7//! A JSON line carries a top-level `severity` field. Cloud Logging maps that
8//! field. `TRACE` and `DEBUG` map to `DEBUG`. `WARN` maps to `WARNING`.
9//!
10//! When `OTEL_EXPORTER_OTLP_ENDPOINT` is set, the tracing subscriber also
11//! gains an OpenTelemetry layer that batches spans and ships them via OTLP
12//! gRPC to the configured collector. `OTEL_SERVICE_NAME` overrides the
13//! service name (defaults to the caller's argument). Dev runs without
14//! that env var pay zero overhead — the layer is not constructed.
15//!
16//! Spans propagate through the in-process `tokio` runtime via
17//! `tracing-opentelemetry`; cross-process propagation across the
18//! `connectrpc` boundary uses the W3C `traceparent` header — wire it through
19//! the dialer when running in production.
20
21use 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
38/// How many recent log lines the in-process broadcast ring buffers for a slow
39/// subscriber before it starts dropping the oldest. The local dashboard's log
40/// stream is a live tail, not an audit trail, so lagging is fine.
41const LOG_BROADCAST_CAPACITY: usize = 512;
42
43/// Process-wide sender for the live log stream, set once by [`init`]. A
44/// `OnceLock` (not passed through every call site) so any in-process surface —
45/// the local dashboard's log route in particular — can `subscribe` without the
46/// binary threading a handle down to it.
47static LOG_BROADCAST: OnceLock<broadcast::Sender<String>> = OnceLock::new();
48
49/// Subscribe to the live process-log stream.
50///
51/// Returns a receiver that yields one formatted log line per emitted `tracing`
52/// event (the same text written to stderr), or `None` when [`init`] has not run
53/// yet. A receiver that falls behind the bounded ring drops the oldest lines
54/// rather than stalling the writer.
55#[must_use]
56pub fn subscribe_logs() -> Option<broadcast::Receiver<String>> {
57    LOG_BROADCAST.get().map(broadcast::Sender::subscribe)
58}
59
60/// A `MakeWriter` that tees each formatted log line to stderr and to the live
61/// broadcast. Cloning shares the same sender.
62#[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
78/// One event's worth of formatted bytes: written straight through to stderr and
79/// accumulated so the whole line can be broadcast when the writer drops (the
80/// `fmt` layer makes a fresh writer per event and drops it once written).
81struct 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        // stdout is reserved for program data; all logs go to stderr.
89        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        // A lossy decode keeps a stray non-UTF-8 byte from dropping the line.
105        let line = String::from_utf8_lossy(&self.buf).trim_end().to_owned();
106        if !line.is_empty() {
107            // `send` errs only when there are no receivers; the buffered ring
108            // means that is the sole failure mode, and it is not one to log.
109            let _ = self.tx.send(line);
110        }
111    }
112}
113
114/// Handle returned from [`init`] — keep it alive for the process lifetime;
115/// `drop` flushes traces before exit. `SdkTracerProvider::shutdown` blocks
116/// until exporters drain.
117pub 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            // Best-effort: log to stderr and continue if flush fails.
125            if let Err(err) = provider.shutdown() {
126                eprintln!("opentelemetry shutdown failed: {err}");
127            }
128        }
129    }
130}
131
132/// Initialise the global tracing subscriber + (optional) `OTel` layer.
133///
134/// Also installs the panic hook. Call once, early in `main`. The returned
135/// guard flushes traces on drop — keep it alive for the process lifetime.
136#[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    // Install the W3C `traceparent` propagator unconditionally. The format
145    // half is cheap (a unit struct), and installing it always means the
146    // dialer's `inject_*` calls and the server's `extract_*` calls behave
147    // consistently whether or not OTLP export is enabled — without it, the
148    // global getter returns a no-op propagator and `traceparent` headers
149    // silently disappear, which is the worst-of-both-worlds debug state.
150    // See `crate::propagation` for the matching inject/extract helpers.
151    global::set_text_map_propagator(TraceContextPropagator::new());
152
153    // Tee every formatted log line to stderr and to a process-wide broadcast the
154    // local dashboard tails. The ring is bounded, so a stalled subscriber never
155    // slows logging. Setting the sender always (not only when subscribed) means
156    // the very first lines are already flowing before anyone connects.
157    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
178/// Build the OpenTelemetry tracer (the thing the tracing layer wraps), or
179/// `(None, None)` if no OTLP endpoint is configured. Also installs the
180/// global tracer provider so libraries that talk to
181/// `opentelemetry::global` see it.
182///
183/// Caveat: "the global" means *this* otel version's static. The lock
184/// currently carries a second otel stack (0.31, pinned by
185/// commonware-runtime — see the lockstep comment in the workspace
186/// Cargo.toml) whose own `global` stays the no-op default; only libraries
187/// emitting through the shared `tracing` facade are version-proof.
188fn 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
219/// Emit a single structured `error` event on panic, then defer to the previous
220/// hook (which, under `panic = "abort"`, terminates the process).
221fn 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
232/// Wraps the stock JSON event formatter and inserts a Cloud Logging `severity`.
233///
234/// The inner formatter still writes `level` and every event field. This wrapper
235/// does not parse the line. It inserts the mapped name immediately after `{`.
236struct 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
266/// Maps a tracing level to the Cloud Logging severity name.
267fn 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
276/// Inserts `"severity":"<mapped>",` immediately after the opening `{`.
277///
278/// The stock JSON formatter always starts the line with `{`. A line that does
279/// not is written unchanged, so a log is never dropped.
280fn 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}