std/asyncdispatch
Исходный кодРедактироватьЭтот модуль реализует асинхронный ввод-вывод. Он включает в себя диспетчер, реализацию типа Future и макрос async, который позволяет писать асинхронный код в синхронном стиле с помощью ключевого слова await.
Диспетчер действует как своего рода цикл обработки событий. Вы должны вызвать poll на нем (или функцию, которая это делает за вас, например, waitFor или runForever), чтобы проверить наличие ожидающих событий. Подлежащая реализация основана на epoll в Linux, IO Completion Ports в Windows и select в других операционных системах.
Функция poll сама по себе не вернет никаких событий. Вместо этого будет завершен соответствующий объект Future. Future — это тип, который хранит значение, которое еще недоступно, но может стать доступным в будущем. Вы можете проверить, завершилось ли будущее, используя функцию finished. Когда будущее завершается, это означает, что либо значение, которое оно хранит, теперь доступно, либо оно содержит ошибку. Последний случай возникает, когда операция завершения будущего завершается с исключением. Вы можете отличить эти два случая с помощью функции failed.
Объекты Future также могут хранить процедуру обратного вызова, которая будет вызываться автоматически после завершения Future.
Таким образом, Future можно рассматривать как реализацию паттерна проактора. В этом паттерне вы запрашиваете выполнение действия, и после того, как это действие выполнено, Future завершается результатом этого действия. Запросы можно делать, вызывая соответствующие функции. Например, вызов функции recv создаст запрос на чтение определенного количества данных из сокета. Future, который возвращает функция recv, завершится, когда запрошенное количество данных будет считано или произойдет исключение.
Код для чтения данных из сокета может выглядеть примерно так:
var future = socket.recv(100)
future.addCallback(
proc () =
echo(future.read)
) Все асинхронные функции, возвращающие Future, не будут блокировать. Однако они не вернут результат немедленно. У асинхронной функции будет код, который будет выполнен до создания асинхронного запроса, в большинстве случаев этот код настраивает запрос.
В приведенном выше примере функция recv вернет новый экземпляр Future, как только будет сделан запрос на чтение данных из сокета. Этот экземпляр Future завершится, когда запрошенное количество данных будет считано, в данном случае 100 байт. На второй строке устанавливается обратный вызов на это будущее, который будет вызван после завершения будущего. Всё, что делает обратный вызов, это записывает данные, хранящиеся в будущем, в stdout. Для этого используется функция read, которая проверяет, завершилось ли будущее с ошибкой (если да, она просто поднимет ошибку), если же ошибки нет, она возвращает значение будущего.
Асинхронные процедуры
Асинхронные процедуры устраняют сложность работы с обратными вызовами. Они делают это, позволяя вам писать асинхронный код так же, как вы пишете синхронный код.
Асинхронная процедура отмечается с помощью директивы {.async.}. При маркировке процедуры директивой {.async.} она должна иметь тип возвращаемого значения Future[T] или не иметь возвращаемого типа вообще. Если вы не укажете тип возвращаемого значения, то используется тип Future[void].
Внутри асинхронных процедур можно использовать await для вызова любых процедур, возвращающих Future; это включает в себя асинхронные процедуры. Когда процедура «ожидается», асинхронная процедура, в которой она ожидается, приостановит свое выполнение, пока Future ожидаемой процедуры не завершится. В этот момент асинхронная процедура возобновит своё выполнение. В период приостановки асинхронной процедуры другие асинхронные процедуры будут выполняться диспетчером.
Вызов await может быть использован во многих контекстах. Он может быть использован в правой части объявления переменной: var data = await socket.recv(100), в этом случае переменная будет автоматически установлена в значение будущего. Он может быть использован для ожидания объекта Future, и он может быть использован для ожидания процедуры, возвращающей Future[void]: await socket.send("foobar").
Если ожидаемое будущее завершится с ошибкой, то await повторно поднимет эту ошибку. Чтобы избежать этого, вы можете использовать ключевое слово yield вместо await. Следующий раздел демонстрирует различные способы обработки исключений в асинхронных процедурах.
var int. Вместо них следует использовать ссылки, такие как ref int.Обработка исключений
Вы можете обрабатывать исключения так же, как и в обычном коде Nim; с помощью оператора try:
try:
let data = await sock.recv(100)
echo("Received ", data)
except:
# Handle exception Альтернативный подход к обработке исключений заключается в использовании yield на будущем, а затем проверке свойства failed будущего. Например:
var future = sock.recv(100) yield future if future.failed: # Handle exception
Отбрасывание будущих значений
Будущие значения никогда не должны отбрасываться напрямую, так как они могут содержать ошибки. Если вам не нужно значение Future, используйте процедуру asyncCheck вместо ключевого слова discard. Обратите внимание, что это не ожидает завершения, и вы должны использовать waitFor или await для этой цели.
await также проверяет, завершается ли будущее с ошибкой, поэтому вы можете безопасно отбросить его результат.Обработка будущих значений
Существует множество различных операций, применимых к будущему значению. Три основные высокоуровневые операции — asyncCheck, waitFor и await.
-
asyncCheck: Поднимает исключение, если будущее завершилось с ошибкой. Она не ожидает завершения будущего и не возвращает результат будущего. -
waitFor: Проверяет цикл событий и блокирует текущую нить до завершения будущего. Это часто используется для вызова асинхронной процедуры из синхронного контекста и никогда не должно использоваться в процедуреasync. -
await: Приостанавливает выполнение текущей асинхронной процедуры до завершения будущего. Пока текущая процедура приостановлена, другие асинхронные процедуры продолжают выполняться. Должна использоваться вместоwaitForв асинхронной процедуре.
Вот удобная сводная таблица, показывающая их основные различия:
| Процедура | Контекст | Блокировка |
|---|---|---|
asyncCheck |
не-асинхронный и асинхронный | неблокирующая |
waitFor |
не-асинхронный | блокирует текущую нить |
await |
асинхронный | приостанавливает текущую процедуру |
Примеры
Для примеров обратитесь к документации по модулям, реализующим асинхронный ввод-вывод. Хорошим началом является модуль asyncnet.
Исследование ожидающих будущих значений
Возможна ситуация, когда асинхронная процедура, или точнее Future[T], застревает и никогда не завершается. Это может произойти по различным причинам и может привести к серьезным утечкам памяти. Когда это происходит, трудно определить застрявшую процедуру.
К счастью, существует механизм, который отслеживает количество каждого ожидающего будущего. Всё что вам нужно сделать, чтобы его включить, это скомпилировать с -d:futureLogging и использовать процедуру getFuturesInProgress для получения списка ожидающих будущих значений вместе с трассировками стека на момент их создания.
Вам также может быть полезно использовать этот пакет prometheus, который будет регистрировать ожидающие Future в prometheus, что позволит вам проанализировать их с помощью удобной диаграммы.
Ограничения/Ошибки
- Система эффектов (
raises: []) не работает с асинхронными процедурами. - Изменяемые параметры не поддерживаются асинхронными процедурами.
Поддержка нескольких асинхронных бэкендов
Благодаря мощной поддержке макросов Nim позволяет реализовывать async/await в библиотеках с минимальной поддержкой языка — таким образом, существуют несколько библиотек async, включая asyncdispatch и chronos, и в будущем могут быть разработаны другие.
Библиотеки, построенные поверх async/await, могут захотеть поддержать несколько асинхронных бэкендов — лучший способ сделать это — создать отдельные модули для каждого бэкенда, которые можно импортировать рядом.
Альтернативный способ — выбрать бэкенд с помощью глобального флага компиляции — этот метод затрудняет составление приложений, использующих оба бэкенда, как это может произойти с транзитивными зависимостями, но может быть уместен в некоторых случаях — библиотеки, выбирающие этот путь, должны вызвать флаг asyncBackend, позволяя приложениям выбирать бэкенд с -d:asyncBackend=<backend_name>.
Известные async бэкенды включают:
-
-d:asyncBackend=none: полностью отключить поддержкуasync -
-d:asyncBackend=asyncdispatch: https://nim-lang.org/docs/asyncdispatch.html -
-d:asyncBackend=chronos: https://github.com/status-im/nim-chronos/
none может быть использован, когда библиотека поддерживает как синхронный, так и асинхронный API, чтобы отключить последний.
Импорты
- os, таблицы, strutils, времена, heapqueue, options, asyncstreams, математика, monotimes, asyncfutures, nativesockets, сеть, deques, winlean, множества, hashes, asyncmacro
Типы
AsyncEvent = ptr AsyncEventImpl
- Источник Редактировать
AsyncFD = distinct int
- Источник Редактировать
Callback = proc (fd: AsyncFD): bool {.closure, ...gcsafe.}- Источник Редактировать
CompletionData = object fd*: AsyncFD cb*: owned(proc (fd: AsyncFD; bytesTransferred: DWORD; errcode: OSErrorCode) {. closure, ...gcsafe.}) cell*: ForeignCell- Источник Редактировать
CustomRef = ref CustomObj
- Источник Редактировать
PDispatcher = ref object of PDispatcherBase handles*: HashSet[AsyncFD]
- Источник Редактировать
Процедуры
proc `==`(x: AsyncFD; y: AsyncFD): bool {.borrow, ...raises: [], tags: [], forbids: [].}- Исходный код Изменить
proc accept(socket: AsyncFD; flags = {SafeDisconn}; inheritable = defined(nimInheritHandles)): owned(Future[AsyncFD]) {. ...raises: [ValueError, OSError, Exception], tags: [RootEffect], forbids: [].}-
Принимает новое соединение. Возвращает будущее, содержащее сокет клиента, соответствующий этому соединению.
Если
inheritableложно (по умолчанию), результирующий сокет клиента не будет наследуемым для дочерних процессов.Будущее завершится, когда соединение будет успешно принято.
Исходный код Изменить proc acceptAddr(socket: AsyncFD; flags = {SafeDisconn}; inheritable = defined(nimInheritHandles)): owned( Future[tuple[address: string, client: AsyncFD]]) {....gcsafe, raises: [ValueError, OSError, Exception, ValueError, Exception], tags: [RootEffect], forbids: [].}-
Принимает новое соединение. Возвращает будущее, содержащее сокет клиента, соответствующий этому соединению, и удалённый адрес клиента. Будущее завершится, когда соединение будет успешно принято.
Результирующий сокет клиента автоматически регистрируется в диспетчере.
Если
inheritableложно (по умолчанию), результирующий сокет клиента не будет наследуемым для дочерних процессов.Вызов
Исходный код Изменитьacceptможет привести к ошибке, если соединяющийся сокет отключится в течение действияaccept. Если указан флагSafeDisconn, эта ошибка не будет поднята, а вместо этого будет вызван accept ещё раз. proc activeDescriptors(): int {.inline, ...raises: [], tags: [], forbids: [].}- Возвращает текущее количество активных дескрипторов файлов для текущего цикла событий. Это операция с низкой стоимостью, не требующая системного вызова. Исходный код Изменить
proc addEvent(ev: AsyncEvent; cb: Callback) {....raises: [OSError], tags: [], forbids: [].}- Регистрирует обратный вызов
cbдля вызова, когдаevбудет сигнализировать Исходный код Изменить proc addProcess(pid: int; cb: Callback) {....raises: [OSError], tags: [], forbids: [].}- Регистрирует обратный вызов
cbдля вызова, когда процесс с идентификатором процессаpidзавершится. Исходный код Изменить proc addRead(fd: AsyncFD; cb: Callback) {....raises: [OSError], tags: [], forbids: [].}-
Начать наблюдение за дескриптором файла на доступность для чтения и затем вызвать обратный вызов
cb.Это не механизм
pureдля Windows Completion Ports (IOCP), поэтому, если можно избежать этого, пожалуйста, сделайте это. ИспользуйтеaddReadтолько если действительно нужно (основной случай использования - адаптация unix-подобных библиотек для асинхронной работы на Windows).Если вы используете эту функцию, вам не нужно использовать asyncdispatch.recv() или asyncdispatch.accept(), потому что они используют IOCP, пожалуйста, используйте nativesockets.recv() и nativesockets.accept() вместо этого.
Убедитесь, что ваш обратный вызов
Исходный код Изменитьcbвозвращаетtrue, если вы хотите удалить наблюдение за уведомлениямиread, иfalse, если вы хотите продолжить получение уведомлений. proc addTimer(timeout: int; oneshot: bool; cb: Callback) {....raises: [OSError], tags: [], forbids: [].}-
Регистрирует обратный вызов
cbдля вызова при истечении таймера.Параметры:
-
timeout- значение таймаута в миллисекундах. -
oneshot-
true- генерировать только одно событие таймаута -
false- генерировать события таймаута периодически
-
-
proc addWrite(fd: AsyncFD; cb: Callback) {....raises: [OSError], tags: [], forbids: [].}-
Начать наблюдение за дескриптором файла на доступность для записи и затем вызвать обратный вызов
cb.Это не механизм
pureдля Windows Completion Ports (IOCP), поэтому, если можно избежать этого, пожалуйста, сделайте это. ИспользуйтеaddWriteтолько если действительно нужно (основной случай использования - адаптация unix-подобных библиотек для асинхронной работы на Windows).Если вы используете эту функцию, вам не нужно использовать asyncdispatch.send() или asyncdispatch.connect(), потому что они используют IOCP, пожалуйста, используйте nativesockets.send() и nativesockets.connect() вместо этого.
Убедитесь, что ваш обратный вызов
Исходный код Изменитьcbвозвращаетtrue, если вы хотите удалить наблюдение за уведомлениямиwrite, иfalse, если вы хотите продолжить получение уведомлений. proc callSoon(cbproc: proc () {....gcsafe.}) {....gcsafe, raises: [], tags: [], forbids: [].}- Распланировать вызов
cbprocкак можно скорее. Обратный вызов вызывается, когда управление возвращается в цикл событий. Исходный код Изменить proc close(ev: AsyncEvent) {....raises: [OSError], tags: [], forbids: [].}- Закрывает событие
ev. Исходный код Изменить proc closeSocket(socket: AsyncFD) {....raises: [], tags: [], forbids: [].}- Закрывает сокет и гарантирует, что он будет снят с регистрации. Исходный код Изменить
proc connect(socket: AsyncFD; address: string; port: Port; domain = Domain.AF_INET): owned(Future[void]) {. ...raises: [ValueError, OSError, Exception], tags: [RootEffect], forbids: [].}- Исходный код Изменить
proc contains(disp: PDispatcher; fd: AsyncFD): bool {....raises: [], tags: [], forbids: [].}- Исходный код Изменить
proc createAsyncNativeSocket(domain: cint; sockType: cint; protocol: cint; inheritable = defined(nimInheritHandles)): AsyncFD {. ...raises: [OSError], tags: [], forbids: [].}- Исходный код Изменить
proc createAsyncNativeSocket(domain: Domain = Domain.AF_INET; sockType: SockType = SOCK_STREAM; protocol: Protocol = IPPROTO_TCP; inheritable = defined(nimInheritHandles)): AsyncFD {. ...raises: [OSError], tags: [], forbids: [].}- Исходный код Изменить
proc dial(address: string; port: Port; protocol: Protocol = IPPROTO_TCP): owned( Future[AsyncFD]) {....raises: [OSError, ValueError, Exception], tags: [RootEffect], forbids: [].}- Устанавливает соединение с указанной парой
address:portчерез указанный протокол. Процедура перебирает возможные разрешенияaddress, пока не добьётся успеха, что означает бесшовную работу как с IPv4, так и с IPv6. Возвращает асинхронный дескриптор файла, зарегистрированный в диспетчере текущей нити, готовый к отправке или приёму данных. Исходный код Изменить proc drain(timeout = 500) {....raises: [ValueError, Exception, OSError], tags: [TimeEffect, RootEffect], forbids: [].}- Ожидает завершения всех событий и обрабатывает их. Поднимает
ValueError, если нет ожидающих операций. В отличие отpoll, эта функция обрабатывает столько событий, сколько доступно, до истечения таймаута. Исходный код Изменить proc getGlobalDispatcher(): PDispatcher {....raises: [], tags: [], forbids: [].}- Исходный код Изменить
proc getIoHandler(disp: PDispatcher): Handle {....raises: [], tags: [], forbids: [].}- Возвращает основной обработчик ввода-вывода (Windows) или селектор (Unix) для указанного диспетчера. Исходный код Изменить
proc maxDescriptors(): int {....raises: OSError, tags: [], forbids: [].}- Возвращает максимальное количество активных дескрипторов файла для текущего процесса. Это предполагает системный вызов. Пока
maxDescriptorsподдерживается в следующих операционных системах: Windows, Linux, OSX, BSD, Solaris. Источник Изменить proc poll(timeout = 500) {....raises: [ValueError, Exception, OSError], tags: [TimeEffect, RootEffect], forbids: [].}- Ожидает событий завершения и обрабатывает их. Вызывает исключение
ValueError, если нет ожидающих операций. Выполняет базовую операцию ОC epoll или kqueue только один раз. Источник Изменить proc recv(socket: AsyncFD; size: int; flags = {SafeDisconn}): owned( Future[string]) {....raises: [ValueError, Exception], tags: [RootEffect], forbids: [].}- Читает до
sizeбайт изsocket. Возвращаемое будущее завершится, когда будут прочитаны все запрошенные данные, часть данных была прочитана или сокет был отключён, в этом случае будущее завершится со значением"".Предупреждение: Флаг сокетаИсточник ИзменитьPeekне поддерживается в Windows. proc recvFromInto(socket: AsyncFD; data: pointer; size: int; saddr: ptr SockAddr; saddrLen: ptr SockLen; flags = {SafeDisconn}): owned(Future[int]) {. ...raises: [ValueError, Exception], tags: [RootEffect], forbids: [].}- Принимает данные дейтаграммы от
socketвbuf, размер которого должен быть как минимумsize, адрес отправителя дейтаграммы будет сохранён вsaddrиsaddrLen. Возвращаемое будущее завершится, как только будет получена одна дейтаграмма, и вернёт размер полученного пакета. Источник Изменить proc recvInto(socket: AsyncFD; buf: pointer; size: int; flags = {SafeDisconn}): owned( Future[int]) {....raises: [ValueError, Exception], tags: [RootEffect], forbids: [].}- Читает до
sizeбайт изsocketвbuf, размер которого должен быть как минимум таким. Возвращаемое будущее завершится, когда будут прочитаны все запрошенные данные, часть данных была прочитана или сокет был отключён, в этом случае будущее завершится со значением0.Предупреждение: Флаг сокетаИсточник ИзменитьPeekне поддерживается в Windows. proc send(socket: AsyncFD; buf: pointer; size: int; flags = {SafeDisconn}): owned( Future[void]) {....raises: [ValueError, Exception], tags: [RootEffect], forbids: [].}- Отправляет
sizeбайт изbufвsocket. Возвращаемое будущее завершится, когда все данные будут отправлены.Предупреждение: Используйте с осторожностью. ЕслиИсточник Изменитьbufссылается на удалённый объект GC, вы должны использовать вызовы GC_ref/GC_unref, чтобы избежать преждевременного освобождения буфера. proc sendTo(socket: AsyncFD; data: pointer; size: int; saddr: ptr SockAddr; saddrLen: SockLen; flags = {SafeDisconn}): owned(Future[void]) {. ...raises: [ValueError, Exception], tags: [RootEffect], forbids: [].}- Отправляет
dataпо указанному адресу назначенияsaddr, используя сокетsocket. Возвращаемое будущее завершится, когда все данные будут отправлены. Источник Изменить proc setInheritable(fd: AsyncFD; inheritable: bool): bool {....raises: [], tags: [], forbids: [].}-
Управляет тем, может ли дескриптор файла унаследоваться дочерними процессами. Возвращает
trueпри успехе.Эта процедура не гарантируется для всех платформ. Проверьте доступность с помощью declared().
Источник Изменить
proc withTimeout[T](fut: Future[T]; timeout: int): owned(Future[bool])
-
Возвращает будущее, которое завершится, когда
futзавершится или по истеченииtimeoutмиллисекунд.Если
Исходный код Редактироватьfutзавершится первой, возвращаемое будущее будет содержать true, иначе, если сначала истечётtimeoutмиллисекунд, возвращаемое будущее будет содержать false.
© 2006–2024 Andreas Rumpf
Licensed under the MIT License.
https://nim-lang.org/docs/asyncdispatch.html