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

Dispatcher Module

Dispatcher(router_id=None, storage=MemoryContext, *, use_create_task=False, event_isolation=None, **storage_kwargs)

Bases: BotMixin

Основной класс для обработки событий бота.

Обеспечивает запуск поллинга и вебхука, маршрутизацию событий, применение middleware, фильтров и вызов соответствующих обработчиков.

Инициализация диспетчера.

Parameters:

Name Type Description Default
router_id str | None

Идентификатор роутера для логов.

None
use_create_task bool

Флаг, отвечающий за параллелизацию обработок событий.

False
event_isolation BaseEventIsolation | None

Изоляция обработки событий: сериализует конкурентные апдейты одного пользователя (см. :class:~maxapi.context.SimpleEventIsolation). По умолчанию отключена (:class:~maxapi.context.DisabledEventIsolation).

None
storage Any

Класс контекста для хранения данных (MemoryContext, RedisContext и т.д.).

MemoryContext
**storage_kwargs Any

Дополнительные аргументы для инициализации хранилища.

{}
Source code in maxapi/dispatcher.py
def __init__(
    self,
    router_id: str | None = None,
    storage: Any = MemoryContext,
    *,
    use_create_task: bool = False,
    event_isolation: BaseEventIsolation | None = None,
    **storage_kwargs: Any,
) -> None:
    """
    Инициализация диспетчера.

    Args:
        router_id: Идентификатор роутера для логов.
        use_create_task: Флаг, отвечающий за параллелизацию
            обработок событий.
        event_isolation: Изоляция обработки событий: сериализует
            конкурентные апдейты одного пользователя
            (см. :class:`~maxapi.context.SimpleEventIsolation`).
            По умолчанию отключена
            (:class:`~maxapi.context.DisabledEventIsolation`).
        storage: Класс контекста для хранения
            данных (MemoryContext, RedisContext и т.д.).
        **storage_kwargs: Дополнительные аргументы для
            инициализации хранилища.
    """

    self.router_id = router_id
    self.storage = storage
    self.storage_kwargs = storage_kwargs
    self.event_isolation: BaseEventIsolation = (
        event_isolation
        if event_isolation is not None
        else DisabledEventIsolation()
    )
    self._fsm = ContextManager(self, self.__get_context)

    self.event_handlers: list[Handler] = []
    self.error_handlers: list[ErrorHandler] = []
    self.handlers_by_type: dict[UpdateType, list[Handler]] | None = None
    self.contexts: OrderedDict[
        tuple[int | None, int | None], BaseContext
    ] = OrderedDict()
    self.routers: list[Router | Dispatcher] = []
    self.filters: list[MagicFilter] = []
    self.base_filters: list[BaseFilter] = []
    self.outer_middlewares: list[BaseMiddleware] = []
    self.inner_middlewares: list[BaseMiddleware] = []

    self.bot: Bot | None = None
    self.on_started_func: Callable | None = None
    self.polling = False
    self.use_create_task = use_create_task
    self._cached_router_entries: list[_DispatchEntry] | None = None
    self._global_mw_chain: HandlerCallable | None = None
    self._background_tasks: set[asyncio.Task] = set()
    self._closing: bool = False
    self._deferred_shutdown: bool = False
    self._polling_task: asyncio.Task | None = None
    self._polling_active: bool = False
    self._lifecycle_holders: int = 0
    self._cleanup_done: asyncio.Event | None = None
    self._loop_done: asyncio.Event | None = None
    self._polling_error: BaseException | None = None
    self._stop_event: asyncio.Event | None = None
    self._ready: bool = False
    self._running_on_started: bool = False
    self._parents: weakref.WeakSet[Dispatcher] = weakref.WeakSet()
    self._handlers_dirty: bool = False
    self._warned_duplicate_routers: weakref.WeakSet[
        Router | Dispatcher
    ] = weakref.WeakSet()

    self.message_created = Event(
        update_type=UpdateType.MESSAGE_CREATED, router=self
    )
    self.errors = ErrorEventObserver(router=self)
    self.error = self.errors
    self.bot_added = Event(update_type=UpdateType.BOT_ADDED, router=self)
    self.bot_removed = Event(
        update_type=UpdateType.BOT_REMOVED, router=self
    )
    self.bot_started = Event(
        update_type=UpdateType.BOT_STARTED, router=self
    )
    self.bot_stopped = Event(
        update_type=UpdateType.BOT_STOPPED, router=self
    )
    self.dialog_cleared = Event(
        update_type=UpdateType.DIALOG_CLEARED, router=self
    )
    self.dialog_muted = Event(
        update_type=UpdateType.DIALOG_MUTED, router=self
    )
    self.dialog_unmuted = Event(
        update_type=UpdateType.DIALOG_UNMUTED, router=self
    )
    self.dialog_removed = Event(
        update_type=UpdateType.DIALOG_REMOVED, router=self
    )
    self.raw_api_response = Event(
        update_type=UpdateType.RAW_API_RESPONSE, router=self
    )
    self.chat_title_changed = Event(
        update_type=UpdateType.CHAT_TITLE_CHANGED, router=self
    )
    self.message_callback = Event(
        update_type=UpdateType.MESSAGE_CALLBACK, router=self
    )
    self.message_chat_created = Event(
        update_type=UpdateType.MESSAGE_CHAT_CREATED,
        router=self,
        deprecated=True,
    )
    self.message_edited = Event(
        update_type=UpdateType.MESSAGE_EDITED, router=self
    )
    self.message_removed = Event(
        update_type=UpdateType.MESSAGE_REMOVED, router=self
    )
    self.user_added = Event(update_type=UpdateType.USER_ADDED, router=self)
    self.user_removed = Event(
        update_type=UpdateType.USER_REMOVED, router=self
    )
    self.on_started = Event(update_type=UpdateType.ON_STARTED, router=self)

fsm property

Менеджер FSM-контекстов диспетчера.

middlewares property writable

Список outer-middleware.

.. deprecated:: Используйте :attr:outer_middlewares.

check_me() async

Проверяет и логирует информацию о боте.

Source code in maxapi/dispatcher.py
async def check_me(self) -> None:
    """
    Проверяет и логирует информацию о боте.
    """

    bot = self._ensure_bot()
    me = await bot.get_me()

    bot.me = me

    logger_dp.info(
        "Бот: @%s first_name=%s id=%s",
        me.username,
        me.first_name,
        me.user_id,
    )

build_middleware_chain(middlewares, handler) staticmethod

Формирует цепочку вызова middleware вокруг хендлера.

Parameters:

Name Type Description Default
middlewares list[BaseMiddleware]

Список middleware.

required
handler HandlerCallable

Финальный обработчик.

required

Returns:

Name Type Description
Callable HandlerCallable

Обёрнутый обработчик.

Source code in maxapi/dispatcher.py
@staticmethod
def build_middleware_chain(
    middlewares: list[BaseMiddleware],
    handler: HandlerCallable,
) -> HandlerCallable:
    """
    Формирует цепочку вызова middleware вокруг хендлера.

    Args:
        middlewares: Список middleware.
        handler: Финальный обработчик.

    Returns:
        Callable: Обёрнутый обработчик.
    """

    for mw in reversed(middlewares):
        handler = functools.partial(mw, handler)

    return handler

include_routers(*routers)

Добавляет указанные роутеры в диспетчер.

Можно вызывать и после старта: индекс обработчиков будет перестроен перед следующей диспетчеризацией, тогда же у добавленного роутера появится router.bot (до этого он остаётся None).

Порядок обхода сохраняется: включённые роутеры проверяются раньше собственных обработчиков диспетчера — в том числе при позднем включении, когда сам диспетчер уже добавлен в конец self.routers (см. :meth:__ready).

