mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-07 00:00:08 +02:00
Fix tests
This commit is contained in:
+16
-7
@@ -41,10 +41,11 @@ impl Serialize for LineSerialize {
|
|||||||
|
|
||||||
pub struct EchoFramed<T> {
|
pub struct EchoFramed<T> {
|
||||||
inner: T,
|
inner: T,
|
||||||
|
eof: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T, U> Future for EchoFramed<T>
|
impl<T, U> Future for EchoFramed<T>
|
||||||
where T: FramedIo<In = U, Out = U>,
|
where T: FramedIo<In = U, Out = Option<U>>,
|
||||||
{
|
{
|
||||||
type Item = ();
|
type Item = ();
|
||||||
type Error = ();
|
type Error = ();
|
||||||
@@ -56,9 +57,15 @@ impl<T, U> Future for EchoFramed<T>
|
|||||||
}
|
}
|
||||||
|
|
||||||
// Wait until we can simultaneously read and write a message
|
// Wait until we can simultaneously read and write a message
|
||||||
while self.inner.poll_read().is_ready() && self.inner.poll_write().is_ready() {
|
while !self.eof &&
|
||||||
|
self.inner.poll_read().is_ready() &&
|
||||||
|
self.inner.poll_write().is_ready() {
|
||||||
let frame = match self.inner.read() {
|
let frame = match self.inner.read() {
|
||||||
Ok(Async::Ready(frame)) => frame,
|
Ok(Async::Ready(Some(frame))) => frame,
|
||||||
|
Ok(Async::Ready(None)) => {
|
||||||
|
self.eof = true;
|
||||||
|
break
|
||||||
|
}
|
||||||
Ok(Async::NotReady) => break,
|
Ok(Async::NotReady) => break,
|
||||||
Err(e) => panic!("error in read: {}", e),
|
Err(e) => panic!("error in read: {}", e),
|
||||||
};
|
};
|
||||||
@@ -72,9 +79,11 @@ impl<T, U> Future for EchoFramed<T>
|
|||||||
|
|
||||||
// If we wrote some frames try to flush again. Ignore whether this is
|
// If we wrote some frames try to flush again. Ignore whether this is
|
||||||
// ready to finish or not as we're going to continue to return NotReady
|
// ready to finish or not as we're going to continue to return NotReady
|
||||||
drop(self.inner.flush().expect("flush error"));
|
if self.inner.flush().expect("flush error").is_ready() && self.eof {
|
||||||
|
Ok(().into())
|
||||||
Ok(Async::NotReady)
|
} else {
|
||||||
|
Ok(Async::NotReady)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -89,7 +98,7 @@ fn echo() {
|
|||||||
let addr = listener.local_addr().unwrap();
|
let addr = listener.local_addr().unwrap();
|
||||||
let srv = listener.incoming().for_each(move |(socket, _)| {
|
let srv = listener.incoming().for_each(move |(socket, _)| {
|
||||||
let framed = EasyFramed::new(socket, LineParser, LineSerialize);
|
let framed = EasyFramed::new(socket, LineParser, LineSerialize);
|
||||||
handle.spawn(EchoFramed { inner: framed });
|
handle.spawn(EchoFramed { inner: framed, eof: false });
|
||||||
Ok(())
|
Ok(())
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user