use std::error::Error;
use ruststream::codec::{Codec, JsonCodec};
use ruststream::memory::{MemoryBroker, MemoryPublish, MemoryPublisher};
use ruststream::runtime::{
App, AppInfo, HandlerResult, Out, Outgoing, PublishLayer, PublishNext, PublishPipeline,
PublishTransform, RustStream, Transactional, TypedPublisher,
};
use ruststream::{OutgoingMessage, Publisher, TransactionalPublisher, subscriber};
use serde::{Deserialize, Serialize};
#[derive(Debug, Deserialize)]
struct Request {
id: u64,
}
#[derive(Debug, Serialize)]
struct Response {
ok: bool,
}
#[derive(Debug, Deserialize, Serialize)]
struct Event {
id: u64,
}
#[subscriber("requests", publish("responses"))]
async fn respond(req: &Request) -> Response {
println!("responding to request {}", req.id);
Response { ok: true }
}
#[subscriber("validated-requests", publish("responses"))]
async fn validate(req: &Request) -> Result<Response, HandlerResult> {
if req.id == 0 {
return Err(HandlerResult::drop());
}
Ok(Response { ok: true })
}
#[subscriber("ingress")]
async fn forward(event: &Event, Out(out): Out<MemoryPublisher>) -> HandlerResult {
let payload = JsonCodec.encode(event).expect("serializable");
let msg = OutgoingMessage::new("egress", payload.as_ref());
if out.publish(msg).await.is_err() {
return HandlerResult::retry();
}
HandlerResult::Ack
}
#[subscriber("gateway-requests", publish("gateway-responses"))]
async fn gateway(req: &Request, Out(out): Out<MemoryPublisher>) -> Result<Response, HandlerResult> {
let audit = JsonCodec
.encode(&Event { id: req.id })
.expect("serializable");
if out
.publish(OutgoingMessage::new("gateway-audit", audit.as_ref()))
.await
.is_err()
{
return Err(HandlerResult::retry());
}
Ok(Response { ok: true })
}
struct EnvelopeTransform;
impl<C> PublishTransform<C> for EnvelopeTransform {
fn apply(&self, out: &mut Outgoing<'_>, _cx: &ruststream::runtime::PublishContext<'_, C>) {
out.headers_mut().insert("x-envelope", b"1".to_vec());
}
}
#[derive(Clone)]
struct AuditPublish;
impl PublishLayer for AuditPublish {
async fn on_publish<'a, N: PublishPipeline, P: Publisher>(
&'a self,
out: &'a mut Outgoing<'a>,
next: PublishNext<'a, N, P>,
) -> Result<(), Box<dyn Error + Send + Sync>> {
println!("publishing to {}", out.name());
next.run(out).await
}
}
#[subscriber(batch("orders"), publish("confirmations"))]
async fn confirm(orders: &[Event]) -> Result<Vec<Event>, HandlerResult> {
if orders.is_empty() {
return Err(HandlerResult::drop()); }
Ok(orders.iter().map(|o| Event { id: o.id }).collect())
}
async fn seed_events<P>(
seeder: Transactional<P, JsonCodec>,
) -> Result<(), Box<dyn Error + Send + Sync>>
where
P: TransactionalPublisher,
{
let mut scope = seeder.begin().await?;
scope.publish("events", &Event { id: 1 }).await?;
scope.publish("events", &Event { id: 2 }).await?;
scope.commit().await?;
Ok(())
}
#[ruststream::app]
fn app() -> impl App {
let broker = MemoryBroker::new();
RustStream::new(AppInfo::new("publishing", "0.1.0"))
.publish_layer(AuditPublish)
.with_broker(broker, |b| {
b.after_startup(
TypedPublisher::with_codec(MemoryPublish, JsonCodec).transactional(),
async move |seeder| seed_events(seeder).await.map_err(std::io::Error::other),
);
b.include(respond)
.publisher(TypedPublisher::new(MemoryPublish).transform(EnvelopeTransform));
b.include(validate);
b.include(forward).publisher(MemoryPublish);
b.include(gateway).out(MemoryPublish);
b.include_batch(confirm)
.publisher(TypedPublisher::new(MemoryPublish).transactional());
})
}