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 }
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 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 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
198pub 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;