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, 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 #[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 #[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 #[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 #[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 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 pub fn register_rtp_worker(&self) -> Arc<RtpMetricsRecorder> {
360 self.rtp_metrics.register_worker(None)
361 }
362
363 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 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}