Skip to main content

o_sfu_telemetry/metrics/
catalog.rs

1//! Defines process-local metric storage and recording APIs.
2
3use std::{
4    sync::Arc,
5    time::{Duration, Instant},
6};
7
8use o_sfu_model::WebSocketCloseCode;
9
10use super::{
11    counter::{
12        Counter, CounterFamily, Histogram, HistogramFamily, UpDownCounter, UpDownCounterFamily,
13    },
14    labels::{
15        BudgetSolverOutcome, ControlPlaneDurationBucket, HttpDisconnectResponseStatus,
16        HttpRoomResponseStatus, HttpRoute, MediaQualityLossDirection, MediaQualityRttBucket,
17        MediaQualitySample, RecordingActionOutcome, RtpRelayDropKind, SourceSelectionKind,
18        TransportHealthState, TransportHealthTransition, TransportIceState,
19        TransportUserLifetimeBucket, WsBusClientFrameKind, WsBusDirection, WsBusFailureKind,
20        WsConnectionStage, WsSessionLoopExitReason, WsStartupFailureKind,
21    },
22    rtc::{RtcMetrics, RtcMetricsRecorder},
23    rtp::{RtpMetrics, RtpMetricsRecorder},
24};
25
26#[derive(Debug, Default)]
27pub struct RuntimeMetrics {
28    pub(super) http_requests: CounterFamily<HttpRoute>,
29    pub(super) http_room_responses: CounterFamily<HttpRoomResponseStatus>,
30    pub(super) http_disconnect_responses: CounterFamily<HttpDisconnectResponseStatus>,
31    pub(super) http_inflight_requests: UpDownCounterFamily<HttpRoute>,
32    pub(super) http_request_duration: HistogramFamily<HttpRoute, ControlPlaneDurationBucket>,
33    pub(super) ws_connections: CounterFamily<WsConnectionStage>,
34    pub(super) ws_handshake_rejections: CounterFamily<WebSocketCloseCode>,
35    pub(super) ws_handshake_rejections_other: Counter,
36    pub(super) ws_startup_failures: CounterFamily<WsStartupFailureKind>,
37    pub(super) ws_user_loops_started: Counter,
38    pub(super) ws_user_loop_exits: CounterFamily<WsSessionLoopExitReason>,
39    pub(super) ws_bus_batches: CounterFamily<WsBusDirection>,
40    pub(super) ws_bus_envelopes: CounterFamily<WsBusDirection>,
41    pub(super) ws_bus_parse_failures: Counter,
42    pub(super) ws_bus_failures: CounterFamily<WsBusFailureKind>,
43    pub(super) ws_bus_client_frames: CounterFamily<WsBusClientFrameKind>,
44    pub(super) ws_outbound_queued_messages: UpDownCounter,
45    pub(super) ws_outbound_queue_overflows: Counter,
46    pub(super) ws_handshake_duration: Histogram<ControlPlaneDurationBucket>,
47    pub(super) ws_auth_duration: Histogram<ControlPlaneDurationBucket>,
48    pub(super) ws_user_initialize_duration: Histogram<ControlPlaneDurationBucket>,
49    pub(super) active_transport_users: UpDownCounter,
50    pub(super) transport_health_users: UpDownCounterFamily<TransportHealthState>,
51    pub(super) recording_actions: CounterFamily<RecordingActionOutcome>,
52    pub(super) recording_captured_packets: Counter,
53    pub(super) recording_captured_streams: Counter,
54    pub(super) rtp_metrics: RtpMetrics,
55    pub(super) rtc_metrics: RtcMetrics,
56    pub(super) rtp_relay_overload_drops: CounterFamily<RtpRelayDropKind>,
57    pub(super) transport_health_transitions: CounterFamily<TransportHealthTransition>,
58    pub(super) transport_ice_state_changes: CounterFamily<TransportIceState>,
59    pub(super) transport_dtls_connected: Counter,
60    pub(super) transport_user_lifetime_buckets: CounterFamily<TransportUserLifetimeBucket>,
61    pub(super) transport_user_lifetime_count: Counter,
62    pub(super) transport_user_lifetime_sum_micros: Counter,
63    pub(super) media_quality_samples: CounterFamily<MediaQualitySample>,
64    pub(super) media_quality_rtt: HistogramFamily<MediaQualitySample, MediaQualityRttBucket>,
65    pub(super) media_quality_loss_ppm_observed: CounterFamily<MediaQualityLossDirection>,
66    pub(super) media_quality_loss_observations: CounterFamily<MediaQualityLossDirection>,
67    pub(super) media_quality_bwe_bps_observed: Counter,
68    pub(super) media_quality_bwe_observations: Counter,
69    pub(super) media_quality_jitter_rtp_timestamp_units_observed: Counter,
70    pub(super) media_quality_jitter_observations: Counter,
71    pub(super) transport_cleanup_failures: Counter,
72    pub(super) source_selection_updates: CounterFamily<SourceSelectionKind>,
73    pub(super) budget_solver_outcomes: CounterFamily<BudgetSolverOutcome>,
74}
75
76struct MetricGuard<'a, F>
77where
78    F: Fn(&RuntimeMetrics, Duration),
79{
80    metrics: &'a RuntimeMetrics,
81    started_at: Instant,
82    finish: F,
83}
84
85impl<F> Drop for MetricGuard<'_, F>
86where
87    F: Fn(&RuntimeMetrics, Duration),
88{
89    fn drop(&mut self) {
90        (self.finish)(self.metrics, self.started_at.elapsed());
91    }
92}
93
94impl RuntimeMetrics {
95    /// counts one HTTP request then records duration and releases inflight state on drop
96    #[must_use = "keep the guard until the HTTP request finishes"]
97    pub fn track_http_request(&self, route: HttpRoute) -> impl Drop + '_ {
98        self.http_requests.increment(route);
99        self.http_inflight_requests.add(route, 1);
100        self.track(move |metrics, duration| {
101            metrics.http_inflight_requests.add(route, -1);
102            metrics.http_request_duration.observe(route, duration);
103        })
104    }
105
106    pub fn record_http_room_success(&self) {
107        self.http_room_responses
108            .increment(HttpRoomResponseStatus::Success);
109    }
110
111    pub fn record_http_room_unauthorized(&self) {
112        self.http_room_responses
113            .increment(HttpRoomResponseStatus::Unauthorized);
114    }
115
116    pub fn record_http_room_forbidden(&self) {
117        self.http_room_responses
118            .increment(HttpRoomResponseStatus::Forbidden);
119    }
120
121    pub fn record_http_room_bad_request(&self) {
122        self.http_room_responses
123            .increment(HttpRoomResponseStatus::BadRequest);
124    }
125
126    pub fn record_http_room_conflict(&self) {
127        self.http_room_responses
128            .increment(HttpRoomResponseStatus::Conflict);
129    }
130
131    pub fn record_http_disconnect_success(&self) {
132        self.http_disconnect_responses
133            .increment(HttpDisconnectResponseStatus::Success);
134    }
135
136    pub fn record_http_disconnect_bad_request(&self) {
137        self.http_disconnect_responses
138            .increment(HttpDisconnectResponseStatus::BadRequest);
139    }
140
141    pub fn record_http_disconnect_unprocessable_entity(&self) {
142        self.http_disconnect_responses
143            .increment(HttpDisconnectResponseStatus::UnprocessableEntity);
144    }
145
146    pub fn record_ws_connection_accepted(&self) {
147        self.ws_connections.increment(WsConnectionStage::Accepted);
148    }
149
150    pub fn record_ws_handshake_credentials_received(&self) {
151        self.ws_connections
152            .increment(WsConnectionStage::CredentialsReceived);
153    }
154
155    pub fn record_ws_handshake_rejection(&self, close_code: Option<WebSocketCloseCode>) {
156        match close_code {
157            Some(
158                close_code @ (WebSocketCloseCode::AuthTimeout
159                | WebSocketCloseCode::AuthFailed
160                | WebSocketCloseCode::ProtocolError
161                | WebSocketCloseCode::RoomFull),
162            ) => self.ws_handshake_rejections.increment(close_code),
163            Some(
164                WebSocketCloseCode::Error
165                | WebSocketCloseCode::Clean
166                | WebSocketCloseCode::Leaving
167                | WebSocketCloseCode::Kicked,
168            )
169            | None => self.ws_handshake_rejections_other.increment(),
170        }
171    }
172
173    pub fn record_ws_user_joined(&self) {
174        self.ws_connections.increment(WsConnectionStage::Joined);
175    }
176
177    pub fn record_ws_startup_send_failure(&self) {
178        self.ws_startup_failures
179            .increment(WsStartupFailureKind::StartupSend);
180    }
181
182    pub fn record_ws_user_initialize_failure(&self) {
183        self.ws_startup_failures
184            .increment(WsStartupFailureKind::SessionInitialize);
185    }
186
187    pub fn record_ws_user_loop_started(&self) {
188        self.ws_user_loops_started.increment();
189    }
190
191    pub fn record_ws_user_loop_exit(&self, reason: WsSessionLoopExitReason) {
192        self.ws_user_loop_exits.increment(reason);
193    }
194
195    pub fn record_ws_bus_batch_received(&self, envelope_count: usize) {
196        self.ws_bus_batches.increment(WsBusDirection::Received);
197        self.ws_bus_envelopes
198            .add(WsBusDirection::Received, envelope_count);
199    }
200
201    pub fn record_ws_bus_invalid_input_failure(&self) {
202        self.ws_bus_parse_failures.increment();
203        self.ws_bus_failures
204            .increment(WsBusFailureKind::InvalidInput);
205    }
206
207    pub fn record_ws_bus_unsupported_feature_failure(&self) {
208        self.ws_bus_parse_failures.increment();
209        self.ws_bus_failures
210            .increment(WsBusFailureKind::UnsupportedFeature);
211    }
212
213    pub fn record_ws_bus_client_request(&self) {
214        self.ws_bus_client_frames
215            .increment(WsBusClientFrameKind::Request);
216    }
217
218    pub fn record_ws_bus_client_message(&self) {
219        self.ws_bus_client_frames
220            .increment(WsBusClientFrameKind::Message);
221    }
222
223    pub fn record_ws_bus_batch_sent(&self, envelope_count: usize) {
224        self.record_ws_bus_batches_sent(1, envelope_count);
225    }
226
227    pub fn record_ws_bus_batches_sent(&self, batch_count: usize, envelope_count: usize) {
228        if batch_count == 0 {
229            return;
230        }
231        self.ws_bus_batches.add(WsBusDirection::Sent, batch_count);
232        self.ws_bus_envelopes
233            .add(WsBusDirection::Sent, envelope_count);
234    }
235
236    pub fn record_ws_bus_send_failure(&self) {
237        self.ws_bus_failures.increment(WsBusFailureKind::Send);
238    }
239
240    pub fn add_ws_outbound_queued_messages(&self, delta: i64) {
241        self.ws_outbound_queued_messages.add(delta);
242    }
243
244    pub fn record_ws_outbound_queue_overflow(&self) {
245        self.ws_outbound_queue_overflows.increment();
246    }
247
248    /// records handshake duration when the guard is dropped
249    #[must_use = "keep the guard until the WebSocket handshake finishes"]
250    pub fn track_ws_handshake(&self) -> impl Drop + '_ {
251        self.track(|metrics, duration| metrics.ws_handshake_duration.observe(duration))
252    }
253
254    /// records authentication duration when the guard is dropped
255    #[must_use = "keep the guard until WebSocket authentication finishes"]
256    pub fn track_ws_authentication(&self) -> impl Drop + '_ {
257        self.track(|metrics, duration| metrics.ws_auth_duration.observe(duration))
258    }
259
260    /// records user initialization duration when the guard is dropped
261    #[must_use = "keep the guard until WebSocket user initialization finishes"]
262    pub fn track_ws_user_initialization(&self) -> impl Drop + '_ {
263        self.track(|metrics, duration| metrics.ws_user_initialize_duration.observe(duration))
264    }
265
266    fn track<'a>(&'a self, finish: impl Fn(&Self, Duration) + 'a) -> impl Drop + 'a {
267        MetricGuard {
268            metrics: self,
269            started_at: Instant::now(),
270            finish,
271        }
272    }
273
274    pub fn add_active_transport_users(&self, delta: i64) {
275        self.active_transport_users.add(delta);
276    }
277
278    /// records one transport-health edge and keeps current-state gauges balanced
279    ///
280    /// callers must pass the previously recorded health value
281    /// `None -> Some` joins the gauge
282    /// `Some -> None` leaves it
283    pub fn record_transport_health_transition(
284        &self,
285        previous: Option<TransportHealthState>,
286        next: Option<TransportHealthState>,
287    ) {
288        if previous == next {
289            return;
290        }
291        match (previous, next) {
292            (None, Some(TransportHealthState::Connected)) => self
293                .transport_health_transitions
294                .increment(TransportHealthTransition::UnsetToConnected),
295            (None, Some(TransportHealthState::Disconnected)) => self
296                .transport_health_transitions
297                .increment(TransportHealthTransition::UnsetToDisconnected),
298            (Some(TransportHealthState::Connected), Some(TransportHealthState::Disconnected)) => {
299                self.transport_health_transitions
300                    .increment(TransportHealthTransition::ConnectedToDisconnected);
301            }
302            (Some(TransportHealthState::Disconnected), Some(TransportHealthState::Connected)) => {
303                self.transport_health_transitions
304                    .increment(TransportHealthTransition::DisconnectedToConnected);
305            }
306            (Some(TransportHealthState::Connected), None) => self
307                .transport_health_transitions
308                .increment(TransportHealthTransition::ConnectedToUnset),
309            (Some(TransportHealthState::Disconnected), None) => self
310                .transport_health_transitions
311                .increment(TransportHealthTransition::DisconnectedToUnset),
312            (None, None)
313            | (Some(TransportHealthState::Connected), Some(TransportHealthState::Connected))
314            | (
315                Some(TransportHealthState::Disconnected),
316                Some(TransportHealthState::Disconnected),
317            ) => {}
318        }
319        if let Some(health) = previous {
320            self.transport_health_users.add(health, -1);
321        }
322        if let Some(health) = next {
323            self.transport_health_users.add(health, 1);
324        }
325    }
326
327    pub fn record_recording_start_accepted(&self) {
328        self.recording_actions
329            .increment(RecordingActionOutcome::StartAccepted);
330    }
331
332    pub fn record_recording_start_rejected(&self) {
333        self.recording_actions
334            .increment(RecordingActionOutcome::StartRejected);
335    }
336
337    pub fn record_recording_stop_accepted(&self) {
338        self.recording_actions
339            .increment(RecordingActionOutcome::StopAccepted);
340    }
341
342    pub fn record_recording_stop_rejected(&self) {
343        self.recording_actions
344            .increment(RecordingActionOutcome::StopRejected);
345    }
346
347    pub fn record_recording_captured_packet(&self) {
348        self.recording_captured_packets.increment();
349    }
350
351    pub fn record_recording_captured_stream(&self) {
352        self.recording_captured_streams.increment();
353    }
354
355    /// Registers one worker-local recorder for hot RTP packet metrics.
356    ///
357    /// The caller should keep the returned handle beside the packet loop and
358    /// reuse it for every RTP packet observation owned by that worker.
359    pub fn register_rtp_worker(&self) -> Arc<RtpMetricsRecorder> {
360        self.rtp_metrics.register_worker(None)
361    }
362
363    /// Registers one worker-local recorder for a known media worker.
364    ///
365    /// The media-worker id is exported only as a bounded worker label. Runtime
366    /// code must not pass room or user identity here.
367    pub fn register_rtp_worker_for_media_worker(
368        &self,
369        media_worker_id: usize,
370    ) -> Arc<RtpMetricsRecorder> {
371        self.rtp_metrics.register_worker(Some(media_worker_id))
372    }
373
374    /// Registers one worker-local recorder for RTC packet-loop metrics.
375    ///
376    /// The caller should keep the returned handle beside the packet loop and
377    /// reuse it for UDP datagram and route-control observations owned by that
378    /// worker.
379    pub fn register_rtc_worker(&self) -> Arc<RtcMetricsRecorder> {
380        self.rtc_metrics.register_worker()
381    }
382
383    pub fn record_rtp_relay_overload_drop(&self, destination: RtpRelayDropKind) {
384        self.rtp_relay_overload_drops.increment(destination);
385    }
386
387    pub fn record_transport_ice_state_change(&self, state: TransportIceState) {
388        self.transport_ice_state_changes.increment(state);
389    }
390
391    pub fn record_transport_dtls_connected(&self) {
392        self.transport_dtls_connected.increment();
393    }
394
395    pub fn record_transport_user_lifetime(&self, duration: Duration) {
396        self.transport_user_lifetime_count.increment();
397        self.transport_user_lifetime_sum_micros
398            .add_u64(u64::try_from(duration.as_micros()).unwrap_or(u64::MAX));
399        if duration <= Duration::from_secs(1) {
400            self.transport_user_lifetime_buckets
401                .increment(TransportUserLifetimeBucket::Le1Second);
402        }
403        if duration <= Duration::from_secs(10) {
404            self.transport_user_lifetime_buckets
405                .increment(TransportUserLifetimeBucket::Le10Seconds);
406        }
407        if duration <= Duration::from_mins(1) {
408            self.transport_user_lifetime_buckets
409                .increment(TransportUserLifetimeBucket::Le60Seconds);
410        }
411        if duration <= Duration::from_mins(5) {
412            self.transport_user_lifetime_buckets
413                .increment(TransportUserLifetimeBucket::Le300Seconds);
414        }
415    }
416
417    pub fn record_media_quality_sample(&self, sample: MediaQualitySample) {
418        self.media_quality_samples.increment(sample);
419    }
420
421    pub fn record_media_quality_rtt(&self, sample: MediaQualitySample, duration: Duration) {
422        self.media_quality_rtt.observe(sample, duration);
423    }
424
425    pub fn record_media_quality_loss_ppm(
426        &self,
427        direction: MediaQualityLossDirection,
428        loss_ppm: u64,
429    ) {
430        self.media_quality_loss_ppm_observed
431            .add_u64(direction, loss_ppm);
432        self.media_quality_loss_observations.increment(direction);
433    }
434
435    pub fn record_media_quality_bwe_bps(&self, bwe_bps: u64) {
436        self.media_quality_bwe_bps_observed.add_u64(bwe_bps);
437        self.media_quality_bwe_observations.increment();
438    }
439
440    pub fn record_media_quality_jitter_rtp_timestamp_units(&self, jitter: u64) {
441        self.media_quality_jitter_rtp_timestamp_units_observed
442            .add_u64(jitter);
443        self.media_quality_jitter_observations.increment();
444    }
445
446    pub fn record_transport_cleanup_failure(&self) {
447        self.transport_cleanup_failures.increment();
448    }
449
450    pub fn record_source_selection_update(&self, selector: SourceSelectionKind) {
451        self.source_selection_updates.increment(selector);
452    }
453
454    pub fn record_budget_solver_outcome(&self, outcome: BudgetSolverOutcome) {
455        self.budget_solver_outcomes.increment(outcome);
456    }
457}