use std::env;
use std::fs::File;
use std::io::BufWriter;
use std::time::Duration;
use futures::StreamExt;
use netring::flow::FlowEvent;
use netring::flow::extract::FiveTuple;
use netring::pcap::CaptureWriter;
use netring::{AsyncCapture, Dedup, StreamCapture};
#[tokio::main(flavor = "current_thread")]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let mut args = env::args().skip(1);
let iface = args.next().unwrap_or_else(|| "lo".to_string());
let out_path = args.next().unwrap_or_else(|| "capture.pcap".to_string());
println!("listening on {iface}, recording to {out_path} (Ctrl+C to stop)...");
let writer = CaptureWriter::create(BufWriter::new(File::create(&out_path)?))?;
let cap = AsyncCapture::open(&iface)?;
let mut stream = cap
.flow_stream(FiveTuple::bidirectional())
.with_dedup(Dedup::loopback())
.with_pcap_tap(writer);
let mut stats_tick = tokio::time::interval(Duration::from_secs(1));
stats_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
stats_tick.tick().await;
loop {
tokio::select! {
biased;
_ = stats_tick.tick() => {
if let Ok(stats) = stream.capture_stats() {
eprintln!(
"[stats] packets={} drops={} freeze_count={}",
stats.packets, stats.drops, stats.freeze_count,
);
}
}
evt = stream.next() => match evt {
Some(Ok(FlowEvent::Started { key, l4, .. })) => {
println!("+ {l4:?} {a} <-> {b}", a = key.a, b = key.b);
}
Some(Ok(FlowEvent::Ended { key, reason, stats, .. })) => {
println!(
"- {a} <-> {b} reason={reason:?} pkts={p}",
a = key.a, b = key.b,
p = stats.packets_initiator + stats.packets_responder,
);
}
Some(Ok(_)) => { }
Some(Err(e)) => {
eprintln!("stream error: {e}");
break;
}
None => break,
}
}
}
Ok(())
}