如何为需可变借用自身的结构体实现Stream trait?
Rust中为BinaryReader实现Stream trait的正确方式
我们有一个包含异步方法read_next的BinaryReader结构体,定义如下:
impl<R: AsyncRead + AsyncSeek + Unpin + Send> BinaryReader<R> { // ...... async fn read_next(&mut self) -> Result<Option<OwnedEvent>, Error> { // ...... todo!() } }
需求是为该结构体实现Stream trait,最初的疑问是:
有没有无需包装器直接实现Stream的方法?
尝试引入包装器后,代码无法通过编译:
#[pin_project::pin_project] struct BinaryReaderWrapper<R: AsyncRead + AsyncSeek + Unpin + Send> { #[pin] inner: BinaryReader<R>, #[pin] last_fut: Option<Pin<Box<dyn Future<Output = Result<Option<OwnedEvent>, Error>>>>> } impl<R: AsyncRead + AsyncSeek + Unpin + Send> Stream for BinaryReaderWrapper<R> { type Item = Result<Option<OwnedEvent>, Error>; fn poll_next(self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> { // - 我们称这个引用的生命周期为 `'1` let mut this = self.project(); if this.last_fut.is_none() { let ret = Box::pin(this.inner.read_next()); *this.last_fut = Some(ret); // ^^^ 生命周期可能不够长 // ^^^ 强制转换要求 `'1` 必须比 `'static` 活得更久 } let result = ready!(this.last_fut.as_pin_mut().unwrap().poll(cx)); this.last_fut = None; result } }
报错原因解析
编译器报错的核心是:Box<dyn Future<...>>默认要求trait对象的生命周期为'static,但调用this.inner.read_next()返回的future捕获了this.inner的可变引用(生命周期为poll_next方法中self的引用周期'1),该生命周期远短于'static,因此无法存入last_fut字段。而无接收器的异步函数不捕获外部引用,返回的future满足'static要求,所以能正常编译。
正确实现方案
方案1:用async-stream crate简化实现
这是最简洁的方式,无需手动管理future状态:
use async_stream::stream; use futures_core::Stream; // 直接为BinaryReader实现Stream impl<R: AsyncRead + AsyncSeek + Unpin + Send> Stream for BinaryReader<R> { type Item = Result<Option<OwnedEvent>, Error>; fn poll_next(self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> { let mut this = self; stream! { loop { yield this.as_mut().read_next().await; } }.poll_next(cx) } }
也可以通过函数将BinaryReader转换为Stream,无需修改原结构体:
fn into_stream<R: AsyncRead + AsyncSeek + Unpin + Send>(mut reader: BinaryReader<R>) -> impl Stream<Item = Result<Option<OwnedEvent>, Error>> { stream! { loop { yield reader.read_next().await; } } }
方案2:手动实现包装器,指定生命周期绑定
修改last_fut的类型,明确标注future的生命周期与结构体自身绑定,而非默认的'static:
#[pin_project::pin_project] struct BinaryReaderWrapper<R: AsyncRead + AsyncSeek + Unpin + Send> { #[pin] inner: BinaryReader<R>, #[pin] // 添加 + '_ 表示future生命周期与结构体一致 last_fut: Option<Pin<Box<dyn Future<Output = Result<Option<OwnedEvent>, Error>> + '_>>> } impl<R: AsyncRead + AsyncSeek + Unpin + Send> Stream for BinaryReaderWrapper<R> { type Item = Result<Option<OwnedEvent>, Error>; fn poll_next(self: Pin<&mut Self>, cx: &mut std::task::Context<'_>) -> Poll<Option<Self::Item>> { let mut this = self.project(); if this.last_fut.is_none() { let ret = Box::pin(this.inner.read_next()); *this.last_fut = Some(ret); } let result = ready!(this.last_fut.as_pin_mut().unwrap().poll(cx)); this.last_fut = None; // 处理流结束的情况:read_next返回Ok(None)时终止Stream match result { Ok(Some(event)) => Poll::Ready(Some(Ok(Some(event)))), Ok(None) => Poll::Ready(None), Err(e) => Poll::Ready(Some(Err(e))), } } }
关于“无需包装器直接实现Stream”的问题
可以直接为BinaryReader实现Stream,但需要给原结构体添加一个类似last_fut的字段来存储待poll的future。如果无法修改原结构体的定义,就必须使用包装器;若允许修改原结构体,将last_fut字段加入BinaryReader内部,再按照上述包装器的逻辑实现Stream即可。
内容的提问来源于stack exchange,提问作者zhou plus
相关产品推荐
相关产品推荐

