about summary refs log tree commit diff
path: root/src/libstd/rt/comm.rs
diff options
context:
space:
mode:
authorBen Blum <bblum@andrew.cmu.edu>2013-07-18 20:34:42 -0400
committerBen Blum <bblum@andrew.cmu.edu>2013-07-30 13:19:25 -0400
commitf34fadd126ce9ffcd5a79c3ad5d06ad58c6fd8a8 (patch)
treee8c622fc282ac880f718375d235feb7090424bff /src/libstd/rt/comm.rs
parent7326bc879ed6d9def8a128a2586a6f1708b5010c (diff)
Implement select() for new runtime pipes.
Diffstat (limited to 'src/libstd/rt/comm.rs')
-rw-r--r--src/libstd/rt/comm.rs166
1 files changed, 134 insertions, 32 deletions
diff --git a/src/libstd/rt/comm.rs b/src/libstd/rt/comm.rs
index b1533237b15..87bf5a23b93 100644
--- a/src/libstd/rt/comm.rs
+++ b/src/libstd/rt/comm.rs
@@ -12,13 +12,13 @@
 
 use option::*;
 use cast;
-use util;
 use ops::Drop;
 use rt::kill::BlockedTask;
 use kinds::Send;
 use rt::sched::Scheduler;
 use rt::local::Local;
-use unstable::atomics::{AtomicUint, AtomicOption, Acquire, SeqCst};
+use rt::select::{Select, SelectPort};
+use unstable::atomics::{AtomicUint, AtomicOption, Acquire, Release, SeqCst};
 use unstable::sync::UnsafeAtomicRcBox;
 use util::Void;
 use comm::{GenericChan, GenericSmartChan, GenericPort, Peekable};
