Spec-Zone.ru › Ruby 3.3

класс Fiber::Scheduler

Родитель:
Объект

Это не существующий класс, а документация интерфейса, которому объект Scheduler должен соответствовать, чтобы использоваться в качестве аргумента для Fiber.scheduler и обрабатывать неблокирующие волокна. См. также раздел «Неблокирующие волокна» в документации класса Fiber для объяснений некоторых концепций.

Ожидается, что поведение и использование планировщика будут следующими:

  • Когда выполнение в неблокирующем Fiber достигает какой-либо блокирующей операции (например, sleep, ожидания процесса или недоступного ввода-вывода), оно вызывает некоторые из методов хуков планировщика, перечисленные ниже.

  • Scheduler каким-то образом регистрирует, на что ожидает текущее волокно, и уступает управление другим волокнам с помощью Fiber.yield (так что волокно будет приостановлено, ожидая завершения ожидания, и другие волокна в том же потоке могут выполнять операции).

  • В конце выполнения текущего потока вызывается метод планировщика scheduler_close.

  • Планировщик запускает цикл ожидания, проверяя все заблокированные волокна (которые он зарегистрировал при вызовах хуков) и возобновляя их, когда ожидаемый ресурс становится доступен (например, ввод-вывод готов или истекло время ожидания сна).

Таким образом, параллельное выполнение будет достигаться прозрачно для каждого отдельного кода волокна.

Реализации Scheduler предоставляются gem'ами, например, Async.

Методы хуков:

  • io_wait, io_read, io_write, io_pread, io_pwrite и io_select, io_close

  • process_wait

  • kernel_sleep

  • timeout_after

  • address_resolve

  • block и unblock

  • (список расширяется по мере того, как Ruby-разработчики создают больше методов с неблокирующими вызовами)

Если не указано иное, реализации хуков обязательны: если они не реализованы, методы, пытающиеся вызвать хук, завершатся ошибкой. Для обеспечения обратной совместимости в будущем хуки будут необязательными (если они не реализованы из-за того, что планировщик создан для более старой версии Ruby, код, которому нужен этот хук, не завершится с ошибкой, а просто будет вести себя блокирующим образом).

Также настоятельно рекомендуется, чтобы планировщик реализовывал метод fiber, которому делегирует Fiber.schedule.

Пример игрушечной реализации планировщика можно найти в коде Ruby в test/fiber/scheduler.rb

Публичные методы экземпляра

address_resolve(hostname) → array_of_strings or nil Show source
VALUE
rb_fiber_scheduler_address_resolve(VALUE scheduler, VALUE hostname)
{
    VALUE arguments[] = {
        hostname
    };

    return rb_check_funcall(scheduler, id_address_resolve, 1, arguments);
}

Вызывается любым методом, который выполняет прямой DNS-запрос. Наиболее заметный метод - Addrinfo.getaddrinfo, но есть и многие другие.

Ожидается, что метод вернет массив строк, соответствующих IP-адресам, к которым разрешается hostname, или nil, если его невозможно разрешить.

Довольно исчерпывающий список всех возможных мест вызова:

  • Addrinfo.getaddrinfo

  • Addrinfo.tcp

  • Addrinfo.udp

  • Addrinfo.ip

  • Addrinfo.new

  • Addrinfo.marshal_load

  • SOCKSSocket.new

  • TCPServer.new

  • TCPSocket.new

  • IPSocket.getaddress

  • TCPSocket.gethostbyname

  • UDPSocket#connect

  • UDPSocket#bind

  • UDPSocket#send

  • Socket.getaddrinfo

  • Socket.gethostbyname

  • Socket.pack_sockaddr_in

  • Socket.sockaddr_in

  • Socket.unpack_sockaddr_in

block(blocker, timeout = nil) Show source
VALUE
rb_fiber_scheduler_block(VALUE scheduler, VALUE blocker, VALUE timeout)
{
    return rb_funcall(scheduler, id_block, 2, blocker, timeout);
}

Вызывается методами, такими как Thread.join, и Mutex, чтобы указать, что текущий Fiber заблокирован до дальнейшего уведомления (например, unblock) или пока не истечет timeout.

blocker - это то, чего мы ждем, только для информации (для отладки и ведения журнала). Гарантий относительно его значения нет.

Ожидается, что вернет логическое значение, указывающее, была ли операция блокировки успешной или нет.

close() Show source
VALUE
rb_fiber_scheduler_close(VALUE scheduler)
{
    VM_ASSERT(ruby_thread_has_gvl_p());

    VALUE result;

    // The reason for calling `scheduler_close` before calling `close` is for
    // legacy schedulers which implement `close` and expect the user to call
    // it. Subsequently, that method would call `Fiber.set_scheduler(nil)`
    // which should call `scheduler_close`. If it were to call `close`, it
    // would create an infinite loop.

    result = rb_check_funcall(scheduler, id_scheduler_close, 0, NULL);
    if (!UNDEF_P(result)) return result;

    result = rb_check_funcall(scheduler, id_close, 0, NULL);
    if (!UNDEF_P(result)) return result;

    return Qnil;
}

