1use core::sync::atomic::Ordering;
27use portable_atomic::AtomicU32;
31
32#[derive(Debug, Default)]
36pub struct PubMonitorCell {
37 pub count: AtomicU32,
38 pub max_latency_us: AtomicU32,
42}
43
44impl PubMonitorCell {
45 pub const fn new() -> Self {
46 Self {
47 count: AtomicU32::new(0),
48 max_latency_us: AtomicU32::new(0),
49 }
50 }
51}
52
53#[derive(Debug, Default)]
57pub struct SubMonitorCell {
58 pub max_age_ms: AtomicU32,
60}
61
62impl SubMonitorCell {
63 pub const fn new() -> Self {
64 Self {
65 max_age_ms: AtomicU32::new(0),
66 }
67 }
68
69 pub fn observe(&self, stamp_us: u64, epoch_now_us: u64) {
73 let age_ms = (epoch_now_us.saturating_sub(stamp_us) / 1_000).min(u32::MAX as u64) as u32;
74 self.max_age_ms.fetch_max(age_ms, Ordering::Relaxed);
75 }
76}
77
78pub fn peek_stamp_us(raw: &[u8], offset: usize) -> Option<u64> {
83 let sec_b = raw.get(offset..offset + 4)?;
84 let nsec_b = raw.get(offset + 4..offset + 8)?;
85 let sec = i32::from_le_bytes([sec_b[0], sec_b[1], sec_b[2], sec_b[3]]);
86 let nsec = u32::from_le_bytes([nsec_b[0], nsec_b[1], nsec_b[2], nsec_b[3]]);
87 if sec <= 0 {
88 return None;
89 }
90 Some(sec as u64 * 1_000_000 + nsec as u64 / 1_000)
91}
92
93#[derive(Debug, Clone, Copy)]
95pub struct MonitorSpec {
96 pub topic: &'static str,
99 pub fqn: &'static str,
102 pub min_rate_hz_milli: u32,
105 pub max_latency_ms: u32,
109 pub cell: &'static PubMonitorCell,
111}
112
113#[derive(Debug, Clone, Copy)]
117pub struct AgeMonitorSpec {
118 pub topic: &'static str,
120 pub fqn: &'static str,
122 pub max_age_ms: u32,
124 pub cell: &'static SubMonitorCell,
126}
127
128pub const RATE_CHECK_INTERVAL_US: u64 = 5_000_000;
131
132pub const MAX_MONITORS: usize = 8;
134pub const MAX_VIOLATIONS: usize = 8;
136
137#[derive(Debug, Clone)]
139pub struct Violation {
140 pub rule: &'static str,
143 pub fqn: &'static str,
146 pub measured: u32,
149 pub declared: u32,
151}
152
153#[derive(Debug, Clone, Copy, Default)]
155pub(crate) struct MonitorState {
156 pub(crate) opened: bool,
160 pub(crate) window_start_us: u64,
161 pub(crate) count_at_window_start: u32,
162 pub(crate) violated_last_window: bool,
164 pub(crate) latency_violated_last_window: bool,
166}
167
168pub(crate) fn check_rate(
175 spec: &MonitorSpec,
176 state: &mut MonitorState,
177 now_us: u64,
178) -> Option<Violation> {
179 if spec.min_rate_hz_milli == 0 {
180 return None;
181 }
182 let count = spec.cell.count.load(Ordering::Relaxed);
183 if !state.opened {
184 state.opened = true;
186 state.window_start_us = now_us;
187 state.count_at_window_start = count;
188 return None;
189 }
190 let window_us = now_us.saturating_sub(state.window_start_us);
191 if window_us < RATE_CHECK_INTERVAL_US {
192 return None;
193 }
194 let published = count.wrapping_sub(state.count_at_window_start) as u64;
195 let measured_milli_hz =
197 (published.saturating_mul(1_000_000_000) / window_us.max(1)).min(u32::MAX as u64) as u32;
198
199 state.window_start_us = now_us;
201 state.count_at_window_start = count;
202
203 if measured_milli_hz < spec.min_rate_hz_milli {
204 if state.violated_last_window {
205 return None; }
207 state.violated_last_window = true;
208 Some(Violation {
209 rule: "rate-hierarchy-runtime",
210 fqn: spec.fqn,
211 measured: measured_milli_hz,
212 declared: spec.min_rate_hz_milli,
213 })
214 } else {
215 state.violated_last_window = false;
216 None
217 }
218}
219
220pub(crate) fn check_latency(spec: &MonitorSpec, state: &mut MonitorState) -> Option<Violation> {
226 if spec.max_latency_ms == 0 {
227 return None;
228 }
229 let max_us = spec.cell.max_latency_us.swap(0, Ordering::Relaxed);
230 let max_ms = max_us / 1_000;
231 if max_ms > spec.max_latency_ms {
232 if state.latency_violated_last_window {
233 return None;
234 }
235 state.latency_violated_last_window = true;
236 Some(Violation {
237 rule: "max-latency-runtime",
238 fqn: spec.fqn,
239 measured: max_ms,
240 declared: spec.max_latency_ms,
241 })
242 } else {
243 state.latency_violated_last_window = false;
246 None
247 }
248}
249
250#[derive(Debug, Clone, Copy, Default)]
252pub(crate) struct AgeState {
253 pub(crate) violated_last_window: bool,
254}
255
256pub(crate) fn check_age(spec: &AgeMonitorSpec, state: &mut AgeState) -> Option<Violation> {
259 if spec.max_age_ms == 0 {
260 return None;
261 }
262 let max_ms = spec.cell.max_age_ms.swap(0, Ordering::Relaxed);
263 if max_ms > spec.max_age_ms {
264 if state.violated_last_window {
265 return None;
266 }
267 state.violated_last_window = true;
268 Some(Violation {
269 rule: "max-age-runtime",
270 fqn: spec.fqn,
271 measured: max_ms,
272 declared: spec.max_age_ms,
273 })
274 } else {
275 state.violated_last_window = false;
276 None
277 }
278}
279
280#[cfg(test)]
281mod tests {
282 use super::*;
283
284 static CELL: PubMonitorCell = PubMonitorCell::new();
285
286 fn spec(min_milli: u32) -> MonitorSpec {
287 MonitorSpec {
288 topic: "/chatter",
289 fqn: "/demo/talker/chatter",
290 min_rate_hz_milli: min_milli,
291 max_latency_ms: 0,
292 cell: &CELL,
293 }
294 }
295
296 #[test]
297 fn slow_publisher_fires_once_until_recovery() {
298 CELL.count.store(0, Ordering::Relaxed);
299 let s = spec(100_000); let mut st = MonitorState::default();
301
302 assert!(check_rate(&s, &mut st, 0).is_none());
304 CELL.count.store(5, Ordering::Relaxed);
306 let v = check_rate(&s, &mut st, RATE_CHECK_INTERVAL_US).expect("fires");
307 assert_eq!(v.rule, "rate-hierarchy-runtime");
308 assert_eq!(v.fqn, "/demo/talker/chatter");
309 assert_eq!(v.measured, 1_000);
310 assert_eq!(v.declared, 100_000);
311 CELL.count.store(10, Ordering::Relaxed);
313 assert!(check_rate(&s, &mut st, 2 * RATE_CHECK_INTERVAL_US).is_none());
314 CELL.count.store(510, Ordering::Relaxed);
316 assert!(check_rate(&s, &mut st, 3 * RATE_CHECK_INTERVAL_US).is_none());
317 CELL.count.store(511, Ordering::Relaxed);
319 assert!(check_rate(&s, &mut st, 4 * RATE_CHECK_INTERVAL_US).is_some());
320 }
321
322 #[test]
323 fn compliant_and_uncontracted_stay_silent() {
324 static C2: PubMonitorCell = PubMonitorCell::new();
325 let s = MonitorSpec {
326 topic: "/t",
327 fqn: "/n/t",
328 min_rate_hz_milli: 500, max_latency_ms: 0,
330 cell: &C2,
331 };
332 let mut st = MonitorState::default();
333 assert!(check_rate(&s, &mut st, 0).is_none());
334 C2.count.store(5, Ordering::Relaxed); assert!(check_rate(&s, &mut st, RATE_CHECK_INTERVAL_US).is_none());
336
337 let s0 = MonitorSpec {
339 topic: "/t",
340 fqn: "/n/t",
341 min_rate_hz_milli: 0,
342 max_latency_ms: 0,
343 cell: &C2,
344 };
345 let mut st0 = MonitorState::default();
346 assert!(check_rate(&s0, &mut st0, 10 * RATE_CHECK_INTERVAL_US).is_none());
347 }
348
349 #[test]
350 fn stale_take_fires_age_once_until_recovery() {
351 static SC: SubMonitorCell = SubMonitorCell::new();
352 let s = AgeMonitorSpec {
353 topic: "/scan",
354 fqn: "/perc/detector/scan",
355 max_age_ms: 100,
356 cell: &SC,
357 };
358 let mut st = AgeState::default();
359
360 SC.observe(1_000_000_000, 1_000_005_000);
362 assert!(check_age(&s, &mut st).is_none());
363 SC.observe(1_000_000_000, 1_000_250_000);
365 let v = check_age(&s, &mut st).expect("fires");
366 assert_eq!(v.rule, "max-age-runtime");
367 assert_eq!(v.fqn, "/perc/detector/scan");
368 assert_eq!(v.measured, 250);
369 assert_eq!(v.declared, 100);
370 SC.observe(1_000_000_000, 1_000_300_000);
372 assert!(check_age(&s, &mut st).is_none());
373 SC.observe(1_000_000_000, 1_000_010_000);
375 assert!(check_age(&s, &mut st).is_none());
376 SC.observe(1_000_000_000, 1_000_999_000);
377 assert!(check_age(&s, &mut st).is_some());
378 }
379
380 #[test]
381 fn peek_stamp_reads_le_time_and_rejects_unstamped() {
382 let mut raw = [0u8; 12];
384 raw[4..8].copy_from_slice(&100i32.to_le_bytes());
385 raw[8..12].copy_from_slice(&5_000u32.to_le_bytes());
386 assert_eq!(peek_stamp_us(&raw, 4), Some(100_000_005));
387 assert_eq!(peek_stamp_us(&[0u8; 12], 4), None);
389 assert_eq!(peek_stamp_us(&raw[..8], 4), None);
391 }
392
393 #[test]
394 fn slow_path_fires_latency_once_until_recovery() {
395 static C3: PubMonitorCell = PubMonitorCell::new();
396 let s = MonitorSpec {
397 topic: "/cmd",
398 fqn: "/ctrl/control/cmd",
399 min_rate_hz_milli: 0,
400 max_latency_ms: 10,
401 cell: &C3,
402 };
403 let mut st = MonitorState::default();
404 C3.max_latency_us.store(4_000, Ordering::Relaxed);
406 assert!(check_latency(&s, &mut st).is_none());
407 assert_eq!(C3.max_latency_us.load(Ordering::Relaxed), 0, "drained");
408 C3.max_latency_us.store(25_000, Ordering::Relaxed);
410 let v = check_latency(&s, &mut st).expect("fires");
411 assert_eq!(v.rule, "max-latency-runtime");
412 assert_eq!(v.measured, 25);
413 assert_eq!(v.declared, 10);
414 C3.max_latency_us.store(30_000, Ordering::Relaxed);
416 assert!(check_latency(&s, &mut st).is_none());
417 C3.max_latency_us.store(1_000, Ordering::Relaxed);
418 assert!(check_latency(&s, &mut st).is_none());
419 C3.max_latency_us.store(30_000, Ordering::Relaxed);
420 assert!(check_latency(&s, &mut st).is_some());
421 }
422}