use ruststream::conformance::harness;
use ruststream::runtime::{AppInfo, HandlerResult, RustStream};
use ruststream::subscriber;
use ruststream::testing::TestApp;
use ruststream_lapin::RabbitQueue;
use ruststream_lapin::testing::LapinTestBroker;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize, PartialEq)]
struct Payment {
amount: u64,
}
#[subscriber(RabbitQueue::new("payments"))]
async fn accept(payment: &Payment) -> HandlerResult {
if payment.amount == 0 {
return HandlerResult::drop();
}
HandlerResult::Ack
}
#[tokio::main(flavor = "multi_thread", worker_threads = 2)]
async fn main() {
let app = RustStream::new(AppInfo::new("payments", "0.1.0")).with_broker(
LapinTestBroker::new(),
|b| {
b.include(accept);
},
);
let tb = TestApp::start(app).await.expect("start");
tb.broker::<LapinTestBroker>()
.publish("payments", &Payment { amount: 100 })
.await
.expect("publish drives the handler to quiescence");
tb.broker::<LapinTestBroker>()
.subscriber("payments")
.assert_called_once()
.with(&Payment { amount: 100 })
.settled(HandlerResult::Ack);
tb.shutdown().await.expect("shutdown");
harness::run_suite(LapinTestBroker::new).await;
println!("all in-process checks passed");
}