Прямая мутация dp.routers (dp.routers.append(...)) индекс устаревшим не помечает: изменения попадут в диспетчеризацию лишь при следующей перестройке, вызванной другой регистрацией, а добавленный так роутер окажется после собственных обработчиков диспетчера. Используйте этот метод.

Роутер должен принадлежать ОДНОМУ дереву — одному диспетчеру. Включение одного и того же роутера в несколько диспетчеров не поддерживается: подготовленное состояние (router.bot, выпеченные цепочки middleware в handler.mw_chain) хранится на самих объектах роутера и обработчиков, и подготовка второго диспетчера перезапишет его для первого — как и bot.commands с bot.dispatcher, которые тоже одни. Внутри одного дерева повторное включение допустимо: обход дедуплицирует роутеры по первому вхождению (и предупреждает о дублях).

Parameters:

Name Type Description Default
*routers Router

Роутеры для добавления.

()
Source code in maxapi/dispatcher.py
def include_routers(self, *routers: Router) -> None:
    """
    Добавляет указанные роутеры в диспетчер.

    Можно вызывать и после старта: индекс обработчиков будет
    перестроен перед следующей диспетчеризацией, тогда же у
    добавленного роутера появится ``router.bot`` (до этого он
    остаётся ``None``).

    Порядок обхода сохраняется: включённые роутеры проверяются
    раньше собственных обработчиков диспетчера — в том числе при
    позднем включении, когда сам диспетчер уже добавлен в конец
    ``self.routers`` (см. :meth:`__ready`).

    Прямая мутация ``dp.routers`` (``dp.routers.append(...)``)
    индекс устаревшим не помечает: изменения попадут в
    диспетчеризацию лишь при следующей перестройке, вызванной
    другой регистрацией, а добавленный так роутер окажется после
    собственных обработчиков диспетчера. Используйте этот метод.

    Роутер должен принадлежать ОДНОМУ дереву — одному диспетчеру.
    Включение одного и того же роутера в несколько диспетчеров не
    поддерживается: подготовленное состояние (``router.bot``,
    выпеченные цепочки middleware в ``handler.mw_chain``) хранится
    на самих объектах роутера и обработчиков, и подготовка второго
    диспетчера перезапишет его для первого — как и ``bot.commands``
    с ``bot.dispatcher``, которые тоже одни. Внутри одного дерева
    повторное включение допустимо: обход дедуплицирует роутеры по
    первому вхождению (и предупреждает о дублях).

    Args:
        *routers: Роутеры для добавления.
    """

    if self in self.routers:
        # Сам диспетчер стоит последним: новые роутеры должны
        # попасть перед ним, иначе поздно включённый роутер
        # получал бы событие после хендлеров самого dp.
        position = self.routers.index(self)
        self.routers[position:position] = routers
    else:
        self.routers.extend(routers)

    for router in routers:
        router._parents.add(self)  # noqa: SLF001

    self._invalidate_handlers()

register_outer_middleware(middleware)

Регистрирует outer middleware (до проверки фильтров handler).

Вызывается для каждого подходящего события ещё до того, как диспетчер узнает, какой именно handler сработает.

Порядок регистрации сохраняется: первый зарегистрированный outer middleware выполняется первым (внешний слой цепочки), что симметрично с :meth:register_inner_middleware.

Parameters:

Name Type Description Default
middleware BaseMiddleware

Middleware.

required
Source code in maxapi/dispatcher.py
def register_outer_middleware(self, middleware: BaseMiddleware) -> None:
    """
    Регистрирует outer middleware (до проверки фильтров handler).

    Вызывается для каждого подходящего события ещё до того, как
    диспетчер узнает, какой именно handler сработает.

    Порядок регистрации сохраняется: первый зарегистрированный
    outer middleware выполняется первым (внешний слой цепочки),
    что симметрично с :meth:`register_inner_middleware`.

    Args:
        middleware: Middleware.
    """
    self.outer_middlewares.append(middleware)
    self._invalidate_handlers()

register_inner_middleware(middleware)

Регистрирует inner middleware (после проверки фильтров handler).

Вызывается только тогда, когда конкретный handler прошёл все свои фильтры и state и будет реально исполнен. На уровне Dispatcher — только для событий, попавших хоть в один handler; на уровне Router — только для handler этого роутера.

Регистрация во время обработки события применяется и к нему, если вызов его обработчика ещё не начался: перестройка индекса переприсваивает handler.mw_chain (см. :meth:_prepare_handlers).

Parameters:

Name Type Description Default
middleware BaseMiddleware

Middleware.

required
Source code in maxapi/dispatcher.py
def register_inner_middleware(self, middleware: BaseMiddleware) -> None:
    """
    Регистрирует inner middleware (после проверки фильтров handler).

    Вызывается только тогда, когда конкретный handler прошёл все
    свои фильтры и state и будет реально исполнен. На уровне
    Dispatcher — только для событий, попавших хоть в один handler;
    на уровне Router — только для handler этого роутера.

    Регистрация во время обработки события применяется и к нему,
    если вызов его обработчика ещё не начался: перестройка индекса
    переприсваивает ``handler.mw_chain``
    (см. :meth:`_prepare_handlers`).

    Args:
        middleware (BaseMiddleware): Middleware.
    """
    self.inner_middlewares.append(middleware)
    self._invalidate_handlers()

outer_middleware(middleware)

Добавляет Middleware на первое место в списке outer_middlewares.

Историческое поведение: insert(0, ...). В новом :meth:register_outer_middleware порядок изменён на append (register order = execution order), поэтому при миграции проверьте порядок вызовов, если он важен.

.. deprecated:: Используйте :meth:register_outer_middleware.

Parameters:

Name Type Description Default
middleware BaseMiddleware

Middleware.

required
Source code in maxapi/dispatcher.py
def outer_middleware(self, middleware: BaseMiddleware) -> None:
    """
    Добавляет Middleware на первое место в списке outer_middlewares.

    Историческое поведение: ``insert(0, ...)``. В новом
    :meth:`register_outer_middleware` порядок изменён на ``append``
    (register order = execution order), поэтому при миграции
    проверьте порядок вызовов, если он важен.

    .. deprecated::
        Используйте :meth:`register_outer_middleware`.

    Args:
        middleware (BaseMiddleware): Middleware.
    """
    warnings.warn(
        f"{type(self).__name__}.outer_middleware() устарел. "
        "Используйте register_outer_middleware().",
        DeprecationWarning,
        stacklevel=2,
    )
    self.outer_middlewares.insert(0, middleware)
    self._invalidate_handlers()

middleware(middleware)

Добавляет Middleware в конец списка.

.. deprecated:: Используйте :meth:register_outer_middleware (текущее поведение — outer, до фильтров handler) или :meth:register_inner_middleware (только когда handler реально вызван).

Parameters:

Name Type Description Default
middleware BaseMiddleware

Middleware.

required
Source code in maxapi/dispatcher.py
def middleware(self, middleware: BaseMiddleware) -> None:
    """
    Добавляет Middleware в конец списка.

    .. deprecated::
        Используйте :meth:`register_outer_middleware` (текущее
        поведение — outer, до фильтров handler) или
        :meth:`register_inner_middleware` (только когда handler
        реально вызван).

    Args:
        middleware: Middleware.
    """
    warnings.warn(
        f"{type(self).__name__}.middleware() устарел. "
        "Используйте register_outer_middleware() (поведение "
        "сохраняется) или register_inner_middleware() для запуска "
        "mw только после выбора handler.",
        DeprecationWarning,
        stacklevel=2,
    )
    self.outer_middlewares.append(middleware)
    self._invalidate_handlers()

filter(base_filter)

Добавляет фильтр уровня роутера.

Принимает как :class:~magic_filter.MagicFilter (F.chat.type == ChatType.DIALOG), так и :class:~maxapi.filters.filter.BaseFilter: тип определяется по значению и фильтр попадает в filters или base_filters соответственно.

