mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-08-07 00:00:09 +02:00
Fix a race in thread wakeup (#459)
This commit is contained in:
committed by
Carl Lerche
parent
dbefa67058
commit
7fb579c667
@@ -41,7 +41,6 @@ use futures2;
|
||||
/// use std::time::Duration;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .pool_size(4)
|
||||
/// .keep_alive(Some(Duration::from_secs(30)))
|
||||
@@ -86,7 +85,6 @@ impl Builder {
|
||||
/// use std::time::Duration;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .pool_size(4)
|
||||
/// .keep_alive(Some(Duration::from_secs(30)))
|
||||
@@ -131,7 +129,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .pool_size(4)
|
||||
/// .build();
|
||||
@@ -164,7 +161,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .max_blocking(200)
|
||||
/// .build();
|
||||
@@ -196,7 +192,6 @@ impl Builder {
|
||||
/// use std::time::Duration;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .keep_alive(Some(Duration::from_secs(30)))
|
||||
/// .build();
|
||||
@@ -224,7 +219,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .name_prefix("my-pool-")
|
||||
/// .build();
|
||||
@@ -251,7 +245,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .stack_size(32 * 1024)
|
||||
/// .build();
|
||||
@@ -265,7 +258,7 @@ impl Builder {
|
||||
/// Execute function `f` on each worker thread.
|
||||
///
|
||||
/// This function is provided a handle to the worker and is expected to call
|
||||
/// `Worker::run`, otherwise the worker thread will shutdown without doing
|
||||
/// [`Worker::run`], otherwise the worker thread will shutdown without doing
|
||||
/// any work.
|
||||
///
|
||||
/// # Examples
|
||||
@@ -276,7 +269,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .around_worker(|worker, _| {
|
||||
/// println!("worker is starting up");
|
||||
@@ -286,6 +278,8 @@ impl Builder {
|
||||
/// .build();
|
||||
/// # }
|
||||
/// ```
|
||||
///
|
||||
/// [`Worker::run`]: struct.Worker.html#method.run
|
||||
pub fn around_worker<F>(&mut self, f: F) -> &mut Self
|
||||
where F: Fn(&Worker, &mut Enter) + Send + Sync + 'static
|
||||
{
|
||||
@@ -306,7 +300,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .after_start(|| {
|
||||
/// println!("thread started");
|
||||
@@ -333,7 +326,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .before_stop(|| {
|
||||
/// println!("thread stopping");
|
||||
@@ -362,7 +354,6 @@ impl Builder {
|
||||
/// # fn decorate<F>(f: F) -> F { f }
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .custom_park(|_| {
|
||||
/// use tokio_threadpool::park::DefaultPark;
|
||||
@@ -402,7 +393,6 @@ impl Builder {
|
||||
/// # use tokio_threadpool::Builder;
|
||||
///
|
||||
/// # pub fn main() {
|
||||
/// // Create a thread pool with default configuration values
|
||||
/// let thread_pool = Builder::new()
|
||||
/// .build();
|
||||
/// # }
|
||||
|
||||
@@ -133,8 +133,8 @@ impl Inner {
|
||||
None => self.condvar.wait(m).unwrap(),
|
||||
};
|
||||
|
||||
// Transition back to idle. If the state has transitions dto `NOTIFY`,
|
||||
// this will consume that notification
|
||||
// Transition back to idle. If the state has transitioned to `NOTIFY`,
|
||||
// this will consume that notification.
|
||||
self.state.store(IDLE, SeqCst);
|
||||
|
||||
// Explicitly drop the mutex guard. There is no real point in doing it
|
||||
@@ -155,10 +155,12 @@ impl Inner {
|
||||
// The other half is sleeping, this requires a lock
|
||||
let _m = self.mutex.lock().unwrap();
|
||||
|
||||
// Transition from SLEEP -> NOTIFY
|
||||
match self.state.compare_and_swap(SLEEP, NOTIFY, SeqCst) {
|
||||
// Transition to NOTIFY
|
||||
match self.state.swap(NOTIFY, SeqCst) {
|
||||
SLEEP => {}
|
||||
_ => return,
|
||||
NOTIFY => return,
|
||||
IDLE => return,
|
||||
_ => unreachable!(),
|
||||
}
|
||||
|
||||
// Wakeup the sleeper
|
||||
|
||||
@@ -29,8 +29,12 @@ use std::time::{Duration, Instant};
|
||||
|
||||
/// Thread worker
|
||||
///
|
||||
/// This is passed to the `around_worker` callback set on `Builder`. This
|
||||
/// callback is only expected to call `run` on it.
|
||||
/// This is passed to the [`around_worker`] callback set on [`Builder`]. This
|
||||
/// callback is only expected to call [`run`] on it.
|
||||
///
|
||||
/// [`Builder`]: struct.Builder.html
|
||||
/// [`around_worker`]: struct.Builder.html#method.around_worker
|
||||
/// [`run`]: struct.Worker.html#method.run
|
||||
#[derive(Debug)]
|
||||
pub struct Worker {
|
||||
// Shared scheduler data
|
||||
|
||||
Reference in New Issue
Block a user