Skip to content

Commit 2df5a69

Browse files
authored
Add Fiber#kill, similar to Thread#kill. (#7823)
1 parent b695f58 commit 2df5a69

6 files changed

Lines changed: 179 additions & 27 deletions

File tree

common.mk

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3453,6 +3453,7 @@ cont.$(OBJEXT): $(top_srcdir)/internal/sanitizers.h
34533453
cont.$(OBJEXT): $(top_srcdir)/internal/serial.h
34543454
cont.$(OBJEXT): $(top_srcdir)/internal/static_assert.h
34553455
cont.$(OBJEXT): $(top_srcdir)/internal/string.h
3456+
cont.$(OBJEXT): $(top_srcdir)/internal/thread.h
34563457
cont.$(OBJEXT): $(top_srcdir)/internal/variable.h
34573458
cont.$(OBJEXT): $(top_srcdir)/internal/vm.h
34583459
cont.$(OBJEXT): $(top_srcdir)/internal/warnings.h

cont.c

Lines changed: 79 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -28,6 +28,7 @@ extern int madvise(caddr_t, size_t, int);
2828
#include "eval_intern.h"
2929
#include "internal.h"
3030
#include "internal/cont.h"
31+
#include "internal/thread.h"
3132
#include "internal/error.h"
3233
#include "internal/gc.h"
3334
#include "internal/proc.h"
@@ -229,18 +230,18 @@ typedef struct rb_context_struct {
229230
struct rb_jit_cont *jit_cont; // Continuation contexts for JITs
230231
} rb_context_t;
231232

232-
233233
/*
234234
* Fiber status:
235-
* [Fiber.new] ------> FIBER_CREATED
236-
* | [Fiber#resume]
237-
* v
238-
* +--> FIBER_RESUMED ----+
239-
* [Fiber#resume] | | [Fiber.yield] |
240-
* | v |
241-
* +-- FIBER_SUSPENDED | [Terminate]
242-
* |
243-
* FIBER_TERMINATED <-+
235+
* [Fiber.new] ------> FIBER_CREATED ----> [Fiber#kill] --> |
236+
* | [Fiber#resume] |
237+
* v |
238+
* +--> FIBER_RESUMED ----> [return] ------> |
239+
* [Fiber#resume] | | [Fiber.yield/transfer] |
240+
* [Fiber#transfer] | v |
241+
* +--- FIBER_SUSPENDED --> [Fiber#kill] --> |
242+
* |
243+
* |
244+
* FIBER_TERMINATED <-------------------+
244245
*/
245246
enum fiber_status {
246247
FIBER_CREATED,
@@ -266,6 +267,8 @@ struct rb_fiber_struct {
266267
unsigned int yielding : 1;
267268
unsigned int blocking : 1;
268269

270+
unsigned int killed : 1;
271+
269272
struct coroutine_context context;
270273
struct fiber_pool_stack stack;
271274
};
@@ -1996,6 +1999,7 @@ fiber_t_alloc(VALUE fiber_value, unsigned int blocking)
19961999
fiber->cont.self = fiber_value;
19972000
fiber->cont.type = FIBER_CONTEXT;
19982001
fiber->blocking = blocking;
2002+
fiber->killed = 0;
19992003
cont_init(&fiber->cont, th);
20002004

20012005
fiber->cont.saved_ec.fiber_ptr = fiber;
@@ -2522,13 +2526,16 @@ rb_fiber_start(rb_fiber_t *fiber)
25222526
if (state == TAG_RAISE) {
25232527
// noop...
25242528
}
2529+
else if (state == TAG_FATAL && err == RUBY_FATAL_FIBER_KILLED) {
2530+
need_interrupt = FALSE;
2531+
err = Qfalse;
2532+
}
25252533
else if (state == TAG_FATAL) {
25262534
rb_threadptr_pending_interrupt_enque(th, err);
25272535
}
25282536
else {
25292537
err = rb_vm_make_jump_tag_but_local_jump(state, err);
25302538
}
2531-
need_interrupt = TRUE;
25322539
}
25332540

