Spec-Zone.ru › Trio

Просмотр и расширение Trio с помощью trio.lowlevel

trio.lowlevel содержит низкоуровневые API для просмотра и расширения Trio. Если вы пишете обычный повседневный код, то можете полностью пропустить этот модуль. Но иногда вам нужно что-то немного более низкого уровня. Вот несколько примеров ситуаций, когда вам потребуется trio.lowlevel:

  • Вы хотите реализовать новую примитив синхронизации, которой Trio (пока) не предоставляет, например, блокировку с чтением-записью.

  • Вы хотите извлечь низкоуровневые метрики для мониторинга состояния вашего приложения.

  • Вы хотите использовать низкоуровневый интерфейс операционной системы, для которого Trio пока не предоставляет собственные обёртки, например, наблюдение за каталогом файловой системы на предмет изменений.

  • Вы хотите реализовать интерфейс для вызовов между Trio и другой циклической обработкой событий в одном процессе.

  • Вы пишете отладчик и хотите визуализировать дерево задач Trio.

  • Вам нужно взаимодействовать с библиотекой C, API которой предоставляет исходные дескрипторы файлов.

Вам не нужно бояться trio.lowlevel, если вы примете надлежащие меры предосторожности. Это реальные публичные API с жёстко определённой и тщательно задокументированной семантикой. Это те же инструменты, которые мы используем для реализации всех удобных высокоуровневых API в пространстве имён trio. Но будьте осторожны. Некоторые из этих строгих семантик имеют острые большие зубы. Если вы допустите ошибку, Trio может не обработать её должным образом; соглашения и гарантии, которые строго соблюдаются в остальной части Trio, не всегда применяются. Когда вы используете этот модуль, вы несёте ответственность за то, как вы будете обрабатывать сложные случаи, чтобы предоставить своим пользователям удобный API в стиле Trio.

Отладка и инструментирование

Trio старается предоставить полезные крючки для отладки и инструментирования. Некоторые из них задокументированы выше (атрибуты интроспекции яслей, trio.Lock.statistics() и т. д.). Вот ещё несколько.

Глобальная статистика

trio.lowlevel.current_statistics() → RunStatistics

Возвращает объект, содержащий отладочную информацию на уровне цикла обработки событий:

class trio.lowlevel.RunStatistics

Объект, содержащий отладочную информацию на уровне цикла обработки событий.

В настоящее время определены следующие поля:

  • tasks_living (целое число): Количество задач, которые были созданы и ещё не завершены.

  • tasks_runnable (целое число): Количество задач, которые в настоящее время находятся в очереди выполнения (по сравнению с заблокированными, ожидающими чего-то).

  • seconds_to_next_deadline (вещественное число): Время до следующей запланированной даты истечения срока действия области отмены. Может быть отрицательным, если срок действия истек, но мы ещё не обработали отмены. Может быть inf, если нет запланированных сроков действия.

  • run_sync_soon_queue_size (целое число): Количество необработанных обратных вызовов, поставленных в очередь с помощью trio.lowlevel.TrioToken.run_sync_soon().

  • io_statistics (объект): Некоторые статистические данные из I/O-бэкенда Trio. У него всегда есть атрибут backend, который является строкой, называющей используемый бэкенд ввода-вывода операционной системы; другие атрибуты варьируются в зависимости от бэкенда.

Проверка наличия Trio

Если вам нужно взаимодействовать с активным запуском Trio — например, вам нужно узнать current_time() или current_task() — Trio нужно иметь определённое состояние, иначе вы получите RuntimeError("must be called from async context"). Это требует, чтобы вы либо:

  • были косвенно внутри (и в том же потоке, что и) вызова trio.run(), для получения информации на уровне запуска, такой как current_time() или current_clock(); или

  • были косвенно внутри задачи Trio, для получения информации на уровне задачи, такой как current_task() или current_effective_deadline().

Внутренне это состояние предоставляется переменными потока, отслеживающими текущий запуск и текущую задачу. Иногда полезно заранее знать, потерпит ли неудачу вызов, или иметь динамичную информацию для защиты от запуска чего-либо внутри или вне Trio. Для этого вызовите trio.lowlevel.in_trio_run() или trio.lowlevel.in_trio_task(), которые дадут ответы в соответствии со следующей таблицей.

ситуация

trio.lowlevel.in_trio_run()

trio.lowlevel.in_trio_task()

внутри асинхронной функции Trio

True

True

в потоке без активного вызова trio.run()

False

False

в хост-цикле гостевого запуска

True

False

внутри вызова инструмента

True

зависит

в потоке, созданном trio.to_thread.run_sync()

False

False

внутри функции прерывания

True

True

trio.lowlevel.in_trio_run() → bool

Проверяет, находимся ли мы в запуске Trio. Возвращает True только в том случае, если current_time() будет успешным.

См. также обсуждение различных способов обнаружения Trio.

trio.lowlevel.in_trio_task() → bool

Проверяет, находимся ли мы в задаче Trio. Возвращает True только в том случае, если current_task() будет успешным.

См. также обсуждение различных способов обнаружения Trio.

Текущий таймер

trio.lowlevel.current_clock() → Clock

Возвращает текущий Clock.

API инструментов

API инструментов предоставляет стандартный способ добавления пользовательской инструментики в цикл выполнения. Хотите построить гистограмму задержек планирования, записать трассировку стека любого задания, которое блокирует цикл выполнения более чем на 50 мс, или измерить процент времени работы процесса, потраченного на ожидание ввода-вывода? Это то место.

Общая идея заключается в том, что в любой момент времени trio.run() поддерживает набор «инструментов», которые представляют собой объекты, реализующие интерфейс trio.abc.Instrument. Когда происходит интересное событие, он перебирает эти инструменты и уведомляет их, вызывая соответствующий метод. Учебник содержит простой пример использования этого для отслеживания.

Поскольку это подключение к Trio на довольно низком уровне, нужно быть осторожным. Обратные вызовы выполняются синхронно, и во многих случаях, если они завершатся ошибкой, нет никакого разумного способа распространить эту ошибку (например, мы можем быть глубоко в механизме обработки исключений...). Поэтому наша текущая стратегия обработки исключений, поднятых инструментами, заключается в (а) записи исключения в журнал "trio.abc.Instrument", который по умолчанию выводит трассировку стека в стандартный поток ошибок, и (б) отключении нарушающего инструмент.

Вы можете зарегистрировать начальный список инструментов, передав их в trio.run(). add_instrument() и remove_instrument() позволяют добавлять и удалять инструменты во время выполнения.

trio.lowlevel.add_instrument(instrument: Instrument) → None

Начать инструментацию текущего цикла выполнения заданным инструментом.

Параметры:

instrument (trio.abc.Instrument) – Инструмент, который нужно активировать.

Если instrument уже активен, ничего не делает.

trio.lowlevel.remove_instrument(instrument: Instrument) → None

Остановка инструментации текущего цикла выполнения заданным инструментом.

Параметры:

instrument (trio.abc.Instrument) – Инструмент, который нужно деактивировать.

Возбуждает:

KeyError – если инструмент не активен. Это может произойти либо потому, что вы его никогда не добавляли, либо потому, что вы его добавили, а затем он поднял необработанное исключение и был автоматически деактивирован.

И вот интерфейс для реализации, если вы хотите создать свой собственный Instrument:

class trio.abc.Instrument

Интерфейс для инструментации цикла выполнения.

Инструменты не обязаны наследоваться от этого абстрактного базового класса, и все эти методы являются необязательными. Этот класс служит в основном для документации.

after_io_wait(timeout: float) → None

Вызывается после обработки ожидающего ввода-вывода.

Параметры:

timeout (float) – Количество секунд, на которое мы были готовы ждать. Это время может или не может истечь, в зависимости от того, был ли готов какой-либо ввод-вывод.

after_run() → None

Вызывается непосредственно перед возвращением trio.run().

after_task_step(task: Task) → None

Вызывается при возвращении в основной цикл выполнения после того, как задание сгенерировало выход.

Параметры:

task (trio.lowlevel.Task) – Задание, которое только что выполнилось.

before_io_wait(timeout: float) → None

Вызывается перед блокированием для ожидания готовности ввода-вывода.

Параметры:

timeout (float) – Количество секунд, на которое мы готовы ждать.

before_run() → None

Вызывается в начале trio.run().

before_task_step(task: Task) → None

Вызывается непосредственно перед возобновлением выполнения заданного задания.

Параметры:

task (trio.lowlevel.Task) – Задание, которое скоро будет выполнено.

task_exited(task: Task) → None

Вызывается, когда заданное задание завершается.

Параметры:

task (trio.lowlevel.Task) – Завершенное задание.

task_scheduled(task: Task) → None

Вызывается, когда заданное задание становится выполнимым.

