use std::time::Duration;
use netring::prelude::*;
#[derive(Default)]
struct FlowCounters {
started: u64,
ended: u64,
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let iface = std::env::args().nth(1).unwrap_or_else(|| "lo".into());
let dur_secs: u64 = std::env::args()
.nth(2)
.and_then(|s| s.parse().ok())
.unwrap_or(30);
eprintln!("monitor_basic: capturing on {iface} for {dur_secs}s");
Monitor::builder()
.interface(&iface)
.protocol::<Tcp>() .state::<FlowCounters>()
.on_ctx::<FlowStarted<Tcp>>(|_evt: &FlowStarted<Tcp>, ctx: &mut Ctx<'_>| {
let ts = ctx.ts;
let (counters, sink) = ctx.split_state_sink::<FlowCounters>();
counters.started += 1;
sink.begin("FlowStartedTcp", Severity::Info, ts)
.with_metric("started_total", counters.started as f64)
.emit();
Ok(())
})
.on_ctx::<FlowEnded<Tcp>>(|_evt: &FlowEnded<Tcp>, ctx: &mut Ctx<'_>| {
ctx.state_mut::<FlowCounters>().ended += 1;
Ok(())
})
.sink(StdoutSink::default())
.build()?
.run_for(Duration::from_secs(dur_secs))
.await?;
eprintln!("monitor_basic: done");
Ok(())
}