Можно вызывать и после старта. Прямая мутация списков (router.filters.append(...)) индекс устаревшим не помечает: добавленный так фильтр начнёт действовать лишь после перестройки, вызванной другой регистрацией. Используйте этот метод.

Parameters:

Name Type Description Default
base_filter MagicFilter | BaseFilter

Фильтр.

required
Source code in maxapi/dispatcher.py
def filter(self, base_filter: MagicFilter | BaseFilter) -> None:
    """
    Добавляет фильтр уровня роутера.

    Принимает как :class:`~magic_filter.MagicFilter`
    (``F.chat.type == ChatType.DIALOG``), так и
    :class:`~maxapi.filters.filter.BaseFilter`: тип определяется по
    значению и фильтр попадает в ``filters`` или ``base_filters``
    соответственно.

    Можно вызывать и после старта. Прямая мутация списков
    (``router.filters.append(...)``) индекс устаревшим не
    помечает: добавленный так фильтр начнёт действовать лишь
    после перестройки, вызванной другой регистрацией. Используйте
    этот метод.

    Args:
        base_filter: Фильтр.
    """

    if isinstance(base_filter, MagicFilter):
        self.filters.append(base_filter)
    else:
        self.base_filters.append(base_filter)

    self._invalidate_handlers()

__ready(bot) async

Подготавливает диспетчер: сохраняет бота, подготавливает обработчики, вызывает on_started.

Флаг _running_on_started взводится ровно на фазу вызова on_started — непосредственно перед чтением on_started_func. По нему :meth:Event.register отличает регистрацию колбэка изнутри самого on_started (в этом запуске он уже не будет вызван) от регистрации в окне подготовки до этой фазы (check_me и проверка подписок): колбэк, зарегистрированный там, штатно сработает в этом же запуске, и предупреждать о нём не о чем.

Parameters:

Name Type Description Default
bot Bot

Экземпляр бота.

required
Source code in maxapi/dispatcher.py
async def __ready(self, bot: Bot) -> None:
    """
    Подготавливает диспетчер: сохраняет бота, подготавливает
    обработчики, вызывает on_started.

    Флаг ``_running_on_started`` взводится ровно на фазу вызова
    ``on_started`` — непосредственно перед чтением
    ``on_started_func``. По нему :meth:`Event.register` отличает
    регистрацию колбэка изнутри самого ``on_started`` (в этом
    запуске он уже не будет вызван) от регистрации в окне
    подготовки до этой фазы (``check_me`` и проверка подписок):
    колбэк, зарегистрированный там, штатно сработает в этом же
    запуске, и предупреждать о нём не о чем.

    Args:
        bot: Экземпляр бота.
    """

    # Сбрасываем признак завершения до раннего выхода: повторный
    # startup() после shutdown() (webhook-сценарий) не проходит
    # подготовку заново, но диспетчер снова принимает события.
    self._closing = False

    if self._ready:
        # Регистрации между shutdown() и повторным startup()
        # должны попасть в индекс сразу: подготовка не
        # повторяется, а bot.commands обязан быть актуален уже
        # до первого события.
        self._ensure_prepared()
        return

    self.bot = bot
    self.bot.dispatcher = self

    # Сам диспетчер добавляем в роутеры до сетевых await'ов
    # ниже: событие, пришедшее в окно подготовки, вызовет
    # перестройку индекса, и без этого его собственные
    # обработчики в неё не попадут.
    if self not in self.routers:
        self.routers.append(self)

    if self.polling and bot.auto_check_subscriptions:
        await self._check_subscriptions(bot)

    await self.check_me()

    self._prepare_handlers(bot)

    self._global_mw_chain = self.build_middleware_chain(
        self.outer_middlewares, self._process_event
    )

    # Флаг взводим до чтения on_started_func: всё, что
    # зарегистрировано позже этой точки, в текущем запуске уже не
    # вызовется.
    self._running_on_started = True
    try:
        if self.on_started_func:
            await self.on_started_func()
    finally:
        self._running_on_started = False

    # Регистрации внутри on_started попадают в индекс сразу,
    # чтобы первое же событие не платило за перестройку.
    self._ensure_prepared()

    self._ready = True

__get_context(chat_id, user_id)

Возвращает существующий или создаёт новый контекст по chat_id и user_id.

Parameters:

Name Type Description Default
chat_id int | None

Идентификатор чата.

required
user_id int | None

Идентификатор пользователя.

required

Returns:

Type Description
BaseContext

Контекст.

Source code in maxapi/dispatcher.py
def __get_context(
    self, chat_id: int | None, user_id: int | None
) -> BaseContext:
    """
    Возвращает существующий или создаёт новый контекст
    по chat_id и user_id.

    Args:
        chat_id: Идентификатор чата.
        user_id: Идентификатор пользователя.

    Returns:
        Контекст.
    """

    key = (chat_id, user_id)
    ctx = self.contexts.get(key)
    if ctx is not None:
        if ctx.is_ttl_expired():
            logger_dp.debug("Истёк TTL контекста %s", key)
            del self.contexts[key]
        else:
            ctx.touch_ttl()
            # Перемещаем в конец, чтобы LRU-вытеснение удаляло
            # самые давно неиспользованные контексты
            self.contexts.move_to_end(key)
            return ctx

    if len(self.contexts) >= CONTEXTS_MAX_SIZE:
        evicted_key = next(iter(self.contexts))
        logger_dp.debug(
            "Вытеснен контекст %s (лимит %d)",
            evicted_key,
            CONTEXTS_MAX_SIZE,
        )
        self.contexts.popitem(last=False)

    new_ctx = self.storage(chat_id, user_id, **self.storage_kwargs)
    new_ctx.touch_ttl()
    self.contexts[key] = new_ctx
    return new_ctx

call_handler(handler, event_object, data) async staticmethod

Вызывает хендлер с нужными аргументами.

Перед вызовом фильтрует data, оставляя только те ключи, которые handler реально принимает (по handler.func_args или параметрам, полученным через :func:inspect.signature). В отличие от get_annotations, signature не включает "return" и не требует eval строковых аннотаций — безопасен при from __future__ import annotations. Несовместимые ключи не дойдут до handler и не приведут к TypeError.

Parameters:

Name Type Description Default
handler Handler

Handler.

required
event_object UpdateUnion | dict[str, Any] | str

Объект события.

required
data dict[str, Any]

Данные, накопленные фильтрами и middleware.

required

Returns:

Type Description
None

None

Source code in maxapi/dispatcher.py
@staticmethod
async def call_handler(
    handler: Handler,
    event_object: UpdateUnion | dict[str, Any] | str,
    data: dict[str, Any],
) -> None:
    """
    Вызывает хендлер с нужными аргументами.

    Перед вызовом фильтрует ``data``, оставляя только те ключи,
    которые handler реально принимает (по ``handler.func_args`` или
    параметрам, полученным через :func:`inspect.signature`).
    В отличие от ``get_annotations``, ``signature`` не включает
    ``"return"`` и не требует eval строковых аннотаций — безопасен
    при ``from __future__ import annotations``. Несовместимые ключи
    не дойдут до handler и не приведут к ``TypeError``.

    Args:
        handler: Handler.
        event_object: Объект события.
        data: Данные, накопленные фильтрами и middleware.

    Returns:
        None
    """
    if data:
        func_args = handler.func_args or frozenset(
            inspect.signature(handler.func_event).parameters,
        )
        kwargs = {k: v for k, v in data.items() if k in func_args}
        if kwargs:
            await handler.func_event(event_object, **kwargs)
            return

    await handler.func_event(event_object)

call_error_handler(handler, event_object, data) async staticmethod

Вызывает обработчик ошибки с подходящими kwargs.

Parameters:

Name Type Description Default
handler ErrorHandler

Обработчик ошибки.

required
event_object ErrorEvent

Событие ошибки.

required
data dict[str, Any]

Данные, накопленные фильтрами.