Может пройти некоторое время, прежде чем оно фактически будет выполнено, если перед ним есть другие выполнимые задания.

Параметры:

task (trio.lowlevel.Task) – Задание, которое стало выполнимым.

task_spawned(task: Task) → None

Вызывается при создании заданного задания.

Параметры:

task (trio.lowlevel.Task) – Новое задание.

В учебнике есть полный пример определения пользовательского инструмента для регистрации внутренних решений планирования Trio.

Запуск процессов низкого уровня

await trio.lowlevel.open_process(command: str | bytes | os.PathLike | Sequence[str | bytes | os.PathLike], *, stdin: int | HasFileno | None = None, stdout: int | HasFileno | None = None, stderr: int | HasFileno | None = None, **options: object) → Process

Выполняет дочернюю программу в новом процессе.

После создания вы можете взаимодействовать с дочерним процессом, записывая данные в его поток stdin (объект SendStream), считывая данные из его потоков stdout и/или stderr (оба являются объектами ReceiveStream), отправляя ему сигналы с помощью terminate, kill или send_signal, и ожидая его завершения с помощью wait. Подробности см. в trio.Process.

Каждый стандартный поток доступен только в том случае, если вы укажете, что для него должен быть создан канал. Например, если вы передадите stdin=subprocess.PIPE, вы можете писать в поток stdin, в противном случае stdin будет None.

В отличие от trio.run_process, эта функция не выполняет автоматическое управление дочерним процессом. Вам необходимо реализовать необходимые семантики самостоятельно.

Параметры:

  • command – Команда для выполнения. Как правило, это последовательность строк или байтов, например, ['ls', '-l', 'directory with spaces'], где первый элемент - имя исполняемого файла, а другие элементы - его аргументы. С shell=True в **options или на Windows command может быть строкой или байтами, которые будут обработаны в соответствии с платформозависимыми правилами цитирования. Во всех случаях command может быть путем или последовательностью путей.

  • stdin – Указывает, к чему должен быть подключен стандартный входной поток дочернего процесса: вывод родительского процесса (subprocess.PIPE), ничего (subprocess.DEVNULL) или открытый файл (передать дескриптор файла или что-либо, чей метод fileno возвращает один). Если stdin не указан, у дочернего процесса будет такой же стандартный входной поток, как у родительского.

  • stdout – Аналогично stdin, но для стандартного выходного потока дочернего процесса.

  • stderr – Аналогично stdin, но для стандартного потока ошибок дочернего процесса. Поддерживается дополнительное значение subprocess.STDOUT, которое вызывает объединение сообщений стандартного вывода и стандартной ошибки дочернего процесса в один стандартный выходной поток, подключенный к тому, к чему подключен опцией stdout.

  • **options – Также принимаются другие общие параметры подпроцессов.

Возвращает:

Новый объект trio.Process.

Возможные исключения:

OSError – если запуск процесса завершился ошибкой, например, потому что указанная команда не найдена.

Примитивы низкоуровневого ввода-вывода

Разные среды предоставляют разные низкоуровневые API для выполнения асинхронного ввода-вывода. trio.lowlevel предоставляет эти API относительно прямым способом, чтобы предоставить максимальную мощность и гибкость для кода более высокого уровня. Однако это означает, что точный предоставляемый API может меняться в зависимости от системы, на которой работает Trio.

Универсально доступный API

Все среды предоставляют следующие функции:

await trio.lowlevel.wait_readable(obj)

Ожидает, пока ядро сообщит, что данный объект готов к чтению.

В системах Unix, obj должен быть целочисленным дескриптором файла или объектом с методом .fileno(), возвращающим целочисленный дескриптор файла. Можно передавать любой тип дескриптора файла, хотя точное поведение будет зависеть от ядра. Например, для файлов на диске это вряд ли будет полезно.

В системах Windows, obj должен быть целочисленным дескриптором SOCKET-ручки или объектом с методом .fileno(), возвращающим целочисленный дескриптор SOCKET-ручки. Дескрипторы файлов не поддерживаются, как и ручки, ссылающиеся на что-либо, кроме SOCKET.

Исключения:

  • trio.BusyResourceError – если другая задача уже ожидает, пока данный сокет станет готовым к чтению.

  • trio.ClosedResourceError – если другая задача вызывает notify_closing(), пока эта функция ещё работает.

await trio.lowlevel.wait_writable(obj)

Ожидает, пока ядро сообщит, что данный объект готов к записи.

См. wait_readable для определения obj.

Исключения:

  • trio.BusyResourceError – если другая задача уже ожидает, пока данный сокет станет готовым к записи.

  • trio.ClosedResourceError – если другая задача вызывает notify_closing(), пока эта функция ещё работает.

trio.lowlevel.notify_closing(obj)

Вызовите эту функцию перед закрытием дескриптора файла (в Unix) или сокета (в Windows). Это заставит любые вызовы wait_readable или wait_writable на данном объекте немедленно проснуться и выбросить ClosedResourceError.

Эта функция не закрывает объект – вам все равно нужно сделать это самостоятельно. Также следует быть осторожным, чтобы не запускать новые задачи, ожидающие объекта между вызовом этой функции и фактическим закрытием. Поэтому, чтобы корректно закрыть что-то, обычно нужно выполнить эти шаги в определённом порядке:

  1. Явно пометить объект как закрытый, чтобы любые новые попытки использовать его прервались до начала.

  2. Вызвать notify_closing, чтобы разбудить всех уже существующих пользователей.

  3. Фактически закрыть объект.

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

API, специфичное для Unix

FdStream поддерживает обертывание Unix-файлов (таких как пайп или TTY) в виде потока.

Если у вас есть два разных дескриптора файлов для отправки и получения и вы хотите объединить их в единый двунаправленный поток Stream, используйте trio.StapledStream:

bidirectional_stream = trio.StapledStream(
    trio.lowlevel.FdStream(write_fd),
    trio.lowlevel.FdStream(read_fd)
)

class trio.lowlevel.FdStream(fd: int)

Базируется на Stream

Представляет поток, используя дескриптор файла пайпа, TTY и т.п.

fd должен ссылаться на файл, открытый для чтения и/или записи, поддерживающий асинхронный ввод-вывод (пайпы и TTY будут работать, файлы на диске — вероятно, нет). Возвращаемый поток берет на себя владение fd, поэтому закрытие потока также закроет fd. Как и в os.fdopen, вы не должны напрямую использовать fd после его обертывания в поток с помощью этой функции.

Для использования в качестве потока Trio открытый файл должен быть помещён в режим без блокировки. К сожалению, это влияет на все ввод-вывод, проходящие через подлежащий открытый файл, включая ввод-вывод, использующий другой дескриптор файла, чем тот, который был передан Trio. Если другие потоки или процессы используют дескрипторы файлов, связанные через os.dup или наследование через os.fork к тому, который использует Trio, они вряд ли будут готовы к внезапному внедрению семантики асинхронного ввода-вывода. Например, можно использовать FdStream(os.dup(sys.stdin.fileno())) для получения потока для чтения со стандартного ввода, но это безопасно только при серьёзных оговорках: ваш стандартный ввод не должен быть общим ни с какими другими процессами, и вы не должны вызывать синхронные методы sys.stdin, пока поток, возвращённый FdStream, не будет закрыт. См. проблему #174 для обсуждения сложностей, связанных с ослаблением этого ограничения.

Параметры:

fd (int) – fd для обертывания.

Возвращаемое значение:

Новый объект FdStream.

API, специфичное для Kqueue

TODO: эти функции реализованы, но на данный момент представляют собой больше набросок, чем что-то реальное. См. #26.

trio.lowlevel.current_kqueue()

await trio.lowlevel.wait_kevent(ident, filter, abort_func)

with trio.lowlevel.monitor_kevent(ident, filter) as queue

API Windows

await trio.lowlevel.WaitForSingleObject(handle)

Асинхронный и отменяемый вариант WaitForSingleObject. Только для Windows.

Параметры:

handle – Движок Win32, как целое число Python.

Исключения:

OSError – Если движок недействителен, например, если он уже закрыт.

TODO: эти функции реализованы, но пока являются скорее набросками, чем чем-то реальным. См. #26 и #52.

trio.lowlevel.register_with_iocp(handle)

await trio.lowlevel.wait_overlapped(handle, lpOverlapped)

await trio.lowlevel.write_overlapped(handle, data)

await trio.lowlevel.readinto_overlapped(handle, data)

trio.lowlevel.current_iocp()

with trio.lowlevel.monitor_completion_key() as queue

Глобальное состояние: системные задачи и локальные переменные выполнения

class trio.lowlevel.RunVar(name: str, default | type[~trio._core._local._NoValue] = ...)

Локальный вариант переменной контекста.

