Skip to main content

o_sfu_telemetry/
setup.rs

1use std::{error::Error as StdError, fmt};
2
3use anyhow::Result;
4use serde_json::{Map, Number, Value};
5use time::format_description::well_known::Rfc3339;
6use tracing::{
7    Span, Subscriber, field,
8    field::{Field, Visit},
9};
10use tracing_subscriber::{
11    EnvFilter, Registry,
12    fmt::{
13        FmtContext,
14        format::{FormatEvent, FormatFields, JsonFields, Writer},
15        layer as fmt_layer,
16    },
17    layer::SubscriberExt,
18    registry::LookupSpan,
19    util::SubscriberInitExt,
20};
21#[cfg(feature = "otel-tracing")]
22use {
23    opentelemetry::{
24        Context, KeyValue, global,
25        trace::{TraceContextExt, TracerProvider as _},
26    },
27    opentelemetry_otlp::{Protocol, WithExportConfig},
28    opentelemetry_sdk::{
29        Resource,
30        trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
31    },
32    tracing::dispatcher,
33    tracing_opentelemetry::{OpenTelemetrySpanExt, get_otel_context},
34    tracing_subscriber::registry::SpanRef,
35};
36
37use crate::{TelemetryConfig, TelemetryLogFormat, schema};
38
39const DEFAULT_ENV_FILTER: &str = "o_sfu=info,o_sfu_core=info,o_sfu_router=info";
40#[cfg(feature = "otel-tracing")]
41const PRODUCTION_ENVIRONMENT_NAME: &str = "production";
42#[cfg(feature = "otel-tracing")]
43const PRODUCTION_TRACE_SAMPLE_RATIO: f64 = 0.05;
44#[cfg(feature = "otel-tracing")]
45const TRACE_EXPORTER_NAME: &str = "o-sfu.runtime";
46
47#[derive(Debug, Default)]
48pub struct TelemetryHandle {
49    #[cfg(feature = "otel-tracing")]
50    tracer_provider: Option<SdkTracerProvider>,
51}
52
53#[derive(Debug, Clone, PartialEq, Eq)]
54struct TelemetryResourceFields {
55    service_name: String,
56    service_version: String,
57    service_instance_id: String,
58    deployment_environment: String,
59}
60
61#[derive(Debug, Clone)]
62struct RuntimeJsonFormatter {
63    resource: TelemetryResourceFields,
64}
65
66#[derive(Debug, Default)]
67struct JsonEventVisitor {
68    fields: Map<String, Value>,
69}
70
71impl TelemetryHandle {
72    #[cfg(feature = "otel-tracing")]
73    fn with_tracer_provider(tracer_provider: Option<SdkTracerProvider>) -> Self {
74        Self { tracer_provider }
75    }
76
77    #[cfg(not(feature = "otel-tracing"))]
78    const fn disabled() -> Self {
79        Self {}
80    }
81}
82
83#[cfg(feature = "otel-tracing")]
84impl Drop for TelemetryHandle {
85    fn drop(&mut self) {
86        if let Some(tracer_provider) = self.tracer_provider.take()
87            && let Err(_error) = tracer_provider.shutdown()
88        {
89            // Drop cannot surface shutdown failures to a caller, and logging here would
90            // recurse through the subscriber that is being torn down.
91        }
92    }
93}
94
95impl RuntimeJsonFormatter {
96    fn new(resource: TelemetryResourceFields) -> Self {
97        Self { resource }
98    }
99}
100
101impl<S, N> FormatEvent<S, N> for RuntimeJsonFormatter
102where
103    S: Subscriber + for<'lookup> LookupSpan<'lookup>,
104    N: for<'fields> FormatFields<'fields> + 'static,
105{
106    fn format_event(
107        &self,
108        ctx: &FmtContext<'_, S, N>,
109        mut writer: Writer<'_>,
110        event: &tracing::Event<'_>,
111    ) -> fmt::Result {
112        let mut visitor = JsonEventVisitor::default();
113        event.record(&mut visitor);
114        let mut payload = visitor.fields;
115        if payload.get(schema::field::EVENT).and_then(Value::as_str)
116            == Some(schema::event::TRANSPORT_HEALTH_CHANGED)
117        {
118            payload.entry("from".to_owned()).or_insert(Value::Null);
119        }
120        payload.insert(
121            schema::field::TIMESTAMP.to_owned(),
122            Value::String(
123                time::OffsetDateTime::now_utc()
124                    .format(&Rfc3339)
125                    .map_err(|_error| fmt::Error)?,
126            ),
127        );
128        payload.insert(
129            schema::field::LEVEL.to_owned(),
130            Value::String(event.metadata().level().to_string()),
131        );
132        payload.insert(
133            schema::field::TARGET.to_owned(),
134            Value::String(event.metadata().target().to_owned()),
135        );
136        payload
137            .entry(schema::field::EVENT.to_owned())
138            .or_insert_with(|| Value::String(schema::event::RUNTIME_LOG.to_owned()));
139        payload.insert(
140            schema::field::SERVICE_NAME.to_owned(),
141            Value::String(self.resource.service_name.clone()),
142        );
143        payload.insert(
144            schema::field::SERVICE_VERSION.to_owned(),
145            Value::String(self.resource.service_version.clone()),
146        );
147        payload.insert(
148            schema::field::SERVICE_INSTANCE_ID.to_owned(),
149            Value::String(self.resource.service_instance_id.clone()),
150        );
151        payload.insert(
152            schema::field::DEPLOYMENT_ENVIRONMENT.to_owned(),
153            Value::String(self.resource.deployment_environment.clone()),
154        );
155        if let Some(trace_id) = current_trace_id(ctx) {
156            payload.insert(schema::field::TRACE_ID.to_owned(), Value::String(trace_id));
157        }
158        let encoded = serde_json::to_string(&payload).map_err(|_error| fmt::Error)?;
159        writer.write_str(encoded.as_str())?;
160        writeln!(writer)
161    }
162}
163
164impl Visit for JsonEventVisitor {
165    fn record_bool(&mut self, field: &Field, value: bool) {
166        self.fields
167            .insert(field.name().to_owned(), Value::Bool(value));
168    }
169
170    fn record_i64(&mut self, field: &Field, value: i64) {
171        self.fields
172            .insert(field.name().to_owned(), Value::Number(Number::from(value)));
173    }
174
175    fn record_u64(&mut self, field: &Field, value: u64) {
176        self.fields
177            .insert(field.name().to_owned(), Value::Number(Number::from(value)));
178    }
179
180    fn record_str(&mut self, field: &Field, value: &str) {
181        self.fields
182            .insert(field.name().to_owned(), Value::String(value.to_owned()));
183    }
184
185    fn record_error(&mut self, field: &Field, value: &(dyn StdError + 'static)) {
186        self.fields
187            .insert(field.name().to_owned(), Value::String(value.to_string()));
188    }
189
190    fn record_f64(&mut self, field: &Field, value: f64) {
191        let json_value =
192            Number::from_f64(value).map_or_else(|| Value::String(value.to_string()), Value::Number);
193        self.fields.insert(field.name().to_owned(), json_value);
194    }
195
196    fn record_debug(&mut self, field: &Field, value: &dyn fmt::Debug) {
197        self.fields
198            .insert(field.name().to_owned(), Value::String(format!("{value:?}")));
199    }
200}
201
202#[cfg(feature = "otel-tracing")]
203/// # Errors
204///
205/// Returns an error when subscriber initialization or OTLP exporter
206/// construction fails.
207pub fn init_tracing(config: &TelemetryConfig, process_id: u32) -> Result<TelemetryHandle> {
208    let env_filter = default_env_filter();
209    let resource = telemetry_resource_fields(config, process_id);
210    let tracer_provider = build_tracer_provider(config, &resource)?;
211    let tracer = tracer_provider
212        .as_ref()
213        .map(|provider| provider.tracer(TRACE_EXPORTER_NAME));
214    match config.log_format {
215        TelemetryLogFormat::Compact => Registry::default()
216            .with(env_filter)
217            .with(fmt_layer().with_target(false).compact())
218            .with(
219                tracer
220                    .as_ref()
221                    .map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer.clone())),
222            )
223            .try_init()?,
224        TelemetryLogFormat::Json => Registry::default()
225            .with(env_filter)
226            .with(
227                fmt_layer()
228                    .fmt_fields(JsonFields::new())
229                    .event_format(RuntimeJsonFormatter::new(resource.clone()))
230                    .with_ansi(false),
231            )
232            .with(
233                tracer
234                    .as_ref()
235                    .map(|tracer| tracing_opentelemetry::layer().with_tracer(tracer.clone())),
236            )
237            .try_init()?,
238    }
239    tracing::info!(
240        event = schema::event::RUNTIME_TELEMETRY_INITIALIZED,
241        service_name = resource.service_name.as_str(),
242        service_version = resource.service_version.as_str(),
243        deployment_environment = resource.deployment_environment.as_str(),
244        service_instance_id = resource.service_instance_id.as_str(),
245        log_format = config.log_format.as_str(),
246        common_fields = ?schema::COMMON_FIELD_NAMES,
247        correlation_fields = ?schema::CORRELATION_FIELD_NAMES,
248        trace_export_otlp_endpoint = config
249            .trace_export
250            .otlp_endpoint
251            .as_deref()
252            .unwrap_or("disabled"),
253        "initialized runtime telemetry"
254    );
255    Ok(TelemetryHandle::with_tracer_provider(tracer_provider))
256}
257
258#[cfg(not(feature = "otel-tracing"))]
259/// # Errors
260///
261/// Returns an error when subscriber initialization fails.
262pub fn init_tracing(config: &TelemetryConfig, process_id: u32) -> Result<TelemetryHandle> {
263    let env_filter = default_env_filter();
264    let resource = telemetry_resource_fields(config, process_id);
265    match config.log_format {
266        TelemetryLogFormat::Compact => Registry::default()
267            .with(env_filter)
268            .with(fmt_layer().with_target(false).compact())
269            .try_init()?,
270        TelemetryLogFormat::Json => Registry::default()
271            .with(env_filter)
272            .with(
273                fmt_layer()
274                    .fmt_fields(JsonFields::new())
275                    .event_format(RuntimeJsonFormatter::new(resource.clone()))
276                    .with_ansi(false),
277            )
278            .try_init()?,
279    }
280    tracing::info!(
281        event = schema::event::RUNTIME_TELEMETRY_INITIALIZED,
282        service_name = resource.service_name.as_str(),
283        service_version = resource.service_version.as_str(),
284        deployment_environment = resource.deployment_environment.as_str(),
285        service_instance_id = resource.service_instance_id.as_str(),
286        log_format = config.log_format.as_str(),
287        common_fields = ?schema::COMMON_FIELD_NAMES,
288        correlation_fields = ?schema::CORRELATION_FIELD_NAMES,
289        trace_export_otlp_endpoint = "feature_disabled",
290        "initialized runtime telemetry"
291    );
292    Ok(TelemetryHandle::disabled())
293}
294
295#[must_use]
296pub fn http_request_span(route: &'static str) -> Span {
297    activated_span(tracing::info_span!(
298        "http.request",
299        "otel.kind" = "server",
300        route,
301        room_id = field::Empty,
302        user_id = field::Empty,
303        connection_id = field::Empty,
304        remote_address = field::Empty
305    ))
306}
307
308#[must_use]
309pub fn ws_upgrade_span() -> Span {
310    activated_span(tracing::info_span!(
311        "ws.upgrade",
312        room_id = field::Empty,
313        user_id = field::Empty,
314        connection_id = field::Empty,
315        remote_address = field::Empty
316    ))
317}
318
319#[must_use]
320pub fn ws_handshake_span() -> Span {
321    activated_span(tracing::info_span!(
322        "ws.handshake",
323        room_id = field::Empty,
324        user_id = field::Empty,
325        connection_id = field::Empty,
326        remote_address = field::Empty
327    ))
328}
329
330#[cfg(feature = "otel-tracing")]
331#[must_use]
332pub fn activated_span(span: Span) -> Span {
333    let _span_context = span.context();
334    span
335}
336
337#[cfg(not(feature = "otel-tracing"))]
338#[must_use]
339pub fn activated_span(span: Span) -> Span {
340    span
341}
342
343#[cfg(feature = "otel-tracing")]
344fn current_trace_id<S, N>(ctx: &FmtContext<'_, S, N>) -> Option<String>
345where
346    S: Subscriber + for<'lookup> LookupSpan<'lookup>,
347    N: for<'fields> FormatFields<'fields> + 'static,
348{
349    dispatcher::get_default(|dispatch| {
350        ctx.event_scope()
351            .and_then(|mut scope| scope.find_map(|span| trace_id_for_span(&span, dispatch)))
352            .or_else(|| {
353                ctx.lookup_current()
354                    .and_then(|span| trace_id_for_span(&span, dispatch))
355            })
356    })
357}
358
359#[cfg(not(feature = "otel-tracing"))]
360fn current_trace_id<S, N>(_ctx: &FmtContext<'_, S, N>) -> Option<String>
361where
362    S: Subscriber + for<'lookup> LookupSpan<'lookup>,
363    N: for<'fields> FormatFields<'fields> + 'static,
364{
365    None
366}
367
368#[cfg(feature = "otel-tracing")]
369fn trace_id_for_span<S>(span: &SpanRef<'_, S>, dispatch: &tracing::Dispatch) -> Option<String>
370where
371    S: for<'lookup> LookupSpan<'lookup>,
372{
373    get_otel_context(&span.id(), dispatch)
374        .and_then(|context| trace_id_from_context(&context))
375        .or_else(|| trace_id_from_context(&Context::current()))
376}
377
378#[cfg(feature = "otel-tracing")]
379fn trace_id_from_context(context: &Context) -> Option<String> {
380    let span_context = context.span().span_context().clone();
381    span_context
382        .is_valid()
383        .then(|| span_context.trace_id().to_string())
384}
385
386fn default_env_filter() -> EnvFilter {
387    EnvFilter::try_from_default_env().unwrap_or_else(|_error| EnvFilter::new(DEFAULT_ENV_FILTER))
388}
389
390fn telemetry_resource_fields(config: &TelemetryConfig, process_id: u32) -> TelemetryResourceFields {
391    TelemetryResourceFields {
392        service_name: config.resource.service_name.clone(),
393        service_version: env!("CARGO_PKG_VERSION").to_owned(),
394        service_instance_id: config.resource.resolved_instance_id(process_id),
395        deployment_environment: config.resource.deployment_environment.clone(),
396    }
397}
398
399#[cfg(feature = "otel-tracing")]
400fn build_tracer_provider(
401    config: &TelemetryConfig,
402    resource: &TelemetryResourceFields,
403) -> Result<Option<SdkTracerProvider>> {
404    let Some(endpoint) = config.trace_export.otlp_endpoint.as_deref() else {
405        return Ok(None);
406    };
407    let exporter = opentelemetry_otlp::SpanExporter::builder()
408        .with_http()
409        .with_protocol(Protocol::HttpBinary)
410        .with_endpoint(normalize_trace_export_endpoint(endpoint))
411        .build()?;
412    let tracer_provider = SdkTracerProvider::builder()
413        .with_batch_exporter(exporter)
414        .with_sampler(default_trace_sampler(
415            resource.deployment_environment.as_str(),
416        ))
417        .with_id_generator(RandomIdGenerator::default())
418        .with_resource(
419            Resource::builder_empty()
420                .with_attributes([
421                    KeyValue::new(schema::field::SERVICE_NAME, resource.service_name.clone()),
422                    KeyValue::new(
423                        schema::field::SERVICE_VERSION,
424                        resource.service_version.clone(),
425                    ),
426                    KeyValue::new(
427                        schema::field::SERVICE_INSTANCE_ID,
428                        resource.service_instance_id.clone(),
429                    ),
430                    KeyValue::new(
431                        schema::field::DEPLOYMENT_ENVIRONMENT,
432                        resource.deployment_environment.clone(),
433                    ),
434                ])
435                .build(),
436        )
437        .build();
438    global::set_tracer_provider(tracer_provider.clone());
439    Ok(Some(tracer_provider))
440}
441
442#[cfg(feature = "otel-tracing")]
443fn default_trace_sampler(deployment_environment: &str) -> Sampler {
444    if deployment_environment == PRODUCTION_ENVIRONMENT_NAME {
445        Sampler::ParentBased(Box::new(Sampler::TraceIdRatioBased(
446            PRODUCTION_TRACE_SAMPLE_RATIO,
447        )))
448    } else {
449        Sampler::AlwaysOn
450    }
451}
452
453#[cfg(feature = "otel-tracing")]
454fn normalize_trace_export_endpoint(endpoint: &str) -> String {
455    if endpoint.ends_with("/v1/traces") {
456        endpoint.to_owned()
457    } else {
458        format!("{}/v1/traces", endpoint.trim_end_matches('/'))
459    }
460}
461
462#[cfg(test)]
463#[path = "TESTS/setup.rs"]
464mod tests;