14static VALUE rb_cRactorPort;
20static void ractor_add_port(
rb_ractor_t *r, st_data_t
id);
31static void ractor_port_note_alive(
const struct ractor_port *rp);
34ractor_port_mark(
void *ptr)
39 rb_gc_mark(rp->r->pub.self);
44 if (rb_gc_single_objspace_p() || rb_gc_during_global_gc_p()) {
45 ractor_port_note_alive(rp);
58 0, 0, RUBY_TYPED_THREAD_SAFE_FREE | RUBY_TYPED_WB_PROTECTED | RUBY_TYPED_FROZEN_SHAREABLE | RUBY_TYPED_EMBEDDABLE,
65 return cr->sync.next_port_id++;
69RACTOR_PORT_PTR(
VALUE self)
71 VM_ASSERT(rb_typeddata_is_kind_of(self, &ractor_port_data_type));
72 return RTYPEDDATA_GET_DATA(self);
77ractor_port_ptr_check(
VALUE self)
81 if (UNLIKELY(rp->r == NULL)) {
89ractor_port_alloc(
VALUE klass)
102 if (!r->sync.ports) {
103 rb_raise(rb_eRactorClosedError,
"The ractor has terminated");
110 rp->id_ = ractor_genid_for_port(r);
112 ractor_add_port(r, ractor_port_id(rp));
126ractor_port_initialize(
VALUE self)
128 return ractor_port_init(self, GET_RACTOR());
133ractor_port_initialize_copy(
VALUE self,
VALUE orig)
136 struct ractor_port *src = ractor_port_ptr_check(orig);
139 dst->id_ = ractor_port_id(src);
147 VALUE rpv = ractor_port_alloc(rb_cRactorPort);
148 ractor_port_init(rpv, r);
153ractor_port_p(
VALUE self)
155 return rb_typeddata_is_kind_of(self, &ractor_port_data_type);
158static const rb_hrtime_t *ractor_timeout_deadline(
VALUE timeout, rb_hrtime_t *storage);
163 const struct ractor_port *rp = ractor_port_ptr_check(self);
165 if (rp->r != rb_ec_ractor_ptr(ec)) {
166 rb_raise(rb_eRactorError,
"only allowed from the creator Ractor of this port");
169 rb_hrtime_t deadline;
170 const rb_hrtime_t *end = ractor_timeout_deadline(timeout, &deadline);
172 VALUE v = ractor_receive(ec, rp, end);
176 return UNDEF_P(v) ?
Qnil : v;
182 const struct ractor_port *rp = ractor_port_ptr_check(self);
183 ractor_send(ec, rp, obj,
RTEST(move));
194 const struct ractor_port *rp = ractor_port_ptr_check(self);
198 if (rb_ec_ractor_ptr(ec) == r) {
201 closed = ractor_closed_port_p(ec, r, rp);
210 closed = ractor_closed_port_p(ec, r, rp);
221 const struct ractor_port *rp = ractor_port_ptr_check(self);
225 rb_raise(rb_eRactorError,
"closing port by other ractors is not allowed");
228 ractor_close_port(ec, cr, rp);
236enum ractor_basket_type {
251 enum ractor_basket_type type;
263 struct ccan_list_node node;
264 struct ccan_list_node off_queue_node;
269ractor_basket_type_p(
const struct ractor_basket *b,
enum ractor_basket_type type)
271 return b->type == type;
277 return ractor_basket_type_p(b, basket_type_none);
284 if (b->p.courier != NULL) {
289 rb_ractor_courier_mark(b->p.courier);
296 rb_gc_mark(b->sender);
302 ractor_off_queue_remove(b);
305 rb_ractor_courier_free(b->p.courier);
312ractor_basket_alloc(
void)
318 b->type = basket_type_none;
322 b->p.exception =
false;
324 ccan_list_node_init(&b->off_queue_node);
334 VM_ASSERT(cr == rb_current_ractor_raw(
false));
335 ccan_list_add_tail(&cr->sync.off_queue_baskets, &b->off_queue_node);
341 ccan_list_del_init(&b->off_queue_node);
348 ccan_list_for_each(&r->sync.off_queue_baskets, b, off_queue_node) {
349 ractor_basket_mark(b);
356 struct ccan_list_head set;
364 ccan_list_head_init(&rq->set);
370ractor_queue_new(
void)
373 ractor_queue_init(rq);
378ractor_port_note_alive(
const struct ractor_port *rp)
382 if (rp->r->sync.ports && st_lookup(rp->r->sync.ports, rp->id_, (st_data_t *)&rq)) {
395 ccan_list_for_each(&rq->set, b, node) {
396 ractor_basket_mark(b);
405 ccan_list_for_each_safe(&rq->set, b, nxt, node) {
406 ccan_list_del_init(&b->node);
407 ractor_basket_free(b);
410 VM_ASSERT(ccan_list_empty(&rq->set));
422 ccan_list_for_each(&rq->set, b, node) {
432 bool closed_now = !rq->closed;
440 struct ccan_list_head *src = &src_rq->set;
441 struct ccan_list_head *dst = &dst_rq->set;
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);
461 return ccan_list_empty(&rq->set);
467 VM_ASSERT(GET_RACTOR() == r);
475 ccan_list_add_tail(&rq->set, &basket->node);
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);
491static void ractor_delete_port(
rb_ractor_t *cr, st_data_t
id,
bool locked);
494ractor_get_queue(
rb_ractor_t *cr, st_data_t
id,
bool locked)
496 VM_ASSERT(cr == GET_RACTOR());
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);
520 ASSERT_ractor_unlocking(r);
522 RUBY_DEBUG_LOG(
"id:%u", (
unsigned int)
id);
526 st_table *
const old_tab = r->sync.ports;
531 inserted = st_insert_no_rebuild(old_tab,
id, (st_data_t)rq) >= 0;
548 const st_index_t entries = st_table_size(old_tab);
550 new_tab = st_copy(old_tab);
552 if (st_table_size(old_tab) == entries) {
553 st_insert(new_tab,
id, (st_data_t)rq);
555 if (st_table_size(old_tab) == entries)
break;
558 st_free_table(new_tab);
563 VM_ASSERT(r->sync.ports == old_tab);
564 r->sync.ports = new_tab;
568 st_free_table(old_tab);
573ractor_delete_port_locked(
rb_ractor_t *cr, st_data_t
id)
575 ASSERT_ractor_locking(cr);
577 RUBY_DEBUG_LOG(
"id:%u", (
unsigned int)
id);
581 if (st_delete(cr->sync.ports, &
id, (st_data_t *)&rq)) {
582 ractor_queue_free(rq);
590ractor_delete_port(
rb_ractor_t *cr, st_data_t
id,
bool locked)
593 ractor_delete_port_locked(cr,
id);
596 RACTOR_LOCK_SELF(cr);
598 ractor_delete_port_locked(cr,
id);
600 RACTOR_UNLOCK_SELF(cr);
607 return RACTOR_PORT_PTR(r->sync.default_port_value);
613 return r->sync.default_port_value;
619 VM_ASSERT(rb_ec_ractor_ptr(ec) == rp->r ? 1 : (ASSERT_ractor_locking(rp->r), 1));
623 if (rp->r->sync.ports && st_lookup(rp->r->sync.ports, ractor_port_id(rp), (st_data_t *)&rq)) {
634static bool ractor_wakeup_all(
rb_ractor_t *r,
enum ractor_wakeup_status wakeup_status);
639 VM_ASSERT(cr == rp->r);
641 bool closed_now =
false;
643 RACTOR_LOCK_SELF(cr);
645 ractor_deliver_incoming_messages(ec, cr);
647 if (st_lookup(rp->r->sync.ports, ractor_port_id(rp), (st_data_t *)&rq)) {
648 closed_now = ractor_queue_close(rq);
650 if (ractor_queue_empty_p(cr, rq)) {
652 ractor_delete_port(cr, ractor_port_id(rp),
true);
658 RACTOR_UNLOCK_SELF(cr);
663 ractor_wakeup_all(cr, wakeup_by_close);
673ractor_reap_dead_ports_i(st_data_t port_id, st_data_t val, st_data_t dat)
682 ractor_queue_free(rq);
694 if (recv_q == NULL)
return;
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);
709 VM_ASSERT(rb_gc_single_objspace_p() || rb_gc_during_global_gc_p());
712 st_foreach(r->sync.ports, ractor_reap_dead_ports_i, 0);
713 ractor_reap_undeliverable_messages(r);
718ractor_free_all_ports_i(st_data_t port_id, st_data_t val, st_data_t dat)
723 ractor_queue_free(rq);
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;
736 if (cr->sync.recv_queue) {
737 ractor_queue_free(cr->sync.recv_queue);
738 cr->sync.recv_queue = NULL;
742#if defined(HAVE_WORKING_FORK)
746 ractor_free_all_ports(r);
747 r->sync.legacy =
Qnil;
755 struct ccan_list_node node;
765 ccan_list_for_each(&r->sync.monitors, rm, node) {
766 rb_gc_mark(rm->port.r->pub.self);
775 if (r->sync.legacy_exc) {
776 RUBY_DEBUG_LOG(
"aborted");
780 RUBY_DEBUG_LOG(
"exited");
789 bool terminated =
false;
790 const struct ractor_port *rp = ractor_port_ptr_check(port);
796 if (UNDEF_P(r->sync.legacy)) {
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);
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));
809 ractor_send_basket(ec, rp, ractor_basket_new_exit(self, ractor_exit_token(r)),
false);
821 const struct ractor_port *rp = ractor_port_ptr_check(port);
825 if (UNDEF_P(r->sync.legacy)) {
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);
848 RUBY_DEBUG_LOG(
"exc:%d", exc);
849 VM_ASSERT(!UNDEF_P(legacy));
850 VM_ASSERT(cr->sync.legacy ==
Qundef);
852 RACTOR_LOCK_SELF(cr);
854 ractor_free_all_ports(cr);
856 cr->sync.legacy = legacy;
857 cr->sync.legacy_exc = exc;
859 RACTOR_UNLOCK_SELF(cr);
868 VALUE token = ractor_exit_token(cr);
871 ccan_list_for_each_safe(&cr->sync.monitors, rm, nxt, node)
873 RUBY_DEBUG_LOG(
"port:%u@r%u", (
unsigned int)ractor_port_id(&rm->port), (
unsigned int)rb_ractor_id(rm->port.r));
875 ractor_send_basket(ec, &rm->port, ractor_basket_new_exit(cr->pub.self, token),
false);
877 ccan_list_del(&rm->node);
881 VM_ASSERT(ccan_list_empty(&cr->sync.monitors));
887ractor_mark_ports_i(st_data_t key, st_data_t val, st_data_t data)
891 ractor_queue_mark(rq);
900 const bool world_stopped = rb_gc_during_global_gc_p();
901 VM_ASSERT(world_stopped || r == rb_current_ractor_raw(
false));
903 rb_gc_mark(r->sync.default_port_value);
907 rb_gc_mark(r->sync.legacy);
915 if (!world_stopped) RACTOR_LOCK_SELF(r);
917 ractor_queue_mark(r->sync.recv_queue);
918 st_foreach(r->sync.ports, ractor_mark_ports_i, 0);
919 ractor_mark_monitors(r);
921 if (!world_stopped) RACTOR_UNLOCK_SELF(r);
932 ractor_mark_off_queue_baskets(r);
937ractor_sync_free_ports_i(st_data_t _key, st_data_t val, st_data_t _args)
941 ractor_queue_free(queue);
949 if (r->sync.recv_queue) {
950 ractor_queue_free(r->sync.recv_queue);
955 st_foreach(r->sync.ports, ractor_sync_free_ports_i, 0);
956 st_free_table(r->sync.ports);
957 r->sync.ports = NULL;
965 return st_memsize(r->sync.ports);
979 ccan_list_head_init(&r->sync.off_queue_baskets);
980 ccan_list_head_init(&r->sync.monitors);
983 ccan_list_head_init(&r->sync.waiters);
986 r->sync.recv_queue = ractor_queue_new();
989 r->sync.ports = st_init_numtable();
992 r->sync.default_port_value =
Qfalse;
1006 VM_ASSERT(r->sync.default_port_value ==
Qfalse);
1007 r->sync.default_port_value = ractor_port_new(r);
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;
1021 return r->sync.successor;
1025ractor_make_remote_exception(
VALUE cause,
VALUE sender)
1029 rb_ec_setup_exception(NULL, err, cause);
1038 rb_ractor_t *sr = ractor_set_successor_once(r, cr);
1041 if (r->sync.legacy_taken) {
1042 rb_raise(rb_eRactorError,
"The value was already taken");
1049 while (!rb_ractor_status_p(r, ractor_terminated)) {
1055 if (r->sync.legacy_taken) {
1056 rb_raise(rb_eRactorError,
"The value was already taken");
1061 rb_ractor_absorb_registered_marks(GET_RACTOR(), r);
1062 rb_ractor_absorb_registered_addrs_without_gc(GET_RACTOR(), r);
1064 rb_gc_objspace_absorb_into_current(&r->objspace);
1068 volatile VALUE legacy_keep = r->sync.legacy;
1073 ractor_local_storage_free(r);
1074 r->local_storage = NULL;
1075 r->idkey_local_storage = NULL;
1080 VALUE legacy = r->sync.legacy;
1081 r->sync.legacy =
Qnil;
1082 r->sync.legacy_taken =
true;
1085 if (r->sync.legacy_exc) {
1086 rb_exc_raise(ractor_make_remote_exception(legacy, self));
1091 rb_raise(rb_eRactorError,
"Only the successor ractor can take a value");
1096ractor_marshal_dump_body(
VALUE obj)
1102ractor_marshal_dump_rescue(
VALUE obj,
VALUE errinfo)
1104 rb_raise(rb_eRactorError,
"can not copy %"PRIsVALUE
" object.",
rb_class_of(obj));
1113 case basket_type_ref:
1117 *ptype = basket_type_ref;
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));
1134#if RBIMPL_COMPILER_IS(GCC) && defined(__OPTIMIZE__)
1141 enum ractor_basket_type type,
bool exc)
1143 b->p.exception = exc;
1144 if (type == basket_type_move) {
1149 rb_ractor_courier_build_move(obj, &b->p.courier);
1154 VALUE v = ractor_prepare_payload(ec, obj, &type, &b->p.courier);
1168 ractor_off_queue_add(cr, b);
1170 enum ruby_tag_type state;
1172 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
1173 ractor_basket_build_payload(ec, b, obj, type, exc);
1176 if (state != TAG_NONE) {
1177 ractor_basket_free(b);
1178 EC_JUMP_TAG(ec, state);
1188 case basket_type_ref:
1190 case basket_type_exit:
1192 return rb_ary_new_from_args(2, b->sender, b->p.v);
1193 case basket_type_copy:
1196 case basket_type_move: {
1214 enum ruby_tag_type state;
1216 if ((state = EC_EXEC_TAG()) == TAG_NONE) {
1217 result = rb_ractor_courier_materialize(courier);
1220 if (state != TAG_NONE) {
1222 ractor_basket_free(b);
1223 EC_JUMP_TAG(ec, state);
1225 rb_ractor_courier_free(courier);
1226 b->p.courier = NULL;
1242 VALUE v = ractor_basket_value(b);
1244 if (b->p.exception) {
1245 VALUE err = ractor_make_remote_exception(v, b->sender);
1246 ractor_basket_free(b);
1250 ractor_basket_free(b);
1256#if VM_CHECK_MODE > 0
1260 ASSERT_ractor_locking(cr);
1264 ccan_list_for_each(&cr->sync.waiters, w, node) {
1274#if USE_RUBY_DEBUG_LOG
1277wakeup_status_str(
enum ractor_wakeup_status wakeup_status)
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";
1285 rb_bug(
"unreachable");
1289basket_type_name(
enum ractor_basket_type
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";
1305ractor_wakeup_all_locked(
rb_ractor_t *r,
enum ractor_wakeup_status wakeup_status)
1307 ASSERT_ractor_locking(r);
1309 RUBY_DEBUG_LOG(
"r:%"PRI_SERIALT_PREFIX
"u wakeup:%s", rb_ractor_id(r), wakeup_status_str(wakeup_status));
1311 bool wakeup_p =
false;
1317 VM_ASSERT(waiter->wakeup_status == wakeup_none);
1319 waiter->wakeup_status = wakeup_status;
1320 rb_ractor_sched_wakeup(r, waiter->th);
1333ractor_wakeup_all(
rb_ractor_t *r,
enum ractor_wakeup_status wakeup_status)
1335 ASSERT_ractor_unlocking(r);
1341 wakeup_p = ractor_wakeup_all_locked(r, wakeup_status);
1349ubf_ractor_wait(
void *ptr)
1358 th->unblock.func = NULL;
1359 th->unblock.arg = NULL;
1365 if (
RUBY_ATOMIC_LOAD(th->unblock.event_serial) == event_serial && waiter->wakeup_status == wakeup_none) {
1366 RUBY_DEBUG_LOG(
"waiter:%p", (
void *)waiter);
1368 waiter->wakeup_status = wakeup_by_interrupt;
1369 ccan_list_del(&waiter->node);
1371 rb_ractor_sched_wakeup(r, waiter->th);
1380static enum ractor_wakeup_status
1386 .wakeup_status = wakeup_none,
1391 RUBY_DEBUG_LOG(
"wait%s",
"");
1393 ASSERT_ractor_locking(cr);
1395 VM_ASSERT(GET_RACTOR() == cr);
1396 VM_ASSERT(!ractor_waiter_included(cr, th));
1398 ccan_list_add_tail(&cr->sync.waiters, &waiter.node);
1401 rb_ractor_sched_wait(ec, cr, ubf_ractor_wait, &waiter);
1403 if (waiter.wakeup_status == wakeup_none) {
1404 ccan_list_del(&waiter.node);
1407 RUBY_DEBUG_LOG(
"wakeup_status:%s", wakeup_status_str(waiter.wakeup_status));
1409 RACTOR_UNLOCK_SELF(cr);
1411 rb_ec_check_ints(ec);
1413 RACTOR_LOCK_SELF(cr);
1415 VM_ASSERT(!ractor_waiter_included(cr, th));
1416 return waiter.wakeup_status;
1422 ASSERT_ractor_locking(cr);
1426 while ((b = ractor_queue_deq(cr, recv_q)) != NULL) {
1427 ractor_queue_enq(cr, ractor_get_queue(cr, b->port_id,
true), b);
1434 struct ractor_queue *received_queue = cr->sync.recv_queue;
1435 bool received =
false;
1437 ASSERT_ractor_locking(cr);
1439 if (ractor_queue_empty_p(cr, received_queue)) {
1440 RUBY_DEBUG_LOG(
"empty");
1446 ractor_queue_init(messages);
1447 ractor_queue_move(messages, received_queue);
1450 VM_ASSERT(ractor_queue_empty_p(cr, received_queue));
1452 RUBY_DEBUG_LOG(
"received:%d", received);
1461ractor_deadline_passed_p(
const rb_hrtime_t *end)
1463 return end != NULL && rb_hrtime_now() >= *end;
1470 bool deliverred =
false;
1471 bool timedout =
false;
1473 RACTOR_LOCK_SELF(cr);
1475 if (ractor_check_received(cr, &messages)) {
1479 ractor_wait(ec, cr, NULL);
1481 else if (*end == 0) {
1487 timedout = ractor_wait(ec, cr, end) == wakeup_none && rb_hrtime_now() >= *end;
1490 RACTOR_UNLOCK_SELF(cr);
1493 VM_ASSERT(!ractor_queue_empty_p(cr, &messages));
1496 while ((b = ractor_queue_deq(cr, &messages)) != NULL) {
1497 ractor_queue_enq(cr, ractor_get_queue(cr, b->port_id,
false), b);
1507 struct ractor_queue *rq = ractor_get_queue(cr, ractor_port_id(rp),
false);
1510 rb_raise(rb_eRactorClosedError,
"The port was already closed");
1515 if (b) ractor_off_queue_add(cr, b);
1517 if (rq->closed && ractor_queue_empty_p(cr, rq)) {
1518 ractor_delete_port(cr, ractor_port_id(rp),
false);
1522 return ractor_basket_accept(b);
1537 VM_ASSERT(cr == rp->r);
1539 RUBY_DEBUG_LOG(
"port:%u", (
unsigned int)ractor_port_id(rp));
1542 VALUE v = ractor_try_receive(ec, cr, rp);
1547 else if (!ractor_wait_receive(ec, cr, end)) {
1550 else if (ractor_deadline_passed_p(end)) {
1553 return ractor_try_receive(ec, cr, rp);
1560static const rb_hrtime_t *
1561ractor_timeout_deadline(
VALUE timeout, rb_hrtime_t *storage)
1563 if (
NIL_P(timeout))
return NULL;
1567 rb_hrtime_t rel = rb_timeval2hrtime(&tv);
1570 *storage = rb_hrtime_add(rb_hrtime_now(), rel);
1584 bool closed =
false;
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);
1590 if (ractor_closed_port_p(ec, rp->r, rp)) {
1594 b->port_id = ractor_port_id(rp);
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);
1601 RACTOR_UNLOCK(rp->r);
1606 RUBY_DEBUG_LOG(
"closed:%u@r%"PRI_SERIALT_PREFIX
"u", (
unsigned int)ractor_port_id(rp), rb_ractor_id(rp->r));
1610 ractor_basket_free(b);
1612 if (raise_on_error) {
1613 rb_raise(rb_eRactorClosedError,
"The port was already closed");
1623ractor_basket_new_exit(
VALUE sender,
VALUE token)
1627 b->type = basket_type_exit;
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);
1640 return rp->r->pub.self;
1646 return ractor_send0(ec, rp, obj, move,
true);
1657ractor_selector_mark_i(st_data_t key, st_data_t val, st_data_t dmy)
1659 rb_gc_mark((
VALUE)key);
1665ractor_selector_mark(
void *ptr)
1670 st_foreach(s->ports, ractor_selector_mark_i, 0);
1675ractor_selector_free(
void *ptr)
1678 st_free_table(s->ports);
1683ractor_selector_memsize(
const void *ptr)
1688 size += st_memsize(s->ports);
1696 ractor_selector_mark,
1697 ractor_selector_free,
1698 ractor_selector_memsize,
1701 0, 0, RUBY_TYPED_THREAD_SAFE_FREE | RUBY_TYPED_WB_PROTECTED,
1705RACTOR_SELECTOR_PTR(
VALUE selv)
1707 VM_ASSERT(rb_typeddata_is_kind_of(selv, &ractor_selector_data_type));
1714ractor_selector_create(
VALUE klass)
1718 s->ports = st_init_numtable();
1734 if (!ractor_port_p(rpv)) {
1735 rb_raise(rb_eArgError,
"Not a Ractor::Port object");
1739 const struct ractor_port *rp = ractor_port_ptr_check(rpv);
1741 if (st_lookup(s->ports, (st_data_t)rpv, NULL)) {
1742 rb_raise(rb_eArgError,
"already added");
1745 st_insert(s->ports, (st_data_t)rpv, (st_data_t)rp);
1762 if (!ractor_port_p(rpv)) {
1763 rb_raise(rb_eArgError,
"Not a Ractor::Port object");
1768 if (!st_lookup(s->ports, (st_data_t)rpv, NULL)) {
1769 rb_raise(rb_eArgError,
"not added yet");
1772 st_delete(s->ports, (st_data_t *)&rpv, NULL);
1786ractor_selector_clear(
VALUE selv)
1800ractor_selector_empty_p(
VALUE selv)
1803 return s->ports->num_entries == 0 ?
Qtrue :
Qfalse;
1817ractor_selector_wait_i(st_data_t key, st_data_t val, st_data_t data)
1822 VALUE v = ractor_try_receive(p->ec, p->cr, rp);
1827 p->rpv = (
VALUE)key;
1848 st_foreach(s->ports, ractor_selector_wait_i, (st_data_t)&data);
1851 return rb_ary_new_from_args(2, data.rpv, data.v);
1853 else if (!ractor_wait_receive(ec, cr, end)) {
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;
1870ractor_selector_wait(
VALUE selector)
1872 return ractor_selector__wait(GET_EC(), selector, NULL);
1876ractor_selector_new(
int argc,
VALUE *ractors,
VALUE klass)
1878 VALUE selector = ractor_selector_create(klass);
1880 for (
int i=0; i<argc; i++) {
1881 ractor_selector_add(selector, ractors[i]);
1890 rb_hrtime_t deadline;
1891 const rb_hrtime_t *end = ractor_timeout_deadline(timeout, &deadline);
1894 VALUE result = ractor_selector__wait(ec, selector, end);
1901#ifndef USE_RACTOR_SELECTOR
1902#define USE_RACTOR_SELECTOR 0
1905RUBY_SYMBOL_EXPORT_BEGIN
1906void rb_init_ractor_selector(
void);
1907RUBY_SYMBOL_EXPORT_END
1916rb_init_ractor_selector(
void)
1925 rb_define_method(rb_cRactorSelector,
"empty?", ractor_selector_empty_p, 0);
1930Init_RactorPort(
void)
1935 rb_define_method(rb_cRactorPort,
"initialize_copy", ractor_port_initialize_copy, 1);
1937#if USE_RACTOR_SELECTOR
1938 rb_init_ractor_selector();
std::atomic< unsigned > rb_atomic_t
Type that is eligible for atomic operations.
#define RUBY_ATOMIC_LOAD(var)
Atomic load.
#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.
#define Qundef
Old name of RUBY_Qundef.
#define ID2SYM
Old name of RB_ID2SYM.
#define UNREACHABLE_RETURN
Old name of RBIMPL_UNREACHABLE_RETURN.
#define T_NONE
Old name of RUBY_T_NONE.
#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.
#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.
VALUE rb_eTypeError
TypeError exception.
VALUE rb_cObject
Object class.
VALUE rb_cRactor
Ractor class.
static VALUE rb_class_of(VALUE obj)
Object to class mapping function.
VALUE rb_obj_class(VALUE obj)
Queries the class of an object.
VALUE rb_obj_freeze(VALUE obj)
Same as RB_OBJ_FREEZE(), but returns the given object.
#define RB_OBJ_WRITTEN(old, oldv, young)
Identical to RB_OBJ_WRITE(), except it doesn't write any values, but only a WB declaration.
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.
#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.
void rb_thread_schedule(void)
Tries to switch to another thread.
struct timeval rb_time_interval(VALUE num)
Creates a "time interval".
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.
void rb_undef_alloc_func(VALUE klass)
Deletes the allocator function of a class.
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().
static bool rb_ractor_shareable_p(VALUE obj)
Queries if multiple Ractors can share the passed object or not.
#define RBIMPL_ATTR_MAYBE_UNUSED()
Wraps (or simulates) [[maybe_unused]]
#define RB_GC_GUARD(v)
Prevents premature destruction of local objects.
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.
#define RARRAY_CONST_PTR
Just another name of rb_array_const_ptr.
#define RUBY_TYPED_DEFAULT_FREE
This is a value you can set to rb_data_type_struct::dfree.
#define DATA_PTR(obj)
Convenient casting macro for backward compatibility.
#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...
#define RTEST
This is an old name of RB_TEST.
This is the struct that holds necessary info for a struct.
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.
static bool RB_TYPE_P(VALUE obj, enum ruby_value_type t)
Queries if the given object is of given type.