Основной функционал Trio
Вход в Trio
Если вы хотите использовать Trio, то первое, что нужно сделать, это вызвать trio.run():
-
Запустить функцию async, предназначенную для Trio, и вернуть результат.
Вызов:
run(async_fn, *args)
эквивалентен:
await async_fn(*args)
за исключением того, что
run()можно (и нужно) вызывать из синхронного контекста.Это основной входной пункт Trio. Почти каждая другая функция в Trio требует, чтобы вы находились внутри вызова
run().-
async_fn – Функция async.
args – Позиционные аргументы, передаваемые в async_fn. Если вам нужно передать именованные аргументы, используйте
functools.partial().clock –
Noneдля использования системного монотонного таймера по умолчанию; в противном случае — объект, реализующий интерфейсtrio.abc.Clock, такой как (например) экземплярtrio.testing.MockClock.instruments (список объектов
trio.abc.Instrument) – Любая инструментировка, которую вы хотите применить к этому запуску. Это также можно изменить во время выполнения; см. Интерфейс инструментов.-
restrict_keyboard_interrupt_to_checkpoints (bool) –
Что произойдёт, если пользователь нажмёт Ctrl+C во время работы
run()? Если этот аргумент равен False (по умолчанию), то происходит стандартное поведение Python: исключениеKeyboardInterruptнемедленно прерывает любую выполняющуюся задачу (или если нет выполняющихся задач, то Trio разбудит задачу для прерывания). В противном случае, если вы установите этот аргумент в True, доставкаKeyboardInterruptбудет отложена: она будет только генерироваться в контрольных точках, как исключениеCancelled.По умолчанию это удобно, потому что это означает, что даже если вы случайно напишете бесконечный цикл, который никогда не выполняет контрольных точек, вы всё равно сможете прервать его с помощью Ctrl+C. Альтернативное поведение удобно, если вы боитесь
KeyboardInterruptв неподходящем месте, оставляющем вашу программу в несогласованном состоянии, потому что это означает, что вам нужно беспокоиться оKeyboardInterruptтолько в тех же местах, где вам уже нужно беспокоиться оCancelled.Это настройка не влияет, если ваша программа зарегистрировала обработчик SIGINT, или если
run()вызывается не из главного потока (это ограничение Python), или если вы используетеopen_signal_receiver()для перехвата SIGINT. strict_exception_groups (bool) – Если не установлено в False, nurseries всегда будут оборачивать даже одно поднятое исключение в группе исключений. Это можно переопределить на уровне отдельных nurseries. Установка в False будет устаревшим методом и в конечном итоге удалена в будущих версиях Trio.
-
Всё, что возвращает
async_fn. -
TrioInternalError – если возникла непредвиденная ошибка внутри внутреннего механизма Trio. Это ошибка, и вы должны сообщить нам об этом.
Другое – если
async_fnгенерирует исключение, тоrun()его распространяет.
Параметры:
Возвращает:
Возбуждает:
-
trio.run(async_fn: Callable[[Unpack[PosArgT]], Awaitable[RetT]], *args: Unpack[PosArgT], clock: Clock | None = None, instruments: Sequence[Instrument] = (), restrict_keyboard_interrupt_to_checkpoints: bool = False, strict_exception_groups: bool = True) → RetT
Общие принципы
Контрольные точки
При написании кода с использованием Trio очень важно понимать концепцию контрольной точки. Многие функции Trio действуют как контрольные точки.
Контрольная точка — это две вещи:
Это точка, в которой Trio проверяет отмену. Например, если код, вызвавший вашу функцию, установил таймаут, и этот таймаут истек, то в следующий раз, когда ваша функция выполнит контрольную точку, Trio поднимет исключение
Cancelled. Более подробную информацию см. в разделе Отмена и таймауты ниже.Это точка, в которой планировщик Trio проверяет свою политику планирования, чтобы определить, подходит ли момент для переключения на другую задачу, и, возможно, делает это. (В настоящее время эта проверка очень проста: планировщик всегда переключается на каждой контрольной точке. Но это может измениться в будущем.)
При написании кода Trio вам нужно отслеживать расположение контрольных точек. Почему? Во-первых, потому что контрольные точки требуют дополнительного внимания: всякий раз, когда вы выполняете контрольную точку, вы должны быть готовы обработать ошибку Cancelled или для выполнения другой задачи, которая изменит некоторые состояния под вами. Во-вторых, потому что вам также необходимо убедиться, что у вас достаточно контрольных точек: если ваш код не проходит через контрольную точку регулярно, то он будет медленно реагировать на отмену и — что гораздо хуже — так как Trio — это кооперативная многозадачная система, где планировщик может переключать задачи только на контрольных точках, это также помешает планировщику справедливо распределять время между различными задачами и негативно повлияет на время отклика всего остального кода, выполняемого в том же процессе. (Неформально мы говорим, что задача, которая делает это, «монополизирует цикл выполнения».)
Итак, когда вы проводите код-ревью проекта, использующего Trio, одним из моментов, которые вы должны обдумать, является наличие достаточного количества контрольных точек и правильная обработка каждой из них. Конечно, это означает, что вам нужен способ распознать контрольные точки. Как это сделать? Основной принцип заключается в том, что любая операция блокирования должна быть контрольной точкой. Это имеет смысл: если операция блокируется, то она может блокироваться долгое время, и вы захотите иметь возможность отменить ее, если истекает таймаут; и в любом случае, в то время как эта задача заблокирована, мы хотим запланировать выполнение другой задачи, чтобы наш код мог полностью использовать ЦП.
Но если мы хотим написать корректный код на практике, то этот принцип слишком небрежен и неточен, чтобы быть полезным. Как мы узнаем, какие функции могут заблокироваться? Что, если функция блокируется иногда, но не в других случаях, в зависимости от передаваемых аргументов/скорости сети/фазы Луны? Как мы выясним, где находятся контрольные точки, когда мы устали и лишены сна, но все равно хотим провести этот код-ревью правильно и хотели бы сохранить свои умственные ресурсы для размышления над фактической логикой, а не беспокоиться о контрольных точках?
Не волнуйтесь — Trio позаботится о вас. Поскольку контрольные точки важны и повсеместны, мы делаем отслеживание их максимально простым. Вот правила:
Обычные (синхронные) функции никогда не содержат контрольных точек.
-
Если вы вызываете асинхронную функцию, предоставленную Trio (
await <something in trio>), и она не вызывает исключение, то она всегда действует как контрольная точка. (Если она вызывает исключение, она может действовать как контрольная точка или нет.)Это включает асинхронные итераторы: Если вы пишете
async for ... in <a trio object>, то в каждой итерации цикла будет хотя бы одна контрольная точка, и она все равно будет контрольной точкой, если итератор пуст.-
Частичное исключение для асинхронных менеджеров контекста: как вход, так и выход из блока
async withопределены как асинхронные функции; но для определенного типа асинхронного менеджера контекста часто бывает так, что только одна из них может блокироваться, что означает, что только одна из них будет действовать как контрольная точка. Это документируется на индивидуальной основе.trio.open_nursery()является дополнительным исключением из этого правила.
Функции/итераторы/менеджеры контекста сторонних производителей могут действовать как контрольные точки; если вы видите
await <something>или одного из его «друзей», то это может быть контрольная точка. Поэтому для безопасности вы должны быть готовы к планированию или отмене там.
Причина, по которой мы различаем функции Trio и другие функции, заключается в том, что мы не можем дать никаких гарантий относительно кода сторонних производителей. Свойство контрольной точки является транзитивным: если функция А действует как контрольная точка, и вы пишете функцию, которая вызывает функцию А, то ваша функция также действует как контрольная точка. Если вы этого не сделаете, то нет. Так что ничего не мешает кому-то написать функцию, например:
# technically legal, but bad style:
async def why_is_this_async():
return 7 которая никогда не вызывает асинхронных функций Trio. Эта функция является асинхронной, но не является контрольной точкой. Но зачем делать функцию асинхронной, если она никогда не вызывает асинхронных функций? Возможно, но это плохая идея. Если у вас есть функция, которая не вызывает никаких асинхронных функций, то вы должны сделать ее синхронной. Пользователи вашей функции поблагодарят вас, потому что это делает очевидным, что ваша функция не является контрольной точкой, и код-ревью пройдет быстрее.
(Помните, как в учебнике мы подчеркнули важность «асинхронного сэндвича» и то, как он означает, что await в итоге становится маркером, который показывает, когда вы вызываете функцию, которая вызывает функцию, которая … в конечном итоге вызывает одну из встроенных асинхронных функций Trio? Транзитивность асинхронности является техническим требованием Python, но поскольку она точно соответствует транзитивности контрольной точки, мы можем использовать это, чтобы помочь вам отслеживать контрольные точки. Довольно хитро, не так ли?)
Несколько более сложный случай — функция, подобная:
async def sleep_or_not(should_sleep):
if should_sleep:
await trio.sleep(1)
else:
pass Здесь функция действует как контрольная точка, если вы вызываете ее с should_sleep, установленным в истинное значение, но не иначе. Вот почему мы подчеркиваем, что собственные асинхронные функции Trio являются безусловными контрольными точками: они всегда проверяют отмену и проверяют планирование независимо от передаваемых аргументов. Если вы обнаружите асинхронную функцию в Trio, которая не следует этому правилу, то это ошибка, и вы должны сообщить нам.
Внутри Trio мы очень требовательны к этому, потому что Trio — основа всей системы, поэтому мы считаем, что дополнительные усилия по повышению предсказуемости стоят того. Вам решать, насколько требовательным вы хотите быть в своем коде. Чтобы дать вам более реалистичный пример того, как выглядит эта проблема в реальной жизни, рассмотрим эту функцию:
async def recv_exactly(sock, nbytes):
data = bytearray()
while nbytes > 0:
# recv() reads up to 'nbytes' bytes each time
chunk = await sock.recv(nbytes)
if not chunk:
raise RuntimeError("socket unexpected closed")
nbytes -= len(chunk)
data += chunk
return data Если она вызывается с nbytes, большим нуля, то она будет вызывать sock.recv по крайней мере один раз, а recv — это асинхронная функция Trio, и, следовательно, безусловная контрольная точка. Таким образом, в данном случае recv_exactly действует как контрольная точка. Но если мы сделаем await
recv_exactly(sock, 0), то она немедленно вернет пустой буфер без выполнения контрольной точки. Если бы это была функция в самом Trio, то это было бы неприемлемо, но вы можете решить, что не хотите беспокоиться об этом типе мелких граничных случаев в собственном коде.
Если вы хотите быть внимательным или у вас есть код, связанный с ЦП, в котором недостаточно контрольных точек, то полезно знать, что await trio.sleep(0) — это идиоматический способ выполнения контрольной точки без выполнения каких-либо других действий, и что trio.testing.assert_checkpoints() можно использовать для проверки того, что произвольный блок кода содержит контрольную точку.
Безопасность потоков
Подавляющее большинство API Trio не является потокобезопасным: его можно использовать только внутри вызова trio.run(). В этом руководстве мы не будем документировать это для отдельных вызовов; если не указано иное, вы должны предполагать, что вызов любых функций Trio из любого места, кроме потока Trio, небезопасен. (Но см. ниже, если вам действительно нужно работать с потоками.)
Время и часы
Каждый вызов run() связан с часами.
По умолчанию Trio использует неуказанные монотонные часы, но это можно изменить, передав объект пользовательских часов в run() (например, для тестирования).
Не следует предполагать, что внутренние часы Trio совпадают с любыми другими часами, к которым у вас есть доступ, включая часы одновременных вызовов trio.run(), происходящих в других процессах или потоках!
Текущие стандартные часы реализованы как time.perf_counter() плюс большое случайное смещение. Цель состоит в том, чтобы поймать код, который случайно использует time.perf_counter() рано, что должно помочь сохранить открытыми возможности для изменения реализации часов в будущем, и (что более важно) убедиться, что вы можете быть уверены, что пользовательские часы, такие как trio.testing.MockClock, будут работать с библиотеками сторонних разработчиков, которыми вы не управляете.
-
Возвращает текущее время по внутренним часам Trio.
-
Текущее время.
-
RuntimeError – если не внутри вызова
trio.run().
Возвращает:
Тип возвращаемого значения:
Исключения:
-
trio.current_time() → float
-
Приостановить выполнение текущей задачи на указанное количество секунд.
-
seconds (float) – Количество секунд для сна. Может быть нулём, чтобы вставить контрольную точку без реального блокирования.
-
ValueError – если seconds отрицательно или NaN.
Параметры:
Исключения:
-
await trio.sleep(seconds: float) → None
-
Приостановить выполнение текущей задачи до указанного времени.
Разница между
sleep()иsleep_until()заключается в том, что первый принимает относительное время, а второй — абсолютное время по внутренним часам Trio (как возвращаетcurrent_time()).-
deadline (float) – Время, по которому мы должны снова проснуться. Может быть в прошлом, в этом случае эта функция выполняет контрольную точку, но не блокирует.
-
ValueError – если deadline является NaN.
Параметры:
Исключения:
-
await trio.sleep_until(deadline: float) → None
-
Приостановить выполнение текущей задачи навсегда (или до отмены).
Эквивалентно вызову
await sleep(math.inf), за исключением того, что при ручном перепланировании это вызоветRuntimeError.-
RuntimeError – если перепланировано
Исключения:
-
await trio.sleep_forever() → NoReturn
Если вы сумасшедший ученый или по какой-либо другой причине чувствуете необходимость прямо управлять ТЕЧЕНИЕМ ВРЕМЕНИ, вы можете реализовать пользовательский класс Clock:
-
Интерфейс для пользовательских часов цикла выполнения.
-
Возвращает текущее время по этим часам.
Используется для реализации функций, таких как
trio.current_time()иtrio.move_on_after().-
Текущее время.
Возвращает:
Тип возвращаемого значения:
-
abstractmethod current_time() → float-
Вычислить реальное время до указанного срока.
Вызывается перед входом в системно-специфическую функцию ожидания, такую как
select.select(), чтобы получить время ожидания для передачи.Для часов, использующих время реального мира, это должно быть примерно так:
return deadline - self.current_time()
но, конечно, это может быть другим, если вы реализуете какие-либо виртуальные часы.
abstractmethod deadline_to_sleep_time(deadline: float) → float-
Выполните все необходимые действия по настройке часов.
Вызывается в начале выполнения.
abstractmethod start_clock() → None -
class trio.abc.Clock
Отмена и таймауты
Trio имеет богатую, композиционную систему отмены задач, как явным образом, так и при истечении таймаута.
Пример простого таймаута
В самом простом случае можно применить таймаут к блоку кода:
with trio.move_on_after(30):
result = await do_http_get("https://...")
print("result is", result)
print("with block finished") Мы называем move_on_after() созданием «области отмены», которая содержит весь код, выполняющийся внутри блока with. Если запрос HTTP займет более 30 секунд, он будет отменен: мы прервем запрос и не увидим result is ... на консоли; вместо этого мы сразу напечатаем сообщение with block
finished.
Примечание
Обратите внимание, что это единственный таймаут в 30 секунд для всего тела оператора
with. Это отличается от того, что вы могли видеть в других библиотеках Python, где таймауты часто относятся к чему-то более сложному. Мы считаем, что такой подход проще для понимания.
Как это работает? Здесь нет магии: Trio построено с использованием обычных возможностей Python, поэтому мы не можем просто бросить код внутри блока with. Вместо этого мы используем стандартный способ Python для прерывания большой и сложной части кода: мы генерируем исключение.
Вот идея: всякий раз, когда вы вызываете функцию, допускающую отмену, такую как await
trio.sleep(...) или await sock.recv(...) – см. Контрольные точки – то первое, что делает эта функция, – проверяет, есть ли окружающая область отмены, таймаут которой истек, или она была отменена. Если это так, то вместо выполнения запрошенной операции функция немедленно завершается с исключением Cancelled. В данном примере это, вероятно, происходит где-то глубоко внутри do_http_get. Исключение затем распространяется как любое обычное исключение (вы даже можете его перехватить, если захотите, но это, как правило, плохая идея), пока не достигнет блока with move_on_after(...):. В этот момент исключение Cancelled выполнило свою работу – успешно размотало всю область отмены – поэтому move_on_after() перехватывает его, и выполнение продолжается как обычно после блока with. И все это работает правильно, даже если у вас есть вложенные области отмены, потому что каждый объект Cancelled несет невидимый маркер, который гарантирует, что область отмены, которая его сгенерировала, является единственной, которая его перехватит.
Обработка отмены
Практически любой код, который вы пишете с использованием Trio, должен иметь какую-то стратегию обработки исключений Cancelled – даже если вы не установили таймаут, ваш вызывающий объект может (и, скорее всего, будет).
Вы можете перехватить исключение Cancelled, но не должны! Или точнее, если вы его перехватываете, то вы должны выполнить некоторую очистку, а затем перебросить его или иначе позволить ему продолжить распространение (если вы не столкнулись с ошибкой, в этом случае вы можете позволить распространиться ошибке). Чтобы напомнить вам об этом факте, Cancelled наследуется от BaseException, как и KeyboardInterrupt и SystemExit, чтобы он не перехватывался блоками перехвата всех исключений except
Exception:.
Также важно в любом длительном коде регулярно проверять отмену, иначе таймауты не будут работать! Это происходит неявно всякий раз, когда вы вызываете операнцию, допускающую отмену; подробности см. в разделе ниже. Если у вас есть задача, которая должна выполнить много работы без ввода-вывода, вы можете использовать await sleep(0) для вставки явного пункта отмены и планирования.
Вот практическое правило для проектирования хороших API в стиле Trio («трионика?»): если вы пишете многоразовую функцию, то не следует принимать параметр timeout=, а вместо этого пусть вызывающий объект позаботится об этом. Это имеет несколько преимуществ. Во-первых, это оставляет вызывающему объекту возможность выбора того, как ему удобнее обрабатывать таймауты — например, ему может быть проще работать с абсолютными крайними сроками, а не с относительными таймаутами. Если они вызывают механизм отмены, то они сами выбирают, и вам не нужно об этом беспокоиться. Во-вторых, и что более важно, это упрощает повторное использование вашего кода. Если вы напишете функцию http_get, а затем я напишу функцию log_in_to_twitter, которая нуждается во внутренней реализации нескольких вызовов http_get, я не хочу иметь дело с настройкой индивидуальных таймаутов для каждого из них — а в системе таймаутов Trio это совершенно не нужно.
Конечно, это правило не относится к API, которые должны накладывать внутренние таймауты. Например, если вы пишете функцию start_http_server, то, вероятно, вы должны предоставить вызывающему объекту способ настройки таймаутов для отдельных запросов.
Семантика отмены
Вы можете свободно вкладывать блоки отмены, и каждое исключение Cancelled «знает», к какому блоку оно относится. Пока вы его не остановите, исключение будет продолжать распространяться, пока не достигнет блока, который его поднял, в этот момент оно автоматически остановится.
Вот пример:
print("starting...")
with trio.move_on_after(5):
with trio.move_on_after(10):
await trio.sleep(20)
print("sleep finished without error")
print("move_on_after(10) finished without error")
print("move_on_after(5) finished without error") В этом коде внешний область действия истечёт через 5 секунд, что приведёт к раннему возврату вызова sleep() с исключением Cancelled. Затем это исключение будет распространяться через строку with
move_on_after(10), пока не будет перехвачено менеджером контекста with
move_on_after(5). Таким образом, этот код выведет:
starting... move_on_after(5) finished without error
В итоге Trio успешно отменил ровно ту работу, которая выполнялась в области действия, которая была отменена.
Глядя на это, вы можете задаться вопросом, как узнать, истекло ли время ожидания внутреннего блока — возможно, вы хотите сделать что-то другое, например, попробовать процедуру резервного копирования или сообщить о сбое нашему вызывающему объекту. Чтобы упростить это, функция move_on_after()´s __enter__ возвращает объект, представляющий эту область отмены, который мы можем использовать для проверки того, перехватила ли эта область исключение Cancelled:
with trio.move_on_after(5) as cancel_scope:
await trio.sleep(10)
print(cancel_scope.cancelled_caught) # prints "True" Объект cancel_scope также позволяет проверять или изменять крайний срок этой области, явно вызывать отмену без ожидания крайнего срока, проверять, была ли область уже отменена и так далее — см. CancelScope ниже для получения подробных сведений.
Отмены в Trio являются «триггерами уровня», что означает, что после отмены блока все отменяемые операции в этом блоке будут продолжать поднимать Cancelled. Это помогает избежать некоторых проблем с очисткой ресурсов. Например, представьте, что у нас есть функция, которая подключается к удалённому серверу и отправляет некоторые сообщения, а затем очищает при выходе:
with trio.move_on_after(TIMEOUT):
conn = make_connection()
try:
await conn.send_hello_msg()
finally:
await conn.send_goodbye_msg() Теперь предположим, что удалённый сервер перестаёт отвечать, поэтому наш вызов к await conn.send_hello_msg() зависает навсегда. К счастью, мы достаточно умны, чтобы поместить таймаут вокруг этого кода, поэтому в конечном итоге таймаут истечёт, и send_hello_msg поднимет Cancelled. Но затем, в блоке finally, мы выполняем ещё одну блокирующую операцию, которая тоже будет зависать навсегда! В этот момент, если бы мы использовали asyncio или другую библиотеку с отменой «триггеров края», у нас возникли бы проблемы: поскольку наш таймаут уже сработал, он больше не сработает, и в этот момент наше приложение зависнет навсегда. Но в Trio этого не происходит: вызов await
conn.send_goodbye_msg() всё ещё находится внутри отменённого блока, поэтому он также поднимет Cancelled.
Конечно, если вы действительно хотите выполнить другой блокирующий вызов в обработчике очистки, Trio позволит вам; он пытается предотвратить случайное нанесение удара себе в ногу. Намеренное нанесение удара себе в ногу не проблема (по крайней мере, это не проблема Trio). Для этого создайте новую область действия и установите атрибут shield в True:
with trio.move_on_after(TIMEOUT):
conn = make_connection()
try:
await conn.send_hello_msg()
finally:
with trio.move_on_after(CLEANUP_TIMEOUT, shield=True) as cleanup_scope:
await conn.send_goodbye_msg() Пока вы находитесь внутри области действия с shield = True установленным, вы будете защищены от внешних отмен. Однако имейте в виду, что это только относится к внешним отменам: если CLEANUP_TIMEOUT истечёт, то await conn.send_goodbye_msg() всё равно будет отменено, и если вызов await conn.send_goodbye_msg() использует какие-либо таймауты внутри, то они также будут работать нормально. Это довольно продвинутая функция, которую большинство людей, вероятно, не будут использовать, но она существует для редких случаев, когда она вам нужна.
Отмена и примитивные операции
Мы много говорили о том, что происходит, когда операция отменяется, и о том, что вам нужно быть к этому готовым при вызове отменяемой операции… но мы не углублялись в подробности о том, какие операции отменяются и как именно они ведут себя при отмене.
Вот правило: если операция находится в пространстве имён trio, и вы используете await для её вызова, то она отменяема (см. Контрольные точки выше). Отменяемая операция означает:
Если вы попытаетесь вызвать её, находясь внутри области отмены, она поднимет
Cancelled.Если она заблокируется, а пока она заблокирована, одна из областей вокруг неё станет отменённой, она вернётся пораньше и поднимет
Cancelled.Поднятие
Cancelledозначает, что операция не произошла. Если методsendсокета Trio подниметCancelled, то данные не были отправлены. Если методrecvсокета Trio подниметCancelled, то данные не были потеряны — они всё ещё находятся в буфере приёма сокета, ожидая, когда вы снова вызоветеrecv. И так далее.
Есть несколько специфичных случаев, когда внешние ограничения делают невозможным полное выполнение этих семантик. Они всегда документированы. Также есть одно систематическое исключение:
Операции асинхронной очистки — такие как методы
__aexit__или асинхронные методы закрытия — отменяются так же, как и всё остальное, за исключением того, что при их отмене они всё равно выполнят минимальный уровень очистки, прежде чем поднятьCancelled.
Например, закрытие сокета, обернутого TLS, обычно включает в себя отправку уведомления удалённому узлу, чтобы они могли быть криптографически уверены в том, что вы действительно намеревались закрыть сокет, а ваша соединение не было просто прервано злоумышленником, находящимся посредине. Но надёжное решение этой задачи немного сложно. Помните наш пример пример выше, где блокирующий send_goodbye_msg вызвал проблемы? Именно так работает закрытие сокета, обернутого TLS: если удалённый узел исчез, то наш код может никогда не смочь фактически отправить наше уведомление о завершении работы, и было бы неплохо, если бы он не блокировался вечно, пытаясь это сделать. Поэтому метод закрытия TLS-обёрнутого сокета попытается отправить это уведомление — и если он будет отменён, то он откажется от отправки сообщения, но закроет лежащий в основе сокет до поднятия Cancelled, поэтому вы хотя бы не потеряете этот ресурс.
Подробности API отмены
move_on_after() и все остальные средства отмены, предоставляемые Trio, в конечном счёте реализованы с помощью объектов CancelScope.
-
Область отмены: связь между блоком отменяемой работы и системой отмены Trio.
Объект
CancelScopeсвязывается с отменяемой работой, когда используется в качестве менеджера контекста, окружающего эту работу:cancel_scope = trio.CancelScope() ... with cancel_scope: await long_running_operation()Внутри блока
withотменаcancel_scope(через вызов его методаcancel()или при истечении срокаdeadline) немедленно прервёт работуlong_running_operation(), вызвав исключениеCancelledв следующей точке проверки проверки.Менеджер контекста
__enter__возвращает сам объектCancelScope, поэтому вы также можете написатьwith trio.CancelScope() as cancel_scope:.Если область отмены отменяется до входа в блок
with, исключениеCancelledбудет поднято в первой точке проверки внутри блокаwith. Это позволяет создатьCancelScopeв одной задаче и передать её другой, чтобы первая задача могла позже отменить работу во второй.Области отмены не могут быть повторно использованы или рекурсивны; то есть каждая область отмены может быть использована не более чем в одном блоке
with. (Вы получитеRuntimeError, если нарушите это правило.)Конструктор
CancelScopeпринимает начальные значения для атрибутов области отменыdeadlineиshield; их можно свободно изменять после создания, независимо от того, был ли введён контекст области, и изменения вступают в силу немедленно.-
Чтение/запись,
float. Абсолютное время на часах текущего выполнения, когда эта область автоматически будет отменена. Вы можете изменить крайний срок, изменив этот атрибут, например:# I need a little more time! cancel_scope.deadline += 30
Обратите внимание, что для повышения эффективности ядро цикла выполнения проверяет истечение крайних сроков лишь время от времени. Это означает, что в некоторых случаях может быть небольшая задержка между тем, когда часы показывают, что крайний срок истек, и когда проверки начинают поднимать исключение
Cancelled. Это очень редкий случай, который вы вряд ли заметите, но мы документируем его для полноты картины. (Конечно, если это всё же вызывает проблемы, то дайте нам знать!)По умолчанию равно
math.inf, что означает «нет крайнего срока», хотя это можно переопределить аргументомdeadline=в конструкторCancelScope.
deadlinerelative_deadline-
Чтение/запись,
bool, по умолчаниюFalse. Пока это значение установлено вTrue, код внутри этой области не будет получать исключенияCancelledот областей, находящихся за пределами этой области. Они всё ещё могут получать исключенияCancelledот (1) этой области или (2) областей внутри этой области. Вы можете изменить этот атрибут:with trio.CancelScope() as cancel_scope: cancel_scope.shield = True # This cannot be interrupted by any means short of # killing the process: await sleep(10) cancel_scope.shield = False # Now this can be cancelled normally: await sleep(10)По умолчанию
False, но это можно переопределить аргументомshield=в конструкторCancelScope.
shield-
Возвращает None после входа. Возвращает False, если и deadline, и relative_deadline равны inf.
is_relative()-
Немедленно отменяет эту область.
Этот метод идемпотентен, т.е. если область уже была отменена, то этот метод ничего не делает.
cancel()-
Только чтение,
bool. Указывает, поймала ли эта область исключениеCancelled. Для этого требуется два условия: (1) блокwithзавершился с исключениемCancelled, и (2) эта область отвечала за срабатывание этого исключенияCancelled.
cancelled_caught-
Только чтение,
bool. Указывает, была ли запрошена отмена для этой области, либо явным вызовом методаcancel(), либо по истечении крайнего срока.То, что этот атрибут имеет значение True, не обязательно означает, что код внутри области был или будет затронут отменой. Например, если
cancel()был вызван после последней точки проверки в блокеwith, когда слишком поздно передать исключениеCancelled, то этот атрибут всё равно будет True.Этот атрибут в основном полезен для отладки и интроспекции. Если вы хотите узнать, был ли отменён какой-либо фрагмент кода, то
cancelled_caughtобычно более уместен.
cancel_called -
class trio.CancelScope(*, relative_deadline: float = inf, deadline: float = inf, shield: bool = False)
Часто нет необходимости создавать объекты CancelScope. Trio уже включает атрибут cancel_scope в связанном с задачей объекте Nursery. Мы рассмотрим ясли позже в руководстве.
Trio также предоставляет несколько удобных функций для распространённой ситуации, когда нужно просто установить тайм-аут для какого-либо кода:
-
Используйте в качестве менеджера контекста для создания области отмены, дедлайн которой установлен на текущее время + seconds.
Дедлайн области отмены вычисляется при входе.
-
ValueError – если значение
secondsменьше нуля или NaN.
Параметры:
Исключения:
-
Используйте в качестве менеджера контекста для создания области отмены с заданным абсолютным дедлайном.
-
ValueError – если deadline является NaN.
Параметры:
Исключения:
-
Создает область отмены с заданным таймаутом и поднимает ошибку, если она фактически отменена.
Эта функция и
move_on_after()аналогичны тем, что обе создают область отмены с заданным таймаутом, и если таймаут истекает, обе вызовутCancelledвнутри области. Разница в том, что когда исключениеCancelledдостигаетmove_on_after(), оно перехватывается и отбрасывается. Когда оно достигаетfail_after(), оно перехватывается, и вместо него поднимаетсяTooSlowError.Дедлайн области отмены вычисляется при входе.
-
TooSlowError – если исключение
Cancelledвозникает в этой области и перехватывается менеджером контекста.ValueError – если seconds меньше нуля или NaN.
Параметры:
Исключения:
-
Создает область отмены с заданным дедлайном и поднимает ошибку, если она фактически отменена.
Эта функция и
move_on_at()аналогичны тем, что обе создают область отмены с заданным абсолютным дедлайном, и если дедлайн истекает, обе вызовутCancelledвнутри области. Разница в том, что когда исключениеCancelledдостигаетmove_on_at(), оно перехватывается и отбрасывается. Когда оно достигаетfail_at(), оно перехватывается, и вместо него поднимаетсяTooSlowError.-
TooSlowError – если исключение
Cancelledвозникает в этой области и перехватывается менеджером контекста.ValueError – если deadline является NaN.
Параметры:
Исключения:
Список команд
-
Если вы хотите установить таймаут для функции, но вам все равно, истек ли он:
with trio.move_on_after(TIMEOUT): await do_whatever() # carry on! -
Если вы хотите установить таймаут для функции и затем выполнить восстановление, если он истек:
with trio.move_on_after(TIMEOUT) as cancel_scope: await do_whatever() if cancel_scope.cancelled_caught: # The operation timed out, try something else try_to_recover() -
Если вы хотите установить таймаут для функции и, если он истек, просто отказаться и выдать ошибку для обработки вызывающей стороной:
with trio.fail_after(TIMEOUT): await do_whatever()
Также можно проверить, какой текущий эффективный дедлайн, что иногда бывает полезно:
-
Возвращает текущий эффективный дедлайн для текущей задачи.
Эта функция проверяет все области отмены, которые в настоящее время действуют (с учетом экранирования), и возвращает дедлайн, который истечет первым.
Один из примеров, где это может быть полезно, — это если ваш код пытается решить, следует ли начать дорогостоящую операцию, например, вызов RPC, но хочет пропустить ее, если знает, что она не может завершиться в доступное время. Другой пример — если вы используете протокол, подобный gRPC, который передает информацию о таймауте удаленному узлу; эта функция предоставляет способ извлечения этой информации, чтобы вы могли передать ее.
Если этот метод вызывается в контексте, где активна отмена (т.е. блокирующий вызов сразу же поднимет
Cancelled), то возвращаемый дедлайн —-inf. Если он вызывается в контексте, где ни у одной области отмены нет установленного дедлайна, он возвращаетinf.-
эффективный дедлайн в виде абсолютного момента времени.
Возвращает:
Тип возвращаемого значения:
-
Задачи позволяют выполнять несколько действий одновременно
Один из основных принципов проектирования Trio: нет неявной конкурентности. Каждая функция выполняется простым образом, сверху вниз, завершая каждое действие перед переходом к следующему — как задумывал Гайдо.
Но, конечно же, вся суть асинхронной библиотеки — позволить выполнять несколько действий одновременно. Единственный способ сделать это в Trio — через интерфейс запуска задач. Так что если вы хотите, чтобы ваша программа и ходила, и жевала жвачку, этот раздел для вас.
Питомники и запуск
Большинство библиотек для параллельного программирования позволяют запускать новые дочерние задачи (или потоки, или что-то ещё) как угодно, когда и где вам захочется. Trio немного отличается: вы не можете запустить дочернюю задачу, пока не готовы быть ответственным родителем. Способ продемонстрировать свою ответственность — создание питомника:
async with trio.open_nursery() as nursery:
... И после того, как у вас есть ссылка на объект питомника, вы можете запускать дочерние задачи в этом питомнике:
async def child():
...
async def parent():
async with trio.open_nursery() as nursery:
# Make two concurrent calls to child()
nursery.start_soon(child)
nursery.start_soon(child) Это означает, что задачи образуют дерево: когда вы вызываете run(), тогда это создаёт начальную задачу, а все ваши другие задачи будут дочерними, внучатыми и т. д. начальной задачи.
По существу, тело блока async with действует как начальная задача, выполняемая внутри питомника, а каждый вызов nursery.start_soon добавляет ещё одну задачу, которая выполняется параллельно. Два важных момента, которые следует помнить:
Если какая-либо задача внутри питомника завершается с необработанным исключением, питомник немедленно отменяет все задачи внутри питомника.
Поскольку все задачи выполняются одновременно внутри блока
async with, блок не завершается, пока не завершатся все задачи. Если вы использовали другие фреймворки для параллельного выполнения, то можно представить, что отступ в конце блокаasync withавтоматически «соединяет» (ожидает завершения) все задачи в питомнике.-
После завершения всех задач:
Питомник помечается как «закрытый», что означает, что новые задачи не могут быть запущены внутри него.
Любые необработанные исключения повторно поднимаются внутри родительской задачи, сгруппированные в одно исключение
BaseExceptionGroupилиExceptionGroup.
Поскольку все задачи являются потомками начальной задачи, следствием этого является то, что run() не может завершиться, пока не завершатся все задачи.
Примечание
Оператор return не отменит питомник, если в нём всё ещё выполняются задачи:
async def main(): async with trio.open_nursery() as nursery: nursery.start_soon(trio.sleep, 5) return trio.run(main)Этот код подождёт 5 секунд (пока завершится дочерняя задача), а затем вернётся.
Дочерние задачи и отмена
В Trio дочерние задачи наследуют области отмены родительского питомника. Итак, в этом примере обе дочерние задачи будут отменены, когда истечёт таймаут:
with trio.move_on_after(TIMEOUT):
async with trio.open_nursery() as nursery:
nursery.start_soon(child1)
nursery.start_soon(child2) Обратите внимание, что здесь важны области отмены, которые были активны при вызове open_nursery(), а не области отмены, активные при вызове start_soon. Например, блок таймаута ниже ничего не делает:
async with trio.open_nursery() as nursery:
with trio.move_on_after(TIMEOUT): # don't do this!
nursery.start_soon(child) Почему так? Ну, start_soon() возвращается как только запланирует запуск новой задачи. Поток выполнения в родительской задаче затем продолжается и выходит из блока with trio.move_on_after(TIMEOUT):, в этот момент Trio полностью забывает о таймауте. Чтобы таймаут применился к дочерней задаче, Trio должен уметь понять, что связанная с ним область отмены останется открытой как минимум до тех пор, пока выполняется дочерняя задача. И Trio может узнать это наверняка только если блок области отмены находится вне блока питомника.
Вы можете задаться вопросом, почему Trio не может просто запомнить «эта задача должна быть отменена через TIMEOUT секунд», даже после того, как блок with trio.move_on_after(TIMEOUT): исчез. Причина связана с реализацией отмены. Вспомните, что отмена представлена исключением Cancelled, которое в конечном итоге должно быть перехвачено областью отмены, которая её вызвала. (В противном случае исключение остановило бы всю вашу программу!) Чтобы иметь возможность отменить дочерние задачи, область отмены должна иметь возможность «видеть» исключения Cancelled, которые они поднимают — а эти исключения возникают из блока async with open_nursery(), а не из вызова start_soon().
Если вы хотите, чтобы таймаут применялся к одной задаче, но не к другой, вам нужно поместить область отмены в функцию этой отдельной задачи — child() в этом примере.
Ошибки в нескольких дочерних задачах
Обычно в Python происходит только одно действие за раз, что означает, что только одна ошибка может возникнуть за раз. Trio не имеет такого ограничения. Рассмотрим код, подобный:
async def broken1():
d = {}
return d["missing"]
async def broken2():
seq = range(10)
return seq[20]
async def parent():
async with trio.open_nursery() as nursery:
nursery.start_soon(broken1)
nursery.start_soon(broken2) broken1 вызывает KeyError. broken2 вызывает IndexError. Очевидно, что parent должен вызвать какую-то ошибку, но какую? Ответ заключается в том, что обе исключения сгруппированы в ExceptionGroup. ExceptionGroup и её родительский класс BaseExceptionGroup используются для инкапсуляции нескольких исключений, возникающих одновременно.
Для перехвата отдельных исключений, инкапсулированных в группе исключений, в Python 3.11 был введён блок except* (PEP 654). Вот как это работает:
try:
async with trio.open_nursery() as nursery:
nursery.start_soon(broken1)
nursery.start_soon(broken2)
except* KeyError as excgroup:
for exc in excgroup.exceptions:
... # handle each KeyError
except* IndexError as excgroup:
for exc in excgroup.exceptions:
... # handle each IndexError Если вы хотите повторно возбудить исключения или возбудить новые, вы можете это сделать, но имейте в виду, что исключения, возбужденные в блоках except*, будут возбуждены вместе в новой группе исключений.
Но что, если вы не можете использовать Python 3.11, а значит except* пока ещё недоступно? Та же библиотека exceptiongroup, которая реализует обратную совместимость с ExceptionGroup, также позволяет приблизительно воспроизвести это поведение с помощью обратного вызова обработчика исключений:
from exceptiongroup import catch
def handle_keyerrors(excgroup):
for exc in excgroup.exceptions:
... # handle each KeyError
def handle_indexerrors(excgroup):
for exc in excgroup.exceptions:
... # handle each IndexError
with catch({
KeyError: handle_keyerrors,
IndexError: handle_indexerrors
}):
async with trio.open_nursery() as nursery:
nursery.start_soon(broken1)
nursery.start_soon(broken2) Семантика функций-обработчиков совпадает с семантикой блоков except*, за исключением установки локальных переменных. Если вам нужно установить локальные переменные, вам нужно объявить их внутри функций-обработчиков с помощью ключевого слова nonlocal:
def handle_keyerrors(excgroup):
nonlocal myflag
myflag = True
myflag = False
with catch({KeyError: handle_keyerrors}):
async with trio.open_nursery() as nursery:
nursery.start_soon(broken1) Проектирование на случай нескольких ошибок
Структурированная конкурентность — это всё ещё молодой шаблон проектирования, но есть несколько шаблонов, которые мы выявили, чтобы вы (или ваши пользователи) могли обработать группы исключений. Обратите внимание, что конечный шаблон — просто возбудить ExceptionGroup, является наиболее распространённым — и nurseries автоматически делают это за вас.
Во-первых, вы можете «отдать приоритет» определённому типу исключения, возбудив именно его, если такой экземпляр есть в группе. Например: KeyboardInterrupt имеет чёткий смысл для окружающего кода, может разумно иметь приоритет над ошибками других типов, и не важно, одна у вас такая ошибка или несколько.
Этот шаблон часто можно реализовать с помощью декоратора или менеджера контекста, например, trio_util.multi_error_defer_to() или trio_util.defer_to_cancelled(). Однако обратите внимание, что повторное возбуждение «листового» исключения отбросит любую часть трассировки, связанную с ExceptionGroup самим по себе, поэтому мы не рекомендуем это для ошибок, которые будут показаны пользователям.
Во-вторых, вы можете рассматривать конвейер внутри своего кода как реализационную деталь, скрытую от ваших пользователей — например, абстрагировать протокол, который включает отправку и получение данных, до простого интерфейса только для получения, или реализовать менеджер контекста, который поддерживает некоторые фоновые задачи на протяжении блока async with.
Простой вариант — raise MySpecificError from group, позволяющий пользователям обработать ошибку вашей библиотеки. Это просто и надёжно, но не полностью скрывает nursery. Не разворачивайте отдельные исключения, если может возникнуть несколько исключений; это всегда приводит к скрытым ошибкам и затем к сбоям.
Более сложный вариант — убедиться, что может произойти только одно исключение за раз. Это очень сложно, например, вам нужно будет как-то обработать KeyboardInterrupt, и мы настоятельно рекомендуем иметь raise PleaseReportBug from group резервную обработку на случай, если вы получите группу, содержащую более одного исключения. Это полезно при написании менеджера контекста, который запускает некоторые фоновые задачи, а затем уступает место коду пользователя, который фактически выполняется «встроенно» в теле блока nursery. В этом случае фоновые задачи можно обернуть, например, с помощью библиотеки outcome, чтобы гарантировать, что может быть возбуждено только одно исключение (из кода конечного пользователя); а затем вы можете либо raise SomeInternalError, если фоновая задача завершилась неудачей, либо развернуть исключение пользователя, если это была единственная ошибка.
В-третьих, и чаще всего, существование nursery в вашем коде не является просто реализационной деталью, и вызывающие функции должны быть готовы обработать несколько исключений в виде ExceptionGroup, будь то с помощью except*, ручного анализа или просто передачей её своим вызывающим функциям. Поскольку это так часто встречается, nurseries по умолчанию ведут себя именно так, и вам ничего не нужно делать.
Историческая заметка: «нестрогие» ExceptionGroup
В ранних версиях Trio синтаксис except* ещё не был придуман, и мы не работали со структурированной конкурентностью долго или в больших кодовых базах. В качестве уступки удобству некоторые API возбуждали одиночные исключения и оборачивали исключения по конкурентности только в старом типе trio.MultiError, если их было два или более.
К сожалению, результаты были неудовлетворительными: вызывающий код часто не понимал, что некоторая функция может возбудить MultiError, и поэтому обрабатывал только общий случай — с результатом, что всё работало хорошо в тестировании, а затем падало под большей нагрузкой (обычно в производстве). asyncio.TaskGroup извлек урок из этого опыта и всегда оборачивает ошибки в ExceptionGroup, как это делает и anyio, а начиная с Trio 0.25 это тоже наше поведение по умолчанию.
В настоящее время мы поддерживаем аргумент совместимости strict_exception_groups=False для trio.run и trio.open_nursery, который восстанавливает старое поведение (хотя сам MultiError был полностью удалён). Мы настоятельно не рекомендуем использовать его в новом коде и рекомендуем существующим пользователям перейти на него — мы считаем этот вариант устаревшим и планируем удалить его после периода документирования, а затем предупреждений во время выполнения.
Запуск задач без установления родительских отношений
Иногда задача, которая запускает дочернюю задачу, не должна нести ответственность за её наблюдение. Например, задача сервера может хотеть запускать новую задачу для каждого подключения, но она не может одновременно прослушивать подключения и контролировать дочерние задачи.
Решение здесь просто, когда вы его видите: нет требования, чтобы объект nursery оставался в той задаче, которая его создала! Мы можем написать код, подобный этому:
async def new_connection_listener(handler, nursery):
while True:
conn = await get_new_connection()
nursery.start_soon(handler, conn)
async def server(handler):
async with trio.open_nursery() as nursery:
nursery.start_soon(new_connection_listener, handler, nursery) Обратите внимание, что server открывает nursery и передаёт его в new_connection_listener, а затем new_connection_listener может запускать новые задачи в качестве «братьев» самой себя. Конечно, в этом случае мы могли бы точно так же написать:
async def server(handler):
async with trio.open_nursery() as nursery:
while True:
conn = await get_new_connection()
nursery.start_soon(handler, conn) …но иногда всё не так просто, и этот трюк оказывается полезным.
Однако следует помнить, что области отмены наследуются от nursery, а не от задачи, которая вызывает start_soon. Итак, в данном примере таймаут не применяется к child (или к чему-либо ещё):
async def do_spawn(nursery):
with trio.move_on_after(TIMEOUT): # don't do this, it has no effect
nursery.start_soon(child)
async with trio.open_nursery() as nursery:
nursery.start_soon(do_spawn, nursery) Настраиваемые супервизоры
Стандартная логика очистки часто достаточна для простых случаев, но что, если вам нужен более сложный супервизор? Например, возможно, у вас есть Erlang envy и вам нужны такие функции, как автоматический перезапуск аварийных задач. Сам Trio не предоставляет таких функций, но вы можете их реализовать сверху; цель Trio — обеспечить базовую гигиену и освободить вас от проблем. (То есть: Trio не позволит вам создать супервизор, который завершит работу и оставит сироты задачи, и если у вас есть необработанное исключение из-за ошибок или лени, то Trio позаботится о том, чтобы они распространились.) А затем вы можете обернуть свой сложный супервизор в библиотеку и разместить её в PyPI, потому что супервизоры сложны, и нет причин, по которым каждый должен писать свой собственный.
Например, вот функция, которая принимает список функций, выполняет их все одновременно и возвращает результат от первой завершённой:
async def race(*async_fns):
if not async_fns:
raise ValueError("must pass at least one argument")
winner = None
async def jockey(async_fn, cancel_scope):
nonlocal winner
winner = await async_fn()
cancel_scope.cancel()
async with trio.open_nursery() as nursery:
for async_fn in async_fns:
nursery.start_soon(jockey, async_fn, nursery.cancel_scope)
return winner Это работает, запуская набор задач, которые каждая пытается выполнить свою функцию. Как только первая функция завершит своё выполнение, задача установит нелокальную переменную winner из внешнего пространства имён в результат функции и отменит другие задачи с помощью переданной области отмены. После того, как все задачи будут отменены (что завершает работу блока nursery), переменная winner будет возвращена.
В этом случае, если одна или несколько из конкурирующих функций вызовут необработанное исключение, стандартная обработка Trio вступит в силу: она отменит остальные и затем распространит исключение. Если вы хотите другое поведение, вы можете получить его, добавив блок try к функции jockey, чтобы перехватить исключения и обработать их как вам нужно.
Хранилище данных задачи
Представьте, что вы пишете сервер, который отвечает на сетевые запросы, и вы записываете некоторую информацию о каждом запросе по мере его обработки. Если сервер загружен и обрабатывается несколько запросов одновременно, то у вас могут появиться логи, подобные этому:
Request handler started Request handler started Request handler finished Request handler finished
В этом журнале трудно понять, какие строки относятся к какому запросу. (Запрос, который начался первым, также завершился первым или нет?) Один из способов решения этой проблемы — назначить каждому запросу уникальный идентификатор и затем включить этот идентификатор в каждое сообщение журнала:
request 1: Request handler started request 2: Request handler started request 2: Request handler finished request 1: Request handler finished
Таким образом, мы видим, что запрос 1 был медленным: он начался до запроса 2, но завершился после него. (Вы также можете сделать намного более сложную систему, но этого достаточно для примера.)
Теперь проблема в том, как код ведения журнала узнает, какой идентификатор запроса нужно использовать? Один подход — явно передавать его каждой функции, которая может генерировать логи… но это фактически каждая функция, потому что никогда не знаешь, когда тебе может понадобиться добавить вызов log.debug(...) к какой-то служебной функции, глубоко вложенной в стеке вызовов, а когда ты находишься в середине отладки сложной проблемы, последнее, чего тебе хочется, это прерваться и переписать всё, чтобы передавать идентификатор запроса! Иногда это правильное решение, но в других случаях гораздо удобнее хранить идентификатор в глобальной переменной, чтобы функция ведения журнала могла его просмотреть всякий раз, когда это необходимо. Но… глобальная переменная может содержать только одно значение за раз, поэтому если у нас работает несколько обработчиков одновременно, то это не сработает. Нам нужно что-то, что похоже на глобальную переменную, но может иметь разные значения в зависимости от того, какой обработчик запросов к ней обращается.
Для решения этой проблемы в стандартной библиотеке Python есть модуль: contextvars.
Вот пример, демонстрирующий, как использовать contextvars:
import random
import trio
import contextvars
request_info = contextvars.ContextVar("request_info")
# Example logging function that tags each line with the request identifier.
def log(msg):
# Read from task-local storage:
request_tag = request_info.get()
print(f"request {request_tag}: {msg}")
# An example "request handler" that does some work itself and also
# spawns some helper tasks to do some concurrent work.
async def handle_request(tag):
# Write to task-local storage:
request_info.set(tag)
log("Request handler started")
await trio.sleep(random.random())
async with trio.open_nursery() as nursery:
nursery.start_soon(concurrent_helper, "a")
nursery.start_soon(concurrent_helper, "b")
await trio.sleep(random.random())
log("Request received finished")
async def concurrent_helper(job):
log(f"Helper task {job} started")
await trio.sleep(random.random())
log(f"Helper task {job} finished")
# Spawn several "request handlers" simultaneously, to simulate a
# busy server handling multiple requests at the same time.
async def main():
async with trio.open_nursery() as nursery:
for i in range(3):
nursery.start_soon(handle_request, i)
trio.run(main) Пример вывода (ваш может незначительно отличаться):
request 1: Request handler started request 2: Request handler started request 0: Request handler started request 2: Helper task a started request 2: Helper task b started request 1: Helper task a started request 1: Helper task b started request 0: Helper task b started request 0: Helper task a started request 2: Helper task b finished request 2: Helper task a finished request 2: Request received finished request 0: Helper task a finished request 1: Helper task a finished request 1: Helper task b finished request 1: Request received finished request 0: Helper task b finished request 0: Request received finished
Для получения дополнительной информации, прочитайте документацию по contextvars.
Синхронизация и обмен данными между задачами
Trio предоставляет стандартный набор примитивов синхронизации и межзадачного обмена данными. API этих объектов в целом моделируется по аналогии с соответствующими классами в стандартной библиотеке, но с некоторыми отличиями.
Блокирующие и неблокирующие методы
Примитивы синхронизации стандартной библиотеки имеют различные механизмы для указания тайм-аутов и поведения блокировки, а также для сигнализации о том, завершилась ли операция из-за успеха или тайм-аута.
В Trio мы придерживаемся следующих соглашений:
Мы не предоставляем аргументы тайм-аута. Если вам нужен тайм-аут, используйте область отмены.
Для операций, имеющих неблокирующий вариант, блокирующий и неблокирующий варианты являются различными методами с именами, например,
XиX_nowaitсоответственно. (Это аналогичноqueue.Queue, но в отличие от большинства классов вthreading.) Мы предпочитаем этот подход, так как он позволяет нам сделать блокирующую версию асинхронной, а неблокирующую — синхронной.Когда неблокирующий метод не может преуспеть (канал пуст, замок уже захвачен и т. п.), он вызывает
trio.WouldBlock. Нет эквивалента различениюqueue.Emptyиqueue.Full— у нас просто есть одно исключение, которое мы используем последовательно.
Справедливость
Эти классы гарантированно «справедливы», то есть, когда наступает время выбора следующего, кто получит блокировку, элемент из очереди и т. п., всегда выбирается задача, которая ожидает дольше всего. Пока не совсем ясно, является ли это наилучшим выбором, но пока работает именно так.
В качестве примера того, что это означает, приведём небольшую программу, в которой две задачи конкурируют за блокировку. Обратите внимание, что задача, которая разблокирует, всегда сразу же пытается повторно захватить её, прежде чем другая задача успеет запуститься. (И помните, что здесь используется кооперативное многозадачность, поэтому на самом деле детерминировано, что задача, освобождающая блокировку, вызовет acquire() до того, как проснётся другая задача; в Trio освобождение блокировки не является контрольной точкой.) При несправедливой блокировке это привело бы к тому, что одна и та же задача будет удерживать блокировку вечно, а другая будет голодать. Но если вы запустите это, вы увидите, что две задачи вежливо по очереди делят блокировку:
# fairness-demo.py
import trio
async def loopy_child(number, lock):
while True:
async with lock:
print(f"Child {number} has the lock!")
await trio.sleep(0.5)
async def main():
async with trio.open_nursery() as nursery:
lock = trio.Lock()
nursery.start_soon(loopy_child, 1, lock)
nursery.start_soon(loopy_child, 2, lock)
trio.run(main) Вещание события с Event
-
Переменная булевого типа с ожиданием, полезная для межзадачной синхронизации, вдохновлённая
threading.Event.Объект события имеет внутренний булевый флаг, обозначающий, произошло ли событие. Флаг изначально False, и метод
wait()ожидает, пока флаг не станет True. Если флаг уже True, тогдаwait()возвращается немедленно. (Если событие уже произошло, ждать нечего.) Методset()устанавливает флаг в True и будит все ожидающие задачи.Это поведение полезно, так как помогает избежать гонок и потерянных пробуждений: неважно, вызывается ли
set()непосредственно перед или послеwait(). Если вам нужен более низкоуровневый примитив пробуждения, не имеющий этой защиты, рассмотритеConditionилиtrio.lowlevel.ParkingLot.Примечание
В отличие от
threading.Event,trio.Eventне имеет методаclear. В Trio, как только событие произошло, оно не может быть отменено. Если вам нужно представить серию событий, создавайте новые объектыEventдля каждого из них (они дешевы!), или используйте другие методы синхронизации, такие как каналы илиtrio.lowlevel.ParkingLot.-
Возвращает текущее значение внутреннего флага.
is_set() → bool-
Устанавливает значение внутреннего флага в True и будит все ожидающие задачи.
set() → None-
Возвращает объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
tasks_waiting: Количество задач, заблокированных в методеwait()этого события.
statistics() → EventStatistics-
Блокируется до тех пор, пока значение внутреннего флага не станет True.
Если оно уже True, метод возвращается немедленно.
await wait() → None -
class trio.Event
-
Объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
tasks_waiting: Количество задач, заблокированных в методеtrio.Event.wait()этого события.
class trio.EventStatistics(tasks_waiting: int)
Использование каналов для передачи значений между задачами
Каналы позволяют безопасно и удобно отправлять объекты между различными задачами. Они особенно полезны для реализации шаблонов «производитель/потребитель».
Основной API каналов определен абстрактными базовыми классами trio.abc.SendChannel и trio.abc.ReceiveChannel. Вы можете использовать их для реализации собственных пользовательских каналов, которые выполняют такие действия, как передача объектов между процессами или по сети. Но во многих случаях вам просто нужно передавать объекты между различными задачами внутри одного процесса, и для этого вы можете использовать trio.open_memory_channel():
-
Открытие канала для передачи объектов между задачами в пределах процесса.
Память каналы лёгкие, недорогие в выделении и полностью находятся в памяти. Они не используют системные ресурсы операционной системы или какой-либо сериализации. Они просто передают объекты Python напрямую между задачами (с возможной остановкой во внутреннем буфере по пути).
Объекты каналов могут быть закрыты вызовом
acloseили с помощьюasync with. Они не автоматически закрываются при сборке мусора. Закрытие каналов памяти не обязательно, но в целом это хорошая идея, так как это помогает избежать ситуаций, когда задачи застревают, ожидая канала, когда с другой стороны никого нет. Подробнее см. Чистая остановка с каналами.Операции канала памяти являются атомными относительно отмены, либо
receiveуспешно вернёт объект, либо он возбудитCancelled, оставив канал неизменным.-
max_buffer_size (int или math.inf) – Максимальное количество элементов, которое может быть буферизовано в канале, прежде чем
send()заблокируется. Выбор разумного значения здесь важен для обеспечения своевременной обратной связи о задержке и предотвращения ненужной задержки; подробнее см. Буферизация в каналах. В случае сомнений используйте 0. -
Пара
(send_channel, receive_channel). Если у вас проблемы с запоминанием порядка, помните: данные текут слева → справа.
Параметры:
Возвращаемое значение:
В дополнение к стандартным методам канала все объекты каналов памяти предоставляют метод
statistics(), который возвращает объект со следующими полями:current_buffer_used: Количество элементов, хранящихся в настоящее время в буфере канала.max_buffer_size: Максимальное количество элементов, разрешённых в буфере, как передано вopen_memory_channel().open_send_channels: Количество открытыхMemorySendChannelконечных точек, указывающих на этот канал. Изначально 1, но может быть увеличено с помощьюMemorySendChannel.clone().open_receive_channels: Аналогично, но для открытыхMemoryReceiveChannelконечных точек.tasks_waiting_send: Количество задач, заблокированных вsendна этом канале (суммируется по всем клонам).tasks_waiting_receive: Количество задач, заблокированных вreceiveна этом канале (суммируется по всем клонам).
-
trio.open_memory_channel(max_buffer_size)
Примечание
Если вы использовали модули
threadingилиasyncio, вам могут быть знакомыqueue.Queueилиasyncio.Queue. В Trio,open_memory_channel()— то, что вы используете, когда ищете очередь. Основное различие заключается в том, что Trio разбивает классический интерфейс очереди на два объекта. Преимущество этого заключается в том, что можно поместить оба конца в разные процессы, не переписывая свой код, и можно закрыть обе стороны по отдельности.
MemorySendChannel и MemoryReceiveChannel также предоставляют несколько дополнительных функций помимо основного интерфейса канала:
-
-
Асинхронно закрыть этот канал отправки.
await aclose() → None-
Создать копию этого канала отправки.
Возвращает новый объект
MemorySendChannel, который является дубликатом оригинала: отправка в новый объект делает ровно то же самое, что и отправка в старый объект. (Если вы знакомы сos.dup, то это аналогичная идея.)Однако закрытие одного из объектов не закрывает другой, и получатели не получат
EndOfChannel, пока не будут закрыты все копии.Это полезно для схем обмена, в которых несколько производителей отправляют объекты в один и тот же пункт назначения. Если вы предоставите каждому производителю свою копию
MemorySendChannel, а затем убедитесь, что вы закрываете каждыйMemorySendChannelпри завершении, получатели автоматически будут уведомлены о завершении всех производителей. См. Управление несколькими производителями и/или несколькими потребителями для примеров.-
trio.ClosedResourceError – если вы уже закрыли этот
MemorySendChannelобъект.
Исключения:
-
clone() → MemorySendChannel[SendType]-
Синхронно закрыть этот канал отправки.
Все объекты каналов имеют асинхронный метод
aclose. Каналы памяти также могут быть закрыты синхронно. Это оказывает такое же влияние на канал и другие задачи, использующие его, ноcloseне является контрольной точкой trio. Это упрощает очистку в отменённых задачах.Использование
with send_channel:закроет объект канала при выходе из блока with.
close() → None-
См.
SendChannel.send.Каналы памяти позволяют нескольким задачам одновременно вызывать
send.
await send(value: SendType) → None-
Аналогично
send, но если буфер канала заполнен, вместо блокирования генерируетWouldBlock.
send_nowait(value: SendType) → None-
Возвращает
MemoryChannelStatisticsдля канала памяти, с которым связан этот объект.
statistics() → MemoryChannelStatistics -
class trio.MemorySendChannel(*args: object, **kwargs: object)
-
-
Асинхронно закрыть этот канал получения.
await aclose() → None-
Создать копию этого канала получения.
Возвращает новый
MemoryReceiveChannelобъект, который является дубликатом оригинала: получение из нового объекта делает ровно то же самое, что и получение из старого объекта.Однако закрытие одного из объектов не закрывает другой, и основной канал не закрывается, пока не будут закрыты все копии. (Если вы знакомы с
os.dup, то это аналогичная идея.)Это полезно для схем обмена, в которых несколько потребителей получают объекты из одного основного канала. См. Управление несколькими производителями и/или несколькими потребителями для примеров.
Предупреждение
Копии делят один и тот же основной канал. Когда копия
receive()значение, оно удаляется из канала, и другие копии не получают это значение. Если вы хотите отправлять несколько копий одного и того же потока значений в несколько пунктов назначения, какitertools.tee(), то вам нужно найти другое решение; этот метод этого не делает.-
trio.ClosedResourceError – если вы уже закрыли этот
MemoryReceiveChannelобъект.
Исключения:
-
clone() → MemoryReceiveChannel[ReceiveType]-
Синхронно закрыть этот канал получения.
Все объекты каналов имеют асинхронный метод
aclose. Каналы памяти также могут быть закрыты синхронно. Это оказывает такое же влияние на канал и другие задачи, использующие его, ноcloseне является контрольной точкой trio. Это упрощает очистку в отменённых задачах.Использование
with receive_channel:закроет объект канала при выходе из блока with.
close() → None-
Каналы памяти позволяют нескольким задачам одновременно вызывать
receive. Первая задача получит первый отправленный элемент, вторая задача получит второй отправленный элемент и так далее.
await receive() → ReceiveType-
Аналогично
receive, но если ничего не готово для получения, генерируетWouldBlockвместо блокирования.
receive_nowait() → ReceiveType-
Возвращает
MemoryChannelStatisticsдля канала памяти, с которым связан этот объект.
statistics() → MemoryChannelStatistics -
class trio.MemoryReceiveChannel(*args: object, **kwargs: object)
class trio.MemoryChannelStatistics(current_buffer_used: int, max_buffer_size: int | float, open_send_channels: int, open_receive_channels: int, tasks_waiting_send: int, tasks_waiting_receive: int)
Простой пример с каналом
Вот простой пример использования каналов памяти:
import trio
async def main():
async with trio.open_nursery() as nursery:
# Open a channel:
send_channel, receive_channel = trio.open_memory_channel(0)
# Start a producer and a consumer, passing one end of the channel to
# each of them:
nursery.start_soon(producer, send_channel)
nursery.start_soon(consumer, receive_channel)
async def producer(send_channel):
# Producer sends 3 messages
for i in range(3):
# The producer sends using 'await send_channel.send(...)'
await send_channel.send(f"message {i}")
async def consumer(receive_channel):
# The consumer uses an 'async for' loop to receive the values:
async for value in receive_channel:
print(f"got value {value!r}")
trio.run(main) Если вы запустите это, вы увидите:
got value "message 0" got value "message 1" got value "message 2"
А затем он будет зависать вечно. (Используйте Ctrl+C, чтобы выйти.)
Чистая остановка с каналами
Конечно, нам обычно не нравится, когда программы зависают. Что случилось? Проблема в том, что производитель отправил 3 сообщения, а затем завершил работу, но потребитель не может определить, что производитель исчез: для него всё равно, что ещё одно сообщение может появиться в любой момент. Поэтому он зависает навсегда, ожидая четвёртого сообщения.
Вот новая версия, которая исправляет это: она производит тот же вывод, что и предыдущая версия, а затем выходит без ошибок. Единственное изменение — добавление блоков async with внутри производителя и потребителя:
import trio
async def main():
async with trio.open_nursery() as nursery:
send_channel, receive_channel = trio.open_memory_channel(0)
nursery.start_soon(producer, send_channel)
nursery.start_soon(consumer, receive_channel)
async def producer(send_channel):
async with send_channel:
for i in range(3):
await send_channel.send(f"message {i}")
async def consumer(receive_channel):
async with receive_channel:
async for value in receive_channel:
print(f"got value {value!r}")
trio.run(main) Самое важное здесь — async with производителя. Когда производитель завершает работу, это закрывает send_channel, и это сообщает потребителю, что больше сообщений не поступает, поэтому он может корректно завершить цикл async for. Затем программа завершается, поскольку обе задачи завершены.
Мы также добавили async with потребителю. Это не так важно, но может помочь нам поймать ошибки или другие проблемы. Например, предположим, что потребитель завершил работу преждевременно по какой-то причине — возможно, из-за ошибки. Тогда производитель будет отправлять сообщения в никуда и может застрять на неопределённое время. Но если потребитель закрывает свой receive_channel, тогда производитель получит BrokenResourceError, чтобы предупредить его о том, что он должен прекратить отправку сообщений, так как никто не слушает.
Если вы хотите увидеть эффект от преждевременного выхода потребителя, попробуйте добавить оператор break в цикл async for — вы должны увидеть BrokenResourceError от производителя.
Управление несколькими производителями и/или несколькими потребителями
У вас также может быть несколько производителей и несколько потребителей, все использующие один и тот же канал. Однако это немного усложняет остановку.
Например, рассмотрим это наивное расширение нашего предыдущего примера, теперь с двумя производителями и двумя потребителями:
# This example usually crashes!
import trio
import random
async def main():
async with trio.open_nursery() as nursery:
send_channel, receive_channel = trio.open_memory_channel(0)
# Start two producers
nursery.start_soon(producer, "A", send_channel)
nursery.start_soon(producer, "B", send_channel)
# And two consumers
nursery.start_soon(consumer, "X", receive_channel)
nursery.start_soon(consumer, "Y", receive_channel)
async def producer(name, send_channel):
async with send_channel:
for i in range(3):
await send_channel.send(f"{i} from producer {name}")
# Random sleeps help trigger the problem more reliably
await trio.sleep(random.random())
async def consumer(name, receive_channel):
async with receive_channel:
async for value in receive_channel:
print(f"consumer {name} got value {value!r}")
# Random sleeps help trigger the problem more reliably
await trio.sleep(random.random())
trio.run(main) Два производителя, A и B, отправляют по 3 сообщения. Затем они случайным образом распределяются между двумя потребителями, X и Y. Итак, мы надеемся увидеть вывод примерно такой:
consumer Y got value '0 from producer B' consumer X got value '0 from producer A' consumer Y got value '1 from producer A' consumer Y got value '1 from producer B' consumer X got value '2 from producer B' consumer X got value '2 from producer A'
Однако, во многих случаях этого не происходит — первая часть вывода нормальная, а затем, когда мы доходим до конца, программа завершается с ClosedResourceError. Если вы запустите программу несколько раз, вы увидите, что иногда в трассировке ошибок показано, что send завершился с ошибкой, а в других случаях — receive, и вы даже можете обнаружить, что в некоторых случаях программа вообще не завершается с ошибкой.
Вот что происходит: предположим, что производитель A завершил работу первым. Он выходит, и его блок async with закрывает send_channel. Но подождите! Производитель B всё ещё использовал этот send_channel… поэтому в следующий раз, когда B вызывает send, он получает ClosedResourceError.
Иногда, если повезёт, два производителя могут завершить работу одновременно (или достаточно близко), поэтому оба выполняют свой последний send, прежде чем любой из них закроет send_channel.
Но даже если это произойдёт, мы всё ещё не избавились от проблем! После того, как производители завершат работу, два потребителя соревнуются, чтобы первыми заметить, что send_channel закрыт. Предположим, что X выиграет гонку. Он завершит свой цикл async for, а затем завершит блок async with… и закроет receive_channel, в то время как Y всё ещё использует его. Опять же, это вызывает сбой.
Мы могли бы этого избежать, используя сложные отслеживания, чтобы убедиться, что только последний производитель и последний потребитель закрывают конечные точки канала… но это было бы утомительно и хрупко. К счастью, есть лучший способ! Вот исправленная версия нашей программы выше:
import trio
import random
async def main():
async with trio.open_nursery() as nursery:
send_channel, receive_channel = trio.open_memory_channel(0)
async with send_channel, receive_channel:
# Start two producers, giving each its own private clone
nursery.start_soon(producer, "A", send_channel.clone())
nursery.start_soon(producer, "B", send_channel.clone())
# And two consumers, giving each its own private clone
nursery.start_soon(consumer, "X", receive_channel.clone())
nursery.start_soon(consumer, "Y", receive_channel.clone())
async def producer(name, send_channel):
async with send_channel:
for i in range(3):
await send_channel.send(f"{i} from producer {name}")
# Random sleeps help trigger the problem more reliably
await trio.sleep(random.random())
async def consumer(name, receive_channel):
async with receive_channel:
async for value in receive_channel:
print(f"consumer {name} got value {value!r}")
# Random sleeps help trigger the problem more reliably
await trio.sleep(random.random())
trio.run(main) Этот пример демонстрирует использование методов MemorySendChannel.clone и MemoryReceiveChannel.clone. Что они делают — создают копии наших конечных точек, которые ведут себя точно так же, как и оригиналы — за исключением того, что они могут быть закрыты независимо. А основной канал закрывается только после того, как все клоны будут закрыты. Таким образом, это полностью решает нашу проблему с остановкой, и если вы запустите эту программу, вы увидите, как она выведет свои шесть строк вывода и затем корректно завершит работу.
Заметьте небольшой трюк, который мы используем: код в main создаёт объекты-клоны для передачи во все дочерние задачи, а затем закрывает исходные объекты с помощью async with. Другой вариант — передавать клоны во все, кроме одной дочерней задачи, а затем передавать исходный объект в последнюю задачу, например:
# Also works, but is more finicky: send_channel, receive_channel = trio.open_memory_channel(0) nursery.start_soon(producer, "A", send_channel.clone()) nursery.start_soon(producer, "B", send_channel) nursery.start_soon(consumer, "X", receive_channel.clone()) nursery.start_soon(consumer, "Y", receive_channel)
Но это более подвержено ошибкам, особенно если вы используете цикл для запуска производителей/потребителей.
Просто убедитесь, что вы не пишете:
# Broken, will cause program to hang: send_channel, receive_channel = trio.open_memory_channel(0) nursery.start_soon(producer, "A", send_channel.clone()) nursery.start_soon(producer, "B", send_channel.clone()) nursery.start_soon(consumer, "X", receive_channel.clone()) nursery.start_soon(consumer, "Y", receive_channel.clone())
Здесь мы передаём клоны в задачи, но никогда не закрываем оригинальные объекты. Это означает, что у нас есть 3 объекта канала отправки (оригинал + два клона), но мы закрываем только 2 из них, поэтому потребители будут зависать вечно, ожидая закрытия последнего.
Буферизация в каналах
Когда вы вызываете open_memory_channel(), вы должны указать, сколько значений можно буферизовать внутри канала. Если буфер полон, любая задача, которая вызывает send(), остановится и будет ждать, пока другая задача вызовет receive(). Это полезно, потому что это создаёт обратную связь: если производители канала работают быстрее, чем потребители, это заставляет производителей замедлиться.
Вы можете полностью отключить буферизацию, выполнив open_memory_channel(0). В этом случае любая задача, которая вызывает send(), будет ждать, пока другая задача вызовет receive(), и наоборот. Это похоже на работу каналов в классической модели процессов, взаимодействующих через сообщения, и является разумным значением по умолчанию, если вы не уверены, какой размер буфера использовать. (Поэтому мы использовали его в примерах выше.)
В другом крайнем случае, вы можете сделать буфер неограниченным, используя open_memory_channel(math.inf). В этом случае, send() всегда возвращается немедленно. Обычно это плохая идея. Чтобы понять почему, рассмотрим программу, где производитель работает быстрее, чем потребитель:
# Simulate a producer that generates values 10x faster than the
# consumer can handle them.
import trio
import math
async def producer(send_channel):
count = 0
while True:
# Pretend that we have to do some work to create this message, and it
# takes 0.1 seconds:
await trio.sleep(0.1)
await send_channel.send(count)
print("Sent message:", count)
count += 1
async def consumer(receive_channel):
async for value in receive_channel:
print("Received message:", value)
# Pretend that we have to do some work to handle this message, and it
# takes 1 second
await trio.sleep(1)
async def main():
send_channel, receive_channel = trio.open_memory_channel(math.inf)
async with trio.open_nursery() as nursery:
nursery.start_soon(producer, send_channel)
nursery.start_soon(consumer, receive_channel)
trio.run(main) Если вы запустите эту программу, вы увидите вывод примерно такой:
Sent message: 0 Received message: 0 Sent message: 1 Sent message: 2 Sent message: 3 Sent message: 4 Sent message: 5 Sent message: 6 Sent message: 7 Sent message: 8 Sent message: 9 Received message: 1 Sent message: 10 Sent message: 11 Sent message: 12 ...
В среднем производитель отправляет десять сообщений в секунду, но потребитель вызывает receive только один раз в секунду. Это означает, что каждая секунда внутренний буфер канала должен увеличиться, чтобы вместить дополнительные девять элементов. Через минуту в буфере будет ~540 элементов; через час — ~32 400. В конечном итоге программа исчерпает память. И задолго до того, как мы исчерпаем память, наша задержка при обработке отдельных сообщений станет ужасной. Например, на отметке одной минуты производитель отправляет сообщение ~600, но потребитель всё ещё обрабатывает сообщение ~60. Сообщение 600 должно просидеть в канале около ~9 минут, прежде чем потребитель догонит его и обработает.
Теперь попробуйте заменить open_memory_channel(math.inf) на open_memory_channel(0) и запустить её снова. Мы получим вывод примерно такой:
Sent message: 0 Received message: 0 Received message: 1 Sent message: 1 Received message: 2 Sent message: 2 Sent message: 3 Received message: 3 ...
Теперь вызовы send ждут завершения вызовов receive, что заставляет производителя замедлиться, чтобы соответствовать скорости потребителя. (Может показаться странным, что некоторые значения сообщаются как «Получено», прежде чем они сообщаются как «Отправлено»; это происходит потому, что фактические отправка/приём происходят одновременно, поэтому какая строка отобразится первой, случайность.)
Теперь давайте попробуем установить небольшой, но ненулевой размер буфера, например, open_memory_channel(3). Что, по вашему мнению, произойдёт?
Я получаю:
Sent message: 0 Received message: 0 Sent message: 1 Sent message: 2 Sent message: 3 Received message: 1 Sent message: 4 Received message: 2 Sent message: 5 ...
Таким образом, вы можете видеть, что производитель обгоняет на 3 сообщения, а затем останавливается, чтобы ждать: когда потребитель читает сообщение 1, он отправляет сообщение 4, а затем, когда потребитель читает сообщение 2, он отправляет сообщение 5 и так далее. После достижения стабильного состояния эта версия ведёт себя точно так же, как и наша предыдущая версия, где мы установили размер буфера в 0, за исключением того, что она использует немного больше памяти, и каждое сообщение находится в буфере немного дольше перед обработкой (т. е. задержка сообщения выше).
Конечно, реальные производители и потребители обычно более сложны, и в некоторых ситуациях небольшое количество буферизации может улучшить пропускную способность. Но слишком много буферизации приводит к потере памяти и увеличению задержек, поэтому, если вы хотите настроить своё приложение, вы должны экспериментировать, чтобы увидеть, какое значение лучше всего подходит для вас.
Почему мы вообще поддерживаем неограниченные буферы? Хороший вопрос! Несмотря на всё вышесказанное, есть моменты, когда вам действительно нужен неограниченный буфер. Например, рассмотрим веб-паука, который использует канал, чтобы отслеживать все URL-адреса, которые ему ещё нужно обработать. Каждый паук выполняет цикл, где он берёт URL-адрес из канала, запрашивает его, проверяет HTML на наличие исходящих ссылок и добавляет новые URL-адреса в канал. Это создаёт круговой поток, где каждый потребитель также является производителем. В этом случае, если буфер канала заполнится, пауки будут блокироваться, когда они будут пытаться добавить новые URL-адреса в канал, и если все пауки заблокируются, то они не будут брать URL-адреса из канала, поэтому они навсегда застрянут в тупике. Использование неограниченного канала избегает этого, потому что это означает, что send() никогда не блокируется.
Примитивы синхронизации низкого уровня
Лично я считаю, что событий и каналов обычно достаточно для реализации большинства задач, и они приводят к более легко читаемому коду, чем примитивы низкого уровня, обсуждаемые в этом разделе. Но если вам они нужны, они здесь. (Если вы прибегаете к ним, пытаясь реализовать новый примитив синхронизации высокого уровня, то, возможно, стоит ознакомиться с возможностями в trio.lowlevel для более непосредственного доступа к внутренней логике синхронизации Trio. Все классы, обсуждаемые в этом разделе, реализованы поверх общедоступных API в trio.lowlevel; они не имеют специального доступа к внутренним данным Trio.)
-
Объект для управления доступом к ресурсу с ограниченной емкостью.
Иногда необходимо ограничить количество задач, которые могут что-то сделать одновременно. Например, вы можете захотеть использовать несколько потоков для одновременного выполнения нескольких блокирующих операций ввода-вывода… но если вы используете слишком много потоков сразу, то ваша система может перегрузиться, и это действительно замедлит работу. Одним из популярных решений является применение политики, например, «запустить до 40 потоков одновременно, но не больше». Но как реализовать такую политику?
Для этого предназначен объект
CapacityLimiter. Можно представить объектCapacityLimiterкак мешок, который изначально содержит определенное количество токенов:limit = trio.CapacityLimiter(40)
Затем задачи могут подойти и взять токен из мешка:
# Borrow a token: async with limit: # We are holding a token! await perform_expensive_operation() # Exiting the 'async with' block puts the token back into the sackИ, что очень важно, если вы попытаетесь взять токен, а мешок пуст, то вам придется подождать, пока другая задача не завершит свою работу и не вернет свой токен обратно, прежде чем вы сможете его взять и продолжить.
Другой способ представить это: объект
CapacityLimiterпохож на диван с фиксированным количеством мест, и если все места заняты, то вам нужно подождать, пока кто-то не встанет, прежде чем вы сможете сесть.По умолчанию,
trio.to_thread.run_sync()используетCapacityLimiterдля ограничения количества работающих потоков; см.trio.to_thread.current_default_thread_limiterдля получения подробной информации.Если вы знакомы с семафорами, то вы можете представить это как ограниченный семафор, специализированный для одного распространенного случая использования, с дополнительной проверкой ошибок. Для более традиционного семафора см.
Semaphore.Примечание
Не путайте это с алгоритмами «дырявого ведра» или «ведра с токенами», используемыми для ограничения использования пропускной способности в сетях. Основная идея использования токенов для отслеживания ограничения ресурса аналогична, но это очень простой мешок, где токены не создаются или не уничтожаются со временем автоматически; они просто берутся взаймы, а затем возвращаются обратно.
-
Взять токен из мешка, блокируя, если необходимо.
-
RuntimeError – если текущая задача уже держит один из токенов этого мешка.
Raises:
-
await acquire() → None-
Взять токен из мешка без блокировки.
-
WouldBlock – если токены недоступны.
RuntimeError – если текущая задача уже держит один из токенов этого мешка.
Raises:
-
acquire_nowait() → None-
Взять токен из мешка от имени
borrower, блокируя, если необходимо.-
borrower – A
trio.lowlevel.Taskor arbitrary opaque object used to record who is borrowing this token; seeacquire_on_behalf_of_nowait()for details. -
RuntimeError – если задача
borrowerуже держит один из токенов этого мешка.
Parameters:
Raises:
-
await acquire_on_behalf_of(borrower: Task | object) → None-
Взять токен из мешка от имени
borrowerбез блокировки.-
borrower – A
trio.lowlevel.Taskили произвольный непрозрачный объект, используемый для записи того, кто берет этот токен. Это используетсяtrio.to_thread.run_sync()для разрешения потокам «держать токены», с намерением в будущем использовать его для обнаружения тупиков и других полезных вещей -
WouldBlock – если токены недоступны.
RuntimeError – если
borrowerуже держит один из токенов этого мешка.
Parameters:
Raises:
-
acquire_on_behalf_of_nowait(borrower: Task | object) → None-
Количество доступной емкости.
property available_tokens: int | float-
Количество емкости, которое используется в данный момент.
property borrowed_tokens: int-
Возвратить токен в мешок.
-
RuntimeError – если текущая задача не получила один из токенов этого мешка.
Raises:
-
release() → None-
Возвратить токен в мешок от имени
borrower.-
RuntimeError – если заданная задача не получила один из токенов этого мешка.
Raises:
-
release_on_behalf_of(borrower: Task | object) → None-
Возвращает объект, содержащий информацию для отладки.
В настоящее время определены следующие поля:
borrowed_tokens: Количество токенов, которые в данный момент взяты из мешка.total_tokens: Общее количество токенов в мешке. Обычно это значение больше, чемborrowed_tokens, но возможно, что оно меньше, еслиtotal_tokensбыло недавно уменьшено.borrowers: Список всех задач или других сущностей, которые в данный момент держат токен.tasks_waiting: Количество задач, заблокированных на этом объектеCapacityLimiterв методахacquire()илиacquire_on_behalf_of().
statistics() → CapacityLimiterStatistics-
Общая доступная емкость.
Вы можете изменить
total_tokens, присвоив значение этому атрибуту. Если вы увеличите его, то соответствующее количество ожидающих задач будет немедленно разбужено, чтобы взять новые токены. Если вы уменьшите total_tokens ниже количества задач, которые в настоящее время используют ресурс, то все текущие задачи будут разрешены для завершения в обычном режиме, но новые задачи не будут разрешены до тех пор, пока общее количество задач не уменьшится до нового total_tokens.
property total_tokens: int | float -
class trio.CapacityLimiter(total_tokens: int | float)
-
Семафор.
Семафор хранит целое значение, которое можно увеличивать, вызывая
release(), и уменьшать, вызываяacquire()– но значение никогда не должно опускаться ниже нуля. Если значение равно нулю, тоacquire()будет блокироваться до тех пор, пока кто-то не вызоветrelease().Если вам нужен семафор для ограничения числа задач, которые могут одновременно получить доступ к некоторому ресурсу, то вместо этого можно использовать
CapacityLimiter.Интерфейс этого объекта похож на, но отличается от, интерфейса
threading.Semaphore.Объект
Semaphoreможет использоваться как контекстный менеджер async; он блокирует выполнение на входе, но не на выходе.Параметры:
-
Уменьшает значение семафора, блокируя выполнение, если это необходимо, чтобы избежать снижения значения ниже нуля.
await acquire() → None-
Попытка уменьшить значение семафора без блокировки.
-
WouldBlock – если значение равно нулю.
Исключения:
-
acquire_nowait() → None-
Максимально допустимое значение. Может быть None, чтобы указать отсутствие ограничения.
property max_value: int | None-
Увеличивает значение семафора, возможно разблокировав задачу, заблокированную в
acquire().-
ValueError – если увеличение значения приведет к превышению
max_value.
Исключения:
-
release() → None-
Возвращает объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
tasks_waiting: Количество задач, заблокированных на методеacquire()этого семафора.
statistics() → ParkingLotStatistics-
Текущее значение семафора.
property value: int
class trio.Semaphore(initial_value: int, *, max_value: int | None = None)
-
Классический мьютекс.
Это нереентерабельный мьютекс, защищающий один ресурс. В отличие от
threading.Lock, только владелец мьютекса может его освободить.Объект
Lockможет использоваться как контекстный менеджер async; он блокирует выполнение на входе, но не на выходе.-
Получить мьютекс, блокируя выполнение, если необходимо.
-
BrokenResourceError – если владелец мьютекса завершается без освобождения.
Исключения:
-
await acquire() → None-
Попытка получить мьютекс без блокировки.
-
WouldBlock – если мьютекс уже захвачен.
Исключения:
-
acquire_nowait() → None-
Проверка, захвачен ли мьютекс.
-
True, если мьютекс захвачен, False в противном случае.
Возвращает:
Тип возвращаемого значения:
-
locked() → bool-
Освободить мьютекс.
-
RuntimeError – если вызывающая задача не держит мьютекс.
Исключения:
-
release() → None-
Возвращает объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
locked: boolean, указывающий, захвачен ли мьютекс.owner: задачаtrio.lowlevel.Task, которая в данный момент держит мьютекс, или None, если мьютекс не захвачен.tasks_waiting: Количество задач, заблокированных на методеacquire()этого мьютекса.
statistics() → LockStatistics -
class trio.Lock
-
Вариант
Lock, где задачи гарантированно получают блокировку в строгом порядке «первым пришёл — первым обслужен».Пример полезного использования — реализация чего-то вроде
trio.SSLStreamили HTTP/2-сервера, использующего h2, когда несколько задач взаимодействуют с общей машиной состояний, и в любой момент машина состояний может запросить отправку фрагмента данных по сети. (Например, при использовании h2 просто чтение входящих данных может иногда создать исходящие данные для отправки.) Задача состоит в том, чтобы гарантировать отправку этих фрагментов в правильном порядке, без искажений.Один из вариантов — использовать обычную
Lockи обернуть ею каждое взаимодействие с машиной состояний:# This approach is sometimes workable but often sub-optimal; see below async with lock: state_machine.do_something() if state_machine.has_data_to_send(): await conn.sendall(state_machine.get_data_to_send())Но это может быть проблематично. Если вы используете h2, то обычно чтение входящих данных не требует отправки данных, поэтому мы не хотим заставлять каждую задачу, пытающуюся прочитать данные из сети, ждать потенциально долгого времени, пока
sendallзавершит свою работу. В некоторых ситуациях это может даже привести к тупиковой ситуации, если удалённый узел ожидает, пока вы прочитаете какие-то данные, прежде чем принять отправляемые вами данные.StrictFIFOLockпредоставляет альтернативу. Мы можем переписать наш пример так:# Note: no awaits between when we start using the state machine and # when we block to take the lock! state_machine.do_something() if state_machine.has_data_to_send(): # Notice that we fetch the data to send out of the state machine # *before* sleeping, so that other tasks won't see it. chunk = state_machine.get_data_to_send() async with strict_fifo_lock: await conn.sendall(chunk)Сначала мы выполняем все взаимодействия с машиной состояний в одном кванте планирования (обратите внимание, что в нём нет
await), поэтому это автоматически атомарно относительно других задач. И только если нам нужно отправить данные, мы встаём в очередь на отправку — иStrictFIFOLockгарантирует, что каждая задача отправит свои данные в том же порядке, в котором их сгенерировала машина состояний.В настоящее время
StrictFIFOLockидентичнаLock, но (а) это может быть не всегда верно в будущем, особенно если Trio когда-либо реализует более сложные политики планирования, и (б) приведенный выше код опирается на довольно тонкое свойство своей блокировки. ИспользованиеStrictFIFOLockслужит напоминанием о том, что вы полагаетесь на это свойство.
class trio.StrictFIFOLock
-
Классическая переменная состояния, аналогичная
threading.Condition.Объект
Conditionможет использоваться как асинхронный менеджер контекста для получения базовой блокировки; он блокируется при входе, но не при выходе.-
lock (Lock) — объект блокировки для использования. Если задано, должно быть
trio.Lock. Если None, будет выделена и использована новаяLock.
Параметры:
-
Получение базовой блокировки, при необходимости блокируя выполнение.
-
BrokenResourceError — если владелец базовой блокировки завершается без освобождения.
Исключения:
-
await acquire() → None-
Попытка получить базовую блокировку без блокирования.
-
WouldBlock — если блокировка в настоящее время используется.
Исключения:
-
acquire_nowait() → None-
Проверка, в настоящее ли время используется базовая блокировка.
-
True, если блокировка используется, False — иначе.
Возвращаемое значение:
Тип возвращаемого значения:
-
locked() → bool-
Разбудить одну или несколько задач, заблокированных в
wait().-
n (int) — число задач для разбуживания.
-
RuntimeError — если вызывающая задача не владеет блокировкой.
Параметры:
Исключения:
-
notify(n: int = 1) → None-
Разбудить все задачи, в настоящее время заблокированные в
wait().-
RuntimeError — если вызывающая задача не владеет блокировкой.
Исключения:
-
notify_all() → None-
Освобождение базовой блокировки.
release() → None-
Возврат объекта, содержащего отладочную информацию.
В настоящее время определены следующие поля:
tasks_waiting: количество задач, заблокированных в методеwait()данной переменной состояния.lock_statistics: результат вызова методаLockstatistics().
statistics() → ConditionStatistics-
Ожидание вызова
notify()илиnotify_all()другой задачей.При вызове этого метода вы должны владеть блокировкой. Она освобождается во время ожидания и затем снова приобретается перед пробуждением.
Существует тонкость в том, как этот метод взаимодействует с отменением: при отмене он будет блокироваться, чтобы повторно приобрести блокировку, прежде чем вызывать
Cancelled. Это может привести к тому, что отмена будет менее оперативной, чем ожидалось. Преимущество заключается в том, что это позволяет работать с кодом такого рода:async with condition: await condition.wait()Если мы не повторно приобретали блокировку перед пробуждением, и
wait()была отменена здесь, то мы бы потерпели сбой вcondition.__aexit__при попытке освободить блокировку, которой мы больше не владели.-
RuntimeError — если вызывающая задача не владеет блокировкой.
BrokenResourceError — если владелец блокировки завершается без освобождения при попытке повторного приобретения.
Исключения:
-
await wait() → None -
class trio.Condition(lock: Lock | None = None)
Эти примитивы возвращают объекты статистики, которые можно просматривать.
-
Объект, содержащий отладочную информацию.
В настоящее время определены следующие поля:
borrowed_tokens: Количество токенов, в настоящее время взятых из резерва.total_tokens: Общее количество токенов в резерве. Обычно это значение больше, чемborrowed_tokens, но возможно, что оно меньше, еслиtrio.CapacityLimiter.total_tokensбыло недавно уменьшено.borrowers: Список всех задач или других сущностей, которые в настоящее время удерживают токен.tasks_waiting: Количество задач, заблокированных на этомCapacityLimiter’strio.CapacityLimiter.acquire()илиtrio.CapacityLimiter.acquire_on_behalf_of()методах.
class trio.CapacityLimiterStatistics(borrowed_tokens: int, total_tokens: int | float, borrowers: list[Task | object], tasks_waiting: int)
-
Объект, содержащий отладочную информацию для блокировки.
В настоящее время определены следующие поля:
locked(boolean): указывает, удерживается ли блокировка.owner:trio.lowlevel.Task, которая в настоящее время удерживает блокировку, или None, если блокировка не удерживается.tasks_waiting(int): Количество задач, заблокированных на методеtrio.Lock.acquire()этой блокировки.
class trio.LockStatistics(locked: bool, owner: Task | None, tasks_waiting: int)
-
Объект, содержащий отладочную информацию для условия.
В настоящее время определены следующие поля:
tasks_waiting(int): Количество задач, заблокированных на методеtrio.Condition.wait()этого условия.lock_statistics: Результат вызова методаLocksstatistics().
class trio.ConditionStatistics(tasks_waiting: int, lock_statistics: LockStatistics)
Примечания об асинхронных генераторах
Python 3.6 добавил поддержку асинхронных генераторов, которые могут использовать await, async for и async with между своими операторами yield. Как можно ожидать, для итерирования по ним используется async for. PEP 525 содержит больше подробностей, если они вам нужны.
Например, следующий код является косвенным способом вывода чисел от 0 до 9 с задержкой в 1 секунду перед каждым числом:
async def range_slowly(*args):
"""Like range(), but adds a 1-second sleep before each value."""
for value in range(*args):
await trio.sleep(1)
yield value
async def use_it():
async for value in range_slowly(10):
print(value)
trio.run(use_it) Trio поддерживает асинхронные генераторы с некоторыми оговорками, описанными в этом разделе.
Заключительная обработка
Если вы итерируетесь по асинхронному генератору полностью, как в примере выше, то выполнение асинхронного генератора произойдёт полностью в контексте кода, который его итерирует, и не будет особых неожиданностей.
Однако, если вы прервёте частично завершённый асинхронный генератор, например, break из итерации, ситуация становится сложнее. Объект итератора асинхронного генератора всё ещё жив, ожидая возобновления итерации, чтобы произвести больше значений. В какой-то момент Python поймёт, что все ссылки на итератор были потеряны, и вызовет Trio, чтобы бросить исключение GeneratorExit, чтобы у оставшегося кода очистки внутри генератора был шанс выполниться: finally блоки, __aexit__ обработчики и так далее.
Всё хорошо. К сожалению, Python не гарантирует, когда это произойдёт. Это может произойти сразу после выхода из цикла async for или через произвольное количество времени. Это может даже произойти после завершения всего выполнения Trio! Практически единственная гарантия состоит в том, что это не произойдёт в задаче, которая использовала генератор. Эта задача продолжит выполнение других операций, а очистка асинхронного генератора произойдёт «в какое-то время позже, где-то ещё»: потенциально с другими переменными контекста, не подлежащими таймаутам, и/или после закрытия всех используемых вами детёнышей (nurseries).
Если вам не нравится эта неопределённость, и вы хотите убедиться, что блоки finally и обработчики __aexit__ генератора выполняются сразу после того, как вы им воспользовались, вам нужно обернуть использование генератора чем-то вроде async_generator.aclosing():
# Instead of this:
async for value in my_generator():
if value == 42:
break
# Do this:
async with aclosing(my_generator()) as aiter:
async for value in aiter:
if value == 42:
break Это неудобно, но, к сожалению, Python не предоставляет других надёжных вариантов. Если вы используете aclosing(), то код очистки вашего генератора выполняется в том же контексте, что и остальная часть его итераций, поэтому таймауты, исключения и переменные контекста работают так, как ожидается.
Если вы не используете aclosing(), то Trio сделает всё возможное, но вам придётся столкнуться со следующей семантикой:
Очистка генератора происходит в отменённом контексте, т.е. все блокирующие вызовы, выполненные во время очистки, вызовут
Cancelled. Это сделано для компенсации того факта, что любые таймауты, связанные с первоначальным использованием генератора, давно забыты.Очистка выполняется без доступа к любым переменным контекста, которые могли быть присутствовать при первоначальном использовании генератора.
Если генератор вызывает исключение во время очистки, то оно выводится в логгер
trio.async_generator_errorsи в противном случае игнорируется.Если асинхронный генератор всё ещё жив в конце всего вызова
trio.run(), то он будет очищен после выхода всех задач и перед возвратомtrio.run(). Поскольку «системный детёныш» (system nursery) к этому моменту уже закрыт, Trio не может поддерживать новые вызовыtrio.lowlevel.spawn_system_task().
Если вы планируете запускать свой код на PyPy, чтобы воспользоваться его лучшей производительностью, вы должны знать, что PyPy гораздо чаще, чем CPython, выполняет очистку асинхронных генераторов в момент, значительно позже последнего использования генератора. (Это следствие того, что PyPy не использует подсчёт ссылок для управления памятью.) Для помощи в обнаружении подобных проблем Trio выведет ResourceWarning (игнорируется по умолчанию, но включено при запуске под python -X dev, например) для каждого асинхронного генератора, которому требуется обработка через резервный путь окончательной обработки.
Области отмены и детёныши
Предупреждение
Вы не можете написать оператор
yield, приостанавливающий асинхронный генератор внутриCancelScopeилиNursery, введённых внутри генератора.
То есть это нормально:
async def some_agen():
with trio.move_on_after(1):
await long_operation()
yield "first"
async with trio.open_nursery() as nursery:
nursery.start_soon(task1)
nursery.start_soon(task2)
yield "second"
... Но это не так:
async def some_agen():
with trio.move_on_after(1):
yield "first"
async with trio.open_nursery() as nursery:
yield "second"
... Асинхронные генераторы, декорированные @asynccontextmanager, чтобы служить шаблоном для асинхронного менеджера контекста, не подпадают под это ограничение, поскольку @asynccontextmanager использует их ограниченным образом, который не создаёт проблем.
Нарушение правила, описанного в этом разделе, иногда даёт вам полезное сообщение об ошибке, но Trio не может обнаружить все такие случаи, поэтому иногда вы получите бесполезное TrioInternalError. (А иногда кажется, что всё работает, что, вероятно, худший из всех вариантов, так как тогда вы можете не заметить проблему до незначительной переработки генератора или кода, его итерирующего, или просто не повезёт. Есть предложенное улучшение Python, которое, по крайней мере, обеспечит постоянное отклонение.)
Причина ограничения по областям отмены связана со сложностью обнаружения моментов приостановки и возобновления генератора. Области отмены внутри генератора не должны влиять на код, работающий за пределами генератора, но Trio не участвует в процессе выхода и повторного входа в генератор, поэтому ему трудно поддерживать целостность своей системы отмены. Детёныши (nurseries) используют область отмены внутри, поэтому у них есть все проблемы областей отмены плюс ряд собственных проблем: например, когда генератор приостановлен, что должны делать фоновые задачи? Нет хорошего способа их приостановить, но если они продолжают работать и вызывают исключение, где это исключение может быть повторно вызвано?
Если у вас есть асинхронный генератор, который хочет yield изнутри детёныша или области отмены, лучше всего переработать его в отдельную задачу, которая взаимодействует через каналы памяти. Пакет trio_util предлагает декоратор, который делает это за вас прозрачно.
Для более подробного обсуждения см. проблемы Trio 264 (особенно этот комментарий) и 638.
Потоки (если это необходимо)
В идеальном мире все сторонние библиотеки и низкоуровневые API были бы нативно асинхронными и интегрированными в Trio, и все было бы счастьем и радугой.
К сожалению, такого мира пока не существует. До тех пор, пока он не появится, вам, возможно, придется взаимодействовать с не-Trio API, которые делают грубые вещи, такие как «блокирование».
Признавая эту реальность, Trio предоставляет две полезные утилиты для работы с реальными потоками на уровне операционной системы, threading-модульный стиль. Во-первых, если вы находитесь в Trio, но вам нужно переместить некоторое блокирующее ввод-вывод в поток, есть trio.to_thread.run_sync. А если вы находитесь в потоке и вам нужно связаться обратно с Trio, вы можете использовать trio.from_thread.run() и trio.from_thread.run_sync().
Философия Trio по управлению рабочими потоками
Если вы использовали другие фреймворки ввода-вывода, вы могли столкнуться с понятием «пула потоков», которое чаще всего реализуется как фиксированное множество потоков, которые ожидают задания задач. Они решают две разные проблемы: Во-первых, повторное использование одних и тех же потоков более эффективно, чем запуск и остановка нового потока для каждой задачи; по сути, пул действует как своего рода кэш для потоков в ожидании. И во-вторых, фиксированный размер предотвращает ситуацию, когда одновременно отправляется 100 000 задач, а затем создается 100 000 потоков, система перегружается и выходит из строя. Вместо этого N потоков начинают выполнять первые N задач, в то время как остальные (100 000 - N) задачи находятся в очереди и ожидают своей очереди. Как правило, этого и нужно, и так по умолчанию работает trio.to_thread.run_sync().
Недостатком такого пула потоков является то, что иногда требуется более сложная логика для управления количеством потоков, запущенных одновременно. Например, вы можете захотеть политику типа «максимум 20 потоков в общей сложности, но не более 3 из них могут запускать задачи, связанные с одной учетной записью пользователя», или вы можете захотеть пул, размер которого динамически изменяется со временем в ответ на условия системы.
Даже фиксированная политика может привести к непредвиденным тупикам. Представьте ситуацию, когда у нас есть два разных типа блокирующих задач, которые вы хотите запустить в пуле потоков, тип A и тип B. Тип A довольно прост: он просто выполняется и завершается довольно быстро. Но тип B более сложен: он должен остановиться посередине и подождать завершения другой работы, и эта другая работа включает выполнение задачи типа A. Теперь предположим, что вы отправляете N задач типа B в пул. Они все начинают выполняться, а затем в конечном итоге отправляют одну или несколько задач типа A. Но поскольку каждый поток в нашем пуле уже занят, задачи типа A фактически не начинаются — они просто находятся в очереди, ожидая завершения задач типа B. Но задачи типа B никогда не завершатся, потому что они ожидают завершения задач типа A. Наша система зависла. Идеальное решение этой проблемы — избежать задач типа B в первую очередь — обычно лучше сохранить сложную логику синхронизации в основном потоке Trio. Но если вы не можете этого сделать, вам нужна настраиваемая политика распределения потоков, которая отслеживает отдельные лимиты для разных типов задач и делает невозможным заполнение всеми задачами типа B слотов, необходимых для выполнения задач типа A.
Таким образом, важно иметь возможность изменять политику, управляющую распределением потоков для задач. Но во многих фреймворках это требует реализации нового пула потоков с нуля, что является очень сложной задачей; и если разные типы задач требуют разных политик, вам, возможно, придется создать несколько пулов, что неэффективно, потому что теперь у вас есть два разных кэша потоков, которые не делят ресурсы.
Решение Trio для этой проблемы заключается в разделении управления рабочими потоками на два уровня. Нижний уровень отвечает за взятие задач блокирующего ввода-вывода и организацию их немедленного выполнения в каком-то рабочем потоке. Он занимается решением сложных проблем конкурентности, связанных с управлением потоками, и отвечает за оптимизации, такие как повторное использование потоков, но не имеет политики допуска: если вы дадите ему 100 000 задач, он создаст 100 000 потоков. Верхний уровень отвечает за обеспечение политики, чтобы этого не произошло — но так как он только должен беспокоиться о политике, он может быть намного проще. Фактически, все, что нужно, это аргумент limiter=, передаваемый в trio.to_thread.run_sync(). По умолчанию это глобальный объект CapacityLimiter, который предоставляет классическое поведение пула потоков фиксированного размера. (См. trio.to_thread.current_default_thread_limiter().) Но если вы хотите использовать «отдельные пулы» для задач типа A и задач типа B, то это всего лишь вопрос создания двух отдельных объектов CapacityLimiter и передачи их при выполнении этих задач. Или вот пример определения пользовательской политики, которая учитывает глобальный лимит потоков, а также гарантирует, что ни один отдельный пользователь не может использовать более 3 потоков одновременно:
class CombinedLimiter:
def __init__(self, first, second):
self._first = first
self._second = second
async def acquire_on_behalf_of(self, borrower):
# Acquire both, being careful to clean up properly on error
await self._first.acquire_on_behalf_of(borrower)
try:
await self._second.acquire_on_behalf_of(borrower)
except:
self._first.release_on_behalf_of(borrower)
raise
def release_on_behalf_of(self, borrower):
# Release both, being careful to clean up properly on error
try:
self._second.release_on_behalf_of(borrower)
finally:
self._first.release_on_behalf_of(borrower)
# Use a weak value dictionary, so that we don't waste memory holding
# limiter objects for users who don't have any worker threads running.
USER_LIMITERS = weakref.WeakValueDictionary()
MAX_THREADS_PER_USER = 3
def get_user_limiter(user_id):
try:
return USER_LIMITERS[user_id]
except KeyError:
per_user_limiter = trio.CapacityLimiter(MAX_THREADS_PER_USER)
global_limiter = trio.current_default_thread_limiter()
# IMPORTANT: acquire the per_user_limiter before the global_limiter.
# If we get 100 jobs for a user at the same time, we want
# to only allow 3 of them at a time to even compete for the
# global thread slots.
combined_limiter = CombinedLimiter(per_user_limiter, global_limiter)
USER_LIMITERS[user_id] = combined_limiter
return combined_limiter
async def run_sync_in_thread_for_user(user_id, sync_fn, *args):
combined_limiter = get_user_limiter(user_id)
return await trio.to_thread.run_sync(sync_fn, *args, limiter=combined_limiter) Выполнение блокирующего ввода-вывода в потоках-работниках
-
Преобразование блокирующей операции в асинхронную с использованием потока.
Эти две строки эквивалентны:
sync_fn(*args) await trio.to_thread.run_sync(sync_fn, *args)
за исключением того, что если
sync_fnзанимает много времени, то первая строка заблокирует цикл Trio во время выполнения, а вторая строка позволяет другим задачам Trio продолжать работу, покаsync_fnвыполняется. Это достигается путём переноса вызоваsync_fn(*args)в поток-работник.Изнутри потока-работника вы можете вернуться в Trio, используя функции в
trio.from_thread.-
sync_fn – Произвольная синхронная вызываемая функция.
*args – Позиционные аргументы, передаваемые в sync_fn. Если вам нужны ключевые аргументы, используйте
functools.partial().abandon_on_cancel (bool) – Прерывать ли этот поток при отмене этой операции. См. обсуждение ниже.
thread_name (str) – Необязательная строка для установки имени потока. Всегда устанавливает
threading.Thread.name, но устанавливает имя в ОС только если доступен pthread.h (т.е. большинство установок POSIX). Имена потоков pthread ограничены 15 символами и могут быть прочитаны из/proc/<PID>/task/<SPID>/commили с помощьюps -eT, среди прочих. По умолчанию{sync_fn.__name__|None} from {trio.lowlevel.current_task().name}.-
limiter (None или объект типа CapacityLimiter) –
Объект, используемый для ограничения числа одновременных потоков. Чаще всего это будет
CapacityLimiter, но это может быть что угодно, предоставляющее совместимые методыacquire_on_behalf_of()иrelease_on_behalf_of(). Эта функция вызоветacquire_on_behalf_ofперед запуском потока иrelease_on_behalf_ofпосле завершения потока.Если None (значение по умолчанию), использует значение по умолчанию
CapacityLimiter, как возвращаетсяcurrent_default_thread_limiter().
Параметры:
Обработка отмены: Отмена – сложная проблема здесь, поскольку ни Python, ни операционные системы, на которых он работает, не предоставляют общего механизма для отмены произвольной синхронной функции, выполняемой в потоке. Эта функция всегда проверяет отмену при входе, прежде чем запускать поток. Но после запуска потока есть два способа обработки отмены:
Если
abandon_on_cancel=False, функция игнорирует отмену и продолжает выполняться так же, как если бы мы вызвалиsync_fnсинхронно. Это поведение по умолчанию.-
Если
abandon_on_cancel=True, функция сразу же вызываетCancelled. В этом случае поток продолжает выполняться в фоновом режиме – мы просто отказываемся от него, чтобы он делал то, что собирался сделать, и молча отбрасываем любое возвращаемое значение или ошибки, которые он вызывает. Используйте только в том случае, если вы знаете, что операция безопасна и не имеет побочных эффектов. (Например:trio.socket.getaddrinfo()использует поток сabandon_on_cancel=True, потому что это не повлияет на что-либо, если запущенный поиск имени хоста будет продолжать работу в фоновом режиме.)Блокировка
limiterснимается только после того, как поток действительно завершится – что в случае отмены может произойти некоторое время после возврата этой функции. Еслиtrio.run()завершится до завершения потока, то метод освобождения ограничителя никогда не будет вызван.
Предупреждение
Не следует использовать эту функцию для вызова длительных CPU-связанных функций! Помимо обычных причин, связанных с GIL, почему использование потоков для CPU-связанной работы неэффективно в Python, есть дополнительная проблема: в CPython CPU-связанные потоки имеют тенденцию «вытеснять» потоки ввода-вывода, поэтому использование потоков для CPU-связанной работы, скорее всего, негативно повлияет на основной поток, выполняющий Trio. Если вам нужно это сделать, лучше использовать процесс-работник или, возможно, PyPy (который все еще имеет GIL, но может лучше распределять время ЦП между потоками).
-
То, что возвращает
sync_fn(*args). -
Исключение – то, что вызывает
sync_fn(*args).
Возвращаемое значение:
Исключения:
-
await trio.to_thread.run_sync(sync_fn: Callable[[Unpack[Ts]], RetT], *args: Unpack[Ts], thread_name: str | None = None, abandon_on_cancel: bool = False, limiter: CapacityLimiter | None = None) → RetT
-
Получение значения по умолчанию
CapacityLimiter, используемого функциейtrio.to_thread.run_sync.Наиболее распространённая причина для вызова этой функции – это изменение значения атрибута
total_tokens.
trio.to_thread.current_default_thread_limiter() → CapacityLimiter
Возвращение в поток Trio из другого потока
-
Выполняет заданную асинхронную функцию в родительском потоке Trio, ожидая её завершения.
-
То, что вернёт
afn(*args).
Возвращает:
Возвращает или поднимает исключение, в зависимости от того, что вернёт или поднимет заданная функция. Также может поднимать собственные исключения:
-
RunFinishedError – если соответствующий вызов
trio.run()уже завершён, или если выполнение достигло финальной фазы очистки и больше не может запускать новые системные задачи.Cancelled – Если исходный вызов
trio.to_thread.run_sync()отменён (если trio_token равно None) или вызовtrio.run()завершился (если trio_token не равно None) во время выполненияafn(*args), то afn, вероятно, подниметtrio.Cancelled.RuntimeError – если вы пытаетесь вызвать это изнутри потока Trio, что приведёт к тупику, или если не был предоставлен
trio_token, и мы не можем определить его из контекста.TypeError – если
afnне является асинхронной функцией.
Возможные исключения:
Поиск TrioToken: Есть два способа указать, какой цикл
trio.runнужно повторно войти:Запустите этот поток из
trio.to_thread.run_sync. Trio автоматически захватывает соответствующий токен Trio и использует его для повторного входа в ту же задачу Trio.Передайте ключевой аргумент,
trio_token, определяющий конкретный циклtrio.runдля повторного входа. Это полезно в случае «внешнего» потока, созданного с помощью другого фреймворка, и вам всё равно нужно войти в Trio, или если вы хотите использовать новую системную задачу для вызоваafn, возможно, чтобы избежать контекста отмены соответствующей задачиtrio.to_thread.run_sync. Этот токен можно получить изtrio.lowlevel.current_trio_token().
-
trio.from_thread.run(afn: Callable[[Unpack[Ts]], Awaitable[RetT]], *args: Unpack[Ts], trio_token: TrioToken | None = None) → RetT
-
Выполняет заданную синхронную функцию в родительском потоке Trio, ожидая её завершения.
-
То, что вернёт
fn(*args).
Возвращает:
Возвращает или поднимает исключение, в зависимости от того, что вернёт или поднимет заданная функция. Также может поднимать собственные исключения:
-
RunFinishedError – если соответствующий вызов
trio.runуже завершён.RuntimeError – если вы пытаетесь вызвать это изнутри потока Trio, что приведёт к тупику, или если не был предоставлен
trio_token, и мы не можем определить его из контекста.TypeError – если
fn— асинхронная функция.
Возможные исключения:
Поиск TrioToken: Есть два способа указать, какой цикл
trio.runнужно повторно войти:Запустите этот поток из
trio.to_thread.run_sync. Trio автоматически захватывает соответствующий токен Trio и использует его, когда нужно повторно войти в Trio.Передайте ключевой аргумент,
trio_token, определяющий конкретный циклtrio.runдля повторного входа. Это полезно в случае «внешнего» потока, созданного с помощью другого фреймворка, и вам всё равно нужно войти в Trio, или если вы хотите использовать новую системную задачу для вызоваfn, возможно, чтобы избежать контекста отмены соответствующей задачиtrio.to_thread.run_sync.
-
trio.from_thread.run_sync(fn: Callable[[Unpack[Ts]], RetT], *args: Unpack[Ts], trio_token: TrioToken | None = None) → RetT
Это будет, вероятно, понятнее с примером. Здесь мы демонстрируем, как запустить дочерний поток, а затем использовать канал памяти для отправки сообщений между потоком и задачей Trio:
import trio
def thread_fn(receive_from_trio, send_to_trio):
while True:
# Since we're in a thread, we can't call methods on Trio
# objects directly -- so we use trio.from_thread to call them.
try:
request = trio.from_thread.run(receive_from_trio.receive)
except trio.EndOfChannel:
trio.from_thread.run(send_to_trio.aclose)
return
else:
response = request + 1
trio.from_thread.run(send_to_trio.send, response)
async def main():
send_to_thread, receive_from_trio = trio.open_memory_channel(0)
send_to_trio, receive_from_thread = trio.open_memory_channel(0)
async with trio.open_nursery() as nursery:
# In a background thread, run:
# thread_fn(receive_from_trio, send_to_trio)
nursery.start_soon(
trio.to_thread.run_sync, thread_fn, receive_from_trio, send_to_trio
)
# prints "1"
await send_to_thread.send(0)
print(await receive_from_thread.receive())
# prints "2"
await send_to_thread.send(1)
print(await receive_from_thread.receive())
# When we close the channel, it signals the thread to exit.
await send_to_thread.aclose()
# When we exit the nursery, it waits for the background thread to
# exit.
trio.run(main)
Примечание
Функции
from_thread.run*повторно используют хост-задачу, которая вызвалаtrio.to_thread.run_sync()для выполнения вашей предоставленной функции, если вы используете стандартныеabandon_on_cancel=False, так Trio может быть уверен, что задача останется для выполнения работы. Если вы передаётеabandon_on_cancel=Trueв самом начале или предоставляетеTrioTokenпри обращении в Trio, ваши функции будут выполняться в новой системной задаче. Поэтому значенияcurrent_task(),current_effective_deadline()и другие специфичные для дерева задач могут отличаться в зависимости от значений ключевых аргументов.
Вы также можете использовать trio.from_thread.check_cancelled() для проверки отмены из потока, запущенного с помощью trio.to_thread.run_sync(). Если вызов run_sync() был отменён, то check_cancelled() поднимет trio.Cancelled(). Это как trio.from_thread.run(trio.sleep, 0), но намного быстрее.
-
Вызвать исключение
trio.Cancelled, если связанная задача Trio перешла в состояние отмены.Применимо только к потокам, созданным с помощью
trio.to_thread.run_sync. Проверяйте, чтобы разрешить потокамabandon_on_cancel=FalseподниматьCancelledв подходящем месте или завершать брошенныеabandon_on_cancel=Trueпотоки раньше, чем это возможно.-
Отмена – Если соответствующий вызов
trio.to_thread.run_syncстолкнулся с попыткой отмены, независимо от значенияabandon_on_cancel, переданного в качестве аргумента.RuntimeError – Если данный поток не создан с помощью
trio.to_thread.run_sync.
Исключения:
Примечание
check_cancelled()проверяет, была ли задача, выполняющаяtrio.to_thread.run_sync(), когда-либо отменена с момента последнего выполнения функцииtrio.from_thread.run()илиtrio.from_thread.run_sync(). Он может вызыватьtrio.Cancelled, даже если произошла отмена, которая позже была скрыта изменениемtrio.CancelScope.shieldмежду отменённымCancelScopeиtrio.to_thread.run_sync(). Это отличается от поведения обычных контрольных точек Trio, которые вызываютCancelledтолько в том случае, если отмена всё ещё активна при выполнении контрольной точки. Это различие крайне маловероятно, чтобы повлиять на ваше приложение, но мы упоминаем об этом для полноты. -
trio.from_thread.check_cancelled() → None
Потоки и локальное хранилище задач
При работе с потоками вы можете использовать те же contextvars, о которых мы говорили выше, потому что их значения сохраняются.
Это делается путём автоматического копирования контекста contextvars при использовании любого из:
Это означает, что значения переменных контекста доступны даже в потоках-воркерах или при отправке функции для выполнения в основном/родительском потоке Trio с помощью trio.from_thread.run из одного из этих потоков-воркеров.
Но это также означает, что, поскольку контекст не тот же, а копия, если вы set значение переменной контекста внутри одной из этих функций, работающих в потоках, новое значение будет доступно только в этом контексте (который был скопирован). Таким образом, новое значение будет доступно для этой функции и других внутренних/дочерних задач, но значение не будет доступно в родительском потоке.
Если вам нужно изменить значения, которые хранятся в переменных контекста, и вам необходимо произвести эти изменения из дочерних потоков, вы можете вместо этого установить изменяемый объект (например, словарь) в переменной контекста верхнего уровня/родительского потока Trio. Затем в дочерних потоках вместо установки переменной контекста вы можете get тот же объект и изменить его значения. Таким образом, вы сохраните тот же объект в переменной контекста и измените его только в дочерних потоках.
Таким образом, вы можете изменять содержимое объекта в дочерних потоках и по-прежнему получать доступ к новому содержимому в родительском потоке.
Вот пример:
import contextvars
import time
import trio
request_state = contextvars.ContextVar("request_state")
# Blocking function that should be run on a thread
# It could be reading or writing files, communicating with a database
# with a driver not compatible with async / await, etc.
def work_in_thread(msg):
# Only use request_state.get() inside the worker thread
state_value = request_state.get()
current_user_id = state_value["current_user_id"]
time.sleep(3) # this would be some blocking call, like reading a file
print(f"Processed user {current_user_id} with message {msg} in a thread worker")
# Modify/mutate the state object, without setting the entire
# contextvar with request_state.set()
state_value["msg"] = msg
# An example "request handler" that does some work itself and also
# spawns some helper tasks in threads to execute blocking code.
async def handle_request(current_user_id):
# Write to task-local storage:
current_state = {"current_user_id": current_user_id, "msg": ""}
request_state.set(current_state)
# Here the current implicit contextvars context will be automatically copied
# inside the worker thread
await trio.to_thread.run_sync(work_in_thread, f"Hello {current_user_id}")
# Extract the value set inside the thread in the same object stored in a contextvar
new_msg = current_state["msg"]
print(
f"New contextvar value from worker thread for user {current_user_id}: {new_msg}"
)
# Spawn several "request handlers" simultaneously, to simulate a
# busy server handling multiple requests at the same time.
async def main():
async with trio.open_nursery() as nursery:
for i in range(3):
nursery.start_soon(handle_request, i)
trio.run(main) Запуск этого скрипта приведёт к выводу:
Processed user 2 with message Hello 2 in a thread worker Processed user 0 with message Hello 0 in a thread worker Processed user 1 with message Hello 1 in a thread worker New contextvar value from worker thread for user 2: Hello 2 New contextvar value from worker thread for user 1: Hello 1 New contextvar value from worker thread for user 0: Hello 0
Если вы используете contextvars или используете библиотеку, которая их использует, теперь вы знаете, как они взаимодействуют при работе с потоками в Trio.
Однако имейте в виду, что во многих случаях гораздо проще не использовать переменные контекста в собственном коде, а вместо этого передавать значения в аргументах, так как это может быть более явным и проще для осмысления.
Примечание
Контекст автоматически копируется вместо использования одного родительского контекста, потому что один контекст не может использоваться в более чем одном потоке, это не поддерживается
contextvars.
Интерактивная отладка
Когда вы запускаете интерактивную сессию Python для отладки любого асинхронной программы (будь то на основе asyncio, Trio или чего-то ещё), каждое выражение await должно находиться внутри асинхронной функции:
$ python
Python 3.10.6
Type "help", "copyright", "credits" or "license" for more information.
>>> import trio
>>> await trio.sleep(1)
File "<stdin>", line 1
SyntaxError: 'await' outside function
>>> async def main():
... print("hello...")
... await trio.sleep(1)
... print("world!")
...
>>> trio.run(main)
hello...
world! Это может затруднить быстрое итеративное изменение, так как вам нужно переопределять весь тело функции всякий раз, когда вы вносите изменение.
Trio предоставляет изменённую интерактивную консоль, позволяющую вам await на верхнем уровне. Вы можете получить доступ к этой консоли, запустив python -m trio:
$ python -m trio
Trio 0.21.0+dev, Python 3.10.6
Use "await" directly instead of "trio.run()".
Type "help", "copyright", "credits" or "license" for more information.
>>> import trio
>>> print("hello..."); await trio.sleep(1); print("world!")
hello...
world! Если вы пользователь IPython, вы можете использовать функцию autoawait IPython. Это можно включить в оболочке IPython, запустив магическую команду %autoawait trio. Чтобы autoawait было включено всякий раз, когда Trio установлен, вы можете добавить следующее в свои файлы запуска IPython. (например, ~/.ipython/profile_default/startup/10-async.py)
try:
import trio
get_ipython().run_line_magic("autoawait", "trio")
except ImportError:
pass Исключения и предупреждения
-
Выбрасывается блокирующими вызовами, если окружающий контекст был отменён.
Вы должны позволить этому исключению распространиться, чтобы его поймал соответствующий контекст отмены. Чтобы напомнить вам об этом, он наследуется от
BaseException, а не отException, как иKeyboardInterruptиSystemExit. Это означает, что если вы напишете что-то вроде:try: ... except Exception: ...то это не перехватит исключение
Cancelled.Вы не можете самостоятельно вызвать
Cancelled. Попытка сделать это приведёт кTypeError. Используйтеcancel_scope.cancel()вместо этого.Примечание
В США также часто используется написание этого слова с одним «l» — «canceled». Это недавнее и специфичное для США нововведение, и даже в США оба варианта часто используются. Для согласованности с остальным миром и с термином «отмена» (который всегда имеет два «l»), Trio везде использует написание с двумя «l».
exception trio.Cancelled(*args: object, **kwargs: object)
-
Выбрасывается
fail_after()иfail_at(), если таймаут истекает.
exception trio.TooSlowError
-
Выбрасывается функциями
X_nowait, еслиXбы заблокировалась.
exception trio.WouldBlock
-
Выбрасывается при попытке получить данные из
trio.abc.ReceiveChannel, у которого больше нет данных для получения.Это аналогично условию «конец файла», но для каналов.
exception trio.EndOfChannel
-
Выбрасывается, когда задача пытается использовать ресурс, который уже используется другой задачей, что может привести к ошибкам и бессмысленным результатам.
Например, если две задачи пытаются отправить данные через один и тот же сокет одновременно, Trio выбросит
BusyResourceErrorвместо того, чтобы позволить данным быть искажёнными.
exception trio.BusyResourceError
-
Выбрасывается при попытке использовать ресурс после его закрытия.
Обратите внимание, что «закрытие» здесь означает, что ваш код закрыл ресурс, обычно вызвав метод с именем, подобным
closeилиaclose, или завершив работу менеджера контекста. Если проблема возникла где-то ещё — например, из-за сбоя сети или из-за закрытия удалённой стороной соединения — это должно указываться другим классом исключения, таким какBrokenResourceErrorили подклассомOSError.
exception trio.ClosedResourceError
-
Выбрасывается, когда попытка использовать ресурс терпит неудачу из-за внешних обстоятельств.
Например, вы можете получить эту ошибку, если вы пытаетесь отправить данные по потоку, где удалённая сторона уже закрыла соединение.
Вы не получаете эту ошибку, если вы закрыли ресурс — в этом случае вы получаете
ClosedResourceError.Атрибут
__cause__этого исключения часто содержит дополнительную информацию об основной ошибке.
exception trio.BrokenResourceError
-
Выбрасывается
trio.from_thread.runи аналогичными функциями, если соответствующий вызовtrio.run()уже завершён.
exception trio.RunFinishedError
-
Выбрасывается
run(), если мы сталкиваемся с ошибкой в Trio или (возможно) неверным использованием низкоуровневых APItrio.lowlevel.Этого никогда не должно происходить! Если вы получили эту ошибку, пожалуйста, отправьте отчёт об ошибке.
К сожалению, если вы получили эту ошибку, это также означает, что все ставки сброшены — Trio не знает, что происходит, и его обычные инварианты могут быть нарушены. (Например, мы можем «потерять след» задачи. Или потерять след всех задач.) Однако, опять же, это не должно происходить.
exception trio.TrioInternalError
-
Bases:
FutureWarningВыводится предупреждение, если вы используете устаревший функционал Trio.
Как молодой проект, Trio в настоящее время довольно агрессивно относится к устареванию и/или удалению функциональности, которую мы понимаем, была плохой идеей. Если вы используете Trio, вы должны подписаться на вопрос #1, чтобы получить информацию о предстоящих устареваниях и других изменениях обратной совместимости.
Несмотря на название, этот класс в настоящее время наследуется от
FutureWarning, а не отDeprecationWarning, потому что, пока мы находимся в режиме молодости и агрессивности, мы хотим, чтобы эти предупреждения были видимы по умолчанию. Вы можете скрыть их, установив фильтр или с помощью переключателя-W: см. документациюwarningsдля получения подробностей.
exception trio.TrioDeprecationWarning
© 2017 Nathaniel J. Smith
Licensed under the MIT License.
https://trio.readthedocs.io/en/v0.29.0/reference-core.html