use std::collections::VecDeque;
use std::sync::Arc;
use std::time::Duration;
use axum::extract::State;
use axum::http::StatusCode;
use axum::routing::{get, post};
use axum::{Json, Router};
use ruststream::codec::{Codec, JsonCodec};
use ruststream::memory::{MemoryBroker, MemoryPublish, MemoryPublisher};
use ruststream::runtime::{AppInfo, HandlerResult, HealthProbe, HealthState, RustStream};
use ruststream::{Broker, OutgoingMessage, Publisher, subscriber};
use serde::{Deserialize, Serialize};
use tokio::sync::Mutex;
#[derive(Debug, Clone, Serialize, Deserialize)]
struct OrderPlaced {
id: u64,
item: String,
}
#[subscriber("orders.placed")]
async fn fulfil(order: &OrderPlaced) -> HandlerResult {
println!("fulfilling order {} ({})", order.id, order.item);
HandlerResult::Ack
}
#[derive(Default)]
struct Store {
orders: Vec<OrderPlaced>,
outbox: VecDeque<OrderPlaced>,
}
async fn place_order(
State(store): State<Arc<Mutex<Store>>>,
Json(order): Json<OrderPlaced>,
) -> &'static str {
let mut store = store.lock().await;
store.orders.push(order.clone());
store.outbox.push_back(order);
"accepted\n"
}
async fn relay_outbox(store: Arc<Mutex<Store>>, egress: MemoryPublisher) {
let mut tick = tokio::time::interval(Duration::from_millis(200));
loop {
tick.tick().await;
loop {
let Some(event) = store.lock().await.outbox.front().cloned() else {
break;
};
let payload = JsonCodec.encode(&event).expect("serializable");
let out = OutgoingMessage::new("orders.placed", payload.as_ref());
if egress.publish(out).await.is_err() {
break;
}
store.lock().await.outbox.pop_front();
}
}
}
#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
async fn healthz(State(health): State<HealthProbe>) -> (StatusCode, String) {
match health.state() {
HealthState::Running => (StatusCode::OK, "running".to_owned()),
state => (StatusCode::SERVICE_UNAVAILABLE, format!("{state:?}")),
}
}
let broker = MemoryBroker::new().bindable();
let egress = broker.bind(MemoryPublish);
let app = RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(broker, |b| {
b.include(fulfil);
});
let running = app.start().await?;
let egress = running.publisher(egress).await?;
let store = Arc::new(Mutex::new(Store::default()));
tokio::spawn(relay_outbox(store.clone(), egress));
let router = Router::new()
.route("/orders", post(place_order))
.with_state(store)
.route("/healthz", get(healthz).with_state(running.health()));
let listener = tokio::net::TcpListener::bind("127.0.0.1:8080").await?;
println!("orders API on http://127.0.0.1:8080/orders");
let stopping = running.stopping();
axum::serve(listener, router)
.with_graceful_shutdown(async move {
tokio::select! {
_ = tokio::signal::ctrl_c() => {}
() = stopping => {}
}
})
.await?;
running.shutdown().await?;
Ok(())
}