use std::time::Duration;
use ruststream::codec::{Codec, JsonCodec};
use ruststream::runtime::{App, AppInfo, HandlerResult, RustStream, State};
use ruststream::{FromRef, IncomingMessage, OutgoingMessage, RequestReply, subscriber};
use ruststream_lapin::{LapinBroker, LapinRequester};
use serde::{Deserialize, Serialize};
#[derive(Debug, Deserialize)]
struct Order {
sku: String,
quantity: u32,
}
#[derive(Debug, Serialize)]
struct CheckStock {
sku: String,
quantity: u32,
}
#[derive(Debug, Deserialize)]
struct Stock {
#[allow(dead_code)]
sku: String,
available: bool,
}
#[derive(Clone)]
struct Inventory {
requester: LapinRequester,
}
impl Inventory {
async fn check(
&self,
sku: &str,
quantity: u32,
) -> Result<Stock, Box<dyn std::error::Error + Send + Sync>> {
let request = CheckStock {
sku: sku.to_owned(),
quantity,
};
let payload = JsonCodec.encode(&request)?;
let reply = self
.requester
.request(
OutgoingMessage::new("inventory.check", payload.as_ref()),
Duration::from_secs(2),
)
.await?;
Ok(JsonCodec.decode(reply.payload())?)
}
}
#[derive(Clone, FromRef)]
struct AppState {
inventory: Inventory,
}
#[subscriber("orders")]
async fn place_order(order: &Order, State(inventory): State<Inventory>) -> HandlerResult {
match inventory.check(&order.sku, order.quantity).await {
Ok(stock) if stock.available => {
println!("order accepted: {} x{}", order.sku, order.quantity);
HandlerResult::Ack
}
Ok(_) => {
println!(
"order rejected, out of stock: {} x{}",
order.sku, order.quantity
);
HandlerResult::Ack
}
Err(err) => {
eprintln!("inventory unavailable, retrying later: {err}");
HandlerResult::retry()
}
}
}
#[ruststream::app]
fn app() -> impl App {
let broker = LapinBroker::new("amqp://localhost:5672").declare_topology(true);
let inventory = Inventory {
requester: broker.requester(),
};
RustStream::new(AppInfo::new("orders", "0.1.0"))
.on_startup(
move |()| async move { Ok::<_, std::convert::Infallible>(AppState { inventory }) },
)
.with_broker(broker, |b| {
b.include(place_order);
})
}