diff options
| author | bors <bors@rust-lang.org> | 2014-12-19 08:28:52 +0000 |
|---|---|---|
| committer | bors <bors@rust-lang.org> | 2014-12-19 08:28:52 +0000 |
| commit | 0efafac398ff7f28c5f0fe756c15b9008b3e0534 (patch) | |
| tree | 2e8279b94829b65868049d2e3df0b9a6c3365a8f /src/libstd/comm | |
| parent | 6bdce25e155d846bb9252fa4a18baef7e74cf8bf (diff) | |
| parent | 903c5a8f69714382ec9fc22745f902c3e219cb68 (diff) | |
auto merge of #19654 : aturon/rust/merge-rt, r=alexcrichton
This PR substantially narrows the notion of a "runtime" in Rust, and allows calling into Rust code directly without any setup or teardown. After this PR, the basic "runtime support" in Rust will consist of: * Unwinding and backtrace support * Stack guards Other support, such as helper threads for timers or the notion of a "current thread" are initialized automatically upon first use. When using Rust in an embedded context, it should now be possible to call a Rust function directly as a C function with absolutely no setup, though in that case panics will cause the process to abort. In this regard, the C/Rust interface will look much like the C/C++ interface. In more detail, this PR: * Merges `librustrt` back into `std::rt`, undoing the facade. While doing so, it removes a substantial amount of redundant functionality (such as mutexes defined in the `rt` module). Code using `librustrt` can now call into `std::rt` to e.g. start executing Rust code with unwinding support. * Allows all runtime data to be initialized lazily, including the "current thread", the "at_exit" infrastructure, and the "args" storage. * Deprecates and largely removes `std::task` along with the widespread requirement that there be a "current task" for many APIs in `std`. The entire task infrastructure is replaced with `std::thread`, which provides a more standard API for manipulating and creating native OS threads. In particular, it's possible to join on a created thread, and to get a handle to the currently-running thread. In addition, threads are equipped with some basic blocking support in the form of `park`/`unpark` operations (following a tradition in some OSes as well as the JVM). See the `std::thread` documentation for more details. * Channels are refactored to use a new internal blocking infrastructure that itself sits on top of `park`/`unpark`. One important change here is that a Rust program ends when its main thread does, following most threading models. On the other hand, threads will often be created with an RAII-style join handle that will re-institute blocking semantics naturally (and with finer control). This is very much a: [breaking-change] Closes #18000 r? @alexcrichton
Diffstat (limited to 'src/libstd/comm')
| -rw-r--r-- | src/libstd/comm/blocking.rs | 83 | ||||
| -rw-r--r-- | src/libstd/comm/mod.rs | 185 | ||||
| -rw-r--r-- | src/libstd/comm/oneshot.rs | 102 | ||||
| -rw-r--r-- | src/libstd/comm/select.rs | 84 | ||||
| -rw-r--r-- | src/libstd/comm/shared.rs | 144 | ||||
| -rw-r--r-- | src/libstd/comm/stream.rs | 58 | ||||
| -rw-r--r-- | src/libstd/comm/sync.rs | 287 |
7 files changed, 474 insertions, 469 deletions
diff --git a/src/libstd/comm/blocking.rs b/src/libstd/comm/blocking.rs new file mode 100644 index 00000000000..c477acd70aa --- /dev/null +++ b/src/libstd/comm/blocking.rs @@ -0,0 +1,83 @@ +// Copyright 2014 The Rust Project Developers. See the COPYRIGHT +// file at the top-level directory of this distribution and at +// http://rust-lang.org/COPYRIGHT. +// +// Licensed under the Apache License, Version 2.0 <LICENSE-APACHE or +// http://www.apache.org/licenses/LICENSE-2.0> or the MIT license +// <LICENSE-MIT or http://opensource.org/licenses/MIT>, at your +// option. This file may not be copied, modified, or distributed +// except according to those terms. + +//! Generic support for building blocking abstractions. + +use thread::Thread; +use sync::atomic::{AtomicBool, INIT_ATOMIC_BOOL, Ordering}; +use sync::Arc; +use kinds::marker::{NoSend, NoSync}; +use mem; +use clone::Clone; + +struct Inner { + thread: Thread, + woken: AtomicBool, +} + +#[deriving(Clone)] +pub struct SignalToken { + inner: Arc<Inner>, +} + +pub struct WaitToken { + inner: Arc<Inner>, + no_send: NoSend, + no_sync: NoSync, +} + +pub fn tokens() -> (WaitToken, SignalToken) { + let inner = Arc::new(Inner { + thread: Thread::current(), + woken: INIT_ATOMIC_BOOL, + }); + let wait_token = WaitToken { + inner: inner.clone(), + no_send: NoSend, + no_sync: NoSync, + }; + let signal_token = SignalToken { + inner: inner + }; + (wait_token, signal_token) +} + +impl SignalToken { + pub fn signal(&self) -> bool { + let wake = !self.inner.woken.compare_and_swap(false, true, Ordering::SeqCst); + if wake { + self.inner.thread.unpark(); + } + wake + } + + /// Convert to an unsafe uint value. Useful for storing in a pipe's state + /// flag. + #[inline] + pub unsafe fn cast_to_uint(self) -> uint { + mem::transmute(self.inner) + } + + /// Convert from an unsafe uint value. Useful for retrieving a pipe's state + /// flag. + #[inline] + pub unsafe fn cast_from_uint(signal_ptr: uint) -> SignalToken { + SignalToken { inner: mem::transmute(signal_ptr) } + } + +} + +impl WaitToken { + pub fn wait(self) { + while !self.inner.woken.load(Ordering::SeqCst) { + Thread::park() + } + } +} diff --git a/src/libstd/comm/mod.rs b/src/libstd/comm/mod.rs index 29a7b0dd0cc..8f945fec4d5 100644 --- a/src/libstd/comm/mod.rs +++ b/src/libstd/comm/mod.rs @@ -54,42 +54,35 @@ //! There are methods on both of senders and receivers to perform their //! respective operations without panicking, however. //! -//! ## Runtime Requirements -//! -//! The channel types defined in this module generally have very few runtime -//! requirements in order to operate. The major requirement they have is for a -//! local rust `Task` to be available if any *blocking* operation is performed. -//! -//! If a local `Task` is not available (for example an FFI callback), then the -//! `send` operation is safe on a `Sender` (as well as a `send_opt`) as well as -//! the `try_send` method on a `SyncSender`, but no other operations are -//! guaranteed to be safe. -//! //! # Example //! //! Simple usage: //! //! ``` +//! use std::thread::Thread; +//! //! // Create a simple streaming channel //! let (tx, rx) = channel(); -//! spawn(move|| { +//! Thread::spawn(move|| { //! tx.send(10i); -//! }); +//! }).detach(); //! assert_eq!(rx.recv(), 10i); //! ``` //! //! Shared usage: //! //! ``` -//! // Create a shared channel that can be sent along from many tasks +//! use std::thread::Thread; +//! +//! // Create a shared channel that can be sent along from many threads //! // where tx is the sending half (tx for transmission), and rx is the receiving //! // half (rx for receiving). //! let (tx, rx) = channel(); //! for i in range(0i, 10i) { //! let tx = tx.clone(); -//! spawn(move|| { +//! Thread::spawn(move|| { //! tx.send(i); -//! }) +//! }).detach() //! } //! //! for _ in range(0i, 10i) { @@ -111,11 +104,13 @@ //! Synchronous channels: //! //! ``` +//! use std::thread::Thread; +//! //! let (tx, rx) = sync_channel::<int>(0); -//! spawn(move|| { +//! Thread::spawn(move|| { //! // This will wait for the parent task to start receiving //! tx.send(53); -//! }); +//! }).detach(); //! rx.recv(); //! ``` //! @@ -327,28 +322,30 @@ use alloc::arc::Arc; use core::kinds::marker; use core::mem; use core::cell::UnsafeCell; -use rustrt::task::BlockedTask; -pub use comm::select::{Select, Handle}; +pub use self::select::{Select, Handle}; +use self::select::StartResult; +use self::select::StartResult::*; +use self::blocking::SignalToken; macro_rules! test { { fn $name:ident() $b:block $(#[$a:meta])*} => ( mod $name { #![allow(unused_imports)] - extern crate rustrt; - use prelude::*; + use rt; use comm::*; use super::*; - use task; + use thread::Thread; $(#[$a])* #[test] fn f() { $b } } ) } +mod blocking; mod oneshot; mod select; mod shared; @@ -460,15 +457,17 @@ impl<T> UnsafeFlavor<T> for Receiver<T> { /// # Example /// /// ``` +/// use std::thread::Thread; +/// /// // tx is is the sending half (tx for transmission), and rx is the receiving /// // half (rx for receiving). /// let (tx, rx) = channel(); /// /// // Spawn off an expensive computation -/// spawn(move|| { +/// Thread::spawn(move|| { /// # fn expensive_computation() {} /// tx.send(expensive_computation()); -/// }); +/// }).detach(); /// /// // Do some useful work for awhile /// @@ -499,15 +498,17 @@ pub fn channel<T: Send>() -> (Sender<T>, Receiver<T>) { /// # Example /// /// ``` +/// use std::thread::Thread; +/// /// let (tx, rx) = sync_channel(1); /// /// // this returns immediately /// tx.send(1i); /// -/// spawn(move|| { +/// Thread::spawn(move|| { /// // this will block until the previous message has been received /// tx.send(2i); -/// }); +/// }).detach(); /// /// assert_eq!(rx.recv(), 1i); /// assert_eq!(rx.recv(), 2i); @@ -604,12 +605,12 @@ impl<T: Send> Sender<T> { (a, ret) } oneshot::UpDisconnected => (a, Err(t)), - oneshot::UpWoke(task) => { - // This send cannot panic because the task is + oneshot::UpWoke(token) => { + // This send cannot panic because the thread is // asleep (we're looking at it), so the receiver // can't go away. (*a.get()).send(t).ok().unwrap(); - task.wake().map(|t| t.reawaken()); + token.signal(); (a, Ok(())) } } @@ -948,34 +949,33 @@ impl<T: Send> select::Packet for Receiver<T> { } } - fn start_selection(&self, mut task: BlockedTask) -> Result<(), BlockedTask>{ + fn start_selection(&self, mut token: SignalToken) -> StartResult { loop { let (t, new_port) = match *unsafe { self.inner() } { Oneshot(ref p) => { - match unsafe { (*p.get()).start_selection(task) } { - oneshot::SelSuccess => return Ok(()), - oneshot::SelCanceled(task) => return Err(task), + match unsafe { (*p.get()).start_selection(token) } { + oneshot::SelSuccess => return Installed, + oneshot::SelCanceled => return Abort, oneshot::SelUpgraded(t, rx) => (t, rx), } } Stream(ref p) => { - match unsafe { (*p.get()).start_selection(task) } { - stream::SelSuccess => return Ok(()), - stream::SelCanceled(task) => return Err(task), + match unsafe { (*p.get()).start_selection(token) } { + stream::SelSuccess => return Installed, + stream::SelCanceled => return Abort, stream::SelUpgraded(t, rx) => (t, rx), } } Shared(ref p) => { - return unsafe { (*p.get()).start_selection(task) }; + return unsafe { (*p.get()).start_selection(token) }; } Sync(ref p) => { - return unsafe { (*p.get()).start_selection(task) }; + return unsafe { (*p.get()).start_selection(token) }; } }; - task = t; + token = t; unsafe { - mem::swap(self.inner_mut(), - new_port.inner_mut()); + mem::swap(self.inner_mut(), new_port.inner_mut()); } } } @@ -1252,11 +1252,11 @@ mod test { test! { fn oneshot_single_thread_recv_chan_close() { // Receiving on a closed chan will panic - let res = task::try(move|| { + let res = Thread::spawn(move|| { let (tx, rx) = channel::<int>(); drop(tx); rx.recv(); - }); + }).join(); // What is our res? assert!(res.is_err()); } } @@ -1324,9 +1324,9 @@ mod test { spawn(move|| { drop(tx); }); - let res = task::try(move|| { + let res = Thread::spawn(move|| { assert!(rx.recv() == box 10); - }); + }).join(); assert!(res.is_err()); } } @@ -1346,9 +1346,9 @@ mod test { spawn(move|| { drop(rx); }); - let _ = task::try(move|| { + let _ = Thread::spawn(move|| { tx.send(1); - }); + }).join(); } } } @@ -1356,9 +1356,9 @@ mod test { for _ in range(0, stress_factor()) { let (tx, rx) = channel::<int>(); spawn(move|| { - let res = task::try(move|| { + let res = Thread::spawn(move|| { rx.recv(); - }); + }).join(); assert!(res.is_err()); }); spawn(move|| { @@ -1507,7 +1507,7 @@ mod test { tx2.send(()); }); // make sure the other task has gone to sleep - for _ in range(0u, 5000) { task::deschedule(); } + for _ in range(0u, 5000) { Thread::yield_now(); } // upgrade to a shared chan and send a message let t = tx.clone(); @@ -1516,45 +1516,7 @@ mod test { // wait for the child task to exit before we exit rx2.recv(); - } } - - test! { fn sends_off_the_runtime() { - use rustrt::thread::Thread; - - let (tx, rx) = channel(); - let t = Thread::start(move|| { - for _ in range(0u, 1000) { - tx.send(()); - } - }); - for _ in range(0u, 1000) { - rx.recv(); - } - t.join(); - } } - - test! { fn try_recvs_off_the_runtime() { - use rustrt::thread::Thread; - - let (tx, rx) = channel(); - let (cdone, pdone) = channel(); - let t = Thread::start(move|| { - let mut hits = 0u; - while hits < 10 { - match rx.try_recv() { - Ok(()) => { hits += 1; } - Err(Empty) => { Thread::yield_now(); } - Err(Disconnected) => return, - } - } - cdone.send(()); - }); - for _ in range(0u, 10) { - tx.send(()); - } - t.join(); - pdone.recv(); - } } + }} } #[cfg(test)] @@ -1712,11 +1674,11 @@ mod sync_tests { test! { fn oneshot_single_thread_recv_chan_close() { // Receiving on a closed chan will panic - let res = task::try(move|| { + let res = Thread::spawn(move|| { let (tx, rx) = sync_channel::<int>(0); drop(tx); rx.recv(); - }); + }).join(); // What is our res? assert!(res.is_err()); } } @@ -1789,9 +1751,9 @@ mod sync_tests { spawn(move|| { drop(tx); }); - let res = task::try(move|| { + let res = Thread::spawn(move|| { assert!(rx.recv() == box 10); - }); + }).join(); assert!(res.is_err()); } } @@ -1811,9 +1773,9 @@ mod sync_tests { spawn(move|| { drop(rx); }); - let _ = task::try(move|| { + let _ = Thread::spawn(move || { tx.send(1); - }); + }).join(); } } } @@ -1821,9 +1783,9 @@ mod sync_tests { for _ in range(0, stress_factor()) { let (tx, rx) = sync_channel::<int>(0); spawn(move|| { - let res = task::try(move|| { + let res = Thread::spawn(move|| { rx.recv(); - }); + }).join(); assert!(res.is_err()); }); spawn(move|| { @@ -1972,7 +1934,7 @@ mod sync_tests { tx2.send(()); }); // make sure the other task has gone to sleep - for _ in range(0u, 5000) { task::deschedule(); } + for _ in range(0u, 5000) { Thread::yield_now(); } // upgrade to a shared chan and send a message let t = tx.clone(); @@ -1983,29 +1945,6 @@ mod sync_tests { rx2.recv(); } } - test! { fn try_recvs_off_the_runtime() { - use rustrt::thread::Thread; - - let (tx, rx) = sync_channel::<()>(0); - let (cdone, pdone) = channel(); - let t = Thread::start(move|| { - let mut hits = 0u; - while hits < 10 { - match rx.try_recv() { - Ok(()) => { hits += 1; } - Err(Empty) => { Thread::yield_now(); } - Err(Disconnected) => return, - } - } - cdone.send(()); - }); - for _ in range(0u, 10) { - tx.send(()); - } - t.join(); - pdone.recv(); - } } - test! { fn send_opt1() { let (tx, rx) = sync_channel::<int>(0); spawn(move|| { rx.recv(); }); @@ -2064,7 +2003,7 @@ mod sync_tests { test! { fn try_send4() { let (tx, rx) = sync_channel::<int>(0); spawn(move|| { - for _ in range(0u, 1000) { task::deschedule(); } + for _ in range(0u, 1000) { Thread::yield_now(); } assert_eq!(tx.try_send(1), Ok(())); }); assert_eq!(rx.recv(), 1); diff --git a/src/libstd/comm/oneshot.rs b/src/libstd/comm/oneshot.rs index bc34c3e8c52..9c5a6518845 100644 --- a/src/libstd/comm/oneshot.rs +++ b/src/libstd/comm/oneshot.rs @@ -39,18 +39,20 @@ use self::MyUpgrade::*; use core::prelude::*; -use alloc::boxed::Box; +use comm::Receiver; +use comm::blocking::{mod, SignalToken}; use core::mem; -use rustrt::local::Local; -use rustrt::task::{Task, BlockedTask}; - use sync::atomic; -use comm::Receiver; // Various states you can find a port in. -const EMPTY: uint = 0; -const DATA: uint = 1; -const DISCONNECTED: uint = 2; +const EMPTY: uint = 0; // initial state: no data, no blocked reciever +const DATA: uint = 1; // data ready for receiver to take +const DISCONNECTED: uint = 2; // channel is disconnected OR upgraded +// Any other value represents a pointer to a SignalToken value. The +// protocol ensures that when the state moves *to* a pointer, +// ownership of the token is given to the packet, and when the state +// moves *from* a pointer, ownership of the token is transferred to +// whoever changed the state. pub struct Packet<T> { // Internal state of the chan/port pair (stores the blocked task as well) @@ -71,12 +73,12 @@ pub enum Failure<T> { pub enum UpgradeResult { UpSuccess, UpDisconnected, - UpWoke(BlockedTask), + UpWoke(SignalToken), } pub enum SelectionResult<T> { - SelCanceled(BlockedTask), - SelUpgraded(BlockedTask, Receiver<T>), + SelCanceled, + SelUpgraded(SignalToken, Receiver<T>), SelSuccess, } @@ -118,12 +120,10 @@ impl<T: Send> Packet<T> { // Not possible, these are one-use channels DATA => unreachable!(), - // Anything else means that there was a task waiting on the other - // end. We leave the 'DATA' state inside so it'll pick it up on the - // other end. - n => unsafe { - let t = BlockedTask::cast_from_uint(n); - t.wake().map(|t| t.reawaken()); + // There is a thread waiting on the other end. We leave the 'DATA' + // state inside so it'll pick it up on the other end. + ptr => unsafe { + SignalToken::cast_from_uint(ptr).signal(); Ok(()) } } @@ -142,23 +142,17 @@ impl<T: Send> Packet<T> { // Attempt to not block the task (it's a little expensive). If it looks // like we're not empty, then immediately go through to `try_recv`. if self.state.load(atomic::SeqCst) == EMPTY { - let t: Box<Task> = Local::take(); - t.deschedule(1, |task| { - let n = unsafe { task.cast_to_uint() }; - match self.state.compare_and_swap(EMPTY, n, atomic::SeqCst) { - // Nothing on the channel, we legitimately block - EMPTY => Ok(()), - - // If there's data or it's a disconnected channel, then we - // failed the cmpxchg, so we just wake ourselves back up - DATA | DISCONNECTED => { - unsafe { Err(BlockedTask::cast_from_uint(n)) } - } - - // Only one thread is allowed to sleep on this port - _ => unreachable!() - } - }); + let (wait_token, signal_token) = blocking::tokens(); + let ptr = unsafe { signal_token.cast_to_uint() }; + + // race with senders to enter the blocking state + if self.state.compare_and_swap(EMPTY, ptr, atomic::SeqCst) == EMPTY { + wait_token.wait(); + debug_assert!(self.state.load(atomic::SeqCst) != EMPTY); + } else { + // drop the signal token, since we never blocked + drop(unsafe { SignalToken::cast_from_uint(ptr) }); + } } self.try_recv() @@ -197,6 +191,9 @@ impl<T: Send> Packet<T> { } } } + + // We are the sole receiver; there cannot be a blocking + // receiver already. _ => unreachable!() } } @@ -223,7 +220,7 @@ impl<T: Send> Packet<T> { DISCONNECTED => { self.upgrade = prev; UpDisconnected } // If someone's waiting, we gotta wake them up - n => UpWoke(unsafe { BlockedTask::cast_from_uint(n) }) + ptr => UpWoke(unsafe { SignalToken::cast_from_uint(ptr) }) } } @@ -232,9 +229,8 @@ impl<T: Send> Packet<T> { DATA | DISCONNECTED | EMPTY => {} // If someone's waiting, we gotta wake them up - n => unsafe { - let t = BlockedTask::cast_from_uint(n); - t.wake().map(|t| t.reawaken()); + ptr => unsafe { + SignalToken::cast_from_uint(ptr).signal(); } } } @@ -286,13 +282,17 @@ impl<T: Send> Packet<T> { // Attempts to start selection on this port. This can either succeed, fail // because there is data, or fail because there is an upgrade pending. - pub fn start_selection(&mut self, task: BlockedTask) -> SelectionResult<T> { - let n = unsafe { task.cast_to_uint() }; - match self.state.compare_and_swap(EMPTY, n, atomic::SeqCst) { + pub fn start_selection(&mut self, token: SignalToken) -> SelectionResult<T> { + let ptr = unsafe { token.cast_to_uint() }; + match self.state.compare_and_swap(EMPTY, ptr, atomic::SeqCst) { EMPTY => SelSuccess, - DATA => SelCanceled(unsafe { BlockedTask::cast_from_uint(n) }), + DATA => { + drop(unsafe { SignalToken::cast_from_uint(ptr) }); + SelCanceled + } DISCONNECTED if self.data.is_some() => { - SelCanceled(unsafe { BlockedTask::cast_from_uint(n) }) + drop(unsafe { SignalToken::cast_from_uint(ptr) }); + SelCanceled } DISCONNECTED => { match mem::replace(&mut self.upgrade, SendUsed) { @@ -300,8 +300,7 @@ impl<T: Send> Packet<T> { // propagate upwards whether the upgrade can receive // data GoUp(upgrade) => { - SelUpgraded(unsafe { BlockedTask::cast_from_uint(n) }, - upgrade) + SelUpgraded(unsafe { SignalToken::cast_from_uint(ptr) }, upgrade) } // If the other end disconnected without sending an @@ -309,7 +308,8 @@ impl<T: Send> Packet<T> { // disconnected). up => { self.upgrade = up; - SelCanceled(unsafe { BlockedTask::cast_from_uint(n) }) + drop(unsafe { SignalToken::cast_from_uint(ptr) }); + SelCanceled } } } @@ -331,7 +331,7 @@ impl<T: Send> Packet<T> { // If we've got a blocked task, then use an atomic to gain ownership // of it (may fail) - n => self.state.compare_and_swap(n, EMPTY, atomic::SeqCst) + ptr => self.state.compare_and_swap(ptr, EMPTY, atomic::SeqCst) }; // Now that we've got ownership of our state, figure out what to do @@ -358,11 +358,9 @@ impl<T: Send> Packet<T> { } } - // We woke ourselves up from select. Assert that the task should be - // trashed and returned that we don't have any data. - n => { - let t = unsafe { BlockedTask::cast_from_uint(n) }; - t.trash(); + // We woke ourselves up from select. + ptr => unsafe { + drop(SignalToken::cast_from_uint(ptr)); Ok(false) } } diff --git a/src/libstd/comm/select.rs b/src/libstd/comm/select.rs index de2b84b083c..690b5861c22 100644 --- a/src/libstd/comm/select.rs +++ b/src/libstd/comm/select.rs @@ -54,15 +54,13 @@ use core::prelude::*; -use alloc::boxed::Box; use core::cell::Cell; use core::kinds::marker; use core::mem; use core::uint; -use rustrt::local::Local; -use rustrt::task::{Task, BlockedTask}; use comm::Receiver; +use comm::blocking::{mod, SignalToken}; /// The "receiver set" of the select interface. This structure is used to manage /// a set of receivers which are being selected over. @@ -94,9 +92,16 @@ pub struct Handle<'rx, T:'rx> { struct Packets { cur: *mut Handle<'static, ()> } #[doc(hidden)] +#[deriving(PartialEq)] +pub enum StartResult { + Installed, + Abort, +} + +#[doc(hidden)] pub trait Packet { fn can_recv(&self) -> bool; - fn start_selection(&self, task: BlockedTask) -> Result<(), BlockedTask>; + fn start_selection(&self, token: SignalToken) -> StartResult; fn abort_selection(&self) -> bool; } @@ -165,36 +170,39 @@ impl Select { // Most notably, the iterations over all of the receivers shouldn't be // necessary. unsafe { - let mut amt = 0; - for p in self.iter() { - amt += 1; - if do_preflight_checks && (*p).packet.can_recv() { - return (*p).id; + // Stage 1: preflight checks. Look for any packets ready to receive + if do_preflight_checks { + for handle in self.iter() { + if (*handle).packet.can_recv() { + return (*handle).id(); + } } } - assert!(amt > 0); - let mut ready_index = amt; - let mut ready_id = uint::MAX; - let mut iter = self.iter().enumerate(); - - // Acquire a number of blocking contexts, and block on each one - // sequentially until one fails. If one fails, then abort - // immediately so we can go unblock on all the other receivers. - let task: Box<Task> = Local::take(); - task.deschedule(amt, |task| { - // Prepare for the block - let (i, handle) = iter.next().unwrap(); - match (*handle).packet.start_selection(task) { - Ok(()) => Ok(()), - Err(task) => { - ready_index = i; - ready_id = (*handle).id; - Err(task) + // Stage 2: begin the blocking process + // + // Create a number of signal tokens, and install each one + // sequentially until one fails. If one fails, then abort the + // selection on the already-installed tokens. + let (wait_token, signal_token) = blocking::tokens(); + for (i, handle) in self.iter().enumerate() { + match (*handle).packet.start_selection(signal_token.clone()) { + StartResult::Installed => {} + StartResult::Abort => { + // Go back and abort the already-begun selections + for handle in self.iter().take(i) { + (*handle).packet.abort_selection(); + } + return (*handle).id; } } - }); + } + + // Stage 3: no messages available, actually block + wait_token.wait(); + // Stage 4: there *must* be message available; find it. + // // Abort the selection process on each receiver. If the abort // process returns `true`, then that means that the receiver is // ready to receive some data. Note that this also means that the @@ -216,12 +224,14 @@ impl Select { // A rewrite should focus on avoiding a yield loop, and for now this // implementation is tying us over to a more efficient "don't // iterate over everything every time" implementation. - for handle in self.iter().take(ready_index) { + let mut ready_id = uint::MAX; + for handle in self.iter() { if (*handle).packet.abort_selection() { ready_id = (*handle).id; } } + // We must have found a ready receiver assert!(ready_id != uint::MAX); return ready_id; } @@ -404,10 +414,10 @@ mod test { let (tx3, rx3) = channel::<int>(); spawn(move|| { - for _ in range(0u, 20) { task::deschedule(); } + for _ in range(0u, 20) { Thread::yield_now(); } tx1.send(1); rx3.recv(); - for _ in range(0u, 20) { task::deschedule(); } + for _ in range(0u, 20) { Thread::yield_now(); } }); select! { @@ -427,7 +437,7 @@ mod test { let (tx3, rx3) = channel::<()>(); spawn(move|| { - for _ in range(0u, 20) { task::deschedule(); } + for _ in range(0u, 20) { Thread::yield_now(); } tx1.send(1); tx2.send(2); rx3.recv(); @@ -528,7 +538,7 @@ mod test { tx3.send(()); }); - for _ in range(0u, 1000) { task::deschedule(); } + for _ in range(0u, 1000) { Thread::yield_now(); } drop(tx1.clone()); tx2.send(()); rx3.recv(); @@ -631,7 +641,7 @@ mod test { tx2.send(()); }); - for _ in range(0u, 100) { task::deschedule() } + for _ in range(0u, 100) { Thread::yield_now() } tx1.send(()); rx2.recv(); } } @@ -650,7 +660,7 @@ mod test { tx2.send(()); }); - for _ in range(0u, 100) { task::deschedule() } + for _ in range(0u, 100) { Thread::yield_now() } tx1.send(()); rx2.recv(); } } @@ -668,7 +678,7 @@ mod test { tx2.send(()); }); - for _ in range(0u, 100) { task::deschedule() } + for _ in range(0u, 100) { Thread::yield_now() } tx1.send(()); rx2.recv(); } } @@ -684,7 +694,7 @@ mod test { test! { fn sync2() { let (tx, rx) = sync_channel::<int>(0); spawn(move|| { - for _ in range(0u, 100) { task::deschedule() } + for _ in range(0u, 100) { Thread::yield_now() } tx.send(1); }); select! { diff --git a/src/libstd/comm/shared.rs b/src/libstd/comm/shared.rs index 13b5e10fcd3..1022694e634 100644 --- a/src/libstd/comm/shared.rs +++ b/src/libstd/comm/shared.rs @@ -22,15 +22,15 @@ pub use self::Failure::*; use core::prelude::*; -use alloc::boxed::Box; use core::cmp; use core::int; -use rustrt::local::Local; -use rustrt::task::{Task, BlockedTask}; -use rustrt::thread::Thread; use sync::{atomic, Mutex, MutexGuard}; use comm::mpsc_queue as mpsc; +use comm::blocking::{mod, SignalToken}; +use comm::select::StartResult; +use comm::select::StartResult::*; +use thread::Thread; const DISCONNECTED: int = int::MIN; const FUDGE: int = 1024; @@ -43,7 +43,7 @@ pub struct Packet<T> { queue: mpsc::Queue<T>, cnt: atomic::AtomicInt, // How many items are on this channel steals: int, // How many times has a port received without blocking? - to_wake: atomic::AtomicUint, // Task to wake up + to_wake: atomic::AtomicUint, // SignalToken for wake up // The number of channels which are currently using this packet. channels: atomic::AtomicInt, @@ -95,41 +95,34 @@ impl<T: Send> Packet<T> { // // This can only be called at channel-creation time pub fn inherit_blocker(&mut self, - task: Option<BlockedTask>, + token: Option<SignalToken>, guard: MutexGuard<()>) { - match task { - Some(task) => { - assert_eq!(self.cnt.load(atomic::SeqCst), 0); - assert_eq!(self.to_wake.load(atomic::SeqCst), 0); - self.to_wake.store(unsafe { task.cast_to_uint() }, - atomic::SeqCst); - self.cnt.store(-1, atomic::SeqCst); - - // This store is a little sketchy. What's happening here is - // that we're transferring a blocker from a oneshot or stream - // channel to this shared channel. In doing so, we never - // spuriously wake them up and rather only wake them up at the - // appropriate time. This implementation of shared channels - // assumes that any blocking recv() will undo the increment of - // steals performed in try_recv() once the recv is complete. - // This thread that we're inheriting, however, is not in the - // middle of recv. Hence, the first time we wake them up, - // they're going to wake up from their old port, move on to the - // upgraded port, and then call the block recv() function. - // - // When calling this function, they'll find there's data - // immediately available, counting it as a steal. This in fact - // wasn't a steal because we appropriately blocked them waiting - // for data. - // - // To offset this bad increment, we initially set the steal - // count to -1. You'll find some special code in - // abort_selection() as well to ensure that this -1 steal count - // doesn't escape too far. - self.steals = -1; - } - None => {} - } + token.map(|token| { + assert_eq!(self.cnt.load(atomic::SeqCst), 0); + assert_eq!(self.to_wake.load(atomic::SeqCst), 0); + self.to_wake.store(unsafe { token.cast_to_uint() }, atomic::SeqCst); + self.cnt.store(-1, atomic::SeqCst); + + // This store is a little sketchy. What's happening here is that + // we're transferring a blocker from a oneshot or stream channel to + // this shared channel. In doing so, we never spuriously wake them + // up and rather only wake them up at the appropriate time. This + // implementation of shared channels assumes that any blocking + // recv() will undo the increment of steals performed in try_recv() + // once the recv is complete. This thread that we're inheriting, + // however, is not in the middle of recv. Hence, the first time we + // wake them up, they're going to wake up from their old port, move + // on to the upgraded port, and then call the block recv() function. + // + // When calling this function, they'll find there's data immediately + // available, counting it as a steal. This in fact wasn't a steal + // because we appropriately blocked them waiting for data. + // + // To offset this bad increment, we initially set the steal count to + // -1. You'll find some special code in abort_selection() as well to + // ensure that this -1 steal count doesn't escape too far. + self.steals = -1; + }); // When the shared packet is constructed, we grabbed this lock. The // purpose of this lock is to ensure that abort_selection() doesn't @@ -175,7 +168,7 @@ impl<T: Send> Packet<T> { self.queue.push(t); match self.cnt.fetch_add(1, atomic::SeqCst) { -1 => { - self.take_to_wake().wake().map(|t| t.reawaken()); + self.take_to_wake().signal(); } // In this case, we have possibly failed to send our data, and @@ -232,10 +225,10 @@ impl<T: Send> Packet<T> { data => return data, } - let task: Box<Task> = Local::take(); - task.deschedule(1, |task| { - self.decrement(task) - }); + let (wait_token, signal_token) = blocking::tokens(); + if self.decrement(signal_token) == Installed { + wait_token.wait() + } match self.try_recv() { data @ Ok(..) => { self.steals -= 1; data } @@ -244,10 +237,11 @@ impl<T: Send> Packet<T> { } // Essentially the exact same thing as the stream decrement function. - fn decrement(&mut self, task: BlockedTask) -> Result<(), BlockedTask> { + // Returns true if blocking should proceed. + fn decrement(&mut self, token: SignalToken) -> StartResult { assert_eq!(self.to_wake.load(atomic::SeqCst), 0); - let n = unsafe { task.cast_to_uint() }; - self.to_wake.store(n, atomic::SeqCst); + let ptr = unsafe { token.cast_to_uint() }; + self.to_wake.store(ptr, atomic::SeqCst); let steals = self.steals; self.steals = 0; @@ -258,12 +252,13 @@ impl<T: Send> Packet<T> { // data, we successfully sleep n => { assert!(n >= 0); - if n - steals <= 0 { return Ok(()) } + if n - steals <= 0 { return Installed } } } self.to_wake.store(0, atomic::SeqCst); - Err(unsafe { BlockedTask::cast_from_uint(n) }) + drop(unsafe { SignalToken::cast_from_uint(ptr) }); + Abort } pub fn try_recv(&mut self) -> Result<T, Failure> { @@ -271,20 +266,19 @@ impl<T: Send> Packet<T> { mpsc::Data(t) => Some(t), mpsc::Empty => None, - // This is a bit of an interesting case. The channel is - // reported as having data available, but our pop() has - // failed due to the queue being in an inconsistent state. - // This means that there is some pusher somewhere which has - // yet to complete, but we are guaranteed that a pop will - // eventually succeed. In this case, we spin in a yield loop - // because the remote sender should finish their enqueue + // This is a bit of an interesting case. The channel is reported as + // having data available, but our pop() has failed due to the queue + // being in an inconsistent state. This means that there is some + // pusher somewhere which has yet to complete, but we are guaranteed + // that a pop will eventually succeed. In this case, we spin in a + // yield loop because the remote sender should finish their enqueue // operation "very quickly". // // Avoiding this yield loop would require a different queue - // abstraction which provides the guarantee that after M - // pushes have succeeded, at least M pops will succeed. The - // current queues guarantee that if there are N active - // pushes, you can pop N times once all N have finished. + // abstraction which provides the guarantee that after M pushes have + // succeeded, at least M pops will succeed. The current queues + // guarantee that if there are N active pushes, you can pop N times + // once all N have finished. mpsc::Inconsistent => { let data; loop { @@ -354,7 +348,7 @@ impl<T: Send> Packet<T> { } match self.cnt.swap(DISCONNECTED, atomic::SeqCst) { - -1 => { self.take_to_wake().wake().map(|t| t.reawaken()); } + -1 => { self.take_to_wake().signal(); } DISCONNECTED => {} n => { assert!(n >= 0); } } @@ -366,8 +360,7 @@ impl<T: Send> Packet<T> { self.port_dropped.store(true, atomic::SeqCst); let mut steals = self.steals; while { - let cnt = self.cnt.compare_and_swap( - steals, DISCONNECTED, atomic::SeqCst); + let cnt = self.cnt.compare_and_swap(steals, DISCONNECTED, atomic::SeqCst); cnt != DISCONNECTED && cnt != steals } { // See the discussion in 'try_recv' for why we yield @@ -382,11 +375,11 @@ impl<T: Send> Packet<T> { } // Consumes ownership of the 'to_wake' field. - fn take_to_wake(&mut self) -> BlockedTask { - let task = self.to_wake.load(atomic::SeqCst); + fn take_to_wake(&mut self) -> SignalToken { + let ptr = self.to_wake.load(atomic::SeqCst); self.to_wake.store(0, atomic::SeqCst); - assert!(task != 0); - unsafe { BlockedTask::cast_from_uint(task) } + assert!(ptr != 0); + unsafe { SignalToken::cast_from_uint(ptr) } } //////////////////////////////////////////////////////////////////////////// @@ -414,19 +407,18 @@ impl<T: Send> Packet<T> { } } - // Inserts the blocked task for selection on this port, returning it back if - // the port already has data on it. + // Inserts the signal token for selection on this port, returning true if + // blocking should proceed. // // The code here is the same as in stream.rs, except that it doesn't need to // peek at the channel to see if an upgrade is pending. - pub fn start_selection(&mut self, - task: BlockedTask) -> Result<(), BlockedTask> { - match self.decrement(task) { - Ok(()) => Ok(()), - Err(task) => { + pub fn start_selection(&mut self, token: SignalToken) -> StartResult { + match self.decrement(token) { + Installed => Installed, + Abort => { let prev = self.bump(1); assert!(prev == DISCONNECTED || prev >= 0); - return Err(task); + Abort } } } @@ -464,7 +456,7 @@ impl<T: Send> Packet<T> { let cur = prev + steals + 1; assert!(cur >= 0); if prev < 0 { - self.take_to_wake().trash(); + drop(self.take_to_wake()); } else { while self.to_wake.load(atomic::SeqCst) != 0 { Thread::yield_now(); diff --git a/src/libstd/comm/stream.rs b/src/libstd/comm/stream.rs index 06ab4f4427a..b68f626060e 100644 --- a/src/libstd/comm/stream.rs +++ b/src/libstd/comm/stream.rs @@ -24,16 +24,14 @@ use self::Message::*; use core::prelude::*; -use alloc::boxed::Box; use core::cmp; use core::int; -use rustrt::local::Local; -use rustrt::task::{Task, BlockedTask}; -use rustrt::thread::Thread; +use thread::Thread; use sync::atomic; use comm::spsc_queue as spsc; use comm::Receiver; +use comm::blocking::{mod, SignalToken}; const DISCONNECTED: int = int::MIN; #[cfg(test)] @@ -46,7 +44,7 @@ pub struct Packet<T> { cnt: atomic::AtomicInt, // How many items are on this channel steals: int, // How many times has a port received without blocking? - to_wake: atomic::AtomicUint, // Task to wake up + to_wake: atomic::AtomicUint, // SignalToken for the blocked thread to wake up port_dropped: atomic::AtomicBool, // flag if the channel has been destroyed. } @@ -60,13 +58,13 @@ pub enum Failure<T> { pub enum UpgradeResult { UpSuccess, UpDisconnected, - UpWoke(BlockedTask), + UpWoke(SignalToken), } pub enum SelectionResult<T> { SelSuccess, - SelCanceled(BlockedTask), - SelUpgraded(BlockedTask, Receiver<T>), + SelCanceled, + SelUpgraded(SignalToken, Receiver<T>), } // Any message could contain an "upgrade request" to a new shared port, so the @@ -89,7 +87,6 @@ impl<T: Send> Packet<T> { } } - pub fn send(&mut self, t: T) -> Result<(), T> { // If the other port has deterministically gone away, then definitely // must return the data back up the stack. Otherwise, the data is @@ -98,10 +95,11 @@ impl<T: Send> Packet<T> { match self.do_send(Data(t)) { UpSuccess | UpDisconnected => {}, - UpWoke(task) => { task.wake().map(|t| t.reawaken()); } + UpWoke(token) => { token.signal(); } } Ok(()) } + pub fn upgrade(&mut self, up: Receiver<T>) -> UpgradeResult { // If the port has gone away, then there's no need to proceed any // further. @@ -144,20 +142,20 @@ impl<T: Send> Packet<T> { } // Consumes ownership of the 'to_wake' field. - fn take_to_wake(&mut self) -> BlockedTask { - let task = self.to_wake.load(atomic::SeqCst); + fn take_to_wake(&mut self) -> SignalToken { + let ptr = self.to_wake.load(atomic::SeqCst); self.to_wake.store(0, atomic::SeqCst); - assert!(task != 0); - unsafe { BlockedTask::cast_from_uint(task) } + assert!(ptr != 0); + unsafe { SignalToken::cast_from_uint(ptr) } } // Decrements the count on the channel for a sleeper, returning the sleeper // back if it shouldn't sleep. Note that this is the location where we take // steals into account. - fn decrement(&mut self, task: BlockedTask) -> Result<(), BlockedTask> { + fn decrement(&mut self, token: SignalToken) -> Result<(), SignalToken> { assert_eq!(self.to_wake.load(atomic::SeqCst), 0); - let n = unsafe { task.cast_to_uint() }; - self.to_wake.store(n, atomic::SeqCst); + let ptr = unsafe { token.cast_to_uint() }; + self.to_wake.store(ptr, atomic::SeqCst); let steals = self.steals; self.steals = 0; @@ -173,7 +171,7 @@ impl<T: Send> Packet<T> { } self.to_wake.store(0, atomic::SeqCst); - Err(unsafe { BlockedTask::cast_from_uint(n) }) + Err(unsafe { SignalToken::cast_from_uint(ptr) }) } pub fn recv(&mut self) -> Result<T, Failure<T>> { @@ -185,10 +183,10 @@ impl<T: Send> Packet<T> { // Welp, our channel has no data. Deschedule the current task and // initiate the blocking protocol. - let task: Box<Task> = Local::take(); - task.deschedule(1, |task| { - self.decrement(task) - }); + let (wait_token, signal_token) = blocking::tokens(); + if self.decrement(signal_token).is_ok() { + wait_token.wait() + } match self.try_recv() { // Messages which actually popped from the queue shouldn't count as @@ -269,7 +267,7 @@ impl<T: Send> Packet<T> { // Dropping a channel is pretty simple, we just flag it as disconnected // and then wakeup a blocker if there is one. match self.cnt.swap(DISCONNECTED, atomic::SeqCst) { - -1 => { self.take_to_wake().wake().map(|t| t.reawaken()); } + -1 => { self.take_to_wake().signal(); } DISCONNECTED => {} n => { assert!(n >= 0); } } @@ -364,19 +362,19 @@ impl<T: Send> Packet<T> { // Attempts to start selecting on this port. Like a oneshot, this can fail // immediately because of an upgrade. - pub fn start_selection(&mut self, task: BlockedTask) -> SelectionResult<T> { - match self.decrement(task) { + pub fn start_selection(&mut self, token: SignalToken) -> SelectionResult<T> { + match self.decrement(token) { Ok(()) => SelSuccess, - Err(task) => { + Err(token) => { let ret = match self.queue.peek() { Some(&GoUp(..)) => { match self.queue.pop() { - Some(GoUp(port)) => SelUpgraded(task, port), + Some(GoUp(port)) => SelUpgraded(token, port), _ => unreachable!(), } } - Some(..) => SelCanceled(task), - None => SelCanceled(task), + Some(..) => SelCanceled, + None => SelCanceled, }; // Undo our decrement above, and we should be guaranteed that the // previous value is positive because we're not going to sleep @@ -439,7 +437,7 @@ impl<T: Send> Packet<T> { // final solution but rather out of necessity for now to get // something working. if prev < 0 { - self.take_to_wake().trash(); + drop(self.take_to_wake()); } else { while self.to_wake.load(atomic::SeqCst) != 0 { Thread::yield_now(); diff --git a/src/libstd/comm/sync.rs b/src/libstd/comm/sync.rs index a2e839e134c..f75186e70e3 100644 --- a/src/libstd/comm/sync.rs +++ b/src/libstd/comm/sync.rs @@ -38,24 +38,19 @@ use core::prelude::*; pub use self::Failure::*; use self::Blocker::*; -use alloc::boxed::Box; use vec::Vec; use core::mem; -use core::cell::UnsafeCell; -use rustrt::local::Local; -use rustrt::mutex::{NativeMutex, LockGuard}; -use rustrt::task::{Task, BlockedTask}; -use sync::atomic; +use sync::{atomic, Mutex, MutexGuard}; +use comm::blocking::{mod, WaitToken, SignalToken}; +use comm::select::StartResult::{mod, Installed, Abort}; pub struct Packet<T> { /// Only field outside of the mutex. Just done for kicks, but mainly because /// the other shared channel already had the code implemented channels: atomic::AtomicUint, - /// The state field is protected by this mutex - lock: NativeMutex, - state: UnsafeCell<State<T>>, + lock: Mutex<State<T>>, } struct State<T> { @@ -74,10 +69,10 @@ struct State<T> { canceled: Option<&'static mut bool>, } -/// Possible flavors of tasks who can be blocked on this channel. +/// Possible flavors of threads who can be blocked on this channel. enum Blocker { - BlockedSender(BlockedTask), - BlockedReceiver(BlockedTask), + BlockedSender(SignalToken), + BlockedReceiver(SignalToken), NoneBlocked } @@ -89,7 +84,7 @@ struct Queue { } struct Node { - task: Option<BlockedTask>, + token: Option<SignalToken>, next: *mut Node, } @@ -106,36 +101,36 @@ pub enum Failure { Disconnected, } -/// Atomically blocks the current task, placing it into `slot`, unlocking `lock` +/// Atomically blocks the current thread, placing it into `slot`, unlocking `lock` /// in the meantime. This re-locks the mutex upon returning. -fn wait(slot: &mut Blocker, f: fn(BlockedTask) -> Blocker, - lock: &NativeMutex) { - let me: Box<Task> = Local::take(); - me.deschedule(1, |task| { - match mem::replace(slot, f(task)) { - NoneBlocked => {} - _ => unreachable!(), - } - unsafe { lock.unlock_noguard(); } - Ok(()) - }); - unsafe { lock.lock_noguard(); } +fn wait<'a, 'b, T: Send>(lock: &'a Mutex<State<T>>, + mut guard: MutexGuard<'b, State<T>>, + f: fn(SignalToken) -> Blocker) + -> MutexGuard<'a, State<T>> +{ + let (wait_token, signal_token) = blocking::tokens(); + match mem::replace(&mut guard.blocker, f(signal_token)) { + NoneBlocked => {} + _ => unreachable!(), + } + drop(guard); // unlock + wait_token.wait(); // block + lock.lock() // relock } -/// Wakes up a task, dropping the lock at the correct time -fn wakeup(task: BlockedTask, guard: LockGuard) { +/// Wakes up a thread, dropping the lock at the correct time +fn wakeup<T>(token: SignalToken, guard: MutexGuard<State<T>>) { // We need to be careful to wake up the waiting task *outside* of the mutex // in case it incurs a context switch. - mem::drop(guard); - task.wake().map(|t| t.reawaken()); + drop(guard); + token.signal(); } impl<T: Send> Packet<T> { pub fn new(cap: uint) -> Packet<T> { Packet { channels: atomic::AtomicUint::new(1), - lock: unsafe { NativeMutex::new() }, - state: UnsafeCell::new(State { + lock: Mutex::new(State { disconnected: false, blocker: NoneBlocked, cap: cap, @@ -153,68 +148,66 @@ impl<T: Send> Packet<T> { } } - // Locks this channel, returning a guard for the state and the mutable state - // itself. Care should be taken to ensure that the state does not escape the - // guard! - // - // Note that we're ok promoting an & reference to an &mut reference because - // the lock ensures that we're the only ones in the world with a pointer to - // the state. - fn lock<'a>(&'a self) -> (LockGuard<'a>, &'a mut State<T>) { - unsafe { - let guard = self.lock.lock(); - (guard, &mut *self.state.get()) + // wait until a send slot is available, returning locked access to + // the channel state. + fn acquire_send_slot(&self) -> MutexGuard<State<T>> { + let mut node = Node { token: None, next: 0 as *mut Node }; + loop { + let mut guard = self.lock.lock(); + // are we ready to go? + if guard.disconnected || guard.buf.size() < guard.buf.cap() { + return guard; + } + // no room; actually block + let wait_token = guard.queue.enqueue(&mut node); + drop(guard); + wait_token.wait(); } } pub fn send(&self, t: T) -> Result<(), T> { - let (guard, state) = self.lock(); - - // wait for a slot to become available, and enqueue the data - while !state.disconnected && state.buf.size() == state.buf.cap() { - state.queue.enqueue(&self.lock); - } - if state.disconnected { return Err(t) } - state.buf.enqueue(t); + let mut guard = self.acquire_send_slot(); + if guard.disconnected { return Err(t) } + guard.buf.enqueue(t); - match mem::replace(&mut state.blocker, NoneBlocked) { + match mem::replace(&mut guard.blocker, NoneBlocked) { // if our capacity is 0, then we need to wait for a receiver to be // available to take our data. After waiting, we check again to make // sure the port didn't go away in the meantime. If it did, we need // to hand back our data. - NoneBlocked if state.cap == 0 => { + NoneBlocked if guard.cap == 0 => { let mut canceled = false; - assert!(state.canceled.is_none()); - state.canceled = Some(unsafe { mem::transmute(&mut canceled) }); - wait(&mut state.blocker, BlockedSender, &self.lock); - if canceled {Err(state.buf.dequeue())} else {Ok(())} + assert!(guard.canceled.is_none()); + guard.canceled = Some(unsafe { mem::transmute(&mut canceled) }); + let mut guard = wait(&self.lock, guard, BlockedSender); + if canceled {Err(guard.buf.dequeue())} else {Ok(())} } // success, we buffered some data NoneBlocked => Ok(()), // success, someone's about to receive our buffered data. - BlockedReceiver(task) => { wakeup(task, guard); Ok(()) } + BlockedReceiver(token) => { wakeup(token, guard); Ok(()) } BlockedSender(..) => panic!("lolwut"), } } pub fn try_send(&self, t: T) -> Result<(), super::TrySendError<T>> { - let (guard, state) = self.lock(); - if state.disconnected { + let mut guard = self.lock.lock(); + if guard.disconnected { Err(super::RecvDisconnected(t)) - } else if state.buf.size() == state.buf.cap() { + } else if guard.buf.size() == guard.buf.cap() { Err(super::Full(t)) - } else if state.cap == 0 { + } else if guard.cap == 0 { // With capacity 0, even though we have buffer space we can't // transfer the data unless there's a receiver waiting. - match mem::replace(&mut state.blocker, NoneBlocked) { + match mem::replace(&mut guard.blocker, NoneBlocked) { NoneBlocked => Err(super::Full(t)), BlockedSender(..) => unreachable!(), - BlockedReceiver(task) => { - state.buf.enqueue(t); - wakeup(task, guard); + BlockedReceiver(token) => { + guard.buf.enqueue(t); + wakeup(token, guard); Ok(()) } } @@ -222,10 +215,10 @@ impl<T: Send> Packet<T> { // If the buffer has some space and the capacity isn't 0, then we // just enqueue the data for later retrieval, ensuring to wake up // any blocked receiver if there is one. - assert!(state.buf.size() < state.buf.cap()); - state.buf.enqueue(t); - match mem::replace(&mut state.blocker, NoneBlocked) { - BlockedReceiver(task) => wakeup(task, guard), + assert!(guard.buf.size() < guard.buf.cap()); + guard.buf.enqueue(t); + match mem::replace(&mut guard.blocker, NoneBlocked) { + BlockedReceiver(token) => wakeup(token, guard), NoneBlocked => {} BlockedSender(..) => unreachable!(), } @@ -238,34 +231,34 @@ impl<T: Send> Packet<T> { // When reading this, remember that there can only ever be one receiver at // time. pub fn recv(&self) -> Result<T, ()> { - let (guard, state) = self.lock(); + let mut guard = self.lock.lock(); // Wait for the buffer to have something in it. No need for a while loop // because we're the only receiver. let mut waited = false; - if !state.disconnected && state.buf.size() == 0 { - wait(&mut state.blocker, BlockedReceiver, &self.lock); + if !guard.disconnected && guard.buf.size() == 0 { + guard = wait(&self.lock, guard, BlockedReceiver); waited = true; } - if state.disconnected && state.buf.size() == 0 { return Err(()) } + if guard.disconnected && guard.buf.size() == 0 { return Err(()) } // Pick up the data, wake up our neighbors, and carry on - assert!(state.buf.size() > 0); - let ret = state.buf.dequeue(); - self.wakeup_senders(waited, guard, state); + assert!(guard.buf.size() > 0); + let ret = guard.buf.dequeue(); + self.wakeup_senders(waited, guard); return Ok(ret); } pub fn try_recv(&self) -> Result<T, Failure> { - let (guard, state) = self.lock(); + let mut guard = self.lock.lock(); // Easy cases first - if state.disconnected { return Err(Disconnected) } - if state.buf.size() == 0 { return Err(Empty) } + if guard.disconnected { return Err(Disconnected) } + if guard.buf.size() == 0 { return Err(Empty) } // Be sure to wake up neighbors - let ret = Ok(state.buf.dequeue()); - self.wakeup_senders(false, guard, state); + let ret = Ok(guard.buf.dequeue()); + self.wakeup_senders(false, guard); return ret; } @@ -275,31 +268,29 @@ impl<T: Send> Packet<T> { // * `waited` - flag if the receiver blocked to receive some data, or if it // just picked up some data on the way out // * `guard` - the lock guard that is held over this channel's lock - fn wakeup_senders(&self, waited: bool, - guard: LockGuard, - state: &mut State<T>) { - let pending_sender1: Option<BlockedTask> = state.queue.dequeue(); + fn wakeup_senders(&self, waited: bool, mut guard: MutexGuard<State<T>>) { + let pending_sender1: Option<SignalToken> = guard.queue.dequeue(); // If this is a no-buffer channel (cap == 0), then if we didn't wait we // need to ACK the sender. If we waited, then the sender waking us up // was already the ACK. - let pending_sender2 = if state.cap == 0 && !waited { - match mem::replace(&mut state.blocker, NoneBlocked) { + let pending_sender2 = if guard.cap == 0 && !waited { + match mem::replace(&mut guard.blocker, NoneBlocked) { NoneBlocked => None, BlockedReceiver(..) => unreachable!(), - BlockedSender(task) => { - state.canceled.take(); - Some(task) + BlockedSender(token) => { + guard.canceled.take(); + Some(token) } } } else { None }; - mem::drop((state, guard)); + mem::drop(guard); // only outside of the lock do we wake up the pending tasks - pending_sender1.map(|t| t.wake().map(|t| t.reawaken())); - pending_sender2.map(|t| t.wake().map(|t| t.reawaken())); + pending_sender1.map(|t| t.signal()); + pending_sender2.map(|t| t.signal()); } // Prepares this shared packet for a channel clone, essentially just bumping @@ -316,54 +307,54 @@ impl<T: Send> Packet<T> { } // Not much to do other than wake up a receiver if one's there - let (guard, state) = self.lock(); - if state.disconnected { return } - state.disconnected = true; - match mem::replace(&mut state.blocker, NoneBlocked) { + let mut guard = self.lock.lock(); + if guard.disconnected { return } + guard.disconnected = true; + match mem::replace(&mut guard.blocker, NoneBlocked) { NoneBlocked => {} BlockedSender(..) => unreachable!(), - BlockedReceiver(task) => wakeup(task, guard), + BlockedReceiver(token) => wakeup(token, guard), } } pub fn drop_port(&self) { - let (guard, state) = self.lock(); + let mut guard = self.lock.lock(); - if state.disconnected { return } - state.disconnected = true; + if guard.disconnected { return } + guard.disconnected = true; // If the capacity is 0, then the sender may want its data back after // we're disconnected. Otherwise it's now our responsibility to destroy // the buffered data. As with many other portions of this code, this // needs to be careful to destroy the data *outside* of the lock to // prevent deadlock. - let _data = if state.cap != 0 { - mem::replace(&mut state.buf.buf, Vec::new()) + let _data = if guard.cap != 0 { + mem::replace(&mut guard.buf.buf, Vec::new()) } else { Vec::new() }; - let mut queue = mem::replace(&mut state.queue, Queue { + let mut queue = mem::replace(&mut guard.queue, Queue { head: 0 as *mut Node, tail: 0 as *mut Node, }); - let waiter = match mem::replace(&mut state.blocker, NoneBlocked) { + let waiter = match mem::replace(&mut guard.blocker, NoneBlocked) { NoneBlocked => None, - BlockedSender(task) => { - *state.canceled.take().unwrap() = true; - Some(task) + BlockedSender(token) => { + *guard.canceled.take().unwrap() = true; + Some(token) } BlockedReceiver(..) => unreachable!(), }; - mem::drop((state, guard)); + mem::drop(guard); loop { match queue.dequeue() { - Some(task) => { task.wake().map(|t| t.reawaken()); } + Some(token) => { token.signal(); } None => break, } } - waiter.map(|t| t.wake().map(|t| t.reawaken())); + waiter.map(|t| t.signal()); } //////////////////////////////////////////////////////////////////////////// @@ -373,23 +364,23 @@ impl<T: Send> Packet<T> { // If Ok, the value is whether this port has data, if Err, then the upgraded // port needs to be checked instead of this one. pub fn can_recv(&self) -> bool { - let (_g, state) = self.lock(); - state.disconnected || state.buf.size() > 0 + let guard = self.lock.lock(); + guard.disconnected || guard.buf.size() > 0 } // Attempts to start selection on this port. This can either succeed or fail // because there is data waiting. - pub fn start_selection(&self, task: BlockedTask) -> Result<(), BlockedTask>{ - let (_g, state) = self.lock(); - if state.disconnected || state.buf.size() > 0 { - Err(task) + pub fn start_selection(&self, token: SignalToken) -> StartResult { + let mut guard = self.lock.lock(); + if guard.disconnected || guard.buf.size() > 0 { + Abort } else { - match mem::replace(&mut state.blocker, BlockedReceiver(task)) { + match mem::replace(&mut guard.blocker, BlockedReceiver(token)) { NoneBlocked => {} BlockedSender(..) => unreachable!(), BlockedReceiver(..) => unreachable!(), } - Ok(()) + Installed } } @@ -398,14 +389,14 @@ impl<T: Send> Packet<T> { // // The return value indicates whether there's data on this port. pub fn abort_selection(&self) -> bool { - let (_g, state) = self.lock(); - match mem::replace(&mut state.blocker, NoneBlocked) { + let mut guard = self.lock.lock(); + match mem::replace(&mut guard.blocker, NoneBlocked) { NoneBlocked => true, - BlockedSender(task) => { - state.blocker = BlockedSender(task); + BlockedSender(token) => { + guard.blocker = BlockedSender(token); true } - BlockedReceiver(task) => { task.trash(); false } + BlockedReceiver(token) => { drop(token); false } } } } @@ -414,9 +405,9 @@ impl<T: Send> Packet<T> { impl<T: Send> Drop for Packet<T> { fn drop(&mut self) { assert_eq!(self.channels.load(atomic::SeqCst), 0); - let (_g, state) = self.lock(); - assert!(state.queue.dequeue().is_none()); - assert!(state.canceled.is_none()); + let mut guard = self.lock.lock(); + assert!(guard.queue.dequeue().is_none()); + assert!(guard.canceled.is_none()); } } @@ -449,31 +440,25 @@ impl<T> Buffer<T> { //////////////////////////////////////////////////////////////////////////////// impl Queue { - fn enqueue(&mut self, lock: &NativeMutex) { - let task: Box<Task> = Local::take(); - let mut node = Node { - task: None, - next: 0 as *mut Node, - }; - task.deschedule(1, |task| { - node.task = Some(task); - if self.tail.is_null() { - self.head = &mut node as *mut Node; - self.tail = &mut node as *mut Node; - } else { - unsafe { - (*self.tail).next = &mut node as *mut Node; - self.tail = &mut node as *mut Node; - } + fn enqueue(&mut self, node: &mut Node) -> WaitToken { + let (wait_token, signal_token) = blocking::tokens(); + node.token = Some(signal_token); + node.next = 0 as *mut Node; + + if self.tail.is_null() { + self.head = node as *mut Node; + self.tail = node as *mut Node; + } else { + unsafe { + (*self.tail).next = node as *mut Node; + self.tail = node as *mut Node; } - unsafe { lock.unlock_noguard(); } - Ok(()) - }); - unsafe { lock.lock_noguard(); } - assert!(node.next.is_null()); + } + + wait_token } - fn dequeue(&mut self) -> Option<BlockedTask> { + fn dequeue(&mut self) -> Option<SignalToken> { if self.head.is_null() { return None } @@ -484,7 +469,7 @@ impl Queue { } unsafe { (*node).next = 0 as *mut Node; - Some((*node).task.take().unwrap()) + Some((*node).token.take().unwrap()) } } } |
