Class: CarbonFiber::Native::Selector
- Inherits:
-
Object
- Object
- CarbonFiber::Native::Selector
- Defined in:
- lib/carbon_fiber/native/fallback.rb
Overview
Pure-Ruby fallback selector using threads and condition variables.
Loaded automatically when the native Zig extension is unavailable. Provides the same Selector API as the native implementation so the Scheduler and Async adapter work unchanged.
Direct Known Subclasses
Instance Attribute Summary collapse
-
#io_contract_v4 ⇒ Object
writeonly
Accepted for interface parity with the native selector; the fallback has no native I/O paths, so the contract generation changes nothing.
Instance Method Summary collapse
-
#block(fiber, timeout = nil) ⇒ Object
Suspend the current fiber until unblocked or timed out.
-
#cancel_block_timer(fiber) ⇒ Object
Called from Scheduler#fiber_done's ensure block.
-
#cancel_timer(token) ⇒ Object
Cancel a pending timer by token.
-
#destroy ⇒ Object
No-op; nothing to release.
-
#initialize(loop_fiber) ⇒ Selector
constructor
A new instance of Selector.
-
#io_close(fd, exception) ⇒ Object
Cancel pending waiters on a closed descriptor.
-
#io_read(_fd, _buffer, _length, _offset) ⇒ Object
Returns nil; the Scheduler handles io_read via background thread.
-
#io_wait(fiber, fd, events) ⇒ Object
Wait for read readiness on a file descriptor via IO.select on a background thread.
-
#io_wait_with_timeout(fiber, fd, events, timeout) ⇒ Object
Like #io_wait but with a timeout.
-
#io_write(_fd, _buffer, _length, _offset) ⇒ Object
Returns nil; the Scheduler handles io_write via background thread.
-
#kernel_sleep(duration = nil) ⇒ Object
Mirrors
Selector#kernel_sleepon the native side soScheduler#kernel_sleepcan delegate to@selector.kernel_sleepin both paths. -
#parked? ⇒ Boolean
True while a fiber is parked in block() or do_io_wait.
-
#pending? ⇒ Boolean
Whether there is pending work.
-
#poll_readable_now(fd) ⇒ Object
Non-destructive check if a descriptor has data available to read.
-
#process_wait(_fiber, _pid, _flags) ⇒ Object
Returns nil; the Scheduler handles process_wait via background thread.
-
#push(fiber) ⇒ Object
Enqueue a fiber into the ready queue.
-
#raise(fiber, exception) ⇒ Object
Enqueue an exception delivery to a fiber.
-
#raise_after(fiber, exception, duration) ⇒ Object
Schedule an exception to be raised on a fiber after
durationseconds. -
#resume(fiber, value) ⇒ Object
Enqueue a fiber with a return value.
-
#select(timeout = nil) ⇒ Object
Run one event loop iteration.
-
#transfer ⇒ Object
Transfer control to the event loop fiber.
-
#unblock(fiber) ⇒ Object
Resume a fiber previously suspended by #block.
-
#wakeup ⇒ Object
Wake the event loop.
-
#yield ⇒ Object
Transfer to the event loop fiber.
Constructor Details
#initialize(loop_fiber) ⇒ Selector
Returns a new instance of Selector.
14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 |
# File 'lib/carbon_fiber/native/fallback.rb', line 14 def initialize(loop_fiber) @loop_fiber = loop_fiber @mutex = Thread::Mutex.new @cv = Thread::ConditionVariable.new @ready = [] @timers = {} @next_timer = 1 @read_waits = {} @next_wait_token = 1 # Fibers voluntarily parked in block() or do_io_wait, mapped to the # sleep timer token (or nil). flush_ready consults this set to decide # whether a fiber whose transfer returned unexpectedly was interrupted # mid-execution (Ruby 4.0 Fiber#raise bypass) and needs re-queueing. @blocked_fibers = {} end |
Instance Attribute Details
#io_contract_v4=(value) ⇒ Object (writeonly)
Accepted for interface parity with the native selector; the fallback has no native I/O paths, so the contract generation changes nothing.
38 39 40 |
# File 'lib/carbon_fiber/native/fallback.rb', line 38 def io_contract_v4=(value) @io_contract_v4 = value end |
Instance Method Details
#block(fiber, timeout = nil) ⇒ Object
Suspend the current fiber until unblocked or timed out.
161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 |
# File 'lib/carbon_fiber/native/fallback.rb', line 161 def block(fiber, timeout = nil) token = nil token = resume_after(fiber, timeout, false) if timeout @mutex.synchronize { @blocked_fibers[fiber] = token } result = @loop_fiber.transfer # Normal wakeup path: drop the tracking entry and cancel any still- # armed sleep timer. If raise() ran first, it already cancelled the # timer and zeroed the token; the raise-unwind path then relies on # cancel_block_timer (invoked from fiber_done) to remove the entry. @mutex.synchronize do stored = @blocked_fibers.delete(fiber) @timers.delete(stored) if stored end result end |
#cancel_block_timer(fiber) ⇒ Object
Called from Scheduler#fiber_done's ensure block. Removes the fiber from the blocked set and cancels any still-armed sleep timer. This is the only cleanup path when a raise() unwinds the fiber past block()'s normal return.
198 199 200 201 202 203 |
# File 'lib/carbon_fiber/native/fallback.rb', line 198 def cancel_block_timer(fiber) @mutex.synchronize do stored = @blocked_fibers.delete(fiber) @timers.delete(stored) if stored end end |
#cancel_timer(token) ⇒ Object
Cancel a pending timer by token.
190 191 192 |
# File 'lib/carbon_fiber/native/fallback.rb', line 190 def cancel_timer(token) @mutex.synchronize { !!@timers.delete(token) } end |
#destroy ⇒ Object
No-op; nothing to release.
32 33 34 |
# File 'lib/carbon_fiber/native/fallback.rb', line 32 def destroy true end |
#io_close(fd, exception) ⇒ Object
Cancel pending waiters on a closed descriptor.
222 223 224 225 226 227 228 229 230 231 232 233 234 235 |
# File 'lib/carbon_fiber/native/fallback.rb', line 222 def io_close(fd, exception) woke = false @mutex.synchronize do wait = @read_waits.delete(fd) if wait @ready << [:raise, wait[:fiber], exception, true] woke = true @cv.signal end end woke end |
#io_read(_fd, _buffer, _length, _offset) ⇒ Object
Returns nil; the Scheduler handles io_read via background thread.
243 244 245 |
# File 'lib/carbon_fiber/native/fallback.rb', line 243 def io_read(_fd, _buffer, _length, _offset) nil end |
#io_wait(fiber, fd, events) ⇒ Object
Wait for read readiness on a file descriptor via IO.select on a background thread. Returns nil for non-READABLE events (handled by the Scheduler fallback).
208 209 210 211 212 |
# File 'lib/carbon_fiber/native/fallback.rb', line 208 def io_wait(fiber, fd, events) return nil unless events == IO::READABLE do_io_wait(fiber, fd, nil) end |
#io_wait_with_timeout(fiber, fd, events, timeout) ⇒ Object
Like #io_wait but with a timeout.
215 216 217 218 219 |
# File 'lib/carbon_fiber/native/fallback.rb', line 215 def io_wait_with_timeout(fiber, fd, events, timeout) return nil unless events == IO::READABLE do_io_wait(fiber, fd, timeout) end |
#io_write(_fd, _buffer, _length, _offset) ⇒ Object
Returns nil; the Scheduler handles io_write via background thread.
248 249 250 |
# File 'lib/carbon_fiber/native/fallback.rb', line 248 def io_write(_fd, _buffer, _length, _offset) nil end |
#kernel_sleep(duration = nil) ⇒ Object
Mirrors Selector#kernel_sleep on the native side so
Scheduler#kernel_sleep can delegate to @selector.kernel_sleep
in both paths. Branches on the duration: nil parks the fiber on
the loop without a timer, non-positive yields, positive parks on
a native timer for duration seconds.
149 150 151 152 153 154 155 156 157 158 |
# File 'lib/carbon_fiber/native/fallback.rb', line 149 def kernel_sleep(duration = nil) if duration.nil? transfer elsif duration <= 0 self.yield else block(Fiber.current, duration) end true end |
#parked? ⇒ Boolean
True while a fiber is parked in block() or do_io_wait. Nothing in the ready list, timers, or read waits refers to it, so pending? is false, yet select must sleep on the condition variable rather than return: the wake-up comes from another thread (a background operation's resume, an unblock, or wakeup), and returning at once makes Scheduler#run spin with the GVL held. On the main Ractor the timer thread preempts that spin; inside a non-main Ractor Ruby's M:N scheduler does not, and the thread that would wake the fiber never runs.
140 141 142 |
# File 'lib/carbon_fiber/native/fallback.rb', line 140 def parked? @mutex.synchronize { @blocked_fibers.any? } end |
#pending? ⇒ Boolean
Returns whether there is pending work.
41 42 43 |
# File 'lib/carbon_fiber/native/fallback.rb', line 41 def pending? @mutex.synchronize { @ready.any? || @timers.any? || @read_waits.any? } end |
#poll_readable_now(fd) ⇒ Object
Non-destructive check if a descriptor has data available to read.
253 254 255 256 257 258 259 260 261 |
# File 'lib/carbon_fiber/native/fallback.rb', line 253 def poll_readable_now(fd) io = IO.new(fd, autoclose: false) ready = IO.select([io], nil, nil, 0) !!ready rescue IOError, SystemCallError false ensure io.close if io && !io.closed? end |
#process_wait(_fiber, _pid, _flags) ⇒ Object
Returns nil; the Scheduler handles process_wait via background thread.
238 239 240 |
# File 'lib/carbon_fiber/native/fallback.rb', line 238 def process_wait(_fiber, _pid, _flags) nil end |
#push(fiber) ⇒ Object
Enqueue a fiber into the ready queue.
46 47 48 49 50 51 52 |
# File 'lib/carbon_fiber/native/fallback.rb', line 46 def push(fiber) @mutex.synchronize do @ready << [:resume, fiber, nil, false] @cv.signal end fiber end |
#raise(fiber, exception) ⇒ Object
Enqueue an exception delivery to a fiber.
64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 |
# File 'lib/carbon_fiber/native/fallback.rb', line 64 def raise(fiber, exception) @mutex.synchronize do # Cancel any armed sleep timer for this fiber so its block() wakeup # doesn't spuriously fire after the raise. Zero the token but leave # the blocked_fibers entry—it's removed by cancel_block_timer when # the fiber's ensure runs, so flush_ready's re-queue check still # correctly treats the fiber as "parked" until it exits. if @blocked_fibers.key?(fiber) token = @blocked_fibers[fiber] @timers.delete(token) if token @blocked_fibers[fiber] = nil end @ready << [:raise, fiber, exception, true] @cv.signal end fiber end |
#raise_after(fiber, exception, duration) ⇒ Object
Schedule an exception to be raised on a fiber after duration seconds.
185 186 187 |
# File 'lib/carbon_fiber/native/fallback.rb', line 185 def raise_after(fiber, exception, duration) schedule_timer(duration, :raise, fiber, exception) end |
#resume(fiber, value) ⇒ Object
Enqueue a fiber with a return value.
55 56 57 58 59 60 61 |
# File 'lib/carbon_fiber/native/fallback.rb', line 55 def resume(fiber, value) @mutex.synchronize do @ready << [:resume, fiber, value, true] @cv.signal end fiber end |
#select(timeout = nil) ⇒ Object
Run one event loop iteration.
105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 |
# File 'lib/carbon_fiber/native/fallback.rb', line 105 def select(timeout = nil) flush_ready return 0 unless pending? || parked? deadline = next_wait_deadline(timeout) @mutex.synchronize do until @ready.any? collect_expired_timers_locked break if @ready.any? if deadline remaining = deadline - monotonic_time break if remaining <= 0 @cv.wait(@mutex, remaining) else @cv.wait(@mutex) end end end collect_expired_timers flush_ready end |
#transfer ⇒ Object
Transfer control to the event loop fiber.
89 90 91 92 93 |
# File 'lib/carbon_fiber/native/fallback.rb', line 89 def transfer return nil if Fiber.current.equal?(@loop_fiber) @loop_fiber.transfer end |
#unblock(fiber) ⇒ Object
Resume a fiber previously suspended by #block.
180 181 182 |
# File 'lib/carbon_fiber/native/fallback.rb', line 180 def unblock(fiber) resume(fiber, true) end |
#wakeup ⇒ Object
Wake the event loop.
83 84 85 86 |
# File 'lib/carbon_fiber/native/fallback.rb', line 83 def wakeup @mutex.synchronize { @cv.signal } true end |
#yield ⇒ Object
Transfer to the event loop fiber. flush_ready's re-queue logic puts us back in the ready queue on the next pass—no explicit self-push needed, and avoiding it prevents duplicate ready entries.
98 99 100 101 102 |
# File 'lib/carbon_fiber/native/fallback.rb', line 98 def yield return nil if Fiber.current.equal?(@loop_fiber) @loop_fiber.transfer end |