mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-07 00:00:08 +02:00
tokio-timer: reset timeout after elapsed in stream (#648)
This commit is contained in:
@@ -214,6 +214,7 @@ where T: Stream,
|
|||||||
match self.delay.poll() {
|
match self.delay.poll() {
|
||||||
Ok(Async::NotReady) => Ok(Async::NotReady),
|
Ok(Async::NotReady) => Ok(Async::NotReady),
|
||||||
Ok(Async::Ready(_)) => {
|
Ok(Async::Ready(_)) => {
|
||||||
|
self.delay.reset_timeout();
|
||||||
Err(Error::elapsed())
|
Err(Error::elapsed())
|
||||||
},
|
},
|
||||||
Err(e) => Err(Error::timer(e)),
|
Err(e) => Err(Error::timer(e)),
|
||||||
|
|||||||
@@ -152,3 +152,28 @@ fn stream_and_timeout_in_future() {
|
|||||||
assert!(item.is_some());
|
assert!(item.is_some());
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn idle_stream_timesout_periodically() {
|
||||||
|
mocked(|timer, _time| {
|
||||||
|
// Not yet complete
|
||||||
|
let (_tx, rx) = mpsc::unbounded::<()>();
|
||||||
|
|
||||||
|
// Wrap it with a deadline
|
||||||
|
let mut stream = Timeout::new(rx, ms(100));
|
||||||
|
|
||||||
|
// Not ready
|
||||||
|
assert_not_ready!(stream);
|
||||||
|
|
||||||
|
// Turn the timer, it runs for the elapsed time
|
||||||
|
advance(timer, ms(100));
|
||||||
|
|
||||||
|
assert_elapsed!(stream);
|
||||||
|
// Stream's timeout should reset
|
||||||
|
assert_not_ready!(stream);
|
||||||
|
|
||||||
|
// Turn the timer, it runs for the elapsed time
|
||||||
|
advance(timer, ms(100));
|
||||||
|
assert_elapsed!(stream);
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user