mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-02 00:00:11 +02:00
stream: add StreamExt::then (#4355)
This commit is contained in:
@@ -1,3 +1,4 @@
|
||||
use core::future::Future;
|
||||
use futures_core::Stream;
|
||||
|
||||
mod all;
|
||||
@@ -39,15 +40,18 @@ use skip::Skip;
|
||||
mod skip_while;
|
||||
use skip_while::SkipWhile;
|
||||
|
||||
mod try_next;
|
||||
use try_next::TryNext;
|
||||
|
||||
mod take;
|
||||
use take::Take;
|
||||
|
||||
mod take_while;
|
||||
use take_while::TakeWhile;
|
||||
|
||||
mod then;
|
||||
use then::Then;
|
||||
|
||||
mod try_next;
|
||||
use try_next::TryNext;
|
||||
|
||||
cfg_time! {
|
||||
mod timeout;
|
||||
use timeout::Timeout;
|
||||
@@ -197,6 +201,51 @@ pub trait StreamExt: Stream {
|
||||
Map::new(self, f)
|
||||
}
|
||||
|
||||
/// Maps this stream's items asynchronously to a different type, returning a
|
||||
/// new stream of the resulting type.
|
||||
///
|
||||
/// The provided closure is executed over all elements of this stream as
|
||||
/// they are made available, and the returned future is executed. Only one
|
||||
/// future is executed at the time.
|
||||
///
|
||||
/// Note that this function consumes the stream passed into it and returns a
|
||||
/// wrapped version of it, similar to the existing `then` methods in the
|
||||
/// standard library.
|
||||
///
|
||||
/// Be aware that if the future is not `Unpin`, then neither is the `Stream`
|
||||
/// returned by this method. To handle this, you can use `tokio::pin!` as in
|
||||
/// the example below or put the stream in a `Box` with `Box::pin(stream)`.
|
||||
///
|
||||
/// # Examples
|
||||
///
|
||||
/// ```
|
||||
/// # #[tokio::main]
|
||||
/// # async fn main() {
|
||||
/// use tokio_stream::{self as stream, StreamExt};
|
||||
///
|
||||
/// async fn do_async_work(value: i32) -> i32 {
|
||||
/// value + 3
|
||||
/// }
|
||||
///
|
||||
/// let stream = stream::iter(1..=3);
|
||||
/// let stream = stream.then(do_async_work);
|
||||
///
|
||||
/// tokio::pin!(stream);
|
||||
///
|
||||
/// assert_eq!(stream.next().await, Some(4));
|
||||
/// assert_eq!(stream.next().await, Some(5));
|
||||
/// assert_eq!(stream.next().await, Some(6));
|
||||
/// # }
|
||||
/// ```
|
||||
fn then<F, Fut>(self, f: F) -> Then<Self, Fut, F>
|
||||
where
|
||||
F: FnMut(Self::Item) -> Fut,
|
||||
Fut: Future,
|
||||
Self: Sized,
|
||||
{
|
||||
Then::new(self, f)
|
||||
}
|
||||
|
||||
/// Combine two streams into one by interleaving the output of both as it
|
||||
/// is produced.
|
||||
///
|
||||
|
||||
Reference in New Issue
Block a user