diff --git a/tokio/src/fs/read_dir.rs b/tokio/src/fs/read_dir.rs index 3b30eb0aa..a144a8637 100644 --- a/tokio/src/fs/read_dir.rs +++ b/tokio/src/fs/read_dir.rs @@ -33,11 +33,11 @@ const CHUNK_SIZE: usize = 32; pub async fn read_dir(path: impl AsRef) -> io::Result { let path = path.as_ref().to_owned(); asyncify(|| -> io::Result { - let mut std = std::fs::read_dir(path)?.fuse(); + let mut std = std::fs::read_dir(path)?; let mut buf = VecDeque::with_capacity(CHUNK_SIZE); - ReadDir::next_chunk(&mut buf, &mut std); + let remain = ReadDir::next_chunk(&mut buf, &mut std); - Ok(ReadDir(State::Idle(Some((buf, std))))) + Ok(ReadDir(State::Idle(Some((buf, std, remain))))) }) .await } @@ -64,12 +64,10 @@ pub async fn read_dir(path: impl AsRef) -> io::Result { #[must_use = "streams do nothing unless polled"] pub struct ReadDir(State); -type StdReadDir = std::iter::Fuse; - #[derive(Debug)] enum State { - Idle(Option<(VecDeque>, StdReadDir)>), - Pending(JoinHandle<(VecDeque>, StdReadDir)>), + Idle(Option<(VecDeque>, std::fs::ReadDir, bool)>), + Pending(JoinHandle<(VecDeque>, std::fs::ReadDir, bool)>), } impl ReadDir { @@ -105,38 +103,35 @@ impl ReadDir { loop { match self.0 { State::Idle(ref mut data) => { - let (buf, _) = data.as_mut().unwrap(); + let (buf, _, ref remain) = data.as_mut().unwrap(); if let Some(ent) = buf.pop_front() { return Poll::Ready(ent.map(Some)); - }; + } else if !remain { + return Poll::Ready(Ok(None)); + } - let (mut buf, mut std) = data.take().unwrap(); + let (mut buf, mut std, _) = data.take().unwrap(); self.0 = State::Pending(spawn_blocking(move || { - ReadDir::next_chunk(&mut buf, &mut std); - (buf, std) + let remain = ReadDir::next_chunk(&mut buf, &mut std); + (buf, std, remain) })); } State::Pending(ref mut rx) => { - let (mut buf, std) = ready!(Pin::new(rx).poll(cx))?; - - let ret = match buf.pop_front() { - Some(Ok(x)) => Ok(Some(x)), - Some(Err(e)) => Err(e), - None => Ok(None), - }; - - self.0 = State::Idle(Some((buf, std))); - - return Poll::Ready(ret); + self.0 = State::Idle(Some(ready!(Pin::new(rx).poll(cx))?)); } } } } - fn next_chunk(buf: &mut VecDeque>, std: &mut StdReadDir) { - for ret in std.by_ref().take(CHUNK_SIZE) { + fn next_chunk(buf: &mut VecDeque>, std: &mut std::fs::ReadDir) -> bool { + for _ in 0..CHUNK_SIZE { + let ret = match std.next() { + Some(ret) => ret, + None => return false, + }; + let success = ret.is_ok(); buf.push_back(ret.map(|std| DirEntry { @@ -154,6 +149,8 @@ impl ReadDir { break; } } + + true } }