about summary refs log tree commit diff
path: root/src/libstd/comm/sync.rs
diff options
context:
space:
mode:
authorAaron Turon <aturon@mozilla.com>2014-12-06 18:34:37 -0800
committerAaron Turon <aturon@mozilla.com>2014-12-18 23:31:51 -0800
commit43ae4b3301cc0605839778ecf59effb32b752e33 (patch)
treeaa111f5adc1eaa1e996847e1437d1b1b40821ce0 /src/libstd/comm/sync.rs
parent14c1a103bc3f78721df1dc860a75a477c8275e3a (diff)
Fallout from new thread API
Diffstat (limited to 'src/libstd/comm/sync.rs')
-rw-r--r--src/libstd/comm/sync.rs40
1 files changed, 19 insertions, 21 deletions
diff --git a/src/libstd/comm/sync.rs b/src/libstd/comm/sync.rs
index 9e4bdb15b00..b24c6d21fba 100644
--- a/src/libstd/comm/sync.rs
+++ b/src/libstd/comm/sync.rs
@@ -38,10 +38,8 @@ 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 sync::{atomic, Mutex, MutexGuard};
 use comm::blocking::{mod, WaitToken, SignalToken};
@@ -105,10 +103,10 @@ pub enum Failure {
 
 /// Atomically blocks the current thread, placing it into `slot`, unlocking `lock`
 /// in the meantime. This re-locks the mutex upon returning.
-fn wait<'a, 'b, T>(lock: &'a Mutex<State<T>>,
-                   guard: MutexGuard<'b, State<T>>,
-                   f: fn(BlockedTask) -> Blocker)
-                   -> MutexGuard<'a, State<T>>
+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 me: Box<Task> = Local::take();
     me.deschedule(1, |task| {
@@ -170,7 +168,7 @@ impl<T: Send> Packet<T> {
     }
 
     pub fn send(&self, t: T) -> Result<(), T> {
-        let guard = self.acquire_send_slot();
+        let mut guard = self.acquire_send_slot();
         if guard.disconnected { return Err(t) }
         guard.buf.enqueue(t);
 
@@ -183,7 +181,7 @@ impl<T: Send> Packet<T> {
                 let mut canceled = false;
                 assert!(guard.canceled.is_none());
                 guard.canceled = Some(unsafe { mem::transmute(&mut canceled) });
-                let guard = wait(&self.lock, guard, BlockedSender);
+                let mut guard = wait(&self.lock, guard, BlockedSender);
                 if canceled {Err(guard.buf.dequeue())} else {Ok(())}
             }
 
@@ -198,7 +196,7 @@ impl<T: Send> Packet<T> {
     }
 
     pub fn try_send(&self, t: T) -> Result<(), super::TrySendError<T>> {
-        let guard = self.lock.lock();
+        let mut guard = self.lock.lock();
         if guard.disconnected {
             Err(super::RecvDisconnected(t))
         } else if guard.buf.size() == guard.buf.cap() {
@@ -235,13 +233,13 @@ 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 = self.lock.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 !guard.disconnected && guard.buf.size() == 0 {
-            wait(&mut guard.blocker, BlockedReceiver, &self.lock);
+            guard = wait(&self.lock, guard, BlockedReceiver);
             waited = true;
         }
         if guard.disconnected && guard.buf.size() == 0 { return Err(()) }
@@ -249,12 +247,12 @@ impl<T: Send> Packet<T> {
         // Pick up the data, wake up our neighbors, and carry on
         assert!(guard.buf.size() > 0);
         let ret = guard.buf.dequeue();
-        self.wakeup_senders(waited, guard, state);
+        self.wakeup_senders(waited, guard);
         return Ok(ret);
     }
 
     pub fn try_recv(&self) -> Result<T, Failure> {
-        let guard = self.lock();
+        let mut guard = self.lock.lock();
 
         // Easy cases first
         if guard.disconnected { return Err(Disconnected) }
@@ -262,7 +260,7 @@ impl<T: Send> Packet<T> {
 
         // Be sure to wake up neighbors
         let ret = Ok(guard.buf.dequeue());
-        self.wakeup_senders(false, guard, state);
+        self.wakeup_senders(false, guard);
 
         return ret;
     }
@@ -272,7 +270,7 @@ 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: MutexGuard<State<T>>) {
+    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
@@ -311,7 +309,7 @@ impl<T: Send> Packet<T> {
         }
 
         // Not much to do other than wake up a receiver if one's there
-        let guard = self.lock();
+        let mut guard = self.lock.lock();
         if guard.disconnected { return }
         guard.disconnected = true;
         match mem::replace(&mut guard.blocker, NoneBlocked) {
@@ -322,7 +320,7 @@ impl<T: Send> Packet<T> {
     }
 
     pub fn drop_port(&self) {
-        let guard = self.lock();
+        let mut guard = self.lock.lock();
 
         if guard.disconnected { return }
         guard.disconnected = true;
@@ -368,14 +366,14 @@ 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 guard = self.lock();
+        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, token: SignalToken) -> StartResult {
-        let guard = self.lock();
+        let mut guard = self.lock.lock();
         if guard.disconnected || guard.buf.size() > 0 {
             Abort
         } else {
@@ -393,7 +391,7 @@ impl<T: Send> Packet<T> {
     //
     // The return value indicates whether there's data on this port.
     pub fn abort_selection(&self) -> bool {
-        let guard = self.lock();
+        let mut guard = self.lock.lock();
         match mem::replace(&mut guard.blocker, NoneBlocked) {
             NoneBlocked => true,
             BlockedSender(token) => {
@@ -409,7 +407,7 @@ 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 guard = self.lock();
+        let mut guard = self.lock.lock();
         assert!(guard.queue.dequeue().is_none());
         assert!(guard.canceled.is_none());
     }