RunVar объекты похожи на объекты переменных контекста, за исключением того, что они используются в рамках одного вызова trio.run(), а не одной задачи.

trio.lowlevel.spawn_system_task(async_fn: Callable[[Unpack[PosArgT]], Awaitable[object]], *args: Unpack[PosArgT], name: object = None, context: contextvars.Context | None = None) → Task

Запустить системную задачу.

Системные задачи отличаются от обычных задач:

  • Им не нужна явная подсистема; вместо этого они попадают во внутреннюю «системную подсистему».

  • Если системная задача вызывает исключение, то оно преобразуется в TrioInternalError и все задачи отменяются. При написании системной задачи следует быть осторожным, чтобы она не вызывала сбои.

  • Системные задачи автоматически отменяются при завершении основной задачи.

  • По умолчанию, для системных задач включена защита от KeyboardInterrupt. Если вам нужна возможность прерывания задачи нажатием Ctrl+C, то нужно явно использовать disable_ki_protection() (и разработать план действий при возникновении KeyboardInterrupt, так как системные задачи не могут поднимать исключения).

  • Системные задачи не наследуют переменные контекста от своего создателя.

По завершении вызова trio.run(), после завершения основной задачи и всех системных задач, системная подсистема закрывается. В этот момент новые вызовы spawn_system_task() будут вызывать RuntimeError("Nursery is closed to new arrivals") вместо создания системной задачи. Такое состояние может быть встречено в блоке finally в асинхронном генераторе или в обработчике, переданном в TrioToken.run_sync_soon() в нужное время.

Параметры:

  • async_fn – Асинхронная функция.

  • args – Позиционные аргументы для async_fn. Если нужно передать именованные аргументы, используйте functools.partial().

  • name – Имя задачи. Используется только для отладки/интроспекции (например, repr(task_obj)). Если это не строка, spawn_system_task() попытается преобразовать её в строку. Часто используется, если вы оборачиваете функцию перед запуском новой задачи, чтобы передать исходную функцию в качестве name=, что облегчит отладку.

  • context – Необязательный contextvars.Context объект с переменными контекста, которые необходимо использовать для этой задачи. Обычно вы копируете текущий контекст с помощью context = contextvars.copy_context(), а затем передаёте этот context объект.

Возвращаемое значение:

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

Тип возвращаемого значения:

Задача

Токены Trio

class trio.lowlevel.TrioToken

Непрозрачный объект, представляющий один вызов trio.run().

У него нет публичного конструктора; вместо этого см. current_trio_token().

Этот объект используется в двух случаях:

  1. Он позволяет повторно войти в цикл выполнения Trio из внешних потоков или обработчиков сигналов. Это низкоуровневый примитив, который использует trio.to_thread() и trio.from_thread для взаимодействия с рабочими потоками, trio.open_signal_receiver использует для получения уведомлений о сигналах и так далее.

  2. Каждый вызов trio.run() имеет ровно один связанный объект TrioToken, поэтому вы можете использовать его для идентификации конкретного вызова.

run_sync_soon(sync_fn: Callable[[Unpack[PosArgsT]], object], *args: Unpack[PosArgsT], idempotent: bool = False) → None

Запланировать вызов sync_fn(*args) в контексте задачи Trio.

Это безопасно вызывать из основного потока, из других потоков и из обработчиков сигналов. Это основной примитив для повторного входа в цикл выполнения Trio извне.

Вызов произойдет «скоро», но нет гарантии относительно точного момента, и нет механизма для определения момента его выполнения. Если вам это нужно, вы должны разработать свой собственный.

Вызов фактически выполняется как часть системной задачи (см. spawn_system_task()). В частности, это означает, что:

  • Защита от KeyboardInterrupt по умолчанию включена; если вы хотите, чтобы sync_fn можно было прервать нажатием Ctrl+C, необходимо явно использовать disable_ki_protection().

  • Если sync_fn вызывает исключение, то оно преобразуется в TrioInternalError, и все задачи отменяются. Следует позаботиться о том, чтобы sync_fn не аварийно завершался.

Все вызовы с idempotent=False обрабатываются в строгой очереди FIFO.

Если idempotent=True, то sync_fn и args должны быть хешируемыми, и Trio предпримет все возможные усилия для отбрасывания любого отправленного вызова, равного уже ожидающему вызову. Trio будет обрабатывать их в порядке очереди FIFO.

Любые гарантии порядка применяются отдельно к вызовам idempotent=False и idempotent=True; нет правил о том, как вызовы разных категорий упорядочиваются друг относительно друга.

Raises:

trio.RunFinishedError – если связанный вызов trio.run() уже завершился. (Любой вызов, не вызывающий это исключение, гарантированно будет полностью обработан до завершения trio.run().)

trio.lowlevel.current_trio_token() → TrioToken

Получить TrioToken для текущего вызова trio.run().

Запуск потоков

trio.lowlevel.start_thread_soon(fn: Callable[[], RetT], deliver: Callable[[outcome.Outcome[RetT]], object], name: str | None = None) → None

Выполняет deliver(outcome.capture(fn)) в рабочем потоке.

Как правило, fn выполняет некоторую блокирующую работу, а deliver возвращает результат тому, кто в этом заинтересован.

Это низкоуровневый интерфейс без излишеств, очень похожий на использование threading.Thread для прямого запуска потока. Основное отличие заключается в том, что эта функция пытается повторно использовать потоки при возможности, поэтому она может быть немного быстрее, чем threading.Thread.

Рабочие потоки имеют флаг daemon, который означает, что если основной поток завершится, рабочие потоки будут автоматически убиты. Если вы хотите убедиться, что ваш fn выполнится до конца, убедитесь, что основной поток остаётся активным до вызова deliver.

Безопасно вызывать эту функцию одновременно из нескольких потоков.

Parameters:

  • fn (функция синхронизации) – Выполняет произвольную блокирующую работу.

  • deliver (функция синхронизации) – Принимает outcome.Outcome fn и предоставляет его. Не должна блокировать.

Поскольку рабочие потоки кэшируются и повторно используются для нескольких вызовов, ни одна из функций не должна изменять состояние уровня потока, например, объекты threading.local – или, если они это делают, должны быть осторожны, чтобы восстановить свои изменения перед возвратом.

Примечание

Разделение между fn и deliver служит двум целям. Во-первых, это удобно, поскольку большинство вызывающих функций так или иначе нуждаются в чём-то подобном.

Во-вторых, это позволяет избежать небольшой проблемы гонки, которая может привести к созданию слишком большого количества потоков. Рассмотрим программу, которая хочет последовательно выполнять несколько задач в потоке, поэтому основной поток отправляет задачу, ждёт её завершения, отправляет другую задачу и так далее. Теоретически, этой программе понадобится только один рабочий поток. Но что может произойти:

  1. Рабочий поток: первая задача завершается и вызывает deliver.

  2. Основной поток: получает уведомление о завершении задачи и вызывает start_thread_soon.

  3. Основной поток: видит, что ни один рабочий поток не помечен как свободный, поэтому создаёт второй рабочий поток.

  4. Исходный рабочий поток: отмечает себя как свободный.

Чтобы этого избежать, потоки отмечают себя как свободные перед вызовом deliver.

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

Обработка прерывания KeyboardInterrupt с повышенной безопасностью

Обработка Trio прерывания Ctrl+C разработана для баланса удобства использования и безопасности. С одной стороны, существуют чувствительные области (например, основной цикл планирования), где просто невозможно обработать произвольные исключения KeyboardInterrupt, сохраняя при этом основные инварианты корректности. С другой стороны, если пользователь случайно напишет бесконечный цикл, мы хотим иметь возможность прервать его. Наше решение заключается в установке обработчика сигнала по умолчанию, который проверяет, безопасно ли вызвать исключение KeyboardInterrupt в месте получения сигнала. Если да, то мы это делаем; в противном случае, мы планируем доставку KeyboardInterrupt в главную задачу в ближайшую доступную возможность (аналогично тому, как доставляется Cancelled).

Итак, это замечательно, но – как мы можем узнать, находимся ли мы в одной из чувствительных частей программы или нет?

Это определяется на основе каждой функции. По умолчанию:

  • Функция верхнего уровня в обычных пользовательских задачах не защищена.

  • Функция верхнего уровня в системных задачах защищена.

  • Если функция не указывает иное, то она наследует состояние защиты своего вызывающего объекта.

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

Эти переходы выполняются с помощью двух декораторов функций:

@trio.lowlevel.disable_ki_protection

Декоратор, который помечает заданную обычную функцию, функцию-генератор, асинхронную функцию или асинхронную функцию-генератор как незащищенную от KeyboardInterrupt, то есть код внутри этой функции может быть грубо прерван KeyboardInterrupt в любой момент.

