From 47cc31590fbc5662241355c86e331c42e6152775 Mon Sep 17 00:00:00 2001 From: stormshield-fabs <161044609+stormshield-fabs@users.noreply.github.com> Date: Tue, 28 Apr 2026 17:33:54 +0200 Subject: [PATCH] tokio-stream: implement FusedStream for Fuse (#8090) --- tokio-stream/src/stream_ext/fuse.rs | 10 ++++++++++ tokio-stream/tests/stream_fuse.rs | 7 +++++++ 2 files changed, 17 insertions(+) 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]