Skip to main content

oms_modbus/
capture.rs

1// SPDX-License-Identifier: MIT OR Apache-2.0
2//!
3//! BusCapture — unified recording + statistics backend.
4//!
5//! Three modes:
6//! - `stats_only()` — atomic counters, zero memory allocation
7//! - `unbounded()` — full recording, never drops
8//! - `bounded(n)` — ring buffer, oldest evicted when full
9//!
10//! Wrap in `Arc` to share between client and test code.
11//! Cloning the `Arc` shares the same counters and records.
12
13use std::collections::VecDeque;
14use std::sync::atomic::{AtomicU64, Ordering};
15use std::sync::Mutex;
16
17use crate::intercept::{PacketData, PacketRecord};
18use crate::wire_tap::WireTap;
19
20/// Shared inner state — atomic counters + Mutex-protected store.
21struct Inner {
22    requests: AtomicU64,
23    responses: AtomicU64,
24    errors: AtomicU64,
25    dropped: AtomicU64,
26    store: Mutex<Store>,
27}
28
29enum Store {
30    StatsOnly,
31    Unbounded(Vec<PacketRecord>),
32    Bounded {
33        buf: VecDeque<PacketRecord>,
34        cap: usize,
35    },
36}
37
38/// Unified bus capture — WireTap implementation with three storage modes.
39///
40/// Wrap in `Arc` to share between client and test code:
41///
42/// # Examples
43///
44/// ```no_run
45/// use oms_modbus::*;
46/// use std::sync::Arc;
47///
48/// let cap = Arc::new(BusCapture::unbounded());
49/// let opts = ClientOptions::default().with_tap(cap.clone());
50/// ```
51pub struct BusCapture {
52    inner: Inner,
53}
54
55impl BusCapture {
56    /// Statistics only — no recording, atomic counters only.
57    /// Zero memory allocation beyond the counters themselves.
58    pub fn stats_only() -> Self {
59        Self {
60            inner: Inner {
61                requests: AtomicU64::new(0),
62                responses: AtomicU64::new(0),
63                errors: AtomicU64::new(0),
64                dropped: AtomicU64::new(0),
65                store: Mutex::new(Store::StatsOnly),
66            },
67        }
68    }
69
70    /// Full recording — stores every record, never evicts.
71    pub fn unbounded() -> Self {
72        Self {
73            inner: Inner {
74                requests: AtomicU64::new(0),
75                responses: AtomicU64::new(0),
76                errors: AtomicU64::new(0),
77                dropped: AtomicU64::new(0),
78                store: Mutex::new(Store::Unbounded(Vec::with_capacity(1024))),
79            },
80        }
81    }
82
83    /// Bounded ring buffer — oldest records evicted when full.
84    /// Records are returned in chronological order by [`Self::drain`] and [`Self::snapshot`].
85    ///
86    /// `capacity` is clamped to a minimum of 1. Passing 0 is equivalent to
87    /// passing 1 — a single-record buffer.
88    pub fn bounded(capacity: usize) -> Self {
89        let cap = capacity.max(1);
90        Self {
91            inner: Inner {
92                requests: AtomicU64::new(0),
93                responses: AtomicU64::new(0),
94                errors: AtomicU64::new(0),
95                dropped: AtomicU64::new(0),
96                store: Mutex::new(Store::Bounded {
97                    buf: VecDeque::with_capacity(cap),
98                    cap,
99                }),
100            },
101        }
102    }
103
104    // ── Statistics ──────────────────────────────────────────────────────
105
106    pub fn count_requests(&self) -> u64 {
107        self.inner.requests.load(Ordering::Relaxed)
108    }
109    pub fn count_responses(&self) -> u64 {
110        self.inner.responses.load(Ordering::Relaxed)
111    }
112    pub fn count_errors(&self) -> u64 {
113        self.inner.errors.load(Ordering::Relaxed)
114    }
115    pub fn dropped(&self) -> u64 {
116        self.inner.dropped.load(Ordering::Relaxed)
117    }
118
119    pub fn reset_stats(&self) {
120        // Acquire the store lock to ensure atomic reset across all counters —
121        // prevents concurrent readers from seeing impossible states like
122        // responses > requests.
123        let _g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
124        self.inner.requests.store(0, Ordering::Relaxed);
125        self.inner.responses.store(0, Ordering::Relaxed);
126        self.inner.errors.store(0, Ordering::Relaxed);
127        self.inner.dropped.store(0, Ordering::Relaxed);
128    }
129
130    // ── Records ─────────────────────────────────────────────────────────
131
132    pub fn len(&self) -> usize {
133        let g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
134        match &*g {
135            Store::StatsOnly => 0,
136            Store::Unbounded(v) => v.len(),
137            Store::Bounded { buf, .. } => buf.len(),
138        }
139    }
140
141    pub fn is_empty(&self) -> bool {
142        self.len() == 0
143    }
144
145    /// Drain all records in chronological order. Stats are NOT reset.
146    pub fn drain(&self) -> Vec<PacketRecord> {
147        let mut g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
148        match &mut *g {
149            Store::StatsOnly => Vec::new(),
150            Store::Unbounded(v) => std::mem::take(v),
151            Store::Bounded { buf, cap } => {
152                // Swap out the VecDeque under the lock, then drain outside.
153                let old = std::mem::replace(buf, VecDeque::with_capacity(*cap));
154                old.into_iter().collect()
155            }
156        }
157    }
158
159    /// Snapshot without clearing. Records are in chronological order.
160    ///
161    /// Clones all records under the lock. For large capture buffers in
162    /// hot paths, prefer [`Self::drain`] which swaps the buffer and releases
163    /// the lock immediately.
164    pub fn snapshot(&self) -> Vec<PacketRecord> {
165        let g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
166        match &*g {
167            Store::StatsOnly => Vec::new(),
168            Store::Unbounded(v) => v.clone(),
169            Store::Bounded { buf, .. } => buf.iter().cloned().collect(),
170        }
171    }
172
173    fn record(&self, ts: u64, data: PacketData) {
174        let mut g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
175        let record = PacketRecord {
176            timestamp_us: ts,
177            data,
178        };
179        match &mut *g {
180            Store::StatsOnly => {}
181            Store::Unbounded(v) => v.push(record),
182            Store::Bounded { buf, cap } => {
183                if buf.len() < *cap {
184                    buf.push_back(record);
185                } else {
186                    // Evict oldest, keep newest — chronological order preserved.
187                    buf.pop_front();
188                    buf.push_back(record);
189                    self.inner.dropped.fetch_add(1, Ordering::Relaxed);
190                }
191            }
192        }
193    }
194}
195
196impl WireTap for BusCapture {
197    fn on_write(&self, bytes: &[u8], ts: u64) {
198        self.inner.requests.fetch_add(1, Ordering::Relaxed);
199        self.record(ts, PacketData::RawTx(bytes.to_vec()));
200    }
201    fn on_read(&self, bytes: &[u8], ts: u64) {
202        self.inner.responses.fetch_add(1, Ordering::Relaxed);
203        self.record(ts, PacketData::RawRx(bytes.to_vec()));
204    }
205    fn on_error(&self, bytes: &[u8], error: &str, ts: u64) {
206        self.inner.errors.fetch_add(1, Ordering::Relaxed);
207        self.record(ts, PacketData::RawError(bytes.to_vec(), error.to_string()));
208    }
209}