Если у вас несколько декораторов на одной функции, то этот декоратор должен находиться внизу стека (ближе к фактической функции).

Пример использования – в реализации чего-то вроде trio.from_thread.run(), которая использует TrioToken.run_sync_soon() для входа в поток Trio. Обработчики run_sync_soon() выполняются с включенной защитой KeyboardInterrupt, а trio.from_thread.run() использует это для безопасной настройки механизма отправки ответа обратно в исходный поток, но затем использует disable_ki_protection() при входе в функцию, предоставленную пользователем.

@trio.lowlevel.enable_ki_protection

Декоратор, который помечает заданную обычную функцию, функцию-генератор, асинхронную функцию или асинхронную функцию-генератор как защищенную от KeyboardInterrupt, то есть код внутри этой функции не будет грубо прерван KeyboardInterrupt. (Хотя, если он содержит какие-либо точки контроля, то он все еще может получить KeyboardInterrupt в этих точках. Это считается вежливым прерыванием.)

Предупреждение

Будьте очень осторожны, чтобы использовать этот декоратор только для функций, которые, как вы знаете, либо завершатся за ограниченное время, либо регулярно будут проходить через точку контроля. (Конечно, все ваши функции должны иметь это свойство, но если вы ошибетесь здесь, то даже не сможете использовать Ctrl+C для выхода!)

Если у вас несколько декораторов на одной функции, то этот декоратор должен находиться внизу стека (ближе к фактической функции).

Пример использования – на реализации __exit__ для чего-то вроде Lock, где плохое время прерывания KeyboardInterrupt могло оставить замок в несогласованном состоянии и привести к тупику.

Поскольку защита от KeyboardInterrupt отслеживается по объектам кода, любая попытка условно защитить один и тот же блок кода различными способами, скорее всего, не будет вести себя так, как вы ожидаете. Если вы попытаетесь условно защитить замыкание, оно будет защищено безусловно:

def example(protect: bool) -> bool:
    def inner() -> bool:
        return trio.lowlevel.currently_ki_protected()
    if protect:
        inner = trio.lowlevel.enable_ki_protection(inner)
    return inner()

async def amain():
    assert example(False) == False
    assert example(True) == True  # once protected ...
    assert example(False) == True  # ... always protected

trio.run(amain)

Если вам действительно нужна условная защита, вы можете получить ее, предоставив каждому защищенному от KI экземпляру замыкания свой собственный объект кода:

def example(protect: bool) -> bool:
    def inner() -> bool:
        return trio.lowlevel.currently_ki_protected()
    if protect:
        inner.__code__ = inner.__code__.replace()
        inner = trio.lowlevel.enable_ki_protection(inner)
    return inner()

async def amain():
    assert example(False) == False
    assert example(True) == True
    assert example(False) == False

trio.run(amain)

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

trio.lowlevel.currently_ki_protected() → bool

Проверить, включена ли защита от KeyboardInterrupt в вызываемом коде.

Удивительно легко думать, что защита от KeyboardInterrupt включена, когда она не включена, или наоборот. Эта функция сообщает вам, что думает об этом Trio, что делает ее полезной для assert и unit-тестов.

Возвращает:

True, если защита включена, и False в противном случае.

Тип возвращаемого значения:

bool

Ожидание и пробуждение

Абстракция очереди ожидания

class trio.lowlevel.ParkingLot

Справедливая очередь ожидания с возможностью отмены и повторного помещения в очередь.

Этот класс обобщает сложные части реализации очереди ожидания. Он полезен для реализации синхронизирующих примитивов более высокого уровня, таких как очереди и блокировки.

В дополнение к методам ниже, вы можете использовать len(parking_lot) для получения количества приостановленных задач и if parking_lot: ... для проверки наличия приостановленных задач.

: list[Task] broken_by

break_lot(task: Task | None = None) → None

Разбить эту парковку, с задачей task отмеченной как задача, которая это сделала.

Это приводит к тому, что все приостановленные задачи генерируют ошибку, а любые будущие попытки парковки также вызовут ошибку. Unpark и repark становятся пустыми операциями, так как парковка пуста.

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

await park() → None

Приостановить текущую задачу до тех пор, пока она не будет разбужена вызовом unpark() или unpark_all().

Исключения:

BrokenResourceError – если попытка парковки выполняется в сломанной парковке или парковка ломается, прежде чем мы дойдём до разбуживания.

repark(new_lot: ParkingLot, *, count: int | float = 1) → None

Переместить приостановленные задачи из одного объекта ParkingLot в другой.

Это удаляет из одной парковки count задачи и переупорядочивает их в другой, сохраняя порядок. Например:

async def parker(lot):
    print("sleeping")
    await lot.park()
    print("woken")

async def main():
    lot1 = trio.lowlevel.ParkingLot()
    lot2 = trio.lowlevel.ParkingLot()
    async with trio.open_nursery() as nursery:
        nursery.start_soon(parker, lot1)
        await trio.testing.wait_all_tasks_blocked()
        assert len(lot1) == 1
        assert len(lot2) == 0
        lot1.repark(lot2)
        assert len(lot1) == 0
        assert len(lot2) == 1
        # This wakes up the task that was originally parked in lot1
        lot2.unpark()

Если приостановленных задач меньше, чем count, тогда перемещаются доступные задачи и возвращается успешный результат.

Параметры:

  • new_lot (ParkingLot) – парковка, в которую нужно переместить задачи.

  • count (int|math.inf) – количество задач для перемещения.

repark_all(new_lot: ParkingLot) → None

Переместить все приостановленные задачи из одного объекта ParkingLot в другой.

См. repark() для деталей.

statistics() → ParkingLotStatistics

Возвращает объект, содержащий отладочную информацию.

В настоящее время определены следующие поля:

  • tasks_waiting: Количество задач, заблокированных в методе park() этой парковки.

unpark(*, count: int | float = 1) → list[Task]

Разбудить одну или несколько задач.

Это разбудит count задачи, заблокированные в park(). Если приостановленных задач меньше, чем count, то разбуживаются доступные задачи, и затем возвращается успешный результат.

Параметры:

count (int | math.inf) – количество задач для разбуживания.

unpark_all() → list[Task]

Разбудить все приостановленные задачи.

class trio.lowlevel.ParkingLotStatistics(tasks_waiting: int)

Объект, содержащий отладочную информацию для ParkingLot.

В настоящее время определены следующие поля:

  • tasks_waiting (int): Количество задач, заблокированных в методе trio.lowlevel.ParkingLot.park() этой парковки.

trio.lowlevel.add_parking_lot_breaker(task: Task, lot: ParkingLot) → None

Регистрирует задачу как разрушитель парковки. См. ParkingLot.break_lot().

Исключения:

trio.BrokenResourceError – если задача уже завершена.

trio.lowlevel.remove_parking_lot_breaker(task: Task, lot: ParkingLot) → None

Удаляет регистрацию задачи как разрушителя парковки. См. ParkingLot.break_lot()

Функции контрольных точек низкого уровня

await trio.lowlevel.checkpoint() → None

Чистая контрольная точка.

Проверяет отмену и позволяет планировать другие задачи без блокировки.

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

Эквивалентно await trio.sleep(0) (которое реализуется вызовом checkpoint().)

Следующие две функции используются вместе для создания контрольной точки:

await trio.lowlevel.checkpoint_if_cancelled() → None

Вызовите контрольную точку, если контекст вызывающего элемента был отменён.

Эквивалентно (но потенциально более эффективному):

if trio.current_effective_deadline() == -inf:
    await trio.lowlevel.checkpoint()

Это либо нет операция, либо она позволяет планировать другие задачи, а затем генерирует исключение trio.Cancelled.

Обычно используется вместе с cancel_shielded_checkpoint().

await trio.lowlevel.cancel_shielded_checkpoint() → None

Введите точку планирования, но не точку отмены.

Это не контрольная точка, но это половина контрольной точки, а в сочетании с checkpoint_if_cancelled() она может образовать полную контрольную точку.

Эквивалентно (но потенциально более эффективному):

with trio.CancelScope(shield=True):
    await trio.lowlevel.checkpoint()

Они часто используются в случаях, когда у вас есть операция, которая может или не может заблокировать выполнение, и вы хотите реализовать стандартную семантику контрольных точек Trio. Пример:

async def operation_that_maybe_blocks():
    await checkpoint_if_cancelled()
    try:
        ret = attempt_operation()
    except BlockingIOError:
        # need to block and then retry, which we do below
        pass
    else:
        # operation succeeded, finish the checkpoint then return
        await cancel_shielded_checkpoint()
        return ret
    while True:
        await wait_for_operation_to_be_ready()
        try:
            return attempt_operation()
        except BlockingIOError:
            pass

