Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 16 additions & 2 deletions core/src/main/java/org/jruby/FiberScheduler.java
Original file line number Diff line number Diff line change
Expand Up @@ -5,6 +5,7 @@
import org.jruby.runtime.ThreadContext;
import org.jruby.runtime.builtin.IRubyObject;
import org.jruby.util.cli.Options;
import org.jruby.util.io.OpenFile;

import java.nio.ByteBuffer;

Expand Down Expand Up @@ -47,9 +48,22 @@ public static IRubyObject unblock(ThreadContext context, IRubyObject scheduler,
return Helpers.invoke(context, scheduler, "unblock", blocker, fiber);
}

// MRI: rb_fiber_scheduler_io_wait
// MRI: rb_fiber_scheduler_io_wait, which registers the wait so closing the IO can interrupt it
public static IRubyObject ioWait(ThreadContext context, IRubyObject scheduler, IRubyObject io, IRubyObject events, IRubyObject timeout) {
return Helpers.invoke(context, scheduler, "io_wait", io, events, timeout);
OpenFile fptr = io instanceof RubyIO rubyIO ? rubyIO.getOpenFile() : null;
if (fptr == null) return Helpers.invoke(context, scheduler, "io_wait", io, events, timeout);

OpenFile.SchedulerWaiter waiter = fptr.addSchedulerWaiter(scheduler, context.getFiberCurrentThread(), context.getFiber());
try {
return Helpers.invoke(context, scheduler, "io_wait", io, events, timeout);
} finally {
fptr.removeSchedulerWaiter(waiter);
}
}

// MRI: rb_fiber_scheduler_fiber_interrupt, which returns undef (null here) if the scheduler has no fiber_interrupt
public static IRubyObject fiberInterrupt(ThreadContext context, IRubyObject scheduler, IRubyObject fiber, IRubyObject exception) {
return Helpers.invokeChecked(context, scheduler, "fiber_interrupt", fiber, exception);
}

// MRI: rb_fiber_scheduler_io_wait_readable
Expand Down
1 change: 1 addition & 0 deletions core/src/main/java/org/jruby/RubyIO.java
Original file line number Diff line number Diff line change
Expand Up @@ -2411,6 +2411,7 @@ protected IRubyObject rbIoClose(ThreadContext context) {
fptr.interruptBlockingThreads(context);
try {
fptr.unlock();
fptr.interruptSchedulerWaiters(context);
fptr.waitForBlockingThreads(context);
} finally {
fptr.lock();
Expand Down
66 changes: 66 additions & 0 deletions core/src/main/java/org/jruby/util/io/OpenFile.java
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,8 @@ public static class Buffer {
private final Ruby runtime;

protected volatile Set<RubyThread> blockingThreads;
// fibers parked in a fiber scheduler's io_wait on this IO. MRI: rb_io's blocking_operations
private volatile Set<SchedulerWaiter> schedulerWaiters;

private final Ptr spPtr = new Ptr();
private final Ptr dpPtr = new Ptr();
Expand Down Expand Up @@ -2941,6 +2943,70 @@ public void removeBlockingThread(RubyThread thread) {
}
}

public record SchedulerWaiter(IRubyObject scheduler, RubyThread thread, IRubyObject fiber) {}

/**
* Record a fiber waiting on this IO through the fiber scheduler, so closing the IO can interrupt it.
*/
public SchedulerWaiter addSchedulerWaiter(IRubyObject scheduler, RubyThread thread, IRubyObject fiber) {
Set<SchedulerWaiter> schedulerWaiters = this.schedulerWaiters;

if (schedulerWaiters == null) {
synchronized (this) {
schedulerWaiters = this.schedulerWaiters;
if (schedulerWaiters == null) {
this.schedulerWaiters = schedulerWaiters = new HashSet<>(1);
}
}
}

SchedulerWaiter waiter = new SchedulerWaiter(scheduler, thread, fiber);
synchronized (schedulerWaiters) {
schedulerWaiters.add(waiter);
}
return waiter;
}

public void removeSchedulerWaiter(SchedulerWaiter waiter) {
Set<SchedulerWaiter> schedulerWaiters = this.schedulerWaiters;

synchronized (schedulerWaiters) {
schedulerWaiters.remove(waiter);
}
}

/**
* Fire an IOError in all fibers waiting on this IO through a fiber scheduler. Call without holding the IO lock,
* since the scheduler may switch to a waiting fiber, which needs the lock to unwind.
*/
// MRI: rb_thread_io_close_interrupt, which hands fibers waiting through a scheduler to fiber_interrupt,
// falling back to a pending interrupt on the waiter's thread
public void interruptSchedulerWaiters(ThreadContext context) {
Set<SchedulerWaiter> schedulerWaiters = this.schedulerWaiters;

if (schedulerWaiters == null) return;

SchedulerWaiter[] waiters;
synchronized (schedulerWaiters) {
waiters = schedulerWaiters.toArray(SchedulerWaiter[]::new);
}

for (SchedulerWaiter waiter : waiters) {
if (waiter.fiber() == context.getFiber()) continue;

// an earlier fiber_interrupt may have switched fibers, letting this one finish its wait
synchronized (schedulerWaiters) {
if (!schedulerWaiters.contains(waiter)) continue;
}

RubyException error = streamClosedInParallelError(runtime);
IRubyObject result = FiberScheduler.fiberInterrupt(context, waiter.scheduler(), waiter.fiber(), error);

// no fiber_interrupt hook, so raise in the waiter's thread instead, as MRI does
if (result == null) waiter.thread().raise(error);
}
}

/**
* Fire an IOError in all threads blocking on this IO object
*/
Expand Down
1 change: 0 additions & 1 deletion test/mri/excludes/TestFiberIOClose.rb
Original file line number Diff line number Diff line change
@@ -1,2 +1 @@
exclude :test_io_close_across_fibers, "hangs: close does not call the scheduler's fiber_interrupt"
exclude :test_io_close_blocking_thread, "on Linux, a thread blocked reading the IO sees EOF or stays blocked instead of raising IOError when it is closed"
Loading