diff options
| author | Eric Holk <eholk@mozilla.com> | 2011-08-25 11:20:43 -0700 |
|---|---|---|
| committer | Eric Holk <eholk@mozilla.com> | 2011-08-25 11:21:25 -0700 |
| commit | 2f7c583bc12c0bddb28e43ea79b593a014811b09 (patch) | |
| tree | dff43b4686d290723f689b6605449fabf19b3622 /src/lib | |
| parent | b31815f8a0b98445d2a82888a290b9543ad4400f (diff) | |
Cleaning up task and comm exports, updating all the test cases.
Diffstat (limited to 'src/lib')
| -rw-r--r-- | src/lib/aio.rs | 89 | ||||
| -rw-r--r-- | src/lib/comm.rs | 29 | ||||
| -rw-r--r-- | src/lib/sio.rs | 48 | ||||
| -rw-r--r-- | src/lib/task.rs | 33 | ||||
| -rw-r--r-- | src/lib/test.rs | 4 |
5 files changed, 103 insertions, 100 deletions
diff --git a/src/lib/aio.rs b/src/lib/aio.rs index c2e9c13dc89..0f0119cdd7c 100644 --- a/src/lib/aio.rs +++ b/src/lib/aio.rs @@ -2,11 +2,10 @@ import task; import vec; import comm; -import comm::_chan; -import comm::_port; -import comm::mk_port; +import comm::chan; +import comm::port; import comm::send; - +import comm::recv; import net; native "rust" mod rustrt { @@ -15,11 +14,11 @@ native "rust" mod rustrt { fn aio_init(); fn aio_run(); fn aio_stop(); - fn aio_connect(host: *u8, port: int, connected: &_chan<socket>); - fn aio_serve(host: *u8, port: int, acceptChan: &_chan<socket>) -> server; - fn aio_writedata(s: socket, buf: *u8, size: uint, status: &_chan<bool>); - fn aio_read(s: socket, reader: &_chan<[u8]>); - fn aio_close_server(s: server, status: &_chan<bool>); + fn aio_connect(host: *u8, port: int, connected: &chan<socket>); + fn aio_serve(host: *u8, port: int, acceptChan: &chan<socket>) -> server; + fn aio_writedata(s: socket, buf: *u8, size: uint, status: &chan<bool>); + fn aio_read(s: socket, reader: &chan<[u8]>); + fn aio_close_server(s: server, status: &chan<bool>); fn aio_close_socket(s: socket); fn aio_is_null_client(s: socket) -> bool; } @@ -32,42 +31,42 @@ tag pending_connection { remote(net::ip_addr, int); incoming(server); } tag socket_event { connected(client); closed; received([u8]); } -tag server_event { pending(_chan<_chan<socket_event>>); } +tag server_event { pending(chan<chan<socket_event>>); } tag request { quit; - connect(pending_connection, _chan<socket_event>); - serve(net::ip_addr, int, _chan<server_event>, _chan<server>); - write(client, [u8], _chan<bool>); - close_server(server, _chan<bool>); + connect(pending_connection, chan<socket_event>); + serve(net::ip_addr, int, chan<server_event>, chan<server>); + write(client, [u8], chan<bool>); + close_server(server, chan<bool>); close_client(client); } -type ctx = _chan<request>; +type ctx = chan<request>; fn ip_to_sbuf(ip: net::ip_addr) -> *u8 { vec::to_ptr(str::bytes(net::format_addr(ip))) } -fn connect_task(ip: net::ip_addr, portnum: int, evt: _chan<socket_event>) { - let connecter: _port<client> = mk_port(); - rustrt::aio_connect(ip_to_sbuf(ip), portnum, connecter.mk_chan()); - let client = connecter.recv(); +fn connect_task(ip: net::ip_addr, portnum: int, evt: chan<socket_event>) { + let connecter = port(); + rustrt::aio_connect(ip_to_sbuf(ip), portnum, chan(connecter)); + let client = recv(connecter); new_client(client, evt); } -fn new_client(client: client, evt: _chan<socket_event>) { +fn new_client(client: client, evt: chan<socket_event>) { // Start the read before notifying about the connect. This avoids a race // condition where the receiver can close the socket before we start // reading. - let reader: _port<[u8]> = mk_port(); - rustrt::aio_read(client, reader.mk_chan()); + let reader: port<[u8]> = port(); + rustrt::aio_read(client, chan(reader)); send(evt, connected(client)); while true { log "waiting for bytes"; - let data: [u8] = reader.recv(); + let data: [u8] = recv(reader); log "got some bytes"; log vec::len::<u8>(data); if vec::len::<u8>(data) == 0u { @@ -83,42 +82,42 @@ fn new_client(client: client, evt: _chan<socket_event>) { log "close message sent"; } -fn accept_task(client: client, events: _chan<server_event>) { +fn accept_task(client: client, events: chan<server_event>) { log "accept task was spawned"; - let p: _port<_chan<socket_event>> = mk_port(); - send(events, pending(p.mk_chan())); - let evt = p.recv(); + let p = port(); + send(events, pending(chan(p))); + let evt = recv(p); new_client(client, evt); log "done accepting"; } -fn server_task(ip: net::ip_addr, portnum: int, events: _chan<server_event>, - server: _chan<server>) { - let accepter: _port<client> = mk_port(); +fn server_task(ip: net::ip_addr, portnum: int, events: chan<server_event>, + server: chan<server>) { + let accepter = port(); send(server, - rustrt::aio_serve(ip_to_sbuf(ip), portnum, accepter.mk_chan())); + rustrt::aio_serve(ip_to_sbuf(ip), portnum, chan(accepter))); let client: client; while true { log "preparing to accept a client"; - client = accepter.recv(); + client = recv(accepter); if rustrt::aio_is_null_client(client) { log "client was actually null, returning"; ret; - } else { task::_spawn(bind accept_task(client, events)); } + } else { task::spawn(bind accept_task(client, events)); } } } -fn request_task(c: _chan<ctx>) { +fn request_task(c: chan<ctx>) { // Create a port to accept IO requests on - let p: _port<request> = mk_port(); + let p = port(); // Hand of its channel to our spawner - send(c, p.mk_chan()); + send(c, chan(p)); log "uv run task spawned"; // Spin for requests let req: request; while true { - req = p.recv(); + req = recv(p); alt req { quit. { log "got quit message"; @@ -127,10 +126,10 @@ fn request_task(c: _chan<ctx>) { ret; } connect(remote(ip, portnum), client) { - task::_spawn(bind connect_task(ip, portnum, client)); + task::spawn(bind connect_task(ip, portnum, client)); } serve(ip, portnum, events, server) { - task::_spawn(bind server_task(ip, portnum, events, server)); + task::spawn(bind server_task(ip, portnum, events, server)); } write(socket, v, status) { rustrt::aio_writedata(socket, vec::to_ptr::<u8>(v), @@ -148,27 +147,27 @@ fn request_task(c: _chan<ctx>) { } } -fn iotask(c: _chan<ctx>) { +fn iotask(c: chan<ctx>) { log "io task spawned"; // Initialize before accepting requests rustrt::aio_init(); log "io task init"; // Spawn our request task - let reqtask = task::_spawn(bind request_task(c)); + let reqtask = task::spawn_joinable(bind request_task(c)); log "uv run task init"; // Enter IO loop. This never returns until aio_stop is called. rustrt::aio_run(); log "waiting for request task to finish"; - task::join_id(reqtask); + task::join(reqtask); } fn new() -> ctx { - let p: _port<ctx> = mk_port(); - task::_spawn(bind iotask(p.mk_chan())); - ret p.recv(); + let p: port<ctx> = port(); + task::spawn(bind iotask(chan(p))); + ret recv(p); } // Local Variables: diff --git a/src/lib/comm.rs b/src/lib/comm.rs index ebe9c75bc39..b5d1d3c7e93 100644 --- a/src/lib/comm.rs +++ b/src/lib/comm.rs @@ -2,12 +2,7 @@ import sys; import ptr; import unsafe; import task; -import task::task_id; -export _chan; -export _port; -export chan_handle; -export mk_port; export send; export recv; export chan; @@ -17,7 +12,8 @@ native "rust" mod rustrt { type void; type rust_port; - fn chan_id_send<~T>(target_task: task_id, target_port: port_id, data: -T); + fn chan_id_send<~T>(target_task: task::task, + target_port: port_id, data: -T); fn new_port(unit_sz: uint) -> *rust_port; fn del_port(po: *rust_port); @@ -31,10 +27,9 @@ native "rust-intrinsic" mod rusti { type port_id = int; -type chan_handle<~T> = {task: task_id, port: port_id}; - -tag chan<~T> { chan_t(chan_handle<T>); } -type _chan<~T> = chan<T>; +// It's critical that this only have one variant, so it has a record +// layout, and will work in the rust_task structure in task.rs. +tag chan<~T> { chan_t(task::task, port_id); } resource port_ptr(po: *rustrt::rust_port) { rustrt::drop_port(po); @@ -43,17 +38,9 @@ resource port_ptr(po: *rustrt::rust_port) { tag port<~T> { port_t(@port_ptr); } -obj port_obj<~T>(raw_port: port<T>) { - fn mk_chan() -> chan<T> { chan(raw_port) } - - fn recv() -> T { recv(raw_port) } -} -type _port<~T> = port_obj<T>; - -fn mk_port<~T>() -> _port<T> { ret port_obj::<T>(port::<T>()); } - fn send<~T>(ch: &chan<T>, data: -T) { - rustrt::chan_id_send(ch.task, ch.port, data); + let chan_t(t, p) = ch; + rustrt::chan_id_send(t, p, data); } fn port<~T>() -> port<T> { @@ -63,5 +50,5 @@ fn port<~T>() -> port<T> { fn recv<~T>(p: &port<T>) -> T { ret rusti::recv(***p) } fn chan<~T>(p: &port<T>) -> chan<T> { - chan_t({task: task::get_task_id(), port: rustrt::get_port_id(***p)}) + chan_t(task::get_task_id(), rustrt::get_port_id(***p)) } diff --git a/src/lib/sio.rs b/src/lib/sio.rs index b506f92a97b..9040210b826 100644 --- a/src/lib/sio.rs +++ b/src/lib/sio.rs @@ -1,21 +1,21 @@ -import comm::_port; -import comm::_chan; -import comm::mk_port; +import comm::port; +import comm::chan; import comm::send; +import comm::recv; import str; import net; type ctx = aio::ctx; -type client = {ctx: ctx, client: aio::client, evt: _port<aio::socket_event>}; -type server = {ctx: ctx, server: aio::server, evt: _port<aio::server_event>}; +type client = {ctx: ctx, client: aio::client, evt: port<aio::socket_event>}; +type server = {ctx: ctx, server: aio::server, evt: port<aio::server_event>}; fn new() -> ctx { ret aio::new(); } fn destroy(ctx: ctx) { send(ctx, aio::quit); } -fn make_socket(ctx: ctx, p: _port<aio::socket_event>) -> client { - let evt: aio::socket_event = p.recv(); +fn make_socket(ctx: ctx, p: port<aio::socket_event>) -> client { + let evt: aio::socket_event = recv(p); alt evt { aio::connected(client) { ret {ctx: ctx, client: client, evt: p}; } _ { fail "Could not connect to client"; } @@ -23,56 +23,56 @@ fn make_socket(ctx: ctx, p: _port<aio::socket_event>) -> client { } fn connect_to(ctx: ctx, ip: net::ip_addr, portnum: int) -> client { - let p: _port<aio::socket_event> = mk_port(); - send(ctx, aio::connect(aio::remote(ip, portnum), p.mk_chan())); + let p: port<aio::socket_event> = port(); + send(ctx, aio::connect(aio::remote(ip, portnum), chan(p))); ret make_socket(ctx, p); } fn read(c: client) -> [u8] { - alt c.evt.recv() { + alt recv(c.evt) { aio::closed. { ret []; } aio::received(buf) { ret buf; } } } fn create_server(ctx: ctx, ip: net::ip_addr, portnum: int) -> server { - let evt: _port<aio::server_event> = mk_port(); - let p: _port<aio::server> = mk_port(); - send(ctx, aio::serve(ip, portnum, evt.mk_chan(), p.mk_chan())); - let srv: aio::server = p.recv(); + let evt: port<aio::server_event> = port(); + let p: port<aio::server> = port(); + send(ctx, aio::serve(ip, portnum, chan(evt), chan(p))); + let srv: aio::server = recv(p); ret {ctx: ctx, server: srv, evt: evt}; } fn accept_from(server: server) -> client { - let evt: aio::server_event = server.evt.recv(); + let evt: aio::server_event = recv(server.evt); alt evt { aio::pending(callback) { - let p: _port<aio::socket_event> = mk_port(); - send(callback, p.mk_chan()); + let p = port(); + send(callback, chan(p)); ret make_socket(server.ctx, p); } } } fn write_data(c: client, data: [u8]) -> bool { - let p: _port<bool> = mk_port(); - send(c.ctx, aio::write(c.client, data, p.mk_chan())); - ret p.recv(); + let p = port(); + send(c.ctx, aio::write(c.client, data, chan(p))); + ret recv(p); } fn close_server(server: server) { // TODO: make this unit once we learn to send those from native code - let p: _port<bool> = mk_port(); - send(server.ctx, aio::close_server(server.server, p.mk_chan())); + let p = port(); + send(server.ctx, aio::close_server(server.server, chan(p))); log "Waiting for close"; - p.recv(); + recv(p); log "Got close"; } fn close_client(client: client) { send(client.ctx, aio::close_client(client.client)); let evt: aio::socket_event; - do { evt = client.evt.recv(); alt evt { aio::closed. { ret; } _ { } } } + do { evt = recv(client.evt); alt evt { aio::closed. { ret; } _ { } } } while true } diff --git a/src/lib/task.rs b/src/lib/task.rs index 52181528b4f..71032bbcd64 100644 --- a/src/lib/task.rs +++ b/src/lib/task.rs @@ -5,6 +5,24 @@ import option::none; import option = option::t; import ptr; +export task; +export joinable_task; +export sleep; +export yield; +export task_notification; +export join; +export unsupervise; +export pin; +export unpin; +export set_min_stack; +export spawn; +export spawn_notify; +export spawn_joinable; +export task_result; +export tr_success; +export tr_failure; +export get_task_id; + native "rust" mod rustrt { fn task_sleep(time_in_us: uint); fn task_yield(); @@ -29,8 +47,8 @@ native "rust" mod rustrt { type rust_task = {id: task, - mutable notify_enabled: u8, - mutable notify_chan: comm::chan_handle<task_notification>, + mutable notify_enabled: u32, + mutable notify_chan: comm::chan<task_notification>, ctx: task_context, stack_ptr: *u8}; @@ -40,6 +58,7 @@ resource rust_task_ptr(task: *rust_task) { rustrt::drop_task(task); } type task = int; type task_id = task; +type joinable_task = (task_id, comm::port<task_notification>); fn get_task_id() -> task_id { rustrt::get_task_id() } @@ -79,15 +98,13 @@ fn unpin() { rustrt::unpin_task(); } fn set_min_stack(stack_size: uint) { rustrt::set_min_stack(stack_size); } -fn _spawn(thunk: -fn()) -> task { spawn(thunk) } - fn spawn(thunk: -fn()) -> task { spawn_inner(thunk, none) } fn spawn_notify(thunk: -fn(), notify: comm::chan<task_notification>) -> task { spawn_inner(thunk, some(notify)) } -fn spawn_joinable(thunk: -fn()) -> (task_id, comm::port<task_notification>) { +fn spawn_joinable(thunk: -fn()) -> joinable_task { let p = comm::port::<task_notification>(); let id = spawn_notify(thunk, comm::chan::<task_notification>(p)); ret (id, p); @@ -105,7 +122,7 @@ fn spawn_inner(thunk: -fn(), notify: option<comm::chan<task_notification>>) -> // set up the task pointer let task_ptr = rust_task_ptr(rustrt::get_task_pointer(id)); let regs = ptr::addr_of((**task_ptr).ctx.regs); - (*regs).edx = cast(*task_ptr);; + (*regs).edx = cast(*task_ptr); (*regs).esp = cast((**task_ptr).stack_ptr); assert (ptr::null() != (**task_ptr).stack_ptr); @@ -116,8 +133,8 @@ fn spawn_inner(thunk: -fn(), notify: option<comm::chan<task_notification>>) -> // set up notifications if they are enabled. alt notify { some(c) { - (**task_ptr).notify_enabled = 1u8;; - (**task_ptr).notify_chan = *c; + (**task_ptr).notify_enabled = 1u32;; + (**task_ptr).notify_chan = c; } none { } }; diff --git a/src/lib/test.rs b/src/lib/test.rs index 09ea61dc5d7..e0d27cd1d52 100644 --- a/src/lib/test.rs +++ b/src/lib/test.rs @@ -4,7 +4,7 @@ // while providing a base that other test frameworks may build off of. import generic_os::getenv; -import task::task_id; +import task::task; export test_name; export test_fn; @@ -88,7 +88,7 @@ fn parse_opts(args: &[str]) : vec::is_not_empty(args) -> opt_res { tag test_result { tr_ok; tr_failed; tr_ignored; } -type joinable = (task_id, comm::port<task::task_notification>); +type joinable = (task, comm::port<task::task_notification>); // To get isolation and concurrency tests have to be run in their own tasks. // In cases where test functions and closures it is not ok to just dump them |
