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}