timeouts

class mersal.timeouts.DisabledTimeoutManager[source]

Bases: TimeoutManager

Used when no timeout manager is configured.

Receiving a deferred message then fails loudly instead of handling the message right away.

async defer(due_time: datetime, message: TransportMessage) → None[source]

Store message until due_time.

get_due_messages() → AbstractAsyncContextManager[Sequence[DueMessage]][source]

Provide the messages that are due.

Only messages marked as completed are removed, the rest are provided again on a later call. The context manager lets a storage hold resources (e.g. a database transaction/lock on the returned rows) while they’re being sent.

class mersal.timeouts.DueMessage[source]

Bases: object

A deferred message whose due time has passed.

__init__(message: TransportMessage, mark_as_completed: Callable[[], Awaitable[None]]) → None
mark_as_completed: Callable[[], Awaitable[None]]

Removes the message from the storage once it has been sent to its recipient.

message: TransportMessage

The stored message, still carrying its deferred_until/deferred_recipient headers.

class mersal.timeouts.DueMessagesSender[source]

Bases: object

Periodically sends due deferred messages to their recipients.

__init__(periodic_task_factory: PeriodicAsyncTaskFactory, transport: Transport, timeout_manager: TimeoutManager, logger: Logger, poll_interval: float = 1) → None[source]

Initialize DueMessagesSender.

Parameters:
  • periodic_task_factory – Creates the PeriodicAsyncTask that runs the periodic check & send.

  • transport – The relevant Transport.

  • timeout_manager – Storage of deferred messages.

  • logger – Logger instance.

  • poll_interval – Period for checking for due messages (in seconds).

class mersal.timeouts.HandleDeferredMessagesStep[source]

Bases: 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.

__init__(timeout_manager: TimeoutManager, transport: Transport, external_timeout_manager_address: str | None = None) → None[source]
class mersal.timeouts.TimeoutManager[source]

Bases: Protocol

Storage for deferred messages until they are due.

async __call__() → None[source]

Called on app startup.

Can be used to run any initialization required by the storage, for example creating the database table or making sure it already exists.

__init__(*args, **kwargs)
async defer(due_time: datetime, message: TransportMessage) → None[source]

Store message until due_time.

get_due_messages() → AbstractAsyncContextManager[Sequence[DueMessage]][source]

Provide the messages that are due.

Only messages marked as completed are removed, the rest are provided again on a later call. The context manager lets a storage hold resources (e.g. a database transaction/lock on the returned rows) while they’re being sent.

class mersal.timeouts.TimeoutsConfig[source]

Bases: object

Configuration for deferred messages on transports without native deferral.

Exactly one of storage and external_timeout_manager_address must be set.

__init__(storage: TimeoutManager | None = None, external_timeout_manager_address: str | None = None, poll_interval: float = 1) → None
external_timeout_manager_address: str | None = None

Address of another app that hosts the timeout manager. Deferred messages (including ones this app receives) are sent there.

poll_interval: float = 1

Seconds between checks for due messages (only with storage).

storage: TimeoutManager | None = None

Makes this app a timeout manager. Deferred messages are sent to this app’s own address, stored here and sent to their recipient once due.

class mersal.timeouts.TimeoutsPlugin[source]

Bases: Plugin

__init__(config: TimeoutsConfig)[source]