diff --git a/tokio-fs/src/blocking.rs b/tokio-fs/src/blocking.rs index 57a982e8c..5cee2d13b 100644 --- a/tokio-fs/src/blocking.rs +++ b/tokio-fs/src/blocking.rs @@ -18,6 +18,8 @@ use self::State::*; pub(crate) struct Blocking { inner: Option, state: State, + /// true if the lower IO layer needs flushing + need_flush: bool, } #[derive(Debug)] @@ -39,6 +41,7 @@ impl Blocking { Blocking { inner: Some(inner), state: State::Idle(Some(Buf::with_capacity(0))), + need_flush: false, } } } @@ -119,6 +122,7 @@ where (res, buf, inner) })); + self.need_flush = true; return Ready(Ok(n)); } @@ -135,16 +139,35 @@ where } fn poll_flush(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - let (res, buf, inner) = match self.state { - Idle(_) => return Ready(Ok(())), - Busy(ref mut rx) => ready!(Pin::new(rx).poll(cx)), - }; + loop { + let need_flush = self.need_flush; + match self.state { + // The buffer is not used here + Idle(ref mut buf_cell) => { + if need_flush { + let buf = buf_cell.take().unwrap(); + let mut inner = self.inner.take().unwrap(); - // The buffer is not used here - self.state = Idle(Some(buf)); - self.inner = Some(inner); + self.state = Busy(sys::run(move || { + let res = inner.flush().map(|_| 0); + (res, buf, inner) + })); - Ready(res.map(|_| ())) + self.need_flush = false; + } else { + return Ready(Ok(())); + } + } + Busy(ref mut rx) => { + let (res, buf, inner) = ready!(Pin::new(rx).poll(cx)); + self.state = Idle(Some(buf)); + self.inner = Some(inner); + + // If error, return + res?; + } + } + } } fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> {