1use 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 pub fn record_subscription_intent_evictions(&self, targets: usize) {
100 self.subscription_intent_evictions.add(targets);
101 }
102
103 pub fn record_http_connection_rejection(&self) {
105 self.http_connection_rejections.increment();
106 }
107
108 #[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 #[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 #[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 #[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 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 pub fn register_rtp_worker(&self) -> Arc<RtpMetricsRecorder> {
378 self.rtp_metrics.register_worker(None)
379 }
380
381 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 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}