diff options
| author | Ben Blum <bblum@andrew.cmu.edu> | 2013-07-18 20:34:42 -0400 |
|---|---|---|
| committer | Ben Blum <bblum@andrew.cmu.edu> | 2013-07-30 13:19:25 -0400 |
| commit | f34fadd126ce9ffcd5a79c3ad5d06ad58c6fd8a8 (patch) | |
| tree | e8c622fc282ac880f718375d235feb7090424bff /src/libstd/rt/comm.rs | |
| parent | 7326bc879ed6d9def8a128a2586a6f1708b5010c (diff) | |
Implement select() for new runtime pipes.
Diffstat (limited to 'src/libstd/rt/comm.rs')
| -rw-r--r-- | src/libstd/rt/comm.rs | 166 |
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>>> |