required
Source code in maxapi/dispatcher.py
@staticmethod
async def call_error_handler(
    handler: ErrorHandler,
    event_object: ErrorEventObject,
    data: dict[str, Any],
) -> None:
    """
    Вызывает обработчик ошибки с подходящими kwargs.

    Args:
        handler: Обработчик ошибки.
        event_object: Событие ошибки.
        data: Данные, накопленные фильтрами.
    """
    if data:
        func_args = handler.func_args or frozenset(
            inspect.signature(handler.func_event).parameters,
        )
        kwargs = {k: v for k, v in data.items() if k in func_args}
        if kwargs:
            await handler.func_event(event_object, **kwargs)
            return

    await handler.func_event(event_object)

process_base_filters(event, filters, data=None) async staticmethod

Асинхронно применяет фильтры к событию.

Parameters:

Name Type Description Default
event Any

Событие.

required
filters list[BaseFilter]

Список фильтров.

required

Returns:

Type Description
dict[str, Any] | None

dict[str, Any] | None: Словарь с результатом или None, если фильтр не прошёл.

Source code in maxapi/dispatcher.py
@staticmethod
async def process_base_filters(
    event: Any,
    filters: list[BaseFilter],
    data: dict[str, Any] | None = None,
) -> dict[str, Any] | None:
    """
    Асинхронно применяет фильтры к событию.

    Args:
        event: Событие.
        filters: Список фильтров.

    Returns:
        dict[str, Any] | None: Словарь с результатом или None,
            если фильтр не прошёл.
    """

    available_data: dict[str, Any] = dict(data or {})
    filter_data: dict[str, Any] = {}

    for _filter in filters:
        kwargs = Dispatcher._resolve_filter_kwargs(_filter, available_data)
        result = await _filter(event, **kwargs)

        if isinstance(result, dict):
            filter_data.update(result)
            available_data.update(result)

        elif not result:
            return None

    return filter_data

handle_raw_response(event_type, raw_data) async

Специальный метод для обработки сырых ответов API.

raw_data — разобранный JSON-объект ответа либо сырой текст, если тело ответа не является JSON-объектом (например, HTML от прокси при 502/503).

Source code in maxapi/dispatcher.py
async def handle_raw_response(
    self, event_type: UpdateType, raw_data: dict[str, Any] | str
) -> None:
    """
    Специальный метод для обработки сырых ответов API.

    ``raw_data`` — разобранный JSON-объект ответа либо сырой текст,
    если тело ответа не является JSON-объектом (например, HTML
    от прокси при 502/503).
    """
    self._ensure_prepared()

    entries: Iterable[_DispatchEntry] = (
        self._cached_router_entries
        if self._cached_router_entries is not None
        else self._iter_dispatch_entries()
    )
    for router, *_, handlers_index in entries:
        matching_handlers = self._find_matching_handlers(
            router=router,
            event_type=event_type,
            handlers_index=handlers_index,
        )
        for handler in matching_handlers:
            try:
                await self.call_handler(
                    handler=handler,
                    event_object=raw_data,
                    data={},
                )
            except Exception as e:  # noqa: PERF203
                logger_dp.exception(
                    "Ошибка в обработчике RAW_API_RESPONSE: %r", e
                )

spawn_handle_task(event_object)

Создаёт фоновую задачу handle() и регистрирует её в пуле.

Единая точка постановки задач для polling (use_create_task=True) и webhook-интеграций: без регистрации в _background_tasks задачу может потерять GC, а :meth:shutdown не дождётся её завершения.

Parameters:

Name Type Description Default
event_object UpdateUnion

Событие.

required

Returns:

Type Description
Task

Созданная задача.

Source code in maxapi/dispatcher.py
def spawn_handle_task(self, event_object: UpdateUnion) -> asyncio.Task:
    """
    Создаёт фоновую задачу ``handle()`` и регистрирует её в пуле.

    Единая точка постановки задач для polling
    (``use_create_task=True``) и webhook-интеграций: без
    регистрации в ``_background_tasks`` задачу может потерять GC,
    а :meth:`shutdown` не дождётся её завершения.

    Args:
        event_object: Событие.

    Returns:
        Созданная задача.
    """
    if self._closing:
        logger_dp.warning(
            "Задача handle() создана во время shutdown: %s",
            event_object.update_type,
        )
    task = asyncio.create_task(self.handle(event_object))
    self._background_tasks.add(task)
    task.add_done_callback(self._on_background_task_done)
    return task

handle(event_object) async

Основной обработчик события. Применяет фильтры, middleware и вызывает нужный handler.

При включённой изоляции (event_isolation) вся обработка — от чтения FSM-состояния до завершения хендлера и обработчиков ошибок — выполняется под блокировкой по ключу (chat_id, user_id): конкурентные апдейты одного пользователя сериализуются.

Parameters:

Name Type Description Default
event_object UpdateUnion

Событие.

required
Source code in maxapi/dispatcher.py
async def handle(self, event_object: UpdateUnion) -> None:
    """
    Основной обработчик события. Применяет фильтры, middleware
    и вызывает нужный handler.

    При включённой изоляции (``event_isolation``) вся обработка —
    от чтения FSM-состояния до завершения хендлера и обработчиков
    ошибок — выполняется под блокировкой по ключу
    ``(chat_id, user_id)``: конкурентные апдейты одного
    пользователя сериализуются.

    Args:
        event_object: Событие.
    """
    process_info = "нет данных"

    # Маркер «эта задача выполняет handle() ЭТОГО диспетчера»:
    # по нему shutdown() распознаёт реентрантный вызов
    # (обработчик остановил диспетчер сам). Без задачи в маркере
    # его не ставим: сравнивать было бы не с чем.
    current_task = asyncio.current_task()
    token = _in_handler.set(
        (self, current_task) if current_task is not None else None
    )
    try:
        self._ensure_prepared()

        ids = event_object.get_ids()
        process_info = (
            f"{event_object.update_type} | "
            f"chat_id: {ids[0]}, user_id: {ids[1]}"
        )
        async with self.event_isolation.lock(ids):
            await self._handle_locked(
                event_object=event_object,
                ids=ids,
                process_info=process_info,
            )
    except Exception as e:
        logger_dp.exception(
            "Ошибка при обработке события: %s | %r",
            process_info,
            e,
        )
    finally:
        _in_handler.reset(token)

start_polling(bot, *, skip_updates=False) async

Запускает цикл получения обновлений (long polling).

Остановить цикл можно методом :meth:stop_polling, который дожидается выхода из самого цикла (а не задачи, вызвавшей этот метод: та может продолжать работу и после возврата отсюда).

Отмена задачи снаружи (task.cancel()) корректной остановкой не является: цикл прервётся, но фоновые задачи обработчиков (use_create_task=True) не будут дожданы, а изоляция событий не будет закрыта. Останавливайте через :meth:stop_polling либо вызовите :meth:shutdown после отмены. Перед новым запуском дождитесь отменённой задачи (await task с подавлением CancelledError): пока она не завершилась, повторный вызов будет отклонён или задержан (см. ниже).

Повторный вызов на ЖИВОМ цикле — RuntimeError. Если же цикл уже вышел, но уборка за прошлым запуском ещё идёт — отложенный дренаж фоновых задач от инлайн-обработчика (см. :meth:shutdown) или shutdown() внешнего :meth:stop_polling, — вызов не отклоняется, а ждёт её окончания и только затем стартует. Иначе уборка прошлого запуска закрыла бы изоляцию уже нового цикла. Благодаря этому идиома while True: await dp.start_polling(bot) переживает остановку снаружи. Перезапускать цикл нужно именно снаружи обработчиков: ожидание уборки из задачи, которую эта же уборка дренирует, замкнуло бы кольцо.

