Source code for mersal.timeouts.due_messages_sender

from mersal.logging import Logger
from mersal.messages import MessageHeaders, TransportMessage
from mersal.threading.periodic_async_task_factory import PeriodicAsyncTaskFactory
from mersal.transport import TransactionScope, Transport

from .timeout_manager import TimeoutManager

__all__ = ("DueMessagesSender",)


[docs] class DueMessagesSender: """Periodically sends due deferred messages to their recipients."""
[docs] def __init__( self, periodic_task_factory: PeriodicAsyncTaskFactory, transport: Transport, timeout_manager: TimeoutManager, logger: Logger, poll_interval: float = 1, ) -> None: """Initialize ``DueMessagesSender``. Args: periodic_task_factory: Creates the :class:`PeriodicAsyncTask <.threading.PeriodicAsyncTask>` that runs the periodic check & send. transport: The relevant :class:`Transport <.transport.Transport>`. timeout_manager: Storage of deferred messages. logger: Logger instance. poll_interval: Period for checking for due messages (in seconds). """ self._transport = transport self._timeout_manager = timeout_manager self._logger = logger self._task = periodic_task_factory("Timeouts-DueMessagesSender", self.send_due_messages, poll_interval)
async def start(self) -> None: await self._task.start() async def stop(self) -> None: await self._task.stop() async def send_due_messages(self) -> None: async with self._timeout_manager.get_due_messages() as due_messages: for due_message in due_messages: try: await self._send(due_message.message) except Exception: # Left in the storage; retried on the next run. self._logger.exception( "timeouts.due_message.send.error", message_id=due_message.message.headers.message_id, ) continue await due_message.mark_as_completed() async def _send(self, message: TransportMessage) -> None: headers = MessageHeaders(message.headers) headers.pop(MessageHeaders.deferred_until_key, None) recipient = headers.pop(MessageHeaders.deferred_recipient_key, None) if recipient is None: raise ValueError(f"Deferred message {headers.message_id} has no recipient") self._logger.debug("timeouts.due_message.send", message_id=headers.message_id, recipient=recipient) async with TransactionScope() as scope: await self._transport.send(recipient, TransportMessage(message.body, headers), scope.transaction_context) await scope.complete()