Source code for mersal.lifespan.autosubscribe.autosubscribe_plugin

from __future__ import annotations

from copy import deepcopy
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Any

from mersal.lifespan import LifespanHandler
from mersal.logging import Logger
from mersal.plugins import Plugin

if TYPE_CHECKING:
    from mersal.configuration import StandardConfigurator
    from mersal.core.app import Mersal
    from mersal.types import LifespanHook

__all__ = (
    "AutosubscribeConfig",
    "AutosubscribePlugin",
)


[docs] @dataclass class AutosubscribeConfig: """Configure the autosubscribe plugin.""" events: set[Any] = field(default_factory=set) "A set of message types to subscribe to." @property def plugin(self) -> AutosubscribePlugin: return AutosubscribePlugin(self)
class AutosubscribePlugin(Plugin): """Autosubscribe to a given event types once the app is started.""" def __init__(self, config: AutosubscribeConfig) -> None: self._events = config.events def __call__(self, configurator: StandardConfigurator) -> None: def decorate(configurator: StandardConfigurator) -> Any: lifespan_handler: LifespanHandler = configurator.get(LifespanHandler) # type: ignore[type-abstract] if configurator.send_only: # In send-only mode there's no worker/input queue to ever receive # what gets routed to this app's address, so subscribing would # just register a subscriber address that nothing drains. lifespan_handler.register_on_startup_hook(self._log_skip(configurator)) else: app: Mersal = configurator.mersal lifespan_handler.register_on_startup_hook(self._subscribe(app)) return lifespan_handler configurator.decorate(LifespanHandler, decorate) def _subscribe(self, app: Mersal) -> LifespanHook: events = deepcopy(self._events) async def subscribe( events: set[Any] = events, ) -> None: for e in events: await app.subscribe(e) return subscribe def _log_skip(self, configurator: StandardConfigurator) -> LifespanHook: async def log_skip() -> None: configurator.get(Logger).info( # type: ignore[type-abstract] "autosubscribe.send_only.skip", reason="app is send_only; there is no worker to receive routed events", ) return log_skip