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
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
//! Type-erased async task and its wake/poll/complete re-entry protocol.
//!
//! `AsyncTask` is one spawned future as the executor sees it: the future itself
//! type-erased into `ErasedTaskFuture`, plus the flags that decide when it may
//! be polled. Tasks are shared as `Arc<AsyncTask>` between the run queue, the
//! wakers handed to the future, and any thread running
//! `AsyncExecutor::process_pending_tasks`, so both the erasure and the flags
//! carry cross-thread contracts.
//!
//! # Wakers
//!
//! The task is its own waker: `impl Wake for AsyncTask` (`waker.rs`) makes every
//! `Waker::from(Arc<AsyncTask>)` a clone of the task's `Arc`, so a poll mints
//! its waker without allocating and every waker of one task has the same data
//! pointer and vtable, which is what `Waker::will_wake` compares. A waker
//! strongly owns its task, so a pending task lives exactly as long as its run
//! queue entry or a registered waker does. The task in turn holds only `Weak`
//! handles to the run queue and the reactor it re-enqueues into: a strong
//! handle would close the cycle run queue -> task -> run queue whenever the
//! executor is dropped with a task still queued, and a stored waker would
//! close task -> waker -> task. A wake after the executor is gone finds the
//! handles dead and does nothing.
//!
//! # Type erasure
//!
//! `ErasedTaskFuture` is a hand-built vtable — a `NonNull<()>` to the boxed
//! future plus `poll`/`drop` function pointers monomorphized for the concrete
//! `F` at construction. This keeps a heterogeneous run queue without boxing
//! every task behind `dyn Future` at each call site. Three invariants make it
//! sound:
//!
//! 1. **Address stability.** `new` heap-allocates the future with `Box` and
//! never moves it again, so the `Pin::new_unchecked` in
//! `poll_erased_future` is honest: the future is pinned for the whole
//! lifetime of the `ErasedTaskFuture` that owns it.
//! 2. **Type agreement.** `ptr`, `poll`, and `drop` are set together from the
//! same `F` and never reassigned, so the `cast::<F>()` inside each function
//! pointer always recovers the type that was actually stored.
//! 3. **Single drop.** The `ErasedTaskFuture` uniquely owns the allocation, and
//! only its `Drop` frees it (via `Box::from_raw` through the stored `drop`
//! pointer). `Send` is justified by the `F: Send` bound at construction —
//! the pointer moves between threads only with the task it belongs to.
//!
//! # Re-entry protocol (`is_queued` / `completed` / `future_lock`)
//!
//! A future must be polled by one thread at a time and must never be polled
//! after it returns `Ready` — a completed `async` block panics with "resumed
//! after completion". Three fields enforce that, and they divide the work:
//!
//! - `is_queued` — enqueue deduplication, owned by `ExecutorWaker::wake_by_ref`
//! (`waker.rs`). A wake enqueues only if it flipped `is_queued` false→true, so
//! a task appears in the run queue at most once per pending wake. The flag is
//! a Relaxed linearization bit rather than a publication channel: the queue's
//! per-slot Release/Acquire sequence publishes the task payload, and the
//! ordering protocol is exhaustively modeled in `loom_wake_dedup.rs`.
//! - `completed` — set once the future returns `Ready`. The waker checks it as
//! an optimization (a reactor may still hold a live waker for a task that
//! finished by another path, e.g. `timeout(read)` completing via the timer
//! while a socket read-waker stays registered), but the authoritative guard is
//! in `process_pending_tasks`.
//! - `future_lock` — serializes access to the `UnsafeCell<ErasedTaskFuture>`.
//! Its guard is held across the `completed` check *and* the poll, which is
//! what makes the guard sound: `process_pending_tasks` clears `is_queued`
//! before polling (so a self-wake during the poll can re-enqueue), so a second
//! polling thread can dequeue the same task while the first still holds the
//! lock. Checking `completed` outside the lock would let that thread pass the
//! check, block, and then poll a future the first thread completed meanwhile.
//!
//! Together these give the `unsafe impl Sync` below its meaning: every mutable
//! touch of the future happens under `future_lock`, and every poll is gated by a
//! `completed` check taken in the same critical section.
use ;
use IoReactor;
use LockFreeQueue;
use Future;
use Pin;
use NonNull;
use AtomicBool;
use ;
use ;
use Instant;
pub
// SAFETY: the only non-Sync field is the `UnsafeCell<ErasedTaskFuture>`;
// mutable access to it is serialized by `future_lock`, whose guard is held
// across both the `completed` check and the poll it guards (see the re-entry
// protocol in the module docs), so no two threads touch the future
// concurrently and none polls it after completion.
unsafe
pub
// SAFETY: the erased pointer owns a `Box<F>` built under an `F: Send` bound in
// `new`, so moving the task (and with it this pointer) across threads moves a
// `Send` value; the vtable entries are plain fn pointers.
unsafe
unsafe
unsafe