use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream};
use ruststream::subscriber;
use ruststream_rdkafka::{KafkaBroker, KafkaTopic, StartOffset};
use serde::Deserialize;
#[derive(Debug, Deserialize)]
struct OrderEvent {
id: u64,
}
#[subscriber(KafkaTopic::new("orders").and_topic("cancellations").group("orders-svc"))]
async fn on_order_event(event: &OrderEvent) -> HandlerResult {
println!("order event {}", event.id);
HandlerResult::Ack
}
#[subscriber(KafkaTopic::pattern("^audit\\..*").group("audit-svc").start(StartOffset::Earliest))]
async fn on_audit(event: &OrderEvent) -> HandlerResult {
println!("audit event {}", event.id);
HandlerResult::Ack
}
#[ruststream::app]
fn app() -> impl App {
let broker = KafkaBroker::new(["localhost:9092"]);
RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(broker, |b| {
b.include(on_order_event);
b.include(on_audit);
})
}