From e15f1d5cad5d35cc640252377344dcdfbec04b22 Mon Sep 17 00:00:00 2001 From: Jeff Olson Date: Wed, 18 Apr 2012 22:43:11 -0700 Subject: std: refactor global_loop::get.. make it reusable --- src/libstd/uv_global_loop.rs | 65 ++++++++++++++++++++++++++++++++++---------- src/libstd/uv_hl.rs | 45 ++++++++++++++++++++++++------ 2 files changed, 88 insertions(+), 22 deletions(-) (limited to 'src/libstd') diff --git a/src/libstd/uv_global_loop.rs b/src/libstd/uv_global_loop.rs index 74cb9e4442c..b220ba9fe19 100644 --- a/src/libstd/uv_global_loop.rs +++ b/src/libstd/uv_global_loop.rs @@ -6,7 +6,7 @@ import ll = uv_ll; import hl = uv_hl; import get_gl = get; -export get; +export get, get_single_task_gl; native mod rustrt { fn rust_uv_get_kernel_global_chan_ptr() -> *libc::uintptr_t; @@ -26,7 +26,44 @@ loop is running. loop. "] fn get() -> hl::high_level_loop { + ret get_single_task_gl(); +} + +#[doc(hidden)] +fn get_single_task_gl() -> hl::high_level_loop { let global_loop_chan_ptr = rustrt::rust_uv_get_kernel_global_chan_ptr(); + ret spawn_global_weak_task( + global_loop_chan_ptr, + {|weak_exit_po, msg_po, loop_ptr, first_msg| + log(debug, "about to enter inner loop"); + unsafe { + single_task_loop_body(weak_exit_po, msg_po, loop_ptr, + copy(first_msg)) + } + }, + {|msg_ch| + log(debug, "after priv::chan_from_global_ptr"); + unsafe { + let handle = get_global_async_handle_native_representation() + as **ll::uv_async_t; + hl::single_task_loop( + { async_handle: handle, op_chan: msg_ch }) + } + } + ); +} + +// INTERNAL API + +fn spawn_global_weak_task( + global_loop_chan_ptr: *libc::uintptr_t, + weak_task_body_cb: fn~( + comm::port<()>, + comm::port, + *libc::c_void, + hl::high_level_msg) -> bool, + after_task_spawn_cb: fn~(comm::chan) + -> hl::high_level_loop) -> hl::high_level_loop { log(debug, #fmt("ENTERING global_loop::get() loop chan: %?", global_loop_chan_ptr)); @@ -44,25 +81,26 @@ fn get() -> hl::high_level_loop { }; unsafe { log(debug, "before priv::chan_from_global_ptr"); - let chan = priv::chan_from_global_ptr::( + let msg_ch = priv::chan_from_global_ptr::( global_loop_chan_ptr, builder_fn) {|port| // the actual body of our global loop lives here log(debug, "initialized global port task!"); log(debug, "GLOBAL initialized global port task!"); - outer_global_loop_body(port); + outer_global_loop_body(port, weak_task_body_cb); }; - log(debug, "after priv::chan_from_global_ptr"); - let handle = get_global_async_handle_native_representation() - as **ll::uv_async_t; - ret { async_handle: handle, op_chan: chan }; + ret after_task_spawn_cb(msg_ch); } } -// INTERNAL API - -unsafe fn outer_global_loop_body(msg_po: comm::port) { +unsafe fn outer_global_loop_body( + msg_po: comm::port, + weak_task_body_cb: fn~( + comm::port<()>, + comm::port, + *libc::c_void, + hl::high_level_msg) -> bool) { // we're going to use a single libuv-generated loop ptr // for the duration of the process let loop_ptr = ll::loop_new(); @@ -87,9 +125,8 @@ unsafe fn outer_global_loop_body(msg_po: comm::port) { left_val)); false }, {|right_val| - log(debug, "about to enter inner loop"); - inner_global_loop_body(weak_exit_po, msg_po, loop_ptr, - copy(right_val)) + weak_task_body_cb(weak_exit_po, msg_po, loop_ptr, + right_val) }, comm::select2(weak_exit_po, msg_po)); log(debug,#fmt("GLOBAL LOOP EXITED, WAITING TO RESTART? %?", continue)); @@ -99,7 +136,7 @@ unsafe fn outer_global_loop_body(msg_po: comm::port) { ll::loop_delete(loop_ptr); } -unsafe fn inner_global_loop_body(weak_exit_po_in: comm::port<()>, +unsafe fn single_task_loop_body(weak_exit_po_in: comm::port<()>, msg_po_in: comm::port, loop_ptr: *libc::c_void, -first_interaction: hl::high_level_msg) -> bool { diff --git a/src/libstd/uv_hl.rs b/src/libstd/uv_hl.rs index 5cc58932c4b..b235cdd5e07 100644 --- a/src/libstd/uv_hl.rs +++ b/src/libstd/uv_hl.rs @@ -6,7 +6,7 @@ provide a high-level, abstracted interface to some set of libuv functionality. "]; -export high_level_loop, high_level_msg; +export high_level_loop, hl_loop_ext, high_level_msg; export run_high_level_loop, interact, ref_handle, unref_handle; import ll = uv_ll; @@ -22,10 +22,39 @@ the C uv loop to process any pending callbacks * op_chan - a channel used to send function callbacks to be processed by the C uv loop "] -type high_level_loop = { - async_handle: **ll::uv_async_t, - op_chan: comm::chan -}; +enum high_level_loop { + single_task_loop({ + async_handle: **ll::uv_async_t, + op_chan: comm::chan + }), + monitor_task_loop({ + op_chan: comm::chan + }) +} + +impl hl_loop_ext for high_level_loop { + fn async_handle() -> **ll::uv_async_t { + alt self { + single_task_loop({async_handle, op_chan}) { + ret async_handle; + } + _ { + fail "variant of hl::high_level_loop that doesn't include" + + "an async_handle field"; + } + } + } + fn op_chan() -> comm::chan { + alt self { + single_task_loop({async_handle, op_chan}) { + ret op_chan; + } + monitor_task_loop({op_chan}) { + ret op_chan; + } + } + } +} #[doc=" Represents the range of interactions with a `high_level_loop` @@ -158,15 +187,15 @@ enum global_loop_data { unsafe fn send_high_level_msg(hl_loop: high_level_loop, -msg: high_level_msg) unsafe { - comm::send(hl_loop.op_chan, msg); + comm::send(hl_loop.op_chan(), msg); // if the global async handle == 0, then that means // the loop isn't active, so we don't need to wake it up, // (the loop's enclosing task should be blocking on a message // receive on this port) - if (*(hl_loop.async_handle) != 0 as *ll::uv_async_t) { + if (*(hl_loop.async_handle()) != 0 as *ll::uv_async_t) { log(debug,"global async handle != 0, waking up loop.."); - ll::async_send(*(hl_loop.async_handle)); + ll::async_send(*(hl_loop.async_handle())); } else { log(debug,"GLOBAL ASYNC handle == 0"); -- cgit 1.4.1-3-g733a5