#[cfg(test)]
mod tests {
use std::future::Future;
use std::sync::{Arc, Mutex};
use async_trait::async_trait;
use futures::executor::block_on;
use serde::{Deserialize, Serialize};
use crate::arc_dyn;
use crate::core::*;
use crate::events::*;
use crate::resolve;
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
#[event(name = "test.created")]
struct TestCreated {
id: u32,
}
struct SyncSpawner;
impl BackgroundSpawner for SyncSpawner {
fn spawn(&self, fut: std::pin::Pin<Box<dyn Future<Output = ()> + Send + 'static>>) {
std::thread::spawn(move || futures::executor::block_on(fut))
.join()
.expect("handler thread panicked");
}
}
struct PanicErrorHandler;
#[async_trait]
impl BackgroundErrorHandler for PanicErrorHandler {
async fn handle(&self, error: BoxDynError, event_name: &str, handler_name: String) {
panic!("unexpected error event={event_name} handler={handler_name}: {error}");
}
}
fn default_spawner() -> Arc<dyn BackgroundSpawner> {
Arc::new(SyncSpawner)
}
fn panic_on_error() -> Arc<dyn BackgroundErrorHandler> {
Arc::new(PanicErrorHandler)
}
mod round_trip {
use super::*;
struct MockWireTransport {
spawner: Arc<dyn BackgroundSpawner>,
errors: Arc<dyn BackgroundErrorHandler>,
outbox: Arc<Mutex<Vec<(String, Vec<u8>)>>>,
handler_log: Arc<Mutex<Vec<u32>>>,
}
impl Injectable<Container> for MockWireTransport {
fn inject(_: &Container) -> Self {
MockWireTransport {
spawner: default_spawner(),
errors: panic_on_error(),
outbox: Arc::new(Mutex::new(Vec::new())),
handler_log: Arc::new(Mutex::new(Vec::new())),
}
}
}
#[async_trait]
impl EventPublishRaw for MockWireTransport {
async fn publish_raw(&self, name: &str, payload: &[u8]) -> NoemaResult<()> {
self.outbox
.lock()
.unwrap()
.push((name.to_string(), payload.to_vec()));
Ok(())
}
}
impl EventDispatcherContext for MockWireTransport {
fn dispatch_context(&self) -> DispatchContext {
DispatchContext::new(self.spawner.clone(), self.errors.clone())
}
}
struct RoundTripHandler {
log: Arc<Mutex<Vec<u32>>>,
}
#[async_trait]
impl EventListener<TestCreated> for RoundTripHandler {
async fn handle(&self, event: Arc<TestCreated>) -> NoemaResult<()> {
self.log.lock().unwrap().push(event.id);
Ok(())
}
}
impl Injectable<Container> for RoundTripHandler {
fn inject(_: &Container) -> Self {
RoundTripHandler {
log: resolve::<MockWireTransport>().handler_log.clone(),
}
}
}
crate::publisher!(MockWireTransport: [TestCreated]);
crate::subscribe!(MockWireTransport, TestCreated: [RoundTripHandler]);
crate::dependency!(singleton, MockWireTransport);
mod publish {
use super::*;
pub async fn test_created(id: u32) {
let publisher: arc_dyn!(EventPublisher<TestCreated>) =
resolve::<MockWireTransport>();
publisher.publish(TestCreated { id }).await.unwrap();
}
}
mod receive {
use super::*;
pub async fn dispatch_wire(name: &str, payload: &[u8]) {
let dispatcher: arc_dyn!(EventDispatch) = resolve::<MockWireTransport>();
dispatcher.dispatch(name, payload).await.unwrap();
}
}
#[test]
fn publish_subscribe_json_round_trip() {
block_on(publish::test_created(42));
let (wire_name, wire_payload) = resolve::<MockWireTransport>()
.outbox
.lock()
.unwrap()
.pop()
.expect("expected one published message");
assert_eq!(wire_name, "test.created");
block_on(receive::dispatch_wire(&wire_name, &wire_payload));
assert_eq!(
*resolve::<MockWireTransport>().handler_log.lock().unwrap(),
vec![42]
);
}
}
mod multi_handler {
use super::*;
struct MockWireTransport {
spawner: Arc<dyn BackgroundSpawner>,
errors: Arc<dyn BackgroundErrorHandler>,
handler_a_log: Arc<Mutex<Vec<u32>>>,
handler_b_log: Arc<Mutex<Vec<u32>>>,
}
impl Injectable<Container> for MockWireTransport {
fn inject(_: &Container) -> Self {
MockWireTransport {
spawner: default_spawner(),
errors: panic_on_error(),
handler_a_log: Arc::new(Mutex::new(Vec::new())),
handler_b_log: Arc::new(Mutex::new(Vec::new())),
}
}
}
impl EventDispatcherContext for MockWireTransport {
fn dispatch_context(&self) -> DispatchContext {
DispatchContext::new(self.spawner.clone(), self.errors.clone())
}
}
struct HandlerA {
log: Arc<Mutex<Vec<u32>>>,
}
struct HandlerB {
log: Arc<Mutex<Vec<u32>>>,
}
#[async_trait]
impl EventListener<TestCreated> for HandlerA {
async fn handle(&self, event: Arc<TestCreated>) -> NoemaResult<()> {
self.log.lock().unwrap().push(event.id);
Ok(())
}
}
#[async_trait]
impl EventListener<TestCreated> for HandlerB {
async fn handle(&self, event: Arc<TestCreated>) -> NoemaResult<()> {
self.log.lock().unwrap().push(event.id + 1000);
Ok(())
}
}
impl Injectable<Container> for HandlerA {
fn inject(_: &Container) -> Self {
HandlerA {
log: resolve::<MockWireTransport>().handler_a_log.clone(),
}
}
}
impl Injectable<Container> for HandlerB {
fn inject(_: &Container) -> Self {
HandlerB {
log: resolve::<MockWireTransport>().handler_b_log.clone(),
}
}
}
crate::subscribe!(MockWireTransport, TestCreated: [HandlerA, HandlerB]);
crate::dependency!(singleton, MockWireTransport);
mod receive {
use super::*;
pub async fn dispatch_test_created(id: u32) {
let payload = json::to_bytes(&TestCreated { id }).unwrap();
let dispatcher: arc_dyn!(EventDispatch) = resolve::<MockWireTransport>();
dispatcher.dispatch("test.created", &payload).await.unwrap();
}
}
#[test]
fn multiple_handlers() {
block_on(receive::dispatch_test_created(7));
let transport = resolve::<MockWireTransport>();
assert_eq!(*transport.handler_a_log.lock().unwrap(), vec![7]);
assert_eq!(*transport.handler_b_log.lock().unwrap(), vec![1007]);
}
}
mod unknown_event {
use super::*;
struct CaptureError {
last: Arc<Mutex<Option<String>>>,
}
#[async_trait]
impl BackgroundErrorHandler for CaptureError {
async fn handle(&self, error: BoxDynError, event_name: &str, _: String) {
*self.last.lock().unwrap() = Some(format!("{event_name}: {error}"));
}
}
struct MockWireTransport {
spawner: Arc<dyn BackgroundSpawner>,
errors: Arc<CaptureError>,
last_dispatch_error: Arc<Mutex<Option<String>>>,
}
impl Injectable<Container> for MockWireTransport {
fn inject(_: &Container) -> Self {
let last = Arc::new(Mutex::new(None));
MockWireTransport {
spawner: default_spawner(),
errors: Arc::new(CaptureError { last: last.clone() }),
last_dispatch_error: last,
}
}
}
impl EventDispatcherContext for MockWireTransport {
fn dispatch_context(&self) -> DispatchContext {
DispatchContext::new(self.spawner.clone(), self.errors.clone())
}
}
struct NoopHandler;
#[async_trait]
impl EventListener<TestCreated> for NoopHandler {
async fn handle(&self, _: Arc<TestCreated>) -> NoemaResult<()> {
Ok(())
}
}
impl Injectable<Container> for NoopHandler {
fn inject(_: &Container) -> Self {
NoopHandler
}
}
crate::subscribe!(MockWireTransport, TestCreated: [NoopHandler]);
crate::dependency!(singleton, MockWireTransport);
mod receive {
use super::*;
pub async fn dispatch_unknown() {
let dispatcher: arc_dyn!(EventDispatch) = resolve::<MockWireTransport>();
dispatcher
.dispatch("unknown.event", b"{}")
.await
.unwrap_err();
}
}
#[test]
fn unknown_event_invokes_error_handler() {
block_on(receive::dispatch_unknown());
assert!(
resolve::<MockWireTransport>()
.last_dispatch_error
.lock()
.unwrap()
.as_ref()
.unwrap()
.contains("unknown.event")
);
}
}
}