12#ifdef THREAD_SYSTEM_DEPENDENT_IMPLEMENTATION
14#include "internal/gc.h"
15#include "internal/sanitizers.h"
17#ifdef HAVE_SYS_RESOURCE_H
18#include <sys/resource.h>
20#ifdef HAVE_THR_STKSEGMENT
23#if defined(HAVE_FCNTL_H)
25#elif defined(HAVE_SYS_FCNTL_H)
28#ifdef HAVE_SYS_PRCTL_H
31#if defined(HAVE_SYS_TIME_H)
38#include <sys/syscall.h>
46# include <AvailabilityMacros.h>
49#if defined(HAVE_SYS_EVENTFD_H) && defined(HAVE_EVENTFD)
50# define USE_EVENTFD (1)
51# include <sys/eventfd.h>
53# define USE_EVENTFD (0)
56#if defined(HAVE_PTHREAD_CONDATTR_SETCLOCK) && \
57 defined(CLOCK_REALTIME) && defined(CLOCK_MONOTONIC) && \
58 defined(HAVE_CLOCK_GETTIME)
59static pthread_condattr_t condattr_mono;
60static pthread_condattr_t *condattr_monotonic = &condattr_mono;
62static const void *
const condattr_monotonic = NULL;
67#ifndef HAVE_SYS_EVENT_H
68#define HAVE_SYS_EVENT_H 0
71#ifndef HAVE_SYS_EPOLL_H
72#define HAVE_SYS_EPOLL_H 0
80 #if defined(__EMSCRIPTEN__) || defined(COROUTINE_PTHREAD_CONTEXT)
83 #define USE_MN_THREADS 0
84 #elif HAVE_SYS_EPOLL_H
85 #include <sys/epoll.h>
86 #define USE_MN_THREADS 1
87 #elif HAVE_SYS_EVENT_H
88 #include <sys/event.h>
89 #define USE_MN_THREADS 1
91 #define USE_MN_THREADS 0
95#ifdef HAVE_SCHED_YIELD
96#define native_thread_yield() (void)sched_yield()
98#define native_thread_yield() ((void)0)
103#define NATIVE_MUTEX_LOCK_DEBUG 0
104#define NATIVE_MUTEX_LOCK_DEBUG_YIELD 0
107mutex_debug(
const char *msg,
void *lock)
109 if (NATIVE_MUTEX_LOCK_DEBUG) {
111 static pthread_mutex_t dbglock = PTHREAD_MUTEX_INITIALIZER;
113 if ((r = pthread_mutex_lock(&dbglock)) != 0) {exit(EXIT_FAILURE);}
114 fprintf(stdout,
"%s: %p\n", msg, lock);
115 if ((r = pthread_mutex_unlock(&dbglock)) != 0) {exit(EXIT_FAILURE);}
123#if NATIVE_MUTEX_LOCK_DEBUG_YIELD
124 native_thread_yield();
126 mutex_debug(
"lock", lock);
127 if ((r = pthread_mutex_lock(lock)) != 0) {
136 mutex_debug(
"unlock", lock);
137 if ((r = pthread_mutex_unlock(lock)) != 0) {
146 mutex_debug(
"trylock", lock);
147 if ((r = pthread_mutex_trylock(lock)) != 0) {
161 int r = pthread_mutex_init(lock, 0);
162 mutex_debug(
"init", lock);
171 int r = pthread_mutex_destroy(lock);
172 mutex_debug(
"destroy", lock);
181 int r = pthread_cond_init(cond, condattr_monotonic);
190 int r = pthread_cond_destroy(cond);
211 r = pthread_cond_signal(cond);
212 }
while (r == EAGAIN);
223 r = pthread_cond_broadcast(cond);
224 }
while (r == EAGAIN);
233 int r = pthread_cond_wait(cond, mutex);
240native_cond_timedwait(rb_nativethread_cond_t *cond, pthread_mutex_t *mutex,
const rb_hrtime_t *abs)
252 rb_hrtime2timespec(&ts, abs);
253 r = pthread_cond_timedwait(cond, mutex, &ts);
254 }
while (r == EINTR);
256 if (r != 0 && r != ETIMEDOUT) {
264native_cond_timeout(rb_nativethread_cond_t *cond,
const rb_hrtime_t rel)
266 if (condattr_monotonic) {
267 return rb_hrtime_add(rb_hrtime_now(), rel);
273 return rb_hrtime_add(rb_timespec2hrtime(&ts), rel);
280 rb_hrtime_t hrmsec = native_cond_timeout(cond, RB_HRTIME_PER_MSEC * msec);
281 native_cond_timedwait(cond, mutex, &hrmsec);
286static rb_internal_thread_event_hook_t *rb_internal_thread_event_hooks = NULL;
308#define RB_INTERNAL_THREAD_HOOK(event, th) \
309 if (UNLIKELY(rb_internal_thread_event_hooks)) { \
310 fprintf(stderr, "[thread=%"PRIxVALUE"] %s in %s (%s:%d)\n", th->self, event_name(event), __func__, __FILE__, __LINE__); \
311 rb_thread_execute_hooks(event, th); \
314#define RB_INTERNAL_THREAD_HOOK(event, th) if (UNLIKELY(rb_internal_thread_event_hooks)) { rb_thread_execute_hooks(event, th); }
317static rb_serial_t current_fork_gen = 1;
319#if defined(SIGVTALRM) && !defined(__EMSCRIPTEN__)
320# define USE_UBF_LIST 1
323static void threadptr_trap_interrupt(
rb_thread_t *);
330static void timer_thread_wakeup(
void);
331static void timer_thread_wakeup_locked(
rb_vm_t *vm);
332static void timer_thread_wakeup_force(
void);
336static void nt_machine_stack_atfork(
void);
344struct rb_thread_context {
356#define thread_sched_dump(s) thread_sched_dump_(__FILE__, __LINE__, s)
362 return th->nt->dedicated > 0;
367thread_sched_dump_(const
char *file,
int line, struct
rb_thread_sched *sched)
369 fprintf(stderr,
"@%s:%d running:%d\n", file, line, sched->running ? (
int)sched->running->serial : -1);
372 ccan_list_for_each(&sched->readyq, th, sched.node.readyq) {
373 i++;
if (i>10) rb_bug(
"too many");
374 fprintf(stderr,
" ready:%d (%sNT:%d)\n", th->serial,
375 th->nt ? (th->nt->dedicated ?
"D" :
"S") :
"x",
376 th->nt ? (int)th->nt->serial : -1);
380#define ractor_sched_dump(s) ractor_sched_dump_(__FILE__, __LINE__, s)
384ractor_sched_dump_(const
char *file,
int line,
rb_vm_t *vm)
388 fprintf(stderr,
"ractor_sched_dump %s:%d\n", file, line);
391 ccan_list_for_each(&vm->ractor.sched.grq, r, threads.sched.grq_node) {
393 if (i>10) rb_bug(
"!!");
394 fprintf(stderr,
" %d ready:%d\n", i, rb_ractor_id(r));
398#define thread_sched_lock(a, b) thread_sched_lock_(a, b, __FILE__, __LINE__)
399#define thread_sched_unlock(a, b) thread_sched_unlock_(a, b, __FILE__, __LINE__)
405 VM_ASSERT(sched->lock_owner == NULL);
407 sched->lock_owner = th;
415 VM_ASSERT(sched->lock_owner == th);
417 sched->lock_owner = NULL;
427 RUBY_DEBUG_LOG2(file, line,
"r:%d th:%u", th ? (
int)rb_ractor_id(th->ractor) : -1, rb_th_serial(th));
429 RUBY_DEBUG_LOG2(file, line,
"th:%u", rb_th_serial(th));
432 thread_sched_set_locked(sched, th);
438 RUBY_DEBUG_LOG2(file, line,
"th:%u", rb_th_serial(th));
440 thread_sched_set_unlocked(sched, th);
452 VM_ASSERT(sched->lock_owner == th);
455 VM_ASSERT(sched->lock_owner != NULL);
460#define ractor_sched_lock(a, b) ractor_sched_lock_(a, b, __FILE__, __LINE__)
461#define ractor_sched_unlock(a, b) ractor_sched_unlock_(a, b, __FILE__, __LINE__)
468 return rb_ractor_id(r);
479 VM_ASSERT(vm->ractor.sched.lock_owner == NULL);
480 VM_ASSERT(vm->ractor.sched.locked ==
false);
482 vm->ractor.sched.lock_owner = cr;
483 vm->ractor.sched.locked =
true;
491 VM_ASSERT(vm->ractor.sched.locked);
492 VM_ASSERT(vm->ractor.sched.lock_owner == cr);
494 vm->ractor.sched.locked =
false;
495 vm->ractor.sched.lock_owner = NULL;
506 RUBY_DEBUG_LOG2(file, line,
"cr:%u prev_owner:%u", rb_ractor_serial(cr), rb_ractor_serial(vm->ractor.sched.lock_owner));
508 RUBY_DEBUG_LOG2(file, line,
"cr:%u", rb_ractor_serial(cr));
511 ractor_sched_set_locked(vm, cr);
517 RUBY_DEBUG_LOG2(file, line,
"cr:%u", rb_ractor_serial(cr));
519 ractor_sched_set_unlocked(vm, cr);
527 VM_ASSERT(vm->ractor.sched.locked);
528 VM_ASSERT(cr == NULL || vm->ractor.sched.lock_owner == cr);
536 ccan_list_for_each(&vm->ractor.sched.running_threads, rth, sched.node.running_threads) {
537 if (rth == th)
return true;
544ractor_sched_running_threads_size(
rb_vm_t *vm)
548 ccan_list_for_each(&vm->ractor.sched.running_threads, th, sched.node.running_threads) {
556ractor_sched_timeslice_threads_size(
rb_vm_t *vm)
560 ccan_list_for_each(&vm->ractor.sched.timeslice_threads, th, sched.node.timeslice_threads) {
571 ccan_list_for_each(&vm->ractor.sched.timeslice_threads, rth, sched.node.timeslice_threads) {
572 if (rth == th)
return true;
577static void ractor_sched_barrier_join_signal_locked(
rb_vm_t *vm);
584#if USE_RUBY_DEBUG_LOG
585 unsigned int prev_running_cnt = vm->ractor.sched.running_cnt;
590 if (del_th && sched->is_running_timeslice) {
591 del_timeslice_th = del_th;
592 sched->is_running_timeslice =
false;
595 del_timeslice_th = NULL;
598 RUBY_DEBUG_LOG(
"+:%u -:%u +ts:%u -ts:%u",
599 rb_th_serial(add_th), rb_th_serial(del_th),
600 rb_th_serial(add_timeslice_th), rb_th_serial(del_timeslice_th));
602 ractor_sched_lock(vm, cr);
606 VM_ASSERT(ractor_sched_running_threads_contain_p(vm, del_th));
607 VM_ASSERT(del_timeslice_th != NULL ||
608 !ractor_sched_timeslice_threads_contain_p(vm, del_th));
610 ccan_list_del_init(&del_th->sched.node.running_threads);
611 vm->ractor.sched.running_cnt--;
613 if (UNLIKELY(vm->ractor.sched.barrier_waiting)) {
614 ractor_sched_barrier_join_signal_locked(vm);
616 sched->is_running =
false;
620 if (vm->ractor.sched.barrier_waiting) {
622 RUBY_DEBUG_LOG(
"barrier_waiting");
623 RUBY_VM_SET_VM_BARRIER_INTERRUPT(add_th->ec);
626 VM_ASSERT(!ractor_sched_running_threads_contain_p(vm, add_th));
627 VM_ASSERT(!ractor_sched_timeslice_threads_contain_p(vm, add_th));
629 ccan_list_add(&vm->ractor.sched.running_threads, &add_th->sched.node.running_threads);
630 vm->ractor.sched.running_cnt++;
631 sched->is_running =
true;
634 if (add_timeslice_th) {
636 int was_empty = ccan_list_empty(&vm->ractor.sched.timeslice_threads);
637 VM_ASSERT(!ractor_sched_timeslice_threads_contain_p(vm, add_timeslice_th));
638 ccan_list_add(&vm->ractor.sched.timeslice_threads, &add_timeslice_th->sched.node.timeslice_threads);
639 sched->is_running_timeslice =
true;
641 timer_thread_wakeup_locked(vm);
645 if (del_timeslice_th) {
646 VM_ASSERT(ractor_sched_timeslice_threads_contain_p(vm, del_timeslice_th));
647 ccan_list_del_init(&del_timeslice_th->sched.node.timeslice_threads);
650 VM_ASSERT(ractor_sched_running_threads_size(vm) == vm->ractor.sched.running_cnt);
651 VM_ASSERT(ractor_sched_timeslice_threads_size(vm) <= vm->ractor.sched.running_cnt);
653 ractor_sched_unlock(vm, cr);
658 RUBY_DEBUG_LOG(
"run:%u->%u", prev_running_cnt, vm->ractor.sched.running_cnt);
664 ASSERT_thread_sched_locked(sched, th);
665 VM_ASSERT(sched->running == th);
668 thread_sched_setup_running_threads(sched, th->ractor, vm, th, NULL, ccan_list_empty(&sched->readyq) ? NULL : th);
674 ASSERT_thread_sched_locked(sched, th);
677 thread_sched_setup_running_threads(sched, th->ractor, vm, NULL, th, NULL);
685 thread_sched_lock(sched, th);
687 thread_sched_add_running_thread(sched, th);
689 thread_sched_unlock(sched, th);
697 thread_sched_lock(sched, th);
699 thread_sched_del_running_thread(sched, th);
701 thread_sched_unlock(sched, th);
711 RUBY_DEBUG_LOG(
"th:%u->th:%u", rb_th_serial(sched->running), rb_th_serial(th));
712 VM_ASSERT(sched->running != th);
714 if (RUBY_DTRACE_RTS_SET_RUNNING_ENABLED()) {
715 RUBY_DTRACE_RTS_SET_RUNNING(sched, sched->running, th);
726 ccan_list_for_each(&sched->readyq, rth, sched.node.readyq) {
728 VM_ASSERT(th->sched.node.is_ready);
732 VM_ASSERT(!th->sched.node.is_ready);
743 ASSERT_thread_sched_locked(sched, NULL);
746 VM_ASSERT(sched->running != NULL);
748 if (ccan_list_empty(&sched->readyq)) {
752 next_th = ccan_list_pop(&sched->readyq,
rb_thread_t, sched.node.readyq);
753 VM_ASSERT(next_th->sched.node.is_ready);
754 next_th->sched.node.is_ready =
false;
756 VM_ASSERT(sched->readyq_cnt > 0);
758 ccan_list_node_init(&next_th->sched.node.readyq);
761 RUBY_DEBUG_LOG(
"next_th:%u readyq_cnt:%d", rb_th_serial(next_th), sched->readyq_cnt);
770 ASSERT_thread_sched_locked(sched, NULL);
771 RUBY_DEBUG_LOG(
"ready_th:%u readyq_cnt:%d", rb_th_serial(ready_th), sched->readyq_cnt);
773 VM_ASSERT(sched->running != NULL);
774 VM_ASSERT(!thread_sched_readyq_contain_p(sched, ready_th));
776 if (sched->is_running) {
777 if (ccan_list_empty(&sched->readyq)) {
779 thread_sched_setup_running_threads(sched, ready_th->ractor, ready_th->vm, NULL, NULL, sched->running);
787 ccan_list_add_tail(&sched->readyq, &ready_th->sched.node.readyq);
788 ready_th->sched.node.is_ready =
true;
797 ASSERT_thread_sched_locked(sched, NULL);
798 VM_ASSERT(sched->running == next_th);
802 if (th_has_dedicated_nt(next_th)) {
803 RUBY_DEBUG_LOG(
"pinning th:%u", next_th->serial);
808 RUBY_DEBUG_LOG(
"th:%u is already running.", next_th->serial);
813 RUBY_DEBUG_LOG(
"th:%u (do nothing)", rb_th_serial(next_th));
816 RUBY_DEBUG_LOG(
"th:%u (enq)", rb_th_serial(next_th));
817 ractor_sched_enq(next_th->vm, next_th->ractor);
822 RUBY_DEBUG_LOG(
"no waiting threads%s",
"");
830 RUBY_DEBUG_LOG(
"th:%u running:%u redyq_cnt:%d", rb_th_serial(th), rb_th_serial(sched->running), sched->readyq_cnt);
832 VM_ASSERT(sched->running != th);
833 VM_ASSERT(!thread_sched_readyq_contain_p(sched, th));
836 if (sched->running == NULL) {
837 thread_sched_set_running(sched, th);
838 if (wakeup) thread_sched_wakeup_running_thread(sched, th, will_switch);
841 thread_sched_enq(sched, th);
853 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
855 thread_sched_lock(sched, th);
857 thread_sched_to_ready_common(sched, th,
true,
false);
859 thread_sched_unlock(sched, th);
866 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
868 ASSERT_thread_sched_locked(sched, th);
869 VM_ASSERT(th == rb_ec_thread_ptr(rb_current_ec_noinline()));
871 if (th != sched->running) {
873 if (th->has_dedicated_nt && th == sched->runnable_hot_th && (sched->running == NULL || sched->running->has_dedicated_nt)) {
874 RUBY_DEBUG_LOG(
"(nt) stealing: hot-th:%u. running:%u", rb_th_serial(th), rb_th_serial(sched->running));
880 ractor_sched_cancel_enq(th->vm, sched);
883 if (sched->running != NULL) {
885 VM_ASSERT(!thread_sched_readyq_contain_p(sched, running));
886 running->sched.node.is_ready =
true;
887 ccan_list_add(&sched->readyq, &running->sched.node.readyq);
892 if (th->sched.node.is_ready) {
893 VM_ASSERT(thread_sched_readyq_contain_p(sched, th));
894 ccan_list_del_init(&th->sched.node.readyq);
895 th->sched.node.is_ready =
false;
898 thread_sched_set_running(sched, th);
899 rb_ractor_thread_switch(th->ractor, th,
false);
901 else if (th == sched->runnable_hot_th) {
905 sched->runnable_hot_th = NULL;
906 sched->runnable_hot_th_waiting = 0;
914 while((next_th = sched->running) != th) {
915 if (th_has_dedicated_nt(th)) {
916 RUBY_DEBUG_LOG(
"(nt) sleep th:%u running:%u", rb_th_serial(th), rb_th_serial(sched->running));
918 thread_sched_set_unlocked(sched, th);
920 RUBY_DEBUG_LOG(
"nt:%d cond:%p", th->nt->serial, &th->nt->cond.readyq);
923 thread_sched_set_locked(sched, th);
925 if (sched->runnable_hot_th != NULL && sched->runnable_hot_th_waiting) {
926 VM_ASSERT(sched->runnable_hot_th != th);
930 thread_sched_unlock(sched, th);
931 thread_sched_lock(sched, th);
934 RUBY_DEBUG_LOG(
"(nt) wakeup %s", sched->running == th ?
"success" :
"failed");
935 if (th == sched->running) {
936 rb_ractor_thread_switch(th->ractor, th,
false);
941 if (can_direct_transfer &&
942 (next_th = sched->running) != NULL &&
946 RUBY_DEBUG_LOG(
"th:%u->%u (direct)", rb_th_serial(th), rb_th_serial(next_th));
948 thread_sched_set_unlocked(sched, th);
950 rb_ractor_set_current_ec(th->ractor, NULL);
951 thread_sched_switch(th, next_th);
953 thread_sched_set_locked(sched, th);
958 native_thread_assign(NULL, th);
960 RUBY_DEBUG_LOG(
"th:%u->%u (ractor scheduling)", rb_th_serial(th), rb_th_serial(next_th));
962 thread_sched_set_unlocked(sched, th);
964 rb_ractor_set_current_ec(th->ractor, NULL);
965 coroutine_transfer0(th->sched.context, nt->nt_context,
false);
967 thread_sched_set_locked(sched, th);
970 VM_ASSERT(rb_current_ec_noinline() == th->ec);
974 VM_ASSERT(th->nt != NULL);
975 VM_ASSERT(rb_current_ec_noinline() == th->ec);
976 VM_ASSERT(th->sched.waiting_reason.flags == thread_sched_waiting_none);
979 thread_sched_add_running_thread(sched, th);
984 sched->runnable_hot_th = NULL;
985 sched->runnable_hot_th_waiting = 0;
995 RUBY_DEBUG_LOG(
"th:%u dedicated:%d", rb_th_serial(th), th_has_dedicated_nt(th));
997 VM_ASSERT(sched->running != th);
998 VM_ASSERT(th_has_dedicated_nt(th));
999 VM_ASSERT(GET_THREAD() == th);
1001 native_thread_dedicated_dec(th->vm, th->ractor, th->nt);
1004 thread_sched_to_ready_common(sched, th,
false,
false);
1006 if (sched->running == th) {
1007 thread_sched_add_running_thread(sched, th);
1011 thread_sched_wait_running_turn(sched, th,
false);
1027 if (sched->runnable_hot_th == th) {
1028 sched->runnable_hot_th_waiting = 1;
1030 thread_sched_lock(sched, th);
1032 thread_sched_to_running_common(sched, th);
1034 thread_sched_unlock(sched, th);
1047 ASSERT_thread_sched_locked(sched, th);
1049 VM_ASSERT(sched->running == th);
1050 VM_ASSERT(sched->running->nt != NULL);
1054 RUBY_DEBUG_LOG(
"next_th:%u", rb_th_serial(next_th));
1055 VM_ASSERT(th != next_th);
1057 thread_sched_set_running(sched, next_th);
1058 VM_ASSERT(next_th == sched->running);
1059 thread_sched_wakeup_running_thread(sched, next_th, will_switch);
1061 if (th != next_th) {
1062 thread_sched_del_running_thread(sched, th);
1070 RUBY_DEBUG_LOG(
"th:%u DNT:%d", rb_th_serial(th), th->nt->dedicated);
1078 thread_sched_wakeup_next_thread(sched, th, !th_has_dedicated_nt(th));
1087 thread_sched_lock(sched, th);
1089 thread_sched_to_dead_common(sched, th);
1091 thread_sched_unlock(sched, th);
1100 RUBY_DEBUG_LOG(
"th:%u DNT:%d", rb_th_serial(th), th->nt->dedicated);
1104 native_thread_dedicated_inc(th->vm, th->ractor, th->nt);
1105 if (!yield_immediately) {
1106 sched->runnable_hot_th = th;
1107 sched->runnable_hot_th_waiting = 0;
1109 thread_sched_wakeup_next_thread(sched, th,
false);
1118 thread_sched_lock(sched, th);
1120 thread_sched_to_waiting_common(sched, th, yield_immediately);
1122 thread_sched_unlock(sched, th);
1130 VM_ASSERT(func != NULL);
1133 if (RUBY_VM_INTERRUPTED(th->ec)) {
1134 RUBY_DEBUG_LOG(
"interrupted:0x%x", th->ec->interrupt_flag);
1140 if (!th->ec->raised_flag && RUBY_VM_INTERRUPTED(th->ec)) {
1145 VM_ASSERT(th->unblock.func == NULL);
1146 th->unblock.func = func;
1147 th->unblock.arg = arg;
1150 *event_serial = prev_serial+1;
1163 th->unblock.func = NULL;
1164 th->unblock.arg = NULL;
1173ubf_waiting(
void *ptr)
1179 th->unblock.func = NULL;
1180 th->unblock.arg = NULL;
1182 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
1184 thread_sched_lock(sched, th);
1186 if (sched->running == th) {
1190 thread_sched_to_ready_common(sched, th,
true,
false);
1193 thread_sched_unlock(sched, th);
1202 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
1204 RB_VM_SAVE_MACHINE_CONTEXT(th);
1209 thread_sched_lock(sched, th);
1212 if (ubf_set(th, ubf_waiting, (
void *)th, NULL)) {
1213 RUBY_DEBUG_LOG(
"th:%u interrupted", rb_th_serial(th));
1216 bool can_direct_transfer = !th_has_dedicated_nt(th);
1218 thread_sched_wakeup_next_thread(sched, th, can_direct_transfer);
1219 thread_sched_wait_running_turn(sched, th, can_direct_transfer);
1222 thread_sched_unlock(sched, th);
1224 ubf_clear(th,
false);
1232 RUBY_DEBUG_LOG(
"th:%d sched->readyq_cnt:%d", (
int)th->serial, sched->readyq_cnt);
1234 thread_sched_lock(sched, th);
1236 if (!ccan_list_empty(&sched->readyq)) {
1238 thread_sched_wakeup_next_thread(sched, th, !th_has_dedicated_nt(th));
1239 bool can_direct_transfer = !th_has_dedicated_nt(th);
1240 thread_sched_to_ready_common(sched, th,
false, can_direct_transfer);
1241 thread_sched_wait_running_turn(sched, th, can_direct_transfer);
1242 th->status = THREAD_RUNNABLE;
1245 VM_ASSERT(sched->readyq_cnt == 0);
1248 thread_sched_unlock(sched, th);
1257 sched->lock_owner = NULL;
1260 ccan_list_head_init(&sched->readyq);
1261 sched->readyq_cnt = 0;
1262 ccan_list_node_init(&sched->grq_node);
1265 if (!atfork) sched->enable_mn_threads =
true;
1272#ifdef RUBY_ASAN_ENABLED
1273 void **fake_stack = to_dead ? NULL : &transfer_from->fake_stack;
1274 __sanitizer_start_switch_fiber(fake_stack, transfer_to->stack_base, transfer_to->stack_size);
1277#if defined(COROUTINE_SANITIZE_THREAD)
1281 __tsan_switch_to_fiber(transfer_to->tsan_fiber, 0);
1285 struct
coroutine_context *returning_from = coroutine_transfer(transfer_from, transfer_to);
1289 VM_ASSERT(!to_dead);
1290#ifdef RUBY_ASAN_ENABLED
1291 __sanitizer_finish_switch_fiber(transfer_from->fake_stack,
1292 (
const void**)&returning_from->stack_base, &returning_from->stack_size);
1299 VM_ASSERT(!nt->dedicated);
1300 VM_ASSERT(next_th->nt == NULL);
1302 RUBY_DEBUG_LOG(
"next_th:%u", rb_th_serial(next_th));
1306 ractor_sched_cancel_enq(next_th->vm, TH_SCHED(next_th));
1308 ruby_thread_set_native(next_th);
1309 native_thread_assign(nt, next_th);
1311 coroutine_transfer0(current_cont, next_th->sched.context, to_dead);
1318 native_thread_assign(NULL, cth);
1319 RUBY_DEBUG_LOG(
"th:%u->%u on nt:%d", rb_th_serial(cth), rb_th_serial(next_th), nt->serial);
1320 thread_sched_switch0(cth->sched.context, next_th, nt, cth->status == THREAD_KILLED);
1323#if VM_CHECK_MODE > 0
1328 ASSERT_ractor_sched_locked(vm, cr);
1333 ccan_list_for_each(&vm->ractor.sched.grq, r, threads.sched.grq_node) {
1336 VM_ASSERT(r != prev_r);
1348rb_thread_sched_wait_winding(
rb_vm_t *vm)
1351 native_thread_yield();
1365 if (sched->grq_node.next != &sched->grq_node) {
1366 ractor_sched_lock(vm, NULL);
1368 if (sched->grq_node.next != &sched->grq_node) {
1369 ccan_list_del_init(&sched->grq_node);
1370 VM_ASSERT(vm->ractor.sched.grq_cnt > 0);
1371 vm->ractor.sched.grq_cnt--;
1374 ractor_sched_unlock(vm, NULL);
1384 VM_ASSERT(sched->running != NULL);
1385 VM_ASSERT(sched->running->nt == NULL);
1387 ractor_sched_lock(vm, cr);
1396 if (sched->grq_node.next != &sched->grq_node) {
1397 rb_bug(
"ractor_sched_enq: already enqueued");
1399 ccan_list_add_tail(&vm->ractor.sched.grq, &sched->grq_node);
1400 vm->ractor.sched.grq_cnt++;
1401 VM_ASSERT(grq_size(vm, cr) == vm->ractor.sched.grq_cnt);
1403 RUBY_DEBUG_LOG(
"r:%u th:%u grq_cnt:%u", rb_ractor_id(r), rb_th_serial(sched->running), vm->ractor.sched.grq_cnt);
1412 if (vm->ractor.sched.snt_cnt == 0) {
1413 timer_thread_wakeup_locked(vm);
1418 ractor_sched_unlock(vm, cr);
1424#define MINIMUM_SNT 0
1431#ifndef SNT_IDLE_RETIRE
1432#define SNT_IDLE_RETIRE 3
1438#define SNT_KEEP_MINIMUM (MINIMUM_SNT > 1 ? MINIMUM_SNT : 1)
1444 int idle_streak = 0;
1446 ractor_sched_lock(vm, cr);
1448 RUBY_DEBUG_LOG(
"empty? %d", ccan_list_empty(&vm->ractor.sched.grq));
1451 VM_ASSERT(rb_current_execution_context(
false) == NULL);
1452 VM_ASSERT(grq_size(vm, cr) == vm->ractor.sched.grq_cnt);
1454 while ((r = ccan_list_pop(&vm->ractor.sched.grq,
rb_ractor_t, threads.sched.grq_node)) == NULL) {
1455 RUBY_DEBUG_LOG(
"wait grq_cnt:%d", (
int)vm->ractor.sched.grq_cnt);
1457 if (SNT_IDLE_RETIRE >= 0 && ++idle_streak > SNT_IDLE_RETIRE &&
1458 (
int)vm->ractor.sched.snt_cnt > SNT_KEEP_MINIMUM) {
1459 vm->ractor.sched.snt_cnt--;
1460 RUBY_DEBUG_LOG(
"retire, snt_cnt:%d", (
int)vm->ractor.sched.snt_cnt);
1464 ractor_sched_set_unlocked(vm, cr);
1466 ractor_sched_set_locked(vm, cr);
1468 RUBY_DEBUG_LOG(
"wakeup grq_cnt:%d", (
int)vm->ractor.sched.grq_cnt);
1471 VM_ASSERT(rb_current_execution_context(
false) == NULL);
1474 ccan_list_node_init(&r->threads.sched.grq_node);
1475 VM_ASSERT(vm->ractor.sched.grq_cnt > 0);
1476 vm->ractor.sched.grq_cnt--;
1477 RUBY_DEBUG_LOG(
"r:%d grq_cnt:%u", (
int)rb_ractor_id(r), vm->ractor.sched.grq_cnt);
1481 VM_ASSERT(idle_streak > SNT_IDLE_RETIRE);
1484 ractor_sched_unlock(vm, cr);
1499 RUBY_DEBUG_LOG(
"start%s",
"");
1505 if (ubf_set(th, ubf, ubf_arg, &waiter->event_serial)) {
1510 thread_sched_lock(sched, th);
1511 rb_ractor_unlock_self(cr);
1514 bool can_direct_transfer = !th_has_dedicated_nt(th);
1515 RB_VM_SAVE_MACHINE_CONTEXT(th);
1516 th->status = THREAD_STOPPED_FOREVER;
1518 thread_sched_wakeup_next_thread(sched, th, can_direct_transfer);
1520 thread_sched_wait_running_turn(sched, th, can_direct_transfer);
1521 th->status = THREAD_RUNNABLE;
1523 thread_sched_unlock(sched, th);
1524 rb_ractor_lock_self(cr);
1526 ubf_clear(th,
true);
1528 RUBY_DEBUG_LOG(
"end%s",
"");
1537 RUBY_DEBUG_LOG(
"r:%u th:%d", (
unsigned int)rb_ractor_id(r), r_th->serial);
1539 thread_sched_lock(sched, r_th);
1541 if (r_th->status == THREAD_STOPPED_FOREVER) {
1543 thread_sched_to_ready_common(sched, r_th,
true,
false);
1546 thread_sched_unlock(sched, r_th);
1550ractor_sched_barrier_completed_p(
rb_vm_t *vm)
1552 RUBY_DEBUG_LOG(
"run:%u wait:%u", vm->ractor.sched.running_cnt, vm->ractor.sched.barrier_waiting_cnt);
1553 VM_ASSERT(vm->ractor.sched.running_cnt - 1 >= vm->ractor.sched.barrier_waiting_cnt);
1555 return (vm->ractor.sched.running_cnt - vm->ractor.sched.barrier_waiting_cnt) == 1;
1561 VM_ASSERT(cr == GET_RACTOR());
1562 VM_ASSERT(vm->ractor.sync.lock_owner == cr);
1563 VM_ASSERT(!vm->ractor.sched.barrier_waiting);
1564 VM_ASSERT(vm->ractor.sched.barrier_waiting_cnt == 0);
1565 VM_ASSERT(vm->ractor.sched.barrier_ractor == NULL);
1566 VM_ASSERT(vm->ractor.sched.barrier_lock_rec == 0);
1568 RUBY_DEBUG_LOG(
"start serial:%u", vm->ractor.sched.barrier_serial);
1570 unsigned int lock_rec;
1572 ractor_sched_lock(vm, cr);
1574 vm->ractor.sched.barrier_waiting =
true;
1575 vm->ractor.sched.barrier_ractor = cr;
1576 vm->ractor.sched.barrier_lock_rec = vm->ractor.sync.lock_rec;
1579 lock_rec = vm->ractor.sync.lock_rec;
1580 vm->ractor.sync.lock_rec = 0;
1581 vm->ractor.sync.lock_owner = NULL;
1586 ccan_list_for_each(&vm->ractor.sched.running_threads, ith, sched.node.running_threads) {
1587 if (ith->ractor != cr) {
1588 RUBY_DEBUG_LOG(
"barrier request to th:%u", rb_th_serial(ith));
1589 RUBY_VM_SET_VM_BARRIER_INTERRUPT(ith->ec);
1594 while (!ractor_sched_barrier_completed_p(vm)) {
1595 ractor_sched_set_unlocked(vm, cr);
1597 ractor_sched_set_locked(vm, cr);
1600 RUBY_DEBUG_LOG(
"completed seirial:%u", vm->ractor.sched.barrier_serial);
1603 vm->ractor.sched.barrier_serial++;
1604 vm->ractor.sched.barrier_waiting_cnt = 0;
1609 vm->ractor.sync.lock_rec = lock_rec;
1610 vm->ractor.sync.lock_owner = cr;
1621 RUBY_DEBUG_LOG(
"serial:%u", (
unsigned int)vm->ractor.sched.barrier_serial - 1);
1622 VM_ASSERT(vm->ractor.sched.barrier_waiting);
1623 VM_ASSERT(vm->ractor.sched.barrier_ractor);
1624 VM_ASSERT(vm->ractor.sched.barrier_lock_rec > 0);
1626 vm->ractor.sched.barrier_waiting =
false;
1627 vm->ractor.sched.barrier_ractor = NULL;
1628 vm->ractor.sched.barrier_lock_rec = 0;
1629 ractor_sched_unlock(vm, cr);
1633ractor_sched_barrier_join_signal_locked(
rb_vm_t *vm)
1635 if (ractor_sched_barrier_completed_p(vm)) {
1643 VM_ASSERT(vm->ractor.sched.barrier_waiting);
1645 unsigned int barrier_serial = vm->ractor.sched.barrier_serial;
1647 while (vm->ractor.sched.barrier_serial == barrier_serial) {
1648 RUBY_DEBUG_LOG(
"sleep serial:%u", barrier_serial);
1649 RB_VM_SAVE_MACHINE_CONTEXT(th);
1652 ractor_sched_set_unlocked(vm, cr);
1654 ractor_sched_set_locked(vm, cr);
1656 RUBY_DEBUG_LOG(
"wakeup serial:%u", barrier_serial);
1663 VM_ASSERT(cr->threads.sched.running != NULL);
1664 VM_ASSERT(cr == GET_RACTOR());
1665 VM_ASSERT(vm->ractor.sync.lock_owner == NULL);
1666 VM_ASSERT(vm->ractor.sched.barrier_waiting);
1668#if USE_RUBY_DEBUG_LOG || VM_CHECK_MODE > 0
1669 unsigned int barrier_serial = vm->ractor.sched.barrier_serial;
1672 RUBY_DEBUG_LOG(
"join");
1676 VM_ASSERT(vm->ractor.sched.barrier_waiting);
1677 VM_ASSERT(vm->ractor.sched.barrier_serial == barrier_serial);
1679 ractor_sched_lock(vm, cr);
1684 VM_ASSERT(ractor_sched_running_threads_contain_p(vm, GET_THREAD()));
1685 vm->ractor.sched.barrier_waiting_cnt++;
1686 RUBY_DEBUG_LOG(
"waiting_cnt:%u serial:%u", vm->ractor.sched.barrier_waiting_cnt, barrier_serial);
1688 ractor_sched_barrier_join_signal_locked(vm);
1689 ractor_sched_barrier_join_wait_locked(vm, cr->threads.sched.running);
1691 ractor_sched_unlock(vm, cr);
1701static void clear_thread_cache_altstack(
void);
1714 clear_thread_cache_altstack();
1718#ifdef RB_THREAD_T_HAS_NATIVE_ID
1720get_native_thread_id(
void)
1723 return (
int)syscall(SYS_gettid);
1724#elif defined(__FreeBSD__)
1725 return pthread_getthreadid_np();
1730#if defined(HAVE_WORKING_FORK)
1731static void rb_internal_thread_event_hooks_rw_lock_atfork(
void);
1737 rb_thread_sched_init(sched,
true);
1741 if (th_has_dedicated_nt(th)) {
1742 vm->ractor.sched.snt_cnt = 0;
1743 vm->ractor.sched.dnt_cnt = 1;
1746 vm->ractor.sched.snt_cnt = 1;
1747 vm->ractor.sched.dnt_cnt = 0;
1749 vm->ractor.sched.running_cnt = 0;
1752#if VM_CHECK_MODE > 0
1753 vm->ractor.sched.lock_owner = NULL;
1754 vm->ractor.sched.locked =
false;
1762 ccan_list_head_init(&vm->ractor.sched.grq);
1763 vm->ractor.sched.grq_cnt = 0;
1766 vm->ractor.sched.barrier_waiting =
false;
1767 vm->ractor.sched.barrier_waiting_cnt = 0;
1768 vm->ractor.sched.barrier_ractor = NULL;
1769 vm->ractor.sched.barrier_lock_rec = 0;
1773 vm->ractor.sched.winding_cnt = 0;
1774 ccan_list_head_init(&vm->ractor.sched.timeslice_threads);
1775 ccan_list_head_init(&vm->ractor.sched.running_threads);
1778 nt_machine_stack_atfork();
1780 rb_internal_thread_event_hooks_rw_lock_atfork();
1782 VM_ASSERT(sched->is_running);
1783 sched->is_running_timeslice =
false;
1785 if (sched->running != th) {
1786 thread_sched_to_running(sched, th);
1789 thread_sched_setup_running_threads(sched, th->ractor, vm, th, NULL, NULL);
1792#ifdef RB_THREAD_T_HAS_NATIVE_ID
1794 th->nt->tid = get_native_thread_id();
1801#ifdef RB_THREAD_LOCAL_SPECIFIER
1802static RB_THREAD_LOCAL_SPECIFIER
rb_thread_t *ruby_native_thread;
1804static pthread_key_t ruby_native_thread_key;
1816ruby_thread_from_native(
void)
1818#ifdef RB_THREAD_LOCAL_SPECIFIER
1819 return ruby_native_thread;
1821 return pthread_getspecific(ruby_native_thread_key);
1830 ccan_list_node_init(&th->sched.node.ubf);
1837 rb_ractor_set_current_ec(th->ractor, th->ec);
1839#ifdef RB_THREAD_LOCAL_SPECIFIER
1840 ruby_native_thread = th;
1843 return pthread_setspecific(ruby_native_thread_key, th) == 0;
1851static size_t RB_THREAD_PAGE_SIZE;
1857 RB_THREAD_PAGE_SIZE = sysconf(_SC_PAGESIZE);
1859#if defined(HAVE_PTHREAD_CONDATTR_SETCLOCK)
1860 if (condattr_monotonic) {
1861 int r = pthread_condattr_init(condattr_monotonic);
1863 r = pthread_condattr_setclock(condattr_monotonic, CLOCK_MONOTONIC);
1865 if (r) condattr_monotonic = NULL;
1869#ifndef RB_THREAD_LOCAL_SPECIFIER
1870 if (pthread_key_create(&ruby_native_thread_key, 0) == EAGAIN) {
1871 rb_bug(
"pthread_key_create failed (ruby_native_thread_key)");
1873 if (pthread_key_create(&ruby_current_ec_key, 0) == EAGAIN) {
1874 rb_bug(
"pthread_key_create failed (ruby_current_ec_key)");
1877 ruby_posix_signal(SIGVTALRM, null_func);
1886 ccan_list_head_init(&vm->ractor.sched.grq);
1887 ccan_list_head_init(&vm->ractor.sched.timeslice_threads);
1888 ccan_list_head_init(&vm->ractor.sched.running_threads);
1891 main_th->nt->thread_id = pthread_self();
1892 main_th->nt->serial = 1;
1893#ifdef RUBY_NT_SERIAL
1896 ruby_thread_set_native(main_th);
1897 native_thread_setup(main_th->nt);
1898 native_thread_setup_on_thread(main_th->nt);
1900 TH_SCHED(main_th)->running = main_th;
1901 main_th->has_dedicated_nt = 1;
1903 thread_sched_setup_running_threads(TH_SCHED(main_th), main_th->ractor, vm, main_th, NULL, NULL);
1906 main_th->nt->dedicated = 1;
1907 main_th->nt->vm = vm;
1910 vm->ractor.sched.dnt_cnt = 1;
1913extern int ruby_mn_threads_enabled;
1916ruby_mn_threads_params(
void)
1921 const char *mn_threads_cstr = getenv(
"RUBY_MN_THREADS");
1922 bool enable_mn_threads =
false;
1924 if (USE_MN_THREADS && mn_threads_cstr && (enable_mn_threads = atoi(mn_threads_cstr) > 0)) {
1926 ruby_mn_threads_enabled = 1;
1928 main_ractor->threads.sched.enable_mn_threads = enable_mn_threads;
1930 const char *max_cpu_cstr = getenv(
"RUBY_MAX_CPU");
1931#if defined(HAVE_SYSCONF) && defined(_SC_NPROCESSORS_ONLN)
1932 long nprocessors = sysconf(_SC_NPROCESSORS_ONLN);
1933 const int default_max_cpu = (nprocessors > 0) ? (
int)nprocessors : 8;
1935 const int default_max_cpu = 8;
1937 int max_cpu = default_max_cpu;
1939 if (USE_MN_THREADS && max_cpu_cstr) {
1940 int given_max_cpu = atoi(max_cpu_cstr);
1941 if (given_max_cpu > 0) {
1942 max_cpu = given_max_cpu;
1946 vm->ractor.sched.max_cpu = max_cpu;
1952 RUBY_DEBUG_LOG(
"nt:%d %d->%d", nt->serial, nt->dedicated, nt->dedicated + 1);
1954 if (nt->dedicated == 0) {
1955 ractor_sched_lock(vm, cr);
1957 vm->ractor.sched.snt_cnt--;
1958 vm->ractor.sched.dnt_cnt++;
1962 if (vm->ractor.sched.snt_cnt == 0 && vm->ractor.sched.grq_cnt > 0) {
1963 timer_thread_wakeup_locked(vm);
1966 ractor_sched_unlock(vm, cr);
1975 RUBY_DEBUG_LOG(
"nt:%d %d->%d", nt->serial, nt->dedicated, nt->dedicated - 1);
1976 VM_ASSERT(nt->dedicated > 0);
1979 if (nt->dedicated == 0) {
1980 ractor_sched_lock(vm, cr);
1985 if (vm->ractor.sched.snt_cnt < vm->ractor.sched.max_cpu ||
1986 (
int)vm->ractor.sched.snt_cnt <= MINIMUM_SNT) {
1987 vm->ractor.sched.snt_cnt++;
1990 nt->retiring =
true;
1992 vm->ractor.sched.dnt_cnt--;
1994 ractor_sched_unlock(vm, cr);
2001#if USE_RUBY_DEBUG_LOG
2004 RUBY_DEBUG_LOG(
"th:%d nt:%d->%d", (
int)th->serial, (
int)th->nt->serial, (
int)nt->serial);
2007 RUBY_DEBUG_LOG(
"th:%d nt:NULL->%d", (
int)th->serial, (
int)nt->serial);
2012 RUBY_DEBUG_LOG(
"th:%d nt:%d->NULL", (
int)th->serial, (
int)th->nt->serial);
2015 RUBY_DEBUG_LOG(
"th:%d nt:NULL->NULL", (
int)th->serial);
2038 RB_ALTSTACK_FREE(nt->altstack);
2039 SIZED_FREE(nt->nt_context);
2050 if (&nt->cond.readyq != &nt->cond.intr) {
2054 native_thread_destroy_atfork(nt);
2064#ifdef USE_SIGALTSTACK
2065 stack_t disable = {0};
2066 disable.ss_flags = SS_DISABLE;
2067 sigaltstack(&disable, NULL);
2069 native_thread_destroy(nt);
2072#if defined HAVE_PTHREAD_GETATTR_NP || defined HAVE_PTHREAD_ATTR_GET_NP
2073#define STACKADDR_AVAILABLE 1
2074#elif defined HAVE_PTHREAD_GET_STACKADDR_NP && defined HAVE_PTHREAD_GET_STACKSIZE_NP
2075#define STACKADDR_AVAILABLE 1
2076#undef MAINSTACKADDR_AVAILABLE
2077#define MAINSTACKADDR_AVAILABLE 1
2078void *pthread_get_stackaddr_np(pthread_t);
2079size_t pthread_get_stacksize_np(pthread_t);
2080#elif defined HAVE_THR_STKSEGMENT || defined HAVE_PTHREAD_STACKSEG_NP
2081#define STACKADDR_AVAILABLE 1
2082#elif defined HAVE_PTHREAD_GETTHRDS_NP
2083#define STACKADDR_AVAILABLE 1
2084#elif defined __HAIKU__
2085#define STACKADDR_AVAILABLE 1
2088#ifndef MAINSTACKADDR_AVAILABLE
2089# ifdef STACKADDR_AVAILABLE
2090# define MAINSTACKADDR_AVAILABLE 1
2092# define MAINSTACKADDR_AVAILABLE 0
2095#if MAINSTACKADDR_AVAILABLE && !defined(get_main_stack)
2096# define get_main_stack(addr, size) get_stack(addr, size)
2099#ifdef STACKADDR_AVAILABLE
2104get_stack(
void **addr,
size_t *size)
2106#define CHECK_ERR(expr) \
2107 {int err = (expr); if (err) return err;}
2108#ifdef HAVE_PTHREAD_GETATTR_NP
2109 pthread_attr_t attr;
2111 STACK_GROW_DIR_DETECTION;
2112 CHECK_ERR(pthread_getattr_np(pthread_self(), &attr));
2113# ifdef HAVE_PTHREAD_ATTR_GETSTACK
2114 CHECK_ERR(pthread_attr_getstack(&attr, addr, size));
2115 STACK_DIR_UPPER((
void)0, (
void)(*addr = (
char *)*addr + *size));
2117 CHECK_ERR(pthread_attr_getstackaddr(&attr, addr));
2118 CHECK_ERR(pthread_attr_getstacksize(&attr, size));
2120# ifdef HAVE_PTHREAD_ATTR_GETGUARDSIZE
2121 CHECK_ERR(pthread_attr_getguardsize(&attr, &guard));
2123 guard = RB_THREAD_PAGE_SIZE;
2126 pthread_attr_destroy(&attr);
2127#elif defined HAVE_PTHREAD_ATTR_GET_NP
2128 pthread_attr_t attr;
2129 CHECK_ERR(pthread_attr_init(&attr));
2130 CHECK_ERR(pthread_attr_get_np(pthread_self(), &attr));
2131# ifdef HAVE_PTHREAD_ATTR_GETSTACK
2132 CHECK_ERR(pthread_attr_getstack(&attr, addr, size));
2134 CHECK_ERR(pthread_attr_getstackaddr(&attr, addr));
2135 CHECK_ERR(pthread_attr_getstacksize(&attr, size));
2137 STACK_DIR_UPPER((
void)0, (
void)(*addr = (
char *)*addr + *size));
2138 pthread_attr_destroy(&attr);
2139#elif (defined HAVE_PTHREAD_GET_STACKADDR_NP && defined HAVE_PTHREAD_GET_STACKSIZE_NP)
2140 pthread_t th = pthread_self();
2141 *addr = pthread_get_stackaddr_np(th);
2142 *size = pthread_get_stacksize_np(th);
2143#elif defined HAVE_THR_STKSEGMENT || defined HAVE_PTHREAD_STACKSEG_NP
2145# if defined HAVE_THR_STKSEGMENT
2146 CHECK_ERR(thr_stksegment(&stk));
2148 CHECK_ERR(pthread_stackseg_np(pthread_self(), &stk));
2151 *size = stk.ss_size;
2152#elif defined HAVE_PTHREAD_GETTHRDS_NP
2153 pthread_t th = pthread_self();
2154 struct __pthrdsinfo thinfo;
2156 int regsiz=
sizeof(reg);
2157 CHECK_ERR(pthread_getthrds_np(&th, PTHRDSINFO_QUERY_ALL,
2158 &thinfo,
sizeof(thinfo),
2160 *addr = thinfo.__pi_stackaddr;
2164 *size = thinfo.__pi_stackend - thinfo.__pi_stackaddr;
2165 STACK_DIR_UPPER((
void)0, (
void)(*addr = (
char *)*addr + *size));
2166#elif defined __HAIKU__
2168 STACK_GROW_DIR_DETECTION;
2169 CHECK_ERR(get_thread_info(find_thread(NULL), &info));
2170 *addr = info.stack_base;
2171 *size = (uintptr_t)info.stack_end - (uintptr_t)info.stack_base;
2172 STACK_DIR_UPPER((
void)0, (
void)(*addr = (
char *)*addr + *size));
2174#error STACKADDR_AVAILABLE is defined but not implemented.
2182 rb_nativethread_id_t id;
2183 size_t stack_maxsize;
2185} native_main_thread;
2187#ifdef STACK_END_ADDRESS
2188extern void *STACK_END_ADDRESS;
2192native_thread_init_main_thread_stack(
void *addr)
2194 native_main_thread.id = pthread_self();
2195#ifdef RUBY_ASAN_ENABLED
2196 addr = asan_get_real_stack_addr((
void *)addr);
2199#if MAINSTACKADDR_AVAILABLE
2200 if (native_main_thread.stack_maxsize)
return;
2204 if (get_main_stack(&stackaddr, &size) == 0) {
2205 native_main_thread.stack_maxsize = size;
2206 native_main_thread.stack_start = stackaddr;
2211#ifdef STACK_END_ADDRESS
2212 native_main_thread.stack_start = STACK_END_ADDRESS;
2214 if (!native_main_thread.stack_start ||
2215 STACK_UPPER((
VALUE *)(
void *)&addr,
2216 native_main_thread.stack_start > (
VALUE *)addr,
2217 native_main_thread.stack_start < (
VALUE *)addr)) {
2218 native_main_thread.stack_start = (
VALUE *)addr;
2222#if defined(HAVE_GETRLIMIT)
2223#if defined(PTHREAD_STACK_DEFAULT)
2224 size_t size = PTHREAD_STACK_DEFAULT;
2226 size_t size = RUBY_VM_THREAD_VM_STACK_SIZE;
2230 STACK_GROW_DIR_DETECTION;
2231 if (getrlimit(RLIMIT_STACK, &rlim) == 0) {
2232 size = (size_t)rlim.rlim_cur;
2234 addr = native_main_thread.stack_start;
2235 if (IS_STACK_DIR_UPPER()) {
2236 space = ((size_t)((
char *)addr + size) / RB_THREAD_PAGE_SIZE) * RB_THREAD_PAGE_SIZE - (size_t)addr;
2239 space = (size_t)addr - ((
size_t)((
char *)addr - size) / RB_THREAD_PAGE_SIZE + 1) * RB_THREAD_PAGE_SIZE;
2241 native_main_thread.stack_maxsize = space;
2245#if MAINSTACKADDR_AVAILABLE
2252 STACK_GROW_DIR_DETECTION;
2254 if (IS_STACK_DIR_UPPER()) {
2255 start = native_main_thread.stack_start;
2256 end = (
char *)native_main_thread.stack_start + native_main_thread.stack_maxsize;
2259 start = (
char *)native_main_thread.stack_start - native_main_thread.stack_maxsize;
2260 end = native_main_thread.stack_start;
2263 if ((
void *)addr < start || (
void *)addr > end) {
2265 native_main_thread.stack_start = (
VALUE *)addr;
2266 native_main_thread.stack_maxsize = 0;
2271#define CHECK_ERR(expr) \
2272 {int err = (expr); if (err) {rb_bug_errno(#expr, err);}}
2275native_thread_init_stack(
rb_thread_t *th,
void *local_in_parent_frame)
2277 rb_nativethread_id_t curr = pthread_self();
2278#ifdef RUBY_ASAN_ENABLED
2279 local_in_parent_frame = asan_get_real_stack_addr(local_in_parent_frame);
2280 th->ec->machine.asan_fake_stack_handle = asan_get_thread_fake_stack_handle();
2283 if (!native_main_thread.id) {
2286 native_thread_init_main_thread_stack(local_in_parent_frame);
2289 if (pthread_equal(curr, native_main_thread.id)) {
2290 th->ec->machine.stack_start = native_main_thread.stack_start;
2291 th->ec->machine.stack_maxsize = native_main_thread.stack_maxsize;
2294#ifdef STACKADDR_AVAILABLE
2295 if (th_has_dedicated_nt(th)) {
2299 if (get_stack(&start, &size) == 0) {
2300 uintptr_t diff = (uintptr_t)start - (uintptr_t)local_in_parent_frame;
2301 th->ec->machine.stack_start = local_in_parent_frame;
2302 th->ec->machine.stack_maxsize = size - diff;
2306 rb_raise(
rb_eNotImpError,
"ruby engine can initialize only in the main thread");
2325 pthread_attr_t attr;
2327 const size_t stack_size = nt->vm->default_params.thread_machine_stack_size;
2329#ifdef USE_SIGALTSTACK
2330 nt->altstack = rb_allocate_sigaltstack();
2333 CHECK_ERR(pthread_attr_init(&attr));
2335# ifdef PTHREAD_STACK_MIN
2336 RUBY_DEBUG_LOG(
"stack size: %lu", (
unsigned long)stack_size);
2337 CHECK_ERR(pthread_attr_setstacksize(&attr, stack_size));
2340# ifdef HAVE_PTHREAD_ATTR_SETINHERITSCHED
2341 CHECK_ERR(pthread_attr_setinheritsched(&attr, PTHREAD_INHERIT_SCHED));
2343 CHECK_ERR(pthread_attr_setdetachstate(&attr, PTHREAD_CREATE_DETACHED));
2345 err = pthread_create(&nt->thread_id, &attr, nt_start, nt);
2347 RUBY_DEBUG_LOG(
"nt:%d err:%d", (
int)nt->serial, err);
2349 CHECK_ERR(pthread_attr_destroy(&attr));
2360 if (&nt->cond.readyq != &nt->cond.intr) {
2369#ifdef RB_THREAD_T_HAS_NATIVE_ID
2370 nt->tid = get_native_thread_id();
2374 RB_ALTSTACK_INIT(nt->altstack, nt->altstack);
2378native_thread_alloc(
void)
2381 native_thread_setup(nt);
2387#if USE_RUBY_DEBUG_LOG
2397 th->nt = native_thread_alloc();
2398 th->nt->vm = th->vm;
2399 th->nt->running_thread = th;
2400 th->nt->dedicated = 1;
2403 size_t vm_stack_word_size = th->vm->default_params.thread_vm_stack_size /
sizeof(
VALUE);
2404 void *vm_stack = ruby_xmalloc(vm_stack_word_size *
sizeof(
VALUE));
2405 th->sched.malloc_stack =
true;
2406 rb_ec_initialize_vm_stack(th->ec, vm_stack, vm_stack_word_size);
2407 th->sched.context_stack = vm_stack;
2408 th->sched.context_stack_size = vm_stack_word_size;
2410 int err = native_thread_create0(th->nt);
2413 thread_sched_to_ready(TH_SCHED(th), th);
2427 VALUE stack_start = 0;
2428 VALUE *stack_start_addr = asan_get_real_stack_addr(&stack_start);
2430 native_thread_init_stack(th, stack_start_addr);
2431 thread_start_func_2(th, th->ec->machine.stack_start);
2440 native_thread_setup_on_thread(nt);
2443#ifdef RB_THREAD_T_HAS_NATIVE_ID
2444 nt->tid = get_native_thread_id();
2447#if USE_RUBY_DEBUG_LOG && defined(RUBY_NT_SERIAL)
2448 ruby_nt_serial = nt->serial;
2451 RUBY_DEBUG_LOG(
"nt:%u", nt->serial);
2453 if (!nt->dedicated) {
2454 coroutine_initialize_main(nt->nt_context);
2457 bool retired =
false;
2460 if (nt->dedicated) {
2465 RUBY_DEBUG_LOG(
"on dedicated th:%u", rb_th_serial(th));
2466 ruby_thread_set_native(th);
2468 thread_sched_lock(sched, th);
2470 if (sched->running == th) {
2471 thread_sched_add_running_thread(sched, th);
2473 thread_sched_wait_running_turn(sched, th,
false);
2475 thread_sched_unlock(sched, th);
2478 call_thread_start_func_2(th);
2482 RUBY_DEBUG_LOG(
"check next");
2495 thread_sched_lock(sched, NULL);
2499 if (next_th && next_th->nt == NULL) {
2500 RUBY_DEBUG_LOG(
"nt:%d next_th:%d", (
int)nt->serial, (
int)next_th->serial);
2502 thread_sched_switch0(nt->nt_context, next_th, nt,
false);
2509 if (thread_sched_reclaim(dead_co)) {
2515 thread_sched_switch0(nt->nt_context, next_th, nt,
false);
2519 RUBY_DEBUG_LOG(
"no schedulable threads -- next_th:%p", next_th);
2523 thread_sched_unlock(sched, NULL);
2532 if (nt->dedicated) {
2541 RUBY_DEBUG_LOG(
"retired nt:%u", nt->serial);
2542 native_thread_destroy_self(nt);
2548static int native_thread_create_shared(
rb_thread_t *th);
2551static void nt_free_stack(
void *mstack);
2562 struct rb_thread_context *tctx = (
struct rb_thread_context *)dead_co;
2564 if (tctx != NULL && tctx->dead) {
2565 nt_free_stack(tctx->stack);
2581 if (th->sched.malloc_stack) {
2583 SIZED_FREE_N((
VALUE *)th->sched.context_stack, th->sched.context_stack_size);
2584 native_thread_destroy(th->nt);
2586 else if (th->sched.context != NULL) {
2590 struct rb_thread_context *tctx = (
struct rb_thread_context *)th->sched.context;
2591 nt_free_stack(tctx->stack);
2593 th->sched.context = NULL;
2597 SIZED_FREE_N((
VALUE *)th->sched.context_stack, th->sched.context_stack_size);
2598 native_thread_destroy(th->nt);
2608 VM_ASSERT(th->nt == 0);
2609 RUBY_DEBUG_LOG(
"th:%d has_dnt:%d", th->serial, th->has_dedicated_nt);
2612 if (!th->ractor->threads.sched.enable_mn_threads) {
2613 th->has_dedicated_nt = 1;
2616 if (th->has_dedicated_nt) {
2617 return native_thread_create_dedicated(th);
2620 return native_thread_create_shared(th);
2624#if USE_NATIVE_THREAD_PRIORITY
2629#if defined(_POSIX_PRIORITY_SCHEDULING) && (_POSIX_PRIORITY_SCHEDULING > 0)
2630 struct sched_param sp;
2632 int priority = 0 - th->priority;
2634 pthread_getschedparam(th->nt->thread_id, &policy, &sp);
2635 max = sched_get_priority_max(policy);
2636 min = sched_get_priority_min(policy);
2638 if (min > priority) {
2641 else if (max < priority) {
2645 sp.sched_priority = priority;
2646 pthread_setschedparam(th->nt->thread_id, policy, &sp);
2657 return rb_fd_select(n, readfds, writefds, exceptfds, timeout);
2661ubf_pthread_cond_signal(
void *ptr)
2664 RUBY_DEBUG_LOG(
"th:%u on nt:%d", rb_th_serial(th), (
int)th->nt->serial);
2669native_cond_sleep(
rb_thread_t *th, rb_hrtime_t *rel)
2671 rb_nativethread_lock_t *lock = &th->interrupt_lock;
2672 rb_nativethread_cond_t *cond = &th->nt->cond.intr;
2682 const rb_hrtime_t max = (rb_hrtime_t)100000000 * RB_HRTIME_PER_SEC;
2684 THREAD_BLOCKING_BEGIN(th);
2687 th->unblock.func = ubf_pthread_cond_signal;
2688 th->unblock.arg = th;
2690 if (RUBY_VM_INTERRUPTED(th->ec)) {
2692 RUBY_DEBUG_LOG(
"interrupted before sleep th:%u", rb_th_serial(th));
2705 end = native_cond_timeout(cond, *rel);
2706 native_cond_timedwait(cond, lock, &end);
2709 th->unblock.func = 0;
2713 THREAD_BLOCKING_END(th);
2715 RUBY_DEBUG_LOG(
"done th:%u", rb_th_serial(th));
2719static CCAN_LIST_HEAD(ubf_list_head);
2720static rb_nativethread_lock_t ubf_list_lock = RB_NATIVETHREAD_LOCK_INIT;
2723ubf_list_atfork(
void)
2725 ccan_list_head_init(&ubf_list_head);
2734 ccan_list_for_each(&ubf_list_head, list_th, sched.node.ubf) {
2735 if (list_th == th)
return true;
2744 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
2745 struct ccan_list_node *node = &th->sched.node.ubf;
2747 VM_ASSERT(th->unblock.func != NULL);
2752 if (ccan_list_empty((
struct ccan_list_head*)node)) {
2753 VM_ASSERT(!ubf_list_contain_p(th));
2754 ccan_list_add(&ubf_list_head, node);
2759 timer_thread_wakeup();
2766 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
2767 struct ccan_list_node *node = &th->sched.node.ubf;
2770 VM_ASSERT(th->unblock.func == NULL);
2772 if (!ccan_list_empty((
struct ccan_list_head*)node)) {
2775 VM_ASSERT(ubf_list_contain_p(th));
2776 ccan_list_del_init(node);
2789 RUBY_DEBUG_LOG(
"th:%u thread_id:%p", rb_th_serial(th), (
void *)th->nt->thread_id);
2791 pthread_kill(th->nt->thread_id, SIGVTALRM);
2795ubf_select(
void *ptr)
2798 RUBY_DEBUG_LOG(
"wakeup th:%u", rb_th_serial(th));
2799 ubf_wakeup_thread(th);
2800 register_ubf_list(th);
2804ubf_threads_empty(
void)
2806 return ccan_list_empty(&ubf_list_head) != 0;
2810ubf_wakeup_all_threads(
void)
2815 ccan_list_for_each(&ubf_list_head, th, sched.node.ubf) {
2816 ubf_wakeup_thread(th);
2823#define register_ubf_list(th) (void)(th)
2824#define unregister_ubf_list(th) (void)(th)
2826static void ubf_wakeup_all_threads(
void) {
return; }
2827static bool ubf_threads_empty(
void) {
return true; }
2828#define ubf_list_atfork() do {} while (0)
2832#define WRITE_CONST(fd, str) (void)(write((fd),(str),sizeof(str)-1)<0)
2835rb_thread_wakeup_timer_thread(
int sig)
2841 timer_thread_wakeup_force();
2852 RUBY_VM_SET_TRAP_INTERRUPT(main_th_ec);
2854 if (vm->ubf_async_safe && main_th->unblock.func) {
2855 (main_th->unblock.func)(main_th->unblock.arg);
2862#define CLOSE_INVALIDATE_PAIR(expr) \
2863 close_invalidate_pair(expr,"close_invalidate: "#expr)
2865close_invalidate(
int *fdp,
const char *msg)
2870 if (close(fd) < 0) {
2871 async_bug_fd(msg,
errno, fd);
2876close_invalidate_pair(
int fds[2],
const char *msg)
2878 if (USE_EVENTFD && fds[0] == fds[1]) {
2880 close_invalidate(&fds[0], msg);
2883 close_invalidate(&fds[1], msg);
2884 close_invalidate(&fds[0], msg);
2894 oflags = fcntl(fd, F_GETFL);
2897 oflags |= O_NONBLOCK;
2898 err = fcntl(fd, F_SETFL, oflags);
2905setup_communication_pipe_internal(
int pipes[2])
2909 if (pipes[0] > 0 || pipes[1] > 0) {
2910 VM_ASSERT(pipes[0] > 0);
2911 VM_ASSERT(pipes[1] > 0);
2919#if USE_EVENTFD && defined(EFD_NONBLOCK) && defined(EFD_CLOEXEC)
2920 pipes[0] = pipes[1] = eventfd(0, EFD_NONBLOCK|EFD_CLOEXEC);
2922 if (pipes[0] >= 0) {
2930 rb_bug(
"can not create communication pipe");
2934 set_nonblock(pipes[0]);
2935 set_nonblock(pipes[1]);
2938#if !defined(SET_CURRENT_THREAD_NAME) && defined(__linux__) && defined(PR_SET_NAME)
2939# define SET_CURRENT_THREAD_NAME(name) prctl(PR_SET_NAME, name)
2944#if defined(__linux__)
2946#elif defined(__APPLE__)
2959#ifdef SET_CURRENT_THREAD_NAME
2961 if (!
NIL_P(loc = th->name)) {
2962 SET_CURRENT_THREAD_NAME(RSTRING_PTR(loc));
2964 else if ((loc = threadptr_invoke_proc_location(th)) !=
Qnil) {
2966 char buf[THREAD_NAME_MAX];
2971 p = strrchr(name,
'/');
2979 if (
len >=
sizeof(buf)) {
2980 buf[
sizeof(buf)-2] =
'*';
2981 buf[
sizeof(buf)-1] =
'\0';
2983 SET_CURRENT_THREAD_NAME(buf);
2989native_set_another_thread_name(rb_nativethread_id_t thread_id,
VALUE name)
2991#if defined SET_ANOTHER_THREAD_NAME || defined SET_CURRENT_THREAD_NAME
2992 char buf[THREAD_NAME_MAX];
2994# if !defined SET_ANOTHER_THREAD_NAME
2995 if (!pthread_equal(pthread_self(), thread_id))
return;
3000 if (n >= (
int)
sizeof(buf)) {
3001 memcpy(buf, s,
sizeof(buf)-1);
3002 buf[
sizeof(buf)-1] =
'\0';
3006# if defined SET_ANOTHER_THREAD_NAME
3007 SET_ANOTHER_THREAD_NAME(thread_id, s);
3008# elif defined SET_CURRENT_THREAD_NAME
3009 SET_CURRENT_THREAD_NAME(s);
3014#if defined(RB_THREAD_T_HAS_NATIVE_ID) || defined(__APPLE__)
3016native_thread_native_thread_id(
rb_thread_t *target_th)
3018 if (!target_th->nt)
return Qnil;
3020#ifdef RB_THREAD_T_HAS_NATIVE_ID
3021 int tid = target_th->nt->tid;
3022 if (tid == 0)
return Qnil;
3024#elif defined(__APPLE__)
3030# if (!defined(MAC_OS_X_VERSION_10_6) || \
3031 (MAC_OS_X_VERSION_MAX_ALLOWED < MAC_OS_X_VERSION_10_6) || \
3032 defined(__POWERPC__) )
3033 const bool no_pthread_threadid_np =
true;
3034# define NO_PTHREAD_MACH_THREAD_NP 1
3035# elif MAC_OS_X_VERSION_MIN_REQUIRED >= MAC_OS_X_VERSION_10_6
3036 const bool no_pthread_threadid_np =
false;
3038# if !(defined(__has_attribute) && __has_attribute(availability))
3040 __attribute__((weak))
int pthread_threadid_np(pthread_t, uint64_t*);
3043 const bool no_pthread_threadid_np = !&pthread_threadid_np;
3045 if (no_pthread_threadid_np) {
3046 return ULL2NUM(pthread_mach_thread_np(pthread_self()));
3048# ifndef NO_PTHREAD_MACH_THREAD_NP
3049 int e = pthread_threadid_np(target_th->nt->thread_id, &tid);
3051 return ULL2NUM((
unsigned long long)tid);
3055# define USE_NATIVE_THREAD_NATIVE_THREAD_ID 1
3057# define USE_NATIVE_THREAD_NATIVE_THREAD_ID 0
3061 rb_serial_t created_fork_gen;
3062 pthread_t pthread_id;
3066#if (HAVE_SYS_EPOLL_H || HAVE_SYS_EVENT_H) && USE_MN_THREADS
3069#if HAVE_SYS_EPOLL_H && USE_MN_THREADS
3070#define EPOLL_EVENTS_MAX 0x10
3071 struct epoll_event finished_events[EPOLL_EVENTS_MAX];
3072#elif HAVE_SYS_EVENT_H && USE_MN_THREADS
3073#define KQUEUE_EVENTS_MAX 0x10
3074 struct kevent finished_events[KQUEUE_EVENTS_MAX];
3078 struct ccan_list_head waiting;
3079 pthread_mutex_t waiting_lock;
3081#if (HAVE_SYS_EPOLL_H || HAVE_SYS_EVENT_H) && USE_MN_THREADS
3085 unsigned int fdmap_nchunks;
3088 .created_fork_gen = 0,
3091#define TIMER_THREAD_CREATED_P() (timer_th.created_fork_gen == current_fork_gen)
3093static void timer_thread_check_timeslice(
rb_vm_t *vm);
3094static int timer_thread_set_timeout(
rb_vm_t *vm);
3095static void timer_thread_wakeup_thread(
rb_thread_t *th, uint32_t event_serial);
3098#include "thread_pthread_mn.c"
3112timer_thread_set_timeout(
rb_vm_t *vm)
3119 ractor_sched_lock(vm, NULL);
3121 if ( !ccan_list_empty(&vm->ractor.sched.timeslice_threads)
3122 || !ubf_threads_empty()
3123 || vm->ractor.sched.grq_cnt > 0
3126 RUBY_DEBUG_LOG(
"timeslice:%d ubf:%d grq:%d",
3127 !ccan_list_empty(&vm->ractor.sched.timeslice_threads),
3128 !ubf_threads_empty(),
3129 (vm->ractor.sched.grq_cnt > 0));
3132 vm->ractor.sched.timeslice_wait_inf =
false;
3135 vm->ractor.sched.timeslice_wait_inf =
true;
3138 ractor_sched_unlock(vm, NULL);
3147 if (th && (th->sched.waiting_reason.flags & thread_sched_waiting_timeout)) {
3148 rb_hrtime_t now = rb_hrtime_now();
3149 rb_hrtime_t hrrel = rb_hrtime_sub(th->sched.waiting_reason.data.timeout, now);
3151 RUBY_DEBUG_LOG(
"th:%u now:%lu rel:%lu", rb_th_serial(th), (
unsigned long)now, (
unsigned long)hrrel);
3153 rb_hrtime_t msec = (hrrel + RB_HRTIME_PER_MSEC - 1) / RB_HRTIME_PER_MSEC;
3156 int thread_timeout = msec > INT_MAX ? INT_MAX : (int)msec;
3159 if (timeout < 0 || thread_timeout < timeout) {
3160 timeout = thread_timeout;
3166 RUBY_DEBUG_LOG(
"timeout:%d inf:%d", timeout, (
int)vm->ractor.sched.timeslice_wait_inf);
3174timer_thread_check_signal(
rb_vm_t *vm)
3178 int signum = rb_signal_buff_size();
3179 if (UNLIKELY(signum > 0) && vm->ractor.main_thread) {
3180 RUBY_DEBUG_LOG(
"signum:%d", signum);
3181 threadptr_trap_interrupt(vm->ractor.main_thread);
3186timer_thread_check_exceed(rb_hrtime_t abs, rb_hrtime_t now)
3192timer_thread_deq_wakeup(
rb_vm_t *vm, rb_hrtime_t now, uint32_t *event_serial)
3197 (w->flags & thread_sched_waiting_timeout) &&
3198 timer_thread_check_exceed(w->data.timeout, now)) {
3200 RUBY_DEBUG_LOG(
"wakeup th:%u", rb_th_serial(thread_sched_waiting_thread(w)));
3203 ccan_list_del_init(&w->node);
3207#if (HAVE_SYS_EPOLL_H || HAVE_SYS_EVENT_H) && USE_MN_THREADS
3209 timer_thread_unregister_waiting(th, w->data.fd, w->flags);
3213 w->flags = thread_sched_waiting_none;
3216 *event_serial = w->data.event_serial;
3226 if (sched->running != th && th->sched.event_serial == event_serial) {
3227 thread_sched_to_ready_common(sched, th,
true,
false);
3232timer_thread_wakeup_thread(
rb_thread_t *th, uint32_t event_serial)
3234 RUBY_DEBUG_LOG(
"th:%u", rb_th_serial(th));
3237 thread_sched_lock(sched, th);
3239 timer_thread_wakeup_thread_locked(sched, th, event_serial);
3241 thread_sched_unlock(sched, th);
3245timer_thread_check_timeout(
rb_vm_t *vm)
3247 rb_hrtime_t now = rb_hrtime_now();
3249 uint32_t event_serial;
3253 while ((th = timer_thread_deq_wakeup(vm, now, &event_serial)) != NULL) {
3255 timer_thread_wakeup_thread(th, event_serial);
3263timer_thread_check_timeslice(
rb_vm_t *vm)
3267 ccan_list_for_each(&vm->ractor.sched.timeslice_threads, th, sched.node.timeslice_threads) {
3268 RUBY_DEBUG_LOG(
"timeslice th:%u", rb_th_serial(th));
3269 RUBY_VM_SET_TIMER_INTERRUPT(th->ec);
3277 pthread_sigmask(0, NULL, &oldmask);
3278 if (sigismember(&oldmask, SIGVTALRM)) {
3282 RUBY_DEBUG_LOG(
"ok");
3287timer_thread_func(
void *ptr)
3290#if defined(RUBY_NT_SERIAL)
3294 RUBY_DEBUG_LOG(
"started%s",
"");
3297 timer_thread_check_signal(vm);
3298 timer_thread_check_timeout(vm);
3299 ubf_wakeup_all_threads();
3302 timer_thread_polling(vm);
3305 RUBY_DEBUG_LOG(
"terminated");
3311signal_communication_pipe(
int fd)
3314 const uint64_t buff = 1;
3316 const char buff =
'!';
3323 if ((result = write(fd, &buff,
sizeof(buff))) <= 0) {
3326 case EINTR:
goto retry;
3328#if defined(EWOULDBLOCK) && EWOULDBLOCK != EAGAIN
3333 async_bug_fd(
"rb_thread_wakeup_timer_thread: write", e, fd);
3336 if (TT_DEBUG) WRITE_CONST(2,
"rb_thread_wakeup_timer_thread: write\n");
3344timer_thread_wakeup_force(
void)
3347 signal_communication_pipe(timer_th.comm_fds[1]);
3351timer_thread_wakeup_locked(
rb_vm_t *vm)
3354 ASSERT_ractor_sched_locked(vm, NULL);
3356 if (timer_th.created_fork_gen == current_fork_gen) {
3357 if (vm->ractor.sched.timeslice_wait_inf) {
3358 RUBY_DEBUG_LOG(
"wakeup with fd:%d", timer_th.comm_fds[1]);
3359 timer_thread_wakeup_force();
3362 RUBY_DEBUG_LOG(
"will be wakeup...");
3368timer_thread_wakeup(
void)
3372 ractor_sched_lock(vm, NULL);
3374 timer_thread_wakeup_locked(vm);
3376 ractor_sched_unlock(vm, NULL);
3380rb_thread_create_timer_thread(
void)
3382 rb_serial_t created_fork_gen = timer_th.created_fork_gen;
3384 RUBY_DEBUG_LOG(
"fork_gen create:%d current:%d", (
int)created_fork_gen, (
int)current_fork_gen);
3386 timer_th.created_fork_gen = current_fork_gen;
3388 if (created_fork_gen != current_fork_gen) {
3389 if (created_fork_gen != 0) {
3390 RUBY_DEBUG_LOG(
"forked child process");
3392 CLOSE_INVALIDATE_PAIR(timer_th.comm_fds);
3393#if HAVE_SYS_EPOLL_H && USE_MN_THREADS
3394 close_invalidate(&timer_th.event_fd,
"close event_fd");
3395#elif HAVE_SYS_EVENT_H && USE_MN_THREADS
3398 timer_th.event_fd = -1;
3406 ccan_list_head_init(&timer_th.waiting);
3410 setup_communication_pipe_internal(timer_th.comm_fds);
3413 timer_thread_setup_mn();
3416 int err = pthread_create(&timer_th.pthread_id, NULL, timer_thread_func, GET_VM());
3425native_stop_timer_thread(
void)
3429 RUBY_DEBUG_LOG(
"wakeup send %d", timer_th.comm_fds[1]);
3430 timer_thread_wakeup_force();
3431 RUBY_DEBUG_LOG(
"wakeup sent");
3432 pthread_join(timer_th.pthread_id, NULL);
3434 if (TT_DEBUG) fprintf(stderr,
"stop timer thread\n");
3440native_reset_timer_thread(
void)
3445#ifdef HAVE_SIGALTSTACK
3447ruby_stack_overflowed_p(
const rb_thread_t *th,
const void *addr)
3451 const size_t water_mark = RB_THREAD_PAGE_SIZE;
3452 STACK_GROW_DIR_DETECTION;
3455 size = th->ec->machine.stack_maxsize;
3456 base = (
char *)th->ec->machine.stack_start - STACK_DIR_UPPER(0, size);
3458#ifdef STACKADDR_AVAILABLE
3459 else if (get_stack(&base, &size) == 0) {
3462 if (pthread_equal(pthread_self(), native_main_thread.id)) {
3464 if (getrlimit(RLIMIT_STACK, &rlim) == 0 && rlim.rlim_cur > size) {
3465 size = (size_t)rlim.rlim_cur;
3469 base = (
char *)base + STACK_DIR_UPPER(+size, -size);
3476 if (size > water_mark) size = water_mark;
3477 if (IS_STACK_DIR_UPPER()) {
3478 if (size > ~(
size_t)base+1) size = ~(
size_t)base+1;
3479 if (addr > base && addr <= (
void *)((
char *)base + size))
return 1;
3482 if (size > (
size_t)base) size = (
size_t)base;
3483 if (addr > (
void *)((
char *)base - size) && addr <= base)
return 1;
3493 if (fd < 0)
return 0;
3495 if (fd == timer_th.comm_fds[0] ||
3496 fd == timer_th.comm_fds[1]
3497#if (HAVE_SYS_EPOLL_H || HAVE_SYS_EVENT_H) && USE_MN_THREADS
3498 || fd == timer_th.event_fd
3501 goto check_fork_gen;
3506 if (timer_th.created_fork_gen == current_fork_gen) {
3518 return pthread_self();
3521#if defined(USE_POLL) && !defined(HAVE_PPOLL)
3524ruby_ppoll(
struct pollfd *fds, nfds_t nfds,
3525 const struct timespec *ts,
const sigset_t *sigmask)
3532 if (ts->tv_sec > INT_MAX/1000)
3533 timeout_ms = INT_MAX;
3535 tmp = (int)(ts->tv_sec * 1000);
3537 tmp2 = (int)((ts->tv_nsec + 999999L) / (1000L * 1000L));
3538 if (INT_MAX - tmp < tmp2)
3539 timeout_ms = INT_MAX;
3541 timeout_ms = (int)(tmp + tmp2);
3547 return poll(fds, nfds, timeout_ms);
3549# define ppoll(fds,nfds,ts,sigmask) ruby_ppoll((fds),(nfds),(ts),(sigmask))
3557 RUBY_DEBUG_LOG(
"rel:%d", rel ? (
int)*rel : 0);
3559 if (th_has_dedicated_nt(th)) {
3560 native_cond_sleep(th, rel);
3563 thread_sched_wait_events(sched, th, -1, thread_sched_waiting_timeout, rel);
3567 thread_sched_to_waiting_until_wakeup(sched, th);
3570 RUBY_DEBUG_LOG(
"wakeup");
3574static pthread_rwlock_t rb_thread_fork_rw_lock = PTHREAD_RWLOCK_INITIALIZER;
3577rb_thread_release_fork_lock(
void)
3580 if ((r = pthread_rwlock_unlock(&rb_thread_fork_rw_lock))) {
3586rb_thread_reset_fork_lock(
void)
3589 if ((r = pthread_rwlock_destroy(&rb_thread_fork_rw_lock))) {
3593 if ((r = pthread_rwlock_init(&rb_thread_fork_rw_lock, NULL))) {
3599rb_thread_prevent_fork(
void *(*func)(
void *),
void *data)
3602 if ((r = pthread_rwlock_rdlock(&rb_thread_fork_rw_lock))) {
3605 void *result = func(data);
3606 rb_thread_release_fork_lock();
3611rb_thread_acquire_fork_lock(
void)
3614 if ((r = pthread_rwlock_wrlock(&rb_thread_fork_rw_lock))) {
3621struct rb_internal_thread_event_hook {
3622 rb_internal_thread_event_callback callback;
3626 struct rb_internal_thread_event_hook *next;
3629static pthread_rwlock_t rb_internal_thread_event_hooks_rw_lock = PTHREAD_RWLOCK_INITIALIZER;
3636rb_thread_event_hooks_registered_p(
void)
3639 if ((r = pthread_rwlock_rdlock(&rb_internal_thread_event_hooks_rw_lock))) {
3642 const bool registered = (rb_internal_thread_event_hooks != NULL);
3643 if ((r = pthread_rwlock_unlock(&rb_internal_thread_event_hooks_rw_lock))) {
3649#if defined(HAVE_WORKING_FORK)
3651rb_internal_thread_event_hooks_rw_lock_atfork(
void)
3660 rb_internal_thread_event_hooks_rw_lock =
3661 (pthread_rwlock_t)PTHREAD_RWLOCK_INITIALIZER;
3665rb_internal_thread_event_hook_t *
3668 rb_internal_thread_event_hook_t *hook =
ALLOC_N(rb_internal_thread_event_hook_t, 1);
3669 hook->callback = callback;
3670 hook->user_data = user_data;
3671 hook->event = internal_event;
3674 if ((r = pthread_rwlock_wrlock(&rb_internal_thread_event_hooks_rw_lock))) {
3678 hook->next = rb_internal_thread_event_hooks;
3679 ATOMIC_PTR_EXCHANGE(rb_internal_thread_event_hooks, hook);
3681 if ((r = pthread_rwlock_unlock(&rb_internal_thread_event_hooks_rw_lock))) {
3691 if ((r = pthread_rwlock_wrlock(&rb_internal_thread_event_hooks_rw_lock))) {
3695 bool success = FALSE;
3697 if (rb_internal_thread_event_hooks == hook) {
3698 ATOMIC_PTR_EXCHANGE(rb_internal_thread_event_hooks, hook->next);
3702 rb_internal_thread_event_hook_t *h = rb_internal_thread_event_hooks;
3705 if (h->next == hook) {
3706 h->next = hook->next;
3710 }
while ((h = h->next));
3713 if ((r = pthread_rwlock_unlock(&rb_internal_thread_event_hooks_rw_lock))) {
3730 if (th->self == 0)
return;
3731 if ((r = pthread_rwlock_rdlock(&rb_internal_thread_event_hooks_rw_lock))) {
3735 if (rb_internal_thread_event_hooks) {
3736 rb_internal_thread_event_hook_t *h = rb_internal_thread_event_hooks;
3738 if (h->event & event) {
3742 (*h->callback)(event, &event_data, h->user_data);
3744 }
while((h = h->next));
3746 if ((r = pthread_rwlock_unlock(&rb_internal_thread_event_hooks_rw_lock))) {
3757 bool is_snt = th->nt->dedicated == 0;
3758 native_thread_dedicated_inc(th->vm, th->ractor, th->nt);
3764rb_thread_malloc_stack_set(
rb_thread_t *th,
void *stack,
size_t stack_size)
3766 th->sched.malloc_stack =
true;
3767 th->sched.context_stack = stack;
3768 th->sched.context_stack_size = stack_size;
std::atomic< unsigned > rb_atomic_t
Type that is eligible for atomic operations.
#define RUBY_ATOMIC_FETCH_ADD(var, val)
Atomically replaces the value pointed by var with the result of addition of val to the old value of v...
#define RUBY_ATOMIC_ADD(var, val)
Identical to RUBY_ATOMIC_FETCH_ADD, except for the return type.
#define RUBY_ATOMIC_DEC(var)
Atomically decrements the value pointed by var.
#define RUBY_ATOMIC_LOAD(var)
Atomic load.
#define RUBY_ATOMIC_SET(var, val)
Identical to RUBY_ATOMIC_EXCHANGE, except for the return type.
uint32_t rb_event_flag_t
Represents event(s).
#define INT2FIX
Old name of RB_INT2FIX.
#define ZALLOC
Old name of RB_ZALLOC.
#define ALLOC_N
Old name of RB_ALLOC_N.
#define ULL2NUM
Old name of RB_ULL2NUM.
#define NUM2INT
Old name of RB_NUM2INT.
#define Qnil
Old name of RUBY_Qnil.
#define NIL_P
Old name of RB_NIL_P.
VALUE rb_eNotImpError
NotImplementedError exception.
void rb_syserr_fail(int e, const char *mesg)
Raises appropriate exception that represents a C errno.
void rb_bug_errno(const char *mesg, int errno_arg)
This is a wrapper of rb_bug() which automatically constructs appropriate message from the passed errn...
int rb_cloexec_pipe(int fildes[2])
Opens a pipe with closing on exec.
void rb_update_max_fd(int fd)
Informs the interpreter that the passed fd can be the max.
int rb_reserved_fd_p(int fd)
Queries if the given FD is reserved or not.
void rb_unblock_function_t(void *)
This is the type of UBFs.
void rb_timespec_now(struct timespec *ts)
Fills the current time into the given struct.
int len
Length of the buffer.
#define RUBY_INTERNAL_THREAD_EVENT_RESUMED
Triggered when a thread successfully acquired the GVL.
rb_internal_thread_event_hook_t * rb_internal_thread_add_event_hook(rb_internal_thread_event_callback func, rb_event_flag_t events, void *data)
Registers a thread event hook function.
#define RUBY_INTERNAL_THREAD_EVENT_EXITED
Triggered when a thread exits.
#define RUBY_INTERNAL_THREAD_EVENT_SUSPENDED
Triggered when a thread released the GVL.
bool rb_thread_lock_native_thread(void)
Declare the current Ruby thread should acquire a dedicated native thread on M:N thread scheduler.
#define RUBY_INTERNAL_THREAD_EVENT_STARTED
Triggered when a new thread is started.
bool rb_internal_thread_remove_event_hook(rb_internal_thread_event_hook_t *hook)
Unregister the passed hook.
#define RUBY_INTERNAL_THREAD_EVENT_READY
Triggered when a thread attempt to acquire the GVL.
#define RBIMPL_ATTR_MAYBE_UNUSED()
Wraps (or simulates) [[maybe_unused]]
#define RB_GC_GUARD(v)
Prevents premature destruction of local objects.
#define rb_fd_select
Waits for multiple file descriptors at once.
#define RARRAY_AREF(a, i)
#define RSTRING_GETMEM(str, ptrvar, lenvar)
Convenient macro to obtain the contents and length at once.
#define errno
Ractor-aware version of errno.
The data structure which wraps the fd_set bitmap used by select(2).
rb_nativethread_id_t rb_nativethread_self(void)
Queries the ID of the native thread that is calling this function.
void rb_native_mutex_lock(rb_nativethread_lock_t *lock)
Just another name of rb_nativethread_lock_lock.
void rb_native_cond_initialize(rb_nativethread_cond_t *cond)
Fills the passed condition variable with an initial value.
int rb_native_mutex_trylock(rb_nativethread_lock_t *lock)
Identical to rb_native_mutex_lock(), except it doesn't block in case rb_native_mutex_lock() would.
void rb_native_cond_broadcast(rb_nativethread_cond_t *cond)
Signals a condition variable.
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.
void rb_native_mutex_destroy(rb_nativethread_lock_t *lock)
Just another name of rb_nativethread_lock_destroy.
void rb_native_cond_destroy(rb_nativethread_cond_t *cond)
Destroys the passed condition variable.
void rb_native_cond_signal(rb_nativethread_cond_t *cond)
Signals a condition variable.
void rb_native_cond_wait(rb_nativethread_cond_t *cond, rb_nativethread_lock_t *mutex)
Waits for the passed condition variable to be signalled.
void rb_native_cond_timedwait(rb_nativethread_cond_t *cond, rb_nativethread_lock_t *mutex, unsigned long msec)
Identical to rb_native_cond_wait(), except it additionally takes timeout in msec resolution.
uintptr_t VALUE
Type that represents a Ruby object.