Spec-Zone.ru › Nim

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. Следующий раздел демонстрирует различные способы обработки исключений в асинхронных процедурах.

Внимание: Процедуры, помеченные {.async.}, не поддерживают изменяемые параметры, такие как 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 hasPendingOperations(): bool {....raises: [], tags: [], forbids: [].}
Возвращает true, если у глобального диспетчера есть ожидающие операции. Источник Изменить
proc maxDescriptors(): int {....raises: OSError, tags: [], forbids: [].}
Возвращает максимальное количество активных дескрипторов файла для текущего процесса. Это предполагает системный вызов. Пока maxDescriptors поддерживается в следующих операционных системах: Windows, Linux, OSX, BSD, Solaris. Источник Изменить
proc newAsyncEvent(): AsyncEvent {....raises: [OSError], tags: [], forbids: [].}

Создаёт новый потокобезопасный объект AsyncEvent.

Новый объект AsyncEvent не регистрируется автоматически в диспетчере, как AsyncSocket.

Источник Изменить
proc newCustom(): CustomRef {....raises: [], tags: [], forbids: [].}
Источник Изменить
proc newDispatcher(): owned PDispatcher {....raises: [], tags: [], forbids: [].}
Создаёт новый экземпляр диспетчера. Источник Изменить
proc poll(timeout = 500) {....raises: [ValueError, Exception, OSError],
                           tags: [TimeEffect, RootEffect], forbids: [].}
Ожидает событий завершения и обрабатывает их. Вызывает исключение ValueError, если нет ожидающих операций. Выполняет базовую операцию ОC epoll или kqueue только один раз. Источник Изменить
proc readAll(future: FutureStream[string]): owned(Future[string]) {.
    ...stackTrace: false, raises: [Exception, ValueError], tags: [RootEffect],
    forbids: [].}
Возвращает будущее, которое завершится, когда будут получены все данные строки из указанного потока будущего. Источник Изменить
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 register(fd: AsyncFD) {....raises: [OSError], tags: [], forbids: [].}
Регистрирует fd в диспетчере. Источник Изменить
proc runForever() {....raises: [ValueError, Exception, OSError],
                    tags: [TimeEffect, RootEffect], forbids: [].}
Начинает бесконечный цикл опроса глобального диспетчера. Источник Изменить
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 send(socket: AsyncFD; data: string; flags = {SafeDisconn}): owned(
    Future[void]) {....raises: [ValueError, Exception], tags: [RootEffect],
                    forbids: [].}
Отправляет data в socket. Возвращаемое будущее завершится, как только все данные будут отправлены. Источник Изменить
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 setGlobalDispatcher(disp: sink PDispatcher) {....raises: [], tags: [],
    forbids: [].}
Источник Изменить
proc setInheritable(fd: AsyncFD; inheritable: bool): bool {....raises: [],
    tags: [], forbids: [].}

Управляет тем, может ли дескриптор файла унаследоваться дочерними процессами. Возвращает true при успехе.

Эта процедура не гарантируется для всех платформ. Проверьте доступность с помощью declared().

Источник Изменить
proc sleepAsync(ms: int | float): owned(Future[void])
Приостанавливает выполнение текущей асинхронной процедуры на следующие ms миллисекунд. Источник Изменить
proc trigger(ev: AsyncEvent) {....raises: [OSError], tags: [], forbids: [].}
Установить событие ev в состояние сигнализации. Источник Изменить
proc unregister(ev: AsyncEvent) {....raises: [OSError], tags: [], forbids: [].}
Отменить регистрацию события ev. Источник Изменить
proc unregister(fd: AsyncFD) {....raises: [], tags: [], forbids: [].}
Отменить регистрацию fd. Источник Изменить
proc waitFor[T](fut: Future[T]): T
Блокирует текущий поток, пока указанное будущее не завершится. Источник Изменить
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

Spec-Zone.ru

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