Вызывается при выходе текущего потока. Планировщик должен реализовать этот метод, чтобы позволить всем ожидающим волокнам завершить свое выполнение.

Предлагаемый шаблон - реализовать основной цикл обработки событий в методе close.

fiber(&block)

Реализация Fiber.schedule. Ожидается, что метод немедленно запустит заданный блок кода в отдельном неблокирующем волокне и вернет это Fiber.

Минимальная предлагаемая реализация:

def fiber(&block)
  fiber = Fiber.new(blocking: false, &block)
  fiber.resume
  fiber
end
io_pread(io, buffer, from, length, offset) → read length or -errno Show source
VALUE
rb_fiber_scheduler_io_pread(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t length, size_t offset)
{
    VALUE arguments[] = {
        io, buffer, OFFT2NUM(from), SIZET2NUM(length), SIZET2NUM(offset)
    };

    return rb_check_funcall(scheduler, id_io_pread, 5, arguments);
}

Вызывается IO#pread или IO::Buffer#pread для чтения length байтов из io с смещением from в указанный buffer (см. IO::Buffer) в заданном offset.

Этот метод семантически идентичен io_read, но позволяет указать смещение для чтения и часто лучше подходит для асинхронного IO в одном и том же файле.

Метод следует считать экспериментальным.

io_pwrite(io, buffer, from, length, offset) → written length or -errno Show source
VALUE
rb_fiber_scheduler_io_pwrite(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t length, size_t offset)
{
    VALUE arguments[] = {
        io, buffer, OFFT2NUM(from), SIZET2NUM(length), SIZET2NUM(offset)
    };

    return rb_check_funcall(scheduler, id_io_pwrite, 5, arguments);
}

Вызывается IO#pwrite или IO::Buffer#pwrite для записи length байтов в io со смещением from в указанный buffer (см. IO::Buffer) в заданном offset.

Этот метод семантически идентичен io_write, но позволяет указать смещение для записи и часто лучше подходит для асинхронного IO в одном и том же файле.

Метод следует считать экспериментальным.

io_read(io, buffer, length, offset) → read length or -errno Show source
VALUE
rb_fiber_scheduler_io_read(VALUE scheduler, VALUE io, VALUE buffer, size_t length, size_t offset)
{
    VALUE arguments[] = {
        io, buffer, SIZET2NUM(length), SIZET2NUM(offset)
    };

    return rb_check_funcall(scheduler, id_io_read, 4, arguments);
}

Вызывается IO#read или IO#Buffer.read для чтения length байтов из io в указанный buffer (см. IO::Buffer) в заданном offset.

Аргумент length - это «минимальная длина для чтения». Если размер буфера IO составляет 8 Кбайт, но length равен 1024 (1 Кбайт), может быть прочитано до 8 Кбайт, но как минимум 1 Кбайт. Как правило, единственный случай, когда будет прочитано меньше данных, чем length, - это если произошла ошибка при чтении данных.

Указание length равного 0 допустимо и означает попытку чтения хотя бы один раз и возвращение любых доступных данных.

Предлагаемая реализация должна попытаться прочитать из io в неблокирующем режиме и вызвать io_wait, если io не готов (что приведет к передаче управления другим волокнам).

См. IO::Buffer для интерфейса, доступного для возврата данных.

Ожидается, что вернет количество прочитанных байтов или, в случае ошибки, -errno (отрицательное число, соответствующее коду ошибки системы).

Метод следует считать экспериментальным.

io_select(readables, writables, exceptables, timeout) Show source
VALUE rb_fiber_scheduler_io_select(VALUE scheduler, VALUE readables, VALUE writables, VALUE exceptables, VALUE timeout)
{
    VALUE arguments[] = {
        readables, writables, exceptables, timeout
    };

    return rb_fiber_scheduler_io_selectv(scheduler, 4, arguments);
}

Вызывается IO.select для проверки готовности указанных дескрипторов к указанным событиям в течение указанного timeout.

Ожидается, что вернет 3-кортеж из Array готовых IOs.

io_wait(io, events, timeout) Show source
VALUE
rb_fiber_scheduler_io_wait(VALUE scheduler, VALUE io, VALUE events, VALUE timeout)
{
    return rb_funcall(scheduler, id_io_wait, 3, io, events, timeout);
}

Вызывается IO#wait, IO#wait_readable, IO#wait_writable для проверки готовности указанного дескриптора к указанным событиям в течение указанного timeout.

events - это битовая маска IO::READABLE, IO::WRITABLE и IO::PRIORITY.

Предлагаемая реализация должна регистрировать, какое Fiber ожидает какие ресурсы, и немедленно вызывать Fiber.yield для передачи управления другим волокнам. Затем, в методе close, планировщик может отправить все ресурсы I/O волокнам, ожидающим их.

Ожидается, что вернет подмножество событий, которые готовы немедленно.