Ручная остановка через dp.polling = False (старый идиом) оставляет висеть текущий запрос get_updates до его таймаута. Если флаг сброшен снаружи между пачками (после получения ответа, но до начала его диспетчеризации), пачка целиком пропускается — маркер не сдвинут, и эти события придут снова при следующем запуске. А вот сброс флага инлайн-обработчиком посреди диспетчеризации самой пачки (use_create_task=False) на неё уже не влияет: цикл по событиям пачки не проверяет self.polling на каждой итерации, поэтому остаток пачки дорабатывается как обычно и маркер сдвигается.

Parameters:

Name Type Description Default
bot Bot

Экземпляр бота.

required
skip_updates bool

Флаг, отвечающий за обработку старых событий.

False

Raises:

Type Description
RuntimeError

Если цикл polling на этом диспетчере ещё жив.

Source code in maxapi/dispatcher.py
async def start_polling(
    self, bot: Bot, *, skip_updates: bool = False
) -> None:
    """
    Запускает цикл получения обновлений (long polling).

    Остановить цикл можно методом :meth:`stop_polling`, который
    дожидается выхода из самого цикла (а не задачи, вызвавшей этот
    метод: та может продолжать работу и после возврата отсюда).

    Отмена задачи снаружи (``task.cancel()``) корректной остановкой
    не является: цикл прервётся, но фоновые задачи обработчиков
    (``use_create_task=True``) не будут дожданы, а изоляция событий
    не будет закрыта. Останавливайте через :meth:`stop_polling`
    либо вызовите :meth:`shutdown` после отмены. Перед новым
    запуском дождитесь отменённой задачи (``await task`` с
    подавлением ``CancelledError``): пока она не завершилась,
    повторный вызов будет отклонён или задержан (см. ниже).

    Повторный вызов на ЖИВОМ цикле — ``RuntimeError``. Если же цикл
    уже вышел, но уборка за прошлым запуском ещё идёт — отложенный
    дренаж фоновых задач от инлайн-обработчика (см.
    :meth:`shutdown`) или ``shutdown()`` внешнего
    :meth:`stop_polling`, — вызов не отклоняется, а ждёт её
    окончания и только затем стартует. Иначе уборка прошлого
    запуска закрыла бы изоляцию уже нового цикла. Благодаря этому
    идиома ``while True: await dp.start_polling(bot)`` переживает
    остановку снаружи. Перезапускать цикл нужно именно снаружи
    обработчиков: ожидание уборки из задачи, которую эта же уборка
    дренирует, замкнуло бы кольцо.

    Ручная остановка через ``dp.polling = False`` (старый идиом)
    оставляет висеть текущий запрос ``get_updates`` до его
    таймаута. Если флаг сброшен снаружи между пачками (после
    получения ответа, но до начала его диспетчеризации), пачка
    целиком пропускается — маркер не сдвинут, и эти события
    придут снова при следующем запуске. А вот сброс флага
    инлайн-обработчиком посреди диспетчеризации самой пачки
    (``use_create_task=False``) на неё уже не влияет: цикл по
    событиям пачки не проверяет ``self.polling`` на каждой
    итерации, поэтому остаток пачки дорабатывается как обычно и
    маркер сдвигается.

    Args:
        bot: Экземпляр бота.
        skip_updates: Флаг, отвечающий за обработку старых событий.

    Raises:
        RuntimeError: Если цикл polling на этом диспетчере ещё жив.
    """
    while True:
        if self._polling_active:
            msg = (
                "Polling уже запущен на этом диспетчере. "
                "Остановите его через stop_polling() перед новым "
                "запуском либо используйте отдельный Dispatcher."
            )
            raise RuntimeError(msg)

        cleanup_done = self._cleanup_done
        if cleanup_done is None or self._lifecycle_holders == 0:
            break

        # Цикл вышел, но за прошлым запуском ещё убирают: ждём,
        # иначе та уборка закрыла бы изоляцию уже нашего цикла.
        # После пробуждения проверяем всё заново: пока мы ждали,
        # старт мог перехватить кто-то другой.
        logger_dp.debug("Жду окончания уборки за прошлым запуском")
        await cleanup_done.wait()

    self._polling_active = True
    self._lifecycle_holders += 1
    self._cleanup_done = asyncio.Event()
    self.polling = True
    self._polling_task = asyncio.current_task()
    self._stop_event = asyncio.Event()
    self._polling_error = None
    # Ожидающие остановки ждут именно это событие, а не задачу
    # вызывающего: она может продолжаться и после выхода из цикла.
    loop_done = self._loop_done = asyncio.Event()

    try:
        try:
            await self.__ready(bot)

            current_timestamp = to_ms(datetime.now())

            while self.polling:
                events = await self._fetch_updates_once(bot)
                if events is None:
                    # Recoverable-ошибка или остановка: пробуем
                    # снова (или выходим по условию цикла).
                    continue
                if not self.polling:
                    # Пачку, полученную уже после остановки, не
                    # обрабатываем: маркер не сдвинут, и эти события
                    # придут снова при следующем запуске.
                    continue
                await self._dispatch_fetched_events(
                    events, current_timestamp, skip_updates=skip_updates
                )
        except BaseException as e:
            # Ошибку запоминаем для stop_polling: он больше не
            # дожидается задачи и не может прочитать её exception().
            # Отмена ошибкой цикла не считается — о ней знает тот,
            # кто отменял.
            if not isinstance(e, asyncio.CancelledError):
                self._polling_error = e
            raise
        finally:
            self.polling = False
            self._polling_task = None
            self._stop_event = None
            # Остановка могла прийтись на подготовку (__ready): та
            # дописывает _ready=True уже после сброса в
            # stop_polling, поэтому сбрасываем здесь — иначе
            # следующий start_polling молча пропустил бы check_me
            # и on_started.
            self._ready = False
            # Цикл больше не жив: с этого момента повторный
            # start_polling не отклоняется, а ждёт уборки.
            self._polling_active = False
            # Будим ожидающих сразу после выхода из цикла: ждать
            # отложенного дренажа им нельзя — среди дренируемых
            # задач может быть та самая, что вызвала stop_polling.
            loop_done.set()
    finally:
        try:
            if self._deferred_shutdown:
                # Инлайн-обработчик остановил диспетчер сам: дренаж
                # был отложен, теперь мы вне handle() и можем
                # дождаться фоновых задач и закрыть изоляцию. При
                # внешней отмене задачи (task.cancel()) флаг обычно
                # не выставлен, и лишнего await здесь нет. Но если
                # инлайн-обработчик успел взвести флаг, а затем
                # задачу всё же отменили, отложенный shutdown всё
                # равно выполнится — этот finally отрабатывает и в
                # процессе отмены.
                self._deferred_shutdown = False
                await self.shutdown()
        finally:
            # Держателя уборки отпускаем последним: до этого
            # момента повторный start_polling ждёт, иначе уборка
            # отсюда задела бы уже новый цикл. Отпускаем и при
            # внешней отмене — этот finally отрабатывает и в
            # процессе отмены, счётчик не залипает.
            self._loop_done = None
            self._release_lifecycle()

stop_polling() async

Останавливает цикл получения обновлений (long polling).

Прерывает висящий запрос get_updates и паузы между попытками, после чего дожидается выхода из цикла :meth:start_polling и всех фоновых задач (use_create_task=True), запущенных до момента остановки. После возврата из метода никакой активности диспетчера не остаётся.

Ожидается именно цикл, а не задача, вызвавшая :meth:start_polling: при await dp.start_polling(bot) внутри более крупной корутины та задача продолжает работу и после остановки — ждать её означало бы дедлок, если её продолжение ждёт останавливающий обработчик.

Сетевые вызовы этапа старта (check_me, проверка подписок) и колбэк on_started не прерываются: остановка дождётся их завершения и только потом вернёт управление.

