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#[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 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 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 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;