use std::collections::VecDeque;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Mutex;
use crate::intercept::{PacketData, PacketRecord};
use crate::wire_tap::WireTap;
struct Inner {
requests: AtomicU64,
responses: AtomicU64,
errors: AtomicU64,
dropped: AtomicU64,
store: Mutex<Store>,
}
enum Store {
StatsOnly,
Unbounded(Vec<PacketRecord>),
Bounded {
buf: VecDeque<PacketRecord>,
cap: usize,
},
}
pub struct BusCapture {
inner: Inner,
}
impl BusCapture {
pub fn stats_only() -> Self {
Self {
inner: Inner {
requests: AtomicU64::new(0),
responses: AtomicU64::new(0),
errors: AtomicU64::new(0),
dropped: AtomicU64::new(0),
store: Mutex::new(Store::StatsOnly),
},
}
}
pub fn unbounded() -> Self {
Self {
inner: Inner {
requests: AtomicU64::new(0),
responses: AtomicU64::new(0),
errors: AtomicU64::new(0),
dropped: AtomicU64::new(0),
store: Mutex::new(Store::Unbounded(Vec::with_capacity(1024))),
},
}
}
pub fn bounded(capacity: usize) -> Self {
let cap = capacity.max(1);
Self {
inner: Inner {
requests: AtomicU64::new(0),
responses: AtomicU64::new(0),
errors: AtomicU64::new(0),
dropped: AtomicU64::new(0),
store: Mutex::new(Store::Bounded {
buf: VecDeque::with_capacity(cap),
cap,
}),
},
}
}
pub fn count_requests(&self) -> u64 {
self.inner.requests.load(Ordering::Relaxed)
}
pub fn count_responses(&self) -> u64 {
self.inner.responses.load(Ordering::Relaxed)
}
pub fn count_errors(&self) -> u64 {
self.inner.errors.load(Ordering::Relaxed)
}
pub fn dropped(&self) -> u64 {
self.inner.dropped.load(Ordering::Relaxed)
}
pub fn reset_stats(&self) {
let _g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
self.inner.requests.store(0, Ordering::Relaxed);
self.inner.responses.store(0, Ordering::Relaxed);
self.inner.errors.store(0, Ordering::Relaxed);
self.inner.dropped.store(0, Ordering::Relaxed);
}
pub fn len(&self) -> usize {
let g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
match &*g {
Store::StatsOnly => 0,
Store::Unbounded(v) => v.len(),
Store::Bounded { buf, .. } => buf.len(),
}
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn drain(&self) -> Vec<PacketRecord> {
let mut g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
match &mut *g {
Store::StatsOnly => Vec::new(),
Store::Unbounded(v) => std::mem::take(v),
Store::Bounded { buf, cap } => {
let old = std::mem::replace(buf, VecDeque::with_capacity(*cap));
old.into_iter().collect()
}
}
}
pub fn snapshot(&self) -> Vec<PacketRecord> {
let g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
match &*g {
Store::StatsOnly => Vec::new(),
Store::Unbounded(v) => v.clone(),
Store::Bounded { buf, .. } => buf.iter().cloned().collect(),
}
}
fn record(&self, ts: u64, data: PacketData) {
let mut g = self.inner.store.lock().unwrap_or_else(|e| e.into_inner());
let record = PacketRecord {
timestamp_us: ts,
data,
};
match &mut *g {
Store::StatsOnly => {}
Store::Unbounded(v) => v.push(record),
Store::Bounded { buf, cap } => {
if buf.len() < *cap {
buf.push_back(record);
} else {
buf.pop_front();
buf.push_back(record);
self.inner.dropped.fetch_add(1, Ordering::Relaxed);
}
}
}
}
}
impl WireTap for BusCapture {
fn on_write(&self, bytes: &[u8], ts: u64) {
self.inner.requests.fetch_add(1, Ordering::Relaxed);
self.record(ts, PacketData::RawTx(bytes.to_vec()));
}
fn on_read(&self, bytes: &[u8], ts: u64) {
self.inner.responses.fetch_add(1, Ordering::Relaxed);
self.record(ts, PacketData::RawRx(bytes.to_vec()));
}
fn on_error(&self, bytes: &[u8], error: &str, ts: u64) {
self.inner.errors.fetch_add(1, Ordering::Relaxed);
self.record(ts, PacketData::RawError(bytes.to_vec(), error.to_string()));
}
}