1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
//! Wake protocol of [`AsyncTask`]: re-enqueue a pending task exactly once.
use crate::executor::task::AsyncTask;
use std::sync::Arc;
use std::sync::atomic::Ordering;
use std::task::Wake;
impl Wake for AsyncTask {
fn wake(self: Arc<Self>) {
self.wake_by_ref();
}
fn wake_by_ref(self: &Arc<Self>) {
// A completed task's waker may still be held live by the reactor (a
// read-waker registered against a socket fd, say) and fire after the
// task finished via another path. Re-enqueuing it would poll a future
// that already returned `Ready` and panic, so drop the wake for a
// completed task. `process_pending_tasks` re-checks `completed` to
// close the wake-races-completion window authoritatively.
if self.completed.load(Ordering::Acquire) {
return;
}
if !self
.is_queued
// `is_queued` is a linearization flag for enqueue deduplication,
// not a publication channel. The queue's per-slot Release/Acquire
// sequence publishes the task itself; this RMW only orders the
// false -> true transition against `process_pending_tasks`'
// corresponding clear. Relaxed is therefore sufficient and avoids
// a global ordering edge on every wake.
.swap(true, Ordering::Relaxed)
{
// The executor may be gone (a reactor or timer thread woke a task
// after its executor dropped); nothing will poll it then.
if let Some(run_queue) = self.run_queue.upgrade() {
run_queue.enqueue(Arc::clone(self));
if let Some(reactor) = self.reactor.upgrade() {
let _ = reactor.wake();
}
}
}
}
}
#[cfg(test)]
mod tests {
use crate::executor::AsyncExecutor;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::task::{Poll, Waker};
#[test]
fn every_poll_of_a_task_presents_a_will_wake_equal_waker() {
use std::sync::Mutex;
const POLLS: usize = 1_000;
// A task that re-wakes itself each poll and keeps every waker it was
// handed. The executor drives all polls from one `process_pending_tasks`
// call because each self-wake re-enqueues the task.
let executor = AsyncExecutor::new().expect("a fresh AsyncExecutor must build");
let seen = Arc::new(Mutex::new(Vec::<Waker>::with_capacity(POLLS)));
let future_seen = Arc::clone(&seen);
let _handle = executor.spawn(futures::future::poll_fn(move |context| {
let mut seen = future_seen.lock().expect("waker log must stay available");
seen.push(context.waker().clone());
if seen.len() == POLLS {
Poll::Ready(())
} else {
context.waker().wake_by_ref();
Poll::Pending
}
}));
executor.process_pending_tasks();
let seen = seen.lock().expect("waker log must stay available");
assert_eq!(seen.len(), POLLS);
assert_eq!(executor.stats().tasks_completed, 1);
assert!(
seen.iter().all(|waker| waker.will_wake(&seen[0])),
"a task must present one stable waker identity across polls"
);
}
/// Flags its drop so a test can observe when the task owning it is freed.
struct DropProbe(Arc<AtomicBool>);
impl Drop for DropProbe {
fn drop(&mut self) {
self.0.store(true, Ordering::SeqCst);
}
}
#[test]
fn dropping_the_executor_frees_a_task_still_in_the_run_queue() {
// The task sits in the run queue, which holds the task strongly. The
// task must hold the queue weakly, else queue -> task -> queue keeps
// both alive after the executor drops.
let dropped = Arc::new(AtomicBool::new(false));
let probe = DropProbe(Arc::clone(&dropped));
let executor = AsyncExecutor::new().expect("a fresh AsyncExecutor must build");
let _handle = executor.spawn(async move {
let _probe = probe;
futures::future::pending::<()>().await;
});
assert!(!dropped.load(Ordering::SeqCst));
drop(executor);
assert!(
dropped.load(Ordering::SeqCst),
"the queued task's future must drop with its executor"
);
}
#[test]
fn a_pending_task_lives_and_dies_with_its_last_waker() {
use std::sync::Mutex;
// After its first poll the task is not queued; only the waker the
// future stashed keeps it alive. Dropping that waker must free it,
// which fails if the task owned a waker of itself (task -> waker ->
// task).
let dropped = Arc::new(AtomicBool::new(false));
let probe = DropProbe(Arc::clone(&dropped));
let stash = Arc::new(Mutex::new(None::<Waker>));
let future_stash = Arc::clone(&stash);
let executor = AsyncExecutor::new().expect("a fresh AsyncExecutor must build");
let _handle = executor.spawn(futures::future::poll_fn(move |context| {
let _probe = &probe;
*future_stash
.lock()
.expect("waker stash must stay available") = Some(context.waker().clone());
Poll::<()>::Pending
}));
executor.process_pending_tasks();
assert!(
!dropped.load(Ordering::SeqCst),
"a stashed waker keeps the task alive"
);
let waker = stash
.lock()
.expect("waker stash must stay available")
.take()
.expect("the first poll must stash its waker");
drop(executor);
waker.wake_by_ref();
assert!(
!dropped.load(Ordering::SeqCst),
"a wake after the executor dropped must neither panic nor free the task"
);
drop(waker);
assert!(
dropped.load(Ordering::SeqCst),
"dropping the last waker must free the pending task"
);
}
}