Spec-Zone.ru › Ruby 4.0
  1. Fiber::
  2. Scheduler

class Fiber::Scheduler

Родитель:
Object

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

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

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

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

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

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

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

Реализации Scheduler предоставляются гемами, например 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

  • blocking_operation_wait

  • fiber_interrupt

  • yield

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

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

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

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

Открытые методы экземпляра

address_resolve(hostname) → array_of_strings or nil Показать исходный код
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) Показать исходный код
VALUE
rb_fiber_scheduler_block(VALUE scheduler, VALUE blocker, VALUE timeout)
{
    return rb_funcall(scheduler, id_block, 2, blocker, timeout);
}

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

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

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

blocking_operation_wait(blocking_operation) Показать исходный код
VALUE rb_fiber_scheduler_blocking_operation_wait(VALUE scheduler, void* (*function)(void *), void *data, rb_unblock_function_t *unblock_function, void *data2, int flags, struct rb_fiber_scheduler_blocking_operation_state *state)
{
    // Check if scheduler supports blocking_operation_wait before creating the object
    if (!rb_respond_to(scheduler, id_blocking_operation_wait)) {
        return Qundef;
    }

    // Create a new BlockingOperation with the blocking operation
    VALUE blocking_operation = rb_fiber_scheduler_blocking_operation_new(function, data, unblock_function, data2, flags, state);

    VALUE result = rb_funcall(scheduler, id_blocking_operation_wait, 1, blocking_operation);

    // Get the operation data to check if it was executed
    rb_fiber_scheduler_blocking_operation_t *operation = get_blocking_operation(blocking_operation);
    rb_atomic_t current_status = RUBY_ATOMIC_LOAD(operation->status);

    // Invalidate the operation now that we're done with it
    operation->function = NULL;
    operation->state = NULL;
    operation->data = NULL;
    operation->data2 = NULL;
    operation->unblock_function = NULL;

    // If the blocking operation was never executed, return Qundef to signal the caller to use rb_nogvl instead
    if (current_status == RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED) {
        return Qundef;
    }

    return result;
}

Вызывается основными методами Ruby для выполнения блокирующей операции неблокирующим способом. blocking_operation — это непрозрачный объект, инкапсулирующий блокирующую операцию и имеющий метод call без аргументов.

Если планировщик не реализует этот метод или не выполняет блокирующую операцию, Ruby использует реализацию без планировщика.

Минимальная рекомендуемая реализация:

def blocking_operation_wait(blocking_operation)
  Thread.new { blocking_operation.call }.join
end
close () Показать исходный код
VALUE
rb_fiber_scheduler_close(VALUE scheduler)
{
    RUBY_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) Показать исходный код
VALUE
rb_fiber_scheduler_fiber(VALUE scheduler, int argc, VALUE *argv, int kw_splat)
{
    return rb_funcall_passing_block_kw(scheduler, id_fiber_schedule, argc, argv, kw_splat);
}

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

Минимальная рекомендуемая реализация:

def fiber(&block)
  fiber = Fiber.new(blocking: false, &block)
  fiber.resume
  fiber
end
fiber_interrupt(fiber, exception) Показать исходный код
VALUE rb_fiber_scheduler_fiber_interrupt(VALUE scheduler, VALUE fiber, VALUE exception)
{
    VALUE arguments[] = {
        fiber, exception
    };

    VALUE result;
    enum ruby_tag_type state;

    // We must prevent interrupts while invoking the fiber_interrupt method, because otherwise fibers can be left permanently blocked if an interrupt occurs during the execution of user code. See also `rb_fiber_scheduler_unblock`.
    rb_execution_context_t *ec = GET_EC();
    int saved_interrupt_mask = ec->interrupt_mask;
    ec->interrupt_mask |= PENDING_INTERRUPT_MASK;

    EC_PUSH_TAG(ec);
    if ((state = EC_EXEC_TAG()) == TAG_NONE) {
        result = rb_check_funcall(scheduler, id_fiber_interrupt, 2, arguments);
    }
    EC_POP_TAG();

    ec->interrupt_mask = saved_interrupt_mask;

    if (state) {
        EC_JUMP_TAG(ec, state);
    }

    RUBY_VM_CHECK_INTS(ec);

    return result;
}

Вызывается основными методами Ruby, чтобы сообщить планировщику, что заблокированное волокно следует прервать исключением. Например, IO#close использует этот метод, чтобы прерывать волокна, выполняющие блокирующие операции IO.

