#![cfg(all(
feature = "macros",
feature = "memory",
feature = "json",
feature = "testing"
))]
mod common;
use std::sync::{
Arc,
atomic::{AtomicU32, Ordering},
};
use common::Order;
use ruststream::memory::MemoryBroker;
use ruststream::runtime::{AppInfo, HandlerOutcome, RustStream, SubscriberSettings};
use ruststream::testing::{Outcome, TestApp};
use ruststream::{nonzero, subscriber};
#[derive(Clone, Default)]
struct Counters {
ack: Arc<AtomicU32>,
dropped: Arc<AtomicU32>,
retried: Arc<AtomicU32>,
settle: Arc<AtomicU32>,
handled: Arc<AtomicU32>,
}
impl Counters {
fn read(counter: &AtomicU32) -> u32 {
counter.load(Ordering::SeqCst)
}
}
#[subscriber("orders")]
async fn handle_order(order: &Order, ctx: &mut Context<'_, (), Counters>) -> HandlerOutcome {
let c = ctx.state().clone();
let outcome = if order.id % 2 == 1 {
HandlerOutcome::ack()
} else {
HandlerOutcome::drop()
};
let ack = Arc::clone(&c.ack);
ctx.after(HandlerOutcome::ack()).then(async move {
ack.fetch_add(1, Ordering::SeqCst);
});
let dropped = Arc::clone(&c.dropped);
ctx.after(HandlerOutcome::drop()).then(async move {
dropped.fetch_add(1, Ordering::SeqCst);
});
let retried = Arc::clone(&c.retried);
ctx.after(HandlerOutcome::retry()).then(async move {
retried.fetch_add(1, Ordering::SeqCst);
});
let settle = Arc::clone(&c.settle);
ctx.after_settle(async move {
settle.fetch_add(1, Ordering::SeqCst);
});
c.handled.fetch_add(1, Ordering::SeqCst);
outcome
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn outcome_gated_and_ungated_hooks_fire_per_settlement() {
let counters = Counters::default();
let startup_counters = counters.clone();
let app = RustStream::new(AppInfo::new("orders", "0.1.0"))
.on_startup(move |()| async move { Ok::<_, std::convert::Infallible>(startup_counters) })
.with_broker(MemoryBroker::new(), |b| b.include(handle_order));
let tb = TestApp::start(app).await.expect("startup failed");
for id in [1u32, 2] {
tb.message(&Order { id })
.to("orders")
.publish()
.await
.expect("publish");
}
assert_eq!(
tb.broker::<MemoryBroker>().subscriber("orders").outcomes(),
[Outcome::Ack, Outcome::Drop],
);
tb.drain().await;
assert_eq!(Counters::read(&counters.handled), 2);
assert_eq!(Counters::read(&counters.ack), 1);
assert_eq!(Counters::read(&counters.dropped), 1);
assert_eq!(Counters::read(&counters.settle), 2);
assert_eq!(
Counters::read(&counters.retried),
0,
"a retry-gated hook must not fire when messages are dropped",
);
}
#[subscriber("slow")]
async fn handle_slow(_order: &Order, ctx: &mut Context<'_, (), Counters>) -> HandlerOutcome {
let done = Arc::clone(&ctx.state().ack);
ctx.after_ack(async move {
tokio::task::yield_now().await;
done.fetch_add(1, Ordering::SeqCst);
});
HandlerOutcome::ack()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn hooks_drain_on_graceful_shutdown() {
let counters = Counters::default();
let startup_counters = counters.clone();
let app = RustStream::new(AppInfo::new("slow", "0.1.0"))
.on_startup(move |()| async move { Ok::<_, std::convert::Infallible>(startup_counters) })
.with_broker(MemoryBroker::new(), |b| b.include(handle_slow));
let tb = TestApp::start(app).await.expect("startup failed");
tb.message(&Order { id: 1 })
.to("slow")
.publish()
.await
.expect("publish");
tb.broker::<MemoryBroker>()
.subscriber("slow")
.assert_called_once()
.settled(HandlerOutcome::ack());
tb.shutdown().await.expect("graceful shutdown failed");
assert_eq!(Counters::read(&counters.ack), 1, "hook was not drained");
}
#[subscriber("batched")]
async fn handle_batch(orders: &[Order], ctx: &mut Context<'_, (), Counters>) -> HandlerOutcome {
let _ = orders.len();
let c = ctx.state().clone();
let settle = Arc::clone(&c.settle);
ctx.after_settle(async move {
settle.fetch_add(1, Ordering::SeqCst);
});
let gated = Arc::clone(&c.ack);
ctx.after(HandlerOutcome::ack()).then(async move {
gated.fetch_add(1, Ordering::SeqCst);
});
HandlerOutcome::ack()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn batch_runs_after_settle_drops_outcome_gated() {
let counters = Counters::default();
let startup_counters = counters.clone();
let app = RustStream::new(AppInfo::new("batched", "0.1.0"))
.on_startup(move |()| async move { Ok::<_, std::convert::Infallible>(startup_counters) })
.with_broker(MemoryBroker::new(), |b| {
b.include(handle_batch.batch(nonzero!(64)));
});
let tb = TestApp::start(app).await.expect("startup failed");
for id in 0..3u32 {
tb.message(&Order { id })
.to("batched")
.publish()
.await
.expect("publish");
}
tb.broker::<MemoryBroker>()
.subscriber("batched")
.assert_called(3)
.settled(HandlerOutcome::ack());
tb.drain().await;
assert_eq!(Counters::read(&counters.settle), 3);
assert_eq!(
Counters::read(&counters.ack),
0,
"outcome-gated hooks must not run on the batch path",
);
}