about summary refs log tree commit diff
path: root/src/libstd/sync/mpsc/sync.rs
diff options
context:
space:
mode:
authorbors <bors@rust-lang.org>2019-11-30 12:42:44 +0000
committerbors <bors@rust-lang.org>2019-11-30 12:42:44 +0000
commitd8bdb3fdcbd88eb16e1a6669236122c41ed2aed3 (patch)
tree33974ee0e3d5976f284b056e03e6ef529d15e563 /src/libstd/sync/mpsc/sync.rs
parent8f1bbd69e13c9e04a4c2b75612bc0c31af972439 (diff)
parentb14d9c21203ea79035bf4a8a8a68ad34658a265f (diff)
downloadrust-d8bdb3fdcbd88eb16e1a6669236122c41ed2aed3.tar.gz
rust-d8bdb3fdcbd88eb16e1a6669236122c41ed2aed3.zip
Auto merge of #66887 - dtolnay:rollup-uxowp8d, r=Centril
Rollup of 4 pull requests

Successful merges:

 - #66818 (Format libstd/os with rustfmt)
 - #66819 (Format libstd/sys with rustfmt)
 - #66820 (Format libstd with rustfmt)
 - #66847 (Allow any identifier as format arg name)

Failed merges:

r? @ghost
Diffstat (limited to 'src/libstd/sync/mpsc/sync.rs')
-rw-r--r--src/libstd/sync/mpsc/sync.rs121
1 files changed, 66 insertions, 55 deletions
diff --git a/src/libstd/sync/mpsc/sync.rs b/src/libstd/sync/mpsc/sync.rs
index 58a4b716afb..79e86817154 100644
--- a/src/libstd/sync/mpsc/sync.rs
+++ b/src/libstd/sync/mpsc/sync.rs
@@ -1,3 +1,4 @@
+use self::Blocker::*;
 /// Synchronous channels/ports
 ///
 /// This channel implementation differs significantly from the asynchronous
@@ -22,17 +23,15 @@
 /// implementation shares almost all code for the buffered and unbuffered cases
 /// of a synchronous channel. There are a few branches for the unbuffered case,
 /// but they're mostly just relevant to blocking senders.
-
 pub use self::Failure::*;
-use self::Blocker::*;
 
 use core::intrinsics::abort;
 use core::isize;
 use core::mem;
 use core::ptr;
 
-use crate::sync::atomic::{Ordering, AtomicUsize};
-use crate::sync::mpsc::blocking::{self, WaitToken, SignalToken};
+use crate::sync::atomic::{AtomicUsize, Ordering};
+use crate::sync::mpsc::blocking::{self, SignalToken, WaitToken};
 use crate::sync::{Mutex, MutexGuard};
 use crate::time::Instant;
 
@@ -46,9 +45,9 @@ pub struct Packet<T> {
     lock: Mutex<State<T>>,
 }
 
-unsafe impl<T: Send> Send for Packet<T> { }
+unsafe impl<T: Send> Send for Packet<T> {}
 
