use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
use flowscope::PacketView;
use flowscope::driver::{Driver, Event, SlotMessage};
use flowscope::extract::{FiveTuple, FiveTupleKey};
use flowscope::http::{HttpMessage, HttpParser};
use flowscope::pcap::PcapFlowSource;
fn main() -> Result<(), Box<dyn std::error::Error>> {
let path = std::env::args()
.nth(1)
.unwrap_or_else(|| "tests/data/http_session.pcap".to_string());
let mut builder = Driver::builder(FiveTuple::bidirectional());
let mut http_slot = builder.session_on_ports(HttpParser::default(), [80, 8080]);
let mut driver = builder.build();
let drainer_handle = http_slot.clone();
let done = Arc::new(AtomicBool::new(false));
let done_for_worker = Arc::clone(&done);
let (tx, rx) = mpsc::channel::<SlotMessage<HttpMessage, FiveTupleKey>>();
let worker = thread::spawn(move || {
let mut h = drainer_handle;
let mut buf = Vec::new();
while !done_for_worker.load(Ordering::Acquire) || h.pending() > 0 {
h.drain(&mut buf);
for m in buf.drain(..) {
if tx.send(m).is_err() {
return; }
}
thread::sleep(Duration::from_micros(50));
}
});
let mut events: Vec<Event<FiveTupleKey>> = Vec::new();
let mut flows = 0usize;
for owned in PcapFlowSource::open(&path)?.views() {
let owned = owned?;
events.clear();
driver.track_into(PacketView::from(&owned), &mut events);
for ev in &events {
if matches!(ev, Event::FlowStarted { .. }) {
flows += 1;
}
}
}
driver.finish_into(&mut events);
done.store(true, Ordering::Release);
worker.join().expect("worker panicked");
let mut messages = 0usize;
while let Ok(msg) = rx.try_recv() {
let _ = msg.message; messages += 1;
}
let mut leftover = Vec::new();
http_slot.drain(&mut leftover);
messages += leftover.len();
println!("flows started: {flows}");
println!("http messages observed: {messages}");
Ok(())
}