Skip to content

Message Bus Event Publisher

message_bus_event_publisher

Message-bus-backed EventPublisherPort adapter.

Delegates event publishing to an injected MessageBusPort.

MessageBusEventPublisher

Bases: EventPublisherPort[EventPayloadType]

Infrastructure adapter that publishes events via a MessageBusPort.

Implements EventPublisherPort by delegating publish to
MessageBusPort.dispatch.

Example
class Event[T]:
    def __init__(self) -> None:
        pass


class InMemoryMessageBus[MT, R]:
    async def dispatch(self, message: MT) -> R: ...
    def register(self, message_type, handler): ...


class MyEvent(Event[dict[str, object]]):
    def __init__(self) -> None:
        pass


bus = InMemoryMessageBus[MyEvent, None]()
publisher = MessageBusEventPublisher[dict[str, object]](bus)
await publisher.publish(MyEvent())
Source code in src/forging_blocks/infrastructure/message_bus/message_bus_event_publisher.py
class MessageBusEventPublisher[EventPayloadType](EventPublisherPort[EventPayloadType]):
    """Infrastructure adapter that publishes events via a ``MessageBusPort``.

    Implements ``EventPublisherPort`` by delegating ``publish`` to
    ``MessageBusPort.dispatch``.

    Example:
        ```python
        class Event[T]:
            def __init__(self) -> None:
                pass


        class InMemoryMessageBus[MT, R]:
            async def dispatch(self, message: MT) -> R: ...
            def register(self, message_type, handler): ...


        class MyEvent(Event[dict[str, object]]):
            def __init__(self) -> None:
                pass


        bus = InMemoryMessageBus[MyEvent, None]()
        publisher = MessageBusEventPublisher[dict[str, object]](bus)
        await publisher.publish(MyEvent())
        ```
    """

    def __init__(self, message_bus: MessageBusPort[Event[EventPayloadType], None]) -> None:
        self._message_bus = message_bus

    async def publish(self, event: Event[EventPayloadType]) -> None:
        """Publish a domain event via the message bus."""
        await self._message_bus.dispatch(event)

publish(event: Event[EventPayloadType]) -> None async

Publish a domain event via the message bus.

Source code in src/forging_blocks/infrastructure/message_bus/message_bus_event_publisher.py
async def publish(self, event: Event[EventPayloadType]) -> None:
    """Publish a domain event via the message bus."""
    await self._message_bus.dispatch(event)