-unsafe impl<T: Send> Sync for Packet<T> { }
+unsafe impl<T: Send> Sync for Packet<T> {}
 
 struct State<T> {
     disconnected: bool, // Is the channel disconnected yet?
@@ -72,7 +71,7 @@ unsafe impl<T: Send> Send for State<T> {}
 enum Blocker {
     BlockedSender(SignalToken),
     BlockedReceiver(SignalToken),
-    NoneBlocked
+    NoneBlocked,
 }
 
 /// Simple queue for threading threads together. Nodes are stack-allocated, so
@@ -104,35 +103,35 @@ 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>>,
-                   mut guard: MutexGuard<'b, State<T>>,
-                   f: fn(SignalToken) -> Blocker)
-                   -> MutexGuard<'a, State<T>>
-{
+fn wait<'a, 'b, T>(
+    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
+    drop(guard); // unlock
+    wait_token.wait(); // block
     lock.lock().unwrap() // relock
 }
 
 /// Same as wait, but waiting at most until `deadline`.
-fn wait_timeout_receiver<'a, 'b, T>(lock: &'a Mutex<State<T>>,
-                                    deadline: Instant,
-                                    mut guard: MutexGuard<'b, State<T>>,
-                                    success: &mut bool)
-                                    -> MutexGuard<'a, State<T>>
-{
+fn wait_timeout_receiver<'a, 'b, T>(
+    lock: &'a Mutex<State<T>>,
+    deadline: Instant,
+    mut guard: MutexGuard<'b, State<T>>,
+    success: &mut bool,
+) -> MutexGuard<'a, State<T>> {
     let (wait_token, signal_token) = blocking::tokens();
     match mem::replace(&mut guard.blocker, BlockedReceiver(signal_token)) {
         NoneBlocked => {}
         _ => unreachable!(),
     }
-    drop(guard);         // unlock
-    *success = wait_token.wait_max_until(deadline);   // block
+    drop(guard); // unlock
+    *success = wait_token.wait_max_until(deadline); // block
     let mut new_guard = lock.lock().unwrap(); // relock
     if !*success {
         abort_selection(&mut new_guard);
@@ -147,7 +146,10 @@ fn abort_selection<T>(guard: &mut MutexGuard<'_, State<T>>) -> bool {
             guard.blocker = BlockedSender(token);
             true
         }
-        BlockedReceiver(token) => { drop(token); false }
+        BlockedReceiver(token) => {
+            drop(token);
+            false
+        }
     }
 }
 
@@ -168,12 +170,9 @@ impl<T> Packet<T> {
                 blocker: NoneBlocked,
                 cap: capacity,
                 canceled: None,
-                queue: Queue {
-                    head: ptr::null_mut(),
-                    tail: ptr::null_mut(),
-                },
+                queue: Queue { head: ptr::null_mut(), tail: ptr::null_mut() },
                 buf: Buffer {
-                    buf: (0..capacity + if capacity == 0 {1} else {0}).map(|_| None).collect(),
+                    buf: (0..capacity + if capacity == 0 { 1 } else { 0 }).map(|_| None).collect(),
                     start: 0,
                     size: 0,
                 },
@@ -200,7 +199,9 @@ impl<T> Packet<T> {
 
     pub fn send(&self, t: T) -> Result<(), T> {
         let mut guard = self.acquire_send_slot();
-        if guard.disconnected { return Err(t) }
+        if guard.disconnected {
+            return Err(t);
+        }
         guard.buf.enqueue(t);
 
         match mem::replace(&mut guard.blocker, NoneBlocked) {
@@ -213,14 +214,17 @@ impl<T> Packet<T> {
                 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(())}
+                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(token) => { wakeup(token, guard); Ok(()) }
+            BlockedReceiver(token) => {
+                wakeup(token, guard);
+                Ok(())
+            }
 
             BlockedSender(..) => panic!("lolwut"),
         }
@@ -271,10 +275,8 @@ impl<T> Packet<T> {
         // while loop because we're the only receiver.
         if !guard.disconnected && guard.buf.size() == 0 {
             if let Some(deadline) = deadline {
-                guard = wait_timeout_receiver(&self.lock,
-                                              deadline,
-                                              guard,
-                                              &mut woke_up_after_waiting);
+                guard =
+                    wait_timeout_receiver(&self.lock, deadline, guard, &mut woke_up_after_waiting);
             } else {
                 guard = wait(&self.lock, guard, BlockedReceiver);
                 woke_up_after_waiting = true;
@@ -290,7 +292,9 @@ impl<T> Packet<T> {
         // Pick up the data, wake up our neighbors, and carry on
         assert!(guard.buf.size() > 0 || (deadline.is_some() && !woke_up_after_waiting));
 
-        if guard.buf.size() == 0 { return Err(Empty); }
+        if guard.buf.size() == 0 {
+            return Err(Empty);
+        }
 
         let ret = guard.buf.dequeue();
         self.wakeup_senders(woke_up_after_waiting, guard);
@@ -301,8 +305,12 @@ impl<T> Packet<T> {
         let mut guard = self.lock.lock().unwrap();
 
         // Easy cases first
-        if guard.disconnected && guard.buf.size() == 0 { return Err(Disconnected) }
-        if guard.buf.size() == 0 { return Err(Empty) }
+        if guard.disconnected && guard.buf.size() == 0 {
+            return Err(Disconnected);
+        }
+        if guard.buf.size() == 0 {
+            return Err(Empty);
+        }
 
         // Be sure to wake up neighbors
         let ret = Ok(guard.buf.dequeue());
@@ -357,12 +365,14 @@ impl<T> Packet<T> {
         // Only flag the channel as disconnected if we're the last channel
         match self.channels.fetch_sub(1, Ordering::SeqCst) {
             1 => {}
-            _ => return
+            _ => return,
         }
 
         // Not much to do other than wake up a receiver if one's there
         let mut guard = self.lock.lock().unwrap();
-        if guard.disconnected { return }
+        if guard.disconnected {
+            return;
+        }
         guard.disconnected = true;
         match mem::replace(&mut guard.blocker, NoneBlocked) {
             NoneBlocked => {}
@@ -374,7 +384,9 @@ impl<T> Packet<T> {
     pub fn drop_port(&self) {
         let mut guard = self.lock.lock().unwrap();
 
-        if guard.disconnected { return }
+        if guard.disconnected {
+            return;
+        }
         guard.disconnected = true;
 
         // If the capacity is 0, then the sender may want its data back after
@@ -382,15 +394,9 @@ impl<T> Packet<T> {
         // 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 guard.cap != 0 {
-            mem::take(&mut guard.buf.buf)
-        } else {
-            Vec::new()
-        };
-        let mut queue = mem::replace(&mut guard.queue, Queue {
-            head: ptr::null_mut(),
-            tail: ptr::null_mut(),
-        });
+        let _data = if guard.cap != 0 { mem::take(&mut guard.buf.buf) } else { Vec::new() };
+        let mut queue =
+            mem::replace(&mut guard.queue, Queue { head: ptr::null_mut(), tail: ptr::null_mut() });
 
         let waiter = match mem::replace(&mut guard.blocker, NoneBlocked) {
             NoneBlocked => None,
@@ -402,7 +408,9 @@ impl<T> Packet<T> {
         };
         mem::drop(guard);
 
-        while let Some(token) = queue.dequeue() { token.signal(); }
+        while let Some(token) = queue.dequeue() {
+            token.signal();
+        }
         waiter.map(|t| t.signal());
     }
 }
@@ -416,7 +424,6 @@ impl<T> Drop for Packet<T> {
     }
 }
 
-
 ////////////////////////////////////////////////////////////////////////////////
 // Buffer, a simple ring buffer backed by Vec<T>
 ////////////////////////////////////////////////////////////////////////////////
@@ -437,8 +444,12 @@ impl<T> Buffer<T> {
         result.take().unwrap()
     }
 
-    fn size(&self) -> usize { self.size }
-    fn capacity(&self) -> usize { self.buf.len() }
+    fn size(&self) -> usize {
+        self.size
+    }
+    fn capacity(&self) -> usize {
+        self.buf.len()
+    }
 }
 
 ////////////////////////////////////////////////////////////////////////////////
@@ -466,7 +477,7 @@ impl Queue {
 
     fn dequeue(&mut self) -> Option<SignalToken> {
         if self.head.is_null() {
-            return None
+            return None;
         }
         let node = self.head;
         self.head = unsafe { (*node).next };