diff --git a/.github/workflows/test.yaml b/.github/workflows/test.yaml index 03ea1c90..96c749b7 100644 --- a/.github/workflows/test.yaml +++ b/.github/workflows/test.yaml @@ -8,15 +8,15 @@ permissions: jobs: test: name: ${{matrix.ruby}} on ${{matrix.os}} - runs-on: ${{matrix.os}}-latest + runs-on: ${{matrix.os}} continue-on-error: ${{matrix.experimental}} strategy: matrix: os: - - ubuntu - - macos - - windows + - ubuntu-latest + - macos-latest + - windows-latest ruby: - "3.3" @@ -26,13 +26,15 @@ jobs: experimental: [false] include: - - os: ubuntu + - os: ubuntu-latest ruby: truffleruby experimental: true - - os: ubuntu + - os: ubuntu-latest ruby: jruby experimental: true - - os: ubuntu + # Use newer packaged liburing for Futex coverage on Ruby head. + # Switch back to ubuntu-latest after its 26.04 rollout (2026-11-19). + - os: ubuntu-26.04 ruby: head experimental: true @@ -44,8 +46,8 @@ jobs: bundler-cache: true - name: Install packages (Ubuntu) - if: matrix.os == 'ubuntu' - run: sudo apt-get install -y liburing-dev + if: startsWith(matrix.os, 'ubuntu-') + run: sudo apt-get update && sudo apt-get install -y liburing-dev - name: Run tests timeout-minutes: 10 diff --git a/ext/extconf.rb b/ext/extconf.rb index f9c43d80..2a87886b 100755 --- a/ext/extconf.rb +++ b/ext/extconf.rb @@ -31,9 +31,16 @@ have_func("rb_ext_ractor_safe") have_func("&rb_fiber_transfer") +have_io_buffer = have_header("ruby/io/buffer.h") + +if RUBY_PLATFORM.include?("linux") && have_io_buffer && have_header("linux/futex.h") && have_header("sys/syscall.h") + $srcs << "io/event/futex.c" +end if have_library("uring") and have_header("liburing.h") have_func("io_uring_prep_waitid", "liburing.h") + have_func("io_uring_prep_futex_wait", "liburing.h") + have_func("io_uring_prep_futex_waitv", "liburing.h") $srcs << "io/event/selector/uring.c" end @@ -57,8 +64,6 @@ have_func("&rb_fiber_raise") have_func("epoll_pwait2(0, 0, 0, 0, 0)", "sys/epoll.h") if enable_config("epoll_pwait2", true) -have_header("ruby/io/buffer.h") - # Feature detection for blocking operation support if have_func("rb_fiber_scheduler_blocking_operation_extract") # Feature detection for pthread support (needed for WorkerPool) diff --git a/ext/io/event/event.c b/ext/io/event/event.c index 59492544..eb49c059 100644 --- a/ext/io/event/event.c +++ b/ext/io/event/event.c @@ -3,6 +3,7 @@ #include "event.h" #include "fiber.h" +#include "futex.h" #include "selector/selector.h" void Init_IO_Event(void) @@ -15,6 +16,10 @@ void Init_IO_Event(void) Init_IO_Event_Fiber(IO_Event); + #ifdef IO_EVENT_FUTEX + Init_IO_Event_Futex(IO_Event); + #endif + #ifdef HAVE_IO_EVENT_WORKER_POOL Init_IO_Event_WorkerPool(IO_Event); #endif diff --git a/ext/io/event/futex.c b/ext/io/event/futex.c new file mode 100644 index 00000000..e52b3ecb --- /dev/null +++ b/ext/io/event/futex.c @@ -0,0 +1,489 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +#include "futex.h" + +#ifdef IO_EVENT_FUTEX + +#include +#include +#include +#include +#include + +#include +#include +#include + +// The finalizer owns the allocation lock independently of the Futex. It must +// never retain the Futex or its C struct: either may be reclaimed before the +// finalizer runs. Keeping the buffer here preserves it until it is unlocked. +struct IO_Event_Futex_Finalizer { + VALUE buffer; + uint32_t *address; + size_t waits; +}; + +static VALUE IO_Event_Futex_Finalizer_Class; + +static void IO_Event_Futex_Finalizer_mark(void *_finalizer) { + struct IO_Event_Futex_Finalizer *finalizer = _finalizer; + rb_gc_mark_movable(finalizer->buffer); +} + +static void IO_Event_Futex_Finalizer_compact(void *_finalizer) { + struct IO_Event_Futex_Finalizer *finalizer = _finalizer; + finalizer->buffer = rb_gc_location(finalizer->buffer); +} + +static size_t IO_Event_Futex_Finalizer_size(const void *_finalizer) { + return sizeof(struct IO_Event_Futex_Finalizer); +} + +static const rb_data_type_t IO_Event_Futex_Finalizer_Type = { + .wrap_struct_name = "IO::Event::Futex::Finalizer", + .function = { + .dmark = IO_Event_Futex_Finalizer_mark, + .dcompact = IO_Event_Futex_Finalizer_compact, + .dfree = RUBY_TYPED_DEFAULT_FREE, + .dsize = IO_Event_Futex_Finalizer_size, + }, + .flags = RUBY_TYPED_FREE_IMMEDIATELY | RUBY_TYPED_WB_PROTECTED, +}; + +static VALUE IO_Event_Futex_Finalizer_release(VALUE self) { + struct IO_Event_Futex_Finalizer *finalizer = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + if (!NIL_P(finalizer->buffer)) { + rb_io_buffer_unlock(finalizer->buffer); + RB_OBJ_WRITE(self, &finalizer->buffer, Qnil); + finalizer->address = NULL; + } + return Qnil; +} + +static VALUE IO_Event_Futex_Finalizer_call(VALUE self, VALUE object_id) { + struct IO_Event_Futex_Finalizer *finalizer = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + // Ruby also invokes finalizers at shutdown, even for reachable objects. + // Never unlock memory still used by an outstanding kernel operation. + if (finalizer->waits) return Qnil; + return IO_Event_Futex_Finalizer_release(self); +} + +struct IO_Event_Futex { + VALUE finalizer; +}; + +static void IO_Event_Futex_mark(void *_futex) { + struct IO_Event_Futex *futex = _futex; + rb_gc_mark_movable(futex->finalizer); +} + +static void IO_Event_Futex_compact(void *_futex) { + struct IO_Event_Futex *futex = _futex; + futex->finalizer = rb_gc_location(futex->finalizer); +} + +static size_t IO_Event_Futex_size(const void *_futex) { + return sizeof(struct IO_Event_Futex); +} + +static const rb_data_type_t IO_Event_Futex_Type = { + .wrap_struct_name = "IO::Event::Futex", + .function = { + .dmark = IO_Event_Futex_mark, + .dcompact = IO_Event_Futex_compact, + .dfree = RUBY_TYPED_DEFAULT_FREE, + .dsize = IO_Event_Futex_size, + }, + .flags = RUBY_TYPED_FREE_IMMEDIATELY | RUBY_TYPED_WB_PROTECTED, +}; + +static VALUE IO_Event_Futex_allocate(VALUE klass) { + struct IO_Event_Futex *futex = NULL; + VALUE instance = TypedData_Make_Struct(klass, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + futex->finalizer = Qnil; + return instance; +} + +static ID id_offset; +static ID id_futex_wait; +static ID id_futex_waitv; + +static VALUE IO_Event_Futex_initialize(int argc, VALUE *argv, VALUE self) { + rb_check_frozen(self); + struct IO_Event_Futex *futex = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + if (!NIL_P(futex->finalizer)) rb_raise(rb_eRuntimeError, "Futex is already initialized!"); + + VALUE buffer, options; + rb_scan_args(argc, argv, "1:", &buffer, &options); + + VALUE offset_value = Qundef; + if (!NIL_P(options)) { + ID keys[] = {id_offset}; + VALUE values[1]; + rb_get_kwargs(options, keys, 0, 1, values); + offset_value = values[0]; + } + + size_t offset = offset_value == Qundef ? 0 : NUM2SIZET(offset_value); + // Numeric coercion may have invoked initialize recursively. + if (!NIL_P(futex->finalizer)) rb_raise(rb_eRuntimeError, "Futex is already initialized!"); + + // Register cleanup before acquiring the lock. All allocating operations are + // performed before extracting the pointer; failures leave an inert finalizer. + struct IO_Event_Futex_Finalizer *finalizer = NULL; + VALUE finalizer_value = TypedData_Make_Struct(IO_Event_Futex_Finalizer_Class, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + finalizer->buffer = Qnil; + finalizer->address = NULL; + finalizer->waits = 0; + rb_define_finalizer(self, finalizer_value); + + void *base = NULL; + size_t size = 0; + rb_io_buffer_get_bytes_for_writing(buffer, &base, &size); + + if (offset > size || size - offset < sizeof(uint32_t)) { + rb_raise(rb_eRangeError, "Futex offset exceeds the buffer size!"); + } + + uint32_t *address = (uint32_t *)((char *)base + offset); + if ((uintptr_t)address % sizeof(uint32_t) != 0) { + rb_raise(rb_eArgError, "Futex address must be aligned to 4 bytes!"); + } + + // Nothing after a successful lock can raise or invoke Ruby code. + rb_io_buffer_lock(buffer); + RB_OBJ_WRITE(finalizer_value, &finalizer->buffer, buffer); + finalizer->address = address; + RB_OBJ_WRITE(self, &futex->finalizer, finalizer_value); + + return self; +} + +static VALUE IO_Event_Futex_value(VALUE self) { + return UINT2NUM(__atomic_load_n(IO_Event_Futex_address(self), __ATOMIC_ACQUIRE)); +} + +static VALUE IO_Event_Futex_set_value(VALUE self, VALUE value) { + uint32_t converted = NUM2UINT(value); + __atomic_store_n(IO_Event_Futex_address(self), converted, __ATOMIC_RELEASE); + return value; +} + +static VALUE IO_Event_Futex_increment(int argc, VALUE *argv, VALUE self) { + VALUE amount_value; + rb_scan_args(argc, argv, "01", &amount_value); + uint32_t amount = NIL_P(amount_value) ? 1 : NUM2UINT(amount_value); + + uint32_t value = __atomic_add_fetch(IO_Event_Futex_address(self), amount, __ATOMIC_ACQ_REL); + return UINT2NUM(value); +} + +static VALUE IO_Event_Futex_decrement(int argc, VALUE *argv, VALUE self) { + VALUE amount_value; + rb_scan_args(argc, argv, "01", &amount_value); + uint32_t amount = NIL_P(amount_value) ? 1 : NUM2UINT(amount_value); + + uint32_t value = __atomic_sub_fetch(IO_Event_Futex_address(self), amount, __ATOMIC_ACQ_REL); + return UINT2NUM(value); +} + +static VALUE IO_Event_Futex_compare_exchange(VALUE self, VALUE expected_value, VALUE desired_value) { + uint32_t expected = NUM2UINT(expected_value); + uint32_t desired = NUM2UINT(desired_value); + + bool exchanged = __atomic_compare_exchange_n( + IO_Event_Futex_address(self), + &expected, + desired, + false, + __ATOMIC_ACQ_REL, + __ATOMIC_ACQUIRE + ); + + return exchanged ? Qtrue : Qfalse; +} + +static VALUE IO_Event_Futex_wake(int argc, VALUE *argv, VALUE self) { + VALUE count_value; + rb_scan_args(argc, argv, "01", &count_value); + int count = NIL_P(count_value) ? 1 : NUM2INT(count_value); + if (count < 0) rb_raise(rb_eArgError, "Wake count must be non-negative!"); + + uint32_t *address = IO_Event_Futex_address(self); + // Legacy FUTEX_WAKE can wake one waiter even when count is zero. + if (count == 0) return INT2NUM(0); + + int result = syscall(SYS_futex, address, FUTEX_WAKE, count, NULL, NULL, 0); + if (result < 0) rb_sys_fail("IO_Event_Futex_wake:futex"); + return INT2NUM(result); +} + +static VALUE IO_Event_Futex_signal(int argc, VALUE *argv, VALUE self) { + VALUE count_value; + rb_scan_args(argc, argv, "01", &count_value); + + int count = NIL_P(count_value) ? 1 : NUM2INT(count_value); + if (count < 0) rb_raise(rb_eArgError, "Wake count must be non-negative!"); + uint32_t *address = IO_Event_Futex_address(self); + // A zero-count signal neither changes the word nor wakes any waiters. + if (count == 0) return UINT2NUM(__atomic_load_n(address, __ATOMIC_ACQUIRE)); + + uint32_t value = __atomic_add_fetch(address, 1, __ATOMIC_ACQ_REL); + int result = syscall(SYS_futex, address, FUTEX_WAKE, count, NULL, NULL, 0); + if (result < 0) rb_sys_fail("IO_Event_Futex_signal:futex"); + return UINT2NUM(value); +} + +static struct IO_Event_Futex_Finalizer *IO_Event_Futex_get(VALUE self) { + struct IO_Event_Futex *futex = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + if (NIL_P(futex->finalizer)) rb_raise(rb_eIOError, "Futex is uninitialized!"); + struct IO_Event_Futex_Finalizer *finalizer = NULL; + TypedData_Get_Struct(futex->finalizer, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + if (!finalizer->address) rb_raise(rb_eIOError, "Futex is closed!"); + return finalizer; +} + +uint32_t *IO_Event_Futex_address(VALUE self) { + return IO_Event_Futex_get(self)->address; +} + +uint32_t *IO_Event_Futex_acquire(VALUE self) { + struct IO_Event_Futex_Finalizer *finalizer = IO_Event_Futex_get(self); + if (finalizer->waits == SIZE_MAX) rb_raise(rb_eRuntimeError, "Too many futex waits!"); + finalizer->waits += 1; + return finalizer->address; +} + +void IO_Event_Futex_release(VALUE self) { + struct IO_Event_Futex_Finalizer *finalizer = IO_Event_Futex_get(self); + if (!finalizer->waits) rb_bug("IO_Event_Futex_release: no pending waits"); + finalizer->waits -= 1; +} + +static VALUE IO_Event_Futex_close(VALUE self) { + struct IO_Event_Futex *futex = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + if (NIL_P(futex->finalizer)) return Qnil; + struct IO_Event_Futex_Finalizer *finalizer = NULL; + TypedData_Get_Struct(futex->finalizer, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + if (finalizer->waits) rb_raise(rb_eIOError, "Cannot close a futex with pending waits!"); + return IO_Event_Futex_Finalizer_release(futex->finalizer); +} + +static VALUE IO_Event_Futex_closed_p(VALUE self) { + struct IO_Event_Futex *futex = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + if (NIL_P(futex->finalizer)) return Qtrue; + struct IO_Event_Futex_Finalizer *finalizer = NULL; + TypedData_Get_Struct(futex->finalizer, struct IO_Event_Futex_Finalizer, &IO_Event_Futex_Finalizer_Type, finalizer); + return finalizer->address ? Qfalse : Qtrue; +} + +static VALUE IO_Event_Futex_initialize_copy(VALUE self, VALUE other) { + struct IO_Event_Futex *futex = NULL; + TypedData_Get_Struct(self, struct IO_Event_Futex, &IO_Event_Futex_Type, futex); + // Ruby copies finalizers before invoking initialize_copy. The rejected copy + // must not release the original's allocation lock when it is collected. + if (NIL_P(futex->finalizer)) rb_undefine_finalizer(self); + rb_raise(rb_eTypeError, "Cannot copy a futex!"); +} + +struct IO_Event_Futex_BlockingWait { + VALUE futex; + uint32_t *address; + uint32_t expected; + int result; + int error; +}; + +static void *IO_Event_Futex_blocking_wait_without_gvl(void *_arguments) { + struct IO_Event_Futex_BlockingWait *arguments = _arguments; + arguments->result = syscall(SYS_futex, arguments->address, FUTEX_WAIT, arguments->expected, NULL, NULL, 0); + arguments->error = arguments->result < 0 ? errno : 0; + return NULL; +} + +static VALUE IO_Event_Futex_blocking_wait_body(VALUE _arguments) { + struct IO_Event_Futex_BlockingWait *arguments = (void *)_arguments; + arguments->address = IO_Event_Futex_acquire(arguments->futex); + rb_thread_call_without_gvl(IO_Event_Futex_blocking_wait_without_gvl, arguments, RUBY_UBF_IO, 0); + + if (arguments->result == 0) { + return Qtrue; + } else if (arguments->error == EAGAIN) { + return Qfalse; + } else { + rb_syserr_fail(arguments->error, "IO_Event_Futex_blocking_wait:futex"); + } + + return Qfalse; +} + +static VALUE IO_Event_Futex_blocking_wait_ensure(VALUE _arguments) { + struct IO_Event_Futex_BlockingWait *arguments = (void *)_arguments; + if (arguments->address) IO_Event_Futex_release(arguments->futex); + return Qnil; +} + +static VALUE IO_Event_Futex_blocking_wait(VALUE self, VALUE expected_value) { + struct IO_Event_Futex_BlockingWait arguments = { + .futex = self, + .expected = NUM2UINT(expected_value), + }; + VALUE result = rb_ensure(IO_Event_Futex_blocking_wait_body, (VALUE)&arguments, IO_Event_Futex_blocking_wait_ensure, (VALUE)&arguments); + RB_GC_GUARD(self); + return result; +} + +#ifdef IO_EVENT_FUTEX_WAITV + +// Coerce all arguments before retaining native addresses. The returned array +// owns a private snapshot of the futex references: mutating the caller's entries +// while a wait is pending must not make its futexes collectible. +VALUE IO_Event_Futex_prepare_waitv(VALUE entries, struct futex_waitv *vector) { + entries = rb_ary_dup(rb_Array(entries)); + long count = RARRAY_LEN(entries); + if (count < 1 || count > FUTEX_WAITV_MAX) { + rb_raise(rb_eArgError, "Futex vector must contain between 1 and %d entries!", FUTEX_WAITV_MAX); + } + VALUE futexes = rb_ary_new_capa(count); + for (long index = 0; index < count; index++) { + VALUE entry = rb_Array(RARRAY_AREF(entries, index)); + if (RARRAY_LEN(entry) != 2) { + rb_raise(rb_eArgError, "Each futex vector entry must contain a futex and its expected value!"); + } + rb_ary_push(futexes, RARRAY_AREF(entry, 0)); + vector[index].val = NUM2UINT(RARRAY_AREF(entry, 1)); + vector[index].flags = FUTEX_32; + vector[index].__reserved = 0; + } + return futexes; +} + +struct IO_Event_Futex_BlockingWaitV { + VALUE futexes; + struct futex_waitv *vector; + size_t count; + size_t acquired; + int result; + int error; +}; + +static void *IO_Event_Futex_blocking_waitv_without_gvl(void *_arguments) { + struct IO_Event_Futex_BlockingWaitV *arguments = _arguments; + arguments->result = syscall(SYS_futex_waitv, arguments->vector, arguments->count, 0, NULL, CLOCK_MONOTONIC); + arguments->error = arguments->result < 0 ? errno : 0; + return NULL; +} + +static VALUE IO_Event_Futex_blocking_waitv_body(VALUE _arguments) { + struct IO_Event_Futex_BlockingWaitV *arguments = (void *)_arguments; + while (arguments->acquired < arguments->count) { + size_t index = arguments->acquired; + arguments->vector[index].uaddr = (uintptr_t)IO_Event_Futex_acquire(RARRAY_AREF(arguments->futexes, index)); + arguments->acquired += 1; + } + rb_thread_call_without_gvl(IO_Event_Futex_blocking_waitv_without_gvl, arguments, RUBY_UBF_IO, 0); + + if (arguments->result >= 0) { + return INT2NUM(arguments->result); + } else if (arguments->error == EAGAIN) { + return Qnil; + } else { + rb_syserr_fail(arguments->error, "IO_Event_Futex_blocking_waitv:futex_waitv"); + } + + return Qnil; +} + +static VALUE IO_Event_Futex_blocking_waitv_ensure(VALUE _arguments) { + struct IO_Event_Futex_BlockingWaitV *arguments = (void *)_arguments; + while (arguments->acquired) { + IO_Event_Futex_release(RARRAY_AREF(arguments->futexes, --arguments->acquired)); + } + return Qnil; +} + +static VALUE IO_Event_Futex_blocking_waitv(VALUE entries) { + struct futex_waitv vector[FUTEX_WAITV_MAX]; + VALUE futexes = IO_Event_Futex_prepare_waitv(entries, vector); + struct IO_Event_Futex_BlockingWaitV arguments = { + .futexes = futexes, + .vector = vector, + .count = RARRAY_LEN(futexes), + }; + VALUE result = rb_ensure(IO_Event_Futex_blocking_waitv_body, (VALUE)&arguments, IO_Event_Futex_blocking_waitv_ensure, (VALUE)&arguments); + RB_GC_GUARD(futexes); + return result; +} + +#endif + +static VALUE IO_Event_Futex_wait(VALUE self, VALUE expected_value) { + VALUE scheduler = rb_fiber_scheduler_current(); + if (NIL_P(scheduler)) { + return IO_Event_Futex_blocking_wait(self, expected_value); + } + + if (!rb_respond_to(scheduler, id_futex_wait)) { + rb_raise(rb_eNotImpError, "The current fiber scheduler does not support futex waits!"); + } + + return rb_funcall(scheduler, id_futex_wait, 2, self, expected_value); +} + +#ifdef IO_EVENT_FUTEX_WAITV + +static VALUE IO_Event_Futex_wait_any(VALUE klass, VALUE entries) { + (void)klass; + VALUE scheduler = rb_fiber_scheduler_current(); + if (NIL_P(scheduler)) { + return IO_Event_Futex_blocking_waitv(entries); + } + + if (!rb_respond_to(scheduler, id_futex_waitv)) { + rb_raise(rb_eNotImpError, "The current fiber scheduler does not support vector futex waits!"); + } + + return rb_funcall(scheduler, id_futex_waitv, 1, entries); +} + +#endif + +void Init_IO_Event_Futex(VALUE IO_Event) { + IO_Event_Futex_Finalizer_Class = rb_class_new(rb_cObject); + rb_gc_register_mark_object(IO_Event_Futex_Finalizer_Class); + rb_undef_alloc_func(IO_Event_Futex_Finalizer_Class); + rb_define_method(IO_Event_Futex_Finalizer_Class, "call", IO_Event_Futex_Finalizer_call, 1); + + VALUE IO_Event_Futex = rb_define_class_under(IO_Event, "Futex", rb_cObject); + rb_define_alloc_func(IO_Event_Futex, IO_Event_Futex_allocate); + rb_define_method(IO_Event_Futex, "initialize", IO_Event_Futex_initialize, -1); + rb_define_method(IO_Event_Futex, "initialize_copy", IO_Event_Futex_initialize_copy, 1); + rb_define_method(IO_Event_Futex, "close", IO_Event_Futex_close, 0); + rb_define_method(IO_Event_Futex, "closed?", IO_Event_Futex_closed_p, 0); + rb_define_method(IO_Event_Futex, "value", IO_Event_Futex_value, 0); + rb_define_method(IO_Event_Futex, "value=", IO_Event_Futex_set_value, 1); + rb_define_method(IO_Event_Futex, "increment", IO_Event_Futex_increment, -1); + rb_define_method(IO_Event_Futex, "decrement", IO_Event_Futex_decrement, -1); + rb_define_method(IO_Event_Futex, "compare_exchange", IO_Event_Futex_compare_exchange, 2); + rb_define_method(IO_Event_Futex, "wake", IO_Event_Futex_wake, -1); + rb_define_method(IO_Event_Futex, "signal", IO_Event_Futex_signal, -1); + rb_define_method(IO_Event_Futex, "wait", IO_Event_Futex_wait, 1); + +#ifdef IO_EVENT_FUTEX_WAITV + rb_define_const(IO_Event_Futex, "WAITV_LIMIT", INT2NUM(FUTEX_WAITV_MAX)); + rb_define_singleton_method(IO_Event_Futex, "wait_any", IO_Event_Futex_wait_any, 1); +#endif + + id_offset = rb_intern("offset"); + id_futex_wait = rb_intern("futex_wait"); + id_futex_waitv = rb_intern("futex_waitv"); +} + +#endif diff --git a/ext/io/event/futex.h b/ext/io/event/futex.h new file mode 100644 index 00000000..54b12053 --- /dev/null +++ b/ext/io/event/futex.h @@ -0,0 +1,42 @@ +// Released under the MIT License. +// Copyright, 2026, by Samuel Williams. + +#pragma once + +#include + +#ifdef HAVE_RUBY_IO_BUFFER_H +#include +#endif + +// Version 3 (Ruby 4.1) provides counted allocation locks, allowing independent +// futex words to retain the same buffer without unlocking each other. +#if defined(__linux__) && RUBY_IO_BUFFER_VERSION >= 3 && defined(HAVE_LINUX_FUTEX_H) && defined(HAVE_SYS_SYSCALL_H) + +#define IO_EVENT_FUTEX + +#include +#include +#include + +#ifndef FUTEX2_SIZE_U32 +#define FUTEX2_SIZE_U32 2 +#endif + +#ifndef FUTEX_32 +#define FUTEX_32 2 +#endif + +#if defined(FUTEX_WAITV_MAX) && defined(SYS_futex_waitv) +#define IO_EVENT_FUTEX_WAITV +#endif + +uint32_t *IO_Event_Futex_address(VALUE self); +uint32_t *IO_Event_Futex_acquire(VALUE self); +void IO_Event_Futex_release(VALUE self); +#ifdef IO_EVENT_FUTEX_WAITV +VALUE IO_Event_Futex_prepare_waitv(VALUE entries, struct futex_waitv *vector); +#endif +void Init_IO_Event_Futex(VALUE IO_Event); + +#endif diff --git a/ext/io/event/selector/uring.c b/ext/io/event/selector/uring.c index fe664af6..44a0d70f 100644 --- a/ext/io/event/selector/uring.c +++ b/ext/io/event/selector/uring.c @@ -3,6 +3,7 @@ #include "uring.h" #include "selector.h" +#include "../futex.h" #include "../list.h" #include "../array.h" @@ -804,6 +805,190 @@ VALUE IO_Event_Selector_URing_process_wait(VALUE self, VALUE fiber, VALUE _pid, return rb_ensure(process_wait_transfer, (VALUE)&process_wait_arguments, process_wait_ensure, (VALUE)&process_wait_arguments); } +#if defined(IO_EVENT_FUTEX) && defined(HAVE_IO_URING_PREP_FUTEX_WAIT) + +#pragma mark - Futex Wait + +struct futex_wait_arguments { + struct IO_Event_Selector_URing *selector; + struct IO_Event_Selector_URing_Waiting waiting; + VALUE futex; + uint32_t expected; + bool acquired; + bool submitted; +}; + +static VALUE futex_wait_cancel(VALUE _arguments) { + struct futex_wait_arguments *arguments = (struct futex_wait_arguments *)_arguments; + if (arguments->submitted) { + // Drain the original operation before releasing its futex. The shared + // helper tracks cancellation CQEs to prevent premature completion reuse. + IO_Event_Selector_URing_Waiting_cancel_and_wait(arguments->selector, &arguments->waiting); + } else if (arguments->waiting.completion) { + // Setup failed before an SQE referred to this completion. + IO_Event_Selector_URing_Completion_complete(arguments->selector, arguments->waiting.completion); + } + return Qnil; +} + +static VALUE futex_wait_release(VALUE _arguments) { + struct futex_wait_arguments *arguments = (struct futex_wait_arguments *)_arguments; + if (arguments->acquired) { + IO_Event_Futex_release(arguments->futex); + arguments->acquired = false; + } + return Qnil; +} + +static VALUE futex_wait_ensure(VALUE _arguments) { + return rb_ensure(futex_wait_cancel, _arguments, futex_wait_release, _arguments); +} + +static VALUE futex_wait_transfer(VALUE _arguments) { + struct futex_wait_arguments *arguments = (struct futex_wait_arguments *)_arguments; + uint32_t *address = IO_Event_Futex_acquire(arguments->futex); + arguments->acquired = true; + + // Arguments have been coerced and the futex is now held open. From here + // on, the enclosing ensure owns all completion and cancellation cleanup. + struct IO_Event_Selector_URing_Completion *completion = IO_Event_Selector_URing_Completion_acquire(arguments->selector, &arguments->waiting); + struct io_uring_sqe *sqe = io_get_sqe(arguments->selector); + io_uring_prep_futex_wait(sqe, address, arguments->expected, FUTEX_BITSET_MATCH_ANY, FUTEX2_SIZE_U32, 0); + io_uring_sqe_set_data(sqe, completion); + arguments->submitted = true; + io_uring_submit_pending(arguments->selector); + + IO_Event_Selector_loop_yield(&arguments->selector->backend); + if (arguments->waiting.completion) { + // An out-of-band resume is not a successful futex notification. + IO_Event_Selector_URing_Waiting_cancel_and_wait(arguments->selector, &arguments->waiting); + } + + int32_t result = arguments->waiting.result; + if (result >= 0) { + return Qtrue; + } else if (result == -EAGAIN || result == -ECANCELED) { + return Qfalse; + } else { + rb_syserr_fail(-result, "futex_wait_transfer:io_uring_futex_wait"); + } + + return Qnil; +} + +static VALUE IO_Event_Selector_URing_futex_wait(VALUE self, VALUE fiber, VALUE futex, VALUE expected_value) { + struct IO_Event_Selector_URing *selector = NULL; + TypedData_Get_Struct(self, struct IO_Event_Selector_URing, &IO_Event_Selector_URing_Type, selector); + struct futex_wait_arguments arguments = { + .selector = selector, + .waiting = {.fiber = fiber}, + .futex = futex, + .expected = NUM2UINT(expected_value), + }; + RB_OBJ_WRITTEN(self, Qundef, fiber); + VALUE result = rb_ensure(futex_wait_transfer, (VALUE)&arguments, futex_wait_ensure, (VALUE)&arguments); + RB_GC_GUARD(futex); + return result; +} + +#if defined(IO_EVENT_FUTEX_WAITV) && defined(HAVE_IO_URING_PREP_FUTEX_WAITV) + +#pragma mark - Futex Vector Wait + +struct futex_waitv_arguments { + struct IO_Event_Selector_URing *selector; + struct IO_Event_Selector_URing_Waiting waiting; + VALUE futexes; + struct futex_waitv *vector; + long count; + long acquired; + bool submitted; +}; + +static VALUE futex_waitv_cancel(VALUE _arguments) { + struct futex_waitv_arguments *arguments = (struct futex_waitv_arguments *)_arguments; + if (arguments->submitted) { + // Drain the original operation before releasing its futex references or + // stack-backed wait vector. The shared helper also tracks cancellation + // CQEs so the completion record cannot be reused prematurely. + IO_Event_Selector_URing_Waiting_cancel_and_wait(arguments->selector, &arguments->waiting); + } else if (arguments->waiting.completion) { + // Setup failed before an SQE referred to this completion. + IO_Event_Selector_URing_Completion_complete(arguments->selector, arguments->waiting.completion); + } + return Qnil; +} + +static VALUE futex_waitv_release(VALUE _arguments) { + struct futex_waitv_arguments *arguments = (struct futex_waitv_arguments *)_arguments; + while (arguments->acquired) { + IO_Event_Futex_release(RARRAY_AREF(arguments->futexes, --arguments->acquired)); + } + return Qnil; +} + +static VALUE futex_waitv_ensure(VALUE _arguments) { + return rb_ensure(futex_waitv_cancel, _arguments, futex_waitv_release, _arguments); +} + +static VALUE futex_waitv_transfer(VALUE _arguments) { + struct futex_waitv_arguments *arguments = (struct futex_waitv_arguments *)_arguments; + while (arguments->acquired < arguments->count) { + long index = arguments->acquired; + arguments->vector[index].uaddr = (uintptr_t)IO_Event_Futex_acquire(RARRAY_AREF(arguments->futexes, index)); + arguments->acquired += 1; + } + + // Arguments have been coerced and every futex is now held open. The + // enclosing ensure also releases entries if acquisition fails partway. + struct IO_Event_Selector_URing_Completion *completion = IO_Event_Selector_URing_Completion_acquire(arguments->selector, &arguments->waiting); + struct io_uring_sqe *sqe = io_get_sqe(arguments->selector); + io_uring_prep_futex_waitv(sqe, arguments->vector, arguments->count, 0); + io_uring_sqe_set_data(sqe, completion); + arguments->submitted = true; + io_uring_submit_pending(arguments->selector); + + IO_Event_Selector_loop_yield(&arguments->selector->backend); + if (arguments->waiting.completion) { + // An out-of-band resume is not a successful futex notification. + IO_Event_Selector_URing_Waiting_cancel_and_wait(arguments->selector, &arguments->waiting); + } + + int32_t result = arguments->waiting.result; + if (result >= 0) { + return INT2NUM(result); + } else if (result == -EAGAIN || result == -ECANCELED) { + return Qnil; + } else { + rb_syserr_fail(-result, "futex_waitv_transfer:io_uring_futex_waitv"); + } + + return Qnil; +} + +static VALUE IO_Event_Selector_URing_futex_waitv(VALUE self, VALUE fiber, VALUE entries) { + struct IO_Event_Selector_URing *selector = NULL; + TypedData_Get_Struct(self, struct IO_Event_Selector_URing, &IO_Event_Selector_URing_Type, selector); + + struct futex_waitv vector[FUTEX_WAITV_MAX]; + VALUE futexes = IO_Event_Futex_prepare_waitv(entries, vector); + struct futex_waitv_arguments arguments = { + .selector = selector, + .waiting = {.fiber = fiber}, + .futexes = futexes, + .vector = vector, + .count = RARRAY_LEN(futexes), + }; + RB_OBJ_WRITTEN(self, Qundef, fiber); + VALUE result = rb_ensure(futex_waitv_transfer, (VALUE)&arguments, futex_waitv_ensure, (VALUE)&arguments); + RB_GC_GUARD(futexes); + return result; +} + +#endif + +#endif + #pragma mark - IO#wait static inline @@ -1724,6 +1909,9 @@ VALUE IO_Event_Selector_URing_wakeup(VALUE self) { #pragma mark - Native Methods +static int IO_Event_Selector_URing_futex_supported = 0; +static int IO_Event_Selector_URing_futex_waitv_supported = 0; + static int IO_Event_Selector_URing_supported_p(void) { struct io_uring ring; @@ -1754,6 +1942,17 @@ static int IO_Event_Selector_URing_supported_p(void) { return 0; } + +#if defined(IO_EVENT_FUTEX) && defined(HAVE_IO_URING_PREP_FUTEX_WAIT) + struct io_uring_probe *probe = io_uring_get_probe_ring(&ring); + if (probe) { + IO_Event_Selector_URing_futex_supported = io_uring_opcode_supported(probe, IORING_OP_FUTEX_WAIT); +#if defined(IO_EVENT_FUTEX_WAITV) && defined(HAVE_IO_URING_PREP_FUTEX_WAITV) + IO_Event_Selector_URing_futex_waitv_supported = io_uring_opcode_supported(probe, IORING_OP_FUTEX_WAITV); +#endif + io_uring_free_probe(probe); + } +#endif io_uring_queue_exit(&ring); @@ -1803,4 +2002,16 @@ void Init_IO_Event_Selector_URing(VALUE IO_Event_Selector) { rb_define_method(IO_Event_Selector_URing, "io_close", IO_Event_Selector_URing_io_close, 1); rb_define_method(IO_Event_Selector_URing, "process_wait", IO_Event_Selector_URing_process_wait, 3); + +#if defined(IO_EVENT_FUTEX) && defined(HAVE_IO_URING_PREP_FUTEX_WAIT) + if (IO_Event_Selector_URing_futex_supported) { + rb_define_method(IO_Event_Selector_URing, "futex_wait", IO_Event_Selector_URing_futex_wait, 3); + +#if defined(IO_EVENT_FUTEX_WAITV) && defined(HAVE_IO_URING_PREP_FUTEX_WAITV) + if (IO_Event_Selector_URing_futex_waitv_supported) { + rb_define_method(IO_Event_Selector_URing, "futex_waitv", IO_Event_Selector_URing_futex_waitv, 2); + } +#endif + } +#endif } diff --git a/fixtures/io/event/test_scheduler.rb b/fixtures/io/event/test_scheduler.rb index 6c71315d..e372ea80 100644 --- a/fixtures/io/event/test_scheduler.rb +++ b/fixtures/io/event/test_scheduler.rb @@ -37,6 +37,22 @@ module Forwarders def io_close(descriptor) @selector.io_close(descriptor) end + + # Wait while the futex contains the expected value. + def futex_wait(futex, expected) + @blocked += 1 + @selector.futex_wait(Fiber.current, futex, expected) + ensure + @blocked -= 1 + end + + # Wait until any futex value changes. + def futex_waitv(entries) + @blocked += 1 + @selector.futex_waitv(Fiber.current, entries) + ensure + @blocked -= 1 + end end def initialize(selector: nil, worker_pool: nil, maximum_worker_count: nil) diff --git a/guides/futex/readme.md b/guides/futex/readme.md new file mode 100644 index 00000000..d2a9cd54 --- /dev/null +++ b/guides/futex/readme.md @@ -0,0 +1,208 @@ +# Futex Notifications + +This guide explains how to use `IO::Event::Futex` to notify threads or processes when shared state changes, without continuously polling it. + +Shared state describes what changed; the futex tells you to check it. A wake-up does not transfer a message, reserve a worker, or grant ownership of a permit. + +## Availability + +`IO::Event::Futex` is available on Linux with Ruby 4.1 or later and the required Linux build headers. It uses Ruby's counted buffer locks to retain an aligned, writable 32-bit word. The class itself does not require `io_uring`. + +There are two ways to wait: + +- Without an active fiber scheduler for the calling fiber, waits release the GVL and block the calling thread using Linux futex syscalls. +- With a scheduler, waits invoke its optional `futex_wait` or `futex_waitv` hook. A missing hook raises `NotImplementedError`; it does not silently block the event loop. + +Only the URing selector provides asynchronous futex waits. EPoll, KQueue, and Select do not. URing exposes each method only when the build's liburing and the running kernel support the corresponding operation. Check the selected backend before choosing an IPC mechanism: + +```ruby +require "io/event" + +selector = IO::Event::Selector.new(Fiber.current) +begin + futex_available = IO::Event.const_defined?(:Futex, false) + puts "Single waits: #{futex_available && selector.respond_to?(:futex_wait)}" + puts "Vector waits: #{futex_available && selector.respond_to?(:futex_waitv)}" +ensure + selector.close +end +``` + +Creating a selector does not install a fiber scheduler. A scheduler integration must forward the following hooks to its own selector, and should only expose hooks that selector supports: + +| Scheduler hook | Selector call | +| --- | --- | +| `futex_wait(futex, expected)` | `selector.futex_wait(Fiber.current, futex, expected)` | +| `futex_waitv(entries)` | `selector.futex_waitv(Fiber.current, entries)` | + +For blocking vector waits, `IO::Event::Futex.respond_to?(:wait_any)` indicates build support, not running-kernel support; the syscall can still raise `Errno::ENOSYS`. Deployment restrictions can also prevent kernel operations. Negotiate capabilities before peers begin waiting, and choose a socket, pipe, or an IPC protocol such as `async-bus` when asynchronous futex waits are unavailable. There is no automatic IPC fallback in `Futex`. + +## Binding Shared Memory + +Each futex refers to four bytes at a byte offset within an `IO::Buffer`. The address must be aligned to four bytes, and the buffer must have at least four bytes remaining at that offset. The default offset is zero. Constructing a futex preserves the word's existing value. + +```ruby +require "io/event" + +buffer = IO::Buffer.new(8) +first = IO::Event::Futex.new(buffer) +second = IO::Event::Futex.new(buffer, offset: 4) + +begin + first.value = 0 + second.value = 0 + first.signal + puts first.value # => 1 +ensure + first.close + second.close + buffer.free +end +``` + +An ordinary `IO::Buffer.new` allocation is suitable for threads in one process. Forking does not make that allocation shared between processes. For IPC, map shared storage instead: for example, each process can use `IO::Buffer.map(file, size)` on the same pre-sized file opened for reading and writing, without `IO::Buffer::PRIVATE`. Bind futexes at matching offsets in that mapping; virtual addresses need not match between processes. + +Initialize the shared words once, before peers attach. Attaching peers must not reset live state. Futexes do not discover peers or transport Ruby objects; the application supplies the shared-state layout and synchronization protocol. + +## Atomic Operations and Notifications + +| Operation | Effect and return value | +| --- | --- | +| `value` | Atomically reads the word. | +| `value = n` | Atomically stores `n`. | +| `increment(amount = 1)` | Adds `amount` and returns the new value. | +| `decrement(amount = 1)` | Subtracts `amount` and returns the new value. | +| `compare_exchange(expected, desired)` | Stores `desired` only if the word equals `expected`; returns whether it succeeded. | +| `wake(count = 1)` | Wakes at most `count` waiters without changing the word; returns the number woken. | +| `signal(count = 1)` | For a positive count, increments the word by one, then wakes at most `count` waiters; returns the word value. | + +Loads use acquire ordering, stores use release ordering, and read-modify-write operations use acquire-release ordering. A failed compare-and-exchange uses acquire ordering. Arithmetic wraps modulo `2**32`; a notification counter is not an indefinitely increasing event history. + +Only `wake` and `signal` notify sleeping waiters. Assignment, increment, decrement, and compare-and-exchange do not wake them. For a positive count, `signal` performs an atomic increment followed by a wake syscall, not one indivisible increment-and-wake operation. Its argument is the waiter count, not the increment amount. Wake counts must be nonnegative integers that fit in a C `int`. + +Both `wake(0)` and `signal(0)` are no-ops: they leave the word and waiters unchanged. `wake(0)` returns `0`; `signal(0)` returns the current word value. Closed or uninitialized futexes still raise `IOError`. + +Publish application state before notifying consumers. A futex does not make other memory accesses atomic or provide a queue's synchronization. Use atomic operations or another suitable synchronization protocol for that state, and do not mix concurrent non-atomic buffer accesses with atomic accesses to the futex word. + +## Waiting Without Missing Notifications + +`futex.wait(expected)` requires an explicit expected value and asks the kernel to wait only if the word still equals it. Comparing the word and beginning the wait is atomic with respect to futex operations. If the value has already changed, the call returns `false` without sleeping; a successful wake returns `true`. Omitting `expected` raises `ArgumentError`. + +Always recheck application state after a wait. Wake-ups can be spurious, another consumer may have taken the available work, and waking a waiter does not guarantee FIFO ordering or ownership. System errors raise exceptions. Neither wait API currently accepts a timeout argument. + +When the word is a notification counter for separate application state, the consumer must: + +1. Read the notification counter. +2. Check or attempt to consume the synchronized application state. +3. If no work is available, wait using the saved counter value. +4. Repeat after the wait returns. + +The producer publishes work before incrementing the counter and waking consumers. If a producer publishes between steps 2 and 3, the changed counter prevents the consumer from sleeping. + +### Why the Expected Value Is Required + +Reading the counter *after* checking for work can miss a notification: + +1. The consumer finds the queue empty. The notification counter is `0`. +2. The producer adds work and signals, changing the counter to `1`. Nobody is waiting yet. +3. The consumer reads the counter and gets `1`. +4. The consumer calls `wait(1)`. Since the counter is still `1`, it sleeps despite available work. Without another notification, it can remain asleep indefinitely. + +An argument-free `wait` would hide step 3 inside the method. Requiring `expected` makes the snapshot explicit: the consumer should capture `0` before checking the queue, then call `wait(0)`. In the same sequence, the kernel sees that the counter is now `1` and returns immediately. If the consumer starts waiting before the producer signals, the signal wakes it instead. + +The required argument does not enforce the ordering by itself. Calling `wait(futex.value)` after checking for work has the same race. The intended sequence is **snapshot, check state, then wait using that snapshot**, followed by another state check after the wait returns. + +### Producer and Consumer + +This example uses a thread-safe Ruby `Queue` for application state and a futex for notifications. `Queue` already supports blocking `pop`; the explicit wait here demonstrates the protocol you would use with your own shared state. This particular queue is local to one process, not an IPC queue. + +```ruby +require "io/event" + +buffer = IO::Buffer.new(4) +notification = IO::Event::Futex.new(buffer) +notification.value = 0 +queue = Thread::Queue.new + +producer = Thread.new do + ["first", "second", "third", nil].each do |message| + queue << message + notification.signal + end +end + +begin + loop do + # Snapshot before checking application state: + expected = notification.value + + begin + message = queue.pop(true) + rescue ThreadError + notification.wait(expected) + next + end + + # A nil message marks the end of this example: + break if message.nil? + puts message + end +ensure + producer.join + notification.close + buffer.free +end +``` + +The queue retains work even when a notification wakes nobody. Notifications may be coalesced; consumers must inspect the state rather than count wake-ups. A 32-bit counter can wrap back to a saved value, so protocols must account for wraparound if a consumer could miss an entire counter cycle between its snapshot and wait. + +## Waiting on Several Words + +When supported, `IO::Event::Futex.wait_any(entries)` waits on a vector of `[futex, expected]` pairs. Supply between one and `IO::Event::Futex::WAITV_LIMIT` entries. It returns a zero-based index when woken, or `nil` if any word does not match its expected value. The index is a notification hint, not an exhaustive list of changes or a claim on the associated resource. + +For example, wait until either of two initially-zero words is nonzero: + +```ruby +require "io/event" + +buffer = IO::Buffer.new(8) +first = IO::Event::Futex.new(buffer) +second = IO::Event::Futex.new(buffer, offset: 4) +first.value = second.value = 0 + +producer = Thread.new{second.signal} + +begin + while first.value == 0 && second.value == 0 + IO::Event::Futex.wait_any([[first, 0], [second, 0]]) + end + puts "At least one word changed." +ensure + producer.join + first.close + second.close + buffer.free +end +``` + +For notification counters alongside separate state, snapshot all counters before checking that state, just as with a single wait. Vector waits retain a private snapshot of the supplied futex references; mutating the original entries does not change an already-pending wait. + +## Lifetime and Cancellation + +Each futex owns one counted allocation lock from construction until `close` or finalization. Slices lock their root allocation, and multiple words can independently retain the same allocation. The lock prevents freeing, resizing, or transferring that allocation; it does not prevent writing shared state. + +`close` is idempotent, and `closed?` reports whether the futex has been released. Word operations and waits on closed or uninitialized instances raise `IOError`. Futexes cannot be copied or reinitialized. A finalizer releases the lock if a futex is collected, but explicit `close` makes release deterministic. Avoid application references from the retained buffer back to its futex, which would keep it reachable through the finalizer. + +A futex with pending waits cannot be closed: `close` raises `IOError` rather than cancelling or waking those waits. For orderly shutdown: + +1. Publish shutdown state and notify the necessary waiters, or cancel them through the scheduler. +2. Let the waits finish, including cancellation cleanup. +3. Close all futexes, then free or unmap the buffer. + +URing cancellation drains the original kernel operation before releasing its retained references or wait vector. Keep the event loop running until that cleanup completes. Exceptions still propagate; if a wait is resumed out of band without an exception, it can instead return `false` for a single wait or `nil` for a vector wait. These returns still require rechecking application state. + +Allocation locks cannot protect against an external owner releasing memory or another process truncating a mapped file. The application must preserve the underlying storage for every process that can still access it. + +## Further Reading + +The Linux [futex overview](https://man7.org/linux/man-pages/man2/futex.2.html), [FUTEX_WAIT](https://man7.org/linux/man-pages/man2/FUTEX_WAIT.2const.html), and [FUTEX_WAKE](https://man7.org/linux/man-pages/man2/FUTEX_WAKE.2const.html) documentation describe the underlying shared-memory and notification semantics. diff --git a/guides/getting-started/readme.md b/guides/getting-started/readme.md index 31d54ef8..f8b65852 100644 --- a/guides/getting-started/readme.md +++ b/guides/getting-started/readme.md @@ -76,6 +76,12 @@ puts "[main] Done" # [main] Done ``` +## Shared-Memory Notifications + +On Linux with Ruby 4.1 or later, `IO::Event::Futex` provides atomic operations and notifications on an aligned 32-bit word in an `IO::Buffer`. It supports blocking thread waits and optional asynchronous waits through the URing selector. + +See the [Futex Notifications guide](../futex/index) for capability detection, shared-memory setup, wait/recheck examples, vector waits, and buffer lifetime management. + ## Debugging The {ruby IO::Event::Debug::Selector} class adds extra validations and checks at the expense of performance. It can also log all operations. You can use this by setting the following environment variables: diff --git a/guides/links.yaml b/guides/links.yaml index 7f527b02..4a26f798 100644 --- a/guides/links.yaml +++ b/guides/links.yaml @@ -1,2 +1,4 @@ getting-started: order: 1 +futex: + order: 2 diff --git a/lib/io/event/debug/selector.rb b/lib/io/event/debug/selector.rb index a5eca436..5bcda250 100644 --- a/lib/io/event/debug/selector.rb +++ b/lib/io/event/debug/selector.rb @@ -19,6 +19,18 @@ def io_close(descriptor) log("Closing file descriptor #{descriptor}") @selector.io_close(descriptor) end + + # Wait for a futex value to change, forwarded to the underlying selector. + def futex_wait(fiber, futex, expected) + log("Waiting for futex #{futex.inspect} with value #{expected}") + @selector.futex_wait(fiber, futex, expected) + end + + # Wait for any futex value to change, forwarded to the underlying selector. + def futex_waitv(fiber, entries) + log("Waiting for futex vector #{entries.inspect}") + @selector.futex_waitv(fiber, entries) + end end # Wrap the given selector with debugging. diff --git a/readme.md b/readme.md index 84cbeb3b..c8a39d78 100644 --- a/readme.md +++ b/readme.md @@ -13,6 +13,7 @@ The initial proof-of-concept [Async](https://github.com/socketry/async) was buil Please see the [project documentation](https://socketry.github.io/io-event/) for more details. - [Getting Started](https://socketry.github.io/io-event/guides/getting-started/index) - This guide explains how to use `io-event` for non-blocking IO. + - [Futex Notifications](https://socketry.github.io/io-event/guides/futex/index) - This guide explains shared-memory notifications, capability detection, and safe waiting patterns. ## Releases diff --git a/releases.md b/releases.md index 5b186595..b3a99aec 100644 --- a/releases.md +++ b/releases.md @@ -1,5 +1,9 @@ # Releases +## Unreleased + + - Add `IO::Event::Futex` on Linux with Ruby 4.1+, including atomic value operations and blocking and scheduler-aware single and vector waits over shared memory. Futexes retain a counted allocation lock until explicitly closed or finalized; pending waits prevent closing until completion or cancellation has finished. + ## v1.22.1 - Fix an infinite loop in the ready queue flush when a queued fiber is resumed out of band, e.g. by a stale `unblock` racing a timeout, while another fiber re-queues itself on every iteration. diff --git a/test/io/event/debug/selector.rb b/test/io/event/debug/selector.rb index 5b30badb..5e11eb4a 100644 --- a/test/io/event/debug/selector.rb +++ b/test/io/event/debug/selector.rb @@ -93,6 +93,16 @@ def io_close(descriptor) :calls_io_close end + def futex_wait(fiber, futex, expected) + @calls << [:futex_wait, fiber, futex, expected] + :calls_futex_wait + end + + def futex_waitv(fiber, entries) + @calls << [:futex_waitv, fiber, entries] + :calls_futex_waitv + end + def select(duration = nil) @calls << [:select, duration] :calls_select @@ -153,6 +163,8 @@ def select(duration = nil) end expect(selector.io_close(input.fileno)).to be == :calls_io_close expect(selector.respond_to?(:io_close)).to be == true + expect(selector.futex_wait(fiber, :futex, 1)).to be == :calls_futex_wait + expect(selector.futex_waitv(fiber, [[:futex, 1]])).to be == :calls_futex_waitv expect(selector.select(0)).to be == :calls_select ensure input&.close diff --git a/test/io/event/futex.rb b/test/io/event/futex.rb new file mode 100644 index 00000000..3e59b4d3 --- /dev/null +++ b/test/io/event/futex.rb @@ -0,0 +1,646 @@ +# frozen_string_literal: true + +# Released under the MIT License. +# Copyright, 2026, by Samuel Williams. + +require "io/event" +require "io/event/test_scheduler" + +return unless defined?(IO::Event::Futex) + +describe IO::Event::Futex do + let(:buffer) {IO::Buffer.new(8)} + let(:futex) {subject.new(buffer)} + let(:uring_selector) do + unless defined?(IO::Event::Selector::URing) + skip "io_uring is not available" + end + + selector = IO::Event::Selector::URing.new(Fiber.current) + unless selector.respond_to?(:futex_wait) + selector.close + skip "io_uring futex operations are not available" + end + + selector + end + let(:waitv_selector) do + selector = uring_selector + unless selector.respond_to?(:futex_waitv) + selector.close + skip "io_uring futex waitv operations are not available" + end + + selector + end + + with "buffer lifetime" do + it "locks the allocation until closed" do + instance = subject.new(buffer) + expect(buffer).to be(:locked?) + expect{buffer.free}.to raise_exception(IO::Buffer::LockedError) + expect{buffer.resize(16)}.to raise_exception(IO::Buffer::LockedError) + instance.close + expect(instance).to be(:closed?) + expect(buffer).not.to be(:locked?) + buffer.free + end + + it "owns one independent lock per futex" do + first = subject.new(buffer) + second = subject.new(buffer, offset: 4) + first.close + first.close + expect(buffer).to be(:locked?) + expect(second.increment).to be == 1 + second.close + expect(buffer).not.to be(:locked?) + end + + it "locks the root allocation when bound to a slice" do + instance = subject.new(buffer.slice(4, 4)) + expect{buffer.free}.to raise_exception(IO::Buffer::LockedError) + instance.close + expect(buffer).not.to be(:locked?) + end + + it "does not leak a lock when initialization fails" do + expect{subject.new(buffer, offset: 1)}.to raise_exception(ArgumentError) + expect{subject.new(buffer, offset: 8)}.to raise_exception(RangeError) + expect(buffer).not.to be(:locked?) + end + + it "does not allow reinitialization or copying" do + instance = subject.new(buffer) + expect{instance.send(:initialize, buffer)}.to raise_exception(RuntimeError) + expect{instance.dup}.to raise_exception(TypeError) + expect{instance.clone}.to raise_exception(TypeError) + instance.close + expect(buffer).not.to be(:locked?) + end + + it "does not unlock the original when rejected copies are collected" do + instance = subject.new(buffer) + Thread.new do + expect{instance.dup}.to raise_exception(TypeError) + expect{instance.clone}.to raise_exception(TypeError) + end.join + 3.times{GC.start(full_mark: true, immediate_sweep: true)} + expect(buffer).to be(:locked?) + expect(instance.increment).to be == 1 + instance.close + end + + it "releases locks when futexes are collected" do + # A separate thread removes conservative C-stack references before GC: + Thread.new{subject.new(buffer); nil}.join + 10.times do + GC.start(full_mark: true, immediate_sweep: true) + break unless buffer.locked? + end + expect(buffer).not.to be(:locked?) + end + + it "does not unlock another futex when a closed instance is collected" do + instance = subject.new(buffer) + Thread.new{subject.new(buffer, offset: 4).close}.join + 3.times{GC.start(full_mark: true, immediate_sweep: true)} + expect(buffer).to be(:locked?) + expect(instance.increment).to be == 1 + instance.close + expect(buffer).not.to be(:locked?) + end + + it "retains the buffer across compaction" do + instance = subject.new(buffer) + GC.verify_compaction_references(double_heap: true, toward: :empty) + expect(instance.increment).to be == 1 + instance.close + expect(buffer).not.to be(:locked?) + end + end + + with "closed or uninitialized futexes" do + [ + [:value], [:value=, 1], [:increment], [:decrement], + [:compare_exchange, 0, 1], [:wake], [:wake, 0], [:signal], [:signal, 0], [:wait, 0] + ].each do |arguments| + it "rejects #{arguments.inspect} after close", unique: "closed #{arguments.inspect}" do + futex.close + expect{futex.public_send(*arguments)}.to raise_exception(IOError) + end + + it "rejects #{arguments.inspect} before initialization", unique: "uninitialized #{arguments.inspect}" do + expect{subject.allocate.public_send(*arguments)}.to raise_exception(IOError) + end + end + end + + with "argument coercion" do + it "rechecks the futex after numeric conversion closes it" do + instance = futex + value = Object.new + value.define_singleton_method(:to_int) do + instance.close + 1 + end + expect{instance.value = value}.to raise_exception(IOError) + end + + it "validates the wake count before changing the value" do + expect{futex.signal(-1)}.to raise_exception(ArgumentError) + expect(futex.value).to be == 0 + end + end + + with "#value" do + it "stores and loads the value atomically" do + futex.value = 42 + expect(futex.value).to be == 42 + end + end + + with "zero-count notifications" do + [:wake, :signal].each do |operation| + it "leaves the word and waiters unchanged for #{operation}(0)", unique: operation.to_s do + futex.value = 7 + thread = Thread.new{futex.wait(7)} + Thread.pass while thread.status == "run" + + expect(futex.public_send(operation, 0)).to be == (operation == :wake ? 0 : 7) + expect(futex.value).to be == 7 + expect(thread.join(0.02)).to be_nil + ensure + thread&.kill&.join + end + end + end + + with "#increment" do + it "increments the value" do + expect(futex.increment).to be == 1 + expect(futex.increment(2)).to be == 3 + end + end + + with "#decrement" do + it "decrements the value" do + futex.value = 3 + + expect(futex.decrement).to be == 2 + expect(futex.decrement(2)).to be == 0 + end + end + + with "#compare_exchange" do + it "exchanges a matching value" do + futex.value = 2 + + expect(futex.compare_exchange(2, 1)).to be == true + expect(futex.value).to be == 1 + end + + it "does not exchange a different value" do + futex.value = 2 + + expect(futex.compare_exchange(1, 0)).to be == false + expect(futex.value).to be == 2 + end + end + + with "offset:" do + it "can address independent words in one buffer" do + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + + first.value = 1 + second.value = 2 + + expect(first.value).to be == 1 + expect(second.value).to be == 2 + end + + it "rejects unaligned offsets" do + expect do + subject.new(buffer, offset: 1) + end.to raise_exception(ArgumentError) + end + + it "rejects offsets outside the buffer" do + expect do + subject.new(buffer, offset: 8) + end.to raise_exception(RangeError) + end + end + + with "#wait" do + it "requires an explicit expected value" do + expect{futex.wait}.to raise_exception(ArgumentError) + end + + it "cannot be closed during a blocking wait" do + instance = futex + thread = Thread.new{instance.wait(0)} + Thread.pass while thread.status == "run" + expect{instance.close}.to raise_exception(IOError) + instance.signal + thread.join + instance.close + expect(buffer).not.to be(:locked?) + ensure + thread&.kill&.join + end + + it "releases a blocking wait when its thread is interrupted" do + instance = futex + thread = Thread.new{instance.wait(0)} + Thread.pass while thread.status == "run" + thread.kill.join + instance.close + expect(buffer).not.to be(:locked?) + ensure + thread&.kill&.join + end + + it "cannot be closed while an asynchronous wait is pending" do + selector = uring_selector + fiber = Fiber.new{selector.futex_wait(Fiber.current, futex, 0)} + fiber.transfer + expect{futex.close}.to raise_exception(IOError) + futex.signal + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(fiber).not.to be(:alive?) + futex.close + expect(buffer).not.to be(:locked?) + ensure + selector&.close + end + + it "drains cancellation before allowing close" do + selector = uring_selector + error = RuntimeError.new("cancel futex") + caught = nil + fiber = Fiber.new do + selector.futex_wait(Fiber.current, futex, 0) + rescue RuntimeError => exception + caught = exception + end + fiber.transfer + selector.select(0) + fiber.raise(error) + expect{futex.close}.to raise_exception(IOError) if fiber.alive? + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(caught).to be_equal(error) + expect(fiber).not.to be(:alive?) + futex.close + expect(buffer).not.to be(:locked?) + # Drain the cancellation CQE as well as the original operation: + selector.select(0) + ensure + selector&.close + end + + it "does not report an out-of-band resume as a notification" do + selector = uring_selector + result = :pending + fiber = Fiber.new{result = selector.futex_wait(Fiber.current, futex, 0)} + fiber.transfer + selector.select(0) + fiber.transfer + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(result).to be == false + futex.close + ensure + selector&.close + end + + it "keeps the selector usable after invalid arguments" do + selector = uring_selector + fiber = Fiber.new do + expect{selector.futex_wait(Fiber.current, futex, Object.new)}.to raise_exception(TypeError) + expect{selector.futex_wait(Fiber.current, Object.new, 0)}.to raise_exception(TypeError) + GC.start + futex.value = 1 + expect(selector.futex_wait(Fiber.current, futex, 0)).to be == false + end + fiber.transfer + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(fiber).not.to be(:alive?) + futex.close + ensure + selector&.close + end + + it "waits without blocking other Ruby threads when no scheduler is installed" do + thread = Thread.new do + sleep 0.01 + futex.signal + end + + expect(futex.wait(0)).to be == true + expect(futex.value).to be == 1 + ensure + thread&.join + end + + it "does not wait for a notification published after the snapshot" do + expected = futex.value + futex.signal + expect(futex.wait(expected)).to be == false + end + + it "waits asynchronously for a signal" do + selector = uring_selector + result = nil + + fiber = Fiber.new do + result = selector.futex_wait(Fiber.current, futex, 0) + end + fiber.transfer + + thread = Thread.new do + sleep 0.01 + futex.signal + end + + selector.select(1) + thread.join + + expect(result).to be == true + expect(futex.value).to be == 1 + ensure + selector&.close + thread&.join + end + + it "uses the current scheduler" do + selector = uring_selector + scheduler = IO::Event::TestScheduler.new(selector: selector) + result = nil + + Fiber.set_scheduler(scheduler) + Fiber.schedule do + result = futex.wait(0) + end + + thread = Thread.new do + sleep 0.01 + futex.signal + end + + scheduler.run + + expect(result).to be == true + ensure + Fiber.set_scheduler(nil) + thread&.join + end + + it "does not wait when the value has changed" do + selector = uring_selector + futex.value = 1 + result = nil + + fiber = Fiber.new do + result = selector.futex_wait(Fiber.current, futex, 0) + end + fiber.transfer + selector.select(1) + + expect(result).to be == false + ensure + selector&.close + end + end + + if IO::Event::Futex.respond_to?(:wait_any) + with ".wait_any" do + it "releases earlier entries when a later entry is invalid" do + expect{subject.wait_any([[futex, 0], [Object.new, 0]])}.to raise_exception(TypeError) + futex.close + expect(buffer).not.to be(:locked?) + end + + it "protects every futex during a blocking vector wait" do + first = subject.new(buffer) + second = subject.new(buffer, offset: 4) + thread = Thread.new{subject.wait_any([[first, 0], [second, 0]])} + Thread.pass while thread.status == "run" + expect{first.close}.to raise_exception(IOError) + expect{second.close}.to raise_exception(IOError) + thread.kill.join + first.close + second.close + expect(buffer).not.to be(:locked?) + ensure + thread&.kill&.join + end + + it "retains its own entries until vector cancellation completes" do + selector = waitv_selector + first = subject.new(buffer) + second = subject.new(buffer, offset: 4) + entries = [[first, 0], [second, 0]] + caught = nil + error = RuntimeError.new("cancel vector") + fiber = Fiber.new do + selector.futex_waitv(Fiber.current, entries) + rescue RuntimeError => exception + caught = exception + end + fiber.transfer + entries.clear + GC.verify_compaction_references(double_heap: true, toward: :empty) + expect{first.close}.to raise_exception(IOError) + expect{second.close}.to raise_exception(IOError) + selector.select(0) + fiber.raise(error) + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(caught).to be_equal(error) + expect(fiber).not.to be(:alive?) + first.close + second.close + expect(buffer).not.to be(:locked?) + selector.select(0) + ensure + selector&.close + end + + it "keeps the selector usable after partial vector setup fails" do + selector = waitv_selector + fiber = Fiber.new do + expect{selector.futex_waitv(Fiber.current, [[futex, 0], [Object.new, 0]])}.to raise_exception(TypeError) + futex.close + end + fiber.transfer + GC.start + selector.select(0) + expect(buffer).not.to be(:locked?) + ensure + selector&.close + end + + it "returns nil for an out-of-band vector resume" do + selector = waitv_selector + first = subject.new(buffer) + second = subject.new(buffer, offset: 4) + result = :pending + fiber = Fiber.new{result = selector.futex_waitv(Fiber.current, [[first, 0], [second, 0]])} + fiber.transfer + selector.select(0) + fiber.transfer + 10.times do + selector.select(0.1) + break unless fiber.alive? + end + expect(result).to be_nil + expect(fiber).not.to be(:alive?) + first.close + second.close + expect(buffer).not.to be(:locked?) + selector.select(0) + ensure + selector&.close + end + + it "returns index zero for a single-entry vector notification" do + selector = waitv_selector + result = nil + fiber = Fiber.new{result = selector.futex_waitv(Fiber.current, [[futex, 0]])} + fiber.transfer + thread = Thread.new do + sleep 0.01 + futex.signal + end + selector.select(1) + thread.join + expect(result).to be == 0 + futex.close + expect(buffer).not.to be(:locked?) + ensure + selector&.close + thread&.join + end + + it "exposes the maximum number of wait entries" do + expect(subject::WAITV_LIMIT).to be == 128 + end + + it "rejects more than the maximum number of wait entries" do + entries = Array.new(subject::WAITV_LIMIT + 1){[futex, 0]} + + expect do + subject.wait_any(entries) + end.to raise_exception(ArgumentError) + end + + it "waits without blocking other Ruby threads when no scheduler is installed" do + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + + thread = Thread.new do + sleep 0.01 + second.signal + end + + expect(subject.wait_any([[first, 0], [second, 0]])).to be == 1 + ensure + thread&.join + end + + it "does not wait without a scheduler when a value has changed" do + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + second.value = 1 + + expect(subject.wait_any([[first, 0], [second, 0]])).to be_nil + end + + it "waits asynchronously for any futex to be signalled" do + selector = waitv_selector + + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + result = nil + + fiber = Fiber.new do + result = selector.futex_waitv(Fiber.current, [[first, 0], [second, 0]]) + end + fiber.transfer + + thread = Thread.new do + sleep 0.01 + second.signal + end + + selector.select(1) + thread.join + + expect(result).to be == 1 + ensure + selector&.close + thread&.join + end + + it "uses the current scheduler" do + selector = waitv_selector + + scheduler = IO::Event::TestScheduler.new(selector: selector) + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + result = nil + + Fiber.set_scheduler(scheduler) + Fiber.schedule do + result = subject.wait_any([[first, 0], [second, 0]]) + end + + thread = Thread.new do + sleep 0.01 + second.signal + end + + scheduler.run + + expect(result).to be == 1 + ensure + Fiber.set_scheduler(nil) + thread&.join + end + + it "returns nil when a value has changed" do + selector = waitv_selector + + first = subject.new(buffer, offset: 0) + second = subject.new(buffer, offset: 4) + second.value = 1 + result = :waiting + + fiber = Fiber.new do + result = selector.futex_waitv(Fiber.current, [[first, 0], [second, 0]]) + end + fiber.transfer + selector.select(1) + + expect(result).to be_nil + ensure + selector&.close + end + end + end +end