use std::env;
use std::thread;
use std::time::Duration;
use tokio::net::TcpStream;
use buswatch_tui::{DataSource, StreamSource};
#[tokio::main]
async fn main() {
let addr = env::args().nth(1).unwrap_or_else(|| {
eprintln!("Usage: cargo run --example stream_source -- <host:port>");
eprintln!();
eprintln!("Example: cargo run --example stream_source -- localhost:9090");
std::process::exit(1);
});
println!("Connecting to {}...", addr);
let stream = match TcpStream::connect(&addr).await {
Ok(s) => s,
Err(e) => {
eprintln!("Failed to connect to {}: {}", addr, e);
std::process::exit(1);
}
};
println!("Connected! Waiting for snapshots...\n");
let mut source = StreamSource::spawn(stream, &addr);
loop {
match source.poll() {
Some(snapshot) => {
println!("Received snapshot with {} modules:", snapshot.len());
for (name, state) in snapshot.iter() {
println!(
" - {}: {} read topics, {} write topics",
name,
state.reads.len(),
state.writes.len()
);
}
println!();
}
None => {
if let Some(err) = source.error() {
eprintln!("Error: {}", err);
break;
}
}
}
thread::sleep(Duration::from_millis(100));
}
}