diff --git a/tokio-stream/src/stream_ext/fuse.rs b/tokio-stream/src/stream_ext/fuse.rs index 9d117cf08..c94d5108b 100644 --- a/tokio-stream/src/stream_ext/fuse.rs +++ b/tokio-stream/src/stream_ext/fuse.rs @@ -1,5 +1,6 @@ use crate::Stream; +use futures_core::FusedStream; use pin_project_lite::pin_project; use std::pin::Pin; use std::task::{ready, Context, Poll}; @@ -51,3 +52,12 @@ where } } } + +impl FusedStream for Fuse +where + T: Stream, +{ + fn is_terminated(&self) -> bool { + self.stream.is_none() + } +} diff --git a/tokio-stream/tests/stream_fuse.rs b/tokio-stream/tests/stream_fuse.rs index 8bb56d3de..96bf33638 100644 --- a/tokio-stream/tests/stream_fuse.rs +++ b/tokio-stream/tests/stream_fuse.rs @@ -1,3 +1,4 @@ +use futures_core::FusedStream; use tokio_stream::{Stream, StreamExt}; use std::pin::Pin; @@ -37,16 +38,22 @@ async fn basic_usage() { // however, once it is fused let mut stream = stream.fuse(); + assert!(!stream.is_terminated()); assert_eq!(stream.size_hint(), (0, None)); assert_eq!(stream.next().await, Some(4)); + assert!(!stream.is_terminated()); assert_eq!(stream.size_hint(), (0, None)); assert_eq!(stream.next().await, None); + assert!(stream.is_terminated()); + // it will always return `None` after the first time. assert_eq!(stream.size_hint(), (0, Some(0))); assert_eq!(stream.next().await, None); assert_eq!(stream.size_hint(), (0, Some(0))); + + assert!(stream.is_terminated()); } #[tokio::test]