Просмотр и расширение 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
-
Объект, содержащий отладочную информацию на уровне цикла обработки событий.
В настоящее время определены следующие поля:
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, который является строкой, называющей используемый бэкенд ввода-вывода операционной системы; другие атрибуты варьируются в зависимости от бэкенда.
class trio.lowlevel.RunStatistics
Проверка наличия 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 | ||
в потоке без активного вызова | ||
в хост-цикле гостевого запуска | ||
внутри вызова инструмента | зависит | |
в потоке, созданном | ||
внутри функции прерывания |
-
Проверяет, находимся ли мы в запуске Trio. Возвращает
Trueтолько в том случае, еслиcurrent_time()будет успешным.См. также обсуждение различных способов обнаружения Trio.
trio.lowlevel.in_trio_run() → bool
-
Проверяет, находимся ли мы в задаче Trio. Возвращает
Trueтолько в том случае, еслиcurrent_task()будет успешным.См. также обсуждение различных способов обнаружения Trio.
trio.lowlevel.in_trio_task() → bool
Текущий таймер
-
Возвращает текущий
Clock.
trio.lowlevel.current_clock() → Clock
API инструментов
API инструментов предоставляет стандартный способ добавления пользовательской инструментики в цикл выполнения. Хотите построить гистограмму задержек планирования, записать трассировку стека любого задания, которое блокирует цикл выполнения более чем на 50 мс, или измерить процент времени работы процесса, потраченного на ожидание ввода-вывода? Это то место.
Общая идея заключается в том, что в любой момент времени trio.run() поддерживает набор «инструментов», которые представляют собой объекты, реализующие интерфейс trio.abc.Instrument. Когда происходит интересное событие, он перебирает эти инструменты и уведомляет их, вызывая соответствующий метод. Учебник содержит простой пример использования этого для отслеживания.
Поскольку это подключение к Trio на довольно низком уровне, нужно быть осторожным. Обратные вызовы выполняются синхронно, и во многих случаях, если они завершатся ошибкой, нет никакого разумного способа распространить эту ошибку (например, мы можем быть глубоко в механизме обработки исключений...). Поэтому наша текущая стратегия обработки исключений, поднятых инструментами, заключается в (а) записи исключения в журнал "trio.abc.Instrument", который по умолчанию выводит трассировку стека в стандартный поток ошибок, и (б) отключении нарушающего инструмент.
Вы можете зарегистрировать начальный список инструментов, передав их в trio.run(). add_instrument() и remove_instrument() позволяют добавлять и удалять инструменты во время выполнения.
-
Начать инструментацию текущего цикла выполнения заданным инструментом.
-
instrument (trio.abc.Instrument) – Инструмент, который нужно активировать.
Параметры:
Если
instrumentуже активен, ничего не делает. -
trio.lowlevel.add_instrument(instrument: Instrument) → None
-
Остановка инструментации текущего цикла выполнения заданным инструментом.
-
instrument (trio.abc.Instrument) – Инструмент, который нужно деактивировать.
-
KeyError – если инструмент не активен. Это может произойти либо потому, что вы его никогда не добавляли, либо потому, что вы его добавили, а затем он поднял необработанное исключение и был автоматически деактивирован.
Параметры:
Возбуждает:
-
trio.lowlevel.remove_instrument(instrument: Instrument) → None
И вот интерфейс для реализации, если вы хотите создать свой собственный Instrument:
-
Интерфейс для инструментации цикла выполнения.
Инструменты не обязаны наследоваться от этого абстрактного базового класса, и все эти методы являются необязательными. Этот класс служит в основном для документации.
-
Вызывается после обработки ожидающего ввода-вывода.
-
timeout (float) – Количество секунд, на которое мы были готовы ждать. Это время может или не может истечь, в зависимости от того, был ли готов какой-либо ввод-вывод.
Параметры:
-
after_io_wait(timeout: float) → None-
Вызывается непосредственно перед возвращением
trio.run().
after_run() → None-
Вызывается при возвращении в основной цикл выполнения после того, как задание сгенерировало выход.
-
task (trio.lowlevel.Task) – Задание, которое только что выполнилось.
Параметры:
-
after_task_step(task: Task) → None-
Вызывается перед блокированием для ожидания готовности ввода-вывода.
-
timeout (float) – Количество секунд, на которое мы готовы ждать.
Параметры:
-
before_io_wait(timeout: float) → None-
Вызывается в начале
trio.run().
before_run() → None-
Вызывается непосредственно перед возобновлением выполнения заданного задания.
-
task (trio.lowlevel.Task) – Задание, которое скоро будет выполнено.
Параметры:
-
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 -
class trio.abc.Instrument
В учебнике есть полный пример определения пользовательского инструмента для регистрации внутренних решений планирования Trio.
Запуск процессов низкого уровня
-
Выполняет дочернюю программу в новом процессе.
После создания вы можете взаимодействовать с дочерним процессом, записывая данные в его поток
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или на Windowscommandможет быть строкой или байтами, которые будут обработаны в соответствии с платформозависимыми правилами цитирования. Во всех случаяхcommandможет быть путем или последовательностью путей.stdin – Указывает, к чему должен быть подключен стандартный входной поток дочернего процесса: вывод родительского процесса (
subprocess.PIPE), ничего (subprocess.DEVNULL) или открытый файл (передать дескриптор файла или что-либо, чей методfilenoвозвращает один). Еслиstdinне указан, у дочернего процесса будет такой же стандартный входной поток, как у родительского.stdout – Аналогично
stdin, но для стандартного выходного потока дочернего процесса.stderr – Аналогично
stdin, но для стандартного потока ошибок дочернего процесса. Поддерживается дополнительное значениеsubprocess.STDOUT, которое вызывает объединение сообщений стандартного вывода и стандартной ошибки дочернего процесса в один стандартный выходной поток, подключенный к тому, к чему подключен опциейstdout.**options – Также принимаются другие общие параметры подпроцессов.
-
Новый объект
trio.Process. -
OSError – если запуск процесса завершился ошибкой, например, потому что указанная команда не найдена.
Параметры:
Возвращает:
Возможные исключения:
-
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
Примитивы низкоуровневого ввода-вывода
Разные среды предоставляют разные низкоуровневые API для выполнения асинхронного ввода-вывода. trio.lowlevel предоставляет эти API относительно прямым способом, чтобы предоставить максимальную мощность и гибкость для кода более высокого уровня. Однако это означает, что точный предоставляемый API может меняться в зависимости от системы, на которой работает Trio.
Универсально доступный API
Все среды предоставляют следующие функции:
-
Ожидает, пока ядро сообщит, что данный объект готов к чтению.
В системах Unix,
objдолжен быть целочисленным дескриптором файла или объектом с методом.fileno(), возвращающим целочисленный дескриптор файла. Можно передавать любой тип дескриптора файла, хотя точное поведение будет зависеть от ядра. Например, для файлов на диске это вряд ли будет полезно.В системах Windows,
objдолжен быть целочисленным дескрипторомSOCKET-ручки или объектом с методом.fileno(), возвращающим целочисленный дескрипторSOCKET-ручки. Дескрипторы файлов не поддерживаются, как и ручки, ссылающиеся на что-либо, кромеSOCKET.-
trio.BusyResourceError – если другая задача уже ожидает, пока данный сокет станет готовым к чтению.
trio.ClosedResourceError – если другая задача вызывает
notify_closing(), пока эта функция ещё работает.
Исключения:
-
await trio.lowlevel.wait_readable(obj)
-
Ожидает, пока ядро сообщит, что данный объект готов к записи.
См.
wait_readableдля определенияobj.-
trio.BusyResourceError – если другая задача уже ожидает, пока данный сокет станет готовым к записи.
trio.ClosedResourceError – если другая задача вызывает
notify_closing(), пока эта функция ещё работает.
Исключения:
-
await trio.lowlevel.wait_writable(obj)
-
Вызовите эту функцию перед закрытием дескриптора файла (в Unix) или сокета (в Windows). Это заставит любые вызовы
wait_readableилиwait_writableна данном объекте немедленно проснуться и выброситьClosedResourceError.Эта функция не закрывает объект – вам все равно нужно сделать это самостоятельно. Также следует быть осторожным, чтобы не запускать новые задачи, ожидающие объекта между вызовом этой функции и фактическим закрытием. Поэтому, чтобы корректно закрыть что-то, обычно нужно выполнить эти шаги в определённом порядке:
Явно пометить объект как закрытый, чтобы любые новые попытки использовать его прервались до начала.
Вызвать
notify_closing, чтобы разбудить всех уже существующих пользователей.Фактически закрыть объект.
Также можно выполнить их в другом порядке, если вы гарантируете отсутствие контрольных точек между шагами. Таким образом, все они происходят в едином атомарном шаге, и другие задачи не смогут определить порядок их выполнения.
trio.lowlevel.notify_closing(obj)
API, специфичное для Unix
FdStream поддерживает обертывание Unix-файлов (таких как пайп или TTY) в виде потока.
Если у вас есть два разных дескриптора файлов для отправки и получения и вы хотите объединить их в единый двунаправленный поток Stream, используйте trio.StapledStream:
bidirectional_stream = trio.StapledStream(
trio.lowlevel.FdStream(write_fd),
trio.lowlevel.FdStream(read_fd)
) -
Базируется на
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 для обсуждения сложностей, связанных с ослаблением этого ограничения.
class trio.lowlevel.FdStream(fd: int)
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
-
Асинхронный и отменяемый вариант WaitForSingleObject. Только для Windows.
-
handle – Движок Win32, как целое число Python.
-
OSError – Если движок недействителен, например, если он уже закрыт.
Параметры:
Исключения:
-
await trio.lowlevel.WaitForSingleObject(handle)
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
Глобальное состояние: системные задачи и локальные переменные выполнения
-
Локальный вариант переменной контекста.
RunVarобъекты похожи на объекты переменных контекста, за исключением того, что они используются в рамках одного вызоваtrio.run(), а не одной задачи.
class trio.lowlevel.RunVar(name: str, default | type[~trio._core._local._NoValue] = ...)
-
Запустить системную задачу.
Системные задачи отличаются от обычных задач:
Им не нужна явная подсистема; вместо этого они попадают во внутреннюю «системную подсистему».
Если системная задача вызывает исключение, то оно преобразуется в
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.lowlevel.spawn_system_task(async_fn: Callable[[Unpack[PosArgT]], Awaitable[object]], *args: Unpack[PosArgT], name: object = None, context: contextvars.Context | None = None) → Task
Токены Trio
-
Непрозрачный объект, представляющий один вызов
trio.run().У него нет публичного конструктора; вместо этого см.
current_trio_token().Этот объект используется в двух случаях:
Он позволяет повторно войти в цикл выполнения Trio из внешних потоков или обработчиков сигналов. Это низкоуровневый примитив, который использует
trio.to_thread()иtrio.from_threadдля взаимодействия с рабочими потоками,trio.open_signal_receiverиспользует для получения уведомлений о сигналах и так далее.Каждый вызов
trio.run()имеет ровно один связанный объектTrioToken, поэтому вы можете использовать его для идентификации конкретного вызова.
-
Запланировать вызов
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; нет правил о том, как вызовы разных категорий упорядочиваются друг относительно друга.-
trio.RunFinishedError – если связанный вызов
trio.run()уже завершился. (Любой вызов, не вызывающий это исключение, гарантированно будет полностью обработан до завершенияtrio.run().)
Raises:
run_sync_soon(sync_fn: Callable[[Unpack[PosArgsT]], object], *args: Unpack[PosArgsT], idempotent: bool = False) → None
class trio.lowlevel.TrioToken
-
Получить
TrioTokenдля текущего вызоваtrio.run().
trio.lowlevel.current_trio_token() → TrioToken
Запуск потоков
-
Выполняет
deliver(outcome.capture(fn))в рабочем потоке.Как правило,
fnвыполняет некоторую блокирующую работу, аdeliverвозвращает результат тому, кто в этом заинтересован.Это низкоуровневый интерфейс без излишеств, очень похожий на использование
threading.Threadдля прямого запуска потока. Основное отличие заключается в том, что эта функция пытается повторно использовать потоки при возможности, поэтому она может быть немного быстрее, чемthreading.Thread.Рабочие потоки имеют флаг
daemon, который означает, что если основной поток завершится, рабочие потоки будут автоматически убиты. Если вы хотите убедиться, что вашfnвыполнится до конца, убедитесь, что основной поток остаётся активным до вызоваdeliver.Безопасно вызывать эту функцию одновременно из нескольких потоков.
-
fn (функция синхронизации) – Выполняет произвольную блокирующую работу.
deliver (функция синхронизации) – Принимает
outcome.Outcomefnи предоставляет его. Не должна блокировать.
Parameters:
Поскольку рабочие потоки кэшируются и повторно используются для нескольких вызовов, ни одна из функций не должна изменять состояние уровня потока, например, объекты
threading.local– или, если они это делают, должны быть осторожны, чтобы восстановить свои изменения перед возвратом.Примечание
Разделение между
fnиdeliverслужит двум целям. Во-первых, это удобно, поскольку большинство вызывающих функций так или иначе нуждаются в чём-то подобном.Во-вторых, это позволяет избежать небольшой проблемы гонки, которая может привести к созданию слишком большого количества потоков. Рассмотрим программу, которая хочет последовательно выполнять несколько задач в потоке, поэтому основной поток отправляет задачу, ждёт её завершения, отправляет другую задачу и так далее. Теоретически, этой программе понадобится только один рабочий поток. Но что может произойти:
Рабочий поток: первая задача завершается и вызывает
deliver.Основной поток: получает уведомление о завершении задачи и вызывает
start_thread_soon.Основной поток: видит, что ни один рабочий поток не помечен как свободный, поэтому создаёт второй рабочий поток.
Исходный рабочий поток: отмечает себя как свободный.
Чтобы этого избежать, потоки отмечают себя как свободные перед вызовом
deliver.Является ли этот потенциальный дополнительный поток серьёзной проблемой? Возможно, нет, но его легко избежать, и мы считаем, что если пользователь пытается ограничить количество используемых потоков, то вежливо будет уважать это.
-
trio.lowlevel.start_thread_soon(fn: Callable[[], RetT], deliver: Callable[[outcome.Outcome[RetT]], object], name: str | None = None) → None
Обработка прерывания KeyboardInterrupt с повышенной безопасностью
Обработка Trio прерывания Ctrl+C разработана для баланса удобства использования и безопасности. С одной стороны, существуют чувствительные области (например, основной цикл планирования), где просто невозможно обработать произвольные исключения KeyboardInterrupt, сохраняя при этом основные инварианты корректности. С другой стороны, если пользователь случайно напишет бесконечный цикл, мы хотим иметь возможность прервать его. Наше решение заключается в установке обработчика сигнала по умолчанию, который проверяет, безопасно ли вызвать исключение KeyboardInterrupt в месте получения сигнала. Если да, то мы это делаем; в противном случае, мы планируем доставку KeyboardInterrupt в главную задачу в ближайшую доступную возможность (аналогично тому, как доставляется Cancelled).
Итак, это замечательно, но – как мы можем узнать, находимся ли мы в одной из чувствительных частей программы или нет?
Это определяется на основе каждой функции. По умолчанию:
Функция верхнего уровня в обычных пользовательских задачах не защищена.
Функция верхнего уровня в системных задачах защищена.
Если функция не указывает иное, то она наследует состояние защиты своего вызывающего объекта.
Это означает, что вам нужно переопределить значения по умолчанию только в тех местах, где вы переходите от защищенного кода к незащищенному или наоборот.
Эти переходы выполняются с помощью двух декораторов функций:
Декоратор, который помечает заданную обычную функцию, функцию-генератор, асинхронную функцию или асинхронную функцию-генератор как незащищенную от
KeyboardInterrupt, то есть код внутри этой функции может быть грубо прерванKeyboardInterruptв любой момент.Если у вас несколько декораторов на одной функции, то этот декоратор должен находиться внизу стека (ближе к фактической функции).
Пример использования – в реализации чего-то вроде
trio.from_thread.run(), которая используетTrioToken.run_sync_soon()для входа в поток Trio. Обработчикиrun_sync_soon()выполняются с включенной защитойKeyboardInterrupt, аtrio.from_thread.run()использует это для безопасной настройки механизма отправки ответа обратно в исходный поток, но затем используетdisable_ki_protection()при входе в функцию, предоставленную пользователем.
@trio.lowlevel.disable_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.enable_ki_protection
Проверить, включена ли защита от
KeyboardInterruptв вызываемом коде.Удивительно легко думать, что защита от
KeyboardInterruptвключена, когда она не включена, или наоборот. Эта функция сообщает вам, что думает об этом Trio, что делает ее полезной дляassertи unit-тестов.-
True, если защита включена, и False в противном случае.
Возвращает:
Тип возвращаемого значения:
-
trio.lowlevel.currently_ki_protected() → bool
Ожидание и пробуждение
Абстракция очереди ожидания
-
Справедливая очередь ожидания с возможностью отмены и повторного помещения в очередь.
Этот класс обобщает сложные части реализации очереди ожидания. Он полезен для реализации синхронизирующих примитивов более высокого уровня, таких как очереди и блокировки.
В дополнение к методам ниже, вы можете использовать
len(parking_lot)для получения количества приостановленных задач иif parking_lot: ...для проверки наличия приостановленных задач.: list[Task] broken_by-
Разбить эту парковку, с задачей
taskотмеченной как задача, которая это сделала.Это приводит к тому, что все приостановленные задачи генерируют ошибку, а любые будущие попытки парковки также вызовут ошибку. Unpark и repark становятся пустыми операциями, так как парковка пуста.
Возникающая ошибка содержит ссылку на задачу, переданную в качестве параметра. Задача также сохраняется в парковке в атрибуте
broken_by.
break_lot(task: Task | None = None) → None-
Приостановить текущую задачу до тех пор, пока она не будет разбужена вызовом
unpark()илиunpark_all().-
BrokenResourceError – если попытка парковки выполняется в сломанной парковке или парковка ломается, прежде чем мы дойдём до разбуживания.
Исключения:
-
await park() → 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) – парковка, в которую нужно переместить задачи.
Параметры:
-
repark(new_lot: ParkingLot, *, count: int | float = 1) → None-
Переместить все приостановленные задачи из одного объекта
ParkingLotв другой.См.
repark()для деталей.
repark_all(new_lot: ParkingLot) → None-
Возвращает объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
tasks_waiting: Количество задач, заблокированных в методеpark()этой парковки.
statistics() → ParkingLotStatistics-
Разбудить одну или несколько задач.
Это разбудит
countзадачи, заблокированные вpark(). Если приостановленных задач меньше, чемcount, то разбуживаются доступные задачи, и затем возвращается успешный результат.
unpark(*, count: int | float = 1) → list[Task]-
Разбудить все приостановленные задачи.
unpark_all() → list[Task] -
class trio.lowlevel.ParkingLot
-
Объект, содержащий отладочную информацию для ParkingLot.
В настоящее время определены следующие поля:
tasks_waiting(int): Количество задач, заблокированных в методеtrio.lowlevel.ParkingLot.park()этой парковки.
class trio.lowlevel.ParkingLotStatistics(tasks_waiting: int)
-
Регистрирует задачу как разрушитель парковки. См.
ParkingLot.break_lot().-
trio.BrokenResourceError – если задача уже завершена.
Исключения:
-
trio.lowlevel.add_parking_lot_breaker(task: Task, lot: ParkingLot) → None
-
Удаляет регистрацию задачи как разрушителя парковки. См.
ParkingLot.break_lot()
trio.lowlevel.remove_parking_lot_breaker(task: Task, lot: ParkingLot) → None
Функции контрольных точек низкого уровня
-
Чистая контрольная точка.
Проверяет отмену и позволяет планировать другие задачи без блокировки.
Обратите внимание, что планировщик может проигнорировать это и продолжить выполнение текущей задачи, если посчитает это целесообразным (например, для повышения эффективности).
Эквивалентно
await trio.sleep(0)(которое реализуется вызовомcheckpoint().)
await trio.lowlevel.checkpoint() → None
Следующие две функции используются вместе для создания контрольной точки:
-
Вызовите контрольную точку, если контекст вызывающего элемента был отменён.
Эквивалентно (но потенциально более эффективному):
if trio.current_effective_deadline() == -inf: await trio.lowlevel.checkpoint()Это либо нет операция, либо она позволяет планировать другие задачи, а затем генерирует исключение
trio.Cancelled.Обычно используется вместе с
cancel_shielded_checkpoint().
await trio.lowlevel.checkpoint_if_cancelled() → None
-
Введите точку планирования, но не точку отмены.
Это не контрольная точка, но это половина контрольной точки, а в сочетании с
checkpoint_if_cancelled()она может образовать полную контрольную точку.Эквивалентно (но потенциально более эффективному):
with trio.CancelScope(shield=True): await trio.lowlevel.checkpoint()
await trio.lowlevel.cancel_shielded_checkpoint() → None
Они часто используются в случаях, когда у вас есть операция, которая может или не может заблокировать выполнение, и вы хотите реализовать стандартную семантику контрольных точек 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() перед запуском потока.
Низкоуровневое блокирование
-
Уложить текущую задачу в сон с поддержкой отмены.
Это низкоуровневый API для блокирования в Trio. Каждый раз, когда задача
Taskблокируется, она делает это, вызывая эту функцию (обычно косвенно через какой-либо API более высокого уровня).Это сложный интерфейс без защитных ограждений. Если вы можете использовать
ParkingLotили встроенные функции ожидания ввода-вывода, то вам следует это сделать.В общем случае, перед вызовом этой функции, вы подготавливаете «кого-то», кто вызовет
reschedule()для текущей задачи в какой-то момент позже.Затем вы вызываете
wait_task_rescheduled(), передаваяabort_func, «обработчик отмены».(Терминология: в Trio «отмена» — это процесс попытки прервать заблокированную задачу для передачи отмены.)
Есть два возможных сценария дальнейших действий:
«Кто-то» вызывает
reschedule()для текущей задачи, иwait_task_rescheduled()возвращает или генерирует любое значение или ошибку, которые были переданы вreschedule().-
Контекст вызова переходит в состояние отмены (например, из-за истечения таймаута). В этом случае вызывается
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 более высокого уровня, если это возможно?
await trio.lowlevel.wait_task_rescheduled(abort_func: Callable[[Callable[[], NoReturn]], Abort]) → Any
-
enum.Enumиспользуется в качестве возвращаемого значения от функций отмены.См.
wait_task_rescheduled()для получения подробной информации.SUCCEEDEDFAILED
class trio.lowlevel.Abort(value, names=None, *, module=None, qualname=None, type=None, start=1, boundary=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().
Параметры:
-
trio.lowlevel.reschedule(task: Task, next_send: Outcome[object] =
Вот пример класса блокировки, реализованного с помощью 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 задач
-
Возвращает текущую корневую
Task.Это задача, которая является конечным родителем всех других задач.
trio.lowlevel.current_root_task()
-
Возвращает объект
Task, представляющий текущую задачу.-
объект
Task, который вызвалcurrent_task().
Возвращает:
Тип возвращаемого значения:
-
trio.lowlevel.current_task()
-
Объект
Taskпредставляет собой конкурирующую «поток» выполнения. У него нет публичного конструктора; Trio внутренне создает объектTaskдля каждого вызоваnursery.start(...)илиnursery.start_soon(...).Его публичные члены в основном полезны для интроспекции и отладки:
-
Строка, содержащая имя этой задачи
Task. Обычно это имя функции, в которой выполняется эта задачаTask, но может быть переопределено путем передачиname=вstartилиstart_soon.
name-
Объект корутины этой задачи.
coro-
Рекурсивно итерируется по объектам корутин, на которых ожидает эта задача, и возвращает кадр и номер строки в каждом кадре.
Это аналогично
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()))
for ... in iter_await_frames() → Iterator[tuple[types.FrameType, int]]-
Объект
contextvars.Contextэтой задачи.
context-
Питомник, в котором находится эта задача (или None, если это задача «init»).
Пример использования: построение визуализации дерева задач в отладчике.
parent_nursery-
Питомник, в котором эта задача будет находиться после вызова
task_status.started().Если эта задача уже вызвала
started(), или если она не была запущена с помощьюnursery.start(), то ееeventual_parent_nurseryявляетсяNone.
eventual_parent_nursery-
Питомники, содержащиеся в этой задаче.
Это список, в котором внешние питомники находятся перед внутренними.
child_nurseries-
Trio не присваивает этому переменной никакого значения, за исключением того, что устанавливает ее в
Noneвсякий раз, когда задача перепланируется. Его можно использовать для обмена данными между различными задачами, участвующими в приостановке задачи и ее повторном запуске. (См.wait_task_rescheduled()для получения подробностей.)
custom_sleep_data -
class trio.lowlevel.Task
Использование «гостевого режима» для запуска 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: тот же планировщик, тот же код ввода-вывода, то же самое во всем. Таким образом, вы получаете полный набор функций, и все работает так, как ожидается.
Простая интеграция и широкая совместимость: практически каждый цикл событий предлагает какую-то потокобезопасную операцию «планирования обратного вызова», и этого достаточно, чтобы использовать его как хостовый цикл.
Действительно? Как это возможно?
Примечание
Вы можете использовать гостевой режим, не читая этот раздел. Он включен для тех, кто любит понимать, как работают вещи.
Все циклы событий имеют одинаковую базовую структуру. Они циклически выполняют две операции:
Ожидание уведомления операционной системы о том, что произошло что-то интересное, например, данные, прибывающие на сокет, или истекло время ожидания. Это делается путем вызова специфичного для платформы системного вызова
sleep_until_something_happens()—select,epoll,kqueue,GetQueuedCompletionEventsи т.д.Выполнение всех пользовательских задач, относящихся к тому, что произошло, а затем возвращение к шагу 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 поверх другого цикла событий «хоста».
Каждый цикл хоста может иметь только одно гостевое выполнение одновременно.
Вы всегда должны дождаться завершения выполнения 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. -
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
Передача живых объектов корутин между исполнителями корутин
Внутри синтаксис async/await в Python построен на основе концепции «объектов корутин» и «исполнителей корутин». Объект корутины представляет состояние стека асинхронного вызова. Но сам по себе это просто статический объект, который просто находится там. Если вы хотите, чтобы он что-то делал, вам нужен исполнитель корутины, чтобы продвигать его дальше. Каждая задача Trio имеет связанный с ней объект корутины (см. Task.coro), и планировщик Trio действует как их исполнитель корутины.
Но, конечно, Trio не единственный исполнитель корутин в Python – asyncio имеет своего, другие циклы событий тоже, вы даже можете определить свой собственный.
И в некоторых очень, очень необычных обстоятельствах даже имеет смысл передавать один объект корутины туда и обратно между разными исполнителями корутин. Об этом и говорится в данном разделе. Это чрезвычайно экзотический случай использования и предполагает глубокое понимание того, как Python async/await работает внутри. Для примеров мотивации см. вопрос #42 на GitHub для trio-asyncio и вопрос #649 на GitHub для trio. Для получения более подробной информации о работе корутин мы рекомендуем «Рассказ о циклах событий» Андре Карона или обратиться непосредственно к PEP 492 для получения всех подробностей.
-
Постоянно отсоединить текущую задачу от планировщика Trio.
Обычно задача Trio не завершается, пока не завершится её объект корутины. Когда вы вызываете эту функцию, Trio ведет себя так, как будто объект корутины только что завершился, и задача завершается с заданным результатом. Это полезно, если вы хотите постоянно переключить объект корутины на другой исполнитель корутин.
Когда вызываемая корутина входит в эту функцию, она выполняется в рамках Trio, а когда функция возвращает значение, она выполняется в рамках внешнего исполнителя корутин.
Вы должны убедиться, что объект корутины освободил все ресурсы, специфичные для Trio, которые он получил (например, nurseries).
-
final_outcome (outcome.Outcome) – Trio ведет себя так, как будто текущая задача завершилась с указанным возвращаемым значением или исключением.
Параметры:
Возвращает или вызывает любое значение или исключение, используемые новым исполнителем корутин для возобновления корутины.
-
await trio.lowlevel.permanently_detach_coroutine_object(final_outcome: Outcome[object]) → 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.temporarily_detach_coroutine_object(abort_func: Callable[[Callable[[], NoReturn]], Abort]) → object
-
Подключение объекта корутины, который был отсоединён с помощью
temporarily_detach_coroutine_object().Когда вызываемая корутина входит в эту функцию, она выполняется в рамках внешнего исполнителя корутин, а когда функция возвращает значение, она выполняется в рамках Trio.
Это необходимо вызвать внутри возобновляемой корутины, и она возвращает то значение, которое вы передали. (Предполагается, что вы передадите значение, которое заставит текущий исполнитель корутин прекратить планирование этой задачи.) Затем корутина возобновляется планировщиком Trio в ближайшее время.
await trio.lowlevel.reattach_detached_coroutine_object(task: Task, yield_value: object) → None
© 2017 Nathaniel J. Smith
Licensed under the MIT License.
https://trio.readthedocs.io/en/v0.29.0/reference-lowlevel.html