Skip to main content

cranpose_services/
async_io.rs

1//! Waking a future from another thread.
2//!
3//! A platform's own I/O is blocking — a synchronous HTTP read, a provider call,
4//! a file descriptor a system API hands over — and the framework has no thread
5//! pool to hide that behind. What it has is the waker: work runs on a thread of
6//! its own and, when it has something, wakes whoever was awaiting it.
7//!
8//! Two shapes cover everything here. A [`Signal`] carries one value, which is
9//! what a request's status line is. A [`ChunkChannel`] carries a stream of byte
10//! chunks with the consumer's progress bounding the producer, which is what a
11//! response body is — without the bound, a slow reader and a fast server put the
12//! whole download in memory, which is the thing streaming exists to avoid.
13
14use std::{
15    collections::VecDeque,
16    future::Future,
17    pin::Pin,
18    sync::{Arc, Condvar, Mutex, PoisonError},
19    task::{Context, Poll, Waker},
20};
21
22/// How many chunks may wait ahead of the consumer before the producer stops
23/// reading. Enough to keep a socket busy across one scheduling gap, and far
24/// short of holding a download in memory.
25pub const MAX_PENDING_CHUNKS: usize = 8;
26
27struct SignalState<T> {
28    value: Option<T>,
29    waker: Option<Waker>,
30    closed: bool,
31}
32
33/// A one-value hand-off between a worker and a future.
34///
35/// The worker calls [`Signal::set`]; whoever awaits [`Signal::wait`] is woken
36/// with it. A worker that dies without setting anything closes the signal, and
37/// the wait resolves to `None` rather than hanging for ever.
38pub struct Signal<T> {
39    state: Arc<Mutex<SignalState<T>>>,
40}
41
42impl<T> Clone for Signal<T> {
43    fn clone(&self) -> Self {
44        Self {
45            state: Arc::clone(&self.state),
46        }
47    }
48}
49
50impl<T> Default for Signal<T> {
51    fn default() -> Self {
52        Self::new()
53    }
54}
55
56impl<T> Signal<T> {
57    pub fn new() -> Self {
58        Self {
59            state: Arc::new(Mutex::new(SignalState {
60                value: None,
61                waker: None,
62                closed: false,
63            })),
64        }
65    }
66
67    /// Delivers the value and wakes the waiter. A second call is ignored: one
68    /// signal carries one value.
69    pub fn set(&self, value: T) {
70        let waker = {
71            let mut state = lock(&self.state);
72            if state.closed {
73                return;
74            }
75            state.value = Some(value);
76            state.closed = true;
77            state.waker.take()
78        };
79        if let Some(waker) = waker {
80            waker.wake();
81        }
82    }
83
84    /// Ends the signal with no value, so a waiter stops waiting.
85    pub fn close(&self) {
86        let waker = {
87            let mut state = lock(&self.state);
88            if state.closed {
89                return;
90            }
91            state.closed = true;
92            state.waker.take()
93        };
94        if let Some(waker) = waker {
95            waker.wake();
96        }
97    }
98
99    /// Resolves with the value, or `None` when the signal was closed empty.
100    pub fn wait(&self) -> SignalWait<T> {
101        SignalWait {
102            state: Arc::clone(&self.state),
103        }
104    }
105}
106
107/// The future [`Signal::wait`] returns.
108pub struct SignalWait<T> {
109    state: Arc<Mutex<SignalState<T>>>,
110}
111
112impl<T> Future for SignalWait<T> {
113    type Output = Option<T>;
114
115    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Option<T>> {
116        let mut state = lock(&self.state);
117        if let Some(value) = state.value.take() {
118            return Poll::Ready(Some(value));
119        }
120        if state.closed {
121            return Poll::Ready(None);
122        }
123        state.waker = Some(context.waker().clone());
124        Poll::Pending
125    }
126}
127
128struct ChunkState<E> {
129    ready: VecDeque<Vec<u8>>,
130    error: Option<E>,
131    finished: bool,
132    abandoned: bool,
133    waker: Option<Waker>,
134}
135
136struct ChunkShared<E> {
137    state: Mutex<ChunkState<E>>,
138    room: Condvar,
139}
140
141/// The producing half of a chunked byte stream.
142///
143/// Held by whatever is doing the blocking read. Dropping it without calling
144/// [`ChunkChannel::finish`] or [`ChunkChannel::fail`] ends the stream, so a
145/// worker that panics does not leave a reader waiting for ever.
146pub struct ChunkChannel<E> {
147    shared: Arc<ChunkShared<E>>,
148}
149
150impl<E> ChunkChannel<E> {
151    /// Creates a channel and its reading half.
152    pub fn new() -> (Self, ChunkStream<E>) {
153        let shared = Arc::new(ChunkShared {
154            state: Mutex::new(ChunkState {
155                ready: VecDeque::new(),
156                error: None,
157                finished: false,
158                abandoned: false,
159                waker: None,
160            }),
161            room: Condvar::new(),
162        });
163        (
164            Self {
165                shared: Arc::clone(&shared),
166            },
167            ChunkStream { shared },
168        )
169    }
170
171    /// Publishes one chunk, waiting while the consumer is more than
172    /// [`MAX_PENDING_CHUNKS`] behind.
173    ///
174    /// Returns `false` once the consumer has gone, which is the producer's
175    /// signal to stop reading.
176    pub fn push(&self, chunk: Vec<u8>) -> bool {
177        let waker = {
178            let mut state = lock(&self.shared.state);
179            #[cfg(not(target_arch = "wasm32"))]
180            while state.ready.len() >= MAX_PENDING_CHUNKS && !state.abandoned {
181                state = self
182                    .shared
183                    .room
184                    .wait(state)
185                    .unwrap_or_else(PoisonError::into_inner);
186            }
187            if state.abandoned || state.finished {
188                return false;
189            }
190            state.ready.push_back(chunk);
191            state.waker.take()
192        };
193        if let Some(waker) = waker {
194            waker.wake();
195        }
196        true
197    }
198
199    /// Ends the stream with an error.
200    pub fn fail(&self, error: E) {
201        let waker = {
202            let mut state = lock(&self.shared.state);
203            if state.finished {
204                return;
205            }
206            state.error = Some(error);
207            state.finished = true;
208            state.waker.take()
209        };
210        if let Some(waker) = waker {
211            waker.wake();
212        }
213    }
214
215    /// Ends the stream normally.
216    pub fn finish(&self) {
217        let waker = {
218            let mut state = lock(&self.shared.state);
219            if state.finished {
220                return;
221            }
222            state.finished = true;
223            state.waker.take()
224        };
225        if let Some(waker) = waker {
226            waker.wake();
227        }
228    }
229
230    /// Whether the consumer has gone.
231    pub fn is_abandoned(&self) -> bool {
232        lock(&self.shared.state).abandoned
233    }
234}
235
236impl<E> Drop for ChunkChannel<E> {
237    fn drop(&mut self) {
238        self.finish();
239    }
240}
241
242/// The consuming half of a chunked byte stream.
243pub struct ChunkStream<E> {
244    shared: Arc<ChunkShared<E>>,
245}
246
247impl<E> ChunkStream<E> {
248    /// Resolves with the next chunk, `Ok(None)` at the end of the stream, or
249    /// the error the producer ended with.
250    pub fn next(&self) -> ChunkNext<'_, E> {
251        ChunkNext { stream: self }
252    }
253}
254
255impl<E> Drop for ChunkStream<E> {
256    fn drop(&mut self) {
257        let mut state = lock(&self.shared.state);
258        state.abandoned = true;
259        drop(state);
260        self.shared.room.notify_all();
261    }
262}
263
264/// The future [`ChunkStream::next`] returns.
265pub struct ChunkNext<'a, E> {
266    stream: &'a ChunkStream<E>,
267}
268
269impl<E> Future for ChunkNext<'_, E> {
270    type Output = Result<Option<Vec<u8>>, E>;
271
272    fn poll(self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
273        let shared = Arc::clone(&self.stream.shared);
274        let mut state = lock(&shared.state);
275        if let Some(chunk) = state.ready.pop_front() {
276            drop(state);
277            shared.room.notify_one();
278            return Poll::Ready(Ok(Some(chunk)));
279        }
280        if let Some(error) = state.error.take() {
281            return Poll::Ready(Err(error));
282        }
283        if state.finished {
284            return Poll::Ready(Ok(None));
285        }
286        state.waker = Some(context.waker().clone());
287        Poll::Pending
288    }
289}
290
291fn lock<T>(mutex: &Mutex<T>) -> std::sync::MutexGuard<'_, T> {
292    mutex.lock().unwrap_or_else(PoisonError::into_inner)
293}
294
295#[cfg(test)]
296#[path = "tests/async_io_tests.rs"]
297mod tests;