Если метод вызван из обработчика, выполняющегося прямо в задаче polling (use_create_task=False), дожидаться цикла нельзя — задача не может дождаться саму себя. В этом случае выставляются только флаги, а цикл завершится сразу после возврата из обработчика; дренаж фоновых задач и закрытие изоляции произойдут сразу после выхода из цикла (см. :meth:shutdown).

Отложенного дренажа (инлайн-остановка) метод не ждёт: среди дренируемых задач может быть та самая, из которой вызвана остановка. Пока эта уборка идёт, повторный :meth:start_polling не отклоняется, а ждёт её окончания.

Зато на время ожидания цикла и собственного shutdown() метод сам удерживает уборку за прошлым запуском: новый :meth:start_polling дождётся возврата отсюда и только затем стартует — иначе этот shutdown() дренировал бы фон и закрывал изоляцию уже нового цикла.

Ручная остановка через dp.polling = False полноценной заменой не является: висящий запрос get_updates не прерывается, уже полученная пачка не диспетчеризуется (придёт снова при следующем запуске), а фоновые задачи и изоляция остаются на совести вызывающего.

Вызов до фактического старта цикла (create_task на :meth:start_polling без единого await между ними) — no-op: задачи ещё нет, флаг polling не выставлен, и цикл потом запустится как обычно. Дайте задаче стартовать (например, await asyncio.sleep(0)) перед остановкой.

Source code in maxapi/dispatcher.py
async def stop_polling(self) -> None:
    """
    Останавливает цикл получения обновлений (long polling).

    Прерывает висящий запрос ``get_updates`` и паузы между
    попытками, после чего дожидается выхода из цикла
    :meth:`start_polling` и всех фоновых задач
    (``use_create_task=True``), запущенных до момента остановки.
    После возврата из метода никакой активности диспетчера не
    остаётся.

    Ожидается именно цикл, а не задача, вызвавшая
    :meth:`start_polling`: при ``await dp.start_polling(bot)``
    внутри более крупной корутины та задача продолжает работу и
    после остановки — ждать её означало бы дедлок, если её
    продолжение ждёт останавливающий обработчик.

    Сетевые вызовы этапа старта (``check_me``, проверка подписок) и
    колбэк ``on_started`` не прерываются: остановка дождётся их
    завершения и только потом вернёт управление.

    Если метод вызван из обработчика, выполняющегося прямо в
    задаче polling (``use_create_task=False``), дожидаться цикла
    нельзя — задача не может дождаться саму себя. В этом случае
    выставляются только флаги, а цикл завершится сразу после
    возврата из обработчика; дренаж фоновых задач и закрытие
    изоляции произойдут сразу после выхода из цикла
    (см. :meth:`shutdown`).

    Отложенного дренажа (инлайн-остановка) метод не ждёт: среди
    дренируемых задач может быть та самая, из которой вызвана
    остановка. Пока эта уборка идёт, повторный
    :meth:`start_polling` не отклоняется, а ждёт её окончания.

    Зато на время ожидания цикла и собственного ``shutdown()``
    метод сам удерживает уборку за прошлым запуском: новый
    :meth:`start_polling` дождётся возврата отсюда и только затем
    стартует — иначе этот ``shutdown()`` дренировал бы фон и
    закрывал изоляцию уже нового цикла.

    Ручная остановка через ``dp.polling = False`` полноценной
    заменой не является: висящий запрос ``get_updates`` не
    прерывается, уже полученная пачка не диспетчеризуется (придёт
    снова при следующем запуске), а фоновые задачи и изоляция
    остаются на совести вызывающего.

    Вызов до фактического старта цикла (``create_task`` на
    :meth:`start_polling` без единого await между ними) — no-op:
    задачи ещё нет, флаг ``polling`` не выставлен, и цикл потом
    запустится как обычно. Дайте задаче стартовать (например,
    ``await asyncio.sleep(0)``) перед остановкой.
    """
    if self.polling:
        self.polling = False
        self._ready = False
        if self._stop_event is not None:
            self._stop_event.set()
        logger_dp.info("Останавливаю polling")

    loop_done = self._loop_done
    if self._polling_task is asyncio.current_task():
        # Инлайн-стоп: дождаться цикла из него самого нельзя.
        loop_done = None

    if loop_done is not None:
        # Держим уборку за прошлым запуском, пока ждём цикл и
        # дренируем фон. Инлайн-стоп и стоп без запущенного цикла
        # счётчик не трогают: первый убирает не здесь, а в finally
        # цикла, второму убирать не за кем.
        self._lifecycle_holders += 1

    try:
        if loop_done is not None:
            await loop_done.wait()

            logger_dp.info("Polling остановлен")

            error = self._polling_error
            if error is not None:
                # Забираем ошибку, чтобы конкурентные остановки и
                # следующие вызовы не повторяли одно сообщение.
                self._polling_error = None
                logger_dp.error(
                    "Цикл polling завершился с ошибкой: %r",
                    error,
                )

        await self.shutdown()
    finally:
        if loop_done is not None:
            self._release_lifecycle()

shutdown() async

Завершает работу диспетчера: дожидается фоновых задач (use_create_task=True) и освобождает ресурсы изоляции событий.

Дожидается в цикле до полного опустошения пула: задача, добавленная конкурентным продюсером во время ожидания текущего снимка _background_tasks, будет дождана на следующей итерации, а не останется вне ожидаемого набора. Продюсеры должны быть остановлены до вызова (:meth:stop_polling сбрасывает polling заранее; webhook-интеграции вызывают shutdown после остановки приёма запросов).

Реентрантный вызов (из обработчика — обычно через :meth:stop_polling) только выставляет признак завершения: дренировать фоновые задачи и закрывать изоляцию нельзя. Другие обработчики того же пользователя ждут блокировку event_isolation, которую удерживает вызывающий, — ожидание их завершения замкнуло бы кольцо. Оставшиеся задачи доработают сами; чтобы дождаться их, вызовите shutdown() снаружи обработчика. Исключение — инлайн-обработчик в задаче polling: для него дренаж откладывается до выхода из цикла и выполняется автоматически (см. :meth:start_polling).

Реентрантным считается вызов из ТОЙ ЖЕ задачи, которая прямо сейчас выполняет :meth:handle ЭТОГО диспетчера (маркер _in_handler). Способ запуска обработчика роли не играет: так опознаются и инлайн-обработчик в задаче polling, и задача из :meth:spawn_handle_task, и вебхук с use_create_task=False, где handle() вызывается прямо в задаче HTTP-запроса. Кольцо всё же возможно, если обработчик сам дожидается порождённой им задачи, которая вызывает shutdown(): у дочерней задачи current_task() другой, и её вызов реентрантным не считается.

Отложенный дренаж (см. выше про инлайн-обработчик) выполняется только при выходе из цикла polling. Поэтому shutdown() из инлайн-обработчика без последующей остановки polling ничего не завершает: флаг _deferred_shutdown остаётся взведённым до конца цикла, а _closing=True тем временем даёт warning в :meth:spawn_handle_task при постановке новых задач.

Готовность (_ready) метод не сбрасывает: повторный :meth:startup подготовку не повторяет (не будет ни check_me, ни on_started). Для полного перезапуска используйте :meth:stop_polling.

Вызывается автоматически из :meth:stop_polling и из shutdown-хуков webhook-интеграций (:class:~maxapi.webhook.base.BaseMaxWebhook). Идемпотентен.

