asyncdispatch
Этот модуль реализует асинхронное ввод-вывод. Он включает в себя диспетчер, реализацию типа Future, и макрос async, который позволяет писать асинхронный код в синхронном стиле с ключевым словом await.
Диспетчер выполняет роль цикла событий. Необходимо вызывать poll на нём (или функции, которые это делают за вас, такие как waitFor или runForever) для опроса на наличие ожидающих событий. Реализация основана на epoll в Linux, IO Completion Ports в Windows и select в других операционных системах.
Функция poll сама по себе не вернёт никаких событий. Вместо этого будет завершён соответствующий объект Future . Future — это тип, хранящий значение, которое ещё не доступно, но может стать доступным в будущем. Вы можете проверить, завершён ли future, используя функцию finished . Когда future завершён, это означает, что либо хранимое в нём значение теперь доступно, либо в нём содержится ошибка. Последнее происходит, когда операция по завершению future терпит неудачу с исключением. Вы можете отличить эти ситуации с помощью функции failed.
Объекты future также могут хранить процедуру обратного вызова, которая будет вызвана автоматически после завершения future.
Следовательно, future можно рассматривать как реализацию паттерна proactor. В этом паттерне вы запрашиваете действие, и после выполнения этого действия future завершается результатом этого действия. Запросы могут быть сделаны путём вызова соответствующих функций. Например, вызов функции recv создаст запрос на чтение данных из сокета. Future, возвращаемый функцией recv , будет завершён после того, как будет прочитано запрошенное количество данных или произойдёт исключение.
Код для чтения данных из сокета может выглядеть примерно так:
var future = socket.recv(100)
future.addCallback(
proc () =
echo(future.read)
) Все асинхронные функции, возвращающие Future, не будут блокировать выполнение. Однако они не вернут результат сразу. Асинхронная функция будет иметь код, который будет выполнен до того, как будет сделан асинхронный запрос; в большинстве случаев этот код подготавливает запрос.
В приведённом примере функция recv вернёт новый экземпляр Future после запроса на чтение данных из сокета. Этот экземпляр Future завершится после того, как будет прочитано запрошенное количество данных, в данном случае 100 байт. Вторая строка устанавливает обратный вызов для этого future, который будет вызван после завершения future. Всё, что делает обратный вызов, это записывает данные, хранящиеся в future, в stdout. Функция read используется для этого и проверяет, завершился ли future с ошибкой (если да, то просто поднимает ошибку). В противном случае она возвращает значение future.
Асинхронные процедуры
Асинхронные процедуры устраняют сложность работы с обратными вызовами. Они достигают этого, позволяя писать асинхронный код так же, как вы пишете синхронный код.
Асинхронная процедура помечается с помощью директивы {.async.} . Когда процедура помечается директивой {.async.} , она должна иметь тип возвращаемого значения Future[T] или не иметь его вообще. Если вы не указываете тип возвращаемого значения, то предполагается Future[void].
Внутри асинхронных процедур можно использовать await для вызова любых процедур, которые возвращают Future; это включает асинхронные процедуры. Когда процедура "ожидается", асинхронная процедура, в которой она ожидает, приостановит своё выполнение до тех пор, пока future ожидаемой процедуры не завершится. В этот момент асинхронная процедура возобновит своё выполнение. В период приостановки асинхронной процедуры диспетчер будет запускать другие асинхронные процедуры.
Вызов await может использоваться во многих контекстах. Он может использоваться в правой части объявления переменной: var data = await socket.recv(100), в этом случае переменная будет автоматически установлена в значение future. Он может использоваться для ожидания объекта Future и процедуры, возвращающей Future[void]: await socket.send("foobar").
Если ожидаемое future завершится с ошибкой, то await повторно поднимет эту ошибку. Чтобы избежать этого, можно использовать ключевое слово yield вместо await. Следующий раздел демонстрирует различные способы обработки исключений в асинхронных процедурах.
Обработка исключений
Наиболее надёжный способ обработки исключений — использовать yield на future, а затем проверить свойство failed future. Например:
var future = sock.recv(100) yield future if future.failed: # Handle exception
Процедуры async также предоставляют ограниченную поддержку для оператора try.
try:
let data = await sock.recv(100)
echo("Received ", data)
except:
# Handle exception К сожалению, семантика оператора try может быть некорректной, и иногда компиляция может полностью завершиться неудачей. Поэтому лучше использовать первый стиль, когда это возможно.
Отбрасывание future
Future никогда не следует отбрасывать. Это связано с тем, что они могут содержать ошибки. Если результат Future вас не интересует, следует использовать процедуру asyncCheck вместо ключевого слова discard. Однако имейте в виду, что это не ждёт завершения, и для этого следует использовать waitFor.
Примеры
Примеры смотрите в документации модулей, реализующих асинхронный ввод-вывод. Хорошей отправной точкой является модуль asyncnet.
Исследование ожидающих future
Возможна ситуация, когда асинхронная процедура, или точнее Future[T], зависает и никогда не завершается. Это может произойти по разным причинам и вызвать серьёзные утечки памяти. В таких случаях трудно определить, какая процедура зависла.
К счастью, существует механизм, который отслеживает количество каждого ожидающего future. Для его включения необходимо скомпилировать с -d:futureLogging и использовать процедуру getFuturesInProgress для получения списка ожидающих future вместе со стековыми следами на момент их создания.
Вам также может быть полезно использовать этот пакет prometheus, который будет регистрировать ожидающие future в prometheus, что позволит вам анализировать их с помощью графиков.
Ограничения/Ошибки
- Система эффектов (
raises: []) не работает с асинхронными процедурами.
asyncdispatch модуль зависит от модуля asyncmacro для корректной работы. Импорты
- os, tables, strutils, times, heapqueue, options, asyncstreams, options, math, monotimes, asyncfutures, nativesockets, net, deques, winlean, sets, hashes, macros, strutils, asyncfutures, posix
Типы
CompletionData = object fd*: AsyncFD cb*: owned(proc (fd: AsyncFD; bytesTransferred: DWORD; errcode: OSErrorCode) {...}{. closure, gcsafe.}) cell*: ForeignCell- Исходный код Редактировать
PDispatcher = ref object of PDispatcherBase ioPort: Handle handles*: HashSet[AsyncFD]
- Исходный код Редактировать
CustomRef = ref CustomObj
- Исходный код Редактировать
AsyncFD = distinct int
- Исходный код Редактировать
AsyncEvent = ptr AsyncEventImpl
- Исходный код Редактировать
Callback = proc (fd: AsyncFD): bool {...}{.closure, gcsafe.}- Исходный код Редактировать
Процедуры
proc `==`(x: AsyncFD; y: AsyncFD): bool {...}{.borrow.}- Исходный код Редактировать
proc newDispatcher(): owned PDispatcher {...}{.raises: [], tags: [].}- Создаёт новый экземпляр Dispatcher. Исходный код Редактировать
proc setGlobalDispatcher(disp: sink PDispatcher) {...}{.raises: [Exception], tags: [RootEffect].}- Исходный код Редактировать
proc getGlobalDispatcher(): PDispatcher {...}{.raises: [Exception], tags: [RootEffect].}- Исходный код Редактировать
proc getIoHandler(disp: PDispatcher): Handle {...}{.raises: [], tags: [].}- Возвращает дескриптор порта завершения ввода-вывода (Windows) или селектор (Unix) для указанного диспетчера. Исходный код Редактировать
proc register(fd: AsyncFD) {...}{.raises: [Exception, OSError], tags: [RootEffect].}- Регистрирует
fdв диспетчере. Исходный код Редактировать proc hasPendingOperations(): bool {...}{.raises: [Exception], tags: [RootEffect].}- Возвращает
true, если глобальный диспетчер имеет ожидающие операции. Исходный код Редактировать proc newCustom(): CustomRef {...}{.raises: [], tags: [].}- Исходный код Редактировать
proc recv(socket: AsyncFD; size: int; flags = {SafeDisconn}): owned( Future[string]) {...}{.raises: [Exception, ValueError], tags: [RootEffect].}-
Читает до
sizeбайтов изsocket. Возвращаемое будущее завершится, когда все запрошенные данные будут прочитаны, часть данных будет прочитана или сокет будет отключён, в этом случае будущее завершится со значением"".Предупреждение: Флаг сокета
Исходный код РедактироватьPeekне поддерживается в Windows. proc recvInto(socket: AsyncFD; buf: pointer; size: int; flags = {SafeDisconn}): owned( Future[int]) {...}{.raises: [Exception, ValueError], tags: [RootEffect].}-
Читает до
sizeбайтов изsocketвbuf, размер которого должен быть как минимум таким. Возвращаемое будущее завершится, когда все запрошенные данные будут прочитаны, часть данных будет прочитана или сокет будет отключён, в этом случае будущее завершится со значением0.Предупреждение: Флаг сокета
Исходный код РедактироватьPeekне поддерживается в Windows. proc send(socket: AsyncFD; buf: pointer; size: int; flags = {SafeDisconn}): owned( Future[void]) {...}{.raises: [Exception, ValueError], tags: [RootEffect].}-
Отправляет
sizeбайтов изbufвsocket. Возвращаемое будущее завершится, когда все данные будут отправлены.ВНИМАНИЕ: Используйте с осторожностью. Если
Исходный код Редактироватьbufссылается на удаляемый сборщиком мусора объект, вы должны использовать вызовы GC_ref/GC_unref, чтобы избежать раннего освобождения буфера. proc sendTo(socket: AsyncFD; data: pointer; size: int; saddr: ptr SockAddr; saddrLen: SockLen; flags = {SafeDisconn}): owned(Future[void]) {...}{. raises: [Exception, ValueError], tags: [RootEffect].}- Отправляет
dataв указанный пункт назначенияsaddr, используя сокетsocket. Возвращаемое будущее завершится, когда все данные будут отправлены. Исходный код Редактировать proc recvFromInto(socket: AsyncFD; data: pointer; size: int; saddr: ptr SockAddr; saddrLen: ptr SockLen; flags = {SafeDisconn}): owned(Future[int]) {...}{. raises: [Exception, ValueError], tags: [RootEffect].}- Получает данные датаграммы от
socketвbuf, размер которого должен быть как минимумsize, адрес отправителя датаграммы будет сохранён вsaddrиsaddrLen. Возвращаемое будущее завершится после получения одной датаграммы и вернёт размер полученного пакета. Исходный код Редактировать proc acceptAddr(socket: AsyncFD; flags = {SafeDisconn}; inheritable = defined(nimInheritHandles)): owned( Future[tuple[address: string, client: AsyncFD]]) {...}{. raises: [Exception, ValueError, OSError, ValueError, Exception], tags: [RootEffect].}-
Принимает новое подключение. Возвращает будущее, содержащее сокет клиента, соответствующий этому подключению, и удалённый адрес клиента. Будущее завершится, когда подключение будет успешно принято.
Результат сокет клиента автоматически регистрируется в диспетчере.
Если
inheritableравно false (по умолчанию), получившийся сокет клиента не будет наследуем дочерними процессами.Вызов
Исходный код Редактироватьacceptможет привести к ошибке, если соединяющийся сокет отключится в течение времени выполненияaccept. Если указан флагSafeDisconn, эта ошибка не будет поднята, а вместо этого будет вызван accept ещё раз. proc setInheritable(fd: AsyncFD; inheritable: bool): bool {...}{.raises: [], tags: [].}-
Управляет возможностью наследования дескриптора файла дочерними процессами. Возвращает
trueпри успехе.Эта процедура не гарантируется для всех платформ. Проверьте доступность с помощью declared().
Исходный код Редактировать proc closeSocket(socket: AsyncFD) {...}{.raises: [Exception], tags: [RootEffect].}- Закрывает сокет и гарантирует его дерегистрацию. Исходный код Редактировать
proc unregister(fd: AsyncFD) {...}{.raises: [Exception], tags: [RootEffect].}- Удаляет регистрацию
fd. Исходный код Редактировать proc contains(disp: PDispatcher; fd: AsyncFD): bool {...}{.raises: [], tags: [].}- Исходный код Редактировать
proc addRead(fd: AsyncFD; cb: Callback) {...}{.raises: [Exception, OSError], tags: [RootEffect].}-
Начинает наблюдение за доступностью файла для чтения и затем вызывает обратный вызов
cb.Это не механизм
pureдля портов завершения Windows (IOCP), поэтому, если вы можете избежать этого, пожалуйста, сделайте это. ИспользуйтеaddReadтолько если это действительно необходимо (основной случай использования - адаптация unix-подобных библиотек для асинхронного использования в Windows).Если вы используете эту функцию, вам не нужно использовать asyncdispatch.recv() или asyncdispatch.accept(), поскольку они используют IOCP, пожалуйста, используйте nativesockets.recv() и nativesockets.accept() вместо них.
Убедитесь, что ваш обратный вызов
Исходный код Редактироватьcbвозвращаетtrue, если вы хотите удалить наблюдение за уведомлениямиread, иfalse, если вы хотите продолжить получение уведомлений. proc addWrite(fd: AsyncFD; cb: Callback) {...}{.raises: [Exception, OSError], tags: [RootEffect].}-
Начинает наблюдение за доступностью файла для записи и затем вызывает обратный вызов
cb.Это не механизм
pureдля портов завершения Windows (IOCP), поэтому, если вы можете избежать этого, пожалуйста, сделайте это. ИспользуйтеaddWriteтолько если это действительно необходимо (основной случай использования - адаптация unix-подобных библиотек для асинхронного использования в Windows).Если вы используете эту функцию, вам не нужно использовать asyncdispatch.send() или asyncdispatch.connect(), поскольку они используют IOCP, пожалуйста, используйте nativesockets.send() и nativesockets.connect() вместо них.
Убедитесь, что ваш обратный вызов
Исходный код Редактироватьcbвозвращаетtrue, если вы хотите удалить наблюдение за уведомлениямиwrite, иfalse, если вы хотите продолжить получение уведомлений. proc addTimer(timeout: int; oneshot: bool; cb: Callback) {...}{. raises: [Exception, OSError], tags: [RootEffect].}-
Регистрирует обратный вызов
cbдля вызова при истечении таймера.Параметры:
-
timeout- значение таймаута в миллисекундах. -
oneshot-
true- генерировать только одно событие таймаута -
false- генерировать события таймаута периодически
-
-
proc addProcess(pid: int; cb: Callback) {...}{.raises: [Exception, OSError], tags: [RootEffect].}- Регистрирует обратный вызов
cbдля вызова при завершении процесса с процессом IDpid. Исходный код Редактировать proc newAsyncEvent(): AsyncEvent {...}{.raises: [OSError], tags: [].}-
Создаёт новый потокобезопасный объект
AsyncEvent.Новый объект
Исходный код РедактироватьAsyncEventне регистрируется автоматически в диспетчере, какAsyncSocket. proc trigger(ev: AsyncEvent) {...}{.raises: [OSError], tags: [].}
- Установить событие
evв состояние сигнализации. Исходный код Редактировать proc unregister(ev: AsyncEvent) {...}{.raises: [Exception, OSError], tags: [RootEffect].}- Отменить регистрацию события
ev. Исходный код Редактировать proc close(ev: AsyncEvent) {...}{.raises: [OSError], tags: [].}- Закрыть событие
ev. Исходный код Редактировать proc addEvent(ev: AsyncEvent; cb: Callback) {...}{.raises: [Exception, OSError], tags: [RootEffect].}- Зарегистрировать обратный вызов
cbдля вызова при сигнализацииev. Исходный код Редактировать proc drain(timeout = 500) {...}{.raises: [Exception, ValueError, OSError], tags: [TimeEffect, RootEffect].}- Ожидает завершения всех событий и обрабатывает их. Выбрасывает
ValueError, если нет ожидающих операций. В отличие отpoll, обрабатывает столько событий, сколько доступно, пока не истечёт таймаут. Исходный код Редактировать proc poll(timeout = 500) {...}{.raises: [Exception, ValueError, OSError], tags: [RootEffect, TimeEffect].}- Ожидает завершения событий и обрабатывает их. Выбрасывает
ValueError, если нет ожидающих операций. Выполняет базовый OS epoll или kqueue оператор только один раз. Исходный код Редактировать proc createAsyncNativeSocket(domain: cint; sockType: cint; protocol: cint; inheritable = defined(nimInheritHandles)): AsyncFD {...}{. raises: [OSError, Exception], tags: [RootEffect].}- Исходный код Редактировать
proc createAsyncNativeSocket(domain: Domain = Domain.AF_INET; sockType: SockType = SOCK_STREAM; protocol: Protocol = IPPROTO_TCP; inheritable = defined(nimInheritHandles)): AsyncFD {...}{. raises: [OSError, Exception], tags: [RootEffect].}- Исходный код Редактировать
proc dial(address: string; port: Port; protocol: Protocol = IPPROTO_TCP): owned( Future[AsyncFD]) {...}{.raises: [OSError, ValueError, Exception], tags: [RootEffect].}- Устанавливает подключение к указанной паре
address:portчерез указанный протокол. Процедура перебирает возможные разрешенияaddressдо получения успеха, что означает совместимость с IPv4 и IPv6. Возвращает асинхронный дескриптор файла, зарегистрированный в диспетчере текущей нити, готовый для отправки или получения данных. Исходный код Редактировать proc connect(socket: AsyncFD; address: string; port: Port; domain = Domain.AF_INET): owned(Future[void]) {...}{. raises: [OSError, IOError, ValueError, Exception], tags: [RootEffect].}- Исходный код Редактировать
proc sleepAsync(ms: int | float): owned(Future[void])
- Приостанавливает выполнение текущей асинхронной процедуры на следующие
msмиллисекунд. Исходный код Редактировать proc withTimeout[T](fut: Future[T]; timeout: int): owned(Future[bool])
-
Возвращает будущее, которое завершится, когда
futзавершится или по истеченииtimeoutмиллисекунд.Если
Исходный код Редактироватьfutзавершится раньше, возвращаемое будущее будет содержать true, иначе, если истечётtimeoutмиллисекунд, возвращаемое будущее будет содержать false. proc accept(socket: AsyncFD; flags = {SafeDisconn}; inheritable = defined(nimInheritHandles)): owned(Future[AsyncFD]) {...}{. raises: [Exception, ValueError, OSError], tags: [RootEffect].}-
Принимает новое подключение. Возвращает будущее, содержащее сокет клиента, соответствующий этому подключению.
Если
inheritableравно false (по умолчанию), результирующий сокет клиента не будет унаследован дочерними процессами.Будущее завершится, когда подключение будет успешно принято.
Исходный код Редактировать proc send(socket: AsyncFD; data: string; flags = {SafeDisconn}): owned( Future[void]) {...}{.raises: [Exception, ValueError], tags: [RootEffect].}- Отправляет
dataвsocket. Возвращаемое будущее завершится, когда все данные будут отправлены. Исходный код Редактировать proc readAll(future: FutureStream[string]): owned(Future[string]) {...}{. raises: [Exception, ValueError], tags: [RootEffect].}- Возвращает будущее, которое завершится, когда все строковые данные из указанного потока будущего будут получены. Исходный код Редактировать
proc callSoon(cbproc: proc () {...}{.gcsafe.}) {...}{.gcsafe, raises: [Exception], tags: [RootEffect].}- Расписание
cbprocдля вызова как можно скорее. Обратный вызов вызывается, когда управление возвращается к циклу событий. Исходный код Редактировать proc runForever() {...}{.raises: [Exception, ValueError, OSError], tags: [RootEffect, TimeEffect].}- Запускает цикл опроса глобального диспетчера, который никогда не заканчивается. Исходный код Редактировать
proc waitFor[T](fut: Future[T]): T
- Блокирует текущую нить, пока указанное будущее не завершится. Исходный код Редактировать
proc activeDescriptors(): int {...}{.inline, raises: [], tags: [].}- Возвращает текущее количество активных дескрипторов файлов для текущего цикла событий. Это операция с низкой стоимостью, не включающая системный вызов. Исходный код Редактировать
proc maxDescriptors(): int {...}{.raises: OSError, tags: [].}- Возвращает максимальное количество активных дескрипторов файлов для текущего процесса. Это включает системный вызов. В настоящее время
maxDescriptorsподдерживается на следующих ОС: Windows, Linux, OSX, BSD. Исходный код Редактировать
Макросы
macro async(prc: untyped): untyped
- Макрос, который обрабатывает асинхронные процедуры в соответствующие итераторы и операторы yield. Исходный код Редактировать
macro multisync(prc: untyped): untyped
-
Макрос, который обрабатывает асинхронные процедуры как в асинхронные, так и в синхронные процедуры.
Сгенерированные асинхронные процедуры используют макрос
Исходный код Редактироватьasync, в то время как сгенерированные синхронные процедуры просто удаляют вызовыawait.
Шаблоны
template await(f: typed): untyped {...}{.used.}- Исходный код Редактировать
template await[T](f: Future[T]): auto {...}{.used.}- Исходный код Редактировать
Экспорт
- Порт, ФлагСокета, и, добавитьОбратныйВызов, асинхроннаяПроверка, или, читать, ошибка, установитьОбратныйВызовСразу, добавитьОбратныйВызов, очистить, очиститьОбратныеВызовы, новыйFutureVar, mget, Future, с ошибкой, $, обратныйВызов=, завершить, обратныйВызов=, NimAsyncContinueSuffix, FutureBase, все, завершить, FutureError, получитьОбратныйВызовСразу, FutureVar, включеноВедениеЖурналаFuture, завершить, ошибкаЧтения, завершить, новыйFuture, закончен, длина, обратныйВызов=, ошибка, новыйFutureStream, закончен, записать, завершить, FutureStream, читать, с ошибкой
© 2006–2021 Andreas Rumpf
Licensed under the MIT License.
https://nim-lang.org/docs/asyncdispatch.html