Ruby  2.0.0p594(2014-10-27revision48167)
thread.c
Go to the documentation of this file.
00001 /**********************************************************************
00002 
00003   thread.c -
00004 
00005   $Author: usa $
00006 
00007   Copyright (C) 2004-2007 Koichi Sasada
00008 
00009 **********************************************************************/
00010 
00011 /*
00012   YARV Thread Design
00013 
00014   model 1: Userlevel Thread
00015     Same as traditional ruby thread.
00016 
00017   model 2: Native Thread with Global VM lock
00018     Using pthread (or Windows thread) and Ruby threads run concurrent.
00019 
00020   model 3: Native Thread with fine grain lock
00021     Using pthread and Ruby threads run concurrent or parallel.
00022 
00023 ------------------------------------------------------------------------
00024 
00025   model 2:
00026     A thread has mutex (GVL: Global VM Lock or Giant VM Lock) can run.
00027     When thread scheduling, running thread release GVL.  If running thread
00028     try blocking operation, this thread must release GVL and another
00029     thread can continue this flow.  After blocking operation, thread
00030     must check interrupt (RUBY_VM_CHECK_INTS).
00031 
00032     Every VM can run parallel.
00033 
00034     Ruby threads are scheduled by OS thread scheduler.
00035 
00036 ------------------------------------------------------------------------
00037 
00038   model 3:
00039     Every threads run concurrent or parallel and to access shared object
00040     exclusive access control is needed.  For example, to access String
00041     object or Array object, fine grain lock must be locked every time.
00042  */
00043 
00044 
00045 /*
00046  * FD_SET, FD_CLR and FD_ISSET have a small sanity check when using glibc
00047  * 2.15 or later and set _FORTIFY_SOURCE > 0.
00048  * However, the implementation is wrong. Even though Linux's select(2)
00049  * support large fd size (>FD_SETSIZE), it wrongly assume fd is always
00050  * less than FD_SETSIZE (i.e. 1024). And then when enabling HAVE_RB_FD_INIT,
00051  * it doesn't work correctly and makes program abort. Therefore we need to
00052  * disable FORTY_SOURCE until glibc fixes it.
00053  */
00054 #undef _FORTIFY_SOURCE
00055 #undef __USE_FORTIFY_LEVEL
00056 #define __USE_FORTIFY_LEVEL 0
00057 
00058 /* for model 2 */
00059 
00060 #include "eval_intern.h"
00061 #include "gc.h"
00062 #include "internal.h"
00063 #include "ruby/io.h"
00064 #include "ruby/thread.h"
00065 
00066 #ifndef USE_NATIVE_THREAD_PRIORITY
00067 #define USE_NATIVE_THREAD_PRIORITY 0
00068 #define RUBY_THREAD_PRIORITY_MAX 3
00069 #define RUBY_THREAD_PRIORITY_MIN -3
00070 #endif
00071 
00072 #ifndef THREAD_DEBUG
00073 #define THREAD_DEBUG 0
00074 #endif
00075 
00076 #define TIMET_MAX (~(time_t)0 <= 0 ? (time_t)((~(unsigned_time_t)0) >> 1) : (time_t)(~(unsigned_time_t)0))
00077 #define TIMET_MIN (~(time_t)0 <= 0 ? (time_t)(((unsigned_time_t)1) << (sizeof(time_t) * CHAR_BIT - 1)) : (time_t)0)
00078 
00079 VALUE rb_cMutex;
00080 VALUE rb_cThreadShield;
00081 
00082 static VALUE sym_immediate;
00083 static VALUE sym_on_blocking;
00084 static VALUE sym_never;
00085 
00086 static void sleep_timeval(rb_thread_t *th, struct timeval time, int spurious_check);
00087 static void sleep_wait_for_interrupt(rb_thread_t *th, double sleepsec, int spurious_check);
00088 static void sleep_forever(rb_thread_t *th, int nodeadlock, int spurious_check);
00089 static double timeofday(void);
00090 static int rb_threadptr_dead(rb_thread_t *th);
00091 static void rb_check_deadlock(rb_vm_t *vm);
00092 static int rb_threadptr_pending_interrupt_empty_p(rb_thread_t *th);
00093 
00094 #define eKillSignal INT2FIX(0)
00095 #define eTerminateSignal INT2FIX(1)
00096 static volatile int system_working = 1;
00097 
00098 #define closed_stream_error GET_VM()->special_exceptions[ruby_error_closed_stream]
00099 
00100 inline static void
00101 st_delete_wrap(st_table *table, st_data_t key)
00102 {
00103     st_delete(table, &key, 0);
00104 }
00105 
00106 /********************************************************************************/
00107 
00108 #define THREAD_SYSTEM_DEPENDENT_IMPLEMENTATION
00109 
00110 struct rb_blocking_region_buffer {
00111     enum rb_thread_status prev_status;
00112     struct rb_unblock_callback oldubf;
00113 };
00114 
00115 static int set_unblock_function(rb_thread_t *th, rb_unblock_function_t *func, void *arg,
00116                                 struct rb_unblock_callback *old, int fail_if_interrupted);
00117 static void reset_unblock_function(rb_thread_t *th, const struct rb_unblock_callback *old);
00118 
00119 static inline int blocking_region_begin(rb_thread_t *th, struct rb_blocking_region_buffer *region,
00120                                         rb_unblock_function_t *ubf, void *arg, int fail_if_interrupted);
00121 static inline void blocking_region_end(rb_thread_t *th, struct rb_blocking_region_buffer *region);
00122 
00123 #ifdef __ia64
00124 #define RB_GC_SAVE_MACHINE_REGISTER_STACK(th)          \
00125     do{(th)->machine_register_stack_end = rb_ia64_bsp();}while(0)
00126 #else
00127 #define RB_GC_SAVE_MACHINE_REGISTER_STACK(th)
00128 #endif
00129 #define RB_GC_SAVE_MACHINE_CONTEXT(th)                          \
00130     do {                                                        \
00131         FLUSH_REGISTER_WINDOWS;                                 \
00132         RB_GC_SAVE_MACHINE_REGISTER_STACK(th);                  \
00133         setjmp((th)->machine_regs);                             \
00134         SET_MACHINE_STACK_END(&(th)->machine_stack_end);        \
00135     } while (0)
00136 
00137 #define GVL_UNLOCK_BEGIN() do { \
00138   rb_thread_t *_th_stored = GET_THREAD(); \
00139   RB_GC_SAVE_MACHINE_CONTEXT(_th_stored); \
00140   gvl_release(_th_stored->vm);
00141 
00142 #define GVL_UNLOCK_END() \
00143   gvl_acquire(_th_stored->vm, _th_stored); \
00144   rb_thread_set_current(_th_stored); \
00145 } while(0)
00146 
00147 #ifdef __GNUC__
00148 #define only_if_constant(expr, notconst) (__builtin_constant_p(expr) ? (expr) : (notconst))
00149 #else
00150 #define only_if_constant(expr, notconst) notconst
00151 #endif
00152 #define BLOCKING_REGION(exec, ubf, ubfarg, fail_if_interrupted) do { \
00153     rb_thread_t *__th = GET_THREAD(); \
00154     struct rb_blocking_region_buffer __region; \
00155     if (blocking_region_begin(__th, &__region, (ubf), (ubfarg), fail_if_interrupted) || \
00156         /* always return true unless fail_if_interrupted */ \
00157         !only_if_constant(fail_if_interrupted, TRUE)) { \
00158         exec; \
00159         blocking_region_end(__th, &__region); \
00160     }; \
00161 } while(0)
00162 
00163 #if THREAD_DEBUG
00164 #ifdef HAVE_VA_ARGS_MACRO
00165 void rb_thread_debug(const char *file, int line, const char *fmt, ...);
00166 #define thread_debug(fmt, ...) rb_thread_debug(__FILE__, __LINE__, fmt, ##__VA_ARGS__)
00167 #define POSITION_FORMAT "%s:%d:"
00168 #define POSITION_ARGS ,file, line
00169 #else
00170 void rb_thread_debug(const char *fmt, ...);
00171 #define thread_debug rb_thread_debug
00172 #define POSITION_FORMAT
00173 #define POSITION_ARGS
00174 #endif
00175 
00176 # if THREAD_DEBUG < 0
00177 static int rb_thread_debug_enabled;
00178 
00179 /*
00180  *  call-seq:
00181  *     Thread.DEBUG     -> num
00182  *
00183  *  Returns the thread debug level.  Available only if compiled with
00184  *  THREAD_DEBUG=-1.
00185  */
00186 
00187 static VALUE
00188 rb_thread_s_debug(void)
00189 {
00190     return INT2NUM(rb_thread_debug_enabled);
00191 }
00192 
00193 /*
00194  *  call-seq:
00195  *     Thread.DEBUG = num
00196  *
00197  *  Sets the thread debug level.  Available only if compiled with
00198  *  THREAD_DEBUG=-1.
00199  */
00200 
00201 static VALUE
00202 rb_thread_s_debug_set(VALUE self, VALUE val)
00203 {
00204     rb_thread_debug_enabled = RTEST(val) ? NUM2INT(val) : 0;
00205     return val;
00206 }
00207 # else
00208 # define rb_thread_debug_enabled THREAD_DEBUG
00209 # endif
00210 #else
00211 #define thread_debug if(0)printf
00212 #endif
00213 
00214 #ifndef __ia64
00215 #define thread_start_func_2(th, st, rst) thread_start_func_2(th, st)
00216 #endif
00217 NOINLINE(static int thread_start_func_2(rb_thread_t *th, VALUE *stack_start,
00218                                         VALUE *register_stack_start));
00219 static void timer_thread_function(void *);
00220 
00221 #if   defined(_WIN32)
00222 #include "thread_win32.c"
00223 
00224 #define DEBUG_OUT() \
00225   WaitForSingleObject(&debug_mutex, INFINITE); \
00226   printf(POSITION_FORMAT"%p - %s" POSITION_ARGS, GetCurrentThreadId(), buf); \
00227   fflush(stdout); \
00228   ReleaseMutex(&debug_mutex);
00229 
00230 #elif defined(HAVE_PTHREAD_H)
00231 #include "thread_pthread.c"
00232 
00233 #define DEBUG_OUT() \
00234   pthread_mutex_lock(&debug_mutex); \
00235   printf(POSITION_FORMAT"%#"PRIxVALUE" - %s" POSITION_ARGS, (VALUE)pthread_self(), buf); \
00236   fflush(stdout); \
00237   pthread_mutex_unlock(&debug_mutex);
00238 
00239 #else
00240 #error "unsupported thread type"
00241 #endif
00242 
00243 #if THREAD_DEBUG
00244 static int debug_mutex_initialized = 1;
00245 static rb_thread_lock_t debug_mutex;
00246 
00247 void
00248 rb_thread_debug(
00249 #ifdef HAVE_VA_ARGS_MACRO
00250     const char *file, int line,
00251 #endif
00252     const char *fmt, ...)
00253 {
00254     va_list args;
00255     char buf[BUFSIZ];
00256 
00257     if (!rb_thread_debug_enabled) return;
00258 
00259     if (debug_mutex_initialized == 1) {
00260         debug_mutex_initialized = 0;
00261         native_mutex_initialize(&debug_mutex);
00262     }
00263 
00264     va_start(args, fmt);
00265     vsnprintf(buf, BUFSIZ, fmt, args);
00266     va_end(args);
00267 
00268     DEBUG_OUT();
00269 }
00270 #endif
00271 
00272 void
00273 rb_vm_gvl_destroy(rb_vm_t *vm)
00274 {
00275     gvl_release(vm);
00276     gvl_destroy(vm);
00277     native_mutex_destroy(&vm->thread_destruct_lock);
00278 }
00279 
00280 void
00281 rb_thread_lock_unlock(rb_thread_lock_t *lock)
00282 {
00283     native_mutex_unlock(lock);
00284 }
00285 
00286 void
00287 rb_thread_lock_destroy(rb_thread_lock_t *lock)
00288 {
00289     native_mutex_destroy(lock);
00290 }
00291 
00292 static int
00293 set_unblock_function(rb_thread_t *th, rb_unblock_function_t *func, void *arg,
00294                      struct rb_unblock_callback *old, int fail_if_interrupted)
00295 {
00296   check_ints:
00297     if (fail_if_interrupted) {
00298         if (RUBY_VM_INTERRUPTED_ANY(th)) {
00299             return FALSE;
00300         }
00301     }
00302     else {
00303         RUBY_VM_CHECK_INTS(th);
00304     }
00305 
00306     native_mutex_lock(&th->interrupt_lock);
00307     if (RUBY_VM_INTERRUPTED_ANY(th)) {
00308         native_mutex_unlock(&th->interrupt_lock);
00309         goto check_ints;
00310     }
00311     else {
00312         if (old) *old = th->unblock;
00313         th->unblock.func = func;
00314         th->unblock.arg = arg;
00315     }
00316     native_mutex_unlock(&th->interrupt_lock);
00317 
00318     return TRUE;
00319 }
00320 
00321 static void
00322 reset_unblock_function(rb_thread_t *th, const struct rb_unblock_callback *old)
00323 {
00324     native_mutex_lock(&th->interrupt_lock);
00325     th->unblock = *old;
00326     native_mutex_unlock(&th->interrupt_lock);
00327 }
00328 
00329 static void
00330 rb_threadptr_interrupt_common(rb_thread_t *th, int trap)
00331 {
00332     native_mutex_lock(&th->interrupt_lock);
00333     if (trap)
00334         RUBY_VM_SET_TRAP_INTERRUPT(th);
00335     else
00336         RUBY_VM_SET_INTERRUPT(th);
00337     if (th->unblock.func) {
00338         (th->unblock.func)(th->unblock.arg);
00339     }
00340     else {
00341         /* none */
00342     }
00343     native_mutex_unlock(&th->interrupt_lock);
00344 }
00345 
00346 void
00347 rb_threadptr_interrupt(rb_thread_t *th)
00348 {
00349     rb_threadptr_interrupt_common(th, 0);
00350 }
00351 
00352 void
00353 rb_threadptr_trap_interrupt(rb_thread_t *th)
00354 {
00355     rb_threadptr_interrupt_common(th, 1);
00356 }
00357 
00358 static int
00359 terminate_i(st_data_t key, st_data_t val, rb_thread_t *main_thread)
00360 {
00361     VALUE thval = key;
00362     rb_thread_t *th;
00363     GetThreadPtr(thval, th);
00364 
00365     if (th != main_thread) {
00366         thread_debug("terminate_i: %p\n", (void *)th);
00367         rb_threadptr_pending_interrupt_enque(th, eTerminateSignal);
00368         rb_threadptr_interrupt(th);
00369     }
00370     else {
00371         thread_debug("terminate_i: main thread (%p)\n", (void *)th);
00372     }
00373     return ST_CONTINUE;
00374 }
00375 
00376 typedef struct rb_mutex_struct
00377 {
00378     rb_thread_lock_t lock;
00379     rb_thread_cond_t cond;
00380     struct rb_thread_struct volatile *th;
00381     int cond_waiting;
00382     struct rb_mutex_struct *next_mutex;
00383     int allow_trap;
00384 } rb_mutex_t;
00385 
00386 static void rb_mutex_abandon_all(rb_mutex_t *mutexes);
00387 static void rb_mutex_abandon_keeping_mutexes(rb_thread_t *th);
00388 static void rb_mutex_abandon_locking_mutex(rb_thread_t *th);
00389 static const char* rb_mutex_unlock_th(rb_mutex_t *mutex, rb_thread_t volatile *th);
00390 
00391 void
00392 rb_threadptr_unlock_all_locking_mutexes(rb_thread_t *th)
00393 {
00394     const char *err;
00395     rb_mutex_t *mutex;
00396     rb_mutex_t *mutexes = th->keeping_mutexes;
00397 
00398     while (mutexes) {
00399         mutex = mutexes;
00400         /* rb_warn("mutex #<%p> remains to be locked by terminated thread",
00401                 mutexes); */
00402         mutexes = mutex->next_mutex;
00403         err = rb_mutex_unlock_th(mutex, th);
00404         if (err) rb_bug("invalid keeping_mutexes: %s", err);
00405     }
00406 }
00407 
00408 void
00409 rb_thread_terminate_all(void)
00410 {
00411     rb_thread_t *th = GET_THREAD(); /* main thread */
00412     rb_vm_t *vm = th->vm;
00413 
00414     if (vm->main_thread != th) {
00415         rb_bug("rb_thread_terminate_all: called by child thread (%p, %p)",
00416                (void *)vm->main_thread, (void *)th);
00417     }
00418 
00419     /* unlock all locking mutexes */
00420     rb_threadptr_unlock_all_locking_mutexes(th);
00421 
00422   retry:
00423     thread_debug("rb_thread_terminate_all (main thread: %p)\n", (void *)th);
00424     st_foreach(vm->living_threads, terminate_i, (st_data_t)th);
00425 
00426     while (!rb_thread_alone()) {
00427         int state;
00428 
00429         TH_PUSH_TAG(th);
00430         if ((state = TH_EXEC_TAG()) == 0) {
00431             native_sleep(th, 0);
00432             RUBY_VM_CHECK_INTS_BLOCKING(th);
00433         }
00434         TH_POP_TAG();
00435 
00436         if (state) {
00437             goto retry;
00438         }
00439     }
00440 }
00441 
00442 static void
00443 thread_cleanup_func_before_exec(void *th_ptr)
00444 {
00445     rb_thread_t *th = th_ptr;
00446     th->status = THREAD_KILLED;
00447     th->machine_stack_start = th->machine_stack_end = 0;
00448 #ifdef __ia64
00449     th->machine_register_stack_start = th->machine_register_stack_end = 0;
00450 #endif
00451 }
00452 
00453 static void
00454 thread_cleanup_func(void *th_ptr, int atfork)
00455 {
00456     rb_thread_t *th = th_ptr;
00457 
00458     th->locking_mutex = Qfalse;
00459     thread_cleanup_func_before_exec(th_ptr);
00460 
00461     /*
00462      * Unfortunately, we can't release native threading resource at fork
00463      * because libc may have unstable locking state therefore touching
00464      * a threading resource may cause a deadlock.
00465      */
00466     if (atfork)
00467         return;
00468 
00469     native_mutex_destroy(&th->interrupt_lock);
00470     native_thread_destroy(th);
00471 }
00472 
00473 static VALUE rb_threadptr_raise(rb_thread_t *, int, VALUE *);
00474 
00475 void
00476 ruby_thread_init_stack(rb_thread_t *th)
00477 {
00478     native_thread_init_stack(th);
00479 }
00480 
00481 static int
00482 thread_start_func_2(rb_thread_t *th, VALUE *stack_start, VALUE *register_stack_start)
00483 {
00484     int state;
00485     VALUE args = th->first_args;
00486     rb_proc_t *proc;
00487     rb_thread_list_t *join_list;
00488     rb_thread_t *main_th;
00489     VALUE errinfo = Qnil;
00490 # ifdef USE_SIGALTSTACK
00491     void rb_register_sigaltstack(rb_thread_t *th);
00492 
00493     rb_register_sigaltstack(th);
00494 # endif
00495 
00496     if (th == th->vm->main_thread)
00497         rb_bug("thread_start_func_2 must not used for main thread");
00498 
00499     ruby_thread_set_native(th);
00500 
00501     th->machine_stack_start = stack_start;
00502 #ifdef __ia64
00503     th->machine_register_stack_start = register_stack_start;
00504 #endif
00505     thread_debug("thread start: %p\n", (void *)th);
00506 
00507     gvl_acquire(th->vm, th);
00508     {
00509         thread_debug("thread start (get lock): %p\n", (void *)th);
00510         rb_thread_set_current(th);
00511 
00512         TH_PUSH_TAG(th);
00513         if ((state = EXEC_TAG()) == 0) {
00514             SAVE_ROOT_JMPBUF(th, {
00515                 if (!th->first_func) {
00516                     GetProcPtr(th->first_proc, proc);
00517                     th->errinfo = Qnil;
00518                     th->root_lep = rb_vm_ep_local_ep(proc->block.ep);
00519                     th->root_svar = Qnil;
00520                     EXEC_EVENT_HOOK(th, RUBY_EVENT_THREAD_BEGIN, th->self, 0, 0, Qundef);
00521                     th->value = rb_vm_invoke_proc(th, proc, (int)RARRAY_LEN(args), RARRAY_PTR(args), 0);
00522                     EXEC_EVENT_HOOK(th, RUBY_EVENT_THREAD_END, th->self, 0, 0, Qundef);
00523                 }
00524                 else {
00525                     th->value = (*th->first_func)((void *)args);
00526                 }
00527             });
00528         }
00529         else {
00530             errinfo = th->errinfo;
00531             if (state == TAG_FATAL) {
00532                 /* fatal error within this thread, need to stop whole script */
00533             }
00534             else if (th->safe_level >= 4) {
00535                 /* Ignore it. Main thread shouldn't be harmed from untrusted thread. */
00536                 errinfo = Qnil;
00537             }
00538             else if (rb_obj_is_kind_of(errinfo, rb_eSystemExit)) {
00539                 /* exit on main_thread. */
00540             }
00541             else if (th->vm->thread_abort_on_exception ||
00542                      th->abort_on_exception || RTEST(ruby_debug)) {
00543                 /* exit on main_thread */
00544             }
00545             else {
00546                 errinfo = Qnil;
00547             }
00548             th->value = Qnil;
00549         }
00550 
00551         th->status = THREAD_KILLED;
00552         thread_debug("thread end: %p\n", (void *)th);
00553 
00554         main_th = th->vm->main_thread;
00555         if (main_th == th) {
00556             ruby_stop(0);
00557         }
00558         if (RB_TYPE_P(errinfo, T_OBJECT)) {
00559             /* treat with normal error object */
00560             rb_threadptr_raise(main_th, 1, &errinfo);
00561         }
00562         TH_POP_TAG();
00563 
00564         /* locking_mutex must be Qfalse */
00565         if (th->locking_mutex != Qfalse) {
00566             rb_bug("thread_start_func_2: locking_mutex must not be set (%p:%"PRIxVALUE")",
00567                    (void *)th, th->locking_mutex);
00568         }
00569 
00570         /* delete self other than main thread from living_threads */
00571         st_delete_wrap(th->vm->living_threads, th->self);
00572         if (rb_thread_alone()) {
00573             /* I'm last thread. wake up main thread from rb_thread_terminate_all */
00574             rb_threadptr_interrupt(main_th);
00575         }
00576 
00577         /* wake up joining threads */
00578         join_list = th->join_list;
00579         while (join_list) {
00580             rb_threadptr_interrupt(join_list->th);
00581             switch (join_list->th->status) {
00582               case THREAD_STOPPED: case THREAD_STOPPED_FOREVER:
00583                 join_list->th->status = THREAD_RUNNABLE;
00584               default: break;
00585             }
00586             join_list = join_list->next;
00587         }
00588 
00589         rb_threadptr_unlock_all_locking_mutexes(th);
00590         rb_check_deadlock(th->vm);
00591 
00592         if (!th->root_fiber) {
00593             rb_thread_recycle_stack_release(th->stack);
00594             th->stack = 0;
00595         }
00596     }
00597     native_mutex_lock(&th->vm->thread_destruct_lock);
00598     /* make sure vm->running_thread never point me after this point.*/
00599     th->vm->running_thread = NULL;
00600     native_mutex_unlock(&th->vm->thread_destruct_lock);
00601     thread_cleanup_func(th, FALSE);
00602     gvl_release(th->vm);
00603 
00604     return 0;
00605 }
00606 
00607 static VALUE
00608 thread_create_core(VALUE thval, VALUE args, VALUE (*fn)(ANYARGS))
00609 {
00610     rb_thread_t *th, *current_th = GET_THREAD();
00611     int err;
00612 
00613     if (OBJ_FROZEN(GET_THREAD()->thgroup)) {
00614         rb_raise(rb_eThreadError,
00615                  "can't start a new thread (frozen ThreadGroup)");
00616     }
00617     GetThreadPtr(thval, th);
00618 
00619     /* setup thread environment */
00620     th->first_func = fn;
00621     th->first_proc = fn ? Qfalse : rb_block_proc();
00622     th->first_args = args; /* GC: shouldn't put before above line */
00623 
00624     th->priority = current_th->priority;
00625     th->thgroup = current_th->thgroup;
00626 
00627     th->pending_interrupt_queue = rb_ary_tmp_new(0);
00628     th->pending_interrupt_queue_checked = 0;
00629     th->pending_interrupt_mask_stack = rb_ary_dup(current_th->pending_interrupt_mask_stack);
00630     RBASIC(th->pending_interrupt_mask_stack)->klass = 0;
00631 
00632     th->interrupt_mask = 0;
00633 
00634     native_mutex_initialize(&th->interrupt_lock);
00635 
00636     /* kick thread */
00637     err = native_thread_create(th);
00638     if (err) {
00639         th->status = THREAD_KILLED;
00640         rb_raise(rb_eThreadError, "can't create Thread (%d)", err);
00641     }
00642     st_insert(th->vm->living_threads, thval, (st_data_t) th->thread_id);
00643     return thval;
00644 }
00645 
00646 /*
00647  * call-seq:
00648  *  Thread.new { ... }                  -> thread
00649  *  Thread.new(*args, &proc)            -> thread
00650  *  Thread.new(*args) { |args| ... }    -> thread
00651  *
00652  *  Creates a new thread executing the given block.
00653  *
00654  *  Any +args+ given to ::new will be passed to the block:
00655  *
00656  *      arr = []
00657  *      a, b, c = 1, 2, 3
00658  *      Thread.new(a,b,c) { |d,e,f| arr << d << e << f }.join
00659  *      arr #=> [1, 2, 3]
00660  *
00661  *  A ThreadError exception is raised if ::new is called without a block.
00662  *
00663  *  If you're going to subclass Thread, be sure to call super in your
00664  *  +initialize+ method, otherwise a ThreadError will be raised.
00665  */
00666 static VALUE
00667 thread_s_new(int argc, VALUE *argv, VALUE klass)
00668 {
00669     rb_thread_t *th;
00670     VALUE thread = rb_thread_alloc(klass);
00671 
00672     if (GET_VM()->main_thread->status == THREAD_KILLED)
00673         rb_raise(rb_eThreadError, "can't alloc thread");
00674 
00675     rb_obj_call_init(thread, argc, argv);
00676     GetThreadPtr(thread, th);
00677     if (!th->first_args) {
00678         rb_raise(rb_eThreadError, "uninitialized thread - check `%s#initialize'",
00679                  rb_class2name(klass));
00680     }
00681     return thread;
00682 }
00683 
00684 /*
00685  *  call-seq:
00686  *     Thread.start([args]*) {|args| block }   -> thread
00687  *     Thread.fork([args]*) {|args| block }    -> thread
00688  *
00689  *  Basically the same as ::new. However, if class Thread is subclassed, then
00690  *  calling +start+ in that subclass will not invoke the subclass's
00691  *  +initialize+ method.
00692  */
00693 
00694 static VALUE
00695 thread_start(VALUE klass, VALUE args)
00696 {
00697     return thread_create_core(rb_thread_alloc(klass), args, 0);
00698 }
00699 
00700 /* :nodoc: */
00701 static VALUE
00702 thread_initialize(VALUE thread, VALUE args)
00703 {
00704     rb_thread_t *th;
00705     if (!rb_block_given_p()) {
00706         rb_raise(rb_eThreadError, "must be called with a block");
00707     }
00708     GetThreadPtr(thread, th);
00709     if (th->first_args) {
00710         VALUE proc = th->first_proc, line, loc;
00711         const char *file;
00712         if (!proc || !RTEST(loc = rb_proc_location(proc))) {
00713             rb_raise(rb_eThreadError, "already initialized thread");
00714         }
00715         file = RSTRING_PTR(RARRAY_PTR(loc)[0]);
00716         if (NIL_P(line = RARRAY_PTR(loc)[1])) {
00717             rb_raise(rb_eThreadError, "already initialized thread - %s",
00718                      file);
00719         }
00720         rb_raise(rb_eThreadError, "already initialized thread - %s:%d",
00721                  file, NUM2INT(line));
00722     }
00723     return thread_create_core(thread, args, 0);
00724 }
00725 
00726 VALUE
00727 rb_thread_create(VALUE (*fn)(ANYARGS), void *arg)
00728 {
00729     return thread_create_core(rb_thread_alloc(rb_cThread), (VALUE)arg, fn);
00730 }
00731 
00732 
00733 /* +infty, for this purpose */
00734 #define DELAY_INFTY 1E30
00735 
00736 struct join_arg {
00737     rb_thread_t *target, *waiting;
00738     double limit;
00739     int forever;
00740 };
00741 
00742 static VALUE
00743 remove_from_join_list(VALUE arg)
00744 {
00745     struct join_arg *p = (struct join_arg *)arg;
00746     rb_thread_t *target_th = p->target, *th = p->waiting;
00747 
00748     if (target_th->status != THREAD_KILLED) {
00749         rb_thread_list_t **p = &target_th->join_list;
00750 
00751         while (*p) {
00752             if ((*p)->th == th) {
00753                 *p = (*p)->next;
00754                 break;
00755             }
00756             p = &(*p)->next;
00757         }
00758     }
00759 
00760     return Qnil;
00761 }
00762 
00763 static VALUE
00764 thread_join_sleep(VALUE arg)
00765 {
00766     struct join_arg *p = (struct join_arg *)arg;
00767     rb_thread_t *target_th = p->target, *th = p->waiting;
00768     double now, limit = p->limit;
00769 
00770     while (target_th->status != THREAD_KILLED) {
00771         if (p->forever) {
00772             sleep_forever(th, 1, 0);
00773         }
00774         else {
00775             now = timeofday();
00776             if (now > limit) {
00777                 thread_debug("thread_join: timeout (thid: %p)\n",
00778                              (void *)target_th->thread_id);
00779                 return Qfalse;
00780             }
00781             sleep_wait_for_interrupt(th, limit - now, 0);
00782         }
00783         thread_debug("thread_join: interrupted (thid: %p)\n",
00784                      (void *)target_th->thread_id);
00785     }
00786     return Qtrue;
00787 }
00788 
00789 static VALUE
00790 thread_join(rb_thread_t *target_th, double delay)
00791 {
00792     rb_thread_t *th = GET_THREAD();
00793     struct join_arg arg;
00794 
00795     if (th == target_th) {
00796         rb_raise(rb_eThreadError, "Target thread must not be current thread");
00797     }
00798     if (GET_VM()->main_thread == target_th) {
00799         rb_raise(rb_eThreadError, "Target thread must not be main thread");
00800     }
00801 
00802     arg.target = target_th;
00803     arg.waiting = th;
00804     arg.limit = timeofday() + delay;
00805     arg.forever = delay == DELAY_INFTY;
00806 
00807     thread_debug("thread_join (thid: %p)\n", (void *)target_th->thread_id);
00808 
00809     if (target_th->status != THREAD_KILLED) {
00810         rb_thread_list_t list;
00811         list.next = target_th->join_list;
00812         list.th = th;
00813         target_th->join_list = &list;
00814         if (!rb_ensure(thread_join_sleep, (VALUE)&arg,
00815                        remove_from_join_list, (VALUE)&arg)) {
00816             return Qnil;
00817         }
00818     }
00819 
00820     thread_debug("thread_join: success (thid: %p)\n",
00821                  (void *)target_th->thread_id);
00822 
00823     if (target_th->errinfo != Qnil) {
00824         VALUE err = target_th->errinfo;
00825 
00826         if (FIXNUM_P(err)) {
00827             /* */
00828         }
00829         else if (RB_TYPE_P(target_th->errinfo, T_NODE)) {
00830             rb_exc_raise(rb_vm_make_jump_tag_but_local_jump(
00831                 GET_THROWOBJ_STATE(err), GET_THROWOBJ_VAL(err)));
00832         }
00833         else {
00834             /* normal exception */
00835             rb_exc_raise(err);
00836         }
00837     }
00838     return target_th->self;
00839 }
00840 
00841 /*
00842  *  call-seq:
00843  *     thr.join          -> thr
00844  *     thr.join(limit)   -> thr
00845  *
00846  *  The calling thread will suspend execution and run <i>thr</i>. Does not
00847  *  return until <i>thr</i> exits or until <i>limit</i> seconds have passed. If
00848  *  the time limit expires, <code>nil</code> will be returned, otherwise
00849  *  <i>thr</i> is returned.
00850  *
00851  *  Any threads not joined will be killed when the main program exits.  If
00852  *  <i>thr</i> had previously raised an exception and the
00853  *  <code>abort_on_exception</code> and <code>$DEBUG</code> flags are not set
00854  *  (so the exception has not yet been processed) it will be processed at this
00855  *  time.
00856  *
00857  *     a = Thread.new { print "a"; sleep(10); print "b"; print "c" }
00858  *     x = Thread.new { print "x"; Thread.pass; print "y"; print "z" }
00859  *     x.join # Let x thread finish, a will be killed on exit.
00860  *
00861  *  <em>produces:</em>
00862  *
00863  *     axyz
00864  *
00865  *  The following example illustrates the <i>limit</i> parameter.
00866  *
00867  *     y = Thread.new { 4.times { sleep 0.1; puts 'tick... ' }}
00868  *     puts "Waiting" until y.join(0.15)
00869  *
00870  *  <em>produces:</em>
00871  *
00872  *     tick...
00873  *     Waiting
00874  *     tick...
00875  *     Waitingtick...
00876  *
00877  *
00878  *     tick...
00879  */
00880 
00881 static VALUE
00882 thread_join_m(int argc, VALUE *argv, VALUE self)
00883 {
00884     rb_thread_t *target_th;
00885     double delay = DELAY_INFTY;
00886     VALUE limit;
00887 
00888     GetThreadPtr(self, target_th);
00889 
00890     rb_scan_args(argc, argv, "01", &limit);
00891     if (!NIL_P(limit)) {
00892         delay = rb_num2dbl(limit);
00893     }
00894 
00895     return thread_join(target_th, delay);
00896 }
00897 
00898 /*
00899  *  call-seq:
00900  *     thr.value   -> obj
00901  *
00902  *  Waits for <i>thr</i> to complete (via <code>Thread#join</code>) and returns
00903  *  its value.
00904  *
00905  *     a = Thread.new { 2 + 2 }
00906  *     a.value   #=> 4
00907  */
00908 
00909 static VALUE
00910 thread_value(VALUE self)
00911 {
00912     rb_thread_t *th;
00913     GetThreadPtr(self, th);
00914     thread_join(th, DELAY_INFTY);
00915     return th->value;
00916 }
00917 
00918 /*
00919  * Thread Scheduling
00920  */
00921 
00922 static struct timeval
00923 double2timeval(double d)
00924 {
00925     struct timeval time;
00926 
00927     if (isinf(d)) {
00928         time.tv_sec = TIMET_MAX;
00929         time.tv_usec = 0;
00930         return time;
00931     }
00932 
00933     time.tv_sec = (int)d;
00934     time.tv_usec = (int)((d - (int)d) * 1e6);
00935     if (time.tv_usec < 0) {
00936         time.tv_usec += (int)1e6;
00937         time.tv_sec -= 1;
00938     }
00939     return time;
00940 }
00941 
00942 static void
00943 sleep_forever(rb_thread_t *th, int deadlockable, int spurious_check)
00944 {
00945     enum rb_thread_status prev_status = th->status;
00946     enum rb_thread_status status = deadlockable ? THREAD_STOPPED_FOREVER : THREAD_STOPPED;
00947 
00948     th->status = status;
00949     RUBY_VM_CHECK_INTS_BLOCKING(th);
00950     while (th->status == status) {
00951         if (deadlockable) {
00952             th->vm->sleeper++;
00953             rb_check_deadlock(th->vm);
00954         }
00955         native_sleep(th, 0);
00956         if (deadlockable) {
00957             th->vm->sleeper--;
00958         }
00959         RUBY_VM_CHECK_INTS_BLOCKING(th);
00960         if (!spurious_check)
00961             break;
00962     }
00963     th->status = prev_status;
00964 }
00965 
00966 static void
00967 getclockofday(struct timeval *tp)
00968 {
00969 #if defined(HAVE_CLOCK_GETTIME) && defined(CLOCK_MONOTONIC)
00970     struct timespec ts;
00971 
00972     if (clock_gettime(CLOCK_MONOTONIC, &ts) == 0) {
00973         tp->tv_sec = ts.tv_sec;
00974         tp->tv_usec = ts.tv_nsec / 1000;
00975     } else
00976 #endif
00977     {
00978         gettimeofday(tp, NULL);
00979     }
00980 }
00981 
00982 static void
00983 sleep_timeval(rb_thread_t *th, struct timeval tv, int spurious_check)
00984 {
00985     struct timeval to, tvn;
00986     enum rb_thread_status prev_status = th->status;
00987 
00988     getclockofday(&to);
00989     if (TIMET_MAX - tv.tv_sec < to.tv_sec)
00990         to.tv_sec = TIMET_MAX;
00991     else
00992         to.tv_sec += tv.tv_sec;
00993     if ((to.tv_usec += tv.tv_usec) >= 1000000) {
00994         if (to.tv_sec == TIMET_MAX)
00995             to.tv_usec = 999999;
00996         else {
00997             to.tv_sec++;
00998             to.tv_usec -= 1000000;
00999         }
01000     }
01001 
01002     th->status = THREAD_STOPPED;
01003     RUBY_VM_CHECK_INTS_BLOCKING(th);
01004     while (th->status == THREAD_STOPPED) {
01005         native_sleep(th, &tv);
01006         RUBY_VM_CHECK_INTS_BLOCKING(th);
01007         getclockofday(&tvn);
01008         if (to.tv_sec < tvn.tv_sec) break;
01009         if (to.tv_sec == tvn.tv_sec && to.tv_usec <= tvn.tv_usec) break;
01010         thread_debug("sleep_timeval: %ld.%.6ld > %ld.%.6ld\n",
01011                      (long)to.tv_sec, (long)to.tv_usec,
01012                      (long)tvn.tv_sec, (long)tvn.tv_usec);
01013         tv.tv_sec = to.tv_sec - tvn.tv_sec;
01014         if ((tv.tv_usec = to.tv_usec - tvn.tv_usec) < 0) {
01015             --tv.tv_sec;
01016             tv.tv_usec += 1000000;
01017         }
01018         if (!spurious_check)
01019             break;
01020     }
01021     th->status = prev_status;
01022 }
01023 
01024 void
01025 rb_thread_sleep_forever(void)
01026 {
01027     thread_debug("rb_thread_sleep_forever\n");
01028     sleep_forever(GET_THREAD(), 0, 1);
01029 }
01030 
01031 static void
01032 rb_thread_sleep_deadly(void)
01033 {
01034     thread_debug("rb_thread_sleep_deadly\n");
01035     sleep_forever(GET_THREAD(), 1, 1);
01036 }
01037 
01038 static double
01039 timeofday(void)
01040 {
01041 #if defined(HAVE_CLOCK_GETTIME) && defined(CLOCK_MONOTONIC)
01042     struct timespec tp;
01043 
01044     if (clock_gettime(CLOCK_MONOTONIC, &tp) == 0) {
01045         return (double)tp.tv_sec + (double)tp.tv_nsec * 1e-9;
01046     } else
01047 #endif
01048     {
01049         struct timeval tv;
01050         gettimeofday(&tv, NULL);
01051         return (double)tv.tv_sec + (double)tv.tv_usec * 1e-6;
01052     }
01053 }
01054 
01055 static void
01056 sleep_wait_for_interrupt(rb_thread_t *th, double sleepsec, int spurious_check)
01057 {
01058     sleep_timeval(th, double2timeval(sleepsec), spurious_check);
01059 }
01060 
01061 static void
01062 sleep_for_polling(rb_thread_t *th)
01063 {
01064     struct timeval time;
01065     time.tv_sec = 0;
01066     time.tv_usec = 100 * 1000;  /* 0.1 sec */
01067     sleep_timeval(th, time, 1);
01068 }
01069 
01070 void
01071 rb_thread_wait_for(struct timeval time)
01072 {
01073     rb_thread_t *th = GET_THREAD();
01074     sleep_timeval(th, time, 1);
01075 }
01076 
01077 void
01078 rb_thread_polling(void)
01079 {
01080     if (!rb_thread_alone()) {
01081         rb_thread_t *th = GET_THREAD();
01082         RUBY_VM_CHECK_INTS_BLOCKING(th);
01083         sleep_for_polling(th);
01084     }
01085 }
01086 
01087 /*
01088  * CAUTION: This function causes thread switching.
01089  *          rb_thread_check_ints() check ruby's interrupts.
01090  *          some interrupt needs thread switching/invoke handlers,
01091  *          and so on.
01092  */
01093 
01094 void
01095 rb_thread_check_ints(void)
01096 {
01097     RUBY_VM_CHECK_INTS_BLOCKING(GET_THREAD());
01098 }
01099 
01100 /*
01101  * Hidden API for tcl/tk wrapper.
01102  * There is no guarantee to perpetuate it.
01103  */
01104 int
01105 rb_thread_check_trap_pending(void)
01106 {
01107     return rb_signal_buff_size() != 0;
01108 }
01109 
01110 /* This function can be called in blocking region. */
01111 int
01112 rb_thread_interrupted(VALUE thval)
01113 {
01114     rb_thread_t *th;
01115     GetThreadPtr(thval, th);
01116     return (int)RUBY_VM_INTERRUPTED(th);
01117 }
01118 
01119 void
01120 rb_thread_sleep(int sec)
01121 {
01122     rb_thread_wait_for(rb_time_timeval(INT2FIX(sec)));
01123 }
01124 
01125 static void
01126 rb_thread_schedule_limits(unsigned long limits_us)
01127 {
01128     thread_debug("rb_thread_schedule\n");
01129     if (!rb_thread_alone()) {
01130         rb_thread_t *th = GET_THREAD();
01131 
01132         if (th->running_time_us >= limits_us) {
01133             thread_debug("rb_thread_schedule/switch start\n");
01134             RB_GC_SAVE_MACHINE_CONTEXT(th);
01135             gvl_yield(th->vm, th);
01136             rb_thread_set_current(th);
01137             thread_debug("rb_thread_schedule/switch done\n");
01138         }
01139     }
01140 }
01141 
01142 void
01143 rb_thread_schedule(void)
01144 {
01145     rb_thread_t *cur_th = GET_THREAD();
01146     rb_thread_schedule_limits(0);
01147 
01148     if (UNLIKELY(RUBY_VM_INTERRUPTED_ANY(cur_th))) {
01149         rb_threadptr_execute_interrupts(cur_th, 0);
01150     }
01151 }
01152 
01153 /* blocking region */
01154 
01155 static inline int
01156 blocking_region_begin(rb_thread_t *th, struct rb_blocking_region_buffer *region,
01157                       rb_unblock_function_t *ubf, void *arg, int fail_if_interrupted)
01158 {
01159     region->prev_status = th->status;
01160     if (set_unblock_function(th, ubf, arg, &region->oldubf, fail_if_interrupted)) {
01161         th->blocking_region_buffer = region;
01162         th->status = THREAD_STOPPED;
01163         thread_debug("enter blocking region (%p)\n", (void *)th);
01164         RB_GC_SAVE_MACHINE_CONTEXT(th);
01165         gvl_release(th->vm);
01166         return TRUE;
01167     }
01168     else {
01169         return FALSE;
01170     }
01171 }
01172 
01173 static inline void
01174 blocking_region_end(rb_thread_t *th, struct rb_blocking_region_buffer *region)
01175 {
01176     gvl_acquire(th->vm, th);
01177     rb_thread_set_current(th);
01178     thread_debug("leave blocking region (%p)\n", (void *)th);
01179     remove_signal_thread_list(th);
01180     th->blocking_region_buffer = 0;
01181     reset_unblock_function(th, &region->oldubf);
01182     if (th->status == THREAD_STOPPED) {
01183         th->status = region->prev_status;
01184     }
01185 }
01186 
01187 struct rb_blocking_region_buffer *
01188 rb_thread_blocking_region_begin(void)
01189 {
01190     rb_thread_t *th = GET_THREAD();
01191     struct rb_blocking_region_buffer *region = ALLOC(struct rb_blocking_region_buffer);
01192     blocking_region_begin(th, region, ubf_select, th, FALSE);
01193     return region;
01194 }
01195 
01196 void
01197 rb_thread_blocking_region_end(struct rb_blocking_region_buffer *region)
01198 {
01199     int saved_errno = errno;
01200     rb_thread_t *th = ruby_thread_from_native();
01201     blocking_region_end(th, region);
01202     xfree(region);
01203     RUBY_VM_CHECK_INTS_BLOCKING(th);
01204     errno = saved_errno;
01205 }
01206 
01207 static void *
01208 call_without_gvl(void *(*func)(void *), void *data1,
01209                  rb_unblock_function_t *ubf, void *data2, int fail_if_interrupted)
01210 {
01211     void *val = 0;
01212 
01213     rb_thread_t *th = GET_THREAD();
01214     int saved_errno = 0;
01215 
01216     th->waiting_fd = -1;
01217     if (ubf == RUBY_UBF_IO || ubf == RUBY_UBF_PROCESS) {
01218         ubf = ubf_select;
01219         data2 = th;
01220     }
01221 
01222     BLOCKING_REGION({
01223         val = func(data1);
01224         saved_errno = errno;
01225     }, ubf, data2, fail_if_interrupted);
01226 
01227     if (!fail_if_interrupted) {
01228         RUBY_VM_CHECK_INTS_BLOCKING(th);
01229     }
01230 
01231     errno = saved_errno;
01232 
01233     return val;
01234 }
01235 
01236 /*
01237  * rb_thread_call_without_gvl - permit concurrent/parallel execution.
01238  * rb_thread_call_without_gvl2 - permit concurrent/parallel execution
01239  *                               without interrupt proceess.
01240  *
01241  * rb_thread_call_without_gvl() does:
01242  *   (1) Check interrupts.
01243  *   (2) release GVL.
01244  *       Other Ruby threads may run in parallel.
01245  *   (3) call func with data1
01246  *   (4) acquire GVL.
01247  *       Other Ruby threads can not run in parallel any more.
01248  *   (5) Check interrupts.
01249  *
01250  * rb_thread_call_without_gvl2() does:
01251  *   (1) Check interrupt and return if interrupted.
01252  *   (2) release GVL.
01253  *   (3) call func with data1 and a pointer to the flags.
01254  *   (4) acquire GVL.
01255  *
01256  * If another thread interrupts this thread (Thread#kill, signal delivery,
01257  * VM-shutdown request, and so on), `ubf()' is called (`ubf()' means
01258  * "un-blocking function").  `ubf()' should interrupt `func()' execution by
01259  * toggling a cancellation flag, canceling the invocation of a call inside
01260  * `func()' or similar.  Note that `ubf()' may not be called with the GVL.
01261  *
01262  * There are built-in ubfs and you can specify these ubfs:
01263  *
01264  * * RUBY_UBF_IO: ubf for IO operation
01265  * * RUBY_UBF_PROCESS: ubf for process operation
01266  *
01267  * However, we can not guarantee our built-in ubfs interrupt your `func()'
01268  * correctly. Be careful to use rb_thread_call_without_gvl(). If you don't
01269  * provide proper ubf(), your program will not stop for Control+C or other
01270  * shutdown events.
01271  *
01272  * "Check interrupts" on above list means that check asynchronous
01273  * interrupt events (such as Thread#kill, signal delivery, VM-shutdown
01274  * request, and so on) and call corresponding procedures
01275  * (such as `trap' for signals, raise an exception for Thread#raise).
01276  * If `func()' finished and receive interrupts, you may skip interrupt
01277  * checking.  For example, assume the following func() it read data from file.
01278  *
01279  *   read_func(...) {
01280  *                   // (a) before read
01281  *     read(buffer); // (b) reading
01282  *                   // (c) after read
01283  *   }
01284  *
01285  * If an interrupt occurs at (a) or (b), then `ubf()' cancels this
01286  * `read_func()' and interrupts are checked. However, if an interrupt occurs
01287  * at (c), after *read* operation is completed, check intterrupts is harmful
01288  * because it causes irrevocable side-effect, the read data will vanish.  To
01289  * avoid such problem, the `read_func()' should be used with
01290  * `rb_thread_call_without_gvl2()'.
01291  *
01292  * If `rb_thread_call_without_gvl2()' detects interrupt, return its execution
01293  * immediately. This function does not show when the execution was interrupted.
01294  * For example, there are 4 possible timing (a), (b), (c) and before calling
01295  * read_func(). You need to record progress of a read_func() and check
01296  * the progress after `rb_thread_call_without_gvl2()'. You may need to call
01297  * `rb_thread_check_ints()' correctly or your program can not process proper
01298  * process such as `trap' and so on.
01299  *
01300  * NOTE: You can not execute most of Ruby C API and touch Ruby
01301  *       objects in `func()' and `ubf()', including raising an
01302  *       exception, because current thread doesn't acquire GVL
01303  *       (it causes synchronization problems).  If you need to
01304  *       call ruby functions either use rb_thread_call_with_gvl()
01305  *       or read source code of C APIs and confirm safety by
01306  *       yourself.
01307  *
01308  * NOTE: In short, this API is difficult to use safely.  I recommend you
01309  *       use other ways if you have.  We lack experiences to use this API.
01310  *       Please report your problem related on it.
01311  *
01312  * NOTE: Releasing GVL and re-acquiring GVL may be expensive operations
01313  *       for a short running `func()'. Be sure to benchmark and use this
01314  *       mechanism when `func()' consumes enough time.
01315  *
01316  * Safe C API:
01317  * * rb_thread_interrupted() - check interrupt flag
01318  * * ruby_xmalloc(), ruby_xrealloc(), ruby_xfree() -
01319  *   they will work without GVL, and may acquire GVL when GC is needed.
01320  */
01321 void *
01322 rb_thread_call_without_gvl2(void *(*func)(void *), void *data1,
01323                             rb_unblock_function_t *ubf, void *data2)
01324 {
01325     return call_without_gvl(func, data1, ubf, data2, TRUE);
01326 }
01327 
01328 void *
01329 rb_thread_call_without_gvl(void *(*func)(void *data), void *data1,
01330                             rb_unblock_function_t *ubf, void *data2)
01331 {
01332     return call_without_gvl(func, data1, ubf, data2, FALSE);
01333 }
01334 
01335 VALUE
01336 rb_thread_io_blocking_region(rb_blocking_function_t *func, void *data1, int fd)
01337 {
01338     VALUE val = Qundef; /* shouldn't be used */
01339     rb_thread_t *th = GET_THREAD();
01340     int saved_errno = 0;
01341     int state;
01342 
01343     th->waiting_fd = fd;
01344 
01345     TH_PUSH_TAG(th);
01346     if ((state = EXEC_TAG()) == 0) {
01347         BLOCKING_REGION({
01348             val = func(data1);
01349             saved_errno = errno;
01350         }, ubf_select, th, FALSE);
01351     }
01352     TH_POP_TAG();
01353 
01354     /* clear waitinf_fd anytime */
01355     th->waiting_fd = -1;
01356 
01357     if (state) {
01358         JUMP_TAG(state);
01359     }
01360     /* TODO: check func() */
01361     RUBY_VM_CHECK_INTS_BLOCKING(th);
01362 
01363     errno = saved_errno;
01364 
01365     return val;
01366 }
01367 
01368 VALUE
01369 rb_thread_blocking_region(
01370     rb_blocking_function_t *func, void *data1,
01371     rb_unblock_function_t *ubf, void *data2)
01372 {
01373     void *(*f)(void*) = (void *(*)(void*))func;
01374     return (VALUE)rb_thread_call_without_gvl(f, data1, ubf, data2);
01375 }
01376 
01377 /*
01378  * rb_thread_call_with_gvl - re-enter the Ruby world after GVL release.
01379  *
01380  * After releasing GVL using rb_thread_blocking_region() or
01381  * rb_thread_call_without_gvl() you can not access Ruby values or invoke
01382  * methods. If you need to access Ruby you must use this function
01383  * rb_thread_call_with_gvl().
01384  *
01385  * This function rb_thread_call_with_gvl() does:
01386  * (1) acquire GVL.
01387  * (2) call passed function `func'.
01388  * (3) release GVL.
01389  * (4) return a value which is returned at (2).
01390  *
01391  * NOTE: You should not return Ruby object at (2) because such Object
01392  *       will not marked.
01393  *
01394  * NOTE: If an exception is raised in `func', this function DOES NOT
01395  *       protect (catch) the exception.  If you have any resources
01396  *       which should free before throwing exception, you need use
01397  *       rb_protect() in `func' and return a value which represents
01398  *       exception is raised.
01399  *
01400  * NOTE: This function should not be called by a thread which was not
01401  *       created as Ruby thread (created by Thread.new or so).  In other
01402  *       words, this function *DOES NOT* associate or convert a NON-Ruby
01403  *       thread to a Ruby thread.
01404  */
01405 void *
01406 rb_thread_call_with_gvl(void *(*func)(void *), void *data1)
01407 {
01408     rb_thread_t *th = ruby_thread_from_native();
01409     struct rb_blocking_region_buffer *brb;
01410     struct rb_unblock_callback prev_unblock;
01411     void *r;
01412 
01413     if (th == 0) {
01414         /* Error is occurred, but we can't use rb_bug()
01415          * because this thread is not Ruby's thread.
01416          * What should we do?
01417          */
01418 
01419         fprintf(stderr, "[BUG] rb_thread_call_with_gvl() is called by non-ruby thread\n");
01420         exit(EXIT_FAILURE);
01421     }
01422 
01423     brb = (struct rb_blocking_region_buffer *)th->blocking_region_buffer;
01424     prev_unblock = th->unblock;
01425 
01426     if (brb == 0) {
01427         rb_bug("rb_thread_call_with_gvl: called by a thread which has GVL.");
01428     }
01429 
01430     blocking_region_end(th, brb);
01431     /* enter to Ruby world: You can access Ruby values, methods and so on. */
01432     r = (*func)(data1);
01433     /* leave from Ruby world: You can not access Ruby values, etc. */
01434     blocking_region_begin(th, brb, prev_unblock.func, prev_unblock.arg, FALSE);
01435     return r;
01436 }
01437 
01438 /*
01439  * ruby_thread_has_gvl_p - check if current native thread has GVL.
01440  *
01441  ***
01442  *** This API is EXPERIMENTAL!
01443  *** We do not guarantee that this API remains in ruby 1.9.2 or later.
01444  ***
01445  */
01446 
01447 int
01448 ruby_thread_has_gvl_p(void)
01449 {
01450     rb_thread_t *th = ruby_thread_from_native();
01451 
01452     if (th && th->blocking_region_buffer == 0) {
01453         return 1;
01454     }
01455     else {
01456         return 0;
01457     }
01458 }
01459 
01460 /*
01461  * call-seq:
01462  *    Thread.pass   -> nil
01463  *
01464  * Give the thread scheduler a hint to pass execution to another thread.
01465  * A running thread may or may not switch, it depends on OS and processor.
01466  */
01467 
01468 static VALUE
01469 thread_s_pass(VALUE klass)
01470 {
01471     rb_thread_schedule();
01472     return Qnil;
01473 }
01474 
01475 /*****************************************************/
01476 
01477 /*
01478  * rb_threadptr_pending_interrupt_* - manage asynchronous error queue
01479  *
01480  * Async events such as an exception throwed by Thread#raise,
01481  * Thread#kill and thread termination (after main thread termination)
01482  * will be queued to th->pending_interrupt_queue.
01483  * - clear: clear the queue.
01484  * - enque: enque err object into queue.
01485  * - deque: deque err object from queue.
01486  * - active_p: return 1 if the queue should be checked.
01487  *
01488  * All rb_threadptr_pending_interrupt_* functions are called by
01489  * a GVL acquired thread, of course.
01490  * Note that all "rb_" prefix APIs need GVL to call.
01491  */
01492 
01493 void
01494 rb_threadptr_pending_interrupt_clear(rb_thread_t *th)
01495 {
01496     rb_ary_clear(th->pending_interrupt_queue);
01497 }
01498 
01499 void
01500 rb_threadptr_pending_interrupt_enque(rb_thread_t *th, VALUE v)
01501 {
01502     rb_ary_push(th->pending_interrupt_queue, v);
01503     th->pending_interrupt_queue_checked = 0;
01504 }
01505 
01506 enum handle_interrupt_timing {
01507     INTERRUPT_NONE,
01508     INTERRUPT_IMMEDIATE,
01509     INTERRUPT_ON_BLOCKING,
01510     INTERRUPT_NEVER
01511 };
01512 
01513 static enum handle_interrupt_timing
01514 rb_threadptr_pending_interrupt_check_mask(rb_thread_t *th, VALUE err)
01515 {
01516     VALUE mask;
01517     long mask_stack_len = RARRAY_LEN(th->pending_interrupt_mask_stack);
01518     VALUE *mask_stack = RARRAY_PTR(th->pending_interrupt_mask_stack);
01519     VALUE ancestors = rb_mod_ancestors(err); /* TODO: GC guard */
01520     long ancestors_len = RARRAY_LEN(ancestors);
01521     VALUE *ancestors_ptr = RARRAY_PTR(ancestors);
01522     int i, j;
01523 
01524     for (i=0; i<mask_stack_len; i++) {
01525         mask = mask_stack[mask_stack_len-(i+1)];
01526 
01527         for (j=0; j<ancestors_len; j++) {
01528             VALUE klass = ancestors_ptr[j];
01529             VALUE sym;
01530 
01531             /* TODO: remove rb_intern() */
01532             if ((sym = rb_hash_aref(mask, klass)) != Qnil) {
01533                 if (sym == sym_immediate) {
01534                     return INTERRUPT_IMMEDIATE;
01535                 }
01536                 else if (sym == sym_on_blocking) {
01537                     return INTERRUPT_ON_BLOCKING;
01538                 }
01539                 else if (sym == sym_never) {
01540                     return INTERRUPT_NEVER;
01541                 }
01542                 else {
01543                     rb_raise(rb_eThreadError, "unknown mask signature");
01544                 }
01545             }
01546         }
01547         /* try to next mask */
01548     }
01549     return INTERRUPT_NONE;
01550 }
01551 
01552 static int
01553 rb_threadptr_pending_interrupt_empty_p(rb_thread_t *th)
01554 {
01555     return RARRAY_LEN(th->pending_interrupt_queue) == 0;
01556 }
01557 
01558 static int
01559 rb_threadptr_pending_interrupt_include_p(rb_thread_t *th, VALUE err)
01560 {
01561     int i;
01562     for (i=0; i<RARRAY_LEN(th->pending_interrupt_queue); i++) {
01563         VALUE e = RARRAY_PTR(th->pending_interrupt_queue)[i];
01564         if (rb_class_inherited_p(e, err)) {
01565             return TRUE;
01566         }
01567     }
01568     return FALSE;
01569 }
01570 
01571 static VALUE
01572 rb_threadptr_pending_interrupt_deque(rb_thread_t *th, enum handle_interrupt_timing timing)
01573 {
01574 #if 1 /* 1 to enable Thread#handle_interrupt, 0 to ignore it */
01575     int i;
01576 
01577     for (i=0; i<RARRAY_LEN(th->pending_interrupt_queue); i++) {
01578         VALUE err = RARRAY_PTR(th->pending_interrupt_queue)[i];
01579 
01580         enum handle_interrupt_timing mask_timing = rb_threadptr_pending_interrupt_check_mask(th, CLASS_OF(err));
01581 
01582         switch (mask_timing) {
01583           case INTERRUPT_ON_BLOCKING:
01584             if (timing != INTERRUPT_ON_BLOCKING) {
01585                 break;
01586             }
01587             /* fall through */
01588           case INTERRUPT_NONE: /* default: IMMEDIATE */
01589           case INTERRUPT_IMMEDIATE:
01590             rb_ary_delete_at(th->pending_interrupt_queue, i);
01591             return err;
01592           case INTERRUPT_NEVER:
01593             break;
01594         }
01595     }
01596 
01597     th->pending_interrupt_queue_checked = 1;
01598     return Qundef;
01599 #else
01600     VALUE err = rb_ary_shift(th->pending_interrupt_queue);
01601     if (rb_threadptr_pending_interrupt_empty_p(th)) {
01602         th->pending_interrupt_queue_checked = 1;
01603     }
01604     return err;
01605 #endif
01606 }
01607 
01608 int
01609 rb_threadptr_pending_interrupt_active_p(rb_thread_t *th)
01610 {
01611     /*
01612      * For optimization, we don't check async errinfo queue
01613      * if it nor a thread interrupt mask were not changed
01614      * since last check.
01615      */
01616     if (th->pending_interrupt_queue_checked) {
01617         return 0;
01618     }
01619 
01620     if (rb_threadptr_pending_interrupt_empty_p(th)) {
01621         return 0;
01622     }
01623 
01624     return 1;
01625 }
01626 
01627 static int
01628 handle_interrupt_arg_check_i(VALUE key, VALUE val)
01629 {
01630     if (val != sym_immediate && val != sym_on_blocking && val != sym_never) {
01631         rb_raise(rb_eArgError, "unknown mask signature");
01632     }
01633 
01634     return ST_CONTINUE;
01635 }
01636 
01637 /*
01638  * call-seq:
01639  *   Thread.handle_interrupt(hash) { ... } -> result of the block
01640  *
01641  * Changes asynchronous interrupt timing.
01642  *
01643  * _interrupt_ means asynchronous event and corresponding procedure
01644  * by Thread#raise, Thread#kill, signal trap (not supported yet)
01645  * and main thread termination (if main thread terminates, then all
01646  * other thread will be killed).
01647  *
01648  * The given +hash+ has pairs like <code>ExceptionClass =>
01649  * :TimingSymbol</code>. Where the ExceptionClass is the interrupt handled by
01650  * the given block. The TimingSymbol can be one of the following symbols:
01651  *
01652  * [+:immediate+]   Invoke interrupts immediately.
01653  * [+:on_blocking+] Invoke interrupts while _BlockingOperation_.
01654  * [+:never+]       Never invoke all interrupts.
01655  *
01656  * _BlockingOperation_ means that the operation will block the calling thread,
01657  * such as read and write.  On CRuby implementation, _BlockingOperation_ is any
01658  * operation executed without GVL.
01659  *
01660  * Masked asynchronous interrupts are delayed until they are enabled.
01661  * This method is similar to sigprocmask(3).
01662  *
01663  * === NOTE
01664  *
01665  * Asynchronous interrupts are difficult to use.
01666  *
01667  * If you need to communicate between threads, please consider to use another way such as Queue.
01668  *
01669  * Or use them with deep understanding about this method.
01670  *
01671  * === Usage
01672  *
01673  * In this example, we can guard from Thread#raise exceptions.
01674  *
01675  * Using the +:never+ TimingSymbol the RuntimeError exception will always be
01676  * ignored in the first block of the main thread. In the second
01677  * ::handle_interrupt block we can purposefully handle RuntimeError exceptions.
01678  *
01679  *   th = Thread.new do
01680  *     Thead.handle_interrupt(RuntimeError => :never) {
01681  *       begin
01682  *         # You can write resource allocation code safely.
01683  *         Thread.handle_interrupt(RuntimeError => :immediate) {
01684  *           # ...
01685  *         }
01686  *       ensure
01687  *         # You can write resource deallocation code safely.
01688  *       end
01689  *     }
01690  *   end
01691  *   Thread.pass
01692  *   # ...
01693  *   th.raise "stop"
01694  *
01695  * While we are ignoring the RuntimeError exception, it's safe to write our
01696  * resource allocation code. Then, the ensure block is where we can safely
01697  * deallocate your resources.
01698  *
01699  * ==== Guarding from TimeoutError
01700  *
01701  * In the next example, we will guard from the TimeoutError exception. This
01702  * will help prevent from leaking resources when TimeoutError exceptions occur
01703  * during normal ensure clause. For this example we use the help of the
01704  * standard library Timeout, from lib/timeout.rb
01705  *
01706  *   require 'timeout'
01707  *   Thread.handle_interrupt(TimeoutError => :never) {
01708  *     timeout(10){
01709  *       # TimeoutError doesn't occur here
01710  *       Thread.handle_interrupt(TimeoutError => :on_blocking) {
01711  *         # possible to be killed by TimeoutError
01712  *         # while blocking operation
01713  *       }
01714  *       # TimeoutError doesn't occur here
01715  *     }
01716  *   }
01717  *
01718  * In the first part of the +timeout+ block, we can rely on TimeoutError being
01719  * ignored. Then in the <code>TimeoutError => :on_blocking</code> block, any
01720  * operation that will block the calling thread is susceptible to a
01721  * TimeoutError exception being raised.
01722  *
01723  * ==== Stack control settings
01724  *
01725  * It's possible to stack multiple levels of ::handle_interrupt blocks in order
01726  * to control more than one ExceptionClass and TimingSymbol at a time.
01727  *
01728  *   Thread.handle_interrupt(FooError => :never) {
01729  *     Thread.handle_interrupt(BarError => :never) {
01730  *        # FooError and BarError are prohibited.
01731  *     }
01732  *   }
01733  *
01734  * ==== Inheritance with ExceptionClass
01735  *
01736  * All exceptions inherited from the ExceptionClass parameter will be considered.
01737  *
01738  *   Thread.handle_interrupt(Exception => :never) {
01739  *     # all exceptions inherited from Exception are prohibited.
01740  *   }
01741  *
01742  */
01743 static VALUE
01744 rb_thread_s_handle_interrupt(VALUE self, VALUE mask_arg)
01745 {
01746     VALUE mask;
01747     rb_thread_t *th = GET_THREAD();
01748     VALUE r = Qnil;
01749     int state;
01750 
01751     if (!rb_block_given_p()) {
01752         rb_raise(rb_eArgError, "block is needed.");
01753     }
01754 
01755     mask = rb_convert_type(mask_arg, T_HASH, "Hash", "to_hash");
01756     rb_hash_foreach(mask, handle_interrupt_arg_check_i, 0);
01757     rb_ary_push(th->pending_interrupt_mask_stack, mask);
01758     if (!rb_threadptr_pending_interrupt_empty_p(th)) {
01759         th->pending_interrupt_queue_checked = 0;
01760         RUBY_VM_SET_INTERRUPT(th);
01761     }
01762 
01763     TH_PUSH_TAG(th);
01764     if ((state = EXEC_TAG()) == 0) {
01765         r = rb_yield(Qnil);
01766     }
01767     TH_POP_TAG();
01768 
01769     rb_ary_pop(th->pending_interrupt_mask_stack);
01770     if (!rb_threadptr_pending_interrupt_empty_p(th)) {
01771         th->pending_interrupt_queue_checked = 0;
01772         RUBY_VM_SET_INTERRUPT(th);
01773     }
01774 
01775     RUBY_VM_CHECK_INTS(th);
01776 
01777     if (state) {
01778         JUMP_TAG(state);
01779     }
01780 
01781     return r;
01782 }
01783 
01784 /*
01785  * call-seq:
01786  *   target_thread.pending_interrupt?(error = nil) -> true/false
01787  *
01788  * Returns whether or not the asychronous queue is empty for the target thread.
01789  *
01790  * If +error+ is given, then check only for +error+ type deferred events.
01791  *
01792  * See ::pending_interrupt? for more information.
01793  */
01794 static VALUE
01795 rb_thread_pending_interrupt_p(int argc, VALUE *argv, VALUE target_thread)
01796 {
01797     rb_thread_t *target_th;
01798 
01799     GetThreadPtr(target_thread, target_th);
01800 
01801     if (rb_threadptr_pending_interrupt_empty_p(target_th)) {
01802         return Qfalse;
01803     }
01804     else {
01805         if (argc == 1) {
01806             VALUE err;
01807             rb_scan_args(argc, argv, "01", &err);
01808             if (!rb_obj_is_kind_of(err, rb_cModule)) {
01809                 rb_raise(rb_eTypeError, "class or module required for rescue clause");
01810             }
01811             if (rb_threadptr_pending_interrupt_include_p(target_th, err)) {
01812                 return Qtrue;
01813             }
01814             else {
01815                 return Qfalse;
01816             }
01817         }
01818         return Qtrue;
01819     }
01820 }
01821 
01822 /*
01823  * call-seq:
01824  *   Thread.pending_interrupt?(error = nil) -> true/false
01825  *
01826  * Returns whether or not the asynchronous queue is empty.
01827  *
01828  * Since Thread::handle_interrupt can be used to defer asynchronous events.
01829  * This method can be used to determine if there are any deferred events.
01830  *
01831  * If you find this method returns true, then you may finish +:never+ blocks.
01832  *
01833  * For example, the following method processes deferred asynchronous events
01834  * immediately.
01835  *
01836  *   def Thread.kick_interrupt_immediately
01837  *     Thread.handle_interrupt(Object => :immediate) {
01838  *       Thread.pass
01839  *     }
01840  *   end
01841  *
01842  * If +error+ is given, then check only for +error+ type deferred events.
01843  *
01844  * === Usage
01845  *
01846  *   th = Thread.new{
01847  *     Thread.handle_interrupt(RuntimeError => :on_blocking){
01848  *       while true
01849  *         ...
01850  *         # reach safe point to invoke interrupt
01851  *         if Thread.pending_interrupt?
01852  *           Thread.handle_interrupt(Object => :immediate){}
01853  *         end
01854  *         ...
01855  *       end
01856  *     }
01857  *   }
01858  *   ...
01859  *   th.raise # stop thread
01860  *
01861  * This example can also be written as the following, which you should use to
01862  * avoid asynchronous interrupts.
01863  *
01864  *   flag = true
01865  *   th = Thread.new{
01866  *     Thread.handle_interrupt(RuntimeError => :on_blocking){
01867  *       while true
01868  *         ...
01869  *         # reach safe point to invoke interrupt
01870  *         break if flag == false
01871  *         ...
01872  *       end
01873  *     }
01874  *   }
01875  *   ...
01876  *   flag = false # stop thread
01877  */
01878 
01879 static VALUE
01880 rb_thread_s_pending_interrupt_p(int argc, VALUE *argv, VALUE self)
01881 {
01882     return rb_thread_pending_interrupt_p(argc, argv, GET_THREAD()->self);
01883 }
01884 
01885 static void
01886 rb_threadptr_to_kill(rb_thread_t *th)
01887 {
01888     rb_threadptr_pending_interrupt_clear(th);
01889     th->status = THREAD_RUNNABLE;
01890     th->to_kill = 1;
01891     th->errinfo = INT2FIX(TAG_FATAL);
01892     TH_JUMP_TAG(th, TAG_FATAL);
01893 }
01894 
01895 void
01896 rb_threadptr_execute_interrupts(rb_thread_t *th, int blocking_timing)
01897 {
01898     if (th->raised_flag) return;
01899 
01900     while (1) {
01901         rb_atomic_t interrupt;
01902         rb_atomic_t old;
01903         int sig;
01904         int timer_interrupt;
01905         int pending_interrupt;
01906         int finalizer_interrupt;
01907         int trap_interrupt;
01908 
01909         do {
01910             interrupt = th->interrupt_flag;
01911             old = ATOMIC_CAS(th->interrupt_flag, interrupt, interrupt & th->interrupt_mask);
01912         } while (old != interrupt);
01913 
01914         interrupt &= (rb_atomic_t)~th->interrupt_mask;
01915         if (!interrupt)
01916             return;
01917 
01918         timer_interrupt = interrupt & TIMER_INTERRUPT_MASK;
01919         pending_interrupt = interrupt & PENDING_INTERRUPT_MASK;
01920         finalizer_interrupt = interrupt & FINALIZER_INTERRUPT_MASK;
01921         trap_interrupt = interrupt & TRAP_INTERRUPT_MASK;
01922 
01923         /* signal handling */
01924         if (trap_interrupt && (th == th->vm->main_thread)) {
01925             enum rb_thread_status prev_status = th->status;
01926             th->status = THREAD_RUNNABLE;
01927             while ((sig = rb_get_next_signal()) != 0) {
01928                 rb_signal_exec(th, sig);
01929             }
01930             th->status = prev_status;
01931         }
01932 
01933         /* exception from another thread */
01934         if (pending_interrupt && rb_threadptr_pending_interrupt_active_p(th)) {
01935             VALUE err = rb_threadptr_pending_interrupt_deque(th, blocking_timing ? INTERRUPT_ON_BLOCKING : INTERRUPT_NONE);
01936             thread_debug("rb_thread_execute_interrupts: %"PRIdVALUE"\n", err);
01937 
01938             if (err == Qundef) {
01939                 /* no error */
01940             }
01941             else if (err == eKillSignal        /* Thread#kill receieved */  ||
01942                      err == eTerminateSignal   /* Terminate thread */       ||
01943                      err == INT2FIX(TAG_FATAL) /* Thread.exit etc. */         ) {
01944                 rb_threadptr_to_kill(th);
01945             }
01946             else {
01947                 /* set runnable if th was slept. */
01948                 if (th->status == THREAD_STOPPED ||
01949                     th->status == THREAD_STOPPED_FOREVER)
01950                     th->status = THREAD_RUNNABLE;
01951                 rb_exc_raise(err);
01952             }
01953         }
01954 
01955         if (finalizer_interrupt) {
01956             rb_gc_finalize_deferred();
01957         }
01958 
01959         if (timer_interrupt) {
01960             unsigned long limits_us = TIME_QUANTUM_USEC;
01961 
01962             if (th->priority > 0)
01963                 limits_us <<= th->priority;
01964             else
01965                 limits_us >>= -th->priority;
01966 
01967             if (th->status == THREAD_RUNNABLE)
01968                 th->running_time_us += TIME_QUANTUM_USEC;
01969 
01970             EXEC_EVENT_HOOK(th, RUBY_EVENT_SWITCH, th->cfp->self, 0, 0, Qundef);
01971 
01972             rb_thread_schedule_limits(limits_us);
01973         }
01974     }
01975 }
01976 
01977 void
01978 rb_thread_execute_interrupts(VALUE thval)
01979 {
01980     rb_thread_t *th;
01981     GetThreadPtr(thval, th);
01982     rb_threadptr_execute_interrupts(th, 1);
01983 }
01984 
01985 static void
01986 rb_threadptr_ready(rb_thread_t *th)
01987 {
01988     rb_threadptr_interrupt(th);
01989 }
01990 
01991 static VALUE
01992 rb_threadptr_raise(rb_thread_t *th, int argc, VALUE *argv)
01993 {
01994     VALUE exc;
01995 
01996     if (rb_threadptr_dead(th)) {
01997         return Qnil;
01998     }
01999 
02000     if (argc == 0) {
02001         exc = rb_exc_new(rb_eRuntimeError, 0, 0);
02002     }
02003     else {
02004         exc = rb_make_exception(argc, argv);
02005     }
02006     rb_threadptr_pending_interrupt_enque(th, exc);
02007     rb_threadptr_interrupt(th);
02008     return Qnil;
02009 }
02010 
02011 void
02012 rb_threadptr_signal_raise(rb_thread_t *th, int sig)
02013 {
02014     VALUE argv[2];
02015 
02016     argv[0] = rb_eSignal;
02017     argv[1] = INT2FIX(sig);
02018     rb_threadptr_raise(th->vm->main_thread, 2, argv);
02019 }
02020 
02021 void
02022 rb_threadptr_signal_exit(rb_thread_t *th)
02023 {
02024     VALUE argv[2];
02025 
02026     argv[0] = rb_eSystemExit;
02027     argv[1] = rb_str_new2("exit");
02028     rb_threadptr_raise(th->vm->main_thread, 2, argv);
02029 }
02030 
02031 #if defined(POSIX_SIGNAL) && defined(SIGSEGV) && defined(HAVE_SIGALTSTACK)
02032 #define USE_SIGALTSTACK
02033 #endif
02034 
02035 void
02036 ruby_thread_stack_overflow(rb_thread_t *th)
02037 {
02038     th->raised_flag = 0;
02039 #ifdef USE_SIGALTSTACK
02040     rb_exc_raise(sysstack_error);
02041 #else
02042     th->errinfo = sysstack_error;
02043     TH_JUMP_TAG(th, TAG_RAISE);
02044 #endif
02045 }
02046 
02047 int
02048 rb_threadptr_set_raised(rb_thread_t *th)
02049 {
02050     if (th->raised_flag & RAISED_EXCEPTION) {
02051         return 1;
02052     }
02053     th->raised_flag |= RAISED_EXCEPTION;
02054     return 0;
02055 }
02056 
02057 int
02058 rb_threadptr_reset_raised(rb_thread_t *th)
02059 {
02060     if (!(th->raised_flag & RAISED_EXCEPTION)) {
02061         return 0;
02062     }
02063     th->raised_flag &= ~RAISED_EXCEPTION;
02064     return 1;
02065 }
02066 
02067 static int
02068 thread_fd_close_i(st_data_t key, st_data_t val, st_data_t data)
02069 {
02070     int fd = (int)data;
02071     rb_thread_t *th;
02072     GetThreadPtr((VALUE)key, th);
02073 
02074     if (th->waiting_fd == fd) {
02075         VALUE err = th->vm->special_exceptions[ruby_error_closed_stream];
02076         rb_threadptr_pending_interrupt_enque(th, err);
02077         rb_threadptr_interrupt(th);
02078     }
02079     return ST_CONTINUE;
02080 }
02081 
02082 void
02083 rb_thread_fd_close(int fd)
02084 {
02085     st_foreach(GET_THREAD()->vm->living_threads, thread_fd_close_i, (st_index_t)fd);
02086 }
02087 
02088 /*
02089  *  call-seq:
02090  *     thr.raise
02091  *     thr.raise(string)
02092  *     thr.raise(exception [, string [, array]])
02093  *
02094  *  Raises an exception (see <code>Kernel::raise</code>) from <i>thr</i>. The
02095  *  caller does not have to be <i>thr</i>.
02096  *
02097  *     Thread.abort_on_exception = true
02098  *     a = Thread.new { sleep(200) }
02099  *     a.raise("Gotcha")
02100  *
02101  *  <em>produces:</em>
02102  *
02103  *     prog.rb:3: Gotcha (RuntimeError)
02104  *      from prog.rb:2:in `initialize'
02105  *      from prog.rb:2:in `new'
02106  *      from prog.rb:2
02107  */
02108 
02109 static VALUE
02110 thread_raise_m(int argc, VALUE *argv, VALUE self)
02111 {
02112     rb_thread_t *target_th;
02113     rb_thread_t *th = GET_THREAD();
02114     GetThreadPtr(self, target_th);
02115     rb_threadptr_raise(target_th, argc, argv);
02116 
02117     /* To perform Thread.current.raise as Kernel.raise */
02118     if (th == target_th) {
02119         RUBY_VM_CHECK_INTS(th);
02120     }
02121     return Qnil;
02122 }
02123 
02124 
02125 /*
02126  *  call-seq:
02127  *     thr.exit        -> thr or nil
02128  *     thr.kill        -> thr or nil
02129  *     thr.terminate   -> thr or nil
02130  *
02131  *  Terminates <i>thr</i> and schedules another thread to be run. If this thread
02132  *  is already marked to be killed, <code>exit</code> returns the
02133  *  <code>Thread</code>. If this is the main thread, or the last thread, exits
02134  *  the process.
02135  */
02136 
02137 VALUE
02138 rb_thread_kill(VALUE thread)
02139 {
02140     rb_thread_t *th;
02141 
02142     GetThreadPtr(thread, th);
02143 
02144     if (th != GET_THREAD() && th->safe_level < 4) {
02145         rb_secure(4);
02146     }
02147     if (th->to_kill || th->status == THREAD_KILLED) {
02148         return thread;
02149     }
02150     if (th == th->vm->main_thread) {
02151         rb_exit(EXIT_SUCCESS);
02152     }
02153 
02154     thread_debug("rb_thread_kill: %p (%p)\n", (void *)th, (void *)th->thread_id);
02155 
02156     if (th == GET_THREAD()) {
02157         /* kill myself immediately */
02158         rb_threadptr_to_kill(th);
02159     }
02160     else {
02161         rb_threadptr_pending_interrupt_enque(th, eKillSignal);
02162         rb_threadptr_interrupt(th);
02163     }
02164     return thread;
02165 }
02166 
02167 
02168 /*
02169  *  call-seq:
02170  *     Thread.kill(thread)   -> thread
02171  *
02172  *  Causes the given <em>thread</em> to exit (see <code>Thread::exit</code>).
02173  *
02174  *     count = 0
02175  *     a = Thread.new { loop { count += 1 } }
02176  *     sleep(0.1)       #=> 0
02177  *     Thread.kill(a)   #=> #<Thread:0x401b3d30 dead>
02178  *     count            #=> 93947
02179  *     a.alive?         #=> false
02180  */
02181 
02182 static VALUE
02183 rb_thread_s_kill(VALUE obj, VALUE th)
02184 {
02185     return rb_thread_kill(th);
02186 }
02187 
02188 
02189 /*
02190  *  call-seq:
02191  *     Thread.exit   -> thread
02192  *
02193  *  Terminates the currently running thread and schedules another thread to be
02194  *  run. If this thread is already marked to be killed, <code>exit</code>
02195  *  returns the <code>Thread</code>. If this is the main thread, or the last
02196  *  thread, exit the process.
02197  */
02198 
02199 static VALUE
02200 rb_thread_exit(void)
02201 {
02202     rb_thread_t *th = GET_THREAD();
02203     return rb_thread_kill(th->self);
02204 }
02205 
02206 
02207 /*
02208  *  call-seq:
02209  *     thr.wakeup   -> thr
02210  *
02211  *  Marks <i>thr</i> as eligible for scheduling (it may still remain blocked on
02212  *  I/O, however). Does not invoke the scheduler (see <code>Thread#run</code>).
02213  *
02214  *     c = Thread.new { Thread.stop; puts "hey!" }
02215  *     sleep 0.1 while c.status!='sleep'
02216  *     c.wakeup
02217  *     c.join
02218  *
02219  *  <em>produces:</em>
02220  *
02221  *     hey!
02222  */
02223 
02224 VALUE
02225 rb_thread_wakeup(VALUE thread)
02226 {
02227     if (!RTEST(rb_thread_wakeup_alive(thread))) {
02228         rb_raise(rb_eThreadError, "killed thread");
02229     }
02230     return thread;
02231 }
02232 
02233 VALUE
02234 rb_thread_wakeup_alive(VALUE thread)
02235 {
02236     rb_thread_t *th;
02237     GetThreadPtr(thread, th);
02238 
02239     if (th->status == THREAD_KILLED) {
02240         return Qnil;
02241     }
02242     rb_threadptr_ready(th);
02243     if (th->status == THREAD_STOPPED || th->status == THREAD_STOPPED_FOREVER)
02244         th->status = THREAD_RUNNABLE;
02245     return thread;
02246 }
02247 
02248 
02249 /*
02250  *  call-seq:
02251  *     thr.run   -> thr
02252  *
02253  *  Wakes up <i>thr</i>, making it eligible for scheduling.
02254  *
02255  *     a = Thread.new { puts "a"; Thread.stop; puts "c" }
02256  *     sleep 0.1 while a.status!='sleep'
02257  *     puts "Got here"
02258  *     a.run
02259  *     a.join
02260  *
02261  *  <em>produces:</em>
02262  *
02263  *     a
02264  *     Got here
02265  *     c
02266  */
02267 
02268 VALUE
02269 rb_thread_run(VALUE thread)
02270 {
02271     rb_thread_wakeup(thread);
02272     rb_thread_schedule();
02273     return thread;
02274 }
02275 
02276 
02277 /*
02278  *  call-seq:
02279  *     Thread.stop   -> nil
02280  *
02281  *  Stops execution of the current thread, putting it into a ``sleep'' state,
02282  *  and schedules execution of another thread.
02283  *
02284  *     a = Thread.new { print "a"; Thread.stop; print "c" }
02285  *     sleep 0.1 while a.status!='sleep'
02286  *     print "b"
02287  *     a.run
02288  *     a.join
02289  *
02290  *  <em>produces:</em>
02291  *
02292  *     abc
02293  */
02294 
02295 VALUE
02296 rb_thread_stop(void)
02297 {
02298     if (rb_thread_alone()) {
02299         rb_raise(rb_eThreadError,
02300                  "stopping only thread\n\tnote: use sleep to stop forever");
02301     }
02302     rb_thread_sleep_deadly();
02303     return Qnil;
02304 }
02305 
02306 static int
02307 thread_list_i(st_data_t key, st_data_t val, void *data)
02308 {
02309     VALUE ary = (VALUE)data;
02310     rb_thread_t *th;
02311     GetThreadPtr((VALUE)key, th);
02312 
02313     switch (th->status) {
02314       case THREAD_RUNNABLE:
02315       case THREAD_STOPPED:
02316       case THREAD_STOPPED_FOREVER:
02317         rb_ary_push(ary, th->self);
02318       default:
02319         break;
02320     }
02321     return ST_CONTINUE;
02322 }
02323 
02324 /********************************************************************/
02325 
02326 /*
02327  *  call-seq:
02328  *     Thread.list   -> array
02329  *
02330  *  Returns an array of <code>Thread</code> objects for all threads that are
02331  *  either runnable or stopped.
02332  *
02333  *     Thread.new { sleep(200) }
02334  *     Thread.new { 1000000.times {|i| i*i } }
02335  *     Thread.new { Thread.stop }
02336  *     Thread.list.each {|t| p t}
02337  *
02338  *  <em>produces:</em>
02339  *
02340  *     #<Thread:0x401b3e84 sleep>
02341  *     #<Thread:0x401b3f38 run>
02342  *     #<Thread:0x401b3fb0 sleep>
02343  *     #<Thread:0x401bdf4c run>
02344  */
02345 
02346 VALUE
02347 rb_thread_list(void)
02348 {
02349     VALUE ary = rb_ary_new();
02350     st_foreach(GET_THREAD()->vm->living_threads, thread_list_i, ary);
02351     return ary;
02352 }
02353 
02354 VALUE
02355 rb_thread_current(void)
02356 {
02357     return GET_THREAD()->self;
02358 }
02359 
02360 /*
02361  *  call-seq:
02362  *     Thread.current   -> thread
02363  *
02364  *  Returns the currently executing thread.
02365  *
02366  *     Thread.current   #=> #<Thread:0x401bdf4c run>
02367  */
02368 
02369 static VALUE
02370 thread_s_current(VALUE klass)
02371 {
02372     return rb_thread_current();
02373 }
02374 
02375 VALUE
02376 rb_thread_main(void)
02377 {
02378     return GET_THREAD()->vm->main_thread->self;
02379 }
02380 
02381 /*
02382  *  call-seq:
02383  *     Thread.main   -> thread
02384  *
02385  *  Returns the main thread.
02386  */
02387 
02388 static VALUE
02389 rb_thread_s_main(VALUE klass)
02390 {
02391     return rb_thread_main();
02392 }
02393 
02394 
02395 /*
02396  *  call-seq:
02397  *     Thread.abort_on_exception   -> true or false
02398  *
02399  *  Returns the status of the global ``abort on exception'' condition.  The
02400  *  default is <code>false</code>. When set to <code>true</code>, or if the
02401  *  global <code>$DEBUG</code> flag is <code>true</code> (perhaps because the
02402  *  command line option <code>-d</code> was specified) all threads will abort
02403  *  (the process will <code>exit(0)</code>) if an exception is raised in any
02404  *  thread. See also <code>Thread::abort_on_exception=</code>.
02405  */
02406 
02407 static VALUE
02408 rb_thread_s_abort_exc(void)
02409 {
02410     return GET_THREAD()->vm->thread_abort_on_exception ? Qtrue : Qfalse;
02411 }
02412 
02413 
02414 /*
02415  *  call-seq:
02416  *     Thread.abort_on_exception= boolean   -> true or false
02417  *
02418  *  When set to <code>true</code>, all threads will abort if an exception is
02419  *  raised. Returns the new state.
02420  *
02421  *     Thread.abort_on_exception = true
02422  *     t1 = Thread.new do
02423  *       puts  "In new thread"
02424  *       raise "Exception from thread"
02425  *     end
02426  *     sleep(1)
02427  *     puts "not reached"
02428  *
02429  *  <em>produces:</em>
02430  *
02431  *     In new thread
02432  *     prog.rb:4: Exception from thread (RuntimeError)
02433  *      from prog.rb:2:in `initialize'
02434  *      from prog.rb:2:in `new'
02435  *      from prog.rb:2
02436  */
02437 
02438 static VALUE
02439 rb_thread_s_abort_exc_set(VALUE self, VALUE val)
02440 {
02441     rb_secure(4);
02442     GET_THREAD()->vm->thread_abort_on_exception = RTEST(val);
02443     return val;
02444 }
02445 
02446 
02447 /*
02448  *  call-seq:
02449  *     thr.abort_on_exception   -> true or false
02450  *
02451  *  Returns the status of the thread-local ``abort on exception'' condition for
02452  *  <i>thr</i>. The default is <code>false</code>. See also
02453  *  <code>Thread::abort_on_exception=</code>.
02454  */
02455 
02456 static VALUE
02457 rb_thread_abort_exc(VALUE thread)
02458 {
02459     rb_thread_t *th;
02460     GetThreadPtr(thread, th);
02461     return th->abort_on_exception ? Qtrue : Qfalse;
02462 }
02463 
02464 
02465 /*
02466  *  call-seq:
02467  *     thr.abort_on_exception= boolean   -> true or false
02468  *
02469  *  When set to <code>true</code>, causes all threads (including the main
02470  *  program) to abort if an exception is raised in <i>thr</i>. The process will
02471  *  effectively <code>exit(0)</code>.
02472  */
02473 
02474 static VALUE
02475 rb_thread_abort_exc_set(VALUE thread, VALUE val)
02476 {
02477     rb_thread_t *th;
02478     rb_secure(4);
02479 
02480     GetThreadPtr(thread, th);
02481     th->abort_on_exception = RTEST(val);
02482     return val;
02483 }
02484 
02485 
02486 /*
02487  *  call-seq:
02488  *     thr.group   -> thgrp or nil
02489  *
02490  *  Returns the <code>ThreadGroup</code> which contains <i>thr</i>, or nil if
02491  *  the thread is not a member of any group.
02492  *
02493  *     Thread.main.group   #=> #<ThreadGroup:0x4029d914>
02494  */
02495 
02496 VALUE
02497 rb_thread_group(VALUE thread)
02498 {
02499     rb_thread_t *th;
02500     VALUE group;
02501     GetThreadPtr(thread, th);
02502     group = th->thgroup;
02503 
02504     if (!group) {
02505         group = Qnil;
02506     }
02507     return group;
02508 }
02509 
02510 static const char *
02511 thread_status_name(rb_thread_t *th)
02512 {
02513     switch (th->status) {
02514       case THREAD_RUNNABLE:
02515         if (th->to_kill)
02516             return "aborting";
02517         else
02518             return "run";
02519       case THREAD_STOPPED:
02520       case THREAD_STOPPED_FOREVER:
02521         return "sleep";
02522       case THREAD_KILLED:
02523         return "dead";
02524       default:
02525         return "unknown";
02526     }
02527 }
02528 
02529 static int
02530 rb_threadptr_dead(rb_thread_t *th)
02531 {
02532     return th->status == THREAD_KILLED;
02533 }
02534 
02535 
02536 /*
02537  *  call-seq:
02538  *     thr.status   -> string, false or nil
02539  *
02540  *  Returns the status of <i>thr</i>: ``<code>sleep</code>'' if <i>thr</i> is
02541  *  sleeping or waiting on I/O, ``<code>run</code>'' if <i>thr</i> is executing,
02542  *  ``<code>aborting</code>'' if <i>thr</i> is aborting, <code>false</code> if
02543  *  <i>thr</i> terminated normally, and <code>nil</code> if <i>thr</i>
02544  *  terminated with an exception.
02545  *
02546  *     a = Thread.new { raise("die now") }
02547  *     b = Thread.new { Thread.stop }
02548  *     c = Thread.new { Thread.exit }
02549  *     d = Thread.new { sleep }
02550  *     d.kill                  #=> #<Thread:0x401b3678 aborting>
02551  *     a.status                #=> nil
02552  *     b.status                #=> "sleep"
02553  *     c.status                #=> false
02554  *     d.status                #=> "aborting"
02555  *     Thread.current.status   #=> "run"
02556  */
02557 
02558 static VALUE
02559 rb_thread_status(VALUE thread)
02560 {
02561     rb_thread_t *th;
02562     GetThreadPtr(thread, th);
02563 
02564     if (rb_threadptr_dead(th)) {
02565         if (!NIL_P(th->errinfo) && !FIXNUM_P(th->errinfo)
02566             /* TODO */ ) {
02567             return Qnil;
02568         }
02569         return Qfalse;
02570     }
02571     return rb_str_new2(thread_status_name(th));
02572 }
02573 
02574 
02575 /*
02576  *  call-seq:
02577  *     thr.alive?   -> true or false
02578  *
02579  *  Returns <code>true</code> if <i>thr</i> is running or sleeping.
02580  *
02581  *     thr = Thread.new { }
02582  *     thr.join                #=> #<Thread:0x401b3fb0 dead>
02583  *     Thread.current.alive?   #=> true
02584  *     thr.alive?              #=> false
02585  */
02586 
02587 static VALUE
02588 rb_thread_alive_p(VALUE thread)
02589 {
02590     rb_thread_t *th;
02591     GetThreadPtr(thread, th);
02592 
02593     if (rb_threadptr_dead(th))
02594         return Qfalse;
02595     return Qtrue;
02596 }
02597 
02598 /*
02599  *  call-seq:
02600  *     thr.stop?   -> true or false
02601  *
02602  *  Returns <code>true</code> if <i>thr</i> is dead or sleeping.
02603  *
02604  *     a = Thread.new { Thread.stop }
02605  *     b = Thread.current
02606  *     a.stop?   #=> true
02607  *     b.stop?   #=> false
02608  */
02609 
02610 static VALUE
02611 rb_thread_stop_p(VALUE thread)
02612 {
02613     rb_thread_t *th;
02614     GetThreadPtr(thread, th);
02615 
02616     if (rb_threadptr_dead(th))
02617         return Qtrue;
02618     if (th->status == THREAD_STOPPED || th->status == THREAD_STOPPED_FOREVER)
02619         return Qtrue;
02620     return Qfalse;
02621 }
02622 
02623 /*
02624  *  call-seq:
02625  *     thr.safe_level   -> integer
02626  *
02627  *  Returns the safe level in effect for <i>thr</i>. Setting thread-local safe
02628  *  levels can help when implementing sandboxes which run insecure code.
02629  *
02630  *     thr = Thread.new { $SAFE = 3; sleep }
02631  *     Thread.current.safe_level   #=> 0
02632  *     thr.safe_level              #=> 3
02633  */
02634 
02635 static VALUE
02636 rb_thread_safe_level(VALUE thread)
02637 {
02638     rb_thread_t *th;
02639     GetThreadPtr(thread, th);
02640 
02641     return INT2NUM(th->safe_level);
02642 }
02643 
02644 /*
02645  * call-seq:
02646  *   thr.inspect   -> string
02647  *
02648  * Dump the name, id, and status of _thr_ to a string.
02649  */
02650 
02651 static VALUE
02652 rb_thread_inspect(VALUE thread)
02653 {
02654     const char *cname = rb_obj_classname(thread);
02655     rb_thread_t *th;
02656     const char *status;
02657     VALUE str;
02658 
02659     GetThreadPtr(thread, th);
02660     status = thread_status_name(th);
02661     str = rb_sprintf("#<%s:%p %s>", cname, (void *)thread, status);
02662     OBJ_INFECT(str, thread);
02663 
02664     return str;
02665 }
02666 
02667 VALUE
02668 rb_thread_local_aref(VALUE thread, ID id)
02669 {
02670     rb_thread_t *th;
02671     st_data_t val;
02672 
02673     GetThreadPtr(thread, th);
02674     if (rb_safe_level() >= 4 && th != GET_THREAD()) {
02675         rb_raise(rb_eSecurityError, "Insecure: thread locals");
02676     }
02677     if (!th->local_storage) {
02678         return Qnil;
02679     }
02680     if (st_lookup(th->local_storage, id, &val)) {
02681         return (VALUE)val;
02682     }
02683     return Qnil;
02684 }
02685 
02686 /*
02687  *  call-seq:
02688  *      thr[sym]   -> obj or nil
02689  *
02690  *  Attribute Reference---Returns the value of a fiber-local variable (current thread's root fiber
02691  *  if not explicitely inside a Fiber), using either a symbol or a string name.
02692  *  If the specified variable does not exist, returns <code>nil</code>.
02693  *
02694  *     [
02695  *       Thread.new { Thread.current["name"] = "A" },
02696  *       Thread.new { Thread.current[:name]  = "B" },
02697  *       Thread.new { Thread.current["name"] = "C" }
02698  *     ].each do |th|
02699  *       th.join
02700  *       puts "#{th.inspect}: #{th[:name]}"
02701  *     end
02702  *
02703  *  <em>produces:</em>
02704  *
02705  *     #<Thread:0x00000002a54220 dead>: A
02706  *     #<Thread:0x00000002a541a8 dead>: B
02707  *     #<Thread:0x00000002a54130 dead>: C
02708  *
02709  *  Thread#[] and Thread#[]= are not thread-local but fiber-local.
02710  *  This confusion did not exist in Ruby 1.8 because
02711  *  fibers were only available since Ruby 1.9.
02712  *  Ruby 1.9 chooses that the methods behaves fiber-local to save
02713  *  following idiom for dynamic scope.
02714  *
02715  *    def meth(newvalue)
02716  *      begin
02717  *        oldvalue = Thread.current[:name]
02718  *        Thread.current[:name] = newvalue
02719  *        yield
02720  *      ensure
02721  *        Thread.current[:name] = oldvalue
02722  *      end
02723  *    end
02724  *
02725  *  The idiom may not work as dynamic scope if the methods are thread-local
02726  *  and a given block switches fiber.
02727  *
02728  *    f = Fiber.new {
02729  *      meth(1) {
02730  *        Fiber.yield
02731  *      }
02732  *    }
02733  *    meth(2) {
02734  *      f.resume
02735  *    }
02736  *    f.resume
02737  *    p Thread.current[:name]
02738  *    #=> nil if fiber-local
02739  *    #=> 2 if thread-local (The value 2 is leaked to outside of meth method.)
02740  *
02741  *  For thread-local variables, please see <code>Thread#thread_local_get</code>
02742  *  and <code>Thread#thread_local_set</code>.
02743  *
02744  */
02745 
02746 static VALUE
02747 rb_thread_aref(VALUE thread, VALUE id)
02748 {
02749     return rb_thread_local_aref(thread, rb_to_id(id));
02750 }
02751 
02752 VALUE
02753 rb_thread_local_aset(VALUE thread, ID id, VALUE val)
02754 {
02755     rb_thread_t *th;
02756     GetThreadPtr(thread, th);
02757 
02758     if (rb_safe_level() >= 4 && th != GET_THREAD()) {
02759         rb_raise(rb_eSecurityError, "Insecure: can't modify thread locals");
02760     }
02761     if (OBJ_FROZEN(thread)) {
02762         rb_error_frozen("thread locals");
02763     }
02764     if (!th->local_storage) {
02765         th->local_storage = st_init_numtable();
02766     }
02767     if (NIL_P(val)) {
02768         st_delete_wrap(th->local_storage, id);
02769         return Qnil;
02770     }
02771     st_insert(th->local_storage, id, val);
02772     return val;
02773 }
02774 
02775 /*
02776  *  call-seq:
02777  *      thr[sym] = obj   -> obj
02778  *
02779  *  Attribute Assignment---Sets or creates the value of a fiber-local variable,
02780  *  using either a symbol or a string. See also <code>Thread#[]</code>.  For
02781  *  thread-local variables, please see <code>Thread#thread_variable_set</code>
02782  *  and <code>Thread#thread_variable_get</code>.
02783  */
02784 
02785 static VALUE
02786 rb_thread_aset(VALUE self, VALUE id, VALUE val)
02787 {
02788     return rb_thread_local_aset(self, rb_to_id(id), val);
02789 }
02790 
02791 /*
02792  *  call-seq:
02793  *      thr.thread_variable_get(key)  -> obj or nil
02794  *
02795  *  Returns the value of a thread local variable that has been set.  Note that
02796  *  these are different than fiber local values.  For fiber local values,
02797  *  please see Thread#[] and Thread#[]=.
02798  *
02799  *  Thread local values are carried along with threads, and do not respect
02800  *  fibers.  For example:
02801  *
02802  *    Thread.new {
02803  *      Thread.current.thread_variable_set("foo", "bar") # set a thread local
02804  *      Thread.current["foo"] = "bar"                    # set a fiber local
02805  *
02806  *      Fiber.new {
02807  *        Fiber.yield [
02808  *          Thread.current.thread_variable_get("foo"), # get the thread local
02809  *          Thread.current["foo"],                     # get the fiber local
02810  *        ]
02811  *      }.resume
02812  *    }.join.value # => ['bar', nil]
02813  *
02814  *  The value "bar" is returned for the thread local, where nil is returned
02815  *  for the fiber local.  The fiber is executed in the same thread, so the
02816  *  thread local values are available.
02817  *
02818  *  See also Thread#[]
02819  */
02820 
02821 static VALUE
02822 rb_thread_variable_get(VALUE thread, VALUE id)
02823 {
02824     VALUE locals;
02825     rb_thread_t *th;
02826 
02827     GetThreadPtr(thread, th);
02828 
02829     if (rb_safe_level() >= 4 && th != GET_THREAD()) {
02830         rb_raise(rb_eSecurityError, "Insecure: can't modify thread locals");
02831     }
02832 
02833     locals = rb_iv_get(thread, "locals");
02834     return rb_hash_aref(locals, ID2SYM(rb_to_id(id)));
02835 }
02836 
02837 /*
02838  *  call-seq:
02839  *      thr.thread_variable_set(key, value)
02840  *
02841  *  Sets a thread local with +key+ to +value+.  Note that these are local to
02842  *  threads, and not to fibers.  Please see Thread#thread_variable_get and
02843  *  Thread#[] for more information.
02844  */
02845 
02846 static VALUE
02847 rb_thread_variable_set(VALUE thread, VALUE id, VALUE val)
02848 {
02849     VALUE locals;
02850     rb_thread_t *th;
02851 
02852     GetThreadPtr(thread, th);
02853 
02854     if (rb_safe_level() >= 4 && th != GET_THREAD()) {
02855         rb_raise(rb_eSecurityError, "Insecure: can't modify thread locals");
02856     }
02857     if (OBJ_FROZEN(thread)) {
02858         rb_error_frozen("thread locals");
02859     }
02860 
02861     locals = rb_iv_get(thread, "locals");
02862     return rb_hash_aset(locals, ID2SYM(rb_to_id(id)), val);
02863 }
02864 
02865 /*
02866  *  call-seq:
02867  *     thr.key?(sym)   -> true or false
02868  *
02869  *  Returns <code>true</code> if the given string (or symbol) exists as a
02870  *  fiber-local variable.
02871  *
02872  *     me = Thread.current
02873  *     me[:oliver] = "a"
02874  *     me.key?(:oliver)    #=> true
02875  *     me.key?(:stanley)   #=> false
02876  */
02877 
02878 static VALUE
02879 rb_thread_key_p(VALUE self, VALUE key)
02880 {
02881     rb_thread_t *th;
02882     ID id = rb_to_id(key);
02883 
02884     GetThreadPtr(self, th);
02885 
02886     if (!th->local_storage) {
02887         return Qfalse;
02888     }
02889     if (st_lookup(th->local_storage, id, 0)) {
02890         return Qtrue;
02891     }
02892     return Qfalse;
02893 }
02894 
02895 static int
02896 thread_keys_i(ID key, VALUE value, VALUE ary)
02897 {
02898     rb_ary_push(ary, ID2SYM(key));
02899     return ST_CONTINUE;
02900 }
02901 
02902 static int
02903 vm_living_thread_num(rb_vm_t *vm)
02904 {
02905     return (int)vm->living_threads->num_entries;
02906 }
02907 
02908 int
02909 rb_thread_alone(void)
02910 {
02911     int num = 1;
02912     if (GET_THREAD()->vm->living_threads) {
02913         num = vm_living_thread_num(GET_THREAD()->vm);
02914         thread_debug("rb_thread_alone: %d\n", num);
02915     }
02916     return num == 1;
02917 }
02918 
02919 /*
02920  *  call-seq:
02921  *     thr.keys   -> array
02922  *
02923  *  Returns an an array of the names of the fiber-local variables (as Symbols).
02924  *
02925  *     thr = Thread.new do
02926  *       Thread.current[:cat] = 'meow'
02927  *       Thread.current["dog"] = 'woof'
02928  *     end
02929  *     thr.join   #=> #<Thread:0x401b3f10 dead>
02930  *     thr.keys   #=> [:dog, :cat]
02931  */
02932 
02933 static VALUE
02934 rb_thread_keys(VALUE self)
02935 {
02936     rb_thread_t *th;
02937     VALUE ary = rb_ary_new();
02938     GetThreadPtr(self, th);
02939 
02940     if (th->local_storage) {
02941         st_foreach(th->local_storage, thread_keys_i, ary);
02942     }
02943     return ary;
02944 }
02945 
02946 static int
02947 keys_i(VALUE key, VALUE value, VALUE ary)
02948 {
02949     rb_ary_push(ary, key);
02950     return ST_CONTINUE;
02951 }
02952 
02953 /*
02954  *  call-seq:
02955  *     thr.thread_variables   -> array
02956  *
02957  *  Returns an an array of the names of the thread-local variables (as Symbols).
02958  *
02959  *     thr = Thread.new do
02960  *       Thread.current.thread_variable_set(:cat, 'meow')
02961  *       Thread.current.thread_variable_set("dog", 'woof')
02962  *     end
02963  *     thr.join               #=> #<Thread:0x401b3f10 dead>
02964  *     thr.thread_variables   #=> [:dog, :cat]
02965  *
02966  *  Note that these are not fiber local variables.  Please see Thread#[] and
02967  *  Thread#thread_variable_get for more details.
02968  */
02969 
02970 static VALUE
02971 rb_thread_variables(VALUE thread)
02972 {
02973     VALUE locals;
02974     VALUE ary;
02975 
02976     locals = rb_iv_get(thread, "locals");
02977     ary = rb_ary_new();
02978     rb_hash_foreach(locals, keys_i, ary);
02979 
02980     return ary;
02981 }
02982 
02983 /*
02984  *  call-seq:
02985  *     thr.thread_variable?(key)   -> true or false
02986  *
02987  *  Returns <code>true</code> if the given string (or symbol) exists as a
02988  *  thread-local variable.
02989  *
02990  *     me = Thread.current
02991  *     me.thread_variable_set(:oliver, "a")
02992  *     me.thread_variable?(:oliver)    #=> true
02993  *     me.thread_variable?(:stanley)   #=> false
02994  *
02995  *  Note that these are not fiber local variables.  Please see Thread#[] and
02996  *  Thread#thread_variable_get for more details.
02997  */
02998 
02999 static VALUE
03000 rb_thread_variable_p(VALUE thread, VALUE key)
03001 {
03002     VALUE locals;
03003 
03004     locals = rb_iv_get(thread, "locals");
03005 
03006     if (!RHASH(locals)->ntbl)
03007         return Qfalse;
03008 
03009     if (st_lookup(RHASH(locals)->ntbl, ID2SYM(rb_to_id(key)), 0)) {
03010         return Qtrue;
03011     }
03012 
03013     return Qfalse;
03014 }
03015 
03016 /*
03017  *  call-seq:
03018  *     thr.priority   -> integer
03019  *
03020  *  Returns the priority of <i>thr</i>. Default is inherited from the
03021  *  current thread which creating the new thread, or zero for the
03022  *  initial main thread; higher-priority thread will run more frequently
03023  *  than lower-priority threads (but lower-priority threads can also run).
03024  *
03025  *  This is just hint for Ruby thread scheduler.  It may be ignored on some
03026  *  platform.
03027  *
03028  *     Thread.current.priority   #=> 0
03029  */
03030 
03031 static VALUE
03032 rb_thread_priority(VALUE thread)
03033 {
03034     rb_thread_t *th;
03035     GetThreadPtr(thread, th);
03036     return INT2NUM(th->priority);
03037 }
03038 
03039 
03040 /*
03041  *  call-seq:
03042  *     thr.priority= integer   -> thr
03043  *
03044  *  Sets the priority of <i>thr</i> to <i>integer</i>. Higher-priority threads
03045  *  will run more frequently than lower-priority threads (but lower-priority
03046  *  threads can also run).
03047  *
03048  *  This is just hint for Ruby thread scheduler.  It may be ignored on some
03049  *  platform.
03050  *
03051  *     count1 = count2 = 0
03052  *     a = Thread.new do
03053  *           loop { count1 += 1 }
03054  *         end
03055  *     a.priority = -1
03056  *
03057  *     b = Thread.new do
03058  *           loop { count2 += 1 }
03059  *         end
03060  *     b.priority = -2
03061  *     sleep 1   #=> 1
03062  *     count1    #=> 622504
03063  *     count2    #=> 5832
03064  */
03065 
03066 static VALUE
03067 rb_thread_priority_set(VALUE thread, VALUE prio)
03068 {
03069     rb_thread_t *th;
03070     int priority;
03071     GetThreadPtr(thread, th);
03072 
03073     rb_secure(4);
03074 
03075 #if USE_NATIVE_THREAD_PRIORITY
03076     th->priority = NUM2INT(prio);
03077     native_thread_apply_priority(th);
03078 #else
03079     priority = NUM2INT(prio);
03080     if (priority > RUBY_THREAD_PRIORITY_MAX) {
03081         priority = RUBY_THREAD_PRIORITY_MAX;
03082     }
03083     else if (priority < RUBY_THREAD_PRIORITY_MIN) {
03084         priority = RUBY_THREAD_PRIORITY_MIN;
03085     }
03086     th->priority = priority;
03087 #endif
03088     return INT2NUM(th->priority);
03089 }
03090 
03091 /* for IO */
03092 
03093 #if defined(NFDBITS) && defined(HAVE_RB_FD_INIT)
03094 
03095 /*
03096  * several Unix platforms support file descriptors bigger than FD_SETSIZE
03097  * in select(2) system call.
03098  *
03099  * - Linux 2.2.12 (?)
03100  * - NetBSD 1.2 (src/sys/kern/sys_generic.c:1.25)
03101  *   select(2) documents how to allocate fd_set dynamically.
03102  *   http://netbsd.gw.com/cgi-bin/man-cgi?select++NetBSD-4.0
03103  * - FreeBSD 2.2 (src/sys/kern/sys_generic.c:1.19)
03104  * - OpenBSD 2.0 (src/sys/kern/sys_generic.c:1.4)
03105  *   select(2) documents how to allocate fd_set dynamically.
03106  *   http://www.openbsd.org/cgi-bin/man.cgi?query=select&manpath=OpenBSD+4.4
03107  * - HP-UX documents how to allocate fd_set dynamically.
03108  *   http://docs.hp.com/en/B2355-60105/select.2.html
03109  * - Solaris 8 has select_large_fdset
03110  * - Mac OS X 10.7 (Lion)
03111  *   select(2) returns EINVAL if nfds is greater than FD_SET_SIZE and
03112  *   _DARWIN_UNLIMITED_SELECT (or _DARWIN_C_SOURCE) isn't defined.
03113  *   http://developer.apple.com/library/mac/#releasenotes/Darwin/SymbolVariantsRelNotes/_index.html
03114  *
03115  * When fd_set is not big enough to hold big file descriptors,
03116  * it should be allocated dynamically.
03117  * Note that this assumes fd_set is structured as bitmap.
03118  *
03119  * rb_fd_init allocates the memory.
03120  * rb_fd_term free the memory.
03121  * rb_fd_set may re-allocates bitmap.
03122  *
03123  * So rb_fd_set doesn't reject file descriptors bigger than FD_SETSIZE.
03124  */
03125 
03126 void
03127 rb_fd_init(rb_fdset_t *fds)
03128 {
03129     fds->maxfd = 0;
03130     fds->fdset = ALLOC(fd_set);
03131     FD_ZERO(fds->fdset);
03132 }
03133 
03134 void
03135 rb_fd_init_copy(rb_fdset_t *dst, rb_fdset_t *src)
03136 {
03137     size_t size = howmany(rb_fd_max(src), NFDBITS) * sizeof(fd_mask);
03138 
03139     if (size < sizeof(fd_set))
03140         size = sizeof(fd_set);
03141     dst->maxfd = src->maxfd;
03142     dst->fdset = xmalloc(size);
03143     memcpy(dst->fdset, src->fdset, size);
03144 }
03145 
03146 void
03147 rb_fd_term(rb_fdset_t *fds)
03148 {
03149     if (fds->fdset) xfree(fds->fdset);
03150     fds->maxfd = 0;
03151     fds->fdset = 0;
03152 }
03153 
03154 void
03155 rb_fd_zero(rb_fdset_t *fds)
03156 {
03157     if (fds->fdset)
03158         MEMZERO(fds->fdset, fd_mask, howmany(fds->maxfd, NFDBITS));
03159 }
03160 
03161 static void
03162 rb_fd_resize(int n, rb_fdset_t *fds)
03163 {
03164     size_t m = howmany(n + 1, NFDBITS) * sizeof(fd_mask);
03165     size_t o = howmany(fds->maxfd, NFDBITS) * sizeof(fd_mask);
03166 
03167     if (m < sizeof(fd_set)) m = sizeof(fd_set);
03168     if (o < sizeof(fd_set)) o = sizeof(fd_set);
03169 
03170     if (m > o) {
03171         fds->fdset = xrealloc(fds->fdset, m);
03172         memset((char *)fds->fdset + o, 0, m - o);
03173     }
03174     if (n >= fds->maxfd) fds->maxfd = n + 1;
03175 }
03176 
03177 void
03178 rb_fd_set(int n, rb_fdset_t *fds)
03179 {
03180     rb_fd_resize(n, fds);
03181     FD_SET(n, fds->fdset);
03182 }
03183 
03184 void
03185 rb_fd_clr(int n, rb_fdset_t *fds)
03186 {
03187     if (n >= fds->maxfd) return;
03188     FD_CLR(n, fds->fdset);
03189 }
03190 
03191 int
03192 rb_fd_isset(int n, const rb_fdset_t *fds)
03193 {
03194     if (n >= fds->maxfd) return 0;
03195     return FD_ISSET(n, fds->fdset) != 0; /* "!= 0" avoids FreeBSD PR 91421 */
03196 }
03197 
03198 void
03199 rb_fd_copy(rb_fdset_t *dst, const fd_set *src, int max)
03200 {
03201     size_t size = howmany(max, NFDBITS) * sizeof(fd_mask);
03202 
03203     if (size < sizeof(fd_set)) size = sizeof(fd_set);
03204     dst->maxfd = max;
03205     dst->fdset = xrealloc(dst->fdset, size);
03206     memcpy(dst->fdset, src, size);
03207 }
03208 
03209 static void
03210 rb_fd_rcopy(fd_set *dst, rb_fdset_t *src)
03211 {
03212     size_t size = howmany(rb_fd_max(src), NFDBITS) * sizeof(fd_mask);
03213 
03214     if (size > sizeof(fd_set)) {
03215         rb_raise(rb_eArgError, "too large fdsets");
03216     }
03217     memcpy(dst, rb_fd_ptr(src), sizeof(fd_set));
03218 }
03219 
03220 void
03221 rb_fd_dup(rb_fdset_t *dst, const rb_fdset_t *src)
03222 {
03223     size_t size = howmany(rb_fd_max(src), NFDBITS) * sizeof(fd_mask);
03224 
03225     if (size < sizeof(fd_set))
03226         size = sizeof(fd_set);
03227     dst->maxfd = src->maxfd;
03228     dst->fdset = xrealloc(dst->fdset, size);
03229     memcpy(dst->fdset, src->fdset, size);
03230 }
03231 
03232 #ifdef __native_client__
03233 int select(int nfds, fd_set *readfds, fd_set *writefds,
03234            fd_set *exceptfds, struct timeval *timeout);
03235 #endif
03236 
03237 int
03238 rb_fd_select(int n, rb_fdset_t *readfds, rb_fdset_t *writefds, rb_fdset_t *exceptfds, struct timeval *timeout)
03239 {
03240     fd_set *r = NULL, *w = NULL, *e = NULL;
03241     if (readfds) {
03242         rb_fd_resize(n - 1, readfds);
03243         r = rb_fd_ptr(readfds);
03244     }
03245     if (writefds) {
03246         rb_fd_resize(n - 1, writefds);
03247         w = rb_fd_ptr(writefds);
03248     }
03249     if (exceptfds) {
03250         rb_fd_resize(n - 1, exceptfds);
03251         e = rb_fd_ptr(exceptfds);
03252     }
03253     return select(n, r, w, e, timeout);
03254 }
03255 
03256 #undef FD_ZERO
03257 #undef FD_SET
03258 #undef FD_CLR
03259 #undef FD_ISSET
03260 
03261 #define FD_ZERO(f)      rb_fd_zero(f)
03262 #define FD_SET(i, f)    rb_fd_set((i), (f))
03263 #define FD_CLR(i, f)    rb_fd_clr((i), (f))
03264 #define FD_ISSET(i, f)  rb_fd_isset((i), (f))
03265 
03266 #elif defined(_WIN32)
03267 
03268 void
03269 rb_fd_init(rb_fdset_t *set)
03270 {
03271     set->capa = FD_SETSIZE;
03272     set->fdset = ALLOC(fd_set);
03273     FD_ZERO(set->fdset);
03274 }
03275 
03276 void
03277 rb_fd_init_copy(rb_fdset_t *dst, rb_fdset_t *src)
03278 {
03279     rb_fd_init(dst);
03280     rb_fd_dup(dst, src);
03281 }
03282 
03283 static void
03284 rb_fd_rcopy(fd_set *dst, rb_fdset_t *src)
03285 {
03286     int max = rb_fd_max(src);
03287 
03288     /* we assume src is the result of select() with dst, so dst should be
03289      * larger or equal than src. */
03290     if (max > FD_SETSIZE || (UINT)max > dst->fd_count) {
03291         rb_raise(rb_eArgError, "too large fdsets");
03292     }
03293 
03294     memcpy(dst->fd_array, src->fdset->fd_array, max);
03295     dst->fd_count = max;
03296 }
03297 
03298 void
03299 rb_fd_term(rb_fdset_t *set)
03300 {
03301     xfree(set->fdset);
03302     set->fdset = NULL;
03303     set->capa = 0;
03304 }
03305 
03306 void
03307 rb_fd_set(int fd, rb_fdset_t *set)
03308 {
03309     unsigned int i;
03310     SOCKET s = rb_w32_get_osfhandle(fd);
03311 
03312     for (i = 0; i < set->fdset->fd_count; i++) {
03313         if (set->fdset->fd_array[i] == s) {
03314             return;
03315         }
03316     }
03317     if (set->fdset->fd_count >= (unsigned)set->capa) {
03318         set->capa = (set->fdset->fd_count / FD_SETSIZE + 1) * FD_SETSIZE;
03319         set->fdset = xrealloc(set->fdset, sizeof(unsigned int) + sizeof(SOCKET) * set->capa);
03320     }
03321     set->fdset->fd_array[set->fdset->fd_count++] = s;
03322 }
03323 
03324 #undef FD_ZERO
03325 #undef FD_SET
03326 #undef FD_CLR
03327 #undef FD_ISSET
03328 
03329 #define FD_ZERO(f)      rb_fd_zero(f)
03330 #define FD_SET(i, f)    rb_fd_set((i), (f))
03331 #define FD_CLR(i, f)    rb_fd_clr((i), (f))
03332 #define FD_ISSET(i, f)  rb_fd_isset((i), (f))
03333 
03334 #else
03335 #define rb_fd_rcopy(d, s) (*(d) = *(s))
03336 #endif
03337 
03338 static int
03339 do_select(int n, rb_fdset_t *read, rb_fdset_t *write, rb_fdset_t *except,
03340           struct timeval *timeout)
03341 {
03342     int UNINITIALIZED_VAR(result);
03343     int lerrno;
03344     rb_fdset_t UNINITIALIZED_VAR(orig_read);
03345     rb_fdset_t UNINITIALIZED_VAR(orig_write);
03346     rb_fdset_t UNINITIALIZED_VAR(orig_except);
03347     double limit = 0;
03348     struct timeval wait_rest;
03349     rb_thread_t *th = GET_THREAD();
03350 
03351     if (timeout) {
03352         limit = timeofday();
03353         limit += (double)timeout->tv_sec+(double)timeout->tv_usec*1e-6;
03354         wait_rest = *timeout;
03355         timeout = &wait_rest;
03356     }
03357 
03358     if (read)
03359         rb_fd_init_copy(&orig_read, read);
03360     if (write)
03361         rb_fd_init_copy(&orig_write, write);
03362     if (except)
03363         rb_fd_init_copy(&orig_except, except);
03364 
03365   retry:
03366     lerrno = 0;
03367 
03368     BLOCKING_REGION({
03369             result = native_fd_select(n, read, write, except, timeout, th);
03370             if (result < 0) lerrno = errno;
03371         }, ubf_select, th, FALSE);
03372 
03373     RUBY_VM_CHECK_INTS_BLOCKING(th);
03374 
03375     errno = lerrno;
03376 
03377     if (result < 0) {
03378         switch (errno) {
03379           case EINTR:
03380 #ifdef ERESTART
03381           case ERESTART:
03382 #endif
03383             if (read)
03384                 rb_fd_dup(read, &orig_read);
03385             if (write)
03386                 rb_fd_dup(write, &orig_write);
03387             if (except)
03388                 rb_fd_dup(except, &orig_except);
03389 
03390             if (timeout) {
03391                 double d = limit - timeofday();
03392 
03393                 wait_rest.tv_sec = (time_t)d;
03394                 wait_rest.tv_usec = (int)((d-(double)wait_rest.tv_sec)*1e6);
03395                 if (wait_rest.tv_sec < 0)  wait_rest.tv_sec = 0;
03396                 if (wait_rest.tv_usec < 0) wait_rest.tv_usec = 0;
03397             }
03398 
03399             goto retry;
03400           default:
03401             break;
03402         }
03403     }
03404 
03405     if (read)
03406         rb_fd_term(&orig_read);
03407     if (write)
03408         rb_fd_term(&orig_write);
03409     if (except)
03410         rb_fd_term(&orig_except);
03411 
03412     return result;
03413 }
03414 
03415 static void
03416 rb_thread_wait_fd_rw(int fd, int read)
03417 {
03418     int result = 0;
03419     int events = read ? RB_WAITFD_IN : RB_WAITFD_OUT;
03420 
03421     thread_debug("rb_thread_wait_fd_rw(%d, %s)\n", fd, read ? "read" : "write");
03422 
03423     if (fd < 0) {
03424         rb_raise(rb_eIOError, "closed stream");
03425     }
03426 
03427     result = rb_wait_for_single_fd(fd, events, NULL);
03428     if (result < 0) {
03429         rb_sys_fail(0);
03430     }
03431 
03432     thread_debug("rb_thread_wait_fd_rw(%d, %s): done\n", fd, read ? "read" : "write");
03433 }
03434 
03435 void
03436 rb_thread_wait_fd(int fd)
03437 {
03438     rb_thread_wait_fd_rw(fd, 1);
03439 }
03440 
03441 int
03442 rb_thread_fd_writable(int fd)
03443 {
03444     rb_thread_wait_fd_rw(fd, 0);
03445     return TRUE;
03446 }
03447 
03448 int
03449 rb_thread_select(int max, fd_set * read, fd_set * write, fd_set * except,
03450                  struct timeval *timeout)
03451 {
03452     rb_fdset_t fdsets[3];
03453     rb_fdset_t *rfds = NULL;
03454     rb_fdset_t *wfds = NULL;
03455     rb_fdset_t *efds = NULL;
03456     int retval;
03457 
03458     if (read) {
03459         rfds = &fdsets[0];
03460         rb_fd_init(rfds);
03461         rb_fd_copy(rfds, read, max);
03462     }
03463     if (write) {
03464         wfds = &fdsets[1];
03465         rb_fd_init(wfds);
03466         rb_fd_copy(wfds, write, max);
03467     }
03468     if (except) {
03469         efds = &fdsets[2];
03470         rb_fd_init(efds);
03471         rb_fd_copy(efds, except, max);
03472     }
03473 
03474     retval = rb_thread_fd_select(max, rfds, wfds, efds, timeout);
03475 
03476     if (rfds) {
03477         rb_fd_rcopy(read, rfds);
03478         rb_fd_term(rfds);
03479     }
03480     if (wfds) {
03481         rb_fd_rcopy(write, wfds);
03482         rb_fd_term(wfds);
03483     }
03484     if (efds) {
03485         rb_fd_rcopy(except, efds);
03486         rb_fd_term(efds);
03487     }
03488 
03489     return retval;
03490 }
03491 
03492 int
03493 rb_thread_fd_select(int max, rb_fdset_t * read, rb_fdset_t * write, rb_fdset_t * except,
03494                     struct timeval *timeout)
03495 {
03496     if (!read && !write && !except) {
03497         if (!timeout) {
03498             rb_thread_sleep_forever();
03499             return 0;
03500         }
03501         rb_thread_wait_for(*timeout);
03502         return 0;
03503     }
03504 
03505     if (read) {
03506         rb_fd_resize(max - 1, read);
03507     }
03508     if (write) {
03509         rb_fd_resize(max - 1, write);
03510     }
03511     if (except) {
03512         rb_fd_resize(max - 1, except);
03513     }
03514     return do_select(max, read, write, except, timeout);
03515 }
03516 
03517 /*
03518  * poll() is supported by many OSes, but so far Linux is the only
03519  * one we know of that supports using poll() in all places select()
03520  * would work.
03521  */
03522 #if defined(HAVE_POLL) && defined(__linux__)
03523 #  define USE_POLL
03524 #endif
03525 
03526 #ifdef USE_POLL
03527 
03528 /* The same with linux kernel. TODO: make platform independent definition. */
03529 #define POLLIN_SET (POLLRDNORM | POLLRDBAND | POLLIN | POLLHUP | POLLERR)
03530 #define POLLOUT_SET (POLLWRBAND | POLLWRNORM | POLLOUT | POLLERR)
03531 #define POLLEX_SET (POLLPRI)
03532 
03533 #ifndef HAVE_PPOLL
03534 /* TODO: don't ignore sigmask */
03535 int
03536 ppoll(struct pollfd *fds, nfds_t nfds,
03537       const struct timespec *ts, const sigset_t *sigmask)
03538 {
03539     int timeout_ms;
03540 
03541     if (ts) {
03542         int tmp, tmp2;
03543 
03544         if (ts->tv_sec > TIMET_MAX/1000)
03545             timeout_ms = -1;
03546         else {
03547             tmp = ts->tv_sec * 1000;
03548             tmp2 = ts->tv_nsec / (1000 * 1000);
03549             if (TIMET_MAX - tmp < tmp2)
03550                 timeout_ms = -1;
03551             else
03552                 timeout_ms = tmp + tmp2;
03553         }
03554     }
03555     else
03556         timeout_ms = -1;
03557 
03558     return poll(fds, nfds, timeout_ms);
03559 }
03560 #endif
03561 
03562 /*
03563  * returns a mask of events
03564  */
03565 int
03566 rb_wait_for_single_fd(int fd, int events, struct timeval *tv)
03567 {
03568     struct pollfd fds;
03569     int result = 0, lerrno;
03570     double limit = 0;
03571     struct timespec ts;
03572     struct timespec *timeout = NULL;
03573     rb_thread_t *th = GET_THREAD();
03574 
03575     if (tv) {
03576         ts.tv_sec = tv->tv_sec;
03577         ts.tv_nsec = tv->tv_usec * 1000;
03578         limit = timeofday();
03579         limit += (double)tv->tv_sec + (double)tv->tv_usec * 1e-6;
03580         timeout = &ts;
03581     }
03582 
03583     fds.fd = fd;
03584     fds.events = (short)events;
03585 
03586 retry:
03587     lerrno = 0;
03588     BLOCKING_REGION({
03589         result = ppoll(&fds, 1, timeout, NULL);
03590         if (result < 0) lerrno = errno;
03591     }, ubf_select, th, FALSE);
03592 
03593     RUBY_VM_CHECK_INTS_BLOCKING(th);
03594 
03595     if (result < 0) {
03596         errno = lerrno;
03597         switch (errno) {
03598           case EINTR:
03599 #ifdef ERESTART
03600           case ERESTART:
03601 #endif
03602             if (timeout) {
03603                 double d = limit - timeofday();
03604 
03605                 ts.tv_sec = (long)d;
03606                 ts.tv_nsec = (long)((d - (double)ts.tv_sec) * 1e9);
03607                 if (ts.tv_sec < 0)
03608                     ts.tv_sec = 0;
03609                 if (ts.tv_nsec < 0)
03610                     ts.tv_nsec = 0;
03611             }
03612             goto retry;
03613         }
03614         return -1;
03615     }
03616 
03617     if (fds.revents & POLLNVAL) {
03618         errno = EBADF;
03619         return -1;
03620     }
03621 
03622     /*
03623      * POLLIN, POLLOUT have a different meanings from select(2)'s read/write bit.
03624      * Therefore we need fix it up.
03625      */
03626     result = 0;
03627     if (fds.revents & POLLIN_SET)
03628         result |= RB_WAITFD_IN;
03629     if (fds.revents & POLLOUT_SET)
03630         result |= RB_WAITFD_OUT;
03631     if (fds.revents & POLLEX_SET)
03632         result |= RB_WAITFD_PRI;
03633 
03634     return result;
03635 }
03636 #else /* ! USE_POLL - implement rb_io_poll_fd() using select() */
03637 static rb_fdset_t *
03638 init_set_fd(int fd, rb_fdset_t *fds)
03639 {
03640     rb_fd_init(fds);
03641     rb_fd_set(fd, fds);
03642 
03643     return fds;
03644 }
03645 
03646 struct select_args {
03647     union {
03648         int fd;
03649         int error;
03650     } as;
03651     rb_fdset_t *read;
03652     rb_fdset_t *write;
03653     rb_fdset_t *except;
03654     struct timeval *tv;
03655 };
03656 
03657 static VALUE
03658 select_single(VALUE ptr)
03659 {
03660     struct select_args *args = (struct select_args *)ptr;
03661     int r;
03662 
03663     r = rb_thread_fd_select(args->as.fd + 1,
03664                             args->read, args->write, args->except, args->tv);
03665     if (r == -1)
03666         args->as.error = errno;
03667     if (r > 0) {
03668         r = 0;
03669         if (args->read && rb_fd_isset(args->as.fd, args->read))
03670             r |= RB_WAITFD_IN;
03671         if (args->write && rb_fd_isset(args->as.fd, args->write))
03672             r |= RB_WAITFD_OUT;
03673         if (args->except && rb_fd_isset(args->as.fd, args->except))
03674             r |= RB_WAITFD_PRI;
03675     }
03676     return (VALUE)r;
03677 }
03678 
03679 static VALUE
03680 select_single_cleanup(VALUE ptr)
03681 {
03682     struct select_args *args = (struct select_args *)ptr;
03683 
03684     if (args->read) rb_fd_term(args->read);
03685     if (args->write) rb_fd_term(args->write);
03686     if (args->except) rb_fd_term(args->except);
03687 
03688     return (VALUE)-1;
03689 }
03690 
03691 int
03692 rb_wait_for_single_fd(int fd, int events, struct timeval *tv)
03693 {
03694     rb_fdset_t rfds, wfds, efds;
03695     struct select_args args;
03696     int r;
03697     VALUE ptr = (VALUE)&args;
03698 
03699     args.as.fd = fd;
03700     args.read = (events & RB_WAITFD_IN) ? init_set_fd(fd, &rfds) : NULL;
03701     args.write = (events & RB_WAITFD_OUT) ? init_set_fd(fd, &wfds) : NULL;
03702     args.except = (events & RB_WAITFD_PRI) ? init_set_fd(fd, &efds) : NULL;
03703     args.tv = tv;
03704 
03705     r = (int)rb_ensure(select_single, ptr, select_single_cleanup, ptr);
03706     if (r == -1)
03707         errno = args.as.error;
03708 
03709     return r;
03710 }
03711 #endif /* ! USE_POLL */
03712 
03713 /*
03714  * for GC
03715  */
03716 
03717 #ifdef USE_CONSERVATIVE_STACK_END
03718 void
03719 rb_gc_set_stack_end(VALUE **stack_end_p)
03720 {
03721     VALUE stack_end;
03722     *stack_end_p = &stack_end;
03723 }
03724 #endif
03725 
03726 
03727 /*
03728  *
03729  */
03730 
03731 void
03732 rb_threadptr_check_signal(rb_thread_t *mth)
03733 {
03734     /* mth must be main_thread */
03735     if (rb_signal_buff_size() > 0) {
03736         /* wakeup main thread */
03737         rb_threadptr_trap_interrupt(mth);
03738     }
03739 }
03740 
03741 static void
03742 timer_thread_function(void *arg)
03743 {
03744     rb_vm_t *vm = GET_VM(); /* TODO: fix me for Multi-VM */
03745 
03746     /*
03747      * Tricky: thread_destruct_lock doesn't close a race against
03748      * vm->running_thread switch. however it guarantee th->running_thread
03749      * point to valid pointer or NULL.
03750      */
03751     native_mutex_lock(&vm->thread_destruct_lock);
03752     /* for time slice */
03753     if (vm->running_thread)
03754         RUBY_VM_SET_TIMER_INTERRUPT(vm->running_thread);
03755     native_mutex_unlock(&vm->thread_destruct_lock);
03756 
03757     /* check signal */
03758     rb_threadptr_check_signal(vm->main_thread);
03759 
03760 #if 0
03761     /* prove profiler */
03762     if (vm->prove_profile.enable) {
03763         rb_thread_t *th = vm->running_thread;
03764 
03765         if (vm->during_gc) {
03766             /* GC prove profiling */
03767         }
03768     }
03769 #endif
03770 }
03771 
03772 void
03773 rb_thread_stop_timer_thread(int close_anyway)
03774 {
03775     if (timer_thread_id && native_stop_timer_thread(close_anyway)) {
03776         native_reset_timer_thread();
03777     }
03778 }
03779 
03780 void
03781 rb_thread_reset_timer_thread(void)
03782 {
03783     native_reset_timer_thread();
03784 }
03785 
03786 void
03787 rb_thread_start_timer_thread(void)
03788 {
03789     system_working = 1;
03790     rb_thread_create_timer_thread();
03791 }
03792 
03793 static int
03794 clear_coverage_i(st_data_t key, st_data_t val, st_data_t dummy)
03795 {
03796     int i;
03797     VALUE lines = (VALUE)val;
03798 
03799     for (i = 0; i < RARRAY_LEN(lines); i++) {
03800         if (RARRAY_PTR(lines)[i] != Qnil) {
03801             RARRAY_PTR(lines)[i] = INT2FIX(0);
03802         }
03803     }
03804     return ST_CONTINUE;
03805 }
03806 
03807 static void
03808 clear_coverage(void)
03809 {
03810     VALUE coverages = rb_get_coverages();
03811     if (RTEST(coverages)) {
03812         st_foreach(RHASH_TBL(coverages), clear_coverage_i, 0);
03813     }
03814 }
03815 
03816 static void
03817 rb_thread_atfork_internal(int (*atfork)(st_data_t, st_data_t, st_data_t))
03818 {
03819     rb_thread_t *th = GET_THREAD();
03820     rb_vm_t *vm = th->vm;
03821     VALUE thval = th->self;
03822     vm->main_thread = th;
03823 
03824     gvl_atfork(th->vm);
03825     st_foreach(vm->living_threads, atfork, (st_data_t)th);
03826     st_clear(vm->living_threads);
03827     st_insert(vm->living_threads, thval, (st_data_t)th->thread_id);
03828     vm->sleeper = 0;
03829     clear_coverage();
03830 }
03831 
03832 static int
03833 terminate_atfork_i(st_data_t key, st_data_t val, st_data_t current_th)
03834 {
03835     VALUE thval = key;
03836     rb_thread_t *th;
03837     GetThreadPtr(thval, th);
03838 
03839     if (th != (rb_thread_t *)current_th) {
03840         rb_mutex_abandon_keeping_mutexes(th);
03841         rb_mutex_abandon_locking_mutex(th);
03842         thread_cleanup_func(th, TRUE);
03843     }
03844     return ST_CONTINUE;
03845 }
03846 
03847 void
03848 rb_thread_atfork(void)
03849 {
03850     rb_thread_atfork_internal(terminate_atfork_i);
03851     GET_THREAD()->join_list = NULL;
03852 
03853     /* We don't want reproduce CVE-2003-0900. */
03854     rb_reset_random_seed();
03855 }
03856 
03857 static int
03858 terminate_atfork_before_exec_i(st_data_t key, st_data_t val, st_data_t current_th)
03859 {
03860     VALUE thval = key;
03861     rb_thread_t *th;
03862     GetThreadPtr(thval, th);
03863 
03864     if (th != (rb_thread_t *)current_th) {
03865         thread_cleanup_func_before_exec(th);
03866     }
03867     return ST_CONTINUE;
03868 }
03869 
03870 void
03871 rb_thread_atfork_before_exec(void)
03872 {
03873     rb_thread_atfork_internal(terminate_atfork_before_exec_i);
03874 }
03875 
03876 struct thgroup {
03877     int enclosed;
03878     VALUE group;
03879 };
03880 
03881 static size_t
03882 thgroup_memsize(const void *ptr)
03883 {
03884     return ptr ? sizeof(struct thgroup) : 0;
03885 }
03886 
03887 static const rb_data_type_t thgroup_data_type = {
03888     "thgroup",
03889     {NULL, RUBY_TYPED_DEFAULT_FREE, thgroup_memsize,},
03890 };
03891 
03892 /*
03893  * Document-class: ThreadGroup
03894  *
03895  *  <code>ThreadGroup</code> provides a means of keeping track of a number of
03896  *  threads as a group. A <code>Thread</code> can belong to only one
03897  *  <code>ThreadGroup</code> at a time; adding a thread to a new group will
03898  *  remove it from any previous group.
03899  *
03900  *  Newly created threads belong to the same group as the thread from which they
03901  *  were created.
03902  */
03903 
03904 /*
03905  * Document-const: Default
03906  *
03907  *  The default ThreadGroup created when Ruby starts; all Threads belong to it
03908  *  by default.
03909  */
03910 static VALUE
03911 thgroup_s_alloc(VALUE klass)
03912 {
03913     VALUE group;
03914     struct thgroup *data;
03915 
03916     group = TypedData_Make_Struct(klass, struct thgroup, &thgroup_data_type, data);
03917     data->enclosed = 0;
03918     data->group = group;
03919 
03920     return group;
03921 }
03922 
03923 struct thgroup_list_params {
03924     VALUE ary;
03925     VALUE group;
03926 };
03927 
03928 static int
03929 thgroup_list_i(st_data_t key, st_data_t val, st_data_t data)
03930 {
03931     VALUE thread = (VALUE)key;
03932     VALUE ary = ((struct thgroup_list_params *)data)->ary;
03933     VALUE group = ((struct thgroup_list_params *)data)->group;
03934     rb_thread_t *th;
03935     GetThreadPtr(thread, th);
03936 
03937     if (th->thgroup == group) {
03938         rb_ary_push(ary, thread);
03939     }
03940     return ST_CONTINUE;
03941 }
03942 
03943 /*
03944  *  call-seq:
03945  *     thgrp.list   -> array
03946  *
03947  *  Returns an array of all existing <code>Thread</code> objects that belong to
03948  *  this group.
03949  *
03950  *     ThreadGroup::Default.list   #=> [#<Thread:0x401bdf4c run>]
03951  */
03952 
03953 static VALUE
03954 thgroup_list(VALUE group)
03955 {
03956     VALUE ary = rb_ary_new();
03957     struct thgroup_list_params param;
03958 
03959     param.ary = ary;
03960     param.group = group;
03961     st_foreach(GET_THREAD()->vm->living_threads, thgroup_list_i, (st_data_t) & param);
03962     return ary;
03963 }
03964 
03965 
03966 /*
03967  *  call-seq:
03968  *     thgrp.enclose   -> thgrp
03969  *
03970  *  Prevents threads from being added to or removed from the receiving
03971  *  <code>ThreadGroup</code>. New threads can still be started in an enclosed
03972  *  <code>ThreadGroup</code>.
03973  *
03974  *     ThreadGroup::Default.enclose        #=> #<ThreadGroup:0x4029d914>
03975  *     thr = Thread::new { Thread.stop }   #=> #<Thread:0x402a7210 sleep>
03976  *     tg = ThreadGroup::new               #=> #<ThreadGroup:0x402752d4>
03977  *     tg.add thr
03978  *
03979  *  <em>produces:</em>
03980  *
03981  *     ThreadError: can't move from the enclosed thread group
03982  */
03983 
03984 static VALUE
03985 thgroup_enclose(VALUE group)
03986 {
03987     struct thgroup *data;
03988 
03989     TypedData_Get_Struct(group, struct thgroup, &thgroup_data_type, data);
03990     data->enclosed = 1;
03991 
03992     return group;
03993 }
03994 
03995 
03996 /*
03997  *  call-seq:
03998  *     thgrp.enclosed?   -> true or false
03999  *
04000  *  Returns <code>true</code> if <em>thgrp</em> is enclosed. See also
04001  *  ThreadGroup#enclose.
04002  */
04003 
04004 static VALUE
04005 thgroup_enclosed_p(VALUE group)
04006 {
04007     struct thgroup *data;
04008 
04009     TypedData_Get_Struct(group, struct thgroup, &thgroup_data_type, data);
04010     if (data->enclosed)
04011         return Qtrue;
04012     return Qfalse;
04013 }
04014 
04015 
04016 /*
04017  *  call-seq:
04018  *     thgrp.add(thread)   -> thgrp
04019  *
04020  *  Adds the given <em>thread</em> to this group, removing it from any other
04021  *  group to which it may have previously belonged.
04022  *
04023  *     puts "Initial group is #{ThreadGroup::Default.list}"
04024  *     tg = ThreadGroup.new
04025  *     t1 = Thread.new { sleep }
04026  *     t2 = Thread.new { sleep }
04027  *     puts "t1 is #{t1}"
04028  *     puts "t2 is #{t2}"
04029  *     tg.add(t1)
04030  *     puts "Initial group now #{ThreadGroup::Default.list}"
04031  *     puts "tg group now #{tg.list}"
04032  *
04033  *  <em>produces:</em>
04034  *
04035  *     Initial group is #<Thread:0x401bdf4c>
04036  *     t1 is #<Thread:0x401b3c90>
04037  *     t2 is #<Thread:0x401b3c18>
04038  *     Initial group now #<Thread:0x401b3c18>#<Thread:0x401bdf4c>
04039  *     tg group now #<Thread:0x401b3c90>
04040  */
04041 
04042 static VALUE
04043 thgroup_add(VALUE group, VALUE thread)
04044 {
04045     rb_thread_t *th;
04046     struct thgroup *data;
04047 
04048     rb_secure(4);
04049     GetThreadPtr(thread, th);
04050 
04051     if (OBJ_FROZEN(group)) {
04052         rb_raise(rb_eThreadError, "can't move to the frozen thread group");
04053     }
04054     TypedData_Get_Struct(group, struct thgroup, &thgroup_data_type, data);
04055     if (data->enclosed) {
04056         rb_raise(rb_eThreadError, "can't move to the enclosed thread group");
04057     }
04058 
04059     if (!th->thgroup) {
04060         return Qnil;
04061     }
04062 
04063     if (OBJ_FROZEN(th->thgroup)) {
04064         rb_raise(rb_eThreadError, "can't move from the frozen thread group");
04065     }
04066     TypedData_Get_Struct(th->thgroup, struct thgroup, &thgroup_data_type, data);
04067     if (data->enclosed) {
04068         rb_raise(rb_eThreadError,
04069                  "can't move from the enclosed thread group");
04070     }
04071 
04072     th->thgroup = group;
04073     return group;
04074 }
04075 
04076 
04077 /*
04078  *  Document-class: Mutex
04079  *
04080  *  Mutex implements a simple semaphore that can be used to coordinate access to
04081  *  shared data from multiple concurrent threads.
04082  *
04083  *  Example:
04084  *
04085  *    require 'thread'
04086  *    semaphore = Mutex.new
04087  *
04088  *    a = Thread.new {
04089  *      semaphore.synchronize {
04090  *        # access shared resource
04091  *      }
04092  *    }
04093  *
04094  *    b = Thread.new {
04095  *      semaphore.synchronize {
04096  *        # access shared resource
04097  *      }
04098  *    }
04099  *
04100  */
04101 
04102 #define GetMutexPtr(obj, tobj) \
04103     TypedData_Get_Struct((obj), rb_mutex_t, &mutex_data_type, (tobj))
04104 
04105 #define mutex_mark NULL
04106 
04107 static void
04108 mutex_free(void *ptr)
04109 {
04110     if (ptr) {
04111         rb_mutex_t *mutex = ptr;
04112         if (mutex->th) {
04113             /* rb_warn("free locked mutex"); */
04114             const char *err = rb_mutex_unlock_th(mutex, mutex->th);
04115             if (err) rb_bug("%s", err);
04116         }
04117         native_mutex_destroy(&mutex->lock);
04118         native_cond_destroy(&mutex->cond);
04119     }
04120     ruby_xfree(ptr);
04121 }
04122 
04123 static size_t
04124 mutex_memsize(const void *ptr)
04125 {
04126     return ptr ? sizeof(rb_mutex_t) : 0;
04127 }
04128 
04129 static const rb_data_type_t mutex_data_type = {
04130     "mutex",
04131     {mutex_mark, mutex_free, mutex_memsize,},
04132 };
04133 
04134 VALUE
04135 rb_obj_is_mutex(VALUE obj)
04136 {
04137     if (rb_typeddata_is_kind_of(obj, &mutex_data_type)) {
04138         return Qtrue;
04139     }
04140     else {
04141         return Qfalse;
04142     }
04143 }
04144 
04145 static VALUE
04146 mutex_alloc(VALUE klass)
04147 {
04148     VALUE volatile obj;
04149     rb_mutex_t *mutex;
04150 
04151     obj = TypedData_Make_Struct(klass, rb_mutex_t, &mutex_data_type, mutex);
04152     native_mutex_initialize(&mutex->lock);
04153     native_cond_initialize(&mutex->cond, RB_CONDATTR_CLOCK_MONOTONIC);
04154     return obj;
04155 }
04156 
04157 /*
04158  *  call-seq:
04159  *     Mutex.new   -> mutex
04160  *
04161  *  Creates a new Mutex
04162  */
04163 static VALUE
04164 mutex_initialize(VALUE self)
04165 {
04166     return self;
04167 }
04168 
04169 VALUE
04170 rb_mutex_new(void)
04171 {
04172     return mutex_alloc(rb_cMutex);
04173 }
04174 
04175 /*
04176  * call-seq:
04177  *    mutex.locked?  -> true or false
04178  *
04179  * Returns +true+ if this lock is currently held by some thread.
04180  */
04181 VALUE
04182 rb_mutex_locked_p(VALUE self)
04183 {
04184     rb_mutex_t *mutex;
04185     GetMutexPtr(self, mutex);
04186     return mutex->th ? Qtrue : Qfalse;
04187 }
04188 
04189 static void
04190 mutex_locked(rb_thread_t *th, VALUE self)
04191 {
04192     rb_mutex_t *mutex;
04193     GetMutexPtr(self, mutex);
04194 
04195     if (th->keeping_mutexes) {
04196         mutex->next_mutex = th->keeping_mutexes;
04197     }
04198     th->keeping_mutexes = mutex;
04199 }
04200 
04201 /*
04202  * call-seq:
04203  *    mutex.try_lock  -> true or false
04204  *
04205  * Attempts to obtain the lock and returns immediately. Returns +true+ if the
04206  * lock was granted.
04207  */
04208 VALUE
04209 rb_mutex_trylock(VALUE self)
04210 {
04211     rb_mutex_t *mutex;
04212     VALUE locked = Qfalse;
04213     GetMutexPtr(self, mutex);
04214 
04215     native_mutex_lock(&mutex->lock);
04216     if (mutex->th == 0) {
04217         mutex->th = GET_THREAD();
04218         locked = Qtrue;
04219 
04220         mutex_locked(GET_THREAD(), self);
04221     }
04222     native_mutex_unlock(&mutex->lock);
04223 
04224     return locked;
04225 }
04226 
04227 static int
04228 lock_func(rb_thread_t *th, rb_mutex_t *mutex, int timeout_ms)
04229 {
04230     int interrupted = 0;
04231     int err = 0;
04232 
04233     mutex->cond_waiting++;
04234     for (;;) {
04235         if (!mutex->th) {
04236             mutex->th = th;
04237             break;
04238         }
04239         if (RUBY_VM_INTERRUPTED(th)) {
04240             interrupted = 1;
04241             break;
04242         }
04243         if (err == ETIMEDOUT) {
04244             interrupted = 2;
04245             break;
04246         }
04247 
04248         if (timeout_ms) {
04249             struct timespec timeout_rel;
04250             struct timespec timeout;
04251 
04252             timeout_rel.tv_sec = 0;
04253             timeout_rel.tv_nsec = timeout_ms * 1000 * 1000;
04254             timeout = native_cond_timeout(&mutex->cond, timeout_rel);
04255             err = native_cond_timedwait(&mutex->cond, &mutex->lock, &timeout);
04256         }
04257         else {
04258             native_cond_wait(&mutex->cond, &mutex->lock);
04259             err = 0;
04260         }
04261     }
04262     mutex->cond_waiting--;
04263 
04264     return interrupted;
04265 }
04266 
04267 static void
04268 lock_interrupt(void *ptr)
04269 {
04270     rb_mutex_t *mutex = (rb_mutex_t *)ptr;
04271     native_mutex_lock(&mutex->lock);
04272     if (mutex->cond_waiting > 0)
04273         native_cond_broadcast(&mutex->cond);
04274     native_mutex_unlock(&mutex->lock);
04275 }
04276 
04277 /*
04278  * At maximum, only one thread can use cond_timedwait and watch deadlock
04279  * periodically. Multiple polling thread (i.e. concurrent deadlock check)
04280  * introduces new race conditions. [Bug #6278] [ruby-core:44275]
04281  */
04282 static const rb_thread_t *patrol_thread = NULL;
04283 
04284 /*
04285  * call-seq:
04286  *    mutex.lock  -> self
04287  *
04288  * Attempts to grab the lock and waits if it isn't available.
04289  * Raises +ThreadError+ if +mutex+ was locked by the current thread.
04290  */
04291 VALUE
04292 rb_mutex_lock(VALUE self)
04293 {
04294     rb_thread_t *th = GET_THREAD();
04295     rb_mutex_t *mutex;
04296     GetMutexPtr(self, mutex);
04297 
04298     /* When running trap handler */
04299     if (!mutex->allow_trap && th->interrupt_mask & TRAP_INTERRUPT_MASK) {
04300         rb_raise(rb_eThreadError, "can't be called from trap context");
04301     }
04302 
04303     if (rb_mutex_trylock(self) == Qfalse) {
04304         if (mutex->th == GET_THREAD()) {
04305             rb_raise(rb_eThreadError, "deadlock; recursive locking");
04306         }
04307 
04308         while (mutex->th != th) {
04309             int interrupted;
04310             enum rb_thread_status prev_status = th->status;
04311             volatile int timeout_ms = 0;
04312             struct rb_unblock_callback oldubf;
04313 
04314             set_unblock_function(th, lock_interrupt, mutex, &oldubf, FALSE);
04315             th->status = THREAD_STOPPED_FOREVER;
04316             th->locking_mutex = self;
04317 
04318             native_mutex_lock(&mutex->lock);
04319             th->vm->sleeper++;
04320             /*
04321              * Carefully! while some contended threads are in lock_func(),
04322              * vm->sleepr is unstable value. we have to avoid both deadlock
04323              * and busy loop.
04324              */
04325             if ((vm_living_thread_num(th->vm) == th->vm->sleeper) &&
04326                 !patrol_thread) {
04327                 timeout_ms = 100;
04328                 patrol_thread = th;
04329             }
04330 
04331             GVL_UNLOCK_BEGIN();
04332             interrupted = lock_func(th, mutex, (int)timeout_ms);
04333             native_mutex_unlock(&mutex->lock);
04334             GVL_UNLOCK_END();
04335 
04336             if (patrol_thread == th)
04337                 patrol_thread = NULL;
04338 
04339             reset_unblock_function(th, &oldubf);
04340 
04341             th->locking_mutex = Qfalse;
04342             if (mutex->th && interrupted == 2) {
04343                 rb_check_deadlock(th->vm);
04344             }
04345             if (th->status == THREAD_STOPPED_FOREVER) {
04346                 th->status = prev_status;
04347             }
04348             th->vm->sleeper--;
04349 
04350             if (mutex->th == th) mutex_locked(th, self);
04351 
04352             if (interrupted) {
04353                 RUBY_VM_CHECK_INTS_BLOCKING(th);
04354             }
04355         }
04356     }
04357     return self;
04358 }
04359 
04360 /*
04361  * call-seq:
04362  *    mutex.owned?  -> true or false
04363  *
04364  * Returns +true+ if this lock is currently held by current thread.
04365  * <em>This API is experimental, and subject to change.</em>
04366  */
04367 VALUE
04368 rb_mutex_owned_p(VALUE self)
04369 {
04370     VALUE owned = Qfalse;
04371     rb_thread_t *th = GET_THREAD();
04372     rb_mutex_t *mutex;
04373 
04374     GetMutexPtr(self, mutex);
04375 
04376     if (mutex->th == th)
04377         owned = Qtrue;
04378 
04379     return owned;
04380 }
04381 
04382 static const char *
04383 rb_mutex_unlock_th(rb_mutex_t *mutex, rb_thread_t volatile *th)
04384 {
04385     const char *err = NULL;
04386 
04387     native_mutex_lock(&mutex->lock);
04388 
04389     if (mutex->th == 0) {
04390         err = "Attempt to unlock a mutex which is not locked";
04391     }
04392     else if (mutex->th != th) {
04393         err = "Attempt to unlock a mutex which is locked by another thread";
04394     }
04395     else {
04396         mutex->th = 0;
04397         if (mutex->cond_waiting > 0)
04398             native_cond_signal(&mutex->cond);
04399     }
04400 
04401     native_mutex_unlock(&mutex->lock);
04402 
04403     if (!err) {
04404         rb_mutex_t *volatile *th_mutex = &th->keeping_mutexes;
04405         while (*th_mutex != mutex) {
04406             th_mutex = &(*th_mutex)->next_mutex;
04407         }
04408         *th_mutex = mutex->next_mutex;
04409         mutex->next_mutex = NULL;
04410     }
04411 
04412     return err;
04413 }
04414 
04415 /*
04416  * call-seq:
04417  *    mutex.unlock    -> self
04418  *
04419  * Releases the lock.
04420  * Raises +ThreadError+ if +mutex+ wasn't locked by the current thread.
04421  */
04422 VALUE
04423 rb_mutex_unlock(VALUE self)
04424 {
04425     const char *err;
04426     rb_mutex_t *mutex;
04427     GetMutexPtr(self, mutex);
04428 
04429     err = rb_mutex_unlock_th(mutex, GET_THREAD());
04430     if (err) rb_raise(rb_eThreadError, "%s", err);
04431 
04432     return self;
04433 }
04434 
04435 static void
04436 rb_mutex_abandon_keeping_mutexes(rb_thread_t *th)
04437 {
04438     if (th->keeping_mutexes) {
04439         rb_mutex_abandon_all(th->keeping_mutexes);
04440     }
04441     th->keeping_mutexes = NULL;
04442 }
04443 
04444 static void
04445 rb_mutex_abandon_locking_mutex(rb_thread_t *th)
04446 {
04447     rb_mutex_t *mutex;
04448 
04449     if (!th->locking_mutex) return;
04450 
04451     GetMutexPtr(th->locking_mutex, mutex);
04452     if (mutex->th == th)
04453         rb_mutex_abandon_all(mutex);
04454     th->locking_mutex = Qfalse;
04455 }
04456 
04457 static void
04458 rb_mutex_abandon_all(rb_mutex_t *mutexes)
04459 {
04460     rb_mutex_t *mutex;
04461 
04462     while (mutexes) {
04463         mutex = mutexes;
04464         mutexes = mutex->next_mutex;
04465         mutex->th = 0;
04466         mutex->next_mutex = 0;
04467     }
04468 }
04469 
04470 static VALUE
04471 rb_mutex_sleep_forever(VALUE time)
04472 {
04473     sleep_forever(GET_THREAD(), 1, 0); /* permit spurious check */
04474     return Qnil;
04475 }
04476 
04477 static VALUE
04478 rb_mutex_wait_for(VALUE time)
04479 {
04480     struct timeval *t = (struct timeval *)time;
04481     sleep_timeval(GET_THREAD(), *t, 0); /* permit spurious check */
04482     return Qnil;
04483 }
04484 
04485 VALUE
04486 rb_mutex_sleep(VALUE self, VALUE timeout)
04487 {
04488     time_t beg, end;
04489     struct timeval t;
04490 
04491     if (!NIL_P(timeout)) {
04492         t = rb_time_interval(timeout);
04493     }
04494     rb_mutex_unlock(self);
04495     beg = time(0);
04496     if (NIL_P(timeout)) {
04497         rb_ensure(rb_mutex_sleep_forever, Qnil, rb_mutex_lock, self);
04498     }
04499     else {
04500         rb_ensure(rb_mutex_wait_for, (VALUE)&t, rb_mutex_lock, self);
04501     }
04502     end = time(0) - beg;
04503     return INT2FIX(end);
04504 }
04505 
04506 /*
04507  * call-seq:
04508  *    mutex.sleep(timeout = nil)    -> number
04509  *
04510  * Releases the lock and sleeps +timeout+ seconds if it is given and
04511  * non-nil or forever.  Raises +ThreadError+ if +mutex+ wasn't locked by
04512  * the current thread.
04513  *
04514  * Note that this method can wakeup without explicit Thread#wakeup call.
04515  * For example, receiving signal and so on.
04516  */
04517 static VALUE
04518 mutex_sleep(int argc, VALUE *argv, VALUE self)
04519 {
04520     VALUE timeout;
04521 
04522     rb_scan_args(argc, argv, "01", &timeout);
04523     return rb_mutex_sleep(self, timeout);
04524 }
04525 
04526 /*
04527  * call-seq:
04528  *    mutex.synchronize { ... }    -> result of the block
04529  *
04530  * Obtains a lock, runs the block, and releases the lock when the block
04531  * completes.  See the example under +Mutex+.
04532  */
04533 
04534 VALUE
04535 rb_mutex_synchronize(VALUE mutex, VALUE (*func)(VALUE arg), VALUE arg)
04536 {
04537     rb_mutex_lock(mutex);
04538     return rb_ensure(func, arg, rb_mutex_unlock, mutex);
04539 }
04540 
04541 /*
04542  * call-seq:
04543  *    mutex.synchronize { ... }    -> result of the block
04544  *
04545  * Obtains a lock, runs the block, and releases the lock when the block
04546  * completes.  See the example under +Mutex+.
04547  */
04548 static VALUE
04549 rb_mutex_synchronize_m(VALUE self, VALUE args)
04550 {
04551     if (!rb_block_given_p()) {
04552         rb_raise(rb_eThreadError, "must be called with a block");
04553     }
04554 
04555     return rb_mutex_synchronize(self, rb_yield, Qundef);
04556 }
04557 
04558 void rb_mutex_allow_trap(VALUE self, int val)
04559 {
04560     rb_mutex_t *m;
04561     GetMutexPtr(self, m);
04562 
04563     m->allow_trap = val;
04564 }
04565 
04566 /*
04567  * Document-class: ThreadShield
04568  */
04569 static void
04570 thread_shield_mark(void *ptr)
04571 {
04572     rb_gc_mark((VALUE)ptr);
04573 }
04574 
04575 static const rb_data_type_t thread_shield_data_type = {
04576     "thread_shield",
04577     {thread_shield_mark, 0, 0,},
04578 };
04579 
04580 static VALUE
04581 thread_shield_alloc(VALUE klass)
04582 {
04583     return TypedData_Wrap_Struct(klass, &thread_shield_data_type, (void *)mutex_alloc(0));
04584 }
04585 
04586 #define GetThreadShieldPtr(obj) ((VALUE)rb_check_typeddata((obj), &thread_shield_data_type))
04587 #define THREAD_SHIELD_WAITING_MASK (FL_USER0|FL_USER1|FL_USER2|FL_USER3|FL_USER4|FL_USER5|FL_USER6|FL_USER7|FL_USER8|FL_USER9|FL_USER10|FL_USER11|FL_USER12|FL_USER13|FL_USER14|FL_USER15|FL_USER16|FL_USER17|FL_USER18|FL_USER19)
04588 #define THREAD_SHIELD_WAITING_SHIFT (FL_USHIFT)
04589 #define rb_thread_shield_waiting(b) (int)((RBASIC(b)->flags&THREAD_SHIELD_WAITING_MASK)>>THREAD_SHIELD_WAITING_SHIFT)
04590 
04591 static inline void
04592 rb_thread_shield_waiting_inc(VALUE b)
04593 {
04594     unsigned int w = rb_thread_shield_waiting(b);
04595     w++;
04596     if (w > (unsigned int)(THREAD_SHIELD_WAITING_MASK>>THREAD_SHIELD_WAITING_SHIFT))
04597         rb_raise(rb_eRuntimeError, "waiting count overflow");
04598     RBASIC(b)->flags &= ~THREAD_SHIELD_WAITING_MASK;
04599     RBASIC(b)->flags |= ((VALUE)w << THREAD_SHIELD_WAITING_SHIFT);
04600 }
04601 
04602 static inline void
04603 rb_thread_shield_waiting_dec(VALUE b)
04604 {
04605     unsigned int w = rb_thread_shield_waiting(b);
04606     if (!w) rb_raise(rb_eRuntimeError, "waiting count underflow");
04607     w--;
04608     RBASIC(b)->flags &= ~THREAD_SHIELD_WAITING_MASK;
04609     RBASIC(b)->flags |= ((VALUE)w << THREAD_SHIELD_WAITING_SHIFT);
04610 }
04611 
04612 VALUE
04613 rb_thread_shield_new(void)
04614 {
04615     VALUE thread_shield = thread_shield_alloc(rb_cThreadShield);
04616     rb_mutex_lock((VALUE)DATA_PTR(thread_shield));
04617     return thread_shield;
04618 }
04619 
04620 /*
04621  * Wait a thread shield.
04622  *
04623  * Returns
04624  *  true:  acquired the thread shield
04625  *  false: the thread shield was destroyed and no other threads waiting
04626  *  nil:   the thread shield was destroyed but still in use
04627  */
04628 VALUE
04629 rb_thread_shield_wait(VALUE self)
04630 {
04631     VALUE mutex = GetThreadShieldPtr(self);
04632     rb_mutex_t *m;
04633 
04634     if (!mutex) return Qfalse;
04635     GetMutexPtr(mutex, m);
04636     if (m->th == GET_THREAD()) return Qnil;
04637     rb_thread_shield_waiting_inc(self);
04638     rb_mutex_lock(mutex);
04639     rb_thread_shield_waiting_dec(self);
04640     if (DATA_PTR(self)) return Qtrue;
04641     rb_mutex_unlock(mutex);
04642     return rb_thread_shield_waiting(self) > 0 ? Qnil : Qfalse;
04643 }
04644 
04645 /*
04646  * Release a thread shield, and return true if it has waiting threads.
04647  */
04648 VALUE
04649 rb_thread_shield_release(VALUE self)
04650 {
04651     VALUE mutex = GetThreadShieldPtr(self);
04652     rb_mutex_unlock(mutex);
04653     return rb_thread_shield_waiting(self) > 0 ? Qtrue : Qfalse;
04654 }
04655 
04656 /*
04657  * Release and destroy a thread shield, and return true if it has waiting threads.
04658  */
04659 VALUE
04660 rb_thread_shield_destroy(VALUE self)
04661 {
04662     VALUE mutex = GetThreadShieldPtr(self);
04663     DATA_PTR(self) = 0;
04664     rb_mutex_unlock(mutex);
04665     return rb_thread_shield_waiting(self) > 0 ? Qtrue : Qfalse;
04666 }
04667 
04668 /* variables for recursive traversals */
04669 static ID recursive_key;
04670 
04671 /*
04672  * Returns the current "recursive list" used to detect recursion.
04673  * This list is a hash table, unique for the current thread and for
04674  * the current __callee__.
04675  */
04676 
04677 static VALUE
04678 recursive_list_access(void)
04679 {
04680     volatile VALUE hash = rb_thread_local_aref(rb_thread_current(), recursive_key);
04681     VALUE sym = ID2SYM(rb_frame_this_func());
04682     VALUE list;
04683     if (NIL_P(hash) || !RB_TYPE_P(hash, T_HASH)) {
04684         hash = rb_hash_new();
04685         OBJ_UNTRUST(hash);
04686         rb_thread_local_aset(rb_thread_current(), recursive_key, hash);
04687         list = Qnil;
04688     }
04689     else {
04690         list = rb_hash_aref(hash, sym);
04691     }
04692     if (NIL_P(list) || !RB_TYPE_P(list, T_HASH)) {
04693         list = rb_hash_new();
04694         OBJ_UNTRUST(list);
04695         rb_hash_aset(hash, sym, list);
04696     }
04697     return list;
04698 }
04699 
04700 /*
04701  * Returns Qtrue iff obj_id (or the pair <obj, paired_obj>) is already
04702  * in the recursion list.
04703  * Assumes the recursion list is valid.
04704  */
04705 
04706 static VALUE
04707 recursive_check(VALUE list, VALUE obj_id, VALUE paired_obj_id)
04708 {
04709 #if SIZEOF_LONG == SIZEOF_VOIDP
04710   #define OBJ_ID_EQL(obj_id, other) ((obj_id) == (other))
04711 #elif SIZEOF_LONG_LONG == SIZEOF_VOIDP
04712   #define OBJ_ID_EQL(obj_id, other) (RB_TYPE_P((obj_id), T_BIGNUM) ? \
04713     rb_big_eql((obj_id), (other)) : ((obj_id) == (other)))
04714 #endif
04715 
04716     VALUE pair_list = rb_hash_lookup2(list, obj_id, Qundef);
04717     if (pair_list == Qundef)
04718         return Qfalse;
04719     if (paired_obj_id) {
04720         if (!RB_TYPE_P(pair_list, T_HASH)) {
04721             if (!OBJ_ID_EQL(paired_obj_id, pair_list))
04722                 return Qfalse;
04723         }
04724         else {
04725             if (NIL_P(rb_hash_lookup(pair_list, paired_obj_id)))
04726                 return Qfalse;
04727         }
04728     }
04729     return Qtrue;
04730 }
04731 
04732 /*
04733  * Pushes obj_id (or the pair <obj_id, paired_obj_id>) in the recursion list.
04734  * For a single obj_id, it sets list[obj_id] to Qtrue.
04735  * For a pair, it sets list[obj_id] to paired_obj_id if possible,
04736  * otherwise list[obj_id] becomes a hash like:
04737  *   {paired_obj_id_1 => true, paired_obj_id_2 => true, ... }
04738  * Assumes the recursion list is valid.
04739  */
04740 
04741 static void
04742 recursive_push(VALUE list, VALUE obj, VALUE paired_obj)
04743 {
04744     VALUE pair_list;
04745 
04746     if (!paired_obj) {
04747         rb_hash_aset(list, obj, Qtrue);
04748     }
04749     else if ((pair_list = rb_hash_lookup2(list, obj, Qundef)) == Qundef) {
04750         rb_hash_aset(list, obj, paired_obj);
04751     }
04752     else {
04753         if (!RB_TYPE_P(pair_list, T_HASH)){
04754             VALUE other_paired_obj = pair_list;
04755             pair_list = rb_hash_new();
04756             OBJ_UNTRUST(pair_list);
04757             rb_hash_aset(pair_list, other_paired_obj, Qtrue);
04758             rb_hash_aset(list, obj, pair_list);
04759         }
04760         rb_hash_aset(pair_list, paired_obj, Qtrue);
04761     }
04762 }
04763 
04764 /*
04765  * Pops obj_id (or the pair <obj_id, paired_obj_id>) from the recursion list.
04766  * For a pair, if list[obj_id] is a hash, then paired_obj_id is
04767  * removed from the hash and no attempt is made to simplify
04768  * list[obj_id] from {only_one_paired_id => true} to only_one_paired_id
04769  * Assumes the recursion list is valid.
04770  */
04771 
04772 static void
04773 recursive_pop(VALUE list, VALUE obj, VALUE paired_obj)
04774 {
04775     if (paired_obj) {
04776         VALUE pair_list = rb_hash_lookup2(list, obj, Qundef);
04777         if (pair_list == Qundef) {
04778             VALUE symname = rb_inspect(ID2SYM(rb_frame_this_func()));
04779             VALUE thrname = rb_inspect(rb_thread_current());
04780             rb_raise(rb_eTypeError, "invalid inspect_tbl pair_list for %s in %s",
04781                      StringValuePtr(symname), StringValuePtr(thrname));
04782         }
04783         if (RB_TYPE_P(pair_list, T_HASH)) {
04784             rb_hash_delete(pair_list, paired_obj);
04785             if (!RHASH_EMPTY_P(pair_list)) {
04786                 return; /* keep hash until is empty */
04787             }
04788         }
04789     }
04790     rb_hash_delete(list, obj);
04791 }
04792 
04793 struct exec_recursive_params {
04794     VALUE (*func) (VALUE, VALUE, int);
04795     VALUE list;
04796     VALUE obj;
04797     VALUE objid;
04798     VALUE pairid;
04799     VALUE arg;
04800 };
04801 
04802 static VALUE
04803 exec_recursive_i(VALUE tag, struct exec_recursive_params *p)
04804 {
04805     VALUE result = Qundef;
04806     int state;
04807 
04808     recursive_push(p->list, p->objid, p->pairid);
04809     PUSH_TAG();
04810     if ((state = EXEC_TAG()) == 0) {
04811         result = (*p->func)(p->obj, p->arg, FALSE);
04812     }
04813     POP_TAG();
04814     recursive_pop(p->list, p->objid, p->pairid);
04815     if (state)
04816         JUMP_TAG(state);
04817     return result;
04818 }
04819 
04820 /*
04821  * Calls func(obj, arg, recursive), where recursive is non-zero if the
04822  * current method is called recursively on obj, or on the pair <obj, pairid>
04823  * If outer is 0, then the innermost func will be called with recursive set
04824  * to Qtrue, otherwise the outermost func will be called. In the latter case,
04825  * all inner func are short-circuited by throw.
04826  * Implementation details: the value thrown is the recursive list which is
04827  * proper to the current method and unlikely to be catched anywhere else.
04828  * list[recursive_key] is used as a flag for the outermost call.
04829  */
04830 
04831 static VALUE
04832 exec_recursive(VALUE (*func) (VALUE, VALUE, int), VALUE obj, VALUE pairid, VALUE arg, int outer)
04833 {
04834     VALUE result = Qundef;
04835     struct exec_recursive_params p;
04836     int outermost;
04837     p.list = recursive_list_access();
04838     p.objid = rb_obj_id(obj);
04839     p.obj = obj;
04840     p.pairid = pairid;
04841     p.arg = arg;
04842     outermost = outer && !recursive_check(p.list, ID2SYM(recursive_key), 0);
04843 
04844     if (recursive_check(p.list, p.objid, pairid)) {
04845         if (outer && !outermost) {
04846             rb_throw_obj(p.list, p.list);
04847         }
04848         return (*func)(obj, arg, TRUE);
04849     }
04850     else {
04851         p.func = func;
04852 
04853         if (outermost) {
04854             recursive_push(p.list, ID2SYM(recursive_key), 0);
04855             result = rb_catch_obj(p.list, exec_recursive_i, (VALUE)&p);
04856             recursive_pop(p.list, ID2SYM(recursive_key), 0);
04857             if (result == p.list) {
04858                 result = (*func)(obj, arg, TRUE);
04859             }
04860         }
04861         else {
04862             result = exec_recursive_i(0, &p);
04863         }
04864     }
04865     *(volatile struct exec_recursive_params *)&p;
04866     return result;
04867 }
04868 
04869 /*
04870  * Calls func(obj, arg, recursive), where recursive is non-zero if the
04871  * current method is called recursively on obj
04872  */
04873 
04874 VALUE
04875 rb_exec_recursive(VALUE (*func) (VALUE, VALUE, int), VALUE obj, VALUE arg)
04876 {
04877     return exec_recursive(func, obj, 0, arg, 0);
04878 }
04879 
04880 /*
04881  * Calls func(obj, arg, recursive), where recursive is non-zero if the
04882  * current method is called recursively on the ordered pair <obj, paired_obj>
04883  */
04884 
04885 VALUE
04886 rb_exec_recursive_paired(VALUE (*func) (VALUE, VALUE, int), VALUE obj, VALUE paired_obj, VALUE arg)
04887 {
04888     return exec_recursive(func, obj, rb_obj_id(paired_obj), arg, 0);
04889 }
04890 
04891 /*
04892  * If recursion is detected on the current method and obj, the outermost
04893  * func will be called with (obj, arg, Qtrue). All inner func will be
04894  * short-circuited using throw.
04895  */
04896 
04897 VALUE
04898 rb_exec_recursive_outer(VALUE (*func) (VALUE, VALUE, int), VALUE obj, VALUE arg)
04899 {
04900     return exec_recursive(func, obj, 0, arg, 1);
04901 }
04902 
04903 /*
04904  * If recursion is detected on the current method, obj and paired_obj,
04905  * the outermost func will be called with (obj, arg, Qtrue). All inner
04906  * func will be short-circuited using throw.
04907  */
04908 
04909 VALUE
04910 rb_exec_recursive_paired_outer(VALUE (*func) (VALUE, VALUE, int), VALUE obj, VALUE paired_obj, VALUE arg)
04911 {
04912     return exec_recursive(func, obj, rb_obj_id(paired_obj), arg, 1);
04913 }
04914 
04915 /*
04916  *  call-seq:
04917  *     thr.backtrace     -> array
04918  *
04919  *  Returns the current backtrace of the target thread.
04920  *
04921  */
04922 
04923 static VALUE
04924 rb_thread_backtrace_m(int argc, VALUE *argv, VALUE thval)
04925 {
04926     return vm_thread_backtrace(argc, argv, thval);
04927 }
04928 
04929 /* call-seq:
04930  *  thr.backtrace_locations(*args)      -> array or nil
04931  *
04932  * Returns the execution stack for the target thread---an array containing
04933  * backtrace location objects.
04934  *
04935  * See Thread::Backtrace::Location for more information.
04936  *
04937  * This method behaves similarly to Kernel#caller_locations except it applies
04938  * to a specific thread.
04939  */
04940 static VALUE
04941 rb_thread_backtrace_locations_m(int argc, VALUE *argv, VALUE thval)
04942 {
04943     return vm_thread_backtrace_locations(argc, argv, thval);
04944 }
04945 
04946 /*
04947  *  Document-class: ThreadError
04948  *
04949  *  Raised when an invalid operation is attempted on a thread.
04950  *
04951  *  For example, when no other thread has been started:
04952  *
04953  *     Thread.stop
04954  *
04955  *  <em>raises the exception:</em>
04956  *
04957  *     ThreadError: stopping only thread
04958  */
04959 
04960 /*
04961  *  +Thread+ encapsulates the behavior of a thread of
04962  *  execution, including the main thread of the Ruby script.
04963  *
04964  *  In the descriptions of the methods in this class, the parameter _sym_
04965  *  refers to a symbol, which is either a quoted string or a
04966  *  +Symbol+ (such as <code>:name</code>).
04967  */
04968 
04969 void
04970 Init_Thread(void)
04971 {
04972 #undef rb_intern
04973 #define rb_intern(str) rb_intern_const(str)
04974 
04975     VALUE cThGroup;
04976     rb_thread_t *th = GET_THREAD();
04977 
04978     sym_never = ID2SYM(rb_intern("never"));
04979     sym_immediate = ID2SYM(rb_intern("immediate"));
04980     sym_on_blocking = ID2SYM(rb_intern("on_blocking"));
04981 
04982     rb_define_singleton_method(rb_cThread, "new", thread_s_new, -1);
04983     rb_define_singleton_method(rb_cThread, "start", thread_start, -2);
04984     rb_define_singleton_method(rb_cThread, "fork", thread_start, -2);
04985     rb_define_singleton_method(rb_cThread, "main", rb_thread_s_main, 0);
04986     rb_define_singleton_method(rb_cThread, "current", thread_s_current, 0);
04987     rb_define_singleton_method(rb_cThread, "stop", rb_thread_stop, 0);
04988     rb_define_singleton_method(rb_cThread, "kill", rb_thread_s_kill, 1);
04989     rb_define_singleton_method(rb_cThread, "exit", rb_thread_exit, 0);
04990     rb_define_singleton_method(rb_cThread, "pass", thread_s_pass, 0);
04991     rb_define_singleton_method(rb_cThread, "list", rb_thread_list, 0);
04992     rb_define_singleton_method(rb_cThread, "abort_on_exception", rb_thread_s_abort_exc, 0);
04993     rb_define_singleton_method(rb_cThread, "abort_on_exception=", rb_thread_s_abort_exc_set, 1);
04994 #if THREAD_DEBUG < 0
04995     rb_define_singleton_method(rb_cThread, "DEBUG", rb_thread_s_debug, 0);
04996     rb_define_singleton_method(rb_cThread, "DEBUG=", rb_thread_s_debug_set, 1);
04997 #endif
04998     rb_define_singleton_method(rb_cThread, "handle_interrupt", rb_thread_s_handle_interrupt, 1);
04999     rb_define_singleton_method(rb_cThread, "pending_interrupt?", rb_thread_s_pending_interrupt_p, -1);
05000     rb_define_method(rb_cThread, "pending_interrupt?", rb_thread_pending_interrupt_p, -1);
05001 
05002     rb_define_method(rb_cThread, "initialize", thread_initialize, -2);
05003     rb_define_method(rb_cThread, "raise", thread_raise_m, -1);
05004     rb_define_method(rb_cThread, "join", thread_join_m, -1);
05005     rb_define_method(rb_cThread, "value", thread_value, 0);
05006     rb_define_method(rb_cThread, "kill", rb_thread_kill, 0);
05007     rb_define_method(rb_cThread, "terminate", rb_thread_kill, 0);
05008     rb_define_method(rb_cThread, "exit", rb_thread_kill, 0);
05009     rb_define_method(rb_cThread, "run", rb_thread_run, 0);
05010     rb_define_method(rb_cThread, "wakeup", rb_thread_wakeup, 0);
05011     rb_define_method(rb_cThread, "[]", rb_thread_aref, 1);
05012     rb_define_method(rb_cThread, "[]=", rb_thread_aset, 2);
05013     rb_define_method(rb_cThread, "key?", rb_thread_key_p, 1);
05014     rb_define_method(rb_cThread, "keys", rb_thread_keys, 0);
05015     rb_define_method(rb_cThread, "priority", rb_thread_priority, 0);
05016     rb_define_method(rb_cThread, "priority=", rb_thread_priority_set, 1);
05017     rb_define_method(rb_cThread, "status", rb_thread_status, 0);
05018     rb_define_method(rb_cThread, "thread_variable_get", rb_thread_variable_get, 1);
05019     rb_define_method(rb_cThread, "thread_variable_set", rb_thread_variable_set, 2);
05020     rb_define_method(rb_cThread, "thread_variables", rb_thread_variables, 0);
05021     rb_define_method(rb_cThread, "thread_variable?", rb_thread_variable_p, 1);
05022     rb_define_method(rb_cThread, "alive?", rb_thread_alive_p, 0);
05023     rb_define_method(rb_cThread, "stop?", rb_thread_stop_p, 0);
05024     rb_define_method(rb_cThread, "abort_on_exception", rb_thread_abort_exc, 0);
05025     rb_define_method(rb_cThread, "abort_on_exception=", rb_thread_abort_exc_set, 1);
05026     rb_define_method(rb_cThread, "safe_level", rb_thread_safe_level, 0);
05027     rb_define_method(rb_cThread, "group", rb_thread_group, 0);
05028     rb_define_method(rb_cThread, "backtrace", rb_thread_backtrace_m, -1);
05029     rb_define_method(rb_cThread, "backtrace_locations", rb_thread_backtrace_locations_m, -1);
05030 
05031     rb_define_method(rb_cThread, "inspect", rb_thread_inspect, 0);
05032 
05033     closed_stream_error = rb_exc_new2(rb_eIOError, "stream closed");
05034     OBJ_TAINT(closed_stream_error);
05035     OBJ_FREEZE(closed_stream_error);
05036 
05037     cThGroup = rb_define_class("ThreadGroup", rb_cObject);
05038     rb_define_alloc_func(cThGroup, thgroup_s_alloc);
05039     rb_define_method(cThGroup, "list", thgroup_list, 0);
05040     rb_define_method(cThGroup, "enclose", thgroup_enclose, 0);
05041     rb_define_method(cThGroup, "enclosed?", thgroup_enclosed_p, 0);
05042     rb_define_method(cThGroup, "add", thgroup_add, 1);
05043 
05044     {
05045         th->thgroup = th->vm->thgroup_default = rb_obj_alloc(cThGroup);
05046         rb_define_const(cThGroup, "Default", th->thgroup);
05047     }
05048 
05049     rb_cMutex = rb_define_class("Mutex", rb_cObject);
05050     rb_define_alloc_func(rb_cMutex, mutex_alloc);
05051     rb_define_method(rb_cMutex, "initialize", mutex_initialize, 0);
05052     rb_define_method(rb_cMutex, "locked?", rb_mutex_locked_p, 0);
05053     rb_define_method(rb_cMutex, "try_lock", rb_mutex_trylock, 0);
05054     rb_define_method(rb_cMutex, "lock", rb_mutex_lock, 0);
05055     rb_define_method(rb_cMutex, "unlock", rb_mutex_unlock, 0);
05056     rb_define_method(rb_cMutex, "sleep", mutex_sleep, -1);
05057     rb_define_method(rb_cMutex, "synchronize", rb_mutex_synchronize_m, 0);
05058     rb_define_method(rb_cMutex, "owned?", rb_mutex_owned_p, 0);
05059 
05060     recursive_key = rb_intern("__recursive_key__");
05061     rb_eThreadError = rb_define_class("ThreadError", rb_eStandardError);
05062 
05063     /* init thread core */
05064     {
05065         /* main thread setting */
05066         {
05067             /* acquire global vm lock */
05068             gvl_init(th->vm);
05069             gvl_acquire(th->vm, th);
05070             native_mutex_initialize(&th->vm->thread_destruct_lock);
05071             native_mutex_initialize(&th->interrupt_lock);
05072 
05073             th->pending_interrupt_queue = rb_ary_tmp_new(0);
05074             th->pending_interrupt_queue_checked = 0;
05075             th->pending_interrupt_mask_stack = rb_ary_tmp_new(0);
05076 
05077             th->interrupt_mask = 0;
05078         }
05079     }
05080 
05081     rb_thread_create_timer_thread();
05082 
05083     /* suppress warnings on cygwin, mingw and mswin.*/
05084     (void)native_mutex_trylock;
05085 }
05086 
05087 int
05088 ruby_native_thread_p(void)
05089 {
05090     rb_thread_t *th = ruby_thread_from_native();
05091 
05092     return th != 0;
05093 }
05094 
05095 static int
05096 check_deadlock_i(st_data_t key, st_data_t val, int *found)
05097 {
05098     VALUE thval = key;
05099     rb_thread_t *th;
05100     GetThreadPtr(thval, th);
05101 
05102     if (th->status != THREAD_STOPPED_FOREVER || RUBY_VM_INTERRUPTED(th)) {
05103         *found = 1;
05104     }
05105     else if (th->locking_mutex) {
05106         rb_mutex_t *mutex;
05107         GetMutexPtr(th->locking_mutex, mutex);
05108 
05109         native_mutex_lock(&mutex->lock);
05110         if (mutex->th == th || (!mutex->th && mutex->cond_waiting)) {
05111             *found = 1;
05112         }
05113         native_mutex_unlock(&mutex->lock);
05114     }
05115 
05116     return (*found) ? ST_STOP : ST_CONTINUE;
05117 }
05118 
05119 #ifdef DEBUG_DEADLOCK_CHECK
05120 static int
05121 debug_i(st_data_t key, st_data_t val, int *found)
05122 {
05123     VALUE thval = key;
05124     rb_thread_t *th;
05125     GetThreadPtr(thval, th);
05126 
05127     printf("th:%p %d %d", th, th->status, th->interrupt_flag);
05128     if (th->locking_mutex) {
05129         rb_mutex_t *mutex;
05130         GetMutexPtr(th->locking_mutex, mutex);
05131 
05132         native_mutex_lock(&mutex->lock);
05133         printf(" %p %d\n", mutex->th, mutex->cond_waiting);
05134         native_mutex_unlock(&mutex->lock);
05135     }
05136     else
05137         puts("");
05138 
05139     return ST_CONTINUE;
05140 }
05141 #endif
05142 
05143 static void
05144 rb_check_deadlock(rb_vm_t *vm)
05145 {
05146     int found = 0;
05147 
05148     if (vm_living_thread_num(vm) > vm->sleeper) return;
05149     if (vm_living_thread_num(vm) < vm->sleeper) rb_bug("sleeper must not be more than vm_living_thread_num(vm)");
05150     if (patrol_thread && patrol_thread != GET_THREAD()) return;
05151 
05152     st_foreach(vm->living_threads, check_deadlock_i, (st_data_t)&found);
05153 
05154     if (!found) {
05155         VALUE argv[2];
05156         argv[0] = rb_eFatal;
05157         argv[1] = rb_str_new2("No live threads left. Deadlock?");
05158 #ifdef DEBUG_DEADLOCK_CHECK
05159         printf("%d %d %p %p\n", vm->living_threads->num_entries, vm->sleeper, GET_THREAD(), vm->main_thread);
05160         st_foreach(vm->living_threads, debug_i, (st_data_t)0);
05161 #endif
05162         vm->sleeper--;
05163         rb_threadptr_raise(vm->main_thread, 2, argv);
05164     }
05165 }
05166 
05167 static void
05168 update_coverage(rb_event_flag_t event, VALUE proc, VALUE self, ID id, VALUE klass)
05169 {
05170     VALUE coverage = GET_THREAD()->cfp->iseq->coverage;
05171     if (coverage && RBASIC(coverage)->klass == 0) {
05172         long line = rb_sourceline() - 1;
05173         long count;
05174         if (RARRAY_PTR(coverage)[line] == Qnil) {
05175             return;
05176         }
05177         count = FIX2LONG(RARRAY_PTR(coverage)[line]) + 1;
05178         if (POSFIXABLE(count)) {
05179             RARRAY_PTR(coverage)[line] = LONG2FIX(count);
05180         }
05181     }
05182 }
05183 
05184 VALUE
05185 rb_get_coverages(void)
05186 {
05187     return GET_VM()->coverages;
05188 }
05189 
05190 void
05191 rb_set_coverages(VALUE coverages)
05192 {
05193     GET_VM()->coverages = coverages;
05194     rb_add_event_hook(update_coverage, RUBY_EVENT_COVERAGE, Qnil);
05195 }
05196 
05197 void
05198 rb_reset_coverages(void)
05199 {
05200     GET_VM()->coverages = Qfalse;
05201     rb_remove_event_hook(update_coverage);
05202 }
05203 
05204 VALUE
05205 rb_uninterruptible(VALUE (*b_proc)(ANYARGS), VALUE data)
05206 {
05207     VALUE interrupt_mask = rb_hash_new();
05208     rb_thread_t *cur_th = GET_THREAD();
05209 
05210     rb_hash_aset(interrupt_mask, rb_cObject, sym_never);
05211     rb_ary_push(cur_th->pending_interrupt_mask_stack, interrupt_mask);
05212 
05213     return rb_ensure(b_proc, data, rb_ary_pop, cur_th->pending_interrupt_mask_stack);
05214 }
05215