use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream, TypedPublisher};
use ruststream::{OutgoingMessage, OwnedTransactions, Transaction, subscriber};
use ruststream_fred::{RedisBroker, RedisPublish};
use serde::{Deserialize, Serialize};
#[derive(Debug, Deserialize, Serialize)]
struct Order {
id: u64,
}
#[subscriber(batch("orders"), publish("processed"))]
async fn process(orders: &[Order]) -> Result<Vec<Order>, HandlerResult> {
if orders.is_empty() {
return Err(HandlerResult::drop());
}
Ok(orders.iter().map(|o| Order { id: o.id }).collect())
}
#[ruststream::app]
fn app() -> impl App {
let broker = RedisBroker::standalone("redis://localhost:6379").default_group("workers");
RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(broker, |b| {
b.include_batch(process)
.publisher(TypedPublisher::new(RedisPublish).transactional());
b.after_startup(RedisPublish, async move |publisher| {
let mut seed = publisher.transaction().await?;
seed.publish(OutgoingMessage::new("processed", br#"{"id":0}"#.as_slice()))
.await?;
seed.commit().await
});
})
}