Ruby 4.1.0dev (2026-10-02 revision bac5fe37358bdf2c0dc72c676eec8ca34392b970)
ractor_sync.c (bac5fe37358bdf2c0dc72c676eec8ca34392b970)
1// this file is included by ractor.c
2
3struct ractor_port {
5 st_data_t id_;
6};
7
8static st_data_t
9ractor_port_id(const struct ractor_port *rp)
10{
11 return rp->id_;
12}
13
14static VALUE rb_cRactorPort;
15
16static VALUE ractor_receive(rb_execution_context_t *ec, const struct ractor_port *rp, const rb_hrtime_t *end);
17static VALUE ractor_send(rb_execution_context_t *ec, const struct ractor_port *rp, VALUE obj, VALUE move);
18static struct ractor_basket *ractor_basket_new_exit(VALUE sender, VALUE token);
19static void ractor_send_basket(rb_execution_context_t *ec, const struct ractor_port *rp, struct ractor_basket *b, bool raise_on_error);
20static void ractor_add_port(rb_ractor_t *r, st_data_t id);
21
22// The off-heap courier a copy or a move payload travels in. Defined in ractor.c.
23struct rb_ractor_courier *rb_ractor_courier_build_move(VALUE obj, struct rb_ractor_courier **slot);
24VALUE rb_ractor_courier_materialize(struct rb_ractor_courier *c);
25void rb_ractor_courier_free(struct rb_ractor_courier *c);
26static void ractor_off_queue_add(rb_ractor_t *cr, struct ractor_basket *b);
27static void ractor_off_queue_remove(struct ractor_basket *b);
28void rb_ractor_courier_mark(struct rb_ractor_courier *c);
29struct rb_ractor_courier *rb_ractor_courier_build_copy(VALUE obj, struct rb_ractor_courier **slot);
30
31static void ractor_port_note_alive(const struct ractor_port *rp);
32
33static void
34ractor_port_mark(void *ptr)
35{
36 const struct ractor_port *rp = (struct ractor_port *)ptr;
37
38 if (rp->r) {
39 rb_gc_mark(rp->r->pub.self);
40
41 /* Only a mark that covers every objspace can call a port dead. Ask the single
42 * objspace first: it answers without the VM, which a GC worker thread cannot
43 * reach (mmtk marks from several of them). */
44 if (rb_gc_single_objspace_p() || rb_gc_during_global_gc_p()) {
45 ractor_port_note_alive(rp);
46 }
47 }
48}
49
50static const rb_data_type_t ractor_port_data_type = {
51 "ractor/port",
52 {
53 ractor_port_mark,
55 NULL, // memsize
56 NULL, // update
57 },
58 0, 0, RUBY_TYPED_THREAD_SAFE_FREE | RUBY_TYPED_WB_PROTECTED | RUBY_TYPED_FROZEN_SHAREABLE | RUBY_TYPED_EMBEDDABLE,
59};
60
61static st_data_t
62ractor_genid_for_port(rb_ractor_t *cr)
63{
64 // TODO: enough?
65 return cr->sync.next_port_id++;
66}
67
68static struct ractor_port *
69RACTOR_PORT_PTR(VALUE self)
70{
71 VM_ASSERT(rb_typeddata_is_kind_of(self, &ractor_port_data_type));
72 return RTYPEDDATA_GET_DATA(self);
73}
74
75// r is NULL between Ractor::Port.allocate and ractor_port_init()
76static struct ractor_port *
77ractor_port_ptr_check(VALUE self)
78{
79 struct ractor_port *rp = RACTOR_PORT_PTR(self);
80
81 if (UNLIKELY(rp->r == NULL)) {
82 rb_raise(rb_eTypeError, "uninitialized %"PRIsVALUE, rb_obj_class(self));
83 }
84
85 return rp;
86}
87
88static VALUE
89ractor_port_alloc(VALUE klass)
90{
91 struct ractor_port *rp;
92 VALUE rpv = TypedData_Make_Struct(klass, struct ractor_port, &ractor_port_data_type, rp);
93 rb_obj_freeze(rpv);
94 return rpv;
95}
96
97static VALUE
98ractor_port_init(VALUE rpv, rb_ractor_t *r)
99{
100 // Child threads can still run ensure blocks and report exceptions after
101 // ractor_notify_exit has freed the ports.
102 if (!r->sync.ports) {
103 rb_raise(rb_eRactorClosedError, "The ractor has terminated");
104 }
105
106 struct ractor_port *rp = RACTOR_PORT_PTR(rpv);
107
108 rp->r = r;
109 RB_OBJ_WRITTEN(rpv, Qundef, r->pub.self);
110 rp->id_ = ractor_genid_for_port(r);
111
112 ractor_add_port(r, ractor_port_id(rp));
113
114 rb_obj_freeze(rpv);
115
116 return rpv;
117}
118
119/*
120 * call-seq:
121 * Ractor::Port.new -> new_port
122 *
123 * Returns a new Ractor::Port object.
124 */
125static VALUE
126ractor_port_initialize(VALUE self)
127{
128 return ractor_port_init(self, GET_RACTOR());
129}
130
131/* :nodoc: */
132static VALUE
133ractor_port_initialize_copy(VALUE self, VALUE orig)
134{
135 struct ractor_port *dst = RACTOR_PORT_PTR(self); // uninitialized by definition
136 struct ractor_port *src = ractor_port_ptr_check(orig);
137 dst->r = src->r;
138 RB_OBJ_WRITTEN(self, Qundef, dst->r->pub.self);
139 dst->id_ = ractor_port_id(src);
140
141 return self;
142}
143
144static VALUE
145ractor_port_new(rb_ractor_t *r)
146{
147 VALUE rpv = ractor_port_alloc(rb_cRactorPort);
148 ractor_port_init(rpv, r);
149 return rpv;
150}
151
152static bool
153ractor_port_p(VALUE self)
154{
155 return rb_typeddata_is_kind_of(self, &ractor_port_data_type);
156}
157
158static const rb_hrtime_t *ractor_timeout_deadline(VALUE timeout, rb_hrtime_t *storage);
159
160static VALUE
161ractor_port_receive(rb_execution_context_t *ec, VALUE self, VALUE timeout)
162{
163 const struct ractor_port *rp = ractor_port_ptr_check(self);
164
165 if (rp->r != rb_ec_ractor_ptr(ec)) {
166 rb_raise(rb_eRactorError, "only allowed from the creator Ractor of this port");
167 }
168
169 rb_hrtime_t deadline;
170 const rb_hrtime_t *end = ractor_timeout_deadline(timeout, &deadline);
171
172 VALUE v = ractor_receive(ec, rp, end);
173 RB_GC_GUARD(self);
174
175 // no message before the timeout
176 return UNDEF_P(v) ? Qnil : v;
177}
178
179static VALUE
180ractor_port_send(rb_execution_context_t *ec, VALUE self, VALUE obj, VALUE move)
181{
182 const struct ractor_port *rp = ractor_port_ptr_check(self);
183 ractor_send(ec, rp, obj, RTEST(move));
184 RB_GC_GUARD(self);
185 return self;
186}
187
188static bool ractor_closed_port_p(rb_execution_context_t *ec, rb_ractor_t *r, const struct ractor_port *rp);
189static bool ractor_close_port(rb_execution_context_t *ec, rb_ractor_t *r, const struct ractor_port *rp);
190
191static VALUE
192ractor_port_closed_p(rb_execution_context_t *ec, VALUE self)
193{
194 const struct ractor_port *rp = ractor_port_ptr_check(self);
195 rb_ractor_t *r = rp->r;
196 bool closed;
197
198 if (rb_ec_ractor_ptr(ec) == r) {
199 /* The owner's threads are serialized by the ractor GVL, so the ports
200 * table can't change under this lookup. */
201 closed = ractor_closed_port_p(ec, r, rp);
202 }
203 else {
204 /* A foreign Ractor races the owner's st_insert/st_delete on the ports
205 * table; take the lock like every other foreign reader. ractor_closed_port_p
206 * asserts the lock is held for foreign access, and Port#closed? was the
207 * only path reaching it without the lock. */
208 RACTOR_LOCK(r);
209 {
210 closed = ractor_closed_port_p(ec, r, rp);
211 }
212 RACTOR_UNLOCK(r);
213 }
214
215 return closed ? Qtrue : Qfalse;
216}
217
218static VALUE
219ractor_port_close(rb_execution_context_t *ec, VALUE self)
220{
221 const struct ractor_port *rp = ractor_port_ptr_check(self);
222 rb_ractor_t *cr = rb_ec_ractor_ptr(ec);
223
224 if (cr != rp->r) {
225 rb_raise(rb_eRactorError, "closing port by other ractors is not allowed");
226 }
227
228 ractor_close_port(ec, cr, rp);
229 return self;
230}
231
232// ractor-internal
233
234// ractor-internal - ractor_basket
235
236enum ractor_basket_type {
237 // basket is empty
238 basket_type_none,
239
240 // value is available
241 basket_type_ref,
242 basket_type_copy,
243 basket_type_move,
244 /* A ractor's exit token. The pair it becomes is built on the receiving side,
245 * so a terminating ractor adds no shareable object to the vm, and building it
246 * needs neither a tag nor a courier -- which the sender has no stack for. */
247 basket_type_exit,
248};
249
251 enum ractor_basket_type type;
252 VALUE sender;
253 st_data_t port_id;
254
255 struct {
256 VALUE v;
257 bool exception;
258 /* The off-heap (xmalloc) courier the payload graph was serialized into.
259 * Copy and move both use it; when set, v is unused. */
260 struct rb_ractor_courier *courier;
261 } p; // payload
262
263 struct ccan_list_node node; /* the port queue it waits on */
264 struct ccan_list_node off_queue_node; /* or sync.off_queue_baskets, when on none */
265};
266
267#if 0
268static inline bool
269ractor_basket_type_p(const struct ractor_basket *b, enum ractor_basket_type type)
270{
271 return b->type == type;
272}
273
274static inline bool
275ractor_basket_none_p(const struct ractor_basket *b)
276{
277 return ractor_basket_type_p(b, basket_type_none);
278}
279#endif
280
281static void
282ractor_basket_mark(const struct ractor_basket *b)
283{
284 if (b->p.courier != NULL) {
285 /* The payload became this Ractor's to root the moment the message was enqueued
286 * here: the sender's own roots stop at the send. Before and after the queue the
287 * basket is on its holder's off_queue_baskets instead, so a courier is rooted
288 * from the moment it is allocated to the moment it is freed. */
289 rb_ractor_courier_mark(b->p.courier);
290 }
291 else {
292 rb_gc_mark(b->p.v);
293 }
294
295 /* An exit token names a ractor that may be reachable from nothing else. */
296 rb_gc_mark(b->sender);
297}
298
299static void
300ractor_basket_free(struct ractor_basket *b)
301{
302 ractor_off_queue_remove(b);
303 if (b->p.courier) {
304 /* A courier that was never consumed (a queue being torn down, say). */
305 rb_ractor_courier_free(b->p.courier);
306 b->p.courier = NULL;
307 }
308 SIZED_FREE(b);
309}
310
311static struct ractor_basket *
312ractor_basket_alloc(void)
313{
314 struct ractor_basket *b = ALLOC(struct ractor_basket);
315
316 /* Empty and mark-safe from the start: a basket goes on its holder's in-flight list
317 * before it has a payload, so a GC can walk it while it is still being filled. */
318 b->type = basket_type_none;
319 b->sender = Qnil;
320 b->port_id = 0;
321 b->p.v = Qnil;
322 b->p.exception = false;
323 b->p.courier = NULL;
324 ccan_list_node_init(&b->off_queue_node);
325
326 return b;
327}
328
329/* A basket is rooted by whoever holds it: a port queue while it waits there, and its
330 * holder's off-queue list while it is being built or materialized. */
331static void
332ractor_off_queue_add(rb_ractor_t *cr, struct ractor_basket *b)
333{
334 VM_ASSERT(cr == rb_current_ractor_raw(false));
335 ccan_list_add_tail(&cr->sync.off_queue_baskets, &b->off_queue_node);
336}
337
338static void
339ractor_off_queue_remove(struct ractor_basket *b)
340{
341 ccan_list_del_init(&b->off_queue_node);
342}
343
344static void
345ractor_mark_off_queue_baskets(rb_ractor_t *r)
346{
347 struct ractor_basket *b;
348 ccan_list_for_each(&r->sync.off_queue_baskets, b, off_queue_node) {
349 ractor_basket_mark(b);
350 }
351}
352
353// ractor-internal - ractor_queue
354
356 struct ccan_list_head set;
357 bool closed;
358 bool alive; /* its Ractor::Port is still reachable; see the reap below */
359};
360
361static void
362ractor_queue_init(struct ractor_queue *rq)
363{
364 ccan_list_head_init(&rq->set);
365 rq->closed = false;
366 rq->alive = true;
367}
368
369static struct ractor_queue *
370ractor_queue_new(void)
371{
372 struct ractor_queue *rq = ALLOC(struct ractor_queue);
373 ractor_queue_init(rq);
374 return rq;
375}
376
377static void
378ractor_port_note_alive(const struct ractor_port *rp)
379{
380 struct ractor_queue *rq;
381
382 if (rp->r->sync.ports && st_lookup(rp->r->sync.ports, rp->id_, (st_data_t *)&rq)) {
383 /* Several markers can reach the same port at once (mmtk marks from its GC worker
384 * threads), but they all store the same value and the reap reads it once marking
385 * is over. */
386 rq->alive = true;
387 }
388}
389
390static void
391ractor_queue_mark(const struct ractor_queue *rq)
392{
393 const struct ractor_basket *b;
394
395 ccan_list_for_each(&rq->set, b, node) {
396 ractor_basket_mark(b);
397 }
398}
399
400static void
401ractor_queue_free(struct ractor_queue *rq)
402{
403 struct ractor_basket *b, *nxt;
404
405 ccan_list_for_each_safe(&rq->set, b, nxt, node) {
406 ccan_list_del_init(&b->node);
407 ractor_basket_free(b);
408 }
409
410 VM_ASSERT(ccan_list_empty(&rq->set));
411
412 SIZED_FREE(rq);
413}
414
416static size_t
417ractor_queue_size(const struct ractor_queue *rq)
418{
419 size_t size = 0;
420 const struct ractor_basket *b;
421
422 ccan_list_for_each(&rq->set, b, node) {
423 size++;
424 }
425 return size;
426}
427
428// Returns whether this call is the one that closed it.
429static bool
430ractor_queue_close(struct ractor_queue *rq)
431{
432 bool closed_now = !rq->closed;
433 rq->closed = true;
434 return closed_now;
435}
436
437static void
438ractor_queue_move(struct ractor_queue *dst_rq, struct ractor_queue *src_rq)
439{
440 struct ccan_list_head *src = &src_rq->set;
441 struct ccan_list_head *dst = &dst_rq->set;
442
443 dst->n.next = src->n.next;
444 dst->n.prev = src->n.prev;
445 dst->n.next->prev = &dst->n;
446 dst->n.prev->next = &dst->n;
447 ccan_list_head_init(src);
448}
449
450#if 0
451static struct ractor_basket *
452ractor_queue_head(rb_ractor_t *r, struct ractor_queue *rq)
453{
454 return ccan_list_top(&rq->set, struct ractor_basket, node);
455}
456#endif
457
458static bool
459ractor_queue_empty_p(rb_ractor_t *r, const struct ractor_queue *rq)
460{
461 return ccan_list_empty(&rq->set);
462}
463
464static struct ractor_basket *
465ractor_queue_deq(rb_ractor_t *r, struct ractor_queue *rq)
466{
467 VM_ASSERT(GET_RACTOR() == r);
468
469 return ccan_list_pop(&rq->set, struct ractor_basket, node);
470}
471
472static void
473ractor_queue_enq(rb_ractor_t *r, struct ractor_queue *rq, struct ractor_basket *basket)
474{
475 ccan_list_add_tail(&rq->set, &basket->node);
476}
477
478#if 0
479static void
480rq_dump(const struct ractor_queue *rq)
481{
482 int i=0;
483 struct ractor_basket *b;
484 ccan_list_for_each(&rq->set, b, node) {
485 fprintf(stderr, "%d type:%s %p\n", i, basket_type_name(b->type), (void *)b);
486 i++;
487 }
488}
489#endif
490
491static void ractor_delete_port(rb_ractor_t *cr, st_data_t id, bool locked);
492
493static struct ractor_queue *
494ractor_get_queue(rb_ractor_t *cr, st_data_t id, bool locked)
495{
496 VM_ASSERT(cr == GET_RACTOR());
497
498 struct ractor_queue *rq;
499
500 if (cr->sync.ports && st_lookup(cr->sync.ports, id, (st_data_t *)&rq)) {
501 if (rq->closed && ractor_queue_empty_p(cr, rq)) {
502 ractor_delete_port(cr, id, locked);
503 return NULL;
504 }
505 else {
506 return rq;
507 }
508 }
509 else {
510 return NULL;
511 }
512}
513
514// ractor-internal - ports
515
516static void
517ractor_add_port(rb_ractor_t *r, st_data_t id)
518{
519 struct ractor_queue *rq = ractor_queue_new();
520 ASSERT_ractor_unlocking(r);
521
522 RUBY_DEBUG_LOG("id:%u", (unsigned int)id);
523
524 // Rebuilding the table on insertion can run GC by the allocation and the
525 // GC acquires the VM lock, which is prohibited under the ractor lock.
526 st_table *const old_tab = r->sync.ports;
527 bool inserted;
528
529 RACTOR_LOCK(r);
530 {
531 inserted = st_insert_no_rebuild(old_tab, id, (st_data_t)rq) >= 0;
532 }
533 RACTOR_UNLOCK(r);
534
535 if (!inserted) {
536 // The table is full. Rebuild it outside of the ractor lock (mutators
537 // are serialized by the per-ractor GVL) and swap it under the lock
538 // to exclude the readers (other ractors).
539 st_table *new_tab;
540
541 // Those allocations can run a global GC, whose reap frees dead ports and
542 // drops them from old_tab; a copy taken across one would republish the
543 // freed queues. Only the owner inserts into its own table, so a changed
544 // count means a reap ran. Check after each allocation: st_copy fills the
545 // header before it allocates the entries, so a reap in between leaves the
546 // copy counting rows it does not have, which st_insert must not be given.
547 while (1) {
548 const st_index_t entries = st_table_size(old_tab);
549
550 new_tab = st_copy(old_tab);
551
552 if (st_table_size(old_tab) == entries) {
553 st_insert(new_tab, id, (st_data_t)rq);
554
555 if (st_table_size(old_tab) == entries) break;
556 }
557
558 st_free_table(new_tab);
559 }
560
561 RACTOR_LOCK(r);
562 {
563 VM_ASSERT(r->sync.ports == old_tab);
564 r->sync.ports = new_tab;
565 }
566 RACTOR_UNLOCK(r);
567
568 st_free_table(old_tab);
569 }
570}
571
572static void
573ractor_delete_port_locked(rb_ractor_t *cr, st_data_t id)
574{
575 ASSERT_ractor_locking(cr);
576
577 RUBY_DEBUG_LOG("id:%u", (unsigned int)id);
578
579 struct ractor_queue *rq;
580
581 if (st_delete(cr->sync.ports, &id, (st_data_t *)&rq)) {
582 ractor_queue_free(rq);
583 }
584 else {
585 VM_ASSERT(0);
586 }
587}
588
589static void
590ractor_delete_port(rb_ractor_t *cr, st_data_t id, bool locked)
591{
592 if (locked) {
593 ractor_delete_port_locked(cr, id);
594 }
595 else {
596 RACTOR_LOCK_SELF(cr);
597 {
598 ractor_delete_port_locked(cr, id);
599 }
600 RACTOR_UNLOCK_SELF(cr);
601 }
602}
603
604static const struct ractor_port *
605ractor_default_port(rb_ractor_t *r)
606{
607 return RACTOR_PORT_PTR(r->sync.default_port_value);
608}
609
610static VALUE
611ractor_default_port_value(rb_ractor_t *r)
612{
613 return r->sync.default_port_value;
614}
615
616static bool
617ractor_closed_port_p(rb_execution_context_t *ec, rb_ractor_t *r, const struct ractor_port *rp)
618{
619 VM_ASSERT(rb_ec_ractor_ptr(ec) == rp->r ? 1 : (ASSERT_ractor_locking(rp->r), 1));
620
621 const struct ractor_queue *rq;
622
623 if (rp->r->sync.ports && st_lookup(rp->r->sync.ports, ractor_port_id(rp), (st_data_t *)&rq)) {
624 return rq->closed;
625 }
626 else {
627 return true;
628 }
629}
630
631static void ractor_deliver_incoming_messages(rb_execution_context_t *ec, rb_ractor_t *cr);
632static bool ractor_queue_empty_p(rb_ractor_t *r, const struct ractor_queue *rq);
633
634static bool ractor_wakeup_all(rb_ractor_t *r, enum ractor_wakeup_status wakeup_status);
635
636static bool
637ractor_close_port(rb_execution_context_t *ec, rb_ractor_t *cr, const struct ractor_port *rp)
638{
639 VM_ASSERT(cr == rp->r);
640 struct ractor_queue *rq = NULL;
641 bool closed_now = false;
642
643 RACTOR_LOCK_SELF(cr);
644 {
645 ractor_deliver_incoming_messages(ec, cr); // check incoming messages
646
647 if (st_lookup(rp->r->sync.ports, ractor_port_id(rp), (st_data_t *)&rq)) {
648 closed_now = ractor_queue_close(rq);
649
650 if (ractor_queue_empty_p(cr, rq)) {
651 // delete from the table
652 ractor_delete_port(cr, ractor_port_id(rp), true);
653 }
654
655 // TODO: free rq
656 }
657 }
658 RACTOR_UNLOCK_SELF(cr);
659
660 if (closed_now) {
661 // Only when this call closed it: waking the Ractor is not free to the
662 // other waiters, and a re-close of a closed port is news to nobody.
663 ractor_wakeup_all(cr, wakeup_by_close);
664 }
665
666 return rq != NULL;
667}
668
669/* A port is the only way to receive from its queue, so a queue whose port is gone is
670 * unreachable -- but the table is keyed by id, so no sweep finds it. Mark and sweep the
671 * table itself: ractor_port_mark sets the flag, this clears it for the next cycle. */
672static int
673ractor_reap_dead_ports_i(st_data_t port_id, st_data_t val, st_data_t dat)
674{
675 struct ractor_queue *rq = (struct ractor_queue *)val;
676
677 if (rq->alive) {
678 rq->alive = false;
679 return ST_CONTINUE;
680 }
681 else {
682 ractor_queue_free(rq);
683 return ST_DELETE;
684 }
685}
686
687/* A message is only moved from recv_queue to its port queue when the owner receives or
688 * closes, so one addressed to a port that died first is left here, out of the sweep
689 * above. Its port is gone from the table by now: drop it. */
690static void
691ractor_reap_undeliverable_messages(rb_ractor_t *r)
692{
693 struct ractor_queue *recv_q = r->sync.recv_queue;
694 if (recv_q == NULL) return;
695
696 struct ractor_basket *b, *nxt;
697 ccan_list_for_each_safe(&recv_q->set, b, nxt, node) {
698 if (!st_lookup(r->sync.ports, b->port_id, NULL)) {
699 ccan_list_del_init(&b->node);
700 ractor_basket_free(b);
701 }
702 }
703}
704
705void
706rb_ractor_reap_dead_ports(rb_ractor_t *r)
707{
708 /* No sync lock here: the caller's gate is what keeps foreign senders out. */
709 VM_ASSERT(rb_gc_single_objspace_p() || rb_gc_during_global_gc_p());
710
711 if (r->sync.ports) {
712 st_foreach(r->sync.ports, ractor_reap_dead_ports_i, 0);
713 ractor_reap_undeliverable_messages(r);
714 }
715}
716
717static int
718ractor_free_all_ports_i(st_data_t port_id, st_data_t val, st_data_t dat)
719{
720 struct ractor_queue *rq = (struct ractor_queue *)val;
721 // rb_ractor_t *cr = (rb_ractor_t *)dat;
722
723 ractor_queue_free(rq);
724 return ST_CONTINUE;
725}
726
727static void
728ractor_free_all_ports(rb_ractor_t *cr)
729{
730 if (cr->sync.ports) {
731 st_foreach(cr->sync.ports, ractor_free_all_ports_i, (st_data_t)cr);
732 st_free_table(cr->sync.ports);
733 cr->sync.ports = NULL;
734 }
735
736 if (cr->sync.recv_queue) {
737 ractor_queue_free(cr->sync.recv_queue);
738 cr->sync.recv_queue = NULL;
739 }
740}
741
742#if defined(HAVE_WORKING_FORK)
743static void
744ractor_sync_terminate_atfork(rb_vm_t *vm, rb_ractor_t *r)
745{
746 ractor_free_all_ports(r);
747 r->sync.legacy = Qnil;
748}
749#endif
750
751// Ractor#monitor
752
754 struct ractor_port port;
755 struct ccan_list_node node;
756};
757
758/* Mark the Ractors monitoring r. ractor_notify_exit sends the exit token through each
759 * entry's port, so the monitoring Ractor's struct must outlive r, and its wrapper is
760 * what keeps it alive. */
761static void
762ractor_mark_monitors(rb_ractor_t *r)
763{
764 const struct ractor_monitor *rm;
765 ccan_list_for_each(&r->sync.monitors, rm, node) {
766 rb_gc_mark(rm->port.r->pub.self);
767 }
768}
769
770/* Paired with the ractor it is about when it is received, so that many ractors
771 * can report to one port. */
772static VALUE
773ractor_exit_token(const rb_ractor_t *r)
774{
775 if (r->sync.legacy_exc) {
776 RUBY_DEBUG_LOG("aborted");
777 return ID2SYM(idAborted);
778 }
779 else {
780 RUBY_DEBUG_LOG("exited");
781 return ID2SYM(idExited);
782 }
783}
784
785static VALUE
787{
788 rb_ractor_t *r = RACTOR_PTR(self);
789 bool terminated = false;
790 const struct ractor_port *rp = ractor_port_ptr_check(port);
791 struct ractor_monitor *rm = ALLOC(struct ractor_monitor);
792 rm->port = *rp; // copy port information
793
794 RACTOR_LOCK(r);
795 {
796 if (UNDEF_P(r->sync.legacy)) { // not terminated
797 RUBY_DEBUG_LOG("OK/r:%u -> port:%u@r%u", (unsigned int)rb_ractor_id(r), (unsigned int)ractor_port_id(&rm->port), (unsigned int)rb_ractor_id(rm->port.r));
798 ccan_list_add_tail(&r->sync.monitors, &rm->node);
799 }
800 else {
801 RUBY_DEBUG_LOG("NG/r:%u -> port:%u@r%u", (unsigned int)rb_ractor_id(r), (unsigned int)ractor_port_id(&rm->port), (unsigned int)rb_ractor_id(rm->port.r));
802 terminated = true;
803 }
804 }
805 RACTOR_UNLOCK(r);
806
807 if (terminated) {
808 SIZED_FREE(rm);
809 ractor_send_basket(ec, rp, ractor_basket_new_exit(self, ractor_exit_token(r)), false);
810 }
811 /* rp points into port, which is embedded and may be referenced only from here */
812 RB_GC_GUARD(port);
813
814 return terminated ? Qfalse : Qtrue;
815}
816
817static VALUE
818ractor_unmonitor(rb_execution_context_t *ec, VALUE self, VALUE port)
819{
820 rb_ractor_t *r = RACTOR_PTR(self);
821 const struct ractor_port *rp = ractor_port_ptr_check(port);
822
823 RACTOR_LOCK(r);
824 {
825 if (UNDEF_P(r->sync.legacy)) { // not terminated
826 struct ractor_monitor *rm, *nxt;
827
828 ccan_list_for_each_safe(&r->sync.monitors, rm, nxt, node) {
829 if (rm->port.r == rp->r && ractor_port_id(&rm->port) == ractor_port_id(rp)) {
830 RUBY_DEBUG_LOG("r:%u -> port:%u@r%u",
831 (unsigned int)rb_ractor_id(r),
832 (unsigned int)ractor_port_id(&rm->port),
833 (unsigned int)rb_ractor_id(rm->port.r));
834 ccan_list_del(&rm->node);
835 SIZED_FREE(rm);
836 }
837 }
838 }
839 }
840 RACTOR_UNLOCK(r);
841
842 return self;
843}
844
845static void
846ractor_notify_exit(rb_execution_context_t *ec, rb_ractor_t *cr, VALUE legacy, bool exc)
847{
848 RUBY_DEBUG_LOG("exc:%d", exc);
849 VM_ASSERT(!UNDEF_P(legacy));
850 VM_ASSERT(cr->sync.legacy == Qundef);
851
852 RACTOR_LOCK_SELF(cr);
853 {
854 ractor_free_all_ports(cr);
855
856 cr->sync.legacy = legacy;
857 cr->sync.legacy_exc = exc;
858 }
859 RACTOR_UNLOCK_SELF(cr);
860
861}
862
863/* Sent after the dying thread's post-mortem collection: waking a joiner any earlier makes
864 * ractor_value spin for the whole of that collection. */
865static void
866ractor_send_exit_tokens(rb_execution_context_t *ec, rb_ractor_t *cr)
867{
868 VALUE token = ractor_exit_token(cr);
869 struct ractor_monitor *rm, *nxt;
870
871 ccan_list_for_each_safe(&cr->sync.monitors, rm, nxt, node)
872 {
873 RUBY_DEBUG_LOG("port:%u@r%u", (unsigned int)ractor_port_id(&rm->port), (unsigned int)rb_ractor_id(rm->port.r));
874
875 ractor_send_basket(ec, &rm->port, ractor_basket_new_exit(cr->pub.self, token), false);
876
877 ccan_list_del(&rm->node);
878 SIZED_FREE(rm);
879 }
880
881 VM_ASSERT(ccan_list_empty(&cr->sync.monitors));
882}
883
884// ractor-internal - initialize, mark, free, memsize
885
886static int
887ractor_mark_ports_i(st_data_t key, st_data_t val, st_data_t data)
888{
889 // id -> ractor_queue
890 const struct ractor_queue *rq = (struct ractor_queue *)val;
891 ractor_queue_mark(rq);
892 return ST_CONTINUE;
893}
894
895static void
896ractor_sync_mark(rb_ractor_t *r)
897{
898 /* The owner rewrites the queues, the port table and the monitor list under its sync
899 * lock, so only the owner itself or the stopped world may walk them. */
900 const bool world_stopped = rb_gc_during_global_gc_p();
901 VM_ASSERT(world_stopped || r == rb_current_ractor_raw(false));
902
903 rb_gc_mark(r->sync.default_port_value);
904
905 /* Until the value is absorbed this is its only reliable root (Qundef while the
906 * Ractor still runs); after Ractor#value returns it, the Ruby side roots it. */
907 rb_gc_mark(r->sync.legacy);
908
909 /* ractor_sync_init builds the rest, and a root scan reaches the main Ractor before
910 * that: ports is what tells the two apart (the lock and the list heads are still
911 * zeroed, and walking those crashes). Lock out foreign senders while walking them
912 * (self-lock: not recursive, and a held Ractor lock disables malloc-GC, so no GC
913 * nests); a stopped world needs no lock. */
914 if (r->sync.ports) {
915 if (!world_stopped) RACTOR_LOCK_SELF(r);
916 {
917 ractor_queue_mark(r->sync.recv_queue);
918 st_foreach(r->sync.ports, ractor_mark_ports_i, 0);
919 ractor_mark_monitors(r);
920 }
921 if (!world_stopped) RACTOR_UNLOCK_SELF(r);
922
923 /* The baskets on no queue: one being built to send, one being materialized.
924 * Walked in every collection, like the queues. What they hold is shareable, but
925 * "only a global GC frees a shareable" does not hold: pinned_roots_mark, which
926 * roots a shareable from its page bit, is skipped once the process is back to a
927 * single Ractor (rb_gc_single_objspace_p), and then an ordinary local GC frees
928 * one that nothing else names. A payload in flight is named by its basket and
929 * nothing else, so this list has to be a root whenever the queues are. No sync
930 * lock, though: only the owner touches it (the lock above guards the queues,
931 * which a foreign sender writes). */
932 ractor_mark_off_queue_baskets(r);
933 }
934}
935
936static int
937ractor_sync_free_ports_i(st_data_t _key, st_data_t val, st_data_t _args)
938{
939 struct ractor_queue *queue = (struct ractor_queue *)val;
940
941 ractor_queue_free(queue);
942
943 return ST_CONTINUE;
944}
945
946static void
947ractor_sync_free(rb_ractor_t *r)
948{
949 if (r->sync.recv_queue) {
950 ractor_queue_free(r->sync.recv_queue);
951 }
952
953 // maybe NULL
954 if (r->sync.ports) {
955 st_foreach(r->sync.ports, ractor_sync_free_ports_i, 0);
956 st_free_table(r->sync.ports);
957 r->sync.ports = NULL;
958 }
959}
960
961static size_t
962ractor_sync_memsize(const rb_ractor_t *r)
963{
964 if (r->sync.ports) {
965 return st_memsize(r->sync.ports);
966 }
967 else {
968 return 0;
969 }
970}
971
972static void
973ractor_sync_init(rb_ractor_t *r)
974{
975 // lock
976 rb_native_mutex_initialize(&r->sync.lock);
977
978 // monitors
979 ccan_list_head_init(&r->sync.off_queue_baskets);
980 ccan_list_head_init(&r->sync.monitors);
981
982 // waiters
983 ccan_list_head_init(&r->sync.waiters);
984
985 // receiving queue
986 r->sync.recv_queue = ractor_queue_new();
987
988 // ports
989 r->sync.ports = st_init_numtable();
990 /* ractor_setup_default_port creates it only after the Ractor joins
991 * vm->ractor.set, so a global GC cannot free the rootless port in between. */
992 r->sync.default_port_value = Qfalse;
993
994 // legacy
995 r->sync.legacy = Qundef;
996
997 // no receive is rebuilding a payload yet
998
999}
1000
1001/* Create the default port. Call only after the Ractor joined vm->ractor.set, so the
1002 * root scan can mark the shareable port from creation onwards. */
1003void
1004rb_ractor_setup_default_port(rb_ractor_t *r)
1005{
1006 VM_ASSERT(r->sync.default_port_value == Qfalse);
1007 r->sync.default_port_value = ractor_port_new(r);
1008 RB_OBJ_SET_SHAREABLE(r->sync.default_port_value);
1009}
1010
1011// Ractor#value
1012
1013static rb_ractor_t *
1014ractor_set_successor_once(rb_ractor_t *r, rb_ractor_t *cr)
1015{
1016 if (r->sync.successor == NULL) {
1017 rb_ractor_t *successor = ATOMIC_PTR_CAS(r->sync.successor, NULL, cr);
1018 return successor == NULL ? cr : successor;
1019 }
1020
1021 return r->sync.successor;
1022}
1023
1024static VALUE
1025ractor_make_remote_exception(VALUE cause, VALUE sender)
1026{
1027 VALUE err = rb_exc_new_cstr(rb_eRactorRemoteError, "thrown by remote Ractor.");
1028 rb_ivar_set(err, rb_intern("@ractor"), sender);
1029 rb_ec_setup_exception(NULL, err, cause);
1030 return err;
1031}
1032
1033static VALUE
1034ractor_value(rb_execution_context_t *ec, VALUE self)
1035{
1036 rb_ractor_t *cr = rb_ec_ractor_ptr(ec);
1037 rb_ractor_t *r = RACTOR_PTR(self);
1038 rb_ractor_t *sr = ractor_set_successor_once(r, cr);
1039
1040 if (sr == cr) {
1041 if (r->sync.legacy_taken) {
1042 rb_raise(rb_eRactorError, "The value was already taken");
1043 }
1044
1045 /* The value is returned by reference: inherit the dead Ractor's objspace first,
1046 * making it our own object (containment without a copy). Wait for
1047 * ractor_terminated: a monitor-port wakeup arrives before the dying thread
1048 * finishes teardown (vm_remove_ractor still touches the objspace). */
1049 while (!rb_ractor_status_p(r, ractor_terminated)) {
1051 }
1052
1053 /* The wait above yields the GVL, so another thread of this Ractor can take the
1054 * value first: re-check. */
1055 if (r->sync.legacy_taken) {
1056 rb_raise(rb_eRactorError, "The value was already taken");
1057 }
1058
1059 /* Move r's rb_gc_register_mark_object pins to the joiner before the merge
1060 * below sweeps r's objspace, or the objects pinned there lose their root. */
1061 rb_ractor_absorb_registered_marks(GET_RACTOR(), r);
1062 rb_ractor_absorb_registered_addrs_without_gc(GET_RACTOR(), r);
1063
1064 rb_gc_objspace_absorb_into_current(&r->objspace);
1065
1066 /* Keep legacy alive in a C local until it is returned: after the absorb only
1067 * the C struct reaches it, so let the conservative machine-stack mark find it. */
1068 volatile VALUE legacy_keep = r->sync.legacy;
1069
1070 /* A dead Ractor's local storage is unreachable from Ruby (Ractor#[] only works
1071 * from inside), so let the values die and keep ractor_mark and ractor_free from
1072 * walking a stale table later. */
1073 ractor_local_storage_free(r);
1074 r->local_storage = NULL;
1075 r->idkey_local_storage = NULL;
1076
1077 /* The value is returned to the caller and rooted from Ruby afterwards. Drop it
1078 * from the C struct: keeping it would leave a C-only reference into the
1079 * successor's objspace, needing marking and a pin against compaction. */
1080 VALUE legacy = r->sync.legacy;
1081 r->sync.legacy = Qnil;
1082 r->sync.legacy_taken = true;
1083 RB_GC_GUARD(legacy_keep);
1084
1085 if (r->sync.legacy_exc) {
1086 rb_exc_raise(ractor_make_remote_exception(legacy, self));
1087 }
1088 return legacy;
1089 }
1090 else {
1091 rb_raise(rb_eRactorError, "Only the successor ractor can take a value");
1092 }
1093}
1094
1095static VALUE
1096ractor_marshal_dump_body(VALUE obj)
1097{
1098 return rb_marshal_dump(obj, Qnil);
1099}
1100
1101static VALUE
1102ractor_marshal_dump_rescue(VALUE obj, VALUE errinfo)
1103{
1104 rb_raise(rb_eRactorError, "can not copy %"PRIsVALUE" object.", rb_class_of(obj));
1106}
1107
1108static VALUE
1109ractor_prepare_payload(rb_execution_context_t *ec, VALUE obj, enum ractor_basket_type *ptype,
1110 struct rb_ractor_courier **pcourier)
1111{
1112 switch (*ptype) {
1113 case basket_type_ref:
1114 return obj;
1115 default:
1116 if (rb_ractor_shareable_p(obj)) {
1117 *ptype = basket_type_ref;
1118 return obj;
1119 }
1120 else {
1121 /* Snapshot the object on the sender side without calling the user-visible
1122 * #clone. The courier is off-heap, so an in-flight payload is never a GC
1123 * object and needs no pin: nothing of the sender's heap stays alive while
1124 * the message waits. */
1125 *ptype = basket_type_copy;
1126 if (rb_ractor_courier_build_copy(obj, pcourier) == NULL) {
1127 rb_raise(rb_eRactorError, "can not copy %"PRIsVALUE" object.", rb_class_of(obj));
1128 }
1129 return Qundef;
1130 }
1131 }
1132}
1133
1134#if RBIMPL_COMPILER_IS(GCC) && defined(__OPTIMIZE__)
1135/* GCC produces false-positive -Wclobbered warnings after inlining
1136 * this function into ractor_basket_new(). */
1137NOINLINE(static void ractor_basket_build_payload(rb_execution_context_t *ec, struct ractor_basket *b, VALUE obj, enum ractor_basket_type type, bool exc));
1138#endif
1139static void
1140ractor_basket_build_payload(rb_execution_context_t *ec, struct ractor_basket *b, VALUE obj,
1141 enum ractor_basket_type type, bool exc)
1142{
1143 b->p.exception = exc;
1144 if (type == basket_type_move) {
1145 /* Serialize the graph into an off-heap courier; the sources become
1146 * RactorMovedObject. While in flight there is no GC object left for the
1147 * sender's GC to mark, sweep or move. The build publishes the courier into
1148 * the basket as soon as it exists. */
1149 rb_ractor_courier_build_move(obj, &b->p.courier);
1150 b->type = type;
1151 b->p.v = Qfalse;
1152 }
1153 else {
1154 VALUE v = ractor_prepare_payload(ec, obj, &type, &b->p.courier);
1155 b->type = type;
1156 b->p.v = v;
1157 }
1158}
1159
1160static struct ractor_basket *
1161ractor_basket_new(rb_execution_context_t *ec, VALUE obj, enum ractor_basket_type type, bool exc)
1162{
1163 rb_ractor_t *cr = rb_ec_ractor_ptr(ec);
1164 /* Allocate and list the basket before anything is built into it: from here the
1165 * courier it is about to hold is rooted by this Ractor's in-flight list, even
1166 * half-built, and every raise below frees it through one path. */
1167 struct ractor_basket *b = ractor_basket_alloc();
1168 ractor_off_queue_add(cr, b);
1169
1170 enum ruby_tag_type state;
1171 EC_PUSH_TAG(ec);
1172 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
1173 ractor_basket_build_payload(ec, b, obj, type, exc);
1174 }
1175 EC_POP_TAG();
1176 if (state != TAG_NONE) {
1177 ractor_basket_free(b); /* leaves the list and frees a courier already built */
1178 EC_JUMP_TAG(ec, state);
1179 }
1180
1181 return b;
1182}
1183
1184static VALUE
1185ractor_basket_value(struct ractor_basket *b)
1186{
1187 switch (b->type) {
1188 case basket_type_ref:
1189 break;
1190 case basket_type_exit:
1191 /* Allocated here, in the receiving ractor: a copy, not a shared object. */
1192 return rb_ary_new_from_args(2, b->sender, b->p.v);
1193 case basket_type_copy:
1194 /* An off-heap copy courier rebuilds exactly like a move one; only the sources
1195 * differ (still alive here, already shells there). */
1196 case basket_type_move: {
1197 /* Rebuild the moved graph from the off-heap courier into this Ractor's
1198 * objspace. The sources are already RactorMovedObject (set when the courier
1199 * was built), so move's snapshot semantics hold. The courier is xmalloc'd
1200 * rather than a GC object, so the sender's concurrent local GC never touches
1201 * it; the shareable VALUEs it carries are marked through this basket, which is
1202 * on this Ractor's off-queue list until it is freed.
1203 *
1204 * Rebuilding can raise here too (rb_hash_aset on a moved key with a custom
1205 * #hash runs user code, and an async interrupt can arrive). On a raise the
1206 * courier is still owned by the basket, whose teardown frees it. */
1207 rb_execution_context_t *ec = rb_current_ec_noinline();
1208 struct rb_ractor_courier *courier = b->p.courier;
1209 /* Keep the materialized graph on the machine stack (result): it is the only
1210 * root until it reaches the caller. courier_free below runs a long loop, and
1211 * only the malloc'd basket's p.v holding it would give a concurrent global GC a
1212 * wide window. */
1213 VALUE result = Qundef;
1214 enum ruby_tag_type state;
1215 EC_PUSH_TAG(ec);
1216 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
1217 result = rb_ractor_courier_materialize(courier);
1218 }
1219 EC_POP_TAG();
1220 if (state != TAG_NONE) {
1221 /* An unconsumed courier stays in b->p.courier; basket_free frees it. */
1222 ractor_basket_free(b);
1223 EC_JUMP_TAG(ec, state);
1224 }
1225 rb_ractor_courier_free(courier);
1226 b->p.courier = NULL;
1227 b->p.v = result;
1228 RB_GC_GUARD(result);
1229 break;
1230 }
1231 default:
1232 VM_ASSERT(0); // unreachable
1233 }
1234
1235 VM_ASSERT(!RB_TYPE_P(b->p.v, T_NONE));
1236 return b->p.v;
1237}
1238
1239static VALUE
1240ractor_basket_accept(struct ractor_basket *b)
1241{
1242 VALUE v = ractor_basket_value(b);
1243
1244 if (b->p.exception) {
1245 VALUE err = ractor_make_remote_exception(v, b->sender);
1246 ractor_basket_free(b);
1247 rb_exc_raise(err);
1248 }
1249
1250 ractor_basket_free(b);
1251 return v;
1252}
1253
1254// Ractor blocking by receive
1255
1256#if VM_CHECK_MODE > 0
1257static bool
1258ractor_waiter_included(rb_ractor_t *cr, rb_thread_t *th)
1259{
1260 ASSERT_ractor_locking(cr);
1261
1262 struct ractor_waiter *w;
1263
1264 ccan_list_for_each(&cr->sync.waiters, w, node) {
1265 if (w->th == th) {
1266 return true;
1267 }
1268 }
1269
1270 return false;
1271}
1272#endif
1273
1274#if USE_RUBY_DEBUG_LOG
1275
1276static const char *
1277wakeup_status_str(enum ractor_wakeup_status wakeup_status)
1278{
1279 switch (wakeup_status) {
1280 case wakeup_none: return "none";
1281 case wakeup_by_send: return "by_send";
1282 case wakeup_by_interrupt: return "by_interrupt";
1283 case wakeup_by_close: return "by_close";
1284 }
1285 rb_bug("unreachable");
1286}
1287
1288static const char *
1289basket_type_name(enum ractor_basket_type type)
1290{
1291 switch (type) {
1292 case basket_type_none: return "none";
1293 case basket_type_ref: return "ref";
1294 case basket_type_copy: return "copy";
1295 case basket_type_move: return "move";
1296 case basket_type_exit: return "exit";
1297 }
1298 VM_ASSERT(0);
1299 return NULL;
1300}
1301
1302#endif // USE_RUBY_DEBUG_LOG
1303
1304static bool
1305ractor_wakeup_all_locked(rb_ractor_t *r, enum ractor_wakeup_status wakeup_status)
1306{
1307 ASSERT_ractor_locking(r);
1308
1309 RUBY_DEBUG_LOG("r:%"PRI_SERIALT_PREFIX"u wakeup:%s", rb_ractor_id(r), wakeup_status_str(wakeup_status));
1310
1311 bool wakeup_p = false;
1312
1313 while (1) {
1314 struct ractor_waiter *waiter = ccan_list_pop(&r->sync.waiters, struct ractor_waiter, node);
1315
1316 if (waiter) {
1317 VM_ASSERT(waiter->wakeup_status == wakeup_none);
1318
1319 waiter->wakeup_status = wakeup_status;
1320 rb_ractor_sched_wakeup(r, waiter->th);
1321
1322 wakeup_p = true;
1323 }
1324 else {
1325 break;
1326 }
1327 }
1328
1329 return wakeup_p;
1330}
1331
1332static bool
1333ractor_wakeup_all(rb_ractor_t *r, enum ractor_wakeup_status wakeup_status)
1334{
1335 ASSERT_ractor_unlocking(r);
1336
1337 bool wakeup_p;
1338
1339 RACTOR_LOCK(r);
1340 {
1341 wakeup_p = ractor_wakeup_all_locked(r, wakeup_status);
1342 }
1343 RACTOR_UNLOCK(r);
1344
1345 return wakeup_p;
1346}
1347
1348static void
1349ubf_ractor_wait(void *ptr)
1350{
1351 struct ractor_waiter *waiter = (struct ractor_waiter *)ptr;
1352
1353 rb_thread_t *th = waiter->th;
1354 rb_ractor_t *r = th->ractor;
1355 rb_atomic_t event_serial = waiter->event_serial;
1356
1357 // clear ubf and nobody can kick UBF
1358 th->unblock.func = NULL;
1359 th->unblock.arg = NULL;
1360
1361 rb_native_mutex_unlock(&th->interrupt_lock);
1362 {
1363 RACTOR_LOCK(r);
1364 {
1365 if (RUBY_ATOMIC_LOAD(th->unblock.event_serial) == event_serial && waiter->wakeup_status == wakeup_none) {
1366 RUBY_DEBUG_LOG("waiter:%p", (void *)waiter);
1367
1368 waiter->wakeup_status = wakeup_by_interrupt;
1369 ccan_list_del(&waiter->node);
1370
1371 rb_ractor_sched_wakeup(r, waiter->th);
1372 }
1373 }
1374 RACTOR_UNLOCK(r);
1375 }
1376 rb_native_mutex_lock(&th->interrupt_lock);
1377}
1378
1379// Waits for an event on cr. `end` is an absolute deadline, NULL to wait forever.
1380static enum ractor_wakeup_status
1381ractor_wait(rb_execution_context_t *ec, rb_ractor_t *cr, const rb_hrtime_t *end)
1382{
1383 rb_thread_t *th = rb_ec_thread_ptr(ec);
1384
1385 struct ractor_waiter waiter = {
1386 .wakeup_status = wakeup_none,
1387 .th = th,
1388 .end = end,
1389 };
1390
1391 RUBY_DEBUG_LOG("wait%s", "");
1392
1393 ASSERT_ractor_locking(cr);
1394
1395 VM_ASSERT(GET_RACTOR() == cr);
1396 VM_ASSERT(!ractor_waiter_included(cr, th));
1397
1398 ccan_list_add_tail(&cr->sync.waiters, &waiter.node);
1399
1400 // resume another ready thread and wait for an event
1401 rb_ractor_sched_wait(ec, cr, ubf_ractor_wait, &waiter);
1402
1403 if (waiter.wakeup_status == wakeup_none) {
1404 ccan_list_del(&waiter.node);
1405 }
1406
1407 RUBY_DEBUG_LOG("wakeup_status:%s", wakeup_status_str(waiter.wakeup_status));
1408
1409 RACTOR_UNLOCK_SELF(cr);
1410 {
1411 rb_ec_check_ints(ec);
1412 }
1413 RACTOR_LOCK_SELF(cr);
1414
1415 VM_ASSERT(!ractor_waiter_included(cr, th));
1416 return waiter.wakeup_status;
1417}
1418
1419static void
1420ractor_deliver_incoming_messages(rb_execution_context_t *ec, rb_ractor_t *cr)
1421{
1422 ASSERT_ractor_locking(cr);
1423 struct ractor_queue *recv_q = cr->sync.recv_queue;
1424
1425 struct ractor_basket *b;
1426 while ((b = ractor_queue_deq(cr, recv_q)) != NULL) {
1427 ractor_queue_enq(cr, ractor_get_queue(cr, b->port_id, true), b);
1428 }
1429}
1430
1431static bool
1432ractor_check_received(rb_ractor_t *cr, struct ractor_queue *messages)
1433{
1434 struct ractor_queue *received_queue = cr->sync.recv_queue;
1435 bool received = false;
1436
1437 ASSERT_ractor_locking(cr);
1438
1439 if (ractor_queue_empty_p(cr, received_queue)) {
1440 RUBY_DEBUG_LOG("empty");
1441 }
1442 else {
1443 received = true;
1444
1445 // messages <- incoming
1446 ractor_queue_init(messages);
1447 ractor_queue_move(messages, received_queue);
1448 }
1449
1450 VM_ASSERT(ractor_queue_empty_p(cr, received_queue));
1451
1452 RUBY_DEBUG_LOG("received:%d", received);
1453 return received;
1454}
1455
1456// Returns false if the deadline `end` passed with nothing to deliver. Incoming
1457// messages are delivered even then, so the caller retries its queue once more.
1458// A wait can end on a wakeup meant for another port, so the caller keeps the
1459// deadline: a stream of them must not hold a timed receive past its time.
1460static bool
1461ractor_deadline_passed_p(const rb_hrtime_t *end)
1462{
1463 return end != NULL && rb_hrtime_now() >= *end;
1464}
1465
1466static bool
1467ractor_wait_receive(rb_execution_context_t *ec, rb_ractor_t *cr, const rb_hrtime_t *end)
1468{
1469 struct ractor_queue messages;
1470 bool deliverred = false;
1471 bool timedout = false;
1472
1473 RACTOR_LOCK_SELF(cr);
1474 {
1475 if (ractor_check_received(cr, &messages)) {
1476 deliverred = true;
1477 }
1478 else if (!end) {
1479 ractor_wait(ec, cr, NULL); // no timeout: wait until a message arrives
1480 }
1481 else if (*end == 0) {
1482 timedout = true; // `timeout: 0`: over without reading any clock
1483 }
1484 else {
1485 // only a wakeup nobody claimed can be the deadline, so only then look at
1486 // the clock: a send or an interrupt says what woke this thread by itself
1487 timedout = ractor_wait(ec, cr, end) == wakeup_none && rb_hrtime_now() >= *end;
1488 }
1489 }
1490 RACTOR_UNLOCK_SELF(cr);
1491
1492 if (deliverred) {
1493 VM_ASSERT(!ractor_queue_empty_p(cr, &messages));
1494 struct ractor_basket *b;
1495
1496 while ((b = ractor_queue_deq(cr, &messages)) != NULL) {
1497 ractor_queue_enq(cr, ractor_get_queue(cr, b->port_id, false), b);
1498 }
1499 }
1500
1501 return !timedout;
1502}
1503
1504static VALUE
1505ractor_try_receive(rb_execution_context_t *ec, rb_ractor_t *cr, const struct ractor_port *rp)
1506{
1507 struct ractor_queue *rq = ractor_get_queue(cr, ractor_port_id(rp), false);
1508
1509 if (rq == NULL) {
1510 rb_raise(rb_eRactorClosedError, "The port was already closed");
1511 }
1512
1513 struct ractor_basket *b = ractor_queue_deq(cr, rq);
1514 /* Off the queue and not yet freed: this Ractor roots it while it materializes. */
1515 if (b) ractor_off_queue_add(cr, b);
1516
1517 if (rq->closed && ractor_queue_empty_p(cr, rq)) {
1518 ractor_delete_port(cr, ractor_port_id(rp), false);
1519 }
1520
1521 if (b) {
1522 return ractor_basket_accept(b);
1523 }
1524 else {
1525 return Qundef;
1526 }
1527}
1528
1529// Returns Qundef if the deadline passed first. It bounds how long this blocks; it
1530// does not cut delivery off. A message that lands while the timeout is being
1531// reported is still returned, as Thread::Queue#pop(timeout:) does. Either way
1532// nothing is lost: a basket only leaves the queue when it is returned.
1533static VALUE
1534ractor_receive(rb_execution_context_t *ec, const struct ractor_port *rp, const rb_hrtime_t *end)
1535{
1536 rb_ractor_t *cr = rb_ec_ractor_ptr(ec);
1537 VM_ASSERT(cr == rp->r);
1538
1539 RUBY_DEBUG_LOG("port:%u", (unsigned int)ractor_port_id(rp));
1540
1541 while (1) {
1542 VALUE v = ractor_try_receive(ec, cr, rp);
1543
1544 if (v != Qundef) {
1545 return v;
1546 }
1547 else if (!ractor_wait_receive(ec, cr, end)) {
1548 return Qundef;
1549 }
1550 else if (ractor_deadline_passed_p(end)) {
1551 // The wait ended on a wakeup meant for another port, which says
1552 // nothing about the clock. One more look, then the deadline stands.
1553 return ractor_try_receive(ec, cr, rp);
1554 }
1555 }
1556}
1557
1558// A timeout argument becomes an absolute deadline, or 0 for `timeout: 0`, which
1559// every wait reads as "do not wait". Returns NULL when there is no timeout.
1560static const rb_hrtime_t *
1561ractor_timeout_deadline(VALUE timeout, rb_hrtime_t *storage)
1562{
1563 if (NIL_P(timeout)) return NULL;
1564
1565 if (!(FIXNUM_P(timeout) && FIX2LONG(timeout) == 0)) {
1566 struct timeval tv = rb_time_interval(timeout); // raises on a negative timeout
1567 rb_hrtime_t rel = rb_timeval2hrtime(&tv);
1568
1569 if (rel > 0) {
1570 *storage = rb_hrtime_add(rb_hrtime_now(), rel);
1571 return storage;
1572 }
1573 }
1574
1575 *storage = 0;
1576 return storage;
1577}
1578
1579// Ractor#send
1580
1581static void
1582ractor_send_basket(rb_execution_context_t *ec, const struct ractor_port *rp, struct ractor_basket *b, bool raise_on_error)
1583{
1584 bool closed = false;
1585
1586 RUBY_DEBUG_LOG("port:%u@r%"PRI_SERIALT_PREFIX"u b:%s v:%p", (unsigned int)ractor_port_id(rp), rb_ractor_id(rp->r), basket_type_name(b->type), (void *)b->p.v);
1587
1588 RACTOR_LOCK(rp->r);
1589 {
1590 if (ractor_closed_port_p(ec, rp->r, rp)) {
1591 closed = true;
1592 }
1593 else {
1594 b->port_id = ractor_port_id(rp);
1595 /* The receiver's queue roots it from here; drop it from ours. */
1596 ractor_off_queue_remove(b);
1597 ractor_queue_enq(rp->r, rp->r->sync.recv_queue, b);
1598 ractor_wakeup_all_locked(rp->r, wakeup_by_send);
1599 }
1600 }
1601 RACTOR_UNLOCK(rp->r);
1602
1603 // NOTE: ref r -> b->p.v is created, but Ractor is unprotected object, so no problem on that.
1604
1605 if (closed) {
1606 RUBY_DEBUG_LOG("closed:%u@r%"PRI_SERIALT_PREFIX"u", (unsigned int)ractor_port_id(rp), rb_ractor_id(rp->r));
1607
1608 /* Nothing took the basket: it was not enqueued, so free it whether or not the
1609 * caller wants the error raised. */
1610 ractor_basket_free(b);
1611
1612 if (raise_on_error) {
1613 rb_raise(rb_eRactorClosedError, "The port was already closed");
1614 }
1615 }
1616}
1617
1618/* sender is the ractor the token is about; both it and the token are shareable,
1619 * so nothing is copied until the receiver builds the pair. It also skips the tag
1620 * ractor_basket_new pushes: the tokens go out from a thread whose EC has already
1621 * lost its VM stack, and EC_PUSH_TAG reads ec->cfp under ZJIT. */
1622static struct ractor_basket *
1623ractor_basket_new_exit(VALUE sender, VALUE token)
1624{
1625 struct ractor_basket *b = ractor_basket_alloc();
1626
1627 b->type = basket_type_exit;
1628 b->sender = sender;
1629 b->p.v = token;
1630
1631 return b;
1632}
1633
1634static VALUE
1635ractor_send0(rb_execution_context_t *ec, const struct ractor_port *rp, VALUE obj, VALUE move, bool raise_on_error)
1636{
1637 struct ractor_basket *b = ractor_basket_new(ec, obj, RTEST(move) ? basket_type_move : basket_type_none, false);
1638 ractor_send_basket(ec, rp, b, raise_on_error);
1639 RB_GC_GUARD(obj);
1640 return rp->r->pub.self;
1641}
1642
1643static VALUE
1644ractor_send(rb_execution_context_t *ec, const struct ractor_port *rp, VALUE obj, VALUE move)
1645{
1646 return ractor_send0(ec, rp, obj, move, true);
1647}
1648
1649// Ractor::Selector
1650
1652 struct st_table *ports; // rpv -> rp
1653
1654};
1655
1656static int
1657ractor_selector_mark_i(st_data_t key, st_data_t val, st_data_t dmy)
1658{
1659 rb_gc_mark((VALUE)key); // rpv
1660
1661 return ST_CONTINUE;
1662}
1663
1664static void
1665ractor_selector_mark(void *ptr)
1666{
1667 struct ractor_selector *s = ptr;
1668
1669 if (s->ports) {
1670 st_foreach(s->ports, ractor_selector_mark_i, 0);
1671 }
1672}
1673
1674static void
1675ractor_selector_free(void *ptr)
1676{
1677 struct ractor_selector *s = ptr;
1678 st_free_table(s->ports);
1679 SIZED_FREE(s);
1680}
1681
1682static size_t
1683ractor_selector_memsize(const void *ptr)
1684{
1685 const struct ractor_selector *s = ptr;
1686 size_t size = sizeof(struct ractor_selector);
1687 if (s->ports) {
1688 size += st_memsize(s->ports);
1689 }
1690 return size;
1691}
1692
1693static const rb_data_type_t ractor_selector_data_type = {
1694 "ractor/selector",
1695 {
1696 ractor_selector_mark,
1697 ractor_selector_free,
1698 ractor_selector_memsize,
1699 NULL, // update
1700 },
1701 0, 0, RUBY_TYPED_THREAD_SAFE_FREE | RUBY_TYPED_WB_PROTECTED,
1702};
1703
1704static struct ractor_selector *
1705RACTOR_SELECTOR_PTR(VALUE selv)
1706{
1707 VM_ASSERT(rb_typeddata_is_kind_of(selv, &ractor_selector_data_type));
1708 return (struct ractor_selector *)DATA_PTR(selv);
1709}
1710
1711// Ractor::Selector.new
1712
1713static VALUE
1714ractor_selector_create(VALUE klass)
1715{
1716 struct ractor_selector *s;
1717 VALUE selv = TypedData_Make_Struct(klass, struct ractor_selector, &ractor_selector_data_type, s);
1718 s->ports = st_init_numtable(); // TODO
1719 return selv;
1720}
1721
1722// Ractor::Selector#add(r)
1723
1724/*
1725 * call-seq:
1726 * add(ractor) -> ractor
1727 *
1728 * Adds _ractor_ to +self+. Raises an exception if _ractor_ is already added.
1729 * Returns _ractor_.
1730 */
1731static VALUE
1732ractor_selector_add(VALUE selv, VALUE rpv)
1733{
1734 if (!ractor_port_p(rpv)) {
1735 rb_raise(rb_eArgError, "Not a Ractor::Port object");
1736 }
1737
1738 struct ractor_selector *s = RACTOR_SELECTOR_PTR(selv);
1739 const struct ractor_port *rp = ractor_port_ptr_check(rpv);
1740
1741 if (st_lookup(s->ports, (st_data_t)rpv, NULL)) {
1742 rb_raise(rb_eArgError, "already added");
1743 }
1744
1745 st_insert(s->ports, (st_data_t)rpv, (st_data_t)rp);
1746 RB_OBJ_WRITTEN(selv, Qundef, rpv);
1747
1748 return selv;
1749}
1750
1751// Ractor::Selector#remove(r)
1752
1753/* call-seq:
1754 * remove(ractor) -> ractor
1755 *
1756 * Removes _ractor_ from +self+. Raises an exception if _ractor_ is not added.
1757 * Returns the removed _ractor_.
1758 */
1759static VALUE
1760ractor_selector_remove(VALUE selv, VALUE rpv)
1761{
1762 if (!ractor_port_p(rpv)) {
1763 rb_raise(rb_eArgError, "Not a Ractor::Port object");
1764 }
1765
1766 struct ractor_selector *s = RACTOR_SELECTOR_PTR(selv);
1767
1768 if (!st_lookup(s->ports, (st_data_t)rpv, NULL)) {
1769 rb_raise(rb_eArgError, "not added yet");
1770 }
1771
1772 st_delete(s->ports, (st_data_t *)&rpv, NULL);
1773
1774 return selv;
1775}
1776
1777// Ractor::Selector#clear
1778
1779/*
1780 * call-seq:
1781 * clear -> self
1782 *
1783 * Removes all ractors from +self+. Raises +self+.
1784 */
1785static VALUE
1786ractor_selector_clear(VALUE selv)
1787{
1788 struct ractor_selector *s = RACTOR_SELECTOR_PTR(selv);
1789 st_clear(s->ports);
1790 return selv;
1791}
1792
1793/*
1794 * call-seq:
1795 * empty? -> true or false
1796 *
1797 * Returns +true+ if no ractor is added.
1798 */
1799static VALUE
1800ractor_selector_empty_p(VALUE selv)
1801{
1802 struct ractor_selector *s = RACTOR_SELECTOR_PTR(selv);
1803 return s->ports->num_entries == 0 ? Qtrue : Qfalse;
1804}
1805
1806// Ractor::Selector#wait
1807
1809 rb_ractor_t *cr;
1811 bool found;
1812 VALUE v;
1813 VALUE rpv;
1814};
1815
1816static int
1817ractor_selector_wait_i(st_data_t key, st_data_t val, st_data_t data)
1818{
1819 struct ractor_selector_wait_data *p = (struct ractor_selector_wait_data *)data;
1820 const struct ractor_port *rp = (const struct ractor_port *)val;
1821
1822 VALUE v = ractor_try_receive(p->ec, p->cr, rp);
1823
1824 if (v != Qundef) {
1825 p->found = true;
1826 p->v = v;
1827 p->rpv = (VALUE)key;
1828 return ST_STOP;
1829 }
1830 else {
1831 return ST_CONTINUE;
1832 }
1833}
1834
1835static VALUE
1836ractor_selector__wait(rb_execution_context_t *ec, VALUE selector, const rb_hrtime_t *end)
1837{
1838 rb_ractor_t *cr = rb_ec_ractor_ptr(ec);
1839 struct ractor_selector *s = RACTOR_SELECTOR_PTR(selector);
1840
1841 struct ractor_selector_wait_data data = {
1842 .ec = ec,
1843 .cr = cr,
1844 .found = false,
1845 };
1846
1847 while (1) {
1848 st_foreach(s->ports, ractor_selector_wait_i, (st_data_t)&data);
1849
1850 if (data.found) {
1851 return rb_ary_new_from_args(2, data.rpv, data.v);
1852 }
1853 else if (!ractor_wait_receive(ec, cr, end)) {
1854 return Qnil;
1855 }
1856 else if (ractor_deadline_passed_p(end)) {
1857 st_foreach(s->ports, ractor_selector_wait_i, (st_data_t)&data);
1858 return data.found ? rb_ary_new_from_args(2, data.rpv, data.v) : Qnil;
1859 }
1860 }
1861}
1862
1863/*
1864 * call-seq:
1865 * wait(receive: false, yield_value: undef, move: false) -> [ractor, value]
1866 *
1867 * Waits until any ractor in _selector_ can be active.
1868 */
1869static VALUE
1870ractor_selector_wait(VALUE selector)
1871{
1872 return ractor_selector__wait(GET_EC(), selector, NULL);
1873}
1874
1875static VALUE
1876ractor_selector_new(int argc, VALUE *ractors, VALUE klass)
1877{
1878 VALUE selector = ractor_selector_create(klass);
1879
1880 for (int i=0; i<argc; i++) {
1881 ractor_selector_add(selector, ractors[i]);
1882 }
1883
1884 return selector;
1885}
1886
1887static VALUE
1888ractor_select_internal(rb_execution_context_t *ec, VALUE self, VALUE ports, VALUE timeout)
1889{
1890 rb_hrtime_t deadline;
1891 const rb_hrtime_t *end = ractor_timeout_deadline(timeout, &deadline);
1892
1893 VALUE selector = ractor_selector_new(RARRAY_LENINT(ports), (VALUE *)RARRAY_CONST_PTR(ports), rb_cRactorSelector);
1894 VALUE result = ractor_selector__wait(ec, selector, end);
1895
1896 RB_GC_GUARD(selector);
1897 RB_GC_GUARD(ports);
1898 return result;
1899}
1900
1901#ifndef USE_RACTOR_SELECTOR
1902#define USE_RACTOR_SELECTOR 0
1903#endif
1904
1905RUBY_SYMBOL_EXPORT_BEGIN
1906void rb_init_ractor_selector(void);
1907RUBY_SYMBOL_EXPORT_END
1908
1909/*
1910 * Document-class: Ractor::Selector
1911 * :nodoc: currently
1912 *
1913 * Selects multiple Ractors to be activated.
1914 */
1915void
1916rb_init_ractor_selector(void)
1917{
1918 rb_cRactorSelector = rb_define_class_under(rb_cRactor, "Selector", rb_cObject);
1919 rb_undef_alloc_func(rb_cRactorSelector);
1920
1921 rb_define_singleton_method(rb_cRactorSelector, "new", ractor_selector_new , -1);
1922 rb_define_method(rb_cRactorSelector, "add", ractor_selector_add, 1);
1923 rb_define_method(rb_cRactorSelector, "remove", ractor_selector_remove, 1);
1924 rb_define_method(rb_cRactorSelector, "clear", ractor_selector_clear, 0);
1925 rb_define_method(rb_cRactorSelector, "empty?", ractor_selector_empty_p, 0);
1926 rb_define_method(rb_cRactorSelector, "wait", ractor_selector_wait, 0);
1927}
1928
1929static void
1930Init_RactorPort(void)
1931{
1932 rb_cRactorPort = rb_define_class_under(rb_cRactor, "Port", rb_cObject);
1933 rb_define_alloc_func(rb_cRactorPort, ractor_port_alloc);
1934 rb_define_method(rb_cRactorPort, "initialize", ractor_port_initialize, 0);
1935 rb_define_method(rb_cRactorPort, "initialize_copy", ractor_port_initialize_copy, 1);
1936
1937#if USE_RACTOR_SELECTOR
1938 rb_init_ractor_selector();
1939#endif
1940}
std::atomic< unsigned > rb_atomic_t
Type that is eligible for atomic operations.
Definition atomic.h:69
#define RUBY_ATOMIC_LOAD(var)
Atomic load.
Definition atomic.h:175
#define rb_define_method(klass, mid, func, arity)
Defines klass#mid.
#define rb_define_singleton_method(klass, mid, func, arity)
Defines klass.mid.
#define ALLOC
Old name of RB_ALLOC.
Definition memory.h:400
#define Qundef
Old name of RUBY_Qundef.
#define ID2SYM
Old name of RB_ID2SYM.
Definition symbol.h:44
#define UNREACHABLE_RETURN
Old name of RBIMPL_UNREACHABLE_RETURN.
Definition assume.h:29
#define T_NONE
Old name of RUBY_T_NONE.
Definition value_type.h:74
#define Qtrue
Old name of RUBY_Qtrue.
#define Qnil
Old name of RUBY_Qnil.
#define Qfalse
Old name of RUBY_Qfalse.
#define FIX2LONG
Old name of RB_FIX2LONG.
Definition long.h:46
#define NIL_P
Old name of RB_NIL_P.
#define FIXNUM_P
Old name of RB_FIXNUM_P.
void rb_exc_raise(VALUE mesg)
Raises an exception in the current thread.
Definition eval.c:678
VALUE rb_eTypeError
TypeError exception.
Definition error.c:1473
VALUE rb_cObject
Object class.
Definition object.c:60
VALUE rb_cRactor
Ractor class.
Definition ractor.c:38
static VALUE rb_class_of(VALUE obj)
Object to class mapping function.
Definition globals.h:174
VALUE rb_obj_class(VALUE obj)
Queries the class of an object.
Definition object.c:234
VALUE rb_obj_freeze(VALUE obj)
Same as RB_OBJ_FREEZE(), but returns the given object.
Definition object.c:1309
#define RB_OBJ_WRITTEN(old, oldv, young)
Identical to RB_OBJ_WRITE(), except it doesn't write any values, but only a WB declaration.
Definition gc.h:504
VALUE rb_marshal_dump(VALUE obj, VALUE port)
Serialises the given object and all its referring objects, to write them down to the passed port.
Definition marshal.c:2724
#define rb_exc_new_cstr(exc, str)
Identical to rb_exc_new(), except it assumes the passed pointer is a pointer to a C string.
Definition string.h:1695
void rb_thread_schedule(void)
Tries to switch to another thread.
Definition thread.c:1693
struct timeval rb_time_interval(VALUE num)
Creates a "time interval".
Definition time.c:2984
VALUE rb_ivar_set(VALUE obj, ID name, VALUE val)
Identical to rb_iv_set(), except it accepts the name as an ID instead of a C string.
Definition variable.c:2141
void rb_undef_alloc_func(VALUE klass)
Deletes the allocator function of a class.
Definition vm_method.c:1846
void rb_define_alloc_func(VALUE klass, rb_alloc_func_t func)
Sets the allocator function of a class.
#define RB_OBJ_SET_SHAREABLE(obj)
Wrapper of rb_obj_set_shareable().
Definition ractor.h:290
static bool rb_ractor_shareable_p(VALUE obj)
Queries if multiple Ractors can share the passed object or not.
Definition ractor.h:269
#define RBIMPL_ATTR_MAYBE_UNUSED()
Wraps (or simulates) [[maybe_unused]]
#define RB_GC_GUARD(v)
Prevents premature destruction of local objects.
Definition memory.h:167
VALUE type(ANYARGS)
ANYARGS-ed function type.
static int RARRAY_LENINT(VALUE ary)
Identical to rb_array_len(), except it differs for the return type.
Definition rarray.h:280
#define RARRAY_CONST_PTR
Just another name of rb_array_const_ptr.
Definition rarray.h:51
#define RUBY_TYPED_DEFAULT_FREE
This is a value you can set to rb_data_type_struct::dfree.
Definition rtypeddata.h:81
#define DATA_PTR(obj)
Convenient casting macro for backward compatibility.
Definition rtypeddata.h:439
#define TypedData_Make_Struct(klass, type, data_type, sval)
Identical to TypedData_Wrap_Struct, except it allocates a new data region internally instead of takin...
Definition rtypeddata.h:604
#define RTEST
This is an old name of RB_TEST.
This is the struct that holds necessary info for a struct.
Definition rtypeddata.h:242
Definition st.h:79
void rb_native_mutex_lock(rb_nativethread_lock_t *lock)
Just another name of rb_nativethread_lock_lock.
void rb_native_mutex_initialize(rb_nativethread_lock_t *lock)
Just another name of rb_nativethread_lock_initialize.
void rb_native_mutex_unlock(rb_nativethread_lock_t *lock)
Just another name of rb_nativethread_lock_unlock.
uintptr_t VALUE
Type that represents a Ruby object.
Definition value.h:40
static bool RB_TYPE_P(VALUE obj, enum ruby_value_type t)
Queries if the given object is of given type.
Definition value_type.h:376