12#include "eval_intern.h"
20#include "internal/thread.h"
23#include "ruby_atomic.h"
26static ID id_scheduler_close;
33#ifdef FIBER_SCHEDULER_CALL_TIMEOUT_AFTER
34static ID id_timeout_after;
36static ID id_kernel_sleep;
37static ID id_process_wait;
39static ID id_io_read, id_io_pread;
40static ID id_io_write, id_io_pwrite;
42static ID id_io_select;
45static ID id_address_resolve;
47static ID id_blocking_operation_wait;
48static ID id_fiber_interrupt;
50static ID id_fiber_schedule;
53static VALUE rb_cFiberSchedulerBlockingOperation;
62 RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED,
63 RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_EXECUTING,
64 RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_COMPLETED,
65 RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_CANCELLED
66} rb_fiber_blocking_operation_status_t;
69 void *(*function)(
void *);
83blocking_operation_memsize(
const void *ptr)
89 "Fiber::Scheduler::BlockingOperation",
93 blocking_operation_memsize,
95 0, 0, RUBY_TYPED_THREAD_SAFE_FREE | RUBY_TYPED_WB_PROTECTED | RUBY_TYPED_EMBEDDABLE
102blocking_operation_alloc(
VALUE klass)
107 blocking_operation->function = NULL;
108 blocking_operation->data = NULL;
109 blocking_operation->unblock_function = NULL;
110 blocking_operation->data2 = NULL;
111 blocking_operation->flags = 0;
112 blocking_operation->state = NULL;
113 blocking_operation->status = RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED;
122get_blocking_operation(
VALUE obj)
126 return blocking_operation;
138blocking_operation_call(
VALUE self)
142 if (blocking_operation->status != RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED) {
143 rb_raise(
rb_eRuntimeError,
"Blocking operation has already been executed!");
146 if (blocking_operation->function == NULL) {
147 rb_raise(
rb_eRuntimeError,
"Blocking operation has no function to execute!");
150 if (blocking_operation->state == NULL) {
155 blocking_operation->status = RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_EXECUTING;
158 blocking_operation->state->result =
rb_nogvl(blocking_operation->function, blocking_operation->data,
159 blocking_operation->unblock_function, blocking_operation->data2,
160 blocking_operation->flags);
161 blocking_operation->state->saved_errno = rb_errno();
164 blocking_operation->status = RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_COMPLETED;
182 return get_blocking_operation(self);
197 if (blocking_operation == NULL) {
201 if (blocking_operation->function == NULL || blocking_operation->state == NULL) {
206 rb_thread_resolve_unblock_function(&blocking_operation->unblock_function, &blocking_operation->data2, GET_THREAD());
209 rb_atomic_t expected = RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED;
210 if (
RUBY_ATOMIC_CAS(blocking_operation->status, expected, RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_EXECUTING) != expected) {
216 blocking_operation->state->result = blocking_operation->function(blocking_operation->data);
217 blocking_operation->state->saved_errno =
errno;
220 expected = RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_EXECUTING;
221 if (
RUBY_ATOMIC_CAS(blocking_operation->status, expected, RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_COMPLETED) == expected) {
226 blocking_operation->state->saved_errno = EINTR;
238rb_fiber_scheduler_blocking_operation_new(
void *(*function)(
void *),
void *data,
242 VALUE self = blocking_operation_alloc(rb_cFiberSchedulerBlockingOperation);
245 blocking_operation->function = function;
246 blocking_operation->data = data;
247 blocking_operation->unblock_function = unblock_function;
248 blocking_operation->data2 = data2;
249 blocking_operation->flags = flags;
250 blocking_operation->state = state;
307Init_Fiber_Scheduler(
void)
316#ifdef FIBER_SCHEDULER_CALL_TIMEOUT_AFTER
333 id_blocking_operation_wait =
rb_intern_const(
"blocking_operation_wait");
342 rb_define_method(rb_cFiberSchedulerBlockingOperation,
"call", blocking_operation_call, 0);
345 rb_gc_register_mark_object(rb_cFiberSchedulerBlockingOperation);
348 rb_cFiberScheduler = rb_define_class_under(rb_cFiber,
"Scheduler",
rb_cObject);
359 rb_define_method(rb_cFiberScheduler,
"timeout_after", rb_fiber_scheduler_timeout_after, 3);
378 return thread->scheduler;
382verify_interface(
VALUE scheduler)
385 rb_raise(rb_eArgError,
"Scheduler must implement #block");
389 rb_raise(rb_eArgError,
"Scheduler must implement #unblock");
393 rb_raise(rb_eArgError,
"Scheduler must implement #kernel_sleep");
397 rb_raise(rb_eArgError,
"Scheduler must implement #io_wait");
401 rb_raise(rb_eArgError,
"Scheduler must implement #fiber_interrupt");
406fiber_scheduler_close(
VALUE scheduler)
412fiber_scheduler_close_ensure(
VALUE _thread)
415 thread->scheduler =
Qnil;
428 if (scheduler !=
Qnil) {
429 verify_interface(scheduler);
436 if (thread->scheduler !=
Qnil) {
438 rb_ensure(fiber_scheduler_close, thread->scheduler, fiber_scheduler_close_ensure, (
VALUE)thread);
441 thread->scheduler = scheduler;
443 return thread->scheduler;
447fiber_scheduler_current_for_threadptr(
rb_thread_t *thread)
451 if (thread->blocking == 0) {
452 return thread->scheduler;
463 return fiber_scheduler_current_for_threadptr(GET_THREAD());
469 return fiber_scheduler_current_for_threadptr(rb_thread_ptr(thread));
474 return fiber_scheduler_current_for_threadptr(thread);
501 if (!UNDEF_P(result))
return result;
504 if (!UNDEF_P(result))
return result;
513 return rb_float_new((
double)timeout->tv_sec + (0.000001 * timeout->tv_usec));
533 return rb_funcall(scheduler, id_kernel_sleep, 1, timeout);
539 return rb_funcallv(scheduler, id_kernel_sleep, argc, argv);
553 if (!UNDEF_P(result))
return result;
559#ifdef FIBER_SCHEDULER_CALL_TIMEOUT_AFTER
591 VALUE arguments[] = {
592 timeout, exception, message
599rb_fiber_scheduler_timeout_afterv(
VALUE scheduler,
int argc,
VALUE * argv)
626 VALUE arguments[] = {
650 return rb_funcall(scheduler, id_block, 2, blocker, timeout);
672 enum ruby_tag_type state;
677 int saved_errno =
errno;
681 volatile int saved_interrupt_mask = ec->interrupt_mask;
682 ec->interrupt_mask |= PENDING_INTERRUPT_MASK;
686 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
687 result =
rb_funcall(scheduler, id_unblock, 2, blocker, fiber);
690 rb_vm_rewind_cfp(ec, cfp);
694 ec->interrupt_mask = saved_interrupt_mask;
697 EC_JUMP_TAG(ec, state);
700 RUBY_VM_CHECK_INTS(ec);
727fiber_scheduler_io_wait(
VALUE _argument) {
730 return rb_funcallv(arguments[0], id_io_wait, 3, arguments + 1);
736 VALUE arguments[] = {
737 scheduler, io, events, timeout
740 return rb_thread_io_blocking_operation(io, fiber_scheduler_io_wait, (
VALUE)&arguments);
767 VALUE arguments[] = {
768 readables, writables, exceptables, timeout
806fiber_scheduler_io_read(
VALUE _argument) {
809 return rb_funcallv(arguments[0], id_io_read, 4, arguments + 1);
819 VALUE arguments[] = {
823 return rb_thread_io_blocking_operation(io, fiber_scheduler_io_read, (
VALUE)&arguments);
840fiber_scheduler_io_pread(
VALUE _argument) {
843 return rb_funcallv(arguments[0], id_io_pread, 5, arguments + 1);
853 VALUE arguments[] = {
857 return rb_thread_io_blocking_operation(io, fiber_scheduler_io_pread, (
VALUE)&arguments);
883fiber_scheduler_io_write(
VALUE _argument) {
886 return rb_funcallv(arguments[0], id_io_write, 4, arguments + 1);
896 VALUE arguments[] = {
900 return rb_thread_io_blocking_operation(io, fiber_scheduler_io_write, (
VALUE)&arguments);
917fiber_scheduler_io_pwrite(
VALUE _argument) {
920 return rb_funcallv(arguments[0], id_io_pwrite, 5, arguments + 1);
932 VALUE arguments[] = {
936 return rb_thread_io_blocking_operation(io, fiber_scheduler_io_pwrite, (
VALUE)&arguments);
949fiber_scheduler_io_read_memory(
VALUE _arguments)
953 return rb_fiber_scheduler_io_read(arguments->scheduler, arguments->io, arguments->buffer, arguments->offset, arguments->length);
959 VALUE buffer = rb_io_buffer_new_locked(base, size, 0);
962 .scheduler = scheduler,
969 return rb_ensure(fiber_scheduler_io_read_memory, (
VALUE)&arguments, rb_io_buffer_free_locked, buffer);
973fiber_scheduler_io_write_memory(
VALUE _arguments)
983 VALUE buffer = rb_io_buffer_new_locked((
void*)base, size, RB_IO_BUFFER_READONLY);
986 .scheduler = scheduler,
993 return rb_ensure(fiber_scheduler_io_write_memory, (
VALUE)&arguments, rb_io_buffer_free_locked, buffer);
997fiber_scheduler_io_pread_memory(
VALUE _arguments)
1001 return rb_fiber_scheduler_io_pread(arguments->scheduler, arguments->io, arguments->from, arguments->buffer, arguments->offset, arguments->length);
1007 VALUE buffer = rb_io_buffer_new_locked(base, size, 0);
1010 .scheduler = scheduler,
1018 return rb_ensure(fiber_scheduler_io_pread_memory, (
VALUE)&arguments, rb_io_buffer_free_locked, buffer);
1022fiber_scheduler_io_pwrite_memory(
VALUE _arguments)
1026 return rb_fiber_scheduler_io_pwrite(arguments->scheduler, arguments->io, arguments->from, arguments->buffer, arguments->offset, arguments->length);
1032 VALUE buffer = rb_io_buffer_new_locked((
void*)base, size, RB_IO_BUFFER_READONLY);
1035 .scheduler = scheduler,
1043 return rb_ensure(fiber_scheduler_io_pwrite_memory, (
VALUE)&arguments, rb_io_buffer_free_locked, buffer);
1057 VALUE arguments[] = {io};
1097 VALUE arguments[] = {
1124 if (!
rb_respond_to(scheduler, id_blocking_operation_wait)) {
1129 VALUE blocking_operation = rb_fiber_scheduler_blocking_operation_new(function, data, unblock_function, data2, flags, state);
1131 VALUE result =
rb_funcall(scheduler, id_blocking_operation_wait, 1, blocking_operation);
1138 operation->function = NULL;
1139 operation->state = NULL;
1140 operation->data = NULL;
1141 operation->data2 = NULL;
1142 operation->unblock_function = NULL;
1148 if (current_status == RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED) {
1166 VALUE arguments[] = {
1171 enum ruby_tag_type state;
1175 volatile int saved_interrupt_mask = ec->interrupt_mask;
1176 ec->interrupt_mask |= PENDING_INTERRUPT_MASK;
1180 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
1181 result =
rb_funcallv(scheduler, id_fiber_interrupt, 2, arguments);
1184 rb_vm_rewind_cfp(ec, cfp);
1188 ec->interrupt_mask = saved_interrupt_mask;
1191 EC_JUMP_TAG(ec, state);
1194 RUBY_VM_CHECK_INTS(ec);
1232 if (blocking_operation == NULL) {
1238 switch (current_state) {
1239 case RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_QUEUED:
1241 if (
RUBY_ATOMIC_CAS(blocking_operation->status, current_state, RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_CANCELLED) == current_state) {
1247 case RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_EXECUTING:
1249 if (
RUBY_ATOMIC_CAS(blocking_operation->status, current_state, RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_CANCELLED) != current_state) {
1255 if (unblock_function) {
1257 blocking_operation->unblock_function(blocking_operation->data2);
1262 case RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_COMPLETED:
1263 case RB_FIBER_SCHEDULER_BLOCKING_OPERATION_STATUS_CANCELLED:
#define RUBY_ASSERT(...)
Asserts that the given expression is truthy if and only if RUBY_DEBUG is truthy.
#define RUBY_ATOMIC_CAS(var, oldval, newval)
Atomic compare-and-swap.
std::atomic< unsigned > rb_atomic_t
Type that is eligible for atomic operations.
#define RUBY_ATOMIC_LOAD(var)
Atomic load.
#define rb_define_method(klass, mid, func, arity)
Defines klass#mid.
VALUE rb_class_new(VALUE super)
Creates a new, anonymous class.
#define Qundef
Old name of RUBY_Qundef.
#define SIZET2NUM
Old name of RB_SIZE2NUM.
#define Qnil
Old name of RUBY_Qnil.
VALUE rb_eRuntimeError
RuntimeError exception.
VALUE rb_cObject
Object class.
VALUE rb_funcall(VALUE recv, ID mid, int n,...)
Calls a method.
VALUE rb_funcallv(VALUE recv, ID mid, int argc, const VALUE *argv)
Identical to rb_funcall(), except it takes the method arguments as a C array.
VALUE rb_funcall_passing_block_kw(VALUE recv, ID mid, int argc, const VALUE *argv, int kw_splat)
Identical to rb_funcallv_passing_block(), except you can specify how to handle the last element of th...
void rb_unblock_function_t(void *)
This is the type of UBFs.
int rb_respond_to(VALUE obj, ID mid)
Queries if the object responds to the method.
VALUE rb_check_funcall(VALUE recv, ID mid, int argc, const VALUE *argv)
Identical to rb_funcallv(), except it returns RUBY_Qundef instead of raising rb_eNoMethodError.
void rb_define_alloc_func(VALUE klass, rb_alloc_func_t func)
Sets the allocator function of a class.
static ID rb_intern_const(const char *str)
This is a "tiny optimisation" over rb_intern().
VALUE rb_io_timeout(VALUE io)
Get the timeout associated with the specified io object.
@ RUBY_IO_READABLE
IO::READABLE
@ RUBY_IO_WRITABLE
IO::WRITABLE
void * rb_nogvl(void *(*func)(void *), void *data1, rb_unblock_function_t *ubf, void *data2, int flags)
Identical to rb_thread_call_without_gvl(), except it additionally takes "flags" that change the behav...
#define RB_UINT2NUM
Just another name of rb_uint2num_inline.
#define RB_INT2NUM
Just another name of rb_int2num_inline.
#define RB_GC_GUARD(v)
Prevents premature destruction of local objects.
VALUE rb_ensure(type *q, VALUE w, type *e, VALUE r)
An equivalent of ensure clause.
#define OFFT2NUM
Converts a C's off_t into an instance of rb_cInteger.
#define PIDT2NUM
Converts a C's pid_t into an instance of rb_cInteger.
#define RUBY_DEFAULT_FREE
This is a value you can set to RData::dfree.
#define TypedData_Get_Struct(obj, type, data_type, sval)
Obtains a C struct from inside of a wrapper Ruby object.
#define TypedData_Make_Struct(klass, type, data_type, sval)
Identical to TypedData_Wrap_Struct, except it allocates a new data region internally instead of takin...
#define errno
Ractor-aware version of errno.
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)
Defer the execution of the passed function to the scheduler.
VALUE rb_fiber_scheduler_current(void)
Identical to rb_fiber_scheduler_get(), except it also returns RUBY_Qnil in case of a blocking fiber.
VALUE rb_fiber_scheduler_make_timeout(struct timeval *timeout)
Converts the passed timeout to an expression that rb_fiber_scheduler_block() etc.
VALUE rb_fiber_scheduler_io_wait_readable(VALUE scheduler, VALUE io)
Non-blocking wait until the passed IO is ready for reading.
VALUE rb_fiber_scheduler_kernel_sleepv(VALUE scheduler, int argc, VALUE *argv)
Identical to rb_fiber_scheduler_kernel_sleep(), except it can pass multiple arguments.
VALUE rb_fiber_scheduler_fiber_interrupt(VALUE scheduler, VALUE fiber, VALUE exception)
Interrupt a fiber by raising an exception.
VALUE rb_fiber_scheduler_io_wait(VALUE scheduler, VALUE io, VALUE events, VALUE timeout)
Non-blocking version of rb_io_wait().
VALUE rb_fiber_scheduler_io_select(VALUE scheduler, VALUE readables, VALUE writables, VALUE exceptables, VALUE timeout)
Non-blocking version of IO.select.
VALUE rb_fiber_scheduler_io_read(VALUE scheduler, VALUE io, VALUE buffer, size_t offset, size_t length)
Non-blocking read from the passed IO.
int rb_fiber_scheduler_blocking_operation_cancel(rb_fiber_scheduler_blocking_operation_t *blocking_operation)
Cancel a blocking operation.
VALUE rb_fiber_scheduler_io_pwrite(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t offset, size_t length)
Non-blocking write to the passed IO at the specified offset.
VALUE rb_fiber_scheduler_io_pread_memory(VALUE scheduler, VALUE io, rb_off_t from, void *base, size_t size)
Non-blocking pread from the passed IO using a native buffer.
VALUE rb_fiber_scheduler_io_selectv(VALUE scheduler, int argc, VALUE *argv)
Non-blocking version of IO.select, argv variant.
VALUE rb_fiber_scheduler_process_wait(VALUE scheduler, rb_pid_t pid, int flags)
Non-blocking waitpid.
VALUE rb_fiber_scheduler_io_pread(VALUE scheduler, VALUE io, rb_off_t from, VALUE buffer, size_t offset, size_t length)
Non-blocking read from the passed IO at the specified offset.
VALUE rb_fiber_scheduler_block(VALUE scheduler, VALUE blocker, VALUE timeout)
Non-blocking wait for the passed "blocker", which is for instance Thread.join or Mutex....
VALUE rb_fiber_scheduler_io_read_memory(VALUE scheduler, VALUE io, void *base, size_t size)
Non-blocking read from the passed IO using a native buffer.
int rb_fiber_scheduler_blocking_operation_execute(rb_fiber_scheduler_blocking_operation_t *blocking_operation)
Execute blocking operation from handle (GVL not required).
VALUE rb_fiber_scheduler_io_write(VALUE scheduler, VALUE io, VALUE buffer, size_t offset, size_t length)
Non-blocking write to the passed IO.
VALUE rb_fiber_scheduler_close(VALUE scheduler)
Closes the passed scheduler object.
rb_fiber_scheduler_blocking_operation_t * rb_fiber_scheduler_blocking_operation_extract(VALUE self)
Extract the blocking operation handle from a BlockingOperationRuby object.
VALUE rb_fiber_scheduler_current_for_thread(VALUE thread)
Identical to rb_fiber_scheduler_current(), except it queries for that of the passed thread value inst...
VALUE rb_fiber_scheduler_kernel_sleep(VALUE scheduler, VALUE duration)
Non-blocking sleep.
VALUE rb_fiber_scheduler_address_resolve(VALUE scheduler, VALUE hostname)
Non-blocking DNS lookup.
VALUE rb_fiber_scheduler_yield(VALUE scheduler)
Yield to the scheduler, to be resumed on the next scheduling cycle.
VALUE rb_fiber_scheduler_io_pwrite_memory(VALUE scheduler, VALUE io, rb_off_t from, const void *base, size_t size)
Non-blocking pwrite to the passed IO using a native buffer.
VALUE rb_fiber_scheduler_set(VALUE scheduler)
Destructively assigns the passed scheduler to that of the current thread that is calling this functio...
VALUE rb_fiber_scheduler_current_for_threadptr(struct rb_thread_struct *thread)
Identical to rb_fiber_scheduler_current_for_thread(), except it expects a threadptr instead of a thre...
VALUE rb_fiber_scheduler_io_wait_writable(VALUE scheduler, VALUE io)
Non-blocking wait until the passed IO is ready for writing.
VALUE rb_fiber_scheduler_io_close(VALUE scheduler, VALUE io)
Non-blocking close the given IO.
VALUE rb_fiber_scheduler_get(void)
Queries the current scheduler of the current thread that is calling this function.
VALUE rb_fiber_scheduler_unblock(VALUE scheduler, VALUE blocker, VALUE fiber)
Wakes up a fiber previously blocked using rb_fiber_scheduler_block().
VALUE rb_fiber_scheduler_io_write_memory(VALUE scheduler, VALUE io, const void *base, size_t size)
Non-blocking write to the passed IO using a native buffer.
VALUE rb_fiber_scheduler_fiber(VALUE scheduler, int argc, VALUE *argv, int kw_splat)
Create and schedule a non-blocking fiber.
@ RUBY_Qundef
Represents so-called undef.
This is the struct that holds necessary info for a struct.
uintptr_t ID
Type that represents a Ruby identifier such as a variable name.
uintptr_t VALUE
Type that represents a Ruby object.