Эта логика немного запутанная, но выполняет всё следующее:

  • Каждый успешный путь выполнения проходит через контрольную точку (предполагая, что wait_for_operation_to_be_ready является безусловной контрольной точкой)

  • Наши семантика отмены говорят, что Cancelled должно быть возбуждено только в том случае, если операция не произошла. Использование cancel_shielded_checkpoint() на ветви выхода пораньше достигает этой цели.

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

  • Избегает неявного объединения BlockingIOError с любыми ошибками, генерируемыми attempt_operation или wait_for_operation_to_be_ready, сохраняя цикл while True: за пределами блока except BlockingIOError:.

Эти функции также могут быть полезны в других ситуациях. Например, когда trio.to_thread.run_sync() планирует какую-то работу для выполнения в потоке-работнике, она блокируется до завершения работы (поэтому это точка планирования), но по умолчанию она не допускает отмены. Таким образом, чтобы убедиться, что вызов всегда действует как контрольная точка, она вызывает checkpoint_if_cancelled() перед запуском потока.

Низкоуровневое блокирование

await trio.lowlevel.wait_task_rescheduled(abort_func: Callable[[Callable[[], NoReturn]], Abort]) → Any

Уложить текущую задачу в сон с поддержкой отмены.

Это низкоуровневый API для блокирования в Trio. Каждый раз, когда задача Task блокируется, она делает это, вызывая эту функцию (обычно косвенно через какой-либо API более высокого уровня).

Это сложный интерфейс без защитных ограждений. Если вы можете использовать ParkingLot или встроенные функции ожидания ввода-вывода, то вам следует это сделать.

В общем случае, перед вызовом этой функции, вы подготавливаете «кого-то», кто вызовет reschedule() для текущей задачи в какой-то момент позже.

Затем вы вызываете wait_task_rescheduled(), передавая abort_func, «обработчик отмены».

(Терминология: в Trio «отмена» — это процесс попытки прервать заблокированную задачу для передачи отмены.)

Есть два возможных сценария дальнейших действий:

  1. «Кто-то» вызывает reschedule() для текущей задачи, и wait_task_rescheduled() возвращает или генерирует любое значение или ошибку, которые были переданы в reschedule().

  2. Контекст вызова переходит в состояние отмены (например, из-за истечения таймаута). В этом случае вызывается abort_func. Его интерфейс выглядит так:

    def abort_func(raise_cancel):
        ...
        return trio.lowlevel.Abort.SUCCEEDED  # or FAILED

    Он должен попытаться очистить любые состояния, связанные с этим вызовом, и, в частности, организовать, чтобы reschedule() не вызывался позже. Если (и только если!) это удастся, он должен вернуть Abort.SUCCEEDED, в этом случае задача будет автоматически перепланирована с соответствующей ошибкой Cancelled.

    В противном случае, он должен вернуть Abort.FAILED. Это означает, что задача не может быть отменена в данный момент и всё ещё должна убедиться, что «кто-то» в конечном итоге вызовет reschedule().

    На этом этапе снова есть два варианта. Вы можете просто проигнорировать отмену: дождаться завершения операции, а затем перепланировать и продолжить как обычно. (Например, именно это делает trio.to_thread.run_sync(), если отмена отключена.) Другой вариант заключается в том, что abort_func успешно отменяет операцию, но по какой-то причине не может сообщить об этом сразу. (Пример: в Windows можно запросить отмену асинхронной («перекрытой») операции ввода-вывода, но этот запрос также асинхронный — вы узнаете об этом позже, отменилась операция или нет.) Для сообщения об отложенной отмене вам следует самостоятельно перепланировать задачу и вызвать raise_cancel обработчик, переданный в abort_func, для поднятия исключения Cancelled (или, возможно, KeyboardInterrupt) в эту задачу. Любой из описанных ниже подходов может подойти:

    # Option 1:
    # Catch the exception from raise_cancel and inject it into the task.
    # (This is what Trio does automatically for you if you return
    # Abort.SUCCEEDED.)
    trio.lowlevel.reschedule(task, outcome.capture(raise_cancel))
    
    # Option 2:
    # wait to be woken by "someone", and then decide whether to raise
    # the error from inside the task.
    outer_raise_cancel = None
    def abort(inner_raise_cancel):
        nonlocal outer_raise_cancel
        outer_raise_cancel = inner_raise_cancel
        TRY_TO_CANCEL_OPERATION()
        return trio.lowlevel.Abort.FAILED
    await wait_task_rescheduled(abort)
    if OPERATION_WAS_SUCCESSFULLY_CANCELLED:
        # raises the error
        outer_raise_cancel()

    В любом случае гарантируется, что мы вызываем abort_func не более одного раза на каждый вызов wait_task_rescheduled().

Иногда полезно иметь возможность обмениваться некоторыми данными, связанными со сном, между засыпающей задачей, функцией отмены и возобновляющей задачей. Вы можете использовать атрибут custom_sleep_data засыпающей задачи для хранения этих данных, и Trio не будет их трогать, кроме как убедиться, что они очищаются при перепланировке задачи.

Предупреждение

Если ваш abort_func генерирует ошибку или возвращает любое значение, отличное от Abort.SUCCEEDED или Abort.FAILED, Trio потерпит полную неудачу. Будьте внимательны! Аналогично, вполне возможно заблокировать программу Trio, не перепланировав заблокированную задачу, или нанести ущерб, вызвав reschedule() слишком много раз. Помните, что мы говорили выше о том, что следует использовать API более высокого уровня, если это возможно?

class trio.lowlevel.Abort(value, names=None, *, module=None, qualname=None, type=None, start=1, boundary=None)

enum.Enum используется в качестве возвращаемого значения от функций отмены.

См. wait_task_rescheduled() для получения подробной информации.

SUCCEEDED

FAILED

trio.lowlevel.reschedule(task: Task, next_send: Outcome[object] = ) → None

Перепланировать задачу с заданным outcome.Outcome.

См. wait_task_rescheduled() для подробных сведений.

Должен быть ровно один вызов reschedule() для каждого вызова wait_task_rescheduled(). (И при подсчете имейте в виду, что возвращение Abort.SUCCEEDED из обработчика отмены эквивалентно вызову reschedule() один раз.)

Параметры:

  • task (trio.lowlevel.Task) – задача, подлежащая перепланировке. Должна быть заблокирована в вызове wait_task_rescheduled().

  • next_send (outcome.Outcome) – значение (или ошибка), которое должно быть возвращено (или возбуждено) из wait_task_rescheduled().

Вот пример класса блокировки, реализованного с помощью wait_task_rescheduled() напрямую. Эта реализация имеет ряд недостатков, включая отсутствие справедливости, отмену O(n), отсутствие проверки ошибок, отсутствие вставки контрольной точки на пути без блокировок и т. д. Если вы действительно хотите реализовать свою собственную блокировку, изучите реализацию trio.Lock и используйте ParkingLot, который обрабатывает некоторые из этих проблем за вас. Но это служит иллюстрацией основного строения API wait_task_rescheduled():

class NotVeryGoodLock:
    def __init__(self):
        self._blocked_tasks = collections.deque()
        self._held = False

    async def acquire(self):
        # We might have to try several times to acquire the lock.
        while self._held:
            # Someone else has the lock, so we have to wait.
            task = trio.lowlevel.current_task()
            self._blocked_tasks.append(task)
            def abort_fn(_):
                self._blocked_tasks.remove(task)
                return trio.lowlevel.Abort.SUCCEEDED
            await trio.lowlevel.wait_task_rescheduled(abort_fn)
            # At this point the lock was released -- but someone else
            # might have swooped in and taken it again before we
            # woke up. So we loop around to check the 'while' condition
            # again.
        # if we reach this point, it means that the 'while' condition
        # has just failed, so we know no-one is holding the lock, and
        # we can take it.
        self._held = True

    def release(self):
        self._held = False
        if self._blocked_tasks:
            woken_task = self._blocked_tasks.popleft()
            trio.lowlevel.reschedule(woken_task)

API задач

trio.lowlevel.current_root_task()

Возвращает текущую корневую Task.

Это задача, которая является конечным родителем всех других задач.

trio.lowlevel.current_task()

Возвращает объект Task, представляющий текущую задачу.

Возвращает:

объект Task, который вызвал current_task().

Тип возвращаемого значения:

Task

class trio.lowlevel.Task

Объект Task представляет собой конкурирующую «поток» выполнения. У него нет публичного конструктора; Trio внутренне создает объект Task для каждого вызова nursery.start(...) или nursery.start_soon(...).

Его публичные члены в основном полезны для интроспекции и отладки:

name

Строка, содержащая имя этой задачи Task. Обычно это имя функции, в которой выполняется эта задача Task, но может быть переопределено путем передачи name= в start или start_soon.

coro

Объект корутины этой задачи.

for ... in iter_await_frames() → Iterator[tuple[types.FrameType, int]]

