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 }
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")]
203pub 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"))]
259pub 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;