persistence

class mersal.persistence.in_memory.InMemoryMessageTracker[source]

Bases: MessageTracker

Tracks handled messages in memory.

__init__() → None[source]
async is_message_tracked(message_id: UUID, transaction_context: TransactionContext) → bool[source]

Check if message identified by message_id is handled.

async track_message(message_id: UUID, transaction_context: TransactionContext) → None[source]

Set message identified by message_id as handled.

class mersal.persistence.in_memory.InMemorySagaStorage[source]

Bases: SagaStorage

__init__() → None[source]
class mersal.persistence.in_memory.InMemorySubscriptionStorage[source]

Bases: SubscriptionStorage

In memory implementation for storing topics subscriptions.

__init__() → None[source]
async get_subscriber_addresses(topic: str) → set[str][source]

Get addresses subscribed for the given topic.

Parameters:

topic – topic name to get the addresses for.

Returns:

A set of addresses subscribed to this topic.

property is_centralized: bool

Whether this storage is centralized.

Centralized storage means topic subscriptions are stored in one place and registration/unregistration can be done by directly calling register_subscriber and unregister_subscriber, respectively.

Non centralized means each topic handles its own storage and registration/unregistration is performed by sending a message to the topic owner (async).

async register_subscriber(topic: str, subscriber_address: str) → None[source]

Register the given address for the given topic.

async unregister_subscriber(topic: str, subscriber_address: str) → None[source]

Unregister the given address for the given topic.

class mersal.persistence.in_memory.InMemorySubscriptionStore[source]

Bases: defaultdict[str, MutableSet[str]]

__init__(*args: Any, **kwargs: Any) → None[source]
class mersal.persistence.in_memory.InMemoryTimeoutManager[source]

Bases: TimeoutManager

Stores deferred messages in memory; they’re lost when the process stops.

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

Store message until due_time.

get_due_messages() → AsyncIterator[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.