o_sfu_telemetry/metrics/
rtp.rs1use 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#[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}