RFA-036 · Case file with fixtures · Case 8 of 694 · Cargo workspace evidence
Why a Rust Stream Stops Forever After Returning Pending
Poll::Pending carries a wakeup obligation. Reproduce the check-register race, store the latest waker safely, and wake only after publishing state that makes the next poll useful.
- Reviewed
- Rust
- stable Rust, futures-core 0.3
- Targets
- all async targets
- Profiles
- dev, release, test
Direct answer
What this Rust failure means
- Why it happens
- The implementation returned Pending without arranging a future wake, registered the waker after checking state, or retained a stale waker from another task.
- First discriminating check
- Force data arrival between the readiness check and waker registration while counting every registration and wake call.
Poll::Pending does not mean “please try again whenever convenient.” It means the value is not ready and the implementation has arranged for the current task to be notified when progress may be possible.
The Stream contract states this directly for poll_next. If my stream returns Pending without preserving a wakeup path, an executor is allowed to leave it alone forever.
The broken shape
Here is a simplified stream backed by shared state:
use std::collections::VecDeque;
use std::sync::{Arc, Mutex};
use std::task::{Context, Poll, Waker};
struct Shared<T> {
queue: VecDeque<T>,
waker: Option<Waker>,
}
struct Events<T> {
shared: Arc<Mutex<Shared<T>>>,
}
A naive poll_next may check the queue, release the lock, and only later store the waker. If a producer pushes between those two actions, it sees no waker to notify. The consumer then stores its waker and returns Pending, but the state change already happened. Nobody wakes it.
This is the check-then-sleep race in async form.
The failing Stream fixture pins futures-core 0.3.32 and uses barriers to publish item 7 precisely after the empty check but before waker registration. It does not wait for scheduler luck: the wake counter remains zero and the incorrect expectation fails. The repaired Stream protects readiness and registration with one synchronization rule, observes exactly one wake, and returns the published item on the next poll. Both sides use the real Stream trait from the locked dependency.
Register and recheck under one synchronization rule
With one mutex protecting both queue and waker, the basic operation can be atomic:
use std::pin::Pin;
use futures_core::Stream;
impl<T> Stream for Events<T> {
type Item = T;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<T>> {
let mut shared = self.shared.lock().unwrap();
if let Some(item) = shared.queue.pop_front() {
return Poll::Ready(Some(item));
}
let replace = shared
.waker
.as_ref()
.is_none_or(|old| !old.will_wake(cx.waker()));
if replace {
shared.waker = Some(cx.waker().clone());
}
Poll::Pending
}
}
The producer changes readiness and takes the waker while holding the same lock, then wakes outside it:
let wake = {
let mut shared = shared.lock().unwrap();
shared.queue.push_back(item);
shared.waker.take()
};
if let Some(waker) = wake {
waker.wake();
}
Waking outside the lock avoids immediately scheduling code which may contend for that lock. The crucial ordering is that the item becomes visible before the wake.
Store the current waker
The task polling a stream can change. A waker captured during an earlier poll may no longer wake the current task. I compare with will_wake and clone the current cx.waker() when necessary.
Calling wake_by_ref does not replace registration. Waking once and then returning Pending without a future notification path can create a busy repoll or another stall.
For lock-free state, I prefer a reviewed primitive such as an atomic waker rather than inventing memory ordering. The algorithm must handle a producer racing with both the readiness check and registration.
Pending, item, and termination are distinct
A stream has three meanings:
Pending: no item now; arrange a later notification.Ready(Some(item)): deliver one item; the stream may continue.Ready(None): the stream has terminated and should not be polled again.
Returning None because a queue is temporarily empty terminates the stream. Returning Pending after the producer side is permanently gone can wait forever unless closure also wakes the consumer and lets it return None.
I keep an explicit closed flag alongside the queue. The consumer drains queued items first, then returns Ready(None) when closed and empty. The final producer drop or close operation publishes closed = true and wakes the registered task.
The first discriminating fixture
I insert barriers at the dangerous boundary:
consumer checks empty
consumer pauses before registration
producer publishes item and tries to wake
consumer registers and returns Pending
The broken implementation stalls deterministically. The fixed implementation prevents this sequence from losing the notification, either through one lock or a correct register-then-recheck algorithm.
Counting calls also helps:
poll=1 result=Pending waker=A
publish=item-7 wake=A
poll=2 result=Ready(Some(item-7))
If publish happens with no stored waker and there is no recheck, the trace identifies the missing edge.
False repairs
Calling cx.waker().wake_by_ref() immediately before every Pending makes the task repoll continuously. It may hide missed notifications while consuming a CPU core.
Polling the stream from a timer eventually finds the item, but replaces event-driven behavior with latency and load.
Holding a standard mutex guard while invoking arbitrary wake logic can create reentrancy or lock contention. Publish under the lock, take the waker, then wake outside.
The regression proof
I test item arrival before poll, during registration, and after Pending. I test waker replacement by polling from a different task. I test closure with an empty and non-empty queue. Finally, I assert that an idle stream does not self-wake repeatedly.
The stream is correct when every Pending has a future event capable of waking the latest interested task, and every permanent close leads to Ready(None).