25342541
rb_fiber_terminate(fiber, need_interrupt, err);
@@ -2547,6 +2554,7 @@ rb_threadptr_root_fiber_setup(rb_thread_t *th)
25472554
fiber->cont.saved_ec.fiber_ptr = fiber;
25482555
fiber->cont.saved_ec.thread_ptr = th;
25492556
fiber->blocking = 1;
2557+
fiber->killed = 0;
25502558
fiber_status_set(fiber, FIBER_RESUMED); /* skip CREATED */
25512559
th->ec = &fiber->cont.saved_ec;
25522560
// When rb_threadptr_root_fiber_setup is called for the first time, rb_rjit_enabled and
@@ -2649,6 +2657,19 @@ fiber_store(rb_fiber_t *next_fiber, rb_thread_t *th)
26492657
fiber_setcontext(next_fiber, fiber);
26502658
}
26512659

2660+
static void
2661+
fiber_check_killed(rb_fiber_t *fiber)
2662+
{
2663+
VM_ASSERT(fiber == fiber_current());
2664+
2665+
if (fiber->killed) {
2666+
rb_thread_t *thread = fiber->cont.saved_ec.thread_ptr;
2667+
2668+
thread->ec->errinfo = RUBY_FATAL_FIBER_KILLED;
2669+
EC_JUMP_TAG(thread->ec, RUBY_TAG_FATAL);
2670+
}
2671+
}
2672+
26522673
static inline VALUE
26532674
fiber_switch(rb_fiber_t *fiber, int argc, const VALUE *argv, int kw_splat, rb_fiber_t *resuming_fiber, bool yielding)
26542675
{
@@ -2737,7 +2758,14 @@ fiber_switch(rb_fiber_t *fiber, int argc, const VALUE *argv, int kw_splat, rb_fi
27372758

27382759
current_fiber = th->ec->fiber_ptr;
27392760
value = current_fiber->cont.value;
2740-
if (current_fiber->cont.argc == -1) rb_exc_raise(value);
2761+
2762+
fiber_check_killed(current_fiber);
2763+
2764+
if (current_fiber->cont.argc == -1) {
2765+
// Fiber#raise will trigger this path.
2766+
rb_exc_raise(value);
2767+
}
2768+
27412769
return value;
27422770
}
27432771

@@ -3175,14 +3203,9 @@ rb_fiber_s_yield(int argc, VALUE *argv, VALUE klass)
31753203
}
31763204

