use ruststream::memory::{MemoryBroker, MemoryMessage};
use ruststream::runtime::{AppInfo, Ctx, HandlerResult, RustStream};
use ruststream::testing::TestApp;
use ruststream::{BuildContext, ContextField, IncomingMessage, subscriber};
use serde::{Deserialize, Serialize};
#[derive(Serialize, Deserialize, PartialEq, Debug)]
struct Order {
id: u64,
}
struct DeliveryMeta {
payload_len: usize,
}
impl BuildContext<MemoryMessage> for DeliveryMeta {
fn build(msg: &MemoryMessage) -> Self {
Self {
payload_len: msg.payload().len(),
}
}
}
#[derive(Clone, Copy, Default)]
struct PayloadLen;
impl ContextField for PayloadLen {
type Context = DeliveryMeta;
type Value = usize;
fn read(self, src: &DeliveryMeta) -> usize {
src.payload_len
}
}
#[subscriber("orders")]
async fn audit(order: &Order, Ctx(len): Ctx<PayloadLen>) -> HandlerResult {
println!("order {} arrived as {len} bytes", order.id);
HandlerResult::Ack
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
let app = RustStream::new(AppInfo::new("orders", "0.1.0"))
.with_broker(MemoryBroker::new(), |b| b.include(audit));
let tb = TestApp::start(app).await?;
tb.broker::<MemoryBroker>()
.publish("orders", &Order { id: 40 })
.await?;
println!("done");
Ok(())
}