@@ -76,6 +76,7 @@ pub fn oneshot<T: Send>() -> (PortOne<T>, ChanOne<T>) {
 }
 
 impl<T> ChanOne<T> {
+    #[inline]
     fn packet(&self) -> *mut Packet<T> {
         unsafe {
             let p: *mut ~Packet<T> = cast::transmute(&self.void_packet);
@@ -141,7 +142,6 @@ impl<T> ChanOne<T> {
     }
 }
 
-
 impl<T> PortOne<T> {
     fn packet(&self) -> *mut Packet<T> {
         unsafe {
@@ -162,46 +162,115 @@ impl<T> PortOne<T> {
 
     pub fn try_recv(self) -> Option<T> {
         let mut this = self;
-        let packet = this.packet();
 
         // Optimistic check. If data was sent already, we don't even need to block.
         // No release barrier needed here; we're not handing off our task pointer yet.
-        if unsafe { (*packet).state.load(Acquire) } != STATE_ONE {
+        if !this.optimistic_check() {
             // No data available yet.
             // Switch to the scheduler to put the ~Task into the Packet state.
             let sched = Local::take::<Scheduler>();
             do sched.deschedule_running_task_and_then |sched, task| {
-                unsafe {
-                    // Atomically swap the task pointer into the Packet state, issuing
-                    // an acquire barrier to prevent reordering of the subsequent read
-                    // of the payload. Also issues a release barrier to prevent
-                    // reordering of any previous writes to the task structure.
-                    let task_as_state = task.cast_to_uint();
-                    let oldstate = (*packet).state.swap(task_as_state, SeqCst);
-                    match oldstate {
-                        STATE_BOTH => {
-                            // Data has not been sent. Now we're blocked.
-                            rtdebug!("non-rendezvous recv");
-                            sched.metrics.non_rendezvous_recvs += 1;
-                        }
-                        STATE_ONE => {
-                            rtdebug!("rendezvous recv");
-                            sched.metrics.rendezvous_recvs += 1;
-
-                            // Channel is closed. Switch back and check the data.
-                            // NB: We have to drop back into the scheduler event loop here
-                            // instead of switching immediately back or we could end up
-                            // triggering infinite recursion on the scheduler's stack.
-                            let recvr = BlockedTask::cast_from_uint(task_as_state);
-                            sched.enqueue_blocked_task(recvr);
+                this.block_on(sched, task);
+            }
+        }
+
+        // Task resumes.
+        this.recv_ready()
+    }
+}
+
+impl<T> Select for PortOne<T> {
+    #[inline]
+    fn optimistic_check(&mut self) -> bool {
+        unsafe { (*self.packet()).state.load(Acquire) == STATE_ONE }
+    }
+
+    fn block_on(&mut self, sched: &mut Scheduler, task: BlockedTask) -> bool {
+        unsafe {
+            // Atomically swap the task pointer into the Packet state, issuing
+            // an acquire barrier to prevent reordering of the subsequent read
+            // of the payload. Also issues a release barrier to prevent
+            // reordering of any previous writes to the task structure.
+            let task_as_state = task.cast_to_uint();
+            let oldstate = (*self.packet()).state.swap(task_as_state, SeqCst);
+            match oldstate {
+                STATE_BOTH => {
+                    // Data has not been sent. Now we're blocked.
+                    rtdebug!("non-rendezvous recv");
+                    sched.metrics.non_rendezvous_recvs += 1;
+                    false
+                }
+                STATE_ONE => {
+                    // Re-record that we are the only owner of the packet.
+                    // Release barrier needed in case the task gets reawoken
+                    // on a different core (this is analogous to writing a
+                    // payload; a barrier in enqueueing the task protects it).
+                    // NB(#8132). This *must* occur before the enqueue below.
+                    // FIXME(#6842, #8130) This is usually only needed for the
+                    // assertion in recv_ready, except in the case of select().
+                    // This won't actually ever have cacheline contention, but
+                    // maybe should be optimized out with a cfg(test) anyway?
+                    (*self.packet()).state.store(STATE_ONE, Release);
+
+                    rtdebug!("rendezvous recv");
+                    sched.metrics.rendezvous_recvs += 1;
+
+                    // Channel is closed. Switch back and check the data.
+                    // NB: We have to drop back into the scheduler event loop here
+                    // instead of switching immediately back or we could end up
+                    // triggering infinite recursion on the scheduler's stack.
+                    let recvr = BlockedTask::cast_from_uint(task_as_state);
+                    sched.enqueue_blocked_task(recvr);
+                    true
+                }
+                _ => rtabort!("can't block_on; a task is already blocked")
+            }
+        }
+    }
+
+    // This is the only select trait function that's not also used in recv.
+    fn unblock_from(&mut self) -> bool {
+        let packet = self.packet();
+        unsafe {
+            // In case the data is available, the acquire barrier here matches
+            // the release barrier the sender used to release the payload.
+            match (*packet).state.load(Acquire) {
+                // Impossible. We removed STATE_BOTH when blocking on it, and
+                // no self-respecting sender would put it back.
+                STATE_BOTH    => rtabort!("refcount already 2 in unblock_from"),
+                // Here, a sender already tried to wake us up. Perhaps they
+                // even succeeded! Data is available.
+                STATE_ONE     => true,
+                // Still registered as blocked. Need to "unblock" the pointer.
+                task_as_state => {
+                    // In the window between the load and the CAS, a sender
+                    // might take the pointer and set the refcount to ONE. If
+                    // that happens, we shouldn't clobber that with BOTH!
+                    // Acquire barrier again for the same reason as above.
+                    match (*packet).state.compare_and_swap(task_as_state, STATE_BOTH,
+                                                           Acquire) {
+                        STATE_BOTH => rtabort!("refcount became 2 in unblock_from"),
+                        STATE_ONE  => true, // Lost the race. Data available.
+                        same_ptr   => {
+                            // We successfully unblocked our task pointer.
+                            assert!(task_as_state == same_ptr);
+                            let handle = BlockedTask::cast_from_uint(task_as_state);
+                            // Because we are already awake, the handle we
+                            // gave to this port shall already be empty.
+                            handle.assert_already_awake();
+                            false
                         }
-                        _ => util::unreachable()
                     }
                 }
             }
         }
+    }
+}
 
-        // Task resumes.
+impl<T> SelectPort<T> for PortOne<T> {
+    fn recv_ready(self) -> Option<T> {
+        let mut this = self;
+        let packet = this.packet();
 
         // No further memory barrier is needed here to access the
         // payload. Some scenarios:
@@ -213,8 +282,11 @@ impl<T> PortOne<T> {
         // 3) We encountered STATE_BOTH above and blocked, but the receiving task (this task)
         //    is pinned to some other scheduler, so the sending task had to give us to
         //    a different scheduler for resuming. That send synchronized memory.
-
         unsafe {
+            // See corresponding store() above in block_on for rationale.
+            // FIXME(#8130) This can happen only in test builds.
+            assert!((*packet).state.load(Acquire) == STATE_ONE);
+
             let payload = (*packet).payload.take();
 
             // The sender has closed up shop. Drop the packet.
@@ -234,7 +306,7 @@ impl<T> Peekable<T> for PortOne<T> {
             match oldstate {
                 STATE_BOTH => false,
                 STATE_ONE => (*packet).payload.is_some(),
-                _ => util::unreachable()
+                _ => rtabort!("peeked on a blocked task")
             }
         }
     }
@@ -368,6 +440,36 @@ impl<T> Peekable<T> for Port<T> {
     }
 }
 
+impl<T> Select for Port<T> {
+    #[inline]
+    fn optimistic_check(&mut self) -> bool {
+        do self.next.with_mut_ref |pone| { pone.optimistic_check() }
+    }
+
+    #[inline]
+    fn block_on(&mut self, sched: &mut Scheduler, task: BlockedTask) -> bool {
+        let task = Cell::new(task);
+        do self.next.with_mut_ref |pone| { pone.block_on(sched, task.take()) }
+    }
+
+    #[inline]
+    fn unblock_from(&mut self) -> bool {
+        do self.next.with_mut_ref |pone| { pone.unblock_from() }
+    }
+}
+
+impl<T> SelectPort<(T, Port<T>)> for Port<T> {
+    fn recv_ready(self) -> Option<(T, Port<T>)> {
+        match self.next.take().recv_ready() {
+            Some(StreamPayload { val, next }) => {
+                self.next.put_back(next);
+                Some((val, self))
+            }
+            None => None
+        }
+    }
+}
+
 pub struct SharedChan<T> {
     // Just like Chan, but a shared AtomicOption instead of Cell
     priv next: UnsafeAtomicRcBox<AtomicOption<StreamChanOne<T>>>