Source code for mersal.timeouts.plugin
from __future__ import annotations
from collections.abc import Callable
from typing import TYPE_CHECKING
from mersal.configuration.standard_configurator import InvalidConfigurationError
from mersal.lifespan.lifespan_hooks_registration_plugin import (
LifespanHooksRegistrationPluginConfig,
)
from mersal.logging import Logger
from mersal.plugins import Plugin
from mersal.threading.anyio.anyio_periodic_async_task_factory import (
AnyIOPeriodicTaskFactory,
)
from mersal.timeouts.due_messages_sender import DueMessagesSender
from mersal.timeouts.timeout_manager import TimeoutManager
from mersal.transport import Transport
from mersal.utils.sync import AsyncCallable
if TYPE_CHECKING:
from mersal.configuration import StandardConfigurator
from mersal.timeouts.config import TimeoutsConfig
from mersal.types import LifespanHook
__all__ = ("TimeoutsPlugin",)
[docs]
class TimeoutsPlugin(Plugin):
[docs]
def __init__(self, config: TimeoutsConfig):
self._config = config
def __call__(self, configurator: StandardConfigurator) -> None:
from mersal.timeouts.config import TimeoutsConfig
configurator.register(TimeoutsConfig, lambda _: self._config)
storage = self._config.storage
if storage is None:
return
if configurator.send_only:
raise InvalidConfigurationError(
"A send-only app can't host a timeout manager since it never receives the deferred "
"messages sent to it; use TimeoutsConfig(external_timeout_manager_address=...) instead"
)
def register_sender(configurator: StandardConfigurator) -> DueMessagesSender:
logger = configurator.get(Logger) # type: ignore[type-abstract]
return DueMessagesSender(
AnyIOPeriodicTaskFactory(logger=logger),
configurator.get(Transport), # type: ignore[type-abstract]
configurator.get(TimeoutManager), # type: ignore[type-abstract]
logger=logger,
poll_interval=self._config.poll_interval,
)
configurator.register(TimeoutManager, lambda _: storage)
configurator.register(DueMessagesSender, register_sender)
startup_hooks: list[Callable[[StandardConfigurator], LifespanHook]] = [
lambda config: AsyncCallable(config.get(TimeoutManager)), # type: ignore[type-abstract]
lambda config: AsyncCallable(config.get(DueMessagesSender).start),
]
shutdown_hooks: list[Callable[[StandardConfigurator], LifespanHook]] = [
lambda config: AsyncCallable(config.get(DueMessagesSender).stop),
]
LifespanHooksRegistrationPluginConfig(
on_startup_hooks=startup_hooks,
on_shutdown_hooks=shutdown_hooks,
).plugin(configurator)