timeouts¶
- class mersal.timeouts.DisabledTimeoutManager[source]¶
Bases:
TimeoutManagerUsed 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:
objectA deferred message whose due time has passed.
- 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:
objectPeriodically 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.
- class mersal.timeouts.HandleDeferredMessagesStep[source]¶
Bases:
IncomingStepIntercepts 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.
- class mersal.timeouts.TimeoutManager[source]¶
Bases:
ProtocolStorage 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:
objectConfiguration 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.
- 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]¶