Fix race condition in CurrentThread. (#156)

The logic that enables `CurrentThread::turn` to avoid unbounded
iteration was incorrect. It was possible for unfortunate timing to
result in a dead lock.

This patch provides a fix as well as a test.
This commit is contained in:
Carl Lerche
2018-02-26 20:41:06 -08:00
committed by GitHub
parent df3a92532b
commit 2961a2388c
3 changed files with 67 additions and 5 deletions
+14 -3
View File
@@ -202,7 +202,8 @@ where U: Unpark,
where F: FnMut(&mut Self, &mut Scheduled<U>),
{
let mut ret = false;
let tick = self.inner.tick_num.fetch_add(1, SeqCst);
let tick = self.inner.tick_num.fetch_add(1, SeqCst)
.wrapping_add(1);
loop {
let node = match unsafe { self.inner.dequeue(Some(tick)) } {
@@ -439,8 +440,18 @@ impl<U> Inner<U> {
}
if let Some(tick) = tick {
// Only dequeue if the node matches the tick num
if (*tail).notified_at.load(SeqCst) != tick {
let actual = (*tail).notified_at.load(SeqCst);
// Only dequeue if the node was not scheduled during the current
// tick.
if actual == tick {
// Only doing the check above **should** be enough in
// practice. However, technically there is a potential for
// deadlocking if there are `usize::MAX` ticks while the thread
// scheduling the task is frozen.
//
// If, for some reason, this is not enough, calling `unpark`
// here will resolve the issue.
return Dequeue::Empty;
}
}
+6 -1
View File
@@ -256,7 +256,10 @@ impl Reactor {
}
/// Returns true if the reactor is currently idle.
pub(crate) fn is_idle(&self) -> bool {
///
/// Idle is defined as all tasks that have been spawned have completed,
/// either successfully or with an error.
pub fn is_idle(&self) -> bool {
self.inner.io_dispatch
.read().unwrap()
.is_empty()
@@ -313,9 +316,11 @@ impl Reactor {
if let Some(io) = io_dispatch.get(token) {
io.readiness.fetch_or(ready2usize(ready), Relaxed);
if ready.is_writable() {
io.writer.notify();
}
if !(ready & (!mio::Ready::writable())).is_empty() {
io.reader.notify();
}
+47 -1
View File
@@ -258,7 +258,7 @@ fn tasks_are_scheduled_fairly() {
}
#[test]
fn spawn_and_tick() {
fn spawn_and_turn() {
let cnt = Rc::new(Cell::new(0));
let c = cnt.clone();
@@ -293,6 +293,52 @@ fn spawn_and_tick() {
assert_eq!(2, cnt.get());
}
#[test]
fn hammer_turn() {
use futures::sync::mpsc;
const ITER: usize = 100;
const N: usize = 100;
const THREADS: usize = 4;
for _ in 0..ITER {
let mut ths = vec![];
// Add some jitter
for _ in 0..THREADS {
let th = thread::spawn(|| {
let mut current_thread = CurrentThread::new();
let (tx, rx) = mpsc::unbounded();
current_thread.spawn({
rx.for_each(|_| {
Ok(())
})
.map_err(|e| panic!("err={:?}", e))
});
thread::spawn(move || {
for _ in 0..N {
tx.unbounded_send(()).unwrap();
thread::yield_now();
}
});
while !current_thread.is_idle() {
current_thread.turn(None).unwrap();
}
});
ths.push(th);
}
for th in ths {
th.join().unwrap();
}
}
}
fn ok() -> future::FutureResult<(), ()> {
future::ok(())
}