noema 0.4.0

Noema IOC and DI framework for Rust
Documentation
/// Register event types that a transport can publish (typed → JSON → `publish_raw`).
///
/// # Syntax
///
/// ```ignore
/// publisher!(OrdersKafka: [OrderCreated, OrderShipped]);
/// ```
#[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
                }
            }
        };
    };
}

/// Register handlers for events on a transport (single batch per transport type per crate).
///
/// Your transport MUST implement [`EventDispatcherContext`](crate::events::EventDispatcherContext)
/// explicitly (spawner + error handler). `subscribe!` checks that at compile time.
///
/// # Syntax
///
/// ```ignore
/// subscribe!(
///     OrdersKafka,
///     OrderCreated: [EmailHandler, AuditHandler],
///     OrderShipped: [NotifyHandler],
/// );
/// ```
#[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 => {
                                            // Propagate so the caller (WS `connect`) can reply.
                                            // Spawn still uses `on_error` (fire-and-forget).
                                            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
                }
            }
        };
    };
}