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
use crate::{StreamExt, TakeUntilIf, Trigger, Tripwire};
use futures_core::stream::Stream;
use pin_project::pin_project;
use std::future::Future;
use std::pin::Pin;
use std::task::{Context, Poll};
/// A `Valve` is associated with a [`Trigger`], and can be used to wrap one or more
/// asynchronous streams. All streams wrapped by a given `Valve` (or its clones) will be
/// interrupted when [`Trigger::close`] is called on the valve's associated handle.
#[derive(Clone, Debug)]
pub struct Valve(#[pin] Tripwire);
impl Valve {
/// Make a new `Valve` and an associated [`Trigger`].
pub fn new() -> (Trigger, Self) {
let (t, tw) = Tripwire::new();
(t, Valve(tw))
/// Wrap the given `stream` with this `Valve`.
/// When [`Trigger::close`] is called on the handle associated with this valve, the given
/// stream will immediately yield `None`.
pub fn wrap<S>(&self, stream: S) -> Valved<S>
S: Stream,
/// Check if the valve has been closed.
/// If `Ready`, contains `true` if the stream should be considered closed, and `false`
/// if the `Trigger` has been disabled.
pub fn poll_closed(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<bool> {
/// A `Valved` is wrapper around a `Stream` that enables the stream to be turned off remotely to
/// initiate a graceful shutdown. When a new `Valved` is created with [`Valved::new`], a handle to
/// that `Valved` is also produced; when [`Trigger::close`] is called on that handle, the
/// wrapped stream will immediately yield `None` to indicate that it has completed.
#[derive(Clone, Debug)]
pub struct Valved<S>(TakeUntilIf<S, Tripwire>);
impl<S> Valved<S> {
/// Make the given stream cancellable.
/// To cancel the stream, call [`Trigger::close`] on the returned handle.
pub fn new(stream: S) -> (Trigger, Self)
S: Stream,
let (vh, v) = Valve::new();
(vh, v.wrap(stream))
/// Consumes this wrapper, returning the underlying stream.
pub fn into_inner(self) -> S {
impl<S> Stream for Valved<S>
S: Stream,
type Item = S::Item;
fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
// safe since we never move nor leak &mut
let inner = unsafe { self.map_unchecked_mut(|s| &mut s.0) };
mod tests {
use super::*;
use futures_util::stream::empty;
fn valved_stream_may_be_dropped_safely() {
let _orphan = {
let s = empty::<()>();
let (trigger, valve) = Valve::new();
let _wrapped = valve.wrap(s);