|
Ruby
2.0.0p594(2014-10-27revision48167)
|
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, ®ion->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, ®ion->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
1.7.6.1