Skip to main content

prns_runtime/runtime/
metrics.rs

1use alloc::vec::Vec;
2
3use crate::engine::{AnnounceOrigin, EngineMetricsSnapshot};
4use crate::interfaces::{InterfaceId, InterfaceKind};
5use crate::runtime::ReliabilityMetricsSnapshot;
6use crate::units::InstantMillis;
7
8prns_macros::iterable_enum! {
9    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
10    #[repr(u8)]
11    pub enum AnnounceEgressOutcome {
12        Enqueued,
13        InterfaceUnavailable,
14        LaneFull,
15        LaneMissing,
16        IfacRejected,
17        PacerRejected,
18        PacerEvicted,
19        PacerExpired,
20    }
21}
22
23prns_macros::iterable_enum! {
24    #[derive(Debug, Clone, Copy, PartialEq, Eq)]
25    #[repr(u8)]
26    pub enum AnnounceBackpressureEvent {
27        Deferred,
28        Retry,
29        Recovered,
30    }
31}
32
33impl AnnounceEgressOutcome {
34    const fn index(self) -> usize {
35        self as usize
36    }
37}
38
39impl AnnounceBackpressureEvent {
40    const fn index(self) -> usize {
41        self as usize
42    }
43}
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq)]
46pub struct AnnounceEgressCounts {
47    counts: [[u64; AnnounceEgressOutcome::ALL.len()]; AnnounceOrigin::ALL.len()],
48}
49
50impl Default for AnnounceEgressCounts {
51    fn default() -> Self {
52        Self {
53            counts: [[0; AnnounceEgressOutcome::ALL.len()]; AnnounceOrigin::ALL.len()],
54        }
55    }
56}
57
58impl AnnounceEgressCounts {
59    pub const fn get(&self, origin: AnnounceOrigin, outcome: AnnounceEgressOutcome) -> u64 {
60        self.counts[origin.index()][outcome.index()]
61    }
62
63    pub fn iter(&self) -> impl Iterator<Item = (AnnounceOrigin, AnnounceEgressOutcome, u64)> + '_ {
64        AnnounceOrigin::ALL.into_iter().flat_map(move |origin| {
65            AnnounceEgressOutcome::ALL
66                .into_iter()
67                .map(move |outcome| (origin, outcome, self.get(origin, outcome)))
68        })
69    }
70
71    fn record(&mut self, origin: AnnounceOrigin, outcome: AnnounceEgressOutcome) {
72        let count = &mut self.counts[origin.index()][outcome.index()];
73        *count = count.saturating_add(1);
74    }
75}
76
77#[derive(Debug, Clone, Copy, PartialEq, Eq)]
78pub struct AnnounceBackpressureCounts {
79    counts: [[u64; AnnounceBackpressureEvent::ALL.len()]; AnnounceOrigin::ALL.len()],
80}
81
82impl Default for AnnounceBackpressureCounts {
83    fn default() -> Self {
84        Self {
85            counts: [[0; AnnounceBackpressureEvent::ALL.len()]; AnnounceOrigin::ALL.len()],
86        }
87    }
88}
89
90impl AnnounceBackpressureCounts {
91    pub const fn get(&self, origin: AnnounceOrigin, event: AnnounceBackpressureEvent) -> u64 {
92        self.counts[origin.index()][event.index()]
93    }
94
95    pub fn iter(
96        &self,
97    ) -> impl Iterator<Item = (AnnounceOrigin, AnnounceBackpressureEvent, u64)> + '_ {
98        AnnounceOrigin::ALL.into_iter().flat_map(move |origin| {
99            AnnounceBackpressureEvent::ALL
100                .into_iter()
101                .map(move |event| (origin, event, self.get(origin, event)))
102        })
103    }
104
105    fn record(&mut self, origin: AnnounceOrigin, event: AnnounceBackpressureEvent) {
106        let count = &mut self.counts[origin.index()][event.index()];
107        *count = count.saturating_add(1);
108    }
109}
110
111#[derive(Debug, Clone, Copy, PartialEq, Eq)]
112pub struct AnnounceOriginCounts {
113    counts: [u64; AnnounceOrigin::ALL.len()],
114}
115
116impl Default for AnnounceOriginCounts {
117    fn default() -> Self {
118        Self {
119            counts: [0; AnnounceOrigin::ALL.len()],
120        }
121    }
122}
123
124impl AnnounceOriginCounts {
125    pub const fn get(&self, origin: AnnounceOrigin) -> u64 {
126        self.counts[origin.index()]
127    }
128
129    pub fn iter(&self) -> impl ExactSizeIterator<Item = (AnnounceOrigin, u64)> + '_ {
130        AnnounceOrigin::ALL
131            .into_iter()
132            .map(|origin| (origin, self.get(origin)))
133    }
134
135    fn add(&mut self, origin: AnnounceOrigin, value: u64) {
136        let count = &mut self.counts[origin.index()];
137        *count = count.saturating_add(value);
138    }
139}
140
141#[derive(Debug, Clone, Copy, PartialEq, Eq)]
142pub struct EgressInterfaceKindCounts {
143    counts: [u64; InterfaceKind::ALL.len()],
144    unknown: u64,
145}
146
147impl Default for EgressInterfaceKindCounts {
148    fn default() -> Self {
149        Self {
150            counts: [0; InterfaceKind::ALL.len()],
151            unknown: 0,
152        }
153    }
154}
155
156impl EgressInterfaceKindCounts {
157    pub const fn get(&self, kind: InterfaceKind) -> u64 {
158        self.counts[kind as usize]
159    }
160
161    pub const fn unknown(&self) -> u64 {
162        self.unknown
163    }
164
165    pub fn iter(&self) -> impl ExactSizeIterator<Item = (InterfaceKind, u64)> + '_ {
166        InterfaceKind::ALL
167            .into_iter()
168            .map(|kind| (kind, self.get(kind)))
169    }
170
171    fn record(&mut self, kind: Option<InterfaceKind>) {
172        match kind {
173            Some(kind) => {
174                let count = &mut self.counts[kind as usize];
175                *count = count.saturating_add(1);
176            }
177            None => self.unknown = self.unknown.saturating_add(1),
178        }
179    }
180}
181
182#[derive(Debug, Clone, PartialEq, Eq)]
183pub struct InterfaceAnnounceEgressMetricsSnapshot {
184    pub interface: InterfaceId,
185    pub outcomes: AnnounceEgressCounts,
186    pub backpressure: AnnounceBackpressureCounts,
187    pub enqueued_bytes_by_origin: AnnounceOriginCounts,
188    pub pacer_queue_depth: u32,
189    pub pacer_deferred_depth: u32,
190    pub pacer_oldest_deferred_age_ms: u64,
191}
192
193#[derive(Debug, Clone, Default, PartialEq, Eq)]
194pub struct AnnounceEgressMetricsSnapshot {
195    pub outcomes: AnnounceEgressCounts,
196    pub backpressure: AnnounceBackpressureCounts,
197    pub enqueued_by_interface_kind: EgressInterfaceKindCounts,
198    pub enqueued_bytes_by_origin: AnnounceOriginCounts,
199    pub pacer_queue_depth: u32,
200    pub pacer_deferred_depth: u32,
201    pub pacer_oldest_deferred_age_ms: u64,
202    pub interfaces: Vec<InterfaceAnnounceEgressMetricsSnapshot>,
203}
204
205impl AnnounceEgressMetricsSnapshot {
206    pub fn record(
207        &mut self,
208        origin: AnnounceOrigin,
209        interface: InterfaceId,
210        outcome: AnnounceEgressOutcome,
211        bytes: usize,
212    ) {
213        self.outcomes.record(origin, outcome);
214        if outcome == AnnounceEgressOutcome::Enqueued {
215            self.enqueued_by_interface_kind.record(interface.kind());
216            self.enqueued_bytes_by_origin
217                .add(origin, u64::try_from(bytes).unwrap_or(u64::MAX));
218        }
219        let interface_metrics = self.interface_mut(interface);
220        interface_metrics.outcomes.record(origin, outcome);
221        if outcome == AnnounceEgressOutcome::Enqueued {
222            interface_metrics
223                .enqueued_bytes_by_origin
224                .add(origin, u64::try_from(bytes).unwrap_or(u64::MAX));
225        }
226    }
227
228    pub fn register_interface(&mut self, interface: InterfaceId) {
229        let _ = self.interface_mut(interface);
230    }
231
232    pub fn record_backpressure(
233        &mut self,
234        origin: AnnounceOrigin,
235        interface: InterfaceId,
236        event: AnnounceBackpressureEvent,
237    ) {
238        self.backpressure.record(origin, event);
239        self.interface_mut(interface)
240            .backpressure
241            .record(origin, event);
242    }
243
244    pub fn reset_pacer_gauges(&mut self) {
245        self.pacer_queue_depth = 0;
246        self.pacer_deferred_depth = 0;
247        self.pacer_oldest_deferred_age_ms = 0;
248        for metrics in &mut self.interfaces {
249            metrics.pacer_queue_depth = 0;
250            metrics.pacer_deferred_depth = 0;
251            metrics.pacer_oldest_deferred_age_ms = 0;
252        }
253    }
254
255    pub fn add_pacer_gauges(
256        &mut self,
257        interface: InterfaceId,
258        depth: usize,
259        deferred_depth: usize,
260        oldest_deferred_age_ms: u64,
261    ) {
262        let depth = u32::try_from(depth).unwrap_or(u32::MAX);
263        self.pacer_queue_depth = self.pacer_queue_depth.saturating_add(depth);
264        let deferred_depth = u32::try_from(deferred_depth).unwrap_or(u32::MAX);
265        self.pacer_deferred_depth = self.pacer_deferred_depth.saturating_add(deferred_depth);
266        self.pacer_oldest_deferred_age_ms = self
267            .pacer_oldest_deferred_age_ms
268            .max(oldest_deferred_age_ms);
269        let metrics = self.interface_mut(interface);
270        metrics.pacer_queue_depth = metrics.pacer_queue_depth.saturating_add(depth);
271        metrics.pacer_deferred_depth = metrics.pacer_deferred_depth.saturating_add(deferred_depth);
272        metrics.pacer_oldest_deferred_age_ms = metrics
273            .pacer_oldest_deferred_age_ms
274            .max(oldest_deferred_age_ms);
275    }
276
277    fn interface_mut(
278        &mut self,
279        interface: InterfaceId,
280    ) -> &mut InterfaceAnnounceEgressMetricsSnapshot {
281        if let Some(position) = self
282            .interfaces
283            .iter()
284            .position(|metrics| metrics.interface == interface)
285        {
286            return &mut self.interfaces[position];
287        }
288        self.interfaces
289            .push(InterfaceAnnounceEgressMetricsSnapshot {
290                interface,
291                outcomes: AnnounceEgressCounts::default(),
292                backpressure: AnnounceBackpressureCounts::default(),
293                enqueued_bytes_by_origin: AnnounceOriginCounts::default(),
294                pacer_queue_depth: 0,
295                pacer_deferred_depth: 0,
296                pacer_oldest_deferred_age_ms: 0,
297            });
298        let position = self.interfaces.len() - 1;
299        &mut self.interfaces[position]
300    }
301}
302
303#[derive(Debug, Clone, Copy, PartialEq, Eq)]
304pub struct EgressLaneMetricsSnapshot {
305    pub physical_interface: InterfaceId,
306    pub logical_interface: InterfaceId,
307    pub capacity: u32,
308    pub occupancy: u32,
309}
310
311#[derive(Debug, Clone, Default, PartialEq, Eq)]
312pub struct EgressMetricsSnapshot {
313    pub enqueued_frames: u64,
314    pub unavailable_frame_skips: u64,
315    pub full_lane_drops: u64,
316    pub missing_lane_drops: u64,
317    pub ifac_rejected_frames: u64,
318    pub announces: AnnounceEgressMetricsSnapshot,
319    pub lanes: Vec<EgressLaneMetricsSnapshot>,
320}
321
322#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
323pub struct CryptoMetricsSnapshot {
324    pub submitted_jobs: u64,
325    pub completed_jobs: u64,
326    pub queue_depth: u32,
327    pub maximum_queue_depth: u32,
328    pub backpressure_deferrals: u64,
329    pub packet_verdicts_owed: u32,
330}
331
332#[derive(Debug, Clone, PartialEq, Eq)]
333pub struct RuntimeMetricsSnapshot {
334    pub taken_at: InstantMillis,
335    pub engine: EngineMetricsSnapshot,
336    pub egress: EgressMetricsSnapshot,
337    pub crypto: Option<CryptoMetricsSnapshot>,
338    pub reliability: ReliabilityMetricsSnapshot,
339}
340
341#[cfg(test)]
342mod tests {
343    use super::*;
344
345    #[test]
346    fn announce_backpressure_counters_and_pacer_gauges_preserve_dimensions() {
347        let interface = InterfaceId::new([0x31; 8]);
348        let mut metrics = AnnounceEgressMetricsSnapshot::default();
349        metrics.record_backpressure(
350            AnnounceOrigin::Relay,
351            interface,
352            AnnounceBackpressureEvent::Deferred,
353        );
354        metrics.record_backpressure(
355            AnnounceOrigin::Relay,
356            interface,
357            AnnounceBackpressureEvent::Retry,
358        );
359        metrics.add_pacer_gauges(interface, 3, 2, 1_250);
360
361        assert_eq!(
362            metrics
363                .backpressure
364                .get(AnnounceOrigin::Relay, AnnounceBackpressureEvent::Deferred),
365            1
366        );
367        assert_eq!(metrics.pacer_queue_depth, 3);
368        assert_eq!(metrics.pacer_deferred_depth, 2);
369        assert_eq!(metrics.pacer_oldest_deferred_age_ms, 1_250);
370        assert_eq!(metrics.interfaces.len(), 1);
371        assert_eq!(
372            metrics.interfaces[0]
373                .backpressure
374                .get(AnnounceOrigin::Relay, AnnounceBackpressureEvent::Retry),
375            1
376        );
377
378        metrics.reset_pacer_gauges();
379        assert_eq!(metrics.pacer_queue_depth, 0);
380        assert_eq!(metrics.pacer_deferred_depth, 0);
381        assert_eq!(metrics.pacer_oldest_deferred_age_ms, 0);
382        assert_eq!(
383            metrics
384                .backpressure
385                .get(AnnounceOrigin::Relay, AnnounceBackpressureEvent::Deferred),
386            1,
387            "snapshot gauge refreshes must not reset cumulative lifecycle counters"
388        );
389    }
390}