#[macro_export]
macro_rules! publisher {
($transport:ty: [$($event:ty),* $(,)?]) => {
$(
$crate::__publisher_one!($transport, $event);
)*
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __publisher_one {
($transport:ty, $event:ty) => {
const _: () = {
use $crate::events::{Event, EventPublishRaw, EventPublisher, json};
#[async_trait::async_trait]
impl EventPublisher<$event> for $transport {
async fn publish(&self, event: $event) -> $crate::events::NoemaResult<()> {
let payload = json::to_bytes(&event)?;
self.publish_raw(<$event as Event>::WIRE_NAME, &payload)
.await
}
}
};
};
}
#[macro_export]
macro_rules! subscribe {
($transport:ty, $($event:ty: [$($handler:ty),* $(,)?]),* $(,)?) => {
$crate::__subscribe_batch!($transport, [ $( ($event, [$($handler),*]) ),* ]);
};
}
#[doc(hidden)]
#[macro_export]
macro_rules! __subscribe_batch {
($transport:ty, [ $( ($event:ty, [$($handler:ty),*]) ),* $(,)? ]) => {
const _: () = {
use std::future::Future;
use std::pin::Pin;
use std::sync::Arc;
use $crate::paste::paste;
use $crate::core::{Container, Injectable};
use $crate::events::dispatch::{
DispatchContext, EventDispatcherContext, ReceiveFn, SubscriberEntry,
SubscriberRegistry,
};
use $crate::events::{Event, EventListener, json, NoemaResult};
fn __noema_require_dispatcher_context<T: EventDispatcherContext>() {}
let _ = __noema_require_dispatcher_context::<$transport>;
$(
paste! {
#[allow(non_snake_case)]
fn [<__noema_receive_ $event>](
payload: &[u8],
ctx: &DispatchContext,
) -> Pin<Box<dyn Future<Output = NoemaResult<()>> + Send>> {
let payload = payload.to_vec();
let spawner = ctx.spawner.clone();
let error_handler = ctx.error_handler.clone();
let invoke_mode = ctx.invoke_mode;
Box::pin(async move {
let ctx = DispatchContext {
spawner,
error_handler,
invoke_mode,
};
let event: Arc<$event> = Arc::new(json::from_bytes(&payload)?);
$(
{
let handler =
<$handler as Injectable<Container>>::inject(&Container);
let handler = Arc::new(handler)
as Arc<dyn EventListener<$event> + Send + Sync>;
let spawner = ctx.spawner.clone();
let error_handler = ctx.error_handler.clone();
let event_clone = event.clone();
let event_err = event.clone();
match ctx.invoke_mode {
$crate::events::InvokeMode::Spawn => {
spawner.spawn(Box::pin(async move {
if let Err(err) = handler.handle(event_clone).await {
handler
.on_error(error_handler, err, event_err)
.await;
}
}));
}
$crate::events::InvokeMode::Await => {
let _ = error_handler;
let _ = event_err;
handler.handle(event_clone).await?;
}
}
}
)*
Ok(())
})
}
}
)*
static __NOEMA_SUBSCRIBER_ENTRIES: &[SubscriberEntry] = &[
$(
SubscriberEntry {
name: <$event as Event>::WIRE_NAME,
receive: paste! { [<__noema_receive_ $event>] },
},
)*
];
impl SubscriberRegistry for $transport {
fn entries(&self) -> &'static [SubscriberEntry] {
__NOEMA_SUBSCRIBER_ENTRIES
}
}
};
};
}