Source code in maxapi/dispatcher.py
async def shutdown(self) -> None:
    """
    Завершает работу диспетчера: дожидается фоновых задач
    (``use_create_task=True``) и освобождает ресурсы изоляции
    событий.

    Дожидается в цикле до полного опустошения пула: задача,
    добавленная конкурентным продюсером во время ожидания
    текущего снимка ``_background_tasks``, будет дождана на
    следующей итерации, а не останется вне ожидаемого набора.
    Продюсеры должны быть остановлены до вызова
    (:meth:`stop_polling` сбрасывает ``polling`` заранее;
    webhook-интеграции вызывают shutdown после остановки приёма
    запросов).

    Реентрантный вызов (из обработчика — обычно через
    :meth:`stop_polling`) только выставляет признак завершения:
    дренировать фоновые задачи и закрывать изоляцию нельзя. Другие
    обработчики того же пользователя ждут блокировку
    ``event_isolation``, которую удерживает вызывающий, — ожидание
    их завершения замкнуло бы кольцо. Оставшиеся задачи доработают
    сами; чтобы дождаться их, вызовите ``shutdown()`` снаружи
    обработчика. Исключение — инлайн-обработчик в задаче polling:
    для него дренаж откладывается до выхода из цикла и выполняется
    автоматически (см. :meth:`start_polling`).

    Реентрантным считается вызов из ТОЙ ЖЕ задачи, которая прямо
    сейчас выполняет :meth:`handle` ЭТОГО диспетчера (маркер
    ``_in_handler``). Способ запуска обработчика роли не играет:
    так опознаются и инлайн-обработчик в задаче polling, и задача
    из :meth:`spawn_handle_task`, и вебхук с
    ``use_create_task=False``, где ``handle()`` вызывается прямо в
    задаче HTTP-запроса. Кольцо всё же возможно, если обработчик
    сам дожидается порождённой им задачи, которая вызывает
    ``shutdown()``: у дочерней задачи ``current_task()`` другой,
    и её вызов реентрантным не считается.

    Отложенный дренаж (см. выше про инлайн-обработчик) выполняется
    только при выходе из цикла polling. Поэтому ``shutdown()`` из
    инлайн-обработчика без последующей остановки polling ничего не
    завершает: флаг ``_deferred_shutdown`` остаётся взведённым до
    конца цикла, а ``_closing=True`` тем временем даёт warning
    в :meth:`spawn_handle_task` при постановке новых задач.

    Готовность (``_ready``) метод не сбрасывает: повторный
    :meth:`startup` подготовку не повторяет (не будет ни
    ``check_me``, ни ``on_started``). Для полного перезапуска
    используйте :meth:`stop_polling`.

    Вызывается автоматически из :meth:`stop_polling` и из
    shutdown-хуков webhook-интеграций
    (:class:`~maxapi.webhook.base.BaseMaxWebhook`). Идемпотентен.
    """
    self._closing = True

    # Одного признака «мы внутри handle()» мало: ContextVar
    # наследуется в create_task и общий для всех диспетчеров
    # процесса. Поэтому в маркере лежит пара (диспетчер, задача):
    # реентрантен вызов только из той же задачи и для того же
    # диспетчера.
    current = asyncio.current_task()
    marker = _in_handler.get()
    reentrant = (
        marker is not None and marker[0] is self and marker[1] is current
    )

    if reentrant:
        if current is self._polling_task:
            # Инлайн-обработчик (use_create_task=False) выполняется
            # в самой задаче polling: дренаж и закрытие изоляции
            # откладываем до выхода из цикла, там мы уже вне
            # handle() (см. finally в start_polling).
            self._deferred_shutdown = True
            logger_dp.debug(
                "shutdown вызван из инлайн-обработчика: дренаж "
                "отложен до завершения цикла polling",
            )
        elif others := self._background_tasks - {current}:
            logger_dp.warning(
                "shutdown вызван из обработчика: дренаж фоновых "
                "задач (%d) и закрытие изоляции пропущены",
                len(others),
            )
        else:
            logger_dp.debug(
                "shutdown вызван из обработчика: дренировать "
                "нечего, изоляция не закрыта",
            )
        return

    drained = False
    while pending := tuple(self._background_tasks):
        logger_dp.info(
            "Ожидаю завершения %d фоновых задач...",
            len(pending),
        )
        # Именно wait(), а не gather(): задачу из пула удаляет её
        # done-callback, который мог ещё не выполниться (задача
        # завершилась, а нас разбудили раньше). С Python 3.12
        # gather() по уже завершённым задачам не уступает цикл
        # событий, и while крутился бы вечно. wait() всегда
        # уступает, а call_soon выполняет callback'и по порядку —
        # к пробуждению задачи убраны из пула, их ошибки
        # залогированы.
        try:
            await asyncio.wait(pending)
        except asyncio.CancelledError:
            # gather() отменял бы дочерние задачи вместе с собой;
            # wait() этого не делает — сохраняем поведение: отмена
            # shutdown() (например, по таймауту lifespan) отменяет
            # и недождавшиеся фоновые задачи.
            for task in pending:
                task.cancel()
            # cancel() лишь запрашивает отмену: как и gather(),
            # дожидаемся, пока задачи доработают cleanup и уйдут
            # из пула, и только потом пробрасываем отмену.
            await asyncio.wait(pending)
            raise
        drained = True
    if drained:
        logger_dp.info("Все фоновые задачи завершены")

    await self.event_isolation.close()

startup(bot) async

Инициализирует диспетчер: сохраняет бота, подготавливает обработчики и вызывает on_started.

Используется интеграционными модулями (например, maxapi.webhook.fastapi) для инициализации в lifespan веб-фреймворка.

Parameters:

Name Type Description Default
bot Bot

Экземпляр бота.

required
Source code in maxapi/dispatcher.py
async def startup(self, bot: Bot) -> None:
    """
    Инициализирует диспетчер: сохраняет бота, подготавливает
    обработчики и вызывает on_started.

    Используется интеграционными модулями (например,
    maxapi.webhook.fastapi) для инициализации в lifespan
    веб-фреймворка.

    Args:
        bot: Экземпляр бота.
    """
    await self.__ready(bot)

handle_webhook(bot, *, host=DEFAULT_HOST, port=DEFAULT_PORT, path=DEFAULT_PATH, secret=None, webhook_type=AiohttpMaxWebhook, **kwargs) async

Запускает вебхук-сервер (aiohttp) для приёма обновлений.

Удобный метод «всё в одном»: создаёт aiohttp-приложение через :class:~maxapi.webhook.aiohttp.BaseMaxWebhook, регистрирует маршрут и запускает сервер.

Для более гибкого управления жизненным циклом сервера используйте одну из реализаций BaseMaxWebhook напрямую, например :class:~maxapi.webhook.aiohttp.BaseMaxWebhook.

Parameters:

Name Type Description Default
bot Bot

Экземпляр бота.

required
host str

Хост сервера (по умолчанию "0.0.0.0").

DEFAULT_HOST
port int

Порт сервера (по умолчанию 8080).

DEFAULT_PORT
path str

URL-путь для маршрута вебхука.

DEFAULT_PATH
secret str | None

Секрет для проверки заголовка X-Max-Bot-Api-Secret. Должен совпадать со значением, переданным в :meth:~maxapi.Bot.subscribe_webhook.

None
webhook_type type[BaseMaxWebhook]

Класс вебхука.

AiohttpMaxWebhook
**kwargs Any

Дополнительные аргументы для aiohttp.web.AppRunner.

{}
Source code in maxapi/dispatcher.py
async def handle_webhook(
    self,
    bot: Bot,
    *,
    host: str = DEFAULT_HOST,
    port: int = DEFAULT_PORT,
    path: str = DEFAULT_PATH,
    secret: str | None = None,
    webhook_type: type[BaseMaxWebhook] = AiohttpMaxWebhook,
    **kwargs: Any,
) -> None:
    """
    Запускает вебхук-сервер (aiohttp) для приёма обновлений.

    Удобный метод «всё в одном»: создаёт aiohttp-приложение через
    :class:`~maxapi.webhook.aiohttp.BaseMaxWebhook`,
    регистрирует маршрут и запускает сервер.

    Для более гибкого управления жизненным циклом сервера используйте
    одну из реализаций BaseMaxWebhook напрямую, например
    :class:`~maxapi.webhook.aiohttp.BaseMaxWebhook`.

    Args:
        bot: Экземпляр бота.
        host: Хост сервера (по умолчанию ``"0.0.0.0"``).
        port: Порт сервера (по умолчанию ``8080``).
        path: URL-путь для маршрута вебхука.
        secret: Секрет для проверки заголовка
            ``X-Max-Bot-Api-Secret``. Должен совпадать со значением,
            переданным в :meth:`~maxapi.Bot.subscribe_webhook`.
        webhook_type: Класс вебхука.
        **kwargs: Дополнительные аргументы для ``aiohttp.web.AppRunner``.
    """
    webhook = webhook_type(dp=self, bot=bot, secret=secret)
    await webhook.run(host=host, port=port, path=path, **kwargs)

