use std::sync::atomic::Ordering;
use net::{ConsumeRequest, Event, EventBus, EventBusConfig, Filter};
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let bus = EventBus::new(EventBusConfig::default()).await?;
bus.ingest(Event::from_str(r#"{"token": "hello", "index": 0}"#)?)?;
bus.ingest(Event::from_str(r#"{"token": "world", "index": 1}"#)?)?;
bus.flush().await?;
let stats = bus.stats();
println!(
"ingested={} dispatched={} dropped={}",
stats.events_ingested.load(Ordering::Relaxed),
stats.events_dispatched.load(Ordering::Relaxed),
stats.events_dropped.load(Ordering::Relaxed),
);
bus.shutdown().await?;
Ok(())
}
#[allow(dead_code)]
async fn read_back(bus: &EventBus) -> Result<(), Box<dyn std::error::Error>> {
let response = bus.poll(ConsumeRequest::new(100)).await?;
for event in response.events {
println!("{}", String::from_utf8_lossy(&event.raw));
}
Ok(())
}
#[allow(dead_code)]
async fn filtered(bus: &EventBus) -> Result<(), Box<dyn std::error::Error>> {
let request = ConsumeRequest::new(100).filter(Filter::eq("token", serde_json::json!("hello")));
let _response = bus.poll(request).await?;
Ok(())
}
#[cfg(feature = "net")]
#[allow(dead_code)]
async fn on_the_mesh() -> Result<(), Box<dyn std::error::Error>> {
use net::adapter::net::{NetAdapterConfig, StaticKeypair};
use net::AdapterConfig;
let psk = [0x42u8; 32];
let responder = StaticKeypair::generate();
let adapter = NetAdapterConfig::initiator(
"0.0.0.0:7777".parse()?, "10.0.0.2:7777".parse()?, psk,
responder.public,
);
let config = EventBusConfig::builder()
.adapter(AdapterConfig::Net(Box::new(adapter)))
.build()?;
let _bus = EventBus::new(config).await?;
Ok(())
}