Skip to main content

o_sfu_telemetry/metrics/
rtp.rs

1use std::{
2    collections::BTreeMap,
3    sync::{Arc, Mutex},
4};
5
6use super::{
7    counter::{MetricLabel, PaddedCounterFamily},
8    labels::{RtpDecoderRefreshScope, RtpFlowDirection, RtpForwardDestinationKind},
9};
10
11const RTP_FLOW_DIRECTION_COUNT: usize = <RtpFlowDirection as MetricLabel>::COUNT;
12const RTP_FORWARD_DESTINATION_COUNT: usize = <RtpForwardDestinationKind as MetricLabel>::COUNT;
13const RTP_DECODER_REFRESH_SCOPE_COUNT: usize = <RtpDecoderRefreshScope as MetricLabel>::COUNT;
14
15/// Worker-local RTP packet metric recorder.
16///
17/// Packet loops keep one recorder for their full worker lifetime. Updates touch
18/// only this worker's padded atomics while `RuntimeMetrics` aggregates all
19/// registered recorders during scrape capture.
20#[derive(Debug, Default)]
21pub struct RtpMetricsRecorder {
22    packets: PaddedCounterFamily<RtpFlowDirection>,
23    payload_bytes: PaddedCounterFamily<RtpFlowDirection>,
24    forwarded_packets: PaddedCounterFamily<RtpForwardDestinationKind>,
25    forwarded_payload_bytes: PaddedCounterFamily<RtpForwardDestinationKind>,
26    decoder_refreshes: PaddedCounterFamily<RtpDecoderRefreshScope>,
27}
28
29impl RtpMetricsRecorder {
30    pub fn record_ingress(&self, payload_bytes: usize) {
31        self.packets.increment(RtpFlowDirection::Ingress);
32        self.payload_bytes
33            .add(RtpFlowDirection::Ingress, payload_bytes);
34    }
35
36    pub fn record_egress(&self, payload_bytes: usize) {
37        self.packets.increment(RtpFlowDirection::Egress);
38        self.payload_bytes
39            .add(RtpFlowDirection::Egress, payload_bytes);
40    }
41
42    pub fn record_forwarded(&self, destination: RtpForwardDestinationKind, payload_bytes: usize) {
43        self.forwarded_packets.increment(destination);
44        self.forwarded_payload_bytes.add(destination, payload_bytes);
45    }
46
47    pub fn record_decoder_refresh(&self, scope: RtpDecoderRefreshScope) {
48        self.decoder_refreshes.increment(scope);
49    }
50}
51
52#[derive(Debug, Default)]
53pub(super) struct RtpMetrics {
54    worker_recorders: Mutex<Vec<RtpWorkerMetricsRecorder>>,
55}
56
57impl RtpMetrics {
58    pub(super) fn register_worker(
59        &self,
60        media_worker_id: Option<usize>,
61    ) -> Arc<RtpMetricsRecorder> {
62        let recorder = Arc::new(RtpMetricsRecorder::default());
63        {
64            let mut workers = match self.worker_recorders.lock() {
65                Ok(workers) => workers,
66                Err(poisoned) => poisoned.into_inner(),
67            };
68            workers.push(RtpWorkerMetricsRecorder {
69                media_worker_id,
70                recorder: Arc::clone(&recorder),
71            });
72        }
73        recorder
74    }
75
76    pub(super) fn snapshot(&self) -> RtpMetricsSnapshot {
77        let mut snapshot = RtpMetricsSnapshot::default();
78        {
79            let workers = match self.worker_recorders.lock() {
80                Ok(workers) => workers,
81                Err(poisoned) => poisoned.into_inner(),
82            };
83            let mut worker_snapshots = BTreeMap::<usize, RtpWorkerMetricsSnapshot>::new();
84            for worker in workers.iter() {
85                snapshot.add_recorder(&worker.recorder);
86                if let Some(media_worker_id) = worker.media_worker_id {
87                    worker_snapshots
88                        .entry(media_worker_id)
89                        .or_insert_with(|| RtpWorkerMetricsSnapshot::new(media_worker_id))
90                        .traffic
91                        .add_recorder(&worker.recorder);
92                }
93            }
94            drop(workers);
95            snapshot.worker_snapshots = worker_snapshots.into_values().collect();
96        }
97        snapshot
98    }
99}
100
101#[derive(Debug)]
102struct RtpWorkerMetricsRecorder {
103    media_worker_id: Option<usize>,
104    recorder: Arc<RtpMetricsRecorder>,
105}
106
107#[derive(Debug, Default)]
108pub(super) struct RtpMetricsSnapshot {
109    pub(super) traffic: RtpTrafficSnapshot,
110    decoder_refreshes: [u64; RTP_DECODER_REFRESH_SCOPE_COUNT],
111    worker_snapshots: Vec<RtpWorkerMetricsSnapshot>,
112}
113
114impl RtpMetricsSnapshot {
115    pub(super) fn decoder_refreshes(&self, scope: RtpDecoderRefreshScope) -> u64 {
116        self.decoder_refreshes
117            .get(scope.as_index())
118            .copied()
119            .unwrap_or(0)
120    }
121
122    pub(super) fn worker_snapshots(&self) -> &[RtpWorkerMetricsSnapshot] {
123        &self.worker_snapshots
124    }
125
126    fn add_recorder(&mut self, recorder: &RtpMetricsRecorder) {
127        self.traffic.add_recorder(recorder);
128        recorder
129            .decoder_refreshes
130            .accumulate_into(&mut self.decoder_refreshes);
131    }
132}
133
134#[derive(Debug, Default)]
135pub(super) struct RtpWorkerMetricsSnapshot {
136    media_worker_id: usize,
137    pub(super) traffic: RtpTrafficSnapshot,
138}
139
140impl RtpWorkerMetricsSnapshot {
141    fn new(media_worker_id: usize) -> Self {
142        Self {
143            media_worker_id,
144            ..Self::default()
145        }
146    }
147
148    pub(super) const fn media_worker_id(&self) -> usize {
149        self.media_worker_id
150    }
151}
152
153#[derive(Debug, Default)]
154pub(super) struct RtpTrafficSnapshot {
155    packets: [u64; RTP_FLOW_DIRECTION_COUNT],
156    payload_bytes: [u64; RTP_FLOW_DIRECTION_COUNT],
157    forwarded_packets: [u64; RTP_FORWARD_DESTINATION_COUNT],
158    forwarded_payload_bytes: [u64; RTP_FORWARD_DESTINATION_COUNT],
159}
160
161impl RtpTrafficSnapshot {
162    pub(super) fn packets(&self, direction: RtpFlowDirection) -> u64 {
163        self.packets.get(direction.as_index()).copied().unwrap_or(0)
164    }
165
166    pub(super) fn payload_bytes(&self, direction: RtpFlowDirection) -> u64 {
167        self.payload_bytes
168            .get(direction.as_index())
169            .copied()
170            .unwrap_or(0)
171    }
172
173    pub(super) fn forwarded_packets(&self, destination: RtpForwardDestinationKind) -> u64 {
174        self.forwarded_packets
175            .get(destination.as_index())
176            .copied()
177            .unwrap_or(0)
178    }
179
180    pub(super) fn forwarded_payload_bytes(&self, destination: RtpForwardDestinationKind) -> u64 {
181        self.forwarded_payload_bytes
182            .get(destination.as_index())
183            .copied()
184            .unwrap_or(0)
185    }
186
187    fn add_recorder(&mut self, recorder: &RtpMetricsRecorder) {
188        for direction in <RtpFlowDirection as MetricLabel>::VARIANTS {
189            self.add_flow(
190                *direction,
191                recorder.packets.load(*direction),
192                recorder.payload_bytes.load(*direction),
193            );
194        }
195        for destination in <RtpForwardDestinationKind as MetricLabel>::VARIANTS {
196            self.add_forwarded(
197                *destination,
198                recorder.forwarded_packets.load(*destination),
199                recorder.forwarded_payload_bytes.load(*destination),
200            );
201        }
202    }
203
204    fn add_flow(&mut self, direction: RtpFlowDirection, packets: u64, payload_bytes: u64) {
205        if let Some(counter) = self.packets.get_mut(direction.as_index()) {
206            *counter = counter.saturating_add(packets);
207        }
208        if let Some(counter) = self.payload_bytes.get_mut(direction.as_index()) {
209            *counter = counter.saturating_add(payload_bytes);
210        }
211    }
212
213    fn add_forwarded(
214        &mut self,
215        destination: RtpForwardDestinationKind,
216        packets: u64,
217        payload_bytes: u64,
218    ) {
219        if let Some(counter) = self.forwarded_packets.get_mut(destination.as_index()) {
220            *counter = counter.saturating_add(packets);
221        }
222        if let Some(counter) = self.forwarded_payload_bytes.get_mut(destination.as_index()) {
223            *counter = counter.saturating_add(payload_bytes);
224        }
225    }
226}