use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream};
use ruststream::subscriber;
use serde::Deserialize;
use ruststream_lapin::{LapinBroker, RabbitExchange, RabbitQueue};
#[derive(Debug, Deserialize)]
struct Order {
id: u64,
}
#[subscriber(RabbitQueue::new("orders-shard-a")
.bind(RabbitExchange::consistent_hash("orders-by-key"), "1"))]
async fn shard_a(order: &Order) -> HandlerResult {
println!("shard a: order {}", order.id);
HandlerResult::Ack
}
#[subscriber(RabbitQueue::new("orders-shard-b")
.bind(RabbitExchange::consistent_hash("orders-by-key"), "1"))]
async fn shard_b(order: &Order) -> HandlerResult {
println!("shard b: order {}", order.id);
HandlerResult::Ack
}
#[ruststream::app]
fn app() -> impl App {
let broker = LapinBroker::new("amqp://localhost:5672").declare_topology(true);
RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(broker, |b| {
b.include(shard_a);
b.include(shard_b);
})
}