init_serve(bot, host=DEFAULT_HOST, port=DEFAULT_PORT, **kwargs) async

.. deprecated:: Используйте :meth:handle_webhook вместо init_serve. Метод будет удалён в одной из следующих версий.

Parameters:

Name Type Description Default
bot Bot

Экземпляр бота.

required
host str

Хост.

DEFAULT_HOST
port int

Порт.

DEFAULT_PORT
Source code in maxapi/dispatcher.py
async def init_serve(  # pragma: no cover
    self,
    bot: Bot,
    host: str = DEFAULT_HOST,
    port: int = DEFAULT_PORT,
    **kwargs: Any,
) -> None:
    """
    .. deprecated::
        Используйте :meth:`handle_webhook` вместо ``init_serve``.
        Метод будет удалён в одной из следующих версий.

    Args:
        bot: Экземпляр бота.
        host: Хост.
        port: Порт.
    """
    warn(
        "init_serve устарел и будет удалён в следующих версиях. "
        "Используйте handle_webhook вместо него.",
        DeprecationWarning,
        stacklevel=2,
    )
    await self.handle_webhook(bot, host=host, port=port, **kwargs)

Router(router_id=None)

Bases: Dispatcher

Роутер для группировки обработчиков событий.

Инициализация роутера.

Parameters:

Name Type Description Default
router_id str | None

Идентификатор роутера для логов.

None
Source code in maxapi/dispatcher.py
def __init__(self, router_id: str | None = None):
    """
    Инициализация роутера.

    Args:
        router_id: Идентификатор роутера для логов.
    """

    super().__init__(router_id)

fsm property

Роутер не владеет FSM-хранилищем.

ErrorEventObserver(router)

Декоратор для регистрации обработчиков ошибок.

Инициализирует декоратор ошибок.

Parameters:

Name Type Description Default
router Dispatcher | Router

Экземпляр роутера или диспетчера.

required
Source code in maxapi/dispatcher.py
def __init__(self, router: Dispatcher | Router) -> None:
    """
    Инициализирует декоратор ошибок.

    Args:
        router: Экземпляр роутера или диспетчера.
    """
    self.router = router

register(func_event, *args, **_kwargs)

Регистрирует функцию как обработчик ошибки.

Parameters:

Name Type Description Default
func_event Callable

Функция-обработчик ошибки.

required
*args Any

Типы исключений или фильтры.

()

Returns:

Name Type Description
Callable Callable

Исходная функция.

Source code in maxapi/dispatcher.py
def register(
    self, func_event: Callable, *args: Any, **_kwargs: Any
) -> Callable:
    """
    Регистрирует функцию как обработчик ошибки.

    Args:
        func_event: Функция-обработчик ошибки.
        *args: Типы исключений или фильтры.

    Returns:
        Callable: Исходная функция.
    """
    self.router.error_handlers.append(
        ErrorHandler(*args, func_event=func_event)
    )
    # Обработчики ошибок читаются напрямую из router.error_handlers,
    # но инвалидация всё равно нужна: только перестройка заполняет
    # error_handler.func_args (без него call_error_handler на каждой
    # ошибке заново разбирает сигнатуру через inspect).
    self.router._invalidate_handlers()  # noqa: SLF001
    return func_event

__call__(*args, **kwargs)

Регистрирует функцию как обработчик ошибки через декоратор.

Returns:

Name Type Description
Callable Callable

Декоратор.

Source code in maxapi/dispatcher.py
def __call__(self, *args: Any, **kwargs: Any) -> Callable:
    """
    Регистрирует функцию как обработчик ошибки через декоратор.

    Returns:
        Callable: Декоратор.
    """

    def decorator(func_event: Callable) -> Callable:
        return self.register(func_event, *args, **kwargs)

    return decorator

Event(update_type, router, *, deprecated=False)

Декоратор для регистрации обработчиков событий.

Инициализирует событие-декоратор.

Parameters:

Name Type Description Default
update_type UpdateType

Тип события.

required
router Dispatcher | Router

Экземпляр роутера или диспетчера.

required
deprecated bool

Флаг, указывающий на то, что событие устарело.

False
Source code in maxapi/dispatcher.py
def __init__(
    self,
    update_type: UpdateType,
    router: Dispatcher | Router,
    *,
    deprecated: bool = False,
):
    """
    Инициализирует событие-декоратор.

    Args:
        update_type: Тип события.
        router: Экземпляр роутера или диспетчера.
        deprecated: Флаг, указывающий на то, что событие устарело.
    """

    self.update_type = update_type
    self.router = router
    self.deprecated = deprecated

register(func_event, *args, **kwargs)

Регистрирует функцию как обработчик события.

Parameters:

Name Type Description Default
func_event Callable

Функция-обработчик

required
*args Any

Фильтры

()
**kwargs Any

Дополнительные параметры (например, states)

{}

Returns:

Name Type Description
Callable Callable

Исходная функция.

Source code in maxapi/dispatcher.py
def register(
    self, func_event: Callable, *args: Any, **kwargs: Any
) -> Callable:
    """
    Регистрирует функцию как обработчик события.

    Args:
        func_event: Функция-обработчик
        *args: Фильтры
        **kwargs: Дополнительные параметры (например, states)

    Returns:
        Callable: Исходная функция.
    """

    if self.deprecated:
        warnings.warn(
            f"Событие {self.update_type} устарело "
            f"и будет удалено в будущих версиях.",
            DeprecationWarning,
            stacklevel=3,
        )

    if self.update_type == UpdateType.ON_STARTED:
        # Предупреждаем только тогда, когда колбэк действительно
        # опоздал: подготовка уже пройдена и не сброшена
        # (``_ready``) либо колбэк регистрируют изнутри самого
        # on_started (``_running_on_started``). Регистрация в
        # окне подготовки ДО фазы on_started (например из-под
        # долгого check_me) не опоздала: колбэк будет прочитан и
        # вызван в этом же запуске. После stop_polling() бот
        # остаётся привязан, но подготовка сброшена, и следующий
        # start_polling колбэк вызовет — предупреждать там не о
        # чем.
        if (
            self.router._ready  # noqa: SLF001
            or self.router._running_on_started  # noqa: SLF001
        ):
            logger_dp.warning(
                "Колбэк on_started зарегистрирован слишком поздно: "
                "он не будет вызван, подготовка диспетчера уже "
                "выполнена либо колбэк регистрируется изнутри "
                "самого on_started.",
            )
        self.router.on_started_func = func_event

    else:
        self.router.event_handlers.append(
            Handler(
                *args,
                func_event=func_event,
                update_type=self.update_type,
                **kwargs,
            )
        )
        self.router._invalidate_handlers()  # noqa: SLF001
    return func_event

__call__(*args, **kwargs)

Регистрирует функцию как обработчик события через декоратор.

Returns:

Name Type Description
Callable Callable

Декоратор.

Source code in maxapi/dispatcher.py
def __call__(self, *args: Any, **kwargs: Any) -> Callable:
    """
    Регистрирует функцию как обработчик события через декоратор.

    Returns:
        Callable: Декоратор.
    """

    def decorator(func_event: Callable) -> Callable:
        return self.register(func_event, *args, **kwargs)

    return decorator