Рекурсивно итерируется по объектам корутин, на которых ожидает эта задача, и возвращает кадр и номер строки в каждом кадре.

Это аналогично traceback.walk_stack в синхронном контексте. Обратите внимание, что traceback.walk_stack возвращает кадры снизу стека вызовов наверх, в то время как эта функция начинается с Task.coro и работает вниз.

Пример использования: извлечение трассировки стека:

import traceback

def print_stack_for_task(task):
    ss = traceback.StackSummary.extract(task.iter_await_frames())
    print("".join(ss.format()))

context

Объект contextvars.Context этой задачи.

parent_nursery

Питомник, в котором находится эта задача (или None, если это задача «init»).

Пример использования: построение визуализации дерева задач в отладчике.

eventual_parent_nursery

Питомник, в котором эта задача будет находиться после вызова task_status.started().

Если эта задача уже вызвала started(), или если она не была запущена с помощью nursery.start(), то ее eventual_parent_nursery является None.

child_nurseries

Питомники, содержащиеся в этой задаче.

Это список, в котором внешние питомники находятся перед внутренними.

custom_sleep_data

Trio не присваивает этому переменной никакого значения, за исключением того, что устанавливает ее в None всякий раз, когда задача перепланируется. Его можно использовать для обмена данными между различными задачами, участвующими в приостановке задачи и ее повторном запуске. (См. wait_task_rescheduled() для получения подробностей.)

Использование «гостевого режима» для запуска Trio поверх других циклов событий

Что такое «гостевой режим»?

Цикл событий действует как центральный координатор для управления всеми операциями ввода-вывода в вашей программе. Обычно это означает, что ваше приложение должно выбрать один цикл событий и использовать его для всего. Но что, если вам нравится Trio, но вам также нужно использовать фреймворк, например, Qt или PyGame, у которого есть свой собственный цикл событий? Тогда вам нужен способ запустить оба цикла событий одновременно.

Объединение циклов событий возможно, но стандартные подходы имеют существенные недостатки:

  • Опрос: в этом случае вы используете цикл ожидания для ручного проверки ввода-вывода в обоих циклах событий много раз в секунду. Это добавляет задержку и тратит время процессора и электроэнергию.

  • Подключаемые бэкенды ввода-вывода: в этом случае вы повторно реализуете один из API циклов событий поверх другого, поэтому в итоге получаете только один цикл событий. Это требует значительной работы для каждой пары циклов событий, которые вы хотите интегрировать, и различные бэкенды неизбежно ведут к несовместимому поведению, вынуждая пользователей программировать с учетом наименьшего общего знаменателя. И если два цикла событий предоставляют разные наборы функций, может быть даже невозможно реализовать один через другой.

  • Запуск двух циклов событий в отдельных потоках: Это работает, но большинство API циклов событий не являются потокобезопасными, поэтому в этом подходе вам нужно тщательно отслеживать, какой код выполняется в каком цикле событий, и помнить об использовании явного межпоточного обмена сообщениями всякий раз, когда вы взаимодействуете с другим циклом — в противном случае вы рискуете возникновением трудноотслеживаемых гонок и повреждения данных.

Вот почему Trio предлагает четвертый вариант: гостевой режим. Гостевой режим позволяет выполнять trio.run поверх какого-либо другого «хостового» цикла событий, например Qt. Его преимущества:

  • Эффективность: гостевой режим ориентирован на события, вместо использования цикла ожидания, поэтому он имеет низкую задержку и не тратит электроэнергию.

  • Нет необходимости думать о потоках: ваш код Trio выполняется в том же потоке, что и хостовый цикл событий, поэтому вы можете свободно вызывать синхронные API Trio из хоста и вызывать синхронные API хоста из Trio. Например, если вы создаете приложение GUI с Qt в качестве хостового цикла, то создание кнопки отмены и подключение ее к trio.CancelScope так же просто, как написание:

    # Trio code can create Qt objects without any special ceremony...
    my_cancel_button = QPushButton("Cancel")
    # ...and Qt can call back to Trio just as easily
    my_cancel_button.clicked.connect(my_cancel_scope.cancel)

    (Для асинхронных API это не так просто, но вы можете использовать синхронные API для создания явных мостов между двумя мирами, например, передавая асинхронные функции и их результаты друг другу через очереди.)

  • Согласованное поведение: гостевой режим использует тот же код, что и обычный Trio: тот же планировщик, тот же код ввода-вывода, то же самое во всем. Таким образом, вы получаете полный набор функций, и все работает так, как ожидается.

  • Простая интеграция и широкая совместимость: практически каждый цикл событий предлагает какую-то потокобезопасную операцию «планирования обратного вызова», и этого достаточно, чтобы использовать его как хостовый цикл.

Действительно? Как это возможно?

Примечание

Вы можете использовать гостевой режим, не читая этот раздел. Он включен для тех, кто любит понимать, как работают вещи.

Все циклы событий имеют одинаковую базовую структуру. Они циклически выполняют две операции:

  1. Ожидание уведомления операционной системы о том, что произошло что-то интересное, например, данные, прибывающие на сокет, или истекло время ожидания. Это делается путем вызова специфичного для платформы системного вызова sleep_until_something_happens() — select, epoll, kqueue, GetQueuedCompletionEvents и т.д.

  2. Выполнение всех пользовательских задач, относящихся к тому, что произошло, а затем возвращение к шагу 1.

Проблема здесь в шаге 1. Два разных цикла событий в одном потоке могут по очереди выполнять пользовательские задачи в шаге 2, но когда они простаивают и ничего не происходит, они не могут оба вызвать свою функцию sleep_until_something_happens() одновременно.

Стратегии «опрос» и «подключаемые бэкенды» решают эту проблему, модифицируя циклы, чтобы оба шага 1 могли выполняться одновременно в одном потоке. Поддержание всего в одном потоке отлично для шага 2, но правки шага 1 создают проблемы.

Стратегия «отдельные потоки» решает эту проблему, перемещая оба шага в отдельные потоки. Это позволяет выполнить шаг 1, но недостатком является то, что теперь пользовательские задачи в шаге 2 также выполняются в отдельных потоках, поэтому пользователи вынуждены заниматься межпоточной координацией.

Идея гостевого режима заключается в объединении лучших частей каждого подхода: мы перемещаем шаг 1 Trio в отдельный рабочий поток, сохраняя шаг 2 Trio в основном хостовом потоке. Таким образом, когда приложение простаивает, оба цикла событий выполняют свои sleep_until_something_happens() одновременно в своих потоках. Но когда приложение оживает и ваш код фактически выполняется, всё происходит в одном потоке. Сложности с потоками обрабатываются прозрачно внутри Trio.

Конкретно, мы развертываем внутренний цикл событий Trio в цепочку обратных вызовов, и по окончании каждого обратного вызова мы планируем следующий обратный вызов в хостовый цикл или в рабочий поток, в зависимости от ситуации. Поэтому хостовый цикл должен только предоставить способ планирования обратного вызова в основной поток из рабочего потока.

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

Однако есть полезная оптимизация: нам необходим поток только тогда, когда наш sleep_until_something_happens() вызов фактически приостанавливается, то есть когда часть Trio вашей программы простаивает и не имеет ничего общего. Таким образом, прежде чем переключиться в рабочий поток, мы дважды проверяем, простаивает ли он, и если нет, то пропускаем рабочий поток и переходим непосредственно к шагу 2. Это означает, что ваше приложение платит дополнительную цену за переключение потоков только в моменты, когда в противном случае оно приостанавливалось, поэтому оно должно минимально влиять на общую производительность вашего приложения.

Общий накладные расходы будут зависеть от вашего хостового цикла, вашей платформы, вашего приложения и т. д. Но мы ожидаем, что в большинстве случаев приложения, работающие в гостевом режиме, будут только на 5-10% медленнее, чем тот же код, использующий trio.run. Если вы обнаружите, что это не так для вашего приложения, сообщите нам, и мы посмотрим, сможем ли мы это исправить!

Реализация гостевого режима для вашего любимого цикла событий

Давайте пройдемся по тому, что вам нужно сделать, чтобы интегрировать гостевой режим Trio с вашим любимым циклом событий. Рассматривайте этот раздел как список задач.

Начало работы: Первый шаг — заставить что-то базовое работать. Вот минимальный пример запуска Trio поверх asyncio, который вы можете использовать в качестве модели:

import asyncio
import trio

# A tiny Trio program
async def trio_main():
    for _ in range(5):
        print("Hello from Trio!")
        # This is inside Trio, so we have to use Trio APIs
        await trio.sleep(1)
    return "trio done!"