31773205
static VALUE
3178-
fiber_raise(rb_fiber_t *fiber, int argc, const VALUE *argv)
3206+
fiber_raise(rb_fiber_t *fiber, VALUE exception)
31793207
{
3180-
VALUE exception = rb_make_exception(argc, argv);
3181-
3182-
if (fiber->resuming_fiber) {
3183-
rb_raise(rb_eFiberError, "attempt to raise a resuming fiber");
3184-
}
3185-
else if (FIBER_SUSPENDED_P(fiber) && !fiber->yielding) {
3208+
if (FIBER_SUSPENDED_P(fiber) && !fiber->yielding) {
31863209
return fiber_transfer_kw(fiber, -1, &exception, RB_NO_KEYWORDS);
31873210
}
31883211
else {
@@ -3193,7 +3216,9 @@ fiber_raise(rb_fiber_t *fiber, int argc, const VALUE *argv)
31933216
VALUE
31943217
rb_fiber_raise(VALUE fiber, int argc, const VALUE *argv)
31953218
{
3196-
return fiber_raise(fiber_ptr(fiber), argc, argv);
3219+
VALUE exception = rb_make_exception(argc, argv);
3220+
3221+
return fiber_raise(fiber_ptr(fiber), exception);
31973222
}
31983223

31993224
/*
@@ -3223,6 +3248,39 @@ rb_fiber_m_raise(int argc, VALUE *argv, VALUE self)
32233248
return rb_fiber_raise(self, argc, argv);
32243249
}
32253250

3251+
/*
3252+
* call-seq:
3253+
* fiber.kill -> nil
3254+
*
3255+
* Terminates +fiber+ by raising an uncatchable exception, returning
3256+
* the terminated Fiber.
3257+
*
3258+
* If the fiber has not been started, transition directly to the terminated state.
3259+
*
3260+
* If the fiber is already terminated, does nothing.
3261+
*/
3262+
static VALUE
3263+
rb_fiber_m_kill(VALUE self)
3264+
{
3265+
rb_fiber_t *fiber = fiber_ptr(self);
3266+
3267+
if (fiber->killed) return Qfalse;
3268+
fiber->killed = 1;
3269+
3270+
if (fiber->status == FIBER_CREATED) {
3271+
fiber->status = FIBER_TERMINATED;
3272+
}
3273+
else if (fiber->status != FIBER_TERMINATED) {
3274+
if (fiber_current() == fiber) {
3275+
fiber_check_killed(fiber);
3276+
} else {
3277+
fiber_raise(fiber_ptr(self), Qnil);
3278+
}
3279+
}
3280+
3281+
return self;
3282+
}
3283+
32263284
/*
32273285
* call-seq:
32283286
* Fiber.current -> fiber
@@ -3398,6 +3456,7 @@ Init_Cont(void)
33983456
rb_define_method(rb_cFiber, "storage=", rb_fiber_storage_set, 1);
33993457
rb_define_method(rb_cFiber, "resume", rb_fiber_m_resume, -1);
34003458
rb_define_method(rb_cFiber, "raise", rb_fiber_m_raise, -1);
3459+
rb_define_method(rb_cFiber, "kill", rb_fiber_m_kill, 0);
34013460
rb_define_method(rb_cFiber, "backtrace", rb_fiber_backtrace, -1);
34023461
rb_define_method(rb_cFiber, "backtrace_locations", rb_fiber_backtrace_locations, -1);
34033462
rb_define_method(rb_cFiber, "to_s", fiber_to_s, 0);

internal/thread.h

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,10 @@ struct rb_thread_struct; /* in vm_core.h */
2929
#define COVERAGE_TARGET_ONESHOT_LINES 8
3030
#define COVERAGE_TARGET_EVAL 16
3131

32+
#define RUBY_FATAL_THREAD_KILLED INT2FIX(0)
33+
#define RUBY_FATAL_THREAD_TERMINATED INT2FIX(1)
34+
#define RUBY_FATAL_FIBER_KILLED RB_INT2FIX(2)
35+
3236
VALUE rb_obj_is_mutex(VALUE obj);
3337
VALUE rb_suppress_tracing(VALUE (*func)(VALUE), VALUE arg);
3438
void rb_thread_execute_interrupts(VALUE th);

spec/ruby/core/fiber/kill_spec.rb

Lines changed: 90 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,90 @@
1+
require_relative '../../spec_helper'
2+
require_relative 'fixtures/classes'
3+
require_relative '../../shared/kernel/raise'
4+
5+
ruby_version_is "3.3" do
6+
describe "Fiber#kill" do
7+
it "kills a non-resumed fiber" do
8+
fiber = Fiber.new{}
9+
10+
fiber.alive?.should == true
11+
12+
fiber.kill
13+
fiber.alive?.should == false
14+
end
15+
16+
it "kills a resumed fiber" do
17+
fiber = Fiber.new{while true; Fiber.yield; end}
18+
fiber.resume
19+
20+
fiber.alive?.should == true
21+
22+
fiber.kill
23+
fiber.alive?.should == false
24+
end
25+
26+
it "can kill itself" do
27+
fiber = Fiber.new do
28+
Fiber.current.kill
29+
end
30+
31+
fiber.alive?.should == true
32+
33+
fiber.resume
34+
fiber.alive?.should == false
35+
end
36+
37+
it "kills a resumed fiber from a child" do
38+
parent = Fiber.new do
39+
child = Fiber.new do
40+
parent.kill
41+
parent.alive?.should == true
42+
end
43+
44+
child.resume
45+
end
46+
47+
parent.resume
48+
parent.alive?.should == false
49+
end
50+
51+
it "executes the ensure block" do
52+
ensure_executed = false
53+
54+
fiber = Fiber.new do
55+
while true; Fiber.yield; end
56+
ensure
57+
ensure_executed = true
58+
end
59+
60+
fiber.resume
61+
fiber.kill
62+
ensure_executed.should == true
63+
end
64+
65+
it "does not execute rescue block" do
66+
rescue_executed = false
67+
68+
fiber = Fiber.new do
69+
while true; Fiber.yield; end
70+
rescue Exception
71+
rescue_executed = true
72+
end
73+
74+
fiber.resume
75+
fiber.kill
76+
rescue_executed.should == false
77+
end
78+
79+
it "repeatedly kills a fiber" do
80+
fiber = Fiber.new do
81+
while true; Fiber.yield; end
82+
ensure
83+
while true; Fiber.yield; end
84+
end
85+
86+
fiber.kill
87+
fiber.alive?.should == false
88+
end
89+
end
90+
end

thread.c

Lines changed: 4 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -149,8 +149,6 @@ NORETURN(static void async_bug_fd(const char *mesg, int errno_arg, int fd));
149149
static int consume_communication_pipe(int fd);
150150
static int check_signals_nogvl(rb_thread_t *, int sigwait_fd);
151151

152-
#define eKillSignal INT2FIX(0)
153-
#define eTerminateSignal INT2FIX(1)
154152
static volatile int system_working = 1;
155153

156154
struct waiting_fd {
@@ -388,7 +386,7 @@ terminate_all(rb_ractor_t *r, const rb_thread_t *main_thread)
388386
if (th != main_thread) {
389387
RUBY_DEBUG_LOG("terminate start th:%u status:%s", rb_th_serial(th), thread_status_name(th, TRUE));
390388

391-
rb_threadptr_pending_interrupt_enque(th, eTerminateSignal);
389+
rb_threadptr_pending_interrupt_enque(th, RUBY_FATAL_THREAD_TERMINATED);
392390
rb_threadptr_interrupt(th);
393391

394392
RUBY_DEBUG_LOG("terminate done th:%u status:%s", rb_th_serial(th), thread_status_name(th, TRUE));
@@ -2337,8 +2335,8 @@ rb_threadptr_execute_interrupts(rb_thread_t *th, int blocking_timing)
23372335
if (UNDEF_P(err)) {
23382336
/* no error */
23392337
}
2340-
else if (err == eKillSignal /* Thread#kill received */ ||
2341-
err == eTerminateSignal /* Terminate thread */ ||
2338+
else if (err == RUBY_FATAL_THREAD_KILLED /* Thread#kill received */ ||
2339+
err == RUBY_FATAL_THREAD_TERMINATED /* Terminate thread */ ||
23422340
err == INT2FIX(TAG_FATAL) /* Thread.exit etc. */ ) {
23432341
terminate_interrupt = 1;
23442342
}
@@ -2569,7 +2567,7 @@ rb_thread_kill(VALUE thread)
25692567
}
25702568
else {
25712569
threadptr_check_pending_interrupt_queue(target_th);
2572-
rb_threadptr_pending_interrupt_enque(target_th, eKillSignal);
2570+
rb_threadptr_pending_interrupt_enque(target_th, RUBY_FATAL_THREAD_KILLED);
25732571
rb_threadptr_interrupt(target_th);
25742572
}
25752573

vm_insnhelper.c

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1616,7 +1616,7 @@ vm_throw_continue(const rb_execution_context_t *ec, VALUE err)
16161616
/* continue throw */
16171617

16181618
if (FIXNUM_P(err)) {
1619-
ec->tag->state = FIX2INT(err);
1619+
ec->tag->state = RUBY_TAG_FATAL;
16201620
}
16211621
else if (SYMBOL_P(err)) {
16221622
ec->tag->state = TAG_THROW;

0 commit comments

Comments
 (0)