Skip to main content

tokio_stream/stream_ext/
chain.rs

1use crate::stream_ext::Fuse;
2use crate::Stream;
3
4use core::pin::Pin;
5use core::task::{ready, Context, Poll};
6use futures_core::FusedStream;
7use pin_project_lite::pin_project;
8
9pin_project! {
10    /// Stream returned by the [`chain`](super::StreamExt::chain) method.
11    pub struct Chain<T, U> {
12        #[pin]
13        a: Fuse<T>,
14        #[pin]
15        b: U,
16    }
17}
18
19impl<T, U> Chain<T, U> {
20    pub(super) fn new(a: T, b: U) -> Chain<T, U>
21    where
22        T: Stream,
23        U: Stream,
24    {
25        Chain { a: Fuse::new(a), b }
26    }
27}
28
29impl<T, U> Stream for Chain<T, U>
30where
31    T: Stream,
32    U: Stream<Item = T::Item>,
33{
34    type Item = T::Item;
35
36    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<T::Item>> {
37        use Poll::Ready;
38
39        let me = self.project();
40
41        if let Some(v) = ready!(me.a.poll_next(cx)) {
42            return Ready(Some(v));
43        }
44
45        me.b.poll_next(cx)
46    }
47
48    fn size_hint(&self) -> (usize, Option<usize>) {
49        super::merge_size_hints(self.a.size_hint(), self.b.size_hint())
50    }
51}
52
53impl<T, U> FusedStream for Chain<T, U>
54where
55    T: Stream,
56    U: FusedStream<Item = T::Item>,
57{
58    fn is_terminated(&self) -> bool {
59        self.a.is_terminated() && self.b.is_terminated()
60    }
61}