stream: fix panic in ChunksTimeout::new (#5036)

This commit is contained in:
Marek Kuskowski
2022-09-27 22:34:08 +00:00
committed by GitHub
parent aedcec666d
commit 2df45234ef
@@ -1,6 +1,6 @@
use crate::stream_ext::Fuse; use crate::stream_ext::Fuse;
use crate::Stream; use crate::Stream;
use tokio::time::{sleep, Instant, Sleep}; use tokio::time::{sleep, Sleep};
use core::future::Future; use core::future::Future;
use core::pin::Pin; use core::pin::Pin;
@@ -16,7 +16,7 @@ pin_project! {
#[pin] #[pin]
stream: Fuse<S>, stream: Fuse<S>,
#[pin] #[pin]
deadline: Sleep, deadline: Option<Sleep>,
duration: Duration, duration: Duration,
items: Vec<S::Item>, items: Vec<S::Item>,
cap: usize, // https://github.com/rust-lang/futures-rs/issues/1475 cap: usize, // https://github.com/rust-lang/futures-rs/issues/1475
@@ -27,7 +27,7 @@ impl<S: Stream> ChunksTimeout<S> {
pub(super) fn new(stream: S, max_size: usize, duration: Duration) -> Self { pub(super) fn new(stream: S, max_size: usize, duration: Duration) -> Self {
ChunksTimeout { ChunksTimeout {
stream: Fuse::new(stream), stream: Fuse::new(stream),
deadline: sleep(duration), deadline: None,
duration, duration,
items: Vec::with_capacity(max_size), items: Vec::with_capacity(max_size),
cap: max_size, cap: max_size,
@@ -45,7 +45,7 @@ impl<S: Stream> Stream for ChunksTimeout<S> {
Poll::Pending => break, Poll::Pending => break,
Poll::Ready(Some(item)) => { Poll::Ready(Some(item)) => {
if me.items.is_empty() { if me.items.is_empty() {
me.deadline.as_mut().reset(Instant::now() + *me.duration); me.deadline.set(Some(sleep(*me.duration)));
me.items.reserve_exact(*me.cap); me.items.reserve_exact(*me.cap);
} }
me.items.push(item); me.items.push(item);
@@ -67,7 +67,9 @@ impl<S: Stream> Stream for ChunksTimeout<S> {
} }
if !me.items.is_empty() { if !me.items.is_empty() {
ready!(me.deadline.poll(cx)); if let Some(deadline) = me.deadline.as_pin_mut() {
ready!(deadline.poll(cx));
}
return Poll::Ready(Some(std::mem::take(me.items))); return Poll::Ready(Some(std::mem::take(me.items)));
} }