Перейти к содержанию

Isolation Module

Изоляция обработки событий (аналог EventIsolation из aiogram).

Сериализует конкурентную обработку апдейтов одного пользователя, чтобы FSM-переход, завершённый предыдущим апдейтом, был виден следующему. Без изоляции при параллельной обработке (Dispatcher(use_create_task=True) или вебхук) два быстрых сообщения одного пользователя читают один и тот же снимок состояния и одноразовый шаг FSM выполняется дважды.

IsolationKey = tuple[int | None, int | None] module-attribute

Ключ изоляции: (chat_id, user_id).

BaseEventIsolation

Bases: ABC

Базовый класс изоляции обработки событий.

Реализация обязана вернуть асинхронный контекст-менеджер, удерживающий блокировку по ключу (chat_id, user_id) на всё время обработки события.

lock(key) abstractmethod

Возвращает контекст-менеджер блокировки по ключу.

Parameters:

Name Type Description Default
key IsolationKey

Ключ изоляции (chat_id, user_id).

required
Source code in maxapi/context/isolation.py
@abstractmethod
def lock(self, key: IsolationKey) -> AbstractAsyncContextManager[None]:
    """
    Возвращает контекст-менеджер блокировки по ключу.

    Args:
        key: Ключ изоляции ``(chat_id, user_id)``.
    """

close() abstractmethod async

Освобождает ресурсы изоляции.

Вызывается диспетчером при каждом stop_polling(), поэтому реализация должна быть идемпотентной.

Source code in maxapi/context/isolation.py
@abstractmethod
async def close(self) -> None:
    """
    Освобождает ресурсы изоляции.

    Вызывается диспетчером при каждом ``stop_polling()``,
    поэтому реализация должна быть идемпотентной.
    """

DisabledEventIsolation

Bases: BaseEventIsolation

Отключённая изоляция (поведение по умолчанию).

События обрабатываются без сериализации — как до появления механизма изоляции.

lock(key) async

No-op: блокировка не берётся.

Source code in maxapi/context/isolation.py
@asynccontextmanager
async def lock(self, key: IsolationKey) -> AsyncIterator[None]:
    """No-op: блокировка не берётся."""
    yield

close() async

No-op.

Source code in maxapi/context/isolation.py
async def close(self) -> None:
    """No-op."""

SimpleEventIsolation()

Bases: BaseEventIsolation

Изоляция на asyncio.Lock в памяти процесса.

На каждый ключ (chat_id, user_id) создаётся отдельная блокировка; события одного пользователя обрабатываются строго последовательно (в порядке поступления — asyncio.Lock пробуждает ожидающих в FIFO), события разных пользователей — параллельно.

В отличие от aiogram, неиспользуемые блокировки удаляются из словаря, как только их никто не держит и не ожидает.

Подходит только для одного процесса. При нескольких процессах (например, вебхук за балансировщиком) используйте :class:RedisEventIsolation.

Source code in maxapi/context/isolation.py
def __init__(self) -> None:
    self._locks: dict[IsolationKey, Lock] = {}
    self._refcounts: dict[IsolationKey, int] = {}

lock(key) async

Удерживает блокировку по ключу на время контекста.

Parameters:

Name Type Description Default
key IsolationKey

Ключ изоляции (chat_id, user_id).

required
Source code in maxapi/context/isolation.py
@asynccontextmanager
async def lock(self, key: IsolationKey) -> AsyncIterator[None]:
    """
    Удерживает блокировку по ключу на время контекста.

    Args:
        key: Ключ изоляции ``(chat_id, user_id)``.
    """
    lock = self._locks.get(key)
    if lock is None:
        lock = Lock()
        self._locks[key] = lock
    self._refcounts[key] = self._refcounts.get(key, 0) + 1
    try:
        async with lock:
            yield
    finally:
        # get с дефолтом, а не прямой доступ: close() мог
        # очистить словари, пока блокировка удерживалась
        refs = self._refcounts.get(key, 1) - 1
        if refs > 0:
            self._refcounts[key] = refs
        else:
            # Блокировку никто не держит и не ждёт — чистим,
            # чтобы словарь не рос бесконечно
            self._refcounts.pop(key, None)
            self._locks.pop(key, None)

close() async

Очищает словарь блокировок.

Source code in maxapi/context/isolation.py
async def close(self) -> None:
    """Очищает словарь блокировок."""
    self._locks.clear()
    self._refcounts.clear()