io_write(io, buffer, length, offset) → записанная длина или -errno Показать исходный код
VALUE
rb_fiber_scheduler_io_write(VALUE scheduler, VALUE io, VALUE buffer, size_t length, size_t offset)
{
    VALUE arguments[] = {
        io, buffer, SIZET2NUM(length), SIZET2NUM(offset)
    };

    return rb_check_funcall(scheduler, id_io_write, 4, arguments);
}

Вызывается методом IO#write или IO::Buffer#write для записи length байтов в io из указанного buffer (см. IO::Buffer) по заданному offset.

Аргумент length — это «минимальная длина для записи». Если размер буфера IO составляет 8 КБ, но указанная length — 1024 (1 КБ), будет записано не более 8 КБ, но не менее 1 КБ. В целом, данные, меньшие чем length, будут записаны только в случае возникновения ошибки при записи.

Указание length в значении 0 допустимо и означает попытку записи как минимум один раз, по возможности, максимального объёма данных.

В предложенной реализации должна быть попытка записи в io в асинхронном режиме и вызов io_wait, если io не готов (что передаст управление другим волокнам).

См. IO::Buffer для интерфейса, позволяющего эффективно получать данные из буфера.

Ожидается возврат количества записанных байтов или, в случае ошибки, -errno (отрицательное число, соответствующее системному коду ошибки).

Метод следует считать экспериментальным.

kernel_sleep(duration = nil) Показать исходный код
VALUE
rb_fiber_scheduler_kernel_sleep(VALUE scheduler, VALUE timeout)
{
    return rb_funcall(scheduler, id_kernel_sleep, 1, timeout);
}

Вызывается методом Kernel#sleep и Mutex#sleep и ожидается, что предоставит реализацию ожидания в асинхронном режиме. Реализация может зарегистрировать текущее волокно в некотором списке «какое волокно ждёт до какого момента», вызвать Fiber.yield для передачи управления, а затем в close возобновить волокна, период ожидания которых истек.

process_wait(pid, flags) Показать исходный код
VALUE
rb_fiber_scheduler_process_wait(VALUE scheduler, rb_pid_t pid, int flags)
{
    VALUE arguments[] = {
        PIDT2NUM(pid), RB_INT2NUM(flags)
    };

    return rb_check_funcall(scheduler, id_process_wait, 2, arguments);
}

Вызывается методом Process::Status.wait для ожидания указанного процесса. Описание аргументов см. в описании этого метода.

Предлагаемая минимальная реализация:

Thread.new do
  Process::Status.wait(pid, flags)
end.value

Этот обработчик является необязательным: если он отсутствует в текущем планировщике, Process::Status.wait будет вести себя как блокирующий метод.

Ожидается возврат экземпляра Process::Status.

timeout_after(duration, exception_class, *exception_arguments, &block) → результат блока Показать исходный код
VALUE
rb_fiber_scheduler_timeout_after(VALUE scheduler, VALUE timeout, VALUE exception, VALUE message)
{
    VALUE arguments[] = {
        timeout, exception, message
    };

    return rb_check_funcall(scheduler, id_timeout_after, 3, arguments);
}

Вызывается методом Timeout.timeout для выполнения данного block в течение заданного duration. Его также может вызывать планировщик или пользовательский код напрямую.

Попытка ограничить время выполнения данного block заданным duration, если это возможно. Когда время выполнения асинхронной операции превышает заданное duration, эта асинхронная операция должна прерываться путём возбуждения указанного exception_class, созданного с данными exception_arguments.

Общие таймауты выполнения часто считаются рискованными. Данная реализация будет прерывать только асинхронные операции. Это сделано по умолчанию, так как ожидается, что асинхронные операции могут завершиться по ряду непредсказуемых причин, поэтому приложения должны быть устойчивыми к обработке этих условий, а также таймаутов.

Однако, в результате этого дизайна, если block не вызывает никаких асинхронных операций, прервать её будет невозможно. Если требуется обеспечить предсказуемые точки для таймаутов, добавьте +sleep(0)+.

Если блок выполняется успешно, его результат будет возвращён.

Исключение обычно возбуждается с помощью Fiber#raise.

unblock(blocker, fiber) Показать исходный код
VALUE
rb_fiber_scheduler_unblock(VALUE scheduler, VALUE blocker, VALUE fiber)
{
    VM_ASSERT(rb_obj_is_fiber(fiber));

    return rb_funcall(scheduler, id_unblock, 2, blocker, fiber);
}

Вызывается для разблокировки Fiber, ранее заблокированного с помощью block (например, Mutex#lock вызывает block, а Mutex#unlock — unblock). Планировщик должен использовать параметр fiber для понимания, какое волокно разблокируется.

blocker — то, что ожидалось, но это только информационное значение (для отладки и ведения журналов), и нет гарантии, что оно будет таким же, как blocker для block.

Ruby Core © 1993–2022 Yukihiro Matsumoto
Licensed under the Ruby License.
Ruby Standard Library © contributors
Licensed under their own licenses.

Spec-Zone.ru

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