io_close(fd) Показать исходный код
VALUE
rb_fiber_scheduler_io_close(VALUE scheduler, VALUE io)
{
    VALUE arguments[] = {io};

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

Вызывается основными методами Ruby, чтобы сообщить планировщику о закрытии объекта IO. Обратите внимание, что методу будет передан целочисленный дескриптор закрытого файла, а не сам объект.

io_pread(io, buffer, from, length, offset) → read length or -errno Показать исходный код
VALUE
rb_fiber_scheduler_io_pread(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t length, size_t offset)
{
    if (!rb_respond_to(scheduler, id_io_pread)) {
        return RUBY_Qundef;
    }

    VALUE arguments[] = {
        scheduler, io, buffer, OFFT2NUM(from), SIZET2NUM(length), SIZET2NUM(offset)
    };

    if (rb_respond_to(scheduler, id_fiber_interrupt)) {
        return rb_thread_io_blocking_operation(io, fiber_scheduler_io_pread, (VALUE)&arguments);
    } else {
        return fiber_scheduler_io_pread((VALUE)&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 Показать исходный код
VALUE
rb_fiber_scheduler_io_pwrite(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t length, size_t offset)
{


    if (!rb_respond_to(scheduler, id_io_pwrite)) {
        return RUBY_Qundef;
    }

    VALUE arguments[] = {
        scheduler, io, buffer, OFFT2NUM(from), SIZET2NUM(length), SIZET2NUM(offset)
    };

    if (rb_respond_to(scheduler, id_fiber_interrupt)) {
        return rb_thread_io_blocking_operation(io, fiber_scheduler_io_pwrite, (VALUE)&arguments);
    } else {
        return fiber_scheduler_io_pwrite((VALUE)&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 Показать исходный код
VALUE
rb_fiber_scheduler_io_read(VALUE scheduler, VALUE io, VALUE buffer, size_t length, size_t offset)
{
    if (!rb_respond_to(scheduler, id_io_read)) {
        return RUBY_Qundef;
    }

    VALUE arguments[] = {
        scheduler, io, buffer, SIZET2NUM(length), SIZET2NUM(offset)
    };

    if (rb_respond_to(scheduler, id_fiber_interrupt)) {
        return rb_thread_io_blocking_operation(io, fiber_scheduler_io_read, (VALUE)&arguments);
    } else {
        return fiber_scheduler_io_read((VALUE)&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) Показать исходный код
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.

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

io_wait(io, events, timeout) Показать исходный код
VALUE
rb_fiber_scheduler_io_wait(VALUE scheduler, VALUE io, VALUE events, VALUE timeout)
{
    VALUE arguments[] = {
        scheduler, io, events, timeout
    };

    if (rb_respond_to(scheduler, id_fiber_interrupt)) {
        return rb_thread_io_blocking_operation(io, fiber_scheduler_io_wait, (VALUE)&arguments);
    } else {
        return fiber_scheduler_io_wait((VALUE)&arguments);
    }
}

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

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

Рекомендуется зарегистрировать, какое Fiber ожидает какие ресурсы, и сразу вызвать Fiber.yield, чтобы передать управление другим волокнам. Затем в методе close планировщик может распределить все ресурсы ввода-вывода между ожидающими их волокнами.

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

io_write(io, buffer, length, offset) → written length or -errno Показать исходный код
VALUE
rb_fiber_scheduler_io_write(VALUE scheduler, VALUE io, VALUE buffer, size_t length, size_t offset)
{
    if (!rb_respond_to(scheduler, id_io_write)) {
        return RUBY_Qundef;
    }

    VALUE arguments[] = {
        scheduler, io, buffer, SIZET2NUM(length), SIZET2NUM(offset)
    };

    if (rb_respond_to(scheduler, id_fiber_interrupt)) {
        return rb_thread_io_blocking_operation(io, fiber_scheduler_io_write, (VALUE)&arguments);
    } else {
        return fiber_scheduler_io_write((VALUE)&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 и Thread::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) → result of 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 времени выполнения block, эту неблокирующую операцию следует прервать, вызвав указанное exception_class, созданное с переданными exception_arguments.

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

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

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

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

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

    VALUE result;
    enum ruby_tag_type state;

    // `rb_fiber_scheduler_unblock` can be called from points where `errno` is expected to be preserved. Therefore, we should save and restore it. For example `io_binwrite` calls `rb_fiber_scheduler_unblock` and if `errno` is reset to 0 by user code, it will break the error handling in `io_write`.
    //
    // If we explicitly preserve `errno` in `io_binwrite` and other similar functions (e.g. by returning it), this code is no longer needed. I hope in the future we will be able to remove it.
    int saved_errno = errno;

    // We must prevent interrupts while invoking the unblock method, because otherwise fibers can be left permanently blocked if an interrupt occurs during the execution of user code. See also `rb_fiber_scheduler_fiber_interrupt`.
    rb_execution_context_t *ec = GET_EC();
    int saved_interrupt_mask = ec->interrupt_mask;
    ec->interrupt_mask |= PENDING_INTERRUPT_MASK;

    EC_PUSH_TAG(ec);
    if ((state = EC_EXEC_TAG()) == TAG_NONE) {
        result = rb_funcall(scheduler, id_unblock, 2, blocker, fiber);
    }
    EC_POP_TAG();

    ec->interrupt_mask = saved_interrupt_mask;

    if (state) {
        EC_JUMP_TAG(ec, state);
    }

    RUBY_VM_CHECK_INTS(ec);

    errno = saved_errno;

    return result;
}

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

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

yield Показать исходный код
VALUE
rb_fiber_scheduler_yield(VALUE scheduler)
{
    // First try to call the scheduler's yield method, if it exists:
    VALUE result = rb_check_funcall(scheduler, id_yield, 0, NULL);
    if (!UNDEF_P(result)) return result;

    // Otherwise, we can emulate yield by sleeping for 0 seconds:
    return rb_fiber_scheduler_kernel_sleep(scheduler, RB_INT2NUM(0));
}

Передает управление планировщику, чтобы волокно было возобновлено в следующем цикле планирования.

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

Spec-Zone.ru

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