#![cfg(all(
feature = "memory",
feature = "json",
feature = "macros",
feature = "testing"
))]
use ruststream::memory::prelude::*;
use ruststream::testing::TestApp;
use serde::{Deserialize, Serialize};
#[derive(Debug, Outgoing, Serialize, Deserialize, schemars::JsonSchema)]
struct Order {
id: u64,
}
#[derive(Debug, Clone, PartialEq, Eq, Outgoing, Serialize, Deserialize)]
#[outgoing(name = "orders.settled")]
struct Settled {
id: u64,
}
#[derive(OutSlot)]
#[publishes(Settled)]
struct Journal;
#[derive(OutSlot)]
#[publishes(Settled)]
struct Audit;
#[ruststream::subscriber("orders")]
async fn settle(
order: &Order,
Out(journal): Out<impl TransactionalPublisher, Journal, Settled>,
Out(audit): Out<impl Publisher, Audit, Settled>,
) -> HandlerOutcome {
let Ok(scope) = journal.begin().await else {
return HandlerOutcome::retry();
};
if scope
.message(&Settled { id: order.id })
.publish()
.await
.is_err()
|| scope.commit().await.is_err()
|| audit
.message(&Settled { id: order.id })
.publish()
.await
.is_err()
{
return HandlerOutcome::retry();
}
HandlerOutcome::ack()
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn the_mount_site_names_policies_under_the_uniform_names() {
let app =
RustStream::new(AppInfo::new("orders", "0.1.0")).with_broker(MemoryBroker::new(), |b| {
b.include(settle)
.out(Journal, TransactionalPublish)
.out(Audit, Publish)
.build();
});
let tb = TestApp::start(app).await.expect("harness start");
tb.message(&Order { id: 7 })
.to("orders")
.publish()
.await
.expect("publish");
tb.broker::<MemoryBroker>()
.subscriber("orders")
.assert_called_once()
.settled(HandlerOutcome::ack());
tb.out::<Journal>().assert_called_once();
tb.out::<Audit>().assert_called_once();
}
fn publisher_is_the_core_trait<T: Publisher>() {}
#[test]
fn the_policy_names_are_types_and_the_capability_names_are_traits() {
let _: Publish = Publish;
let _: TransactionalPublish = TransactionalPublish;
let _: Request = Request;
publisher_is_the_core_trait::<ruststream::memory::MemoryPublisher>();
publisher_is_the_core_trait::<ruststream::memory::MemoryRequester>();
}