RedisEventIsolation(redis_client, key_prefix='maxapi', lock_timeout=DEFAULT_REDIS_LOCK_TIMEOUT, lock_sleep=DEFAULT_REDIS_LOCK_SLEEP)

Bases: BaseEventIsolation

Распределённая изоляция на блокировках Redis.

Сериализует обработку событий одного пользователя между несколькими процессами/инстансами бота. Парная к :class:~maxapi.context.context.RedisContext — используйте их вместе. Требует установленной библиотеки redis: pip install redis.

Инициализация изоляции.

Parameters:

Name Type Description Default
redis_client Any

Экземпляр redis.asyncio.Redis.

required
key_prefix str

Префикс ключей блокировок. Рекомендуется тот же, что у вашего RedisContext, чтобы все ключи бота жили в одном неймспейсе; для разных ботов на одном Redis префиксы обязаны различаться, иначе их блокировки пересекутся.

'maxapi'
lock_timeout float | None

Максимальное время удержания блокировки в секундах (страховка от вечного лока при падении процесса). Если обработка события длится дольше, блокировка истекает и изоляция для этого события перестаёт действовать. None — без ограничения; не рекомендуется: после падения процесса, державшего блокировку, ключ останется в Redis навсегда и все апдейты этого пользователя зависнут до ручного удаления ключа.

DEFAULT_REDIS_LOCK_TIMEOUT
lock_sleep float

Интервал опроса блокировки в секундах.

DEFAULT_REDIS_LOCK_SLEEP
Source code in maxapi/context/isolation.py
def __init__(
    self,
    redis_client: Any,  # redis.asyncio.Redis
    key_prefix: str = "maxapi",
    lock_timeout: float | None = DEFAULT_REDIS_LOCK_TIMEOUT,
    lock_sleep: float = DEFAULT_REDIS_LOCK_SLEEP,
) -> None:
    """
    Инициализация изоляции.

    Args:
        redis_client: Экземпляр ``redis.asyncio.Redis``.
        key_prefix: Префикс ключей блокировок. Рекомендуется тот
            же, что у вашего ``RedisContext``, чтобы все ключи
            бота жили в одном неймспейсе; для разных ботов на
            одном Redis префиксы обязаны различаться, иначе их
            блокировки пересекутся.
        lock_timeout: Максимальное время удержания блокировки в
            секундах (страховка от вечного лока при падении
            процесса). Если обработка события длится дольше,
            блокировка истекает и изоляция для этого события
            перестаёт действовать. None — без ограничения; не
            рекомендуется: после падения процесса, державшего
            блокировку, ключ останется в Redis навсегда и все
            апдейты этого пользователя зависнут до ручного
            удаления ключа.
        lock_sleep: Интервал опроса блокировки в секундах.
    """
    self.redis = redis_client
    self.key_prefix = key_prefix
    self.lock_timeout = lock_timeout
    self.lock_sleep = lock_sleep

lock(key) async

Удерживает распределённую блокировку по ключу.

Parameters:

Name Type Description Default
key IsolationKey

Ключ изоляции (chat_id, user_id).

required
Source code in maxapi/context/isolation.py
@asynccontextmanager
async def lock(self, key: IsolationKey) -> AsyncIterator[None]:
    """
    Удерживает распределённую блокировку по ключу.

    Args:
        key: Ключ изоляции ``(chat_id, user_id)``.
    """
    chat_id, user_id = key
    name = f"{self.key_prefix}:{chat_id}:{user_id}:lock"
    lock = self.redis.lock(
        name=name,
        timeout=self.lock_timeout,
        sleep=self.lock_sleep,
    )
    acquired = await lock.acquire()
    if not acquired:
        msg = (
            f"Не удалось захватить блокировку изоляции {name}: "
            "acquire() вернул False"
        )
        raise RuntimeError(msg)
    try:
        yield
    finally:
        try:
            await lock.release()
        except Exception as e:
            # Например, LockNotOwnedError: обработка события
            # заняла дольше lock_timeout, ключ истёк (и изоляция
            # на хвосте обработки не действовала). Событие при
            # этом обработано — не превращаем успех в ошибку
            logger_dp.warning(
                "Не удалось освободить блокировку изоляции %s "
                "(обработка дольше lock_timeout=%s?): %r",
                name,
                self.lock_timeout,
                e,
            )

close() async

No-op: соединением Redis владеет вызывающая сторона.

Source code in maxapi/context/isolation.py
async def close(self) -> None:
    """No-op: соединением Redis владеет вызывающая сторона."""