use std::net::{IpAddr, Ipv4Addr, SocketAddr};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::Duration;
use oms_modbus::*;
struct ByteCounter {
tx_bytes: AtomicU64,
rx_bytes: AtomicU64,
error_count: AtomicU64,
}
impl ByteCounter {
fn new() -> Self {
Self {
tx_bytes: AtomicU64::new(0),
rx_bytes: AtomicU64::new(0),
error_count: AtomicU64::new(0),
}
}
fn tx(&self) -> u64 {
self.tx_bytes.load(Ordering::Relaxed)
}
fn rx(&self) -> u64 {
self.rx_bytes.load(Ordering::Relaxed)
}
fn errors(&self) -> u64 {
self.error_count.load(Ordering::Relaxed)
}
}
impl WireTap for ByteCounter {
fn on_write(&self, bytes: &[u8], _ts: u64) {
self.tx_bytes
.fetch_add(bytes.len() as u64, Ordering::Relaxed);
}
fn on_read(&self, bytes: &[u8], _ts: u64) {
self.rx_bytes
.fetch_add(bytes.len() as u64, Ordering::Relaxed);
}
fn on_error(&self, bytes: &[u8], error: &str, _ts: u64) {
self.error_count.fetch_add(1, Ordering::Relaxed);
eprintln!(" WireTap error ({} bytes): {error}", bytes.len());
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("═══ OMS Modbus — Custom WireTap Demo ═══\n");
println!("[1/3] Starting server …");
let addr = SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0);
let server = tcp::TcpServer::bind(addr).await?;
let bind_addr = server.local_addr()?;
let store = Arc::new(SlaveStore::with_holding_registers(&[(0, 100)]));
tokio::spawn(async move {
server.serve_forever(store).await.ok();
});
tokio::time::sleep(Duration::from_millis(50)).await;
println!(" ✓ Server on {bind_addr}\n");
println!("[2/3] Connecting client with ByteCounter tap …");
let tap = Arc::new(ByteCounter::new());
let opts = ClientOptions::default()
.with_timeout(Duration::from_secs(3))
.with_tap(tap.clone());
let client = tcp::with_options(bind_addr, opts).await?;
println!(" ✓ Client connected\n");
println!("[3/3] Running Modbus traffic …\n");
client.read_holding_registers(1, 0, 1).await?;
client.write_single_register(1, 0, 42).await?;
println!("─── Results ───");
println!(" TX bytes transmitted: {}", tap.tx());
println!(" RX bytes received: {}", tap.rx());
println!(" Errors captured: {}", tap.errors());
println!(" (No error output above — successful frames are NOT logged)\n");
println!("─── Triggering error (slave=247) ───");
let _ = client.read_holding_registers(247, 0, 1).await;
println!(" Errors captured: {} ← incremented\n", tap.errors());
println!("═══ Summary ═══");
println!(" ✓ ByteCounter::on_write — count every transmitted byte");
println!(" ✓ ByteCounter::on_read — count every received byte");
println!(" ✓ ByteCounter::on_error — log only failures, not successes");
println!();
println!("Key point: WireTap hooks are driven by SniffIo at the I/O boundary.");
println!("Only implement the hooks you need — default no-ops cover the rest.");
Ok(())
}