Class: IO::Event::Selector::KQueue
- Inherits:
-
Object
- Object
- IO::Event::Selector::KQueue
- Defined in:
- ext/io/event/selector/kqueue.c
Instance Method Summary collapse
- #close ⇒ Object
- #closed? ⇒ Boolean
- #idle_duration ⇒ Object
- #initialize(loop) ⇒ Object constructor
- #io_read(*args) ⇒ Object
- #io_wait(fiber, io, events) ⇒ Object
- #io_write(*args) ⇒ Object
- #loop ⇒ Object
- #process_wait(fiber, _pid, _flags) ⇒ Object
- #push(fiber) ⇒ Object
- #raise(*args) ⇒ Object
- #ready? ⇒ Boolean
- #resume(*args) ⇒ Object
- #select(duration) ⇒ Object
- #transfer ⇒ Object
- #wakeup ⇒ Object
- #yield ⇒ Object
Constructor Details
#initialize(loop) ⇒ Object
323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 |
# File 'ext/io/event/selector/kqueue.c', line 323
VALUE IO_Event_Selector_KQueue_initialize(VALUE self, VALUE loop) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
IO_Event_Selector_initialize(&selector->backend, self, loop);
int result = kqueue();
if (result == -1) {
rb_sys_fail("IO_Event_Selector_KQueue_initialize:kqueue");
} else {
// Make sure the descriptor is closed on exec.
ioctl(result, FIOCLEX);
selector->descriptor = result;
selector->owner = getpid();
rb_update_max_fd(selector->descriptor);
}
#ifdef IO_EVENT_SELECTOR_KQUEUE_USE_INTERRUPT
IO_Event_Interrupt_open(&selector->interrupt);
IO_Event_Interrupt_add(&selector->interrupt, selector);
#endif
return self;
}
|
Instance Method Details
#close ⇒ Object
366 367 368 369 370 371 372 373 |
# File 'ext/io/event/selector/kqueue.c', line 366
VALUE IO_Event_Selector_KQueue_close(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
close_internal(selector);
return Qnil;
}
|
#closed? ⇒ Boolean
375 376 377 378 379 380 |
# File 'ext/io/event/selector/kqueue.c', line 375
VALUE IO_Event_Selector_KQueue_closed_p(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return selector->descriptor < 0 || selector->owner != getpid() ? Qtrue : Qfalse;
}
|
#idle_duration ⇒ Object
357 358 359 360 361 362 363 364 |
# File 'ext/io/event/selector/kqueue.c', line 357
VALUE IO_Event_Selector_KQueue_idle_duration(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
double duration = selector->idle_duration.tv_sec + (selector->idle_duration.tv_nsec / 1000000000.0);
return DBL2NUM(duration);
}
|
#io_read(*args) ⇒ Object
654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 |
# File 'ext/io/event/selector/kqueue.c', line 654
VALUE IO_Event_Selector_KQueue_io_read(VALUE self, VALUE fiber, VALUE io, VALUE buffer, VALUE _first, VALUE _second) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
void *base;
size_t size;
rb_io_buffer_get_bytes_for_writing(buffer, &base, &size);
#if RUBY_FIBER_SCHEDULER_VERSION >= 4
size_t offset = NUM2SIZET(_first);
size_t length = NUM2SIZET(_second);
if (!IO_Event_Selector_valid_buffer_range(size, offset, length)) {
return rb_fiber_scheduler_io_result(-1, EINVAL);
} else if (length == 0) {
return rb_fiber_scheduler_io_result(0, 0);
}
base = (char*)base + offset;
size = length;
#else
size_t length = NUM2SIZET(_first);
size_t offset = NUM2SIZET(_second);
if (offset > size) {
return rb_fiber_scheduler_io_result(-1, EINVAL);
} else if (offset == size) {
return rb_fiber_scheduler_io_result(0, 0);
}
base = (char*)base + offset;
size -= offset;
#endif
int descriptor = IO_Event_Selector_io_descriptor(io);
struct io_read_arguments io_read_arguments = {
#if RUBY_FIBER_SCHEDULER_VERSION < 4
.self = self,
.fiber = fiber,
.io = io,
#endif
.flags = IO_Event_Selector_nonblock_set(descriptor),
.descriptor = descriptor,
.base = base,
.size = size,
#if RUBY_FIBER_SCHEDULER_VERSION < 4
.length = length,
#endif
};
RB_OBJ_WRITTEN(self, Qundef, fiber);
return rb_ensure(io_read_loop, (VALUE)&io_read_arguments, io_read_ensure, (VALUE)&io_read_arguments);
}
|
#io_wait(fiber, io, events) ⇒ Object
553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 |
# File 'ext/io/event/selector/kqueue.c', line 553
VALUE IO_Event_Selector_KQueue_io_wait(VALUE self, VALUE fiber, VALUE io, VALUE events) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
int descriptor = IO_Event_Selector_io_descriptor(io);
struct IO_Event_Selector_KQueue_Waiting waiting = {
.list = {.type = &IO_Event_Selector_KQueue_io_wait_list_type},
.fiber = fiber,
.events = RB_NUM2INT(events),
};
RB_OBJ_WRITTEN(self, Qundef, fiber);
int result = IO_Event_Selector_KQueue_Waiting_register(selector, descriptor, &waiting);
if (result == -1) {
rb_sys_fail("IO_Event_Selector_KQueue_io_wait:IO_Event_Selector_KQueue_Waiting_register");
}
struct io_wait_arguments io_wait_arguments = {
.selector = selector,
.waiting = &waiting,
};
if (DEBUG_IO_WAIT) fprintf(stderr, "IO_Event_Selector_KQueue_io_wait descriptor=%d\n", descriptor);
return rb_ensure(io_wait_transfer, (VALUE)&io_wait_arguments, io_wait_ensure, (VALUE)&io_wait_arguments);
}
|
#io_write(*args) ⇒ Object
796 797 798 799 800 801 802 803 804 805 806 807 808 809 810 811 812 813 814 815 816 817 818 819 820 821 822 823 824 825 826 827 828 829 830 831 832 833 834 835 836 837 838 839 840 841 842 843 844 845 846 847 848 849 850 851 852 853 854 855 |
# File 'ext/io/event/selector/kqueue.c', line 796
VALUE IO_Event_Selector_KQueue_io_write(VALUE self, VALUE fiber, VALUE io, VALUE buffer, VALUE _first, VALUE _second) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
const void *base;
size_t size;
rb_io_buffer_get_bytes_for_reading(buffer, &base, &size);
#if RUBY_FIBER_SCHEDULER_VERSION >= 4
size_t offset = NUM2SIZET(_first);
size_t length = NUM2SIZET(_second);
if (!IO_Event_Selector_valid_buffer_range(size, offset, length)) {
return rb_fiber_scheduler_io_result(-1, EINVAL);
} else if (length == 0) {
return rb_fiber_scheduler_io_result(0, 0);
}
base = (const char*)base + offset;
size = length;
#else
size_t length = NUM2SIZET(_first);
size_t offset = NUM2SIZET(_second);
if (length > size) {
rb_raise(rb_eRuntimeError, "Length exceeds size of buffer!");
}
if (offset > size) {
return rb_fiber_scheduler_io_result(-1, EINVAL);
} else if (offset == size) {
return rb_fiber_scheduler_io_result(0, 0);
}
base = (const char*)base + offset;
size -= offset;
#endif
int descriptor = IO_Event_Selector_io_descriptor(io);
struct io_write_arguments io_write_arguments = {
#if RUBY_FIBER_SCHEDULER_VERSION < 4
.self = self,
.fiber = fiber,
.io = io,
#endif
.flags = IO_Event_Selector_nonblock_set(descriptor),
.descriptor = descriptor,
.base = base,
.size = size,
#if RUBY_FIBER_SCHEDULER_VERSION < 4
.length = length,
#endif
};
RB_OBJ_WRITTEN(self, Qundef, fiber);
return rb_ensure(io_write_loop, (VALUE)&io_write_arguments, io_write_ensure, (VALUE)&io_write_arguments);
}
|
#loop ⇒ Object
350 351 352 353 354 355 |
# File 'ext/io/event/selector/kqueue.c', line 350
VALUE IO_Event_Selector_KQueue_loop(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return selector->backend.loop;
}
|
#process_wait(fiber, _pid, _flags) ⇒ Object
482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 |
# File 'ext/io/event/selector/kqueue.c', line 482
VALUE IO_Event_Selector_KQueue_process_wait(VALUE self, VALUE fiber, VALUE _pid, VALUE _flags) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
pid_t pid = NUM2PIDT(_pid);
int flags = NUM2INT(_flags);
// `EVFILT_PROC` can only refer to a specific process, so waiting for any child or a process group (pid <= 0) is delegated to the threaded fallback:
if (pid <= 0) {
return IO_Event_Selector_process_wait(pid, flags);
}
struct IO_Event_Selector_KQueue_Waiting waiting = {
.list = {.type = &IO_Event_Selector_KQueue_process_wait_list_type},
.fiber = fiber,
.events = IO_EVENT_EXIT,
};
RB_OBJ_WRITTEN(self, Qundef, fiber);
struct process_wait_arguments process_wait_arguments = {
.selector = selector,
.waiting = &waiting,
.pid = pid,
.flags = flags,
};
int result = IO_Event_Selector_KQueue_Waiting_register(selector, pid, &waiting);
if (result == -1) {
// OpenBSD/NetBSD return ESRCH when attempting to register an EVFILT_PROC event for a zombie process.
if (errno == ESRCH) {
process_prewait(pid);
return IO_Event_Selector_process_status_reap(pid, flags);
}
rb_sys_fail("IO_Event_Selector_KQueue_process_wait:IO_Event_Selector_KQueue_Waiting_register");
}
return rb_ensure(process_wait_transfer, (VALUE)&process_wait_arguments, process_wait_ensure, (VALUE)&process_wait_arguments);
}
|
#push(fiber) ⇒ Object
406 407 408 409 410 411 412 413 414 |
# File 'ext/io/event/selector/kqueue.c', line 406
VALUE IO_Event_Selector_KQueue_push(VALUE self, VALUE fiber)
{
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
IO_Event_Selector_ready_push(&selector->backend, fiber);
return Qnil;
}
|
#raise(*args) ⇒ Object
416 417 418 419 420 421 422 |
# File 'ext/io/event/selector/kqueue.c', line 416
VALUE IO_Event_Selector_KQueue_raise(int argc, VALUE *argv, VALUE self)
{
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return IO_Event_Selector_raise(&selector->backend, argc, argv);
}
|
#ready? ⇒ Boolean
424 425 426 427 428 429 |
# File 'ext/io/event/selector/kqueue.c', line 424
VALUE IO_Event_Selector_KQueue_ready_p(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return selector->backend.ready ? Qtrue : Qfalse;
}
|
#resume(*args) ⇒ Object
390 391 392 393 394 395 396 |
# File 'ext/io/event/selector/kqueue.c', line 390
VALUE IO_Event_Selector_KQueue_resume(int argc, VALUE *argv, VALUE self)
{
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return IO_Event_Selector_resume(&selector->backend, argc, argv);
}
|
#select(duration) ⇒ Object
1045 1046 1047 1048 1049 1050 1051 1052 1053 1054 1055 1056 1057 1058 1059 1060 1061 1062 1063 1064 1065 1066 1067 1068 1069 1070 1071 1072 1073 1074 1075 1076 1077 1078 1079 1080 1081 1082 1083 1084 1085 1086 1087 1088 1089 1090 1091 1092 1093 1094 1095 1096 1097 1098 1099 1100 1101 1102 1103 1104 |
# File 'ext/io/event/selector/kqueue.c', line 1045
VALUE IO_Event_Selector_KQueue_select(VALUE self, VALUE duration) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
selector->idle_duration.tv_sec = 0;
selector->idle_duration.tv_nsec = 0;
int ready = IO_Event_Selector_ready_flush(&selector->backend);
struct select_arguments arguments = {
.selector = selector,
.count = KQUEUE_MAX_EVENTS,
.result = 0,
.storage = {
.tv_sec = 0,
.tv_nsec = 0
},
.saved = {},
};
arguments.timeout = &arguments.storage;
// We break this implementation into two parts.
// (1) result = kevent(..., timeout = 0)
// (2) without gvl: kevent(..., timeout = 0) if result == 0 and timeout != 0
// This allows us to avoid releasing and reacquiring the GVL.
// Non-comprehensive testing shows this gives a 1.5x speedup.
// First do the syscall with no timeout to get any immediately available events:
if (DEBUG) fprintf(stderr, "\r\nselect_internal_with_gvl timeout=" IO_EVENT_TIME_PRINTF_TIMESPEC "\r\n", IO_EVENT_TIME_PRINTF_TIMESPEC_ARGUMENTS(arguments.storage));
int result = select_internal_with_gvl(&arguments);
if (DEBUG) fprintf(stderr, "\r\nselect_internal_with_gvl done\r\n");
// If we:
// 1. Didn't process any ready fibers, and
// 2. Didn't process any events from non-blocking select (above), and
// 3. There are no items in the ready list,
// then we can perform a blocking select.
if (!ready && !result && !selector->backend.ready) {
arguments.timeout = make_timeout(duration, &arguments.storage);
if (select_blocking_allowed(arguments.timeout)) {
struct timespec start_time;
IO_Event_Time_current(&start_time);
if (DEBUG) fprintf(stderr, "IO_Event_Selector_KQueue_select timeout=" IO_EVENT_TIME_PRINTF_TIMESPEC "\n", IO_EVENT_TIME_PRINTF_TIMESPEC_ARGUMENTS(arguments.storage));
result = select_internal_without_gvl(&arguments);
struct timespec end_time;
IO_Event_Time_current(&end_time);
IO_Event_Time_elapsed(&start_time, &end_time, &selector->idle_duration);
}
}
if (result) {
return rb_ensure(select_handle_events, (VALUE)&arguments, select_handle_events_ensure, (VALUE)&arguments);
} else {
return RB_INT2NUM(0);
}
}
|
#transfer ⇒ Object
382 383 384 385 386 387 388 |
# File 'ext/io/event/selector/kqueue.c', line 382
VALUE IO_Event_Selector_KQueue_transfer(VALUE self)
{
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return IO_Event_Selector_loop_yield(&selector->backend);
}
|
#wakeup ⇒ Object
1106 1107 1108 1109 1110 1111 1112 1113 1114 1115 1116 1117 1118 1119 1120 1121 1122 1123 1124 1125 1126 1127 1128 1129 1130 1131 1132 1133 1134 1135 1136 1137 1138 1139 1140 |
# File 'ext/io/event/selector/kqueue.c', line 1106
VALUE IO_Event_Selector_KQueue_wakeup(VALUE self) {
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
if (selector->backend.blocked) {
#ifdef IO_EVENT_SELECTOR_KQUEUE_USE_INTERRUPT
IO_Event_Interrupt_signal(&selector->interrupt);
#else
struct kevent trigger = {0};
trigger.filter = EVFILT_USER;
trigger.flags = EV_ADD | EV_CLEAR;
int result = kevent(selector->descriptor, &trigger, 1, NULL, 0, NULL);
if (result == -1) {
rb_sys_fail("IO_Event_Selector_KQueue_wakeup:kevent");
}
// FreeBSD apparently only works if the NOTE_TRIGGER is done as a separate call.
trigger.flags = 0;
trigger.fflags = NOTE_TRIGGER;
result = kevent(selector->descriptor, &trigger, 1, NULL, 0, NULL);
if (result == -1) {
rb_sys_fail("IO_Event_Selector_KQueue_wakeup:kevent");
}
#endif
return Qtrue;
}
return Qfalse;
}
|
#yield ⇒ Object
398 399 400 401 402 403 404 |
# File 'ext/io/event/selector/kqueue.c', line 398
VALUE IO_Event_Selector_KQueue_yield(VALUE self)
{
struct IO_Event_Selector_KQueue *selector = NULL;
TypedData_Get_Struct(self, struct IO_Event_Selector_KQueue, &IO_Event_Selector_KQueue_Type, selector);
return IO_Event_Selector_yield(&selector->backend);
}
|