# The code to run it as a guest inside asyncio
async def asyncio_main():
    asyncio_loop = asyncio.get_running_loop()

    def run_sync_soon_threadsafe(fn):
        asyncio_loop.call_soon_threadsafe(fn)

    def done_callback(trio_main_outcome):
        print(f"Trio program ended with: {trio_main_outcome}")

    # This is where the magic happens:
    trio.lowlevel.start_guest_run(
        trio_main,
        run_sync_soon_threadsafe=run_sync_soon_threadsafe,
        done_callback=done_callback,
    )

    # Let the host loop run for a while to give trio_main time to
    # finish. (WARNING: This is a hack. See below for better
    # approaches.)
    #
    # This function is in asyncio, so we have to use asyncio APIs.
    await asyncio.sleep(10)

asyncio.run(asyncio_main())

Вы можете видеть, что мы используем API-интерфейсы, специфичные для asyncio, для запуска цикла, а затем вызываем trio.lowlevel.start_guest_run. Эта функция очень похожа на trio.run и принимает все те же аргументы. Но у нее есть два отличия:

Во-первых, вместо того, чтобы блокироваться до тех пор, пока trio_main не завершится, он планирует trio_main для запуска поверх цикла хоста и сразу же возвращается. Таким образом, trio_main работает в фоновом режиме — поэтому нам нужно спать и дать ему время завершиться.

И во-вторых, она требует двух дополнительных ключевых аргументов: run_sync_soon_threadsafe и done_callback.

Для run_sync_soon_threadsafe нам нужна функция, которая принимает синхронный обратный вызов и планирует его для выполнения в вашем цикле хоста. И эта функция должна быть «потокобезопасной» в том смысле, что вы можете безопасно вызывать её из любого потока. Поэтому вам нужно разобраться, как написать функцию, которая делает это, используя API вашего цикла хоста. В asyncio это легко, так как call_soon_threadsafe делает именно то, что нам нужно; для вашего цикла это может быть более или менее сложно.

Для done_callback вы передаете функцию, которую Trio автоматически вызовет при завершении выполнения Trio, так что вы знаете, что это завершено и что произошло. В этой базовой стартовой версии мы просто выводим результат; в следующем разделе мы обсудим лучшие альтернативы.

На этом этапе вы должны быть в состоянии запустить простую программу Trio внутри вашего цикла хоста. Теперь мы превратим этот прототип во что-то надежное.

Жизненный цикл циклов событий: Одна из самых сложных вещей в большинстве циклов событий — правильное завершение. А наличие двух циклов событий усложняет это еще больше!

Если возможно, мы рекомендуем следовать этой схеме:

  • Запустите свой цикл хоста

  • Немедленно вызовите start_guest_run для запуска Trio

  • Когда Trio завершит работу и вызов done_callback, завершите цикл хоста

  • Убедитесь, что ничего другого не завершает ваш цикл хоста

Таким образом, ваши два цикла событий будут иметь одинаковый жизненный цикл, и ваша программа автоматически завершится, когда ваша функция Trio завершится.

Вот как мы бы расширили наш пример asyncio, чтобы реализовать эту схему:

# Improved version, that shuts down properly after Trio finishes
async def asyncio_main():
    asyncio_loop = asyncio.get_running_loop()

    def run_sync_soon_threadsafe(fn):
        asyncio_loop.call_soon_threadsafe(fn)

    # Revised 'done' callback: set a Future
    done_fut = asyncio_loop.create_future()
    def done_callback(trio_main_outcome):
        done_fut.set_result(trio_main_outcome)

    trio.lowlevel.start_guest_run(
        trio_main,
        run_sync_soon_threadsafe=run_sync_soon_threadsafe,
        done_callback=done_callback,
    )

    # Wait for the guest run to finish
    trio_main_outcome = await done_fut
    # Pass through the return value or exception from the guest run
    return trio_main_outcome.unwrap()

А затем вы можете инкапсулировать всю эту механику в утилитарную функцию, которая предоставляет API-интерфейс, похожий на trio.run, но запускает оба цикла вместе:

def trio_run_with_asyncio(trio_main, *args, **trio_run_kwargs):
    async def asyncio_main():
        # same as above
        ...

    return asyncio.run(asyncio_main())

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

  • Вы должны позволить программе Trio завершиться. Многие циклы событий позволяют остановить цикл событий в любой момент, и все ожидающие обратные вызовы/задачи и т.д. просто… не выполняются. Trio использует более структурированную систему, где вы можете отменять вещи, но код всегда выполняется до конца, поэтому finally блоки выполняются, ресурсы очищаются и т.д. Если вы остановите свой цикл хоста слишком рано, прежде чем вызов done_callback не будет вызван, это прервет выполнение Trio посредине без возможности очистки. Это может оставить ваш код в несогласованном состоянии и определенно оставит внутренние компоненты Trio в несогласованном состоянии, что вызовет ошибки, если вы снова попытаетесь использовать Trio в этом потоке.

    Некоторым программам нужно иметь возможность завершаться в любой момент, например, в ответ на закрытие окна графического интерфейса или выбор пользователем пункта «Выход» в меню. В таких случаях мы рекомендуем обернуть всю вашу программу в trio.CancelScope и отменить её, когда вы захотите выйти.

  • Каждый цикл хоста может иметь только один start_guest_run одновременно. Если вы попытаетесь запустить второй, вы получите ошибку. Если вам нужно запустить несколько функций Trio одновременно, запустите один запуск Trio, откройте ясли и затем запустите свои функции как дочерние задачи в этих яслях.

  • Если вы или ваш цикл хоста не зарегистрировали обработчик для signal.SIGINT до запуска Trio (это не обычно), Trio возьмет на себя доставку KeyboardInterrupt. И поскольку Trio не может определить, какой код хоста безопасно прервать, он будет только передавать KeyboardInterrupt в часть кода Trio. Это нормально, если ваша программа настроена на завершение, когда часть Trio завершается, потому что KeyboardInterrupt будет распространяться из Trio и затем вызовет завершение вашего цикла хоста, что именно вам и нужно.

Учитывая эти ограничения, мы считаем, что самый простой подход — всегда запускать и останавливать оба цикла вместе.

Управление сигналами: «Сигналы» — это низкоуровневая примитивная межпроцессная коммуникация. Когда вы нажимаете Ctrl+C для завершения программы, используется сигнал. Обработка сигналов в Python имеет много подвижных элементов. Одним из этих элементов является signal.set_wakeup_fd, который циклы событий используют для того, чтобы гарантировать, что они проснутся, когда придёт сигнал, чтобы на него отреагировать. (Если у вас когда-либо был цикл событий, который игнорировал вас, когда вы нажимали Ctrl+C, это, вероятно, было потому, что они не использовали signal.set_wakeup_fd правильно.)

Однако, только один цикл событий может использовать signal.set_wakeup_fd одновременно. В гостевом режиме это может привести к проблемам: Trio и цикл хоста могут начать борьбу за использование signal.set_wakeup_fd.

Некоторые циклы событий, например asyncio, не будут работать правильно, если они не выиграют эту борьбу. К счастью, Trio немного менее избирателен: пока кто-то гарантирует, что программа просыпается, когда приходит сигнал, она должна работать правильно. Поэтому, если ваш цикл хоста хочет использовать signal.set_wakeup_fd, отключите поддержку signal.set_wakeup_fd в Trio, и оба цикла будут работать правильно.

С другой стороны, если ваш цикл хоста не использует signal.set_wakeup_fd, то единственный способ сделать всё правильно — включить поддержку signal.set_wakeup_fd Trio.

По умолчанию Trio предполагает, что ваш цикл хоста не использует signal.set_wakeup_fd. Он пытается обнаружить, когда это приводит к конфликту с циклом хоста, и выводит предупреждение — но, к сожалению, к тому времени, как он его обнаруживает, ущерб уже нанесён. Поэтому, если вы получаете это предупреждение, вам следует отключить поддержку signal.set_wakeup_fd Trio, передав host_uses_signal_set_wakeup_fd=True в start_guest_run.

Если вы не видите никаких предупреждений с вашим начальным прототипом, вы, вероятно, в порядке. Но единственный способ быть уверенным — проверить исходный код вашего цикла хоста. Например, asyncio может или не может использовать signal.set_wakeup_fd в зависимости от версии Python и операционной системы.

Небольшая оптимизация: Наконец, рассмотрите небольшую оптимизацию. Некоторые циклы событий предлагают две версии своего API «вызвать эту функцию вскоре»: одну, которую можно использовать из любого потока, и одну, которая может быть использована только из потока цикла событий, причём последняя является более дешевой. Например, asyncio имеет как call_soon_threadsafe, так и call_soon.

Если у вас есть такой цикл, вы также можете передать ключевой аргумент run_sync_soon_not_threadsafe=... в start_guest_run, и Trio автоматически будет его использовать, когда это необходимо.

