Source code for mersal.timeouts.handle_deferred_messages_step

from mersal.exceptions import MersalExceptionError
from mersal.messages import TransportMessage
from mersal.pipeline import IncomingStepContext
from mersal.pipeline.incoming_step import IncomingStep
from mersal.transport import TransactionContext, Transport
from mersal.types import AsyncAnyCallable

from .timeout_manager import TimeoutManager

__all__ = ("HandleDeferredMessagesStep",)


[docs] class HandleDeferredMessagesStep(IncomingStep): """Intercepts deferred messages before they're handled. A message carrying a `deferred_until` header is stored in the timeout manager (or forwarded to the external timeout manager, when configured) instead of being passed down the pipeline. `DueMessagesSender` sends it to its `deferred_recipient` once due. """
[docs] def __init__( self, timeout_manager: TimeoutManager, transport: Transport, external_timeout_manager_address: str | None = None, ) -> None: self._timeout_manager = timeout_manager self._transport = transport self._external_timeout_manager_address = external_timeout_manager_address
async def __call__(self, context: IncomingStepContext, next_step: AsyncAnyCallable) -> None: message = context.load(TransportMessage) deferred_until = message.headers.deferred_until if deferred_until is None: await next_step() return if message.headers.deferred_recipient is None: raise MersalExceptionError( f"Received message {message.headers.message_id} with the " f"'{message.headers.deferred_until_key}' header but without the " f"'{message.headers.deferred_recipient_key}' header" ) if self._external_timeout_manager_address is not None: transaction_context = context.load(TransactionContext) # type: ignore[type-abstract] await self._transport.send(self._external_timeout_manager_address, message, transaction_context) else: await self._timeout_manager.defer(deferred_until, message)