2012-07-04 16:53:12 -05:00
|
|
|
//! A process-wide libuv event loop for library use.
|
2012-04-17 14:05:04 -05:00
|
|
|
|
2012-05-25 13:02:56 -05:00
|
|
|
export get;
|
2012-05-24 23:38:48 -05:00
|
|
|
|
2012-04-17 14:05:04 -05:00
|
|
|
import ll = uv_ll;
|
2012-05-25 01:42:12 -05:00
|
|
|
import iotask = uv_iotask;
|
2012-04-17 14:05:04 -05:00
|
|
|
import get_gl = get;
|
2012-05-25 01:42:12 -05:00
|
|
|
import iotask::{iotask, spawn_iotask};
|
2012-05-24 23:38:48 -05:00
|
|
|
import priv::{chan_from_global_ptr, weaken_task};
|
|
|
|
import comm::{port, chan, methods, select2, listen};
|
|
|
|
import either::{left, right};
|
2012-04-17 14:05:04 -05:00
|
|
|
|
2012-07-03 18:11:00 -05:00
|
|
|
extern mod rustrt {
|
2012-04-17 14:05:04 -05:00
|
|
|
fn rust_uv_get_kernel_global_chan_ptr() -> *libc::uintptr_t;
|
|
|
|
}
|
|
|
|
|
2012-07-04 16:53:12 -05:00
|
|
|
/**
|
|
|
|
* Race-free helper to get access to a global task where a libuv
|
|
|
|
* loop is running.
|
|
|
|
*
|
|
|
|
* Use `uv::hl::interact` to do operations against the global
|
|
|
|
* loop that this function returns.
|
|
|
|
*
|
|
|
|
* # Return
|
|
|
|
*
|
|
|
|
* * A `hl::high_level_loop` that encapsulates communication with the global
|
|
|
|
* loop.
|
|
|
|
*/
|
2012-05-25 01:42:12 -05:00
|
|
|
fn get() -> iotask {
|
2012-04-19 01:49:20 -05:00
|
|
|
ret get_monitor_task_gl();
|
|
|
|
}
|
|
|
|
|
|
|
|
#[doc(hidden)]
|
2012-05-25 01:42:12 -05:00
|
|
|
fn get_monitor_task_gl() -> iotask unsafe {
|
2012-05-24 23:38:48 -05:00
|
|
|
|
|
|
|
let monitor_loop_chan_ptr = rustrt::rust_uv_get_kernel_global_chan_ptr();
|
|
|
|
|
|
|
|
#debug("ENTERING global_loop::get() loop chan: %?",
|
|
|
|
monitor_loop_chan_ptr);
|
|
|
|
|
2012-06-30 18:19:07 -05:00
|
|
|
let builder_fn = || {
|
2012-04-17 14:05:04 -05:00
|
|
|
let builder = task::builder();
|
2012-07-10 17:10:13 -05:00
|
|
|
task::unsupervise(builder);
|
2012-07-10 10:45:08 -05:00
|
|
|
task::set_sched_mode(builder, task::single_threaded);
|
2012-04-17 14:05:04 -05:00
|
|
|
builder
|
|
|
|
};
|
2012-05-24 23:38:48 -05:00
|
|
|
|
|
|
|
#debug("before priv::chan_from_global_ptr");
|
2012-05-25 01:42:12 -05:00
|
|
|
type monchan = chan<iotask>;
|
2012-05-24 23:38:48 -05:00
|
|
|
|
2012-06-26 15:55:56 -05:00
|
|
|
let monitor_ch = do chan_from_global_ptr::<monchan>(monitor_loop_chan_ptr,
|
2012-06-30 18:19:07 -05:00
|
|
|
builder_fn) |msg_po| {
|
2012-05-24 23:38:48 -05:00
|
|
|
#debug("global monitor task starting");
|
|
|
|
|
|
|
|
// As a weak task the runtime will notify us when to exit
|
2012-06-30 18:19:07 -05:00
|
|
|
do weaken_task() |weak_exit_po| {
|
2012-05-24 23:38:48 -05:00
|
|
|
#debug("global monitor task is now weak");
|
2012-05-25 00:26:30 -05:00
|
|
|
let hl_loop = spawn_loop();
|
2012-05-24 23:38:48 -05:00
|
|
|
loop {
|
|
|
|
#debug("in outer_loop...");
|
|
|
|
alt select2(weak_exit_po, msg_po) {
|
|
|
|
left(weak_exit) {
|
|
|
|
// all normal tasks have ended, tell the
|
|
|
|
// libuv loop to tear_down, then exit
|
|
|
|
#debug("weak_exit_po recv'd msg: %?", weak_exit);
|
2012-05-25 01:42:12 -05:00
|
|
|
iotask::exit(hl_loop);
|
2012-05-24 23:38:48 -05:00
|
|
|
break;
|
|
|
|
}
|
|
|
|
right(fetch_ch) {
|
|
|
|
#debug("hl_loop req recv'd: %?", fetch_ch);
|
|
|
|
fetch_ch.send(hl_loop);
|
2012-04-27 23:42:04 -05:00
|
|
|
}
|
|
|
|
}
|
2012-05-24 23:38:48 -05:00
|
|
|
}
|
|
|
|
#debug("global monitor task is leaving weakend state");
|
2012-04-17 14:05:04 -05:00
|
|
|
};
|
2012-05-24 23:38:48 -05:00
|
|
|
#debug("global monitor task exiting");
|
|
|
|
};
|
|
|
|
|
|
|
|
// once we have a chan to the monitor loop, we ask it for
|
|
|
|
// the libuv loop's async handle
|
2012-06-30 18:19:07 -05:00
|
|
|
do listen |fetch_ch| {
|
2012-05-24 23:38:48 -05:00
|
|
|
monitor_ch.send(fetch_ch);
|
|
|
|
fetch_ch.recv()
|
2012-04-17 14:05:04 -05:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
2012-05-25 01:42:12 -05:00
|
|
|
fn spawn_loop() -> iotask unsafe {
|
2012-05-25 00:26:30 -05:00
|
|
|
let builder = task::builder();
|
2012-06-30 18:19:07 -05:00
|
|
|
do task::add_wrapper(builder) |task_body| {
|
2012-05-25 00:26:30 -05:00
|
|
|
fn~(move task_body) {
|
|
|
|
// The I/O loop task also needs to be weak so it doesn't keep
|
|
|
|
// the runtime alive
|
2012-06-30 18:19:07 -05:00
|
|
|
do weaken_task |weak_exit_po| {
|
2012-05-25 00:26:30 -05:00
|
|
|
#debug("global libuv task is now weak %?", weak_exit_po);
|
|
|
|
task_body();
|
|
|
|
|
|
|
|
// We don't wait for the exit message on weak_exit_po
|
|
|
|
// because the monitor task will tell the uv loop when to
|
|
|
|
// exit
|
|
|
|
|
|
|
|
#debug("global libuv task is leaving weakened state");
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
2012-05-25 01:42:12 -05:00
|
|
|
spawn_iotask(builder)
|
2012-04-17 14:05:04 -05:00
|
|
|
}
|
|
|
|
|
|
|
|
#[cfg(test)]
|
|
|
|
mod test {
|
2012-07-03 18:32:02 -05:00
|
|
|
extern fn simple_timer_close_cb(timer_ptr: *ll::uv_timer_t) unsafe {
|
2012-04-17 14:05:04 -05:00
|
|
|
let exit_ch_ptr = ll::get_data_for_uv_handle(
|
|
|
|
timer_ptr as *libc::c_void) as *comm::chan<bool>;
|
|
|
|
let exit_ch = *exit_ch_ptr;
|
|
|
|
comm::send(exit_ch, true);
|
|
|
|
log(debug, #fmt("EXIT_CH_PTR simple_timer_close_cb exit_ch_ptr: %?",
|
|
|
|
exit_ch_ptr));
|
|
|
|
}
|
2012-07-03 18:32:02 -05:00
|
|
|
extern fn simple_timer_cb(timer_ptr: *ll::uv_timer_t,
|
2012-05-24 22:31:20 -05:00
|
|
|
_status: libc::c_int) unsafe {
|
2012-04-17 14:05:04 -05:00
|
|
|
log(debug, "in simple timer cb");
|
|
|
|
ll::timer_stop(timer_ptr);
|
|
|
|
let hl_loop = get_gl();
|
2012-06-30 18:19:07 -05:00
|
|
|
do iotask::interact(hl_loop) |_loop_ptr| {
|
2012-04-17 14:05:04 -05:00
|
|
|
log(debug, "closing timer");
|
2012-04-27 23:42:04 -05:00
|
|
|
ll::close(timer_ptr, simple_timer_close_cb);
|
2012-04-17 14:05:04 -05:00
|
|
|
log(debug, "about to deref exit_ch_ptr");
|
|
|
|
log(debug, "after msg sent on deref'd exit_ch");
|
|
|
|
};
|
|
|
|
log(debug, "exiting simple timer cb");
|
|
|
|
}
|
|
|
|
|
2012-05-25 01:42:12 -05:00
|
|
|
fn impl_uv_hl_simple_timer(iotask: iotask) unsafe {
|
2012-04-17 14:05:04 -05:00
|
|
|
let exit_po = comm::port::<bool>();
|
|
|
|
let exit_ch = comm::chan(exit_po);
|
|
|
|
let exit_ch_ptr = ptr::addr_of(exit_ch);
|
|
|
|
log(debug, #fmt("EXIT_CH_PTR newly created exit_ch_ptr: %?",
|
|
|
|
exit_ch_ptr));
|
|
|
|
let timer_handle = ll::timer_t();
|
|
|
|
let timer_ptr = ptr::addr_of(timer_handle);
|
2012-06-30 18:19:07 -05:00
|
|
|
do iotask::interact(iotask) |loop_ptr| {
|
2012-04-17 14:05:04 -05:00
|
|
|
log(debug, "user code inside interact loop!!!");
|
|
|
|
let init_status = ll::timer_init(loop_ptr, timer_ptr);
|
|
|
|
if(init_status == 0i32) {
|
|
|
|
ll::set_data_for_uv_handle(
|
|
|
|
timer_ptr as *libc::c_void,
|
|
|
|
exit_ch_ptr as *libc::c_void);
|
|
|
|
let start_status = ll::timer_start(timer_ptr, simple_timer_cb,
|
|
|
|
1u, 0u);
|
|
|
|
if(start_status == 0i32) {
|
|
|
|
}
|
|
|
|
else {
|
|
|
|
fail "failure on ll::timer_start()";
|
|
|
|
}
|
|
|
|
}
|
|
|
|
else {
|
|
|
|
fail "failure on ll::timer_init()";
|
|
|
|
}
|
|
|
|
};
|
|
|
|
comm::recv(exit_po);
|
|
|
|
log(debug, "global_loop timer test: msg recv on exit_po, done..");
|
|
|
|
}
|
2012-04-27 23:42:04 -05:00
|
|
|
|
2012-04-17 14:05:04 -05:00
|
|
|
#[test]
|
2012-04-27 23:42:04 -05:00
|
|
|
fn test_gl_uv_global_loop_high_level_global_timer() unsafe {
|
2012-04-17 14:05:04 -05:00
|
|
|
let hl_loop = get_gl();
|
2012-04-27 23:42:04 -05:00
|
|
|
let exit_po = comm::port::<()>();
|
|
|
|
let exit_ch = comm::chan(exit_po);
|
2012-06-30 18:19:07 -05:00
|
|
|
task::spawn_sched(task::manual_threads(1u), || {
|
2012-04-17 14:05:04 -05:00
|
|
|
impl_uv_hl_simple_timer(hl_loop);
|
2012-04-27 23:42:04 -05:00
|
|
|
comm::send(exit_ch, ());
|
2012-04-17 14:05:04 -05:00
|
|
|
});
|
|
|
|
impl_uv_hl_simple_timer(hl_loop);
|
2012-04-27 23:42:04 -05:00
|
|
|
comm::recv(exit_po);
|
|
|
|
}
|
|
|
|
|
|
|
|
// keeping this test ignored until some kind of stress-test-harness
|
|
|
|
// is set up for the build bots
|
|
|
|
#[test]
|
|
|
|
#[ignore]
|
|
|
|
fn test_stress_gl_uv_global_loop_high_level_global_timer() unsafe {
|
|
|
|
let hl_loop = get_gl();
|
|
|
|
let exit_po = comm::port::<()>();
|
|
|
|
let exit_ch = comm::chan(exit_po);
|
|
|
|
let cycles = 5000u;
|
2012-07-04 14:04:28 -05:00
|
|
|
for iter::repeat(cycles) {
|
2012-06-30 18:19:07 -05:00
|
|
|
task::spawn_sched(task::manual_threads(1u), || {
|
2012-04-27 23:42:04 -05:00
|
|
|
impl_uv_hl_simple_timer(hl_loop);
|
|
|
|
comm::send(exit_ch, ());
|
|
|
|
});
|
|
|
|
};
|
2012-07-04 14:04:28 -05:00
|
|
|
for iter::repeat(cycles) {
|
2012-04-27 23:42:04 -05:00
|
|
|
comm::recv(exit_po);
|
|
|
|
};
|
|
|
|
log(debug, "test_stress_gl_uv_global_loop_high_level_global_timer"+
|
|
|
|
" exiting sucessfully!");
|
2012-04-17 14:05:04 -05:00
|
|
|
}
|
2012-07-04 14:04:28 -05:00
|
|
|
}
|