Если у вашего цикла нет такого разделения, не беспокойтесь об этом; run_sync_soon_not_threadsafe= является необязательным. (Если он не передан, Trio будет просто использовать вашу потокобезопасную версию во всех случаях.)

И всё! Если вы выполнили все эти шаги, вы должны получить чистую интегрированную гибридную схему цикла событий. Идите создавать крутые GUI/игры/что угодно!

Ограничения

В целом, практически все возможности Trio должны работать в режиме гостя. Исключение составляют функции, которые полагаются на то, что Trio имеет полное представление о том, что делает ваша программа, поскольку очевидно, что оно не может контролировать цикл хоста или видеть, что он делает.

Пользовательские таймеры могут использоваться в режиме гостя, но они влияют только на таймауты Trio, а не на таймауты цикла хоста. А таймер мгновенного перехода и связанные с ним trio.testing.wait_all_tasks_blocked технически могут использоваться в режиме гостя, но они будут учитывать только задачи Trio при решении, нужно ли перепрыгивать через таймер или все задачи заблокированы.

Справочник

trio.lowlevel.start_guest_run(async_fn: Callable[..., Awaitable[RetT]], *args: object, run_sync_soon_threadsafe: Callable[[Callable[[], object]], object], done_callback: Callable[[outcome.Outcome[RetT]], object], run_sync_soon_not_threadsafe: Callable[[Callable[[], object]], object] | None = None, host_uses_signal_set_wakeup_fd: bool = False, clock: Clock | None = None, instruments: Sequence[Instrument] = (), restrict_keyboard_interrupt_to_checkpoints: bool = False, strict_exception_groups: bool = True) → None

Запуск «гостевого» выполнения Trio поверх другого цикла событий «хоста».

Каждый цикл хоста может иметь только одно гостевое выполнение одновременно.

Вы всегда должны дождаться завершения выполнения Trio, прежде чем останавливать цикл хоста; в противном случае он может оставить внутренние структуры данных Trio в несогласованном состоянии. Возможно, вы сможете обойтись без этого, если сразу выйдете из программы, но безопаснее всего этого не делать.

В общем случае лучший способ сделать это — обернуть этот код в функцию, которая запускает цикл хоста, а затем сразу же запускает гостевое выполнение, а затем останавливает хост после завершения гостевого выполнения.

После успешного возвращения start_guest_run() гостевое выполнение будет достаточно настроенно для вызова функций Trio с синхронизацией, таких как current_time(), spawn_system_task() и current_trio_token(). Если во время этой ранней настройки гостевого выполнения произойдет TrioInternalError, оно будет поднято из start_guest_run(). Все остальные ошибки, включая все ошибки, поднятые функцией async_fn, будут доставлены в ваш done_callback в какой-то момент после успешного возвращения из start_guest_run().

Параметры:

  • run_sync_soon_threadsafe –

    Произвольное вызываемое значение, которому будет передана функция в качестве единственного аргумента:

    def my_run_sync_soon_threadsafe(fn):
        ...

    Это вызываемое значение должно запланировать fn() для выполнения хостом в следующий проход по циклу. Должен поддерживать вызов из произвольных потоков.

  • done_callback –

    Произвольное вызываемое значение:

    def my_done_callback(run_outcome):
        ...

    Когда выполнение Trio завершится, Trio вызовет этот обратный вызов, чтобы сообщить вам об этом. Аргументом является outcome.Outcome, сообщающий о том, что должно было быть возвращено или поднято trio.run. Эта функция может делать что угодно, но обычно вы захотите остановить цикл хоста, разразобраться с результатом и т. д.

  • run_sync_soon_not_threadsafe – Как run_sync_soon_threadsafe, но будет вызываться только из главного потока цикла хоста. Необязательно, но если ваш цикл хоста позволяет реализовать это более эффективно, чем run_sync_soon_threadsafe, то передача этого значения сделает вещи немного быстрее.

  • host_uses_signal_set_wakeup_fd (bool) – Передайте True, если ваш цикл хоста использует signal.set_wakeup_fd, и False в противном случае. Более подробную информацию см. в разделе Реализация режима гостя для вашего любимого цикла событий.

Для значения других аргументов см. trio.run.

Передача живых объектов корутин между исполнителями корутин

Внутри синтаксис async/await в Python построен на основе концепции «объектов корутин» и «исполнителей корутин». Объект корутины представляет состояние стека асинхронного вызова. Но сам по себе это просто статический объект, который просто находится там. Если вы хотите, чтобы он что-то делал, вам нужен исполнитель корутины, чтобы продвигать его дальше. Каждая задача Trio имеет связанный с ней объект корутины (см. Task.coro), и планировщик Trio действует как их исполнитель корутины.

Но, конечно, Trio не единственный исполнитель корутин в Python – asyncio имеет своего, другие циклы событий тоже, вы даже можете определить свой собственный.

И в некоторых очень, очень необычных обстоятельствах даже имеет смысл передавать один объект корутины туда и обратно между разными исполнителями корутин. Об этом и говорится в данном разделе. Это чрезвычайно экзотический случай использования и предполагает глубокое понимание того, как Python async/await работает внутри. Для примеров мотивации см. вопрос #42 на GitHub для trio-asyncio и вопрос #649 на GitHub для trio. Для получения более подробной информации о работе корутин мы рекомендуем «Рассказ о циклах событий» Андре Карона или обратиться непосредственно к PEP 492 для получения всех подробностей.

await trio.lowlevel.permanently_detach_coroutine_object(final_outcome: Outcome[object]) → object

Постоянно отсоединить текущую задачу от планировщика Trio.

Обычно задача Trio не завершается, пока не завершится её объект корутины. Когда вы вызываете эту функцию, Trio ведет себя так, как будто объект корутины только что завершился, и задача завершается с заданным результатом. Это полезно, если вы хотите постоянно переключить объект корутины на другой исполнитель корутин.

Когда вызываемая корутина входит в эту функцию, она выполняется в рамках Trio, а когда функция возвращает значение, она выполняется в рамках внешнего исполнителя корутин.

Вы должны убедиться, что объект корутины освободил все ресурсы, специфичные для Trio, которые он получил (например, nurseries).

Параметры:

final_outcome (outcome.Outcome) – Trio ведет себя так, как будто текущая задача завершилась с указанным возвращаемым значением или исключением.

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

await trio.lowlevel.temporarily_detach_coroutine_object(abort_func: Callable[[Callable[[], NoReturn]], Abort]) → object

Временно отсоединить текущий объект корутины от планировщика Trio.

Когда вызываемая корутина входит в эту функцию, она выполняется в рамках Trio, а когда функция возвращает значение, она выполняется в рамках внешнего исполнителя корутин.

Задача Trio Task будет продолжать существовать, но будет приостановлена, пока вы не используете reattach_detached_coroutine_object() для её возобновления. Тем временем вы можете использовать другой исполнитель корутин для планирования объекта корутины. Фактически, вы должны – функция не возвращает значение, пока корутина не будет продвинута извне.

Обратите внимание, что вам нужно сохранить текущий объект Task для его последующего возобновления; вы можете получить его с помощью current_task(). Вы также можете использовать этот объект Task для получения объекта корутины — см. Task.coro.

Параметры:

abort_func – То же, что и для wait_task_rescheduled(), за исключением того, что он должен возвращать Abort.FAILED. (Если он возвращает Abort.SUCCEEDED, то Trio попытается перепланировать отсоединённую задачу напрямую, минуя reattach_detached_coroutine_object(), что было бы плохо.) Ваш abort_func должен по-прежнему организовать отмену выполнения объекта корутины, а затем повторно подключиться к Trio и вызвать обратный вызов raise_cancel, если это возможно.

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

await trio.lowlevel.reattach_detached_coroutine_object(task: Task, yield_value: object) → None

Подключение объекта корутины, который был отсоединён с помощью temporarily_detach_coroutine_object().

Когда вызываемая корутина входит в эту функцию, она выполняется в рамках внешнего исполнителя корутин, а когда функция возвращает значение, она выполняется в рамках Trio.

Это необходимо вызвать внутри возобновляемой корутины, и она возвращает то значение, которое вы передали. (Предполагается, что вы передадите значение, которое заставит текущий исполнитель корутин прекратить планирование этой задачи.) Затем корутина возобновляется планировщиком Trio в ближайшее время.

Параметры:

  • task (Task) – Объект задачи Trio, от которого текущая корутина была отсоединена.

  • yield_value (object) – Объект, который нужно передать текущему исполнителю корутин.

© 2017 Nathaniel J. Smith
Licensed under the MIT License.
https://trio.readthedocs.io/en/v0.29.0/reference-lowlevel.html

Spec-Zone.ru

Настройки Оффлайн Что нового Помощь О нас
Spec-Zone .ru
спецификации, руководства, описания, API