Skip to main content

o_sfu_telemetry/
setup.rs

1use std::fmt;
2
3use anyhow::Result;
4use serde::{
5    Serialize, Serializer,
6    ser::{Error as _, SerializeSeq},
7};
8use serde_json::value::RawValue;
9use time::format_description::well_known::Rfc3339;
10use tracing::{Span, Subscriber, field};
11use tracing_serde::fields::AsMap;
12use tracing_subscriber::{
13    EnvFilter, Registry,
14    fmt::{
15        FmtContext, FormattedFields,
16        format::{FormatEvent, JsonFields, Writer},
17        layer as fmt_layer,
18    },
19    layer::SubscriberExt,
20    registry::{LookupSpan, Scope},
21    util::SubscriberInitExt,
22};
23#[cfg(feature = "otel-tracing")]
24use {
25    opentelemetry::{
26        KeyValue, global,
27        trace::{TraceContextExt, TracerProvider as _},
28    },
29    opentelemetry_otlp::{Protocol, WithExportConfig},
30    opentelemetry_sdk::{
31        Resource,
32        trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
33    },
34    std::sync::{Arc, OnceLock},
35    tracing::dispatcher::WeakDispatch,
36    tracing_opentelemetry::{OpenTelemetrySpanExt, get_otel_context},
37    tracing_subscriber::Layer,
38};
39
40use crate::{TelemetryConfig, TelemetryLogFormat, schema};
41
42const DEFAULT_ENV_FILTER: &str = "o_sfu=info,o_sfu_core=info,o_sfu_router=info";
43#[cfg(feature = "otel-tracing")]
44const PRODUCTION_ENVIRONMENT_NAME: &str = "production";
45#[cfg(feature = "otel-tracing")]
46const PRODUCTION_TRACE_SAMPLE_RATIO: f64 = 0.05;
47#[cfg(feature = "otel-tracing")]
48const TRACE_EXPORTER_NAME: &str = "o-sfu.runtime";
49
50#[derive(Debug, Default)]
51pub struct TelemetryHandle {
52    #[cfg(feature = "otel-tracing")]
53    tracer_provider: Option<SdkTracerProvider>,
54}
55
56#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
57struct TelemetryResourceFields {
58    #[serde(rename = "service.name")]
59    service_name: String,
60    #[serde(rename = "service.version")]
61    service_version: String,
62    #[serde(rename = "service.instance.id")]
63    service_instance_id: String,
64    #[serde(rename = "deployment.environment")]
65    deployment_environment: String,
66}
67
68#[derive(Debug, Clone)]
69struct RuntimeJsonFormatter {
70    resource: TelemetryResourceFields,
71    #[cfg(feature = "otel-tracing")]
72    dispatch: Arc<OnceLock<WeakDispatch>>,
73}
74
75#[derive(Serialize)]
76struct JsonEvent<'a, E, S> {
77    timestamp: String,
78    level: &'static str,
79    target: &'static str,
80    #[serde(flatten)]
81    resource: &'a TelemetryResourceFields,
82    #[serde(skip_serializing_if = "Option::is_none")]
83    trace_id: Option<String>,
84    fields: E,
85    spans: S,
86}
87
88#[derive(Serialize)]
89struct JsonSpan<'a> {
90    name: &'static str,
91    fields: &'a RawValue,
92}
93
94struct JsonSpans<'a, 'context, S>(&'a FmtContext<'context, S, JsonFields>);
95
96#[cfg(feature = "otel-tracing")]
97impl Drop for TelemetryHandle {
98    fn drop(&mut self) {
99        if let Some(tracer_provider) = self.tracer_provider.take()
100            && let Err(_error) = tracer_provider.shutdown()
101        {
102            // Drop cannot surface shutdown failures to a caller, and logging here would
103            // recurse through the subscriber that is being torn down.
104        }
105    }
106}
107
108impl RuntimeJsonFormatter {
109    fn new(resource: TelemetryResourceFields) -> Self {
110        Self {
111            resource,
112            #[cfg(feature = "otel-tracing")]
113            dispatch: Arc::default(),
114        }
115    }
116
117    #[cfg(feature = "otel-tracing")]
118    fn trace_id<S>(&self, ctx: &FmtContext<'_, S, JsonFields>) -> Option<String>
119    where
120        S: Subscriber + for<'lookup> LookupSpan<'lookup>,
121    {
122        let dispatch = self.dispatch.get()?.upgrade()?;
123        let parent = ctx.parent_span()?;
124        let context = get_otel_context(&parent.id(), &dispatch)?;
125        let span = context.span();
126        let span_context = span.span_context();
127        span_context
128            .is_valid()
129            .then(|| span_context.trace_id().to_string())
130    }
131
132    #[cfg(not(feature = "otel-tracing"))]
133    fn trace_id<S>(&self, _ctx: &FmtContext<'_, S, JsonFields>) -> Option<String>
134    where
135        S: Subscriber + for<'lookup> LookupSpan<'lookup>,
136    {
137        None
138    }
139}
140
141#[cfg(feature = "otel-tracing")]
142impl<S: Subscriber> Layer<S> for RuntimeJsonFormatter {
143    fn on_register_dispatch(&self, dispatch: &tracing::Dispatch) {
144        // Nested get_default calls hide the subscriber during event dispatch.
145        // The formatter clone in this layer shares a weak reference with the fmt layer.
146        let _ = self.dispatch.set(dispatch.downgrade());
147    }
148}
149
150impl<S> FormatEvent<S, JsonFields> for RuntimeJsonFormatter
151where
152    S: Subscriber + for<'lookup> LookupSpan<'lookup>,
153{
154    fn format_event(
155        &self,
156        ctx: &FmtContext<'_, S, JsonFields>,
157        mut writer: Writer<'_>,
158        event: &tracing::Event<'_>,
159    ) -> fmt::Result {
160        let payload = JsonEvent {
161            timestamp: time::OffsetDateTime::now_utc()
162                .format(&Rfc3339)
163                .map_err(|_error| fmt::Error)?,
164            level: event.metadata().level().as_str(),
165            target: event.metadata().target(),
166            resource: &self.resource,
167            trace_id: self.trace_id(ctx),
168            fields: event.field_map(),
169            spans: JsonSpans(ctx),
170        };
171        let encoded = serde_json::to_string(&payload).map_err(|_error| fmt::Error)?;
172        writeln!(writer, "{encoded}")
173    }
174}
175
176impl<S> Serialize for JsonSpans<'_, '_, S>
177where
178    S: Subscriber + for<'lookup> LookupSpan<'lookup>,
179{
180    fn serialize<T: Serializer>(&self, serializer: T) -> Result<T::Ok, T::Error> {
181        let mut sequence = serializer.serialize_seq(None)?;
182        // Explicit-parent events can belong to a different scope than the entered span.
183        for span in self.0.event_scope().into_iter().flat_map(Scope::from_root) {
184            let extensions = span.extensions();
185            if let Some(fields) = extensions.get::<FormattedFields<JsonFields>>() {
186                let fields = serde_json::from_str::<&RawValue>(fields.fields.as_str())
187                    .map_err(T::Error::custom)?;
188                sequence.serialize_element(&JsonSpan {
189                    name: span.name(),
190                    fields,
191                })?;
192            }
193        }
194        sequence.end()
195    }
196}
197
198/// Installs the configured tracing subscriber and retains its exporter until
199/// the returned handle is dropped.
200///
201/// # Errors
202///
203/// Returns an [`anyhow::Error`] when subscriber initialization fails or when the
204/// `otel-tracing` feature is enabled and OTLP exporter construction fails.
205pub fn init_tracing(config: &TelemetryConfig, process_id: u32) -> Result<TelemetryHandle> {
206    let env_filter = default_env_filter();
207    let resource = telemetry_resource_fields(config, process_id);
208    #[cfg(feature = "otel-tracing")]
209    let tracer_provider = build_tracer_provider(config, &resource)?;
210    #[cfg(feature = "otel-tracing")]
211    let tracer = tracer_provider
212        .as_ref()
213        .map(|provider| provider.tracer(TRACE_EXPORTER_NAME));
214    match config.log_format {
215        TelemetryLogFormat::Compact => {
216            let subscriber = Registry::default()
217                .with(env_filter)
218                .with(fmt_layer().with_target(false).compact());
219            #[cfg(feature = "otel-tracing")]
220            let subscriber = subscriber.with(
221                tracer
222                    .as_ref()
223                    .map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer.clone())),
224            );
225            subscriber.try_init()?;
226        }
227        TelemetryLogFormat::Json => {
228            let formatter = RuntimeJsonFormatter::new(resource.clone());
229            let subscriber = Registry::default().with(env_filter);
230            #[cfg(feature = "otel-tracing")]
231            let subscriber = subscriber.with(formatter.clone());
232            let subscriber = subscriber.with(
233                fmt_layer()
234                    .fmt_fields(JsonFields::new())
235                    .event_format(formatter)
236                    .with_ansi(false),
237            );
238            #[cfg(feature = "otel-tracing")]
239            let subscriber = subscriber.with(
240                tracer
241                    .as_ref()
242                    .map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer.clone())),
243            );
244            subscriber.try_init()?;
245        }
246    }
247    #[cfg(feature = "otel-tracing")]
248    let trace_export_otlp_endpoint = config
249        .trace_export
250        .otlp_endpoint
251        .as_deref()
252        .unwrap_or("disabled");
253    #[cfg(not(feature = "otel-tracing"))]
254    let trace_export_otlp_endpoint = "feature_disabled";
255    tracing::info!(
256        event = schema::event::RUNTIME_TELEMETRY_INITIALIZED,
257        service_name = resource.service_name.as_str(),
258        service_version = resource.service_version.as_str(),
259        deployment_environment = resource.deployment_environment.as_str(),
260        service_instance_id = resource.service_instance_id.as_str(),
261        log_format = config.log_format.as_str(),
262        common_fields = ?schema::COMMON_FIELD_NAMES,
263        correlation_fields = ?schema::CORRELATION_FIELD_NAMES,
264        trace_export_otlp_endpoint,
265        "initialized runtime telemetry"
266    );
267    Ok(TelemetryHandle {
268        #[cfg(feature = "otel-tracing")]
269        tracer_provider,
270    })
271}
272
273#[must_use]
274pub fn http_request_span(route: &'static str) -> Span {
275    activated_span(tracing::info_span!(
276        "http.request",
277        "otel.kind" = "server",
278        route,
279        room_id = field::Empty,
280        user_id = field::Empty,
281        connection_id = field::Empty,
282        remote_address = field::Empty
283    ))
284}
285
286#[must_use]
287pub fn ws_upgrade_span() -> Span {
288    activated_span(tracing::info_span!(
289        "ws.upgrade",
290        room_id = field::Empty,
291        user_id = field::Empty,
292        connection_id = field::Empty,
293        remote_address = field::Empty
294    ))
295}
296
297#[must_use]
298pub fn ws_handshake_span() -> Span {
299    activated_span(tracing::info_span!(
300        "ws.handshake",
301        room_id = field::Empty,
302        user_id = field::Empty,
303        connection_id = field::Empty,
304        remote_address = field::Empty
305    ))
306}
307
308#[cfg(feature = "otel-tracing")]
309#[must_use]
310pub fn activated_span(span: Span) -> Span {
311    let _span_context = span.context();
312    span
313}
314
315#[cfg(not(feature = "otel-tracing"))]
316#[must_use]
317pub fn activated_span(span: Span) -> Span {
318    span
319}
320
321fn default_env_filter() -> EnvFilter {
322    EnvFilter::try_from_default_env().unwrap_or_else(|_error| EnvFilter::new(DEFAULT_ENV_FILTER))
323}
324
325fn telemetry_resource_fields(config: &TelemetryConfig, process_id: u32) -> TelemetryResourceFields {
326    TelemetryResourceFields {
327        service_name: config.resource.service_name.clone(),
328        service_version: env!("CARGO_PKG_VERSION").to_owned(),
329        service_instance_id: config.resource.resolved_instance_id(process_id),
330        deployment_environment: config.resource.deployment_environment.clone(),
331    }
332}
333
334#[cfg(feature = "otel-tracing")]
335fn build_tracer_provider(
336    config: &TelemetryConfig,
337    resource: &TelemetryResourceFields,
338) -> Result<Option<SdkTracerProvider>> {
339    let Some(endpoint) = config.trace_export.otlp_endpoint.as_deref() else {
340        return Ok(None);
341    };
342    let exporter = opentelemetry_otlp::SpanExporter::builder()
343        .with_http()
344        .with_protocol(Protocol::HttpBinary)
345        .with_endpoint(normalize_trace_export_endpoint(endpoint))
346        .build()?;
347    let tracer_provider = SdkTracerProvider::builder()
348        .with_batch_exporter(exporter)
349        .with_sampler(default_trace_sampler(
350            resource.deployment_environment.as_str(),
351        ))
352        .with_id_generator(RandomIdGenerator::default())
353        .with_resource(
354            Resource::builder_empty()
355                .with_attributes([
356                    KeyValue::new(schema::field::SERVICE_NAME, resource.service_name.clone()),
357                    KeyValue::new(
358                        schema::field::SERVICE_VERSION,
359                        resource.service_version.clone(),
360                    ),
361                    KeyValue::new(
362                        schema::field::SERVICE_INSTANCE_ID,
363                        resource.service_instance_id.clone(),
364                    ),
365                    KeyValue::new(
366                        schema::field::DEPLOYMENT_ENVIRONMENT,
367                        resource.deployment_environment.clone(),
368                    ),
369                ])
370                .build(),
371        )
372        .build();
373    global::set_tracer_provider(tracer_provider.clone());
374    Ok(Some(tracer_provider))
375}
376
377#[cfg(feature = "otel-tracing")]
378fn default_trace_sampler(deployment_environment: &str) -> Sampler {
379    if deployment_environment == PRODUCTION_ENVIRONMENT_NAME {
380        Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
381            PRODUCTION_TRACE_SAMPLE_RATIO,
382        )))
383    } else {
384        Sampler::AlwaysOn
385    }
386}
387
388#[cfg(feature = "otel-tracing")]
389fn normalize_trace_export_endpoint(endpoint: &str) -> String {
390    if endpoint.ends_with("/v1/traces") {
391        endpoint.to_owned()
392    } else {
393        format!("{}/v1/traces", endpoint.trim_end_matches('/'))
394    }
395}
396
397#[cfg(test)]
398#[path = "TESTS/setup.rs"]
399mod tests;