use std::fs::File;
use std::io::{BufWriter, Write};
use std::path::Path;
use std::sync::mpsc as std_mpsc;
use crate::intercept::{PacketData, PacketRecord};
use crate::monitor::channel::RecordSink;
pub struct FileRecorder {
sender: std_mpsc::SyncSender<String>,
dropped: std::sync::atomic::AtomicU64,
_handle: tokio::task::JoinHandle<()>,
}
impl FileRecorder {
pub fn new(path: impl AsRef<Path>) -> std::io::Result<Self> {
let file = File::create(path)?;
let writer = BufWriter::with_capacity(8192, file);
let (tx, rx) = std_mpsc::sync_channel::<String>(4096);
let handle = tokio::task::spawn_blocking(move || {
let mut w = writer;
while let Ok(line) = rx.recv() {
if writeln!(w, "{}", line).is_err() {
break;
}
}
let _ = w.flush();
});
Ok(Self {
sender: tx,
dropped: std::sync::atomic::AtomicU64::new(0),
_handle: handle,
})
}
pub fn dropped(&self) -> u64 {
self.dropped.load(std::sync::atomic::Ordering::Relaxed)
}
fn format_record(record: &PacketRecord) -> String {
let ts = crate::intercept::format_timestamp(record.timestamp_us);
let (dir, detail) = format_detail(record);
format!("{ts} {dir} {detail}")
}
}
impl RecordSink for FileRecorder {
fn on_packet(&mut self, record: PacketRecord) {
let line = Self::format_record(&record);
if self.sender.try_send(line).is_err() {
self.dropped
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
}
}
}
fn hex_str(bytes: &[u8]) -> String {
let mut s = String::with_capacity(bytes.len() * 3);
use std::fmt::Write;
for (i, b) in bytes.iter().enumerate() {
if i > 0 {
s.push(' ');
}
let _ = write!(s, "{b:02X}");
}
s
}
fn format_detail(record: &PacketRecord) -> (&'static str, String) {
match &record.data {
PacketData::RawTx(b) => ("TX", format!("{:>3}B [{}]", b.len(), hex_str(b))),
PacketData::RawRx(b) => ("RX", format!("{:>3}B [{}]", b.len(), hex_str(b))),
PacketData::RawError(b, e) => {
("ERR", format!("{:>3}B [{}] — {}", b.len(), hex_str(b), e))
}
}
}