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