Skip to main content

o_sfu_telemetry/metrics/
rtc.rs

1use std::{
2    mem::take,
3    sync::{Arc, Mutex, PoisonError},
4    time::{Duration, Instant},
5};
6
7use super::{
8    counter::{MetricLabel, PaddedCounter, PaddedCounterFamily},
9    labels::{
10        RtcDatagramDropReason, RtcDatagramRoutePath, RtcDrainFailureStage, RtcInputFailure,
11        RtcKeyframeRequestOutcome, RtcNackDirection, RtcOutputBudgetLimit,
12        RtcProducerSsrcBindingOutcome, RtcRelayEnqueueResult, RtcRemoteControlDropKind,
13        RtcRemotePacketGateConvergence, RtcRouteControlOutcome, RtcTransportIoFailure,
14        RtcWorkerObservationKind,
15    },
16};
17
18const RTC_DATAGRAM_ROUTE_PATH_COUNT: usize = <RtcDatagramRoutePath as MetricLabel>::COUNT;
19const RTC_DATAGRAM_DROP_REASON_COUNT: usize = <RtcDatagramDropReason as MetricLabel>::COUNT;
20const RTC_NACK_DIRECTION_COUNT: usize = <RtcNackDirection as MetricLabel>::COUNT;
21const RTC_TRANSPORT_IO_FAILURE_COUNT: usize = <RtcTransportIoFailure as MetricLabel>::COUNT;
22const RTC_INPUT_FAILURE_COUNT: usize = <RtcInputFailure as MetricLabel>::COUNT;
23const RTC_DRAIN_FAILURE_STAGE_COUNT: usize = <RtcDrainFailureStage as MetricLabel>::COUNT;
24const RTC_WORKER_OBSERVATION_KIND_COUNT: usize = <RtcWorkerObservationKind as MetricLabel>::COUNT;
25const RTC_OUTPUT_BUDGET_LIMIT_COUNT: usize = <RtcOutputBudgetLimit as MetricLabel>::COUNT;
26const RTC_PRODUCER_SSRC_BINDING_OUTCOME_COUNT: usize =
27    <RtcProducerSsrcBindingOutcome as MetricLabel>::COUNT;
28const RTC_ROUTE_CONTROL_OUTCOME_COUNT: usize = <RtcRouteControlOutcome as MetricLabel>::COUNT;
29const RTC_KEYFRAME_REQUEST_OUTCOME_COUNT: usize = <RtcKeyframeRequestOutcome as MetricLabel>::COUNT;
30const RTC_RELAY_ENQUEUE_RESULT_COUNT: usize = <RtcRelayEnqueueResult as MetricLabel>::COUNT;
31const RTC_REMOTE_CONTROL_DROP_KIND_COUNT: usize = <RtcRemoteControlDropKind as MetricLabel>::COUNT;
32const RTC_REMOTE_PACKET_GATE_CONVERGENCE_COUNT: usize =
33    <RtcRemotePacketGateConvergence as MetricLabel>::COUNT;
34
35const FAILURE_LOG_INTERVAL: Duration = Duration::from_secs(5);
36
37#[derive(Debug, Default)]
38struct FailureLogSlot {
39    last_log: Option<Instant>,
40    suppressed: u64,
41}
42
43impl FailureLogSlot {
44    fn observe(&mut self, now: Instant) -> Option<u64> {
45        if self
46            .last_log
47            .is_some_and(|last| now.saturating_duration_since(last) < FAILURE_LOG_INTERVAL)
48        {
49            self.suppressed = self.suppressed.saturating_add(1);
50            return None;
51        }
52        let suppressed = take(&mut self.suppressed);
53        self.last_log = Some(now);
54        Some(suppressed)
55    }
56}
57
58/// Worker-local RTC packet-loop metric recorder.
59///
60/// Packet loops keep one recorder for their full worker lifetime. Datagram and
61/// route-control updates touch only this worker's padded atomics while
62/// `RuntimeMetrics` aggregates all registered recorders during scrape capture.
63#[derive(Debug, Default)]
64pub struct RtcMetricsRecorder {
65    datagram_routes: PaddedCounterFamily<RtcDatagramRoutePath>,
66    datagram_drops: PaddedCounterFamily<RtcDatagramDropReason>,
67    datagram_fallback_scans: PaddedCounter,
68    datagram_scan_users: PaddedCounter,
69    nacks: PaddedCounterFamily<RtcNackDirection>,
70    rtx_packets_received_from_publisher: PaddedCounter,
71    rtx_payload_bytes_received_from_publisher: PaddedCounter,
72    rtcp_ingress_budget_drops: PaddedCounter,
73    transport_io_failures: PaddedCounterFamily<RtcTransportIoFailure>,
74    transport_io_log: Mutex<[FailureLogSlot; RTC_TRANSPORT_IO_FAILURE_COUNT]>,
75    input_failures: PaddedCounterFamily<RtcInputFailure>,
76    input_log: Mutex<[FailureLogSlot; RTC_INPUT_FAILURE_COUNT]>,
77    drain_failures: PaddedCounterFamily<RtcDrainFailureStage>,
78    worker_terminal_failures: PaddedCounter,
79    worker_observation_timeouts: PaddedCounterFamily<RtcWorkerObservationKind>,
80    output_budget_exhaustions: PaddedCounterFamily<RtcOutputBudgetLimit>,
81    output_budget_session_closes: PaddedCounter,
82    producer_ssrc_bindings: PaddedCounterFamily<RtcProducerSsrcBindingOutcome>,
83    route_control: PaddedCounterFamily<RtcRouteControlOutcome>,
84    keyframe_requests: PaddedCounterFamily<RtcKeyframeRequestOutcome>,
85    relay_enqueues: PaddedCounterFamily<RtcRelayEnqueueResult>,
86    relay_mailbox_depth_samples: PaddedCounter,
87    relay_mailbox_depth_total: PaddedCounter,
88    relay_drain_batches: PaddedCounter,
89    relay_drained_packets: PaddedCounter,
90    relay_drain_cap_hits: PaddedCounter,
91    remote_control_drops: PaddedCounterFamily<RtcRemoteControlDropKind>,
92    remote_packet_gate_convergence: PaddedCounterFamily<RtcRemotePacketGateConvergence>,
93}
94
95impl RtcMetricsRecorder {
96    pub fn record_rtc_datagram_route(&self, path: RtcDatagramRoutePath) {
97        self.datagram_routes.increment(path);
98    }
99
100    pub fn record_rtc_datagram_drop(&self, reason: RtcDatagramDropReason) {
101        self.datagram_drops.increment(reason);
102    }
103
104    pub fn record_rtc_datagram_fallback_scan(&self, examined_sessions: usize) {
105        self.datagram_fallback_scans.increment();
106        self.datagram_scan_users.add(examined_sessions);
107    }
108
109    pub fn record_rtc_nacks(&self, direction: RtcNackDirection, count: u64) {
110        self.nacks.add_u64(direction, count);
111    }
112
113    pub fn record_rtc_rtx_received_from_publisher(&self, payload_bytes: usize) {
114        self.rtx_packets_received_from_publisher.increment();
115        self.rtx_payload_bytes_received_from_publisher
116            .add(payload_bytes);
117    }
118
119    pub fn record_rtc_rtcp_ingress_budget_drop(&self) {
120        self.rtcp_ingress_budget_drops.increment();
121    }
122
123    /// Counts every failure and returns suppressed occurrences when this category may log.
124    pub fn record_rtc_transport_io_failure(&self, failure: RtcTransportIoFailure) -> Option<u64> {
125        self.transport_io_failures.increment(failure);
126        let mut slots = self
127            .transport_io_log
128            .lock()
129            .unwrap_or_else(PoisonError::into_inner);
130        slots.get_mut(failure.as_index())?.observe(Instant::now())
131    }
132
133    /// Counts every failure and returns suppressed occurrences when this category may log.
134    pub fn record_rtc_input_failure(&self, failure: RtcInputFailure) -> Option<u64> {
135        self.input_failures.increment(failure);
136        let mut slots = self
137            .input_log
138            .lock()
139            .unwrap_or_else(PoisonError::into_inner);
140        slots.get_mut(failure.as_index())?.observe(Instant::now())
141    }
142
143    pub fn record_rtc_drain_failure(&self, stage: RtcDrainFailureStage) {
144        self.drain_failures.increment(stage);
145    }
146
147    pub fn record_rtc_worker_terminal_failure(&self) {
148        self.worker_terminal_failures.increment();
149    }
150
151    pub fn record_rtc_worker_observation_timeout(&self, kind: RtcWorkerObservationKind) {
152        self.worker_observation_timeouts.increment(kind);
153    }
154
155    pub fn record_rtc_output_budget_exhaustion(&self, limit: RtcOutputBudgetLimit) {
156        self.output_budget_exhaustions.increment(limit);
157    }
158
159    pub fn record_rtc_output_budget_session_close(&self) {
160        self.output_budget_session_closes.increment();
161    }
162
163    /// Records one producer binding change or rejection after RTP authentication.
164    ///
165    /// Repeated packets matching an established binding do not change it and
166    /// should not increment this counter.
167    pub fn record_rtc_producer_ssrc_binding(&self, outcome: RtcProducerSsrcBindingOutcome) {
168        self.producer_ssrc_bindings.increment(outcome);
169    }
170
171    pub fn record_rtc_route_control(&self, outcome: RtcRouteControlOutcome) {
172        self.route_control.increment(outcome);
173    }
174
175    pub fn record_rtc_keyframe_request(&self, outcome: RtcKeyframeRequestOutcome) {
176        self.keyframe_requests.increment(outcome);
177    }
178
179    pub fn record_rtc_relay_enqueue(&self, result: RtcRelayEnqueueResult) {
180        self.relay_enqueues.increment(result);
181    }
182
183    pub fn record_rtc_relay_mailbox_depth(&self, depth: usize) {
184        self.relay_mailbox_depth_samples.increment();
185        self.relay_mailbox_depth_total.add(depth);
186    }
187
188    pub fn record_rtc_relay_drain_batch(&self, drained_packets: usize, cap_hit: bool) {
189        if drained_packets == 0 {
190            return;
191        }
192        self.relay_drain_batches.increment();
193        self.relay_drained_packets.add(drained_packets);
194        if cap_hit {
195            self.relay_drain_cap_hits.increment();
196        }
197    }
198
199    pub fn record_rtc_remote_control_drop(&self, kind: RtcRemoteControlDropKind) {
200        self.remote_control_drops.increment(kind);
201    }
202
203    pub fn record_rtc_remote_packet_gate_convergence(
204        &self,
205        outcome: RtcRemotePacketGateConvergence,
206    ) {
207        self.remote_packet_gate_convergence.increment(outcome);
208    }
209}
210
211#[derive(Debug, Default)]
212pub(super) struct RtcMetrics {
213    worker_recorders: Mutex<Vec<Arc<RtcMetricsRecorder>>>,
214}
215
216impl RtcMetrics {
217    pub(super) fn register_worker(&self) -> Arc<RtcMetricsRecorder> {
218        let recorder = Arc::new(RtcMetricsRecorder::default());
219        {
220            let mut workers = match self.worker_recorders.lock() {
221                Ok(workers) => workers,
222                Err(poisoned) => poisoned.into_inner(),
223            };
224            workers.push(Arc::clone(&recorder));
225        }
226        recorder
227    }
228
229    pub(super) fn snapshot(&self) -> RtcMetricsSnapshot {
230        let mut snapshot = RtcMetricsSnapshot::default();
231        {
232            let workers = match self.worker_recorders.lock() {
233                Ok(workers) => workers,
234                Err(poisoned) => poisoned.into_inner(),
235            };
236            for recorder in workers.iter() {
237                snapshot.add_recorder(recorder);
238            }
239        }
240        snapshot
241    }
242}
243
244#[derive(Debug, Default)]
245pub(super) struct RtcMetricsSnapshot {
246    datagram_routes: [u64; RTC_DATAGRAM_ROUTE_PATH_COUNT],
247    datagram_drops: [u64; RTC_DATAGRAM_DROP_REASON_COUNT],
248    datagram_fallback_scans: u64,
249    datagram_scan_users: u64,
250    nacks: [u64; RTC_NACK_DIRECTION_COUNT],
251    rtx_packets_received_from_publisher: u64,
252    rtx_payload_bytes_received_from_publisher: u64,
253    rtcp_ingress_budget_drops: u64,
254    transport_io_failures: [u64; RTC_TRANSPORT_IO_FAILURE_COUNT],
255    input_failures: [u64; RTC_INPUT_FAILURE_COUNT],
256    drain_failures: [u64; RTC_DRAIN_FAILURE_STAGE_COUNT],
257    worker_terminal_failures: u64,
258    worker_observation_timeouts: [u64; RTC_WORKER_OBSERVATION_KIND_COUNT],
259    output_budget_exhaustions: [u64; RTC_OUTPUT_BUDGET_LIMIT_COUNT],
260    output_budget_session_closes: u64,
261    producer_ssrc_bindings: [u64; RTC_PRODUCER_SSRC_BINDING_OUTCOME_COUNT],
262    route_control: [u64; RTC_ROUTE_CONTROL_OUTCOME_COUNT],
263    keyframe_requests: [u64; RTC_KEYFRAME_REQUEST_OUTCOME_COUNT],
264    relay_enqueues: [u64; RTC_RELAY_ENQUEUE_RESULT_COUNT],
265    relay_mailbox_depth_samples: u64,
266    relay_mailbox_depth_total: u64,
267    relay_drain_batches: u64,
268    relay_drained_packets: u64,
269    relay_drain_cap_hits: u64,
270    remote_control_drops: [u64; RTC_REMOTE_CONTROL_DROP_KIND_COUNT],
271    remote_packet_gate_convergence: [u64; RTC_REMOTE_PACKET_GATE_CONVERGENCE_COUNT],
272}
273
274impl RtcMetricsSnapshot {
275    pub(super) fn datagram_routes(&self, path: RtcDatagramRoutePath) -> u64 {
276        self.datagram_routes
277            .get(path.as_index())
278            .copied()
279            .unwrap_or(0)
280    }
281
282    pub(super) fn datagram_drops(&self, reason: RtcDatagramDropReason) -> u64 {
283        self.datagram_drops
284            .get(reason.as_index())
285            .copied()
286            .unwrap_or(0)
287    }
288
289    pub(super) const fn datagram_fallback_scans(&self) -> u64 {
290        self.datagram_fallback_scans
291    }
292
293    pub(super) const fn datagram_scan_users(&self) -> u64 {
294        self.datagram_scan_users
295    }
296
297    pub(super) fn nacks(&self, direction: RtcNackDirection) -> u64 {
298        self.nacks.get(direction.as_index()).copied().unwrap_or(0)
299    }
300
301    pub(super) const fn rtx_packets_received_from_publisher(&self) -> u64 {
302        self.rtx_packets_received_from_publisher
303    }
304
305    pub(super) const fn rtx_payload_bytes_received_from_publisher(&self) -> u64 {
306        self.rtx_payload_bytes_received_from_publisher
307    }
308
309    pub(super) const fn rtcp_ingress_budget_drops(&self) -> u64 {
310        self.rtcp_ingress_budget_drops
311    }
312
313    pub(super) fn transport_io_failures(&self, failure: RtcTransportIoFailure) -> u64 {
314        self.transport_io_failures
315            .get(failure.as_index())
316            .copied()
317            .unwrap_or(0)
318    }
319
320    pub(super) fn input_failures(&self, failure: RtcInputFailure) -> u64 {
321        self.input_failures
322            .get(failure.as_index())
323            .copied()
324            .unwrap_or(0)
325    }
326
327    pub(super) fn drain_failures(&self, stage: RtcDrainFailureStage) -> u64 {
328        self.drain_failures
329            .get(stage.as_index())
330            .copied()
331            .unwrap_or(0)
332    }
333
334    pub(super) const fn worker_terminal_failures(&self) -> u64 {
335        self.worker_terminal_failures
336    }
337
338    pub(super) fn worker_observation_timeouts(&self, kind: RtcWorkerObservationKind) -> u64 {
339        self.worker_observation_timeouts
340            .get(kind.as_index())
341            .copied()
342            .unwrap_or(0)
343    }
344
345    pub(super) fn output_budget_exhaustions(&self, limit: RtcOutputBudgetLimit) -> u64 {
346        self.output_budget_exhaustions
347            .get(limit.as_index())
348            .copied()
349            .unwrap_or(0)
350    }
351
352    pub(super) const fn output_budget_session_closes(&self) -> u64 {
353        self.output_budget_session_closes
354    }
355
356    pub(super) fn producer_ssrc_bindings(&self, outcome: RtcProducerSsrcBindingOutcome) -> u64 {
357        self.producer_ssrc_bindings
358            .get(outcome.as_index())
359            .copied()
360            .unwrap_or(0)
361    }
362
363    pub(super) fn route_control(&self, outcome: RtcRouteControlOutcome) -> u64 {
364        self.route_control
365            .get(outcome.as_index())
366            .copied()
367            .unwrap_or(0)
368    }
369
370    pub(super) fn keyframe_requests(&self, outcome: RtcKeyframeRequestOutcome) -> u64 {
371        self.keyframe_requests
372            .get(outcome.as_index())
373            .copied()
374            .unwrap_or(0)
375    }
376
377    pub(super) fn relay_enqueues(&self, result: RtcRelayEnqueueResult) -> u64 {
378        self.relay_enqueues
379            .get(result.as_index())
380            .copied()
381            .unwrap_or(0)
382    }
383
384    pub(super) const fn relay_mailbox_depth_samples(&self) -> u64 {
385        self.relay_mailbox_depth_samples
386    }
387
388    pub(super) const fn relay_mailbox_depth_total(&self) -> u64 {
389        self.relay_mailbox_depth_total
390    }
391
392    pub(super) const fn relay_drain_batches(&self) -> u64 {
393        self.relay_drain_batches
394    }
395
396    pub(super) const fn relay_drained_packets(&self) -> u64 {
397        self.relay_drained_packets
398    }
399
400    pub(super) const fn relay_drain_cap_hits(&self) -> u64 {
401        self.relay_drain_cap_hits
402    }
403
404    pub(super) fn remote_control_drops(&self, kind: RtcRemoteControlDropKind) -> u64 {
405        self.remote_control_drops
406            .get(kind.as_index())
407            .copied()
408            .unwrap_or(0)
409    }
410
411    pub(super) fn remote_packet_gate_convergence(
412        &self,
413        outcome: RtcRemotePacketGateConvergence,
414    ) -> u64 {
415        self.remote_packet_gate_convergence
416            .get(outcome.as_index())
417            .copied()
418            .unwrap_or(0)
419    }
420
421    fn add_recorder(&mut self, recorder: &RtcMetricsRecorder) {
422        recorder
423            .datagram_routes
424            .accumulate_into(&mut self.datagram_routes);
425        recorder
426            .datagram_drops
427            .accumulate_into(&mut self.datagram_drops);
428        self.datagram_fallback_scans = self
429            .datagram_fallback_scans
430            .saturating_add(recorder.datagram_fallback_scans.load());
431        self.datagram_scan_users = self
432            .datagram_scan_users
433            .saturating_add(recorder.datagram_scan_users.load());
434        recorder.nacks.accumulate_into(&mut self.nacks);
435        self.rtx_packets_received_from_publisher = self
436            .rtx_packets_received_from_publisher
437            .saturating_add(recorder.rtx_packets_received_from_publisher.load());
438        self.rtx_payload_bytes_received_from_publisher = self
439            .rtx_payload_bytes_received_from_publisher
440            .saturating_add(recorder.rtx_payload_bytes_received_from_publisher.load());
441        self.rtcp_ingress_budget_drops = self
442            .rtcp_ingress_budget_drops
443            .saturating_add(recorder.rtcp_ingress_budget_drops.load());
444        recorder
445            .transport_io_failures
446            .accumulate_into(&mut self.transport_io_failures);
447        recorder
448            .input_failures
449            .accumulate_into(&mut self.input_failures);
450        recorder
451            .drain_failures
452            .accumulate_into(&mut self.drain_failures);
453        self.worker_terminal_failures = self
454            .worker_terminal_failures
455            .saturating_add(recorder.worker_terminal_failures.load());
456        recorder
457            .worker_observation_timeouts
458            .accumulate_into(&mut self.worker_observation_timeouts);
459        recorder
460            .output_budget_exhaustions
461            .accumulate_into(&mut self.output_budget_exhaustions);
462        self.output_budget_session_closes = self
463            .output_budget_session_closes
464            .saturating_add(recorder.output_budget_session_closes.load());
465        recorder
466            .producer_ssrc_bindings
467            .accumulate_into(&mut self.producer_ssrc_bindings);
468        recorder
469            .route_control
470            .accumulate_into(&mut self.route_control);
471        recorder
472            .keyframe_requests
473            .accumulate_into(&mut self.keyframe_requests);
474        recorder
475            .relay_enqueues
476            .accumulate_into(&mut self.relay_enqueues);
477        self.relay_mailbox_depth_samples = self
478            .relay_mailbox_depth_samples
479            .saturating_add(recorder.relay_mailbox_depth_samples.load());
480        self.relay_mailbox_depth_total = self
481            .relay_mailbox_depth_total
482            .saturating_add(recorder.relay_mailbox_depth_total.load());
483        self.relay_drain_batches = self
484            .relay_drain_batches
485            .saturating_add(recorder.relay_drain_batches.load());
486        self.relay_drained_packets = self
487            .relay_drained_packets
488            .saturating_add(recorder.relay_drained_packets.load());
489        self.relay_drain_cap_hits = self
490            .relay_drain_cap_hits
491            .saturating_add(recorder.relay_drain_cap_hits.load());
492        recorder
493            .remote_control_drops
494            .accumulate_into(&mut self.remote_control_drops);
495        recorder
496            .remote_packet_gate_convergence
497            .accumulate_into(&mut self.remote_packet_gate_convergence);
498    }
499}
500
501#[cfg(test)]
502#[path = "TESTS/failure_log.rs"]
503mod failure_log_tests;