diff --git a/Cargo.toml b/Cargo.toml index 40f86dccb..5db634722 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -6,7 +6,7 @@ members = [ "tokio-codec", "tokio-current-thread", "tokio-executor", - # "tokio-fs", + "tokio-fs", "tokio-futures", "tokio-io", "tokio-macros", diff --git a/azure-pipelines.yml b/azure-pipelines.yml index 20880950d..2e815d2de 100644 --- a/azure-pipelines.yml +++ b/azure-pipelines.yml @@ -30,7 +30,7 @@ jobs: cross: true rust: $(nightly) crates: -# - tokio-fs + tokio-fs: [] tokio-reactor: [] tokio-signal: [] tokio-tcp: diff --git a/tokio-fs/Cargo.toml b/tokio-fs/Cargo.toml index 59d2ae57a..10090fef7 100644 --- a/tokio-fs/Cargo.toml +++ b/tokio-fs/Cargo.toml @@ -24,9 +24,10 @@ categories = ["asynchronous", "network-programming", "filesystem"] publish = false [dependencies] -futures = "0.1.21" +futures-core-preview = "0.3.0-alpha.17" tokio-threadpool = { version = "0.2.0", path = "../tokio-threadpool" } tokio-io = { version = "0.2.0", path = "../tokio-io" } +tokio-futures = { version = "0.2.0", path = "../tokio-futures" } [dev-dependencies] rand = "0.6" @@ -34,3 +35,6 @@ tempfile = "3" tempdir = "0.3" tokio-codec = { version = "0.2.0", path = "../tokio-codec" } tokio = { version = "0.2.0", path = "../tokio" } +futures-channel-preview = "0.3.0-alpha.17" +futures-preview = { version = "0.3.0-alpha.17" } +futures-util-preview = "0.3.0-alpha.17" diff --git a/tokio-fs/examples/std-echo.rs b/tokio-fs/examples_old/std-echo.rs similarity index 89% rename from tokio-fs/examples/std-echo.rs rename to tokio-fs/examples_old/std-echo.rs index 58028bb37..0b55c1297 100644 --- a/tokio-fs/examples/std-echo.rs +++ b/tokio-fs/examples_old/std-echo.rs @@ -1,15 +1,17 @@ //! Echo everything received on STDIN to STDOUT. #![deny(deprecated, warnings)] +#![feature(async_await)] use tokio_codec::{FramedRead, FramedWrite, LinesCodec}; use tokio_fs::{stderr, stdin, stdout}; use tokio_threadpool::Builder; -use futures::{Future, Sink, Stream}; +use futures_util::sink::SinkExt; use std::io; -pub fn main() -> Result<(), Box> { +#[tokio::main] +async fn main() -> Result<(), Box> { let pool = Builder::new().pool_size(1).build(); pool.spawn({ diff --git a/tokio-fs/src/create_dir.rs b/tokio-fs/src/create_dir.rs index a07c1ad52..ac784075f 100644 --- a/tokio-fs/src/create_dir.rs +++ b/tokio-fs/src/create_dir.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Creates a new, empty directory at the provided path /// @@ -34,10 +37,9 @@ impl

Future for CreateDirFuture

where P: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::create_dir(&self.path)) } } diff --git a/tokio-fs/src/create_dir_all.rs b/tokio-fs/src/create_dir_all.rs index 9f034fa46..914c0b177 100644 --- a/tokio-fs/src/create_dir_all.rs +++ b/tokio-fs/src/create_dir_all.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Recursively create a directory and all of its parent components if they /// are missing. @@ -35,10 +38,9 @@ impl

Future for CreateDirAllFuture

where P: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::create_dir_all(&self.path)) } } diff --git a/tokio-fs/src/file/clone.rs b/tokio-fs/src/file/clone.rs index 1bbb6cc9e..24a443069 100644 --- a/tokio-fs/src/file/clone.rs +++ b/tokio-fs/src/file/clone.rs @@ -1,6 +1,9 @@ use super::File; -use futures::{Future, Poll}; +use std::future::Future; use std::io; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Future returned by `File::try_clone`. /// @@ -21,15 +24,16 @@ impl CloneFuture { } impl Future for CloneFuture { - type Item = (File, File); - type Error = (File, io::Error); + type Output = Result<(File, File), (File, io::Error)>; - fn poll(&mut self) -> Poll { - self.file + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { + let inner_self = Pin::get_mut(self); + inner_self + .file .as_mut() .expect("Cannot poll `CloneFuture` after it resolves") .poll_try_clone() - .map(|inner| inner.map(|cloned| (self.file.take().unwrap(), cloned))) - .map_err(|err| (self.file.take().unwrap(), err)) + .map(|inner| inner.map(|cloned| (inner_self.file.take().unwrap(), cloned))) + .map_err(|err| (inner_self.file.take().unwrap(), err)) } } diff --git a/tokio-fs/src/file/create.rs b/tokio-fs/src/file/create.rs index da03779e8..bc4661cb8 100644 --- a/tokio-fs/src/file/create.rs +++ b/tokio-fs/src/file/create.rs @@ -1,8 +1,11 @@ use super::File; -use futures::{try_ready, Future, Poll}; use std::fs::File as StdFile; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Future returned by `File::create` and resolves to a `File` instance. #[derive(Debug)] @@ -12,7 +15,7 @@ pub struct CreateFuture

{ impl

CreateFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { pub(crate) fn new(path: P) -> Self { CreateFuture { path } @@ -21,15 +24,14 @@ where impl

Future for CreateFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { - type Item = File; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { - let std = try_ready!(crate::blocking_io(|| StdFile::create(&self.path))); + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { + let std = ready!(crate::blocking_io(|| StdFile::create(&self.path)))?; let file = File::from_std(std); - Ok(file.into()) + Poll::Ready(Ok(file.into())) } } diff --git a/tokio-fs/src/file/metadata.rs b/tokio-fs/src/file/metadata.rs index 7a3c5ef91..a5aa4683b 100644 --- a/tokio-fs/src/file/metadata.rs +++ b/tokio-fs/src/file/metadata.rs @@ -1,8 +1,11 @@ use super::File; -use futures::{try_ready, Future, Poll}; use std::fs::File as StdFile; use std::fs::Metadata; +use std::future::Future; use std::io; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; const POLL_AFTER_RESOLVE: &str = "Cannot poll MetadataFuture after it resolves"; @@ -23,13 +26,13 @@ impl MetadataFuture { } impl Future for MetadataFuture { - type Item = (File, Metadata); - type Error = io::Error; + type Output = io::Result<(File, Metadata)>; - fn poll(&mut self) -> Poll { - let metadata = try_ready!(crate::blocking_io(|| StdFile::metadata(self.std()))); + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { + let inner = Pin::get_mut(self); + let metadata = ready!(crate::blocking_io(|| StdFile::metadata(inner.std())))?; - let file = self.file.take().expect(POLL_AFTER_RESOLVE); - Ok((file, metadata).into()) + let file = inner.file.take().expect(POLL_AFTER_RESOLVE); + Poll::Ready(Ok((file, metadata).into())) } } diff --git a/tokio-fs/src/file/mod.rs b/tokio-fs/src/file/mod.rs index c51bf5b5f..76b51ae85 100644 --- a/tokio-fs/src/file/mod.rs +++ b/tokio-fs/src/file/mod.rs @@ -16,10 +16,12 @@ pub use self::open::OpenFuture; pub use self::open_options::OpenOptions; pub use self::seek::SeekFuture; -use futures::Poll; use std::fs::{File as StdFile, Metadata, Permissions}; use std::io::{self, Read, Seek, Write}; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use tokio_io::{AsyncRead, AsyncWrite}; /// A reference to an open file on the filesystem. @@ -103,7 +105,7 @@ impl File { /// ``` pub fn open

(path: P) -> OpenFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { OpenOptions::new().read(true).open(path) } @@ -142,7 +144,7 @@ impl File { /// ``` pub fn create

(path: P) -> CreateFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { CreateFuture::new(path) } @@ -191,7 +193,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_seek(&mut self, pos: io::SeekFrom) -> Poll { + pub fn poll_seek(&mut self, pos: io::SeekFrom) -> Poll> { crate::blocking_io(|| self.std().seek(pos)) } @@ -243,7 +245,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_sync_all(&mut self) -> Poll<(), io::Error> { + pub fn poll_sync_all(&mut self) -> Poll> { crate::blocking_io(|| self.std().sync_all()) } @@ -273,7 +275,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_sync_data(&mut self) -> Poll<(), io::Error> { + pub fn poll_sync_data(&mut self) -> Poll> { crate::blocking_io(|| self.std().sync_data()) } @@ -305,7 +307,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_set_len(&mut self, size: u64) -> Poll<(), io::Error> { + pub fn poll_set_len(&mut self, size: u64) -> Poll> { crate::blocking_io(|| self.std().set_len(size)) } @@ -344,7 +346,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_metadata(&mut self) -> Poll { + pub fn poll_metadata(&mut self) -> Poll> { crate::blocking_io(|| self.std().metadata()) } @@ -366,7 +368,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_try_clone(&mut self) -> Poll { + pub fn poll_try_clone(&mut self) -> Poll> { crate::blocking_io(|| { let std = self.std().try_clone()?; Ok(File::from_std(std)) @@ -437,7 +439,7 @@ impl File { /// /// tokio::run(task); /// ``` - pub fn poll_set_permissions(&mut self, perm: Permissions) -> Poll<(), io::Error> { + pub fn poll_set_permissions(&mut self, perm: Permissions) -> Poll> { crate::blocking_io(|| self.std().set_permissions(perm)) } @@ -479,8 +481,15 @@ impl Read for File { } impl AsyncRead for File { - unsafe fn prepare_uninitialized_buffer(&self, _: &mut [u8]) -> bool { - false + fn poll_read( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &mut [u8], + ) -> Poll> { + match Pin::get_mut(self).read(buf) { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } } } @@ -495,11 +504,26 @@ impl Write for File { } impl AsyncWrite for File { - fn shutdown(&mut self) -> Poll<(), io::Error> { - crate::blocking_io(|| { - self.std = None; - Ok(()) - }) + fn poll_write( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + match Pin::get_mut(self).write(buf) { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + match Pin::get_mut(self).flush() { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) } } diff --git a/tokio-fs/src/file/open.rs b/tokio-fs/src/file/open.rs index b95af1609..5f7915eeb 100644 --- a/tokio-fs/src/file/open.rs +++ b/tokio-fs/src/file/open.rs @@ -1,19 +1,22 @@ use super::File; -use futures::{try_ready, Future, Poll}; use std::fs::OpenOptions as StdOpenOptions; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Future returned by `File::open` and resolves to a `File` instance. #[derive(Debug)] -pub struct OpenFuture

{ +pub struct OpenFuture { options: StdOpenOptions, path: P, } impl

OpenFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { pub(crate) fn new(options: StdOpenOptions, path: P) -> Self { OpenFuture { options, path } @@ -22,15 +25,14 @@ where impl

Future for OpenFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { - type Item = File; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { - let std = try_ready!(crate::blocking_io(|| self.options.open(&self.path))); + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { + let std = ready!(crate::blocking_io(|| self.options.open(&self.path)))?; let file = File::from_std(std); - Ok(file.into()) + Poll::Ready(Ok(file.into())) } } diff --git a/tokio-fs/src/file/open_options.rs b/tokio-fs/src/file/open_options.rs index 70b78211f..41f232e1c 100644 --- a/tokio-fs/src/file/open_options.rs +++ b/tokio-fs/src/file/open_options.rs @@ -90,7 +90,7 @@ impl OpenOptions { /// [`open`]: https://doc.rust-lang.org/std/fs/struct.OpenOptions.html#method.open pub fn open

(&self, path: P) -> OpenFuture

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { OpenFuture::new(self.0.clone(), path) } diff --git a/tokio-fs/src/file/seek.rs b/tokio-fs/src/file/seek.rs index 16f30e1be..be0d997df 100644 --- a/tokio-fs/src/file/seek.rs +++ b/tokio-fs/src/file/seek.rs @@ -1,6 +1,9 @@ use super::File; -use futures::{try_ready, Future, Poll}; +use std::future::Future; use std::io; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Future returned by `File::seek`. #[derive(Debug)] @@ -19,16 +22,16 @@ impl SeekFuture { } impl Future for SeekFuture { - type Item = (File, u64); - type Error = io::Error; + type Output = io::Result<(File, u64)>; - fn poll(&mut self) -> Poll { - let pos = try_ready!(self + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { + let inner_self = Pin::get_mut(self); + let pos = ready!(inner_self .inner .as_mut() .expect("Cannot poll `SeekFuture` after it resolves") - .poll_seek(self.pos)); - let inner = self.inner.take().unwrap(); - Ok((inner, pos).into()) + .poll_seek(inner_self.pos))?; + let inner = inner_self.inner.take().unwrap(); + Poll::Ready(Ok((inner, pos).into())) } } diff --git a/tokio-fs/src/hard_link.rs b/tokio-fs/src/hard_link.rs index 277eedbe7..b4a023b3f 100644 --- a/tokio-fs/src/hard_link.rs +++ b/tokio-fs/src/hard_link.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Creates a new hard link on the filesystem. /// @@ -41,10 +44,9 @@ where P: AsRef, Q: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::hard_link(&self.src, &self.dst)) } } diff --git a/tokio-fs/src/lib.rs b/tokio-fs/src/lib.rs index 89b956232..2070059da 100644 --- a/tokio-fs/src/lib.rs +++ b/tokio-fs/src/lib.rs @@ -30,6 +30,9 @@ //! [`AsyncRead`]: https://docs.rs/tokio-io/0.1/tokio_io/trait.AsyncRead.html //! [tokio-threadpool]: https://docs.rs/tokio-threadpool/0.1/tokio_threadpool +#[macro_use] +extern crate tokio_futures; + mod create_dir; mod create_dir_all; pub mod file; @@ -68,20 +71,19 @@ pub use crate::stdout::{stdout, Stdout}; pub use crate::symlink_metadata::{symlink_metadata, SymlinkMetadataFuture}; pub use crate::write::{write, WriteFile}; -use futures::Async::*; -use futures::Poll; use std::io; use std::io::ErrorKind::{Other, WouldBlock}; +use std::task::Poll; +use std::task::Poll::*; -fn blocking_io(f: F) -> Poll +fn blocking_io(f: F) -> Poll> where F: FnOnce() -> io::Result, { match tokio_threadpool::blocking(f) { - Ok(Ready(Ok(v))) => Ok(v.into()), - Ok(Ready(Err(err))) => Err(err), - Ok(NotReady) => Ok(NotReady), - Err(_) => Err(blocking_err()), + Ready(Ok(v)) => Ready(v), + Ready(Err(_)) => Ready(Err(blocking_err())), + Pending => Pending, } } @@ -90,13 +92,13 @@ where F: FnOnce() -> io::Result, { match tokio_threadpool::blocking(f) { - Ok(Ready(Ok(v))) => Ok(v), - Ok(Ready(Err(err))) => { + Ready(Ok(Ok(v))) => Ok(v), + Ready(Ok(Err(err))) => { debug_assert_ne!(err.kind(), WouldBlock); Err(err) } - Ok(NotReady) => Err(WouldBlock.into()), - Err(_) => Err(blocking_err()), + Ready(Err(_)) => Err(blocking_err()), + Pending => Err(blocking_err()), } } diff --git a/tokio-fs/src/metadata.rs b/tokio-fs/src/metadata.rs index ece1e05ae..11a252494 100644 --- a/tokio-fs/src/metadata.rs +++ b/tokio-fs/src/metadata.rs @@ -1,8 +1,11 @@ use super::blocking_io; -use futures::{Future, Poll}; use std::fs::{self, Metadata}; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Queries the file system metadata for a path. pub fn metadata

(path: P) -> MetadataFuture

@@ -34,10 +37,9 @@ impl

Future for MetadataFuture

where P: AsRef + Send + 'static, { - type Item = Metadata; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { blocking_io(|| fs::metadata(&self.path)) } } diff --git a/tokio-fs/src/os/unix.rs b/tokio-fs/src/os/unix.rs index 84b62c8f9..b6f8e0be0 100644 --- a/tokio-fs/src/os/unix.rs +++ b/tokio-fs/src/os/unix.rs @@ -1,9 +1,12 @@ //! Unix-specific extensions to primitives in the `tokio_fs` module. -use futures::{Future, Poll}; +use std::future::Future; use std::io; use std::os::unix::fs; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Creates a new symbolic link on the filesystem. /// @@ -42,10 +45,9 @@ where P: AsRef, Q: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::symlink(&self.src, &self.dst)) } } diff --git a/tokio-fs/src/os/windows/symlink_dir.rs b/tokio-fs/src/os/windows/symlink_dir.rs index 90f4a04f5..e322deb99 100644 --- a/tokio-fs/src/os/windows/symlink_dir.rs +++ b/tokio-fs/src/os/windows/symlink_dir.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; +use std::future::Future; use std::io; use std::os::windows::fs; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Creates a new directory symlink on the filesystem. /// @@ -41,10 +44,9 @@ where P: AsRef, Q: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::symlink_dir(&self.src, &self.dst)) } } diff --git a/tokio-fs/src/os/windows/symlink_file.rs b/tokio-fs/src/os/windows/symlink_file.rs index eb7c9d04e..afd6b2297 100644 --- a/tokio-fs/src/os/windows/symlink_file.rs +++ b/tokio-fs/src/os/windows/symlink_file.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; +use std::future::Future; use std::io; use std::os::windows::fs; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Creates a new file symbolic link on the filesystem. /// @@ -41,10 +44,9 @@ where P: AsRef, Q: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::symlink_file(&self.src, &self.dst)) } } diff --git a/tokio-fs/src/read.rs b/tokio-fs/src/read.rs index fbcce0354..cfd0afe3f 100644 --- a/tokio-fs/src/read.rs +++ b/tokio-fs/src/read.rs @@ -1,7 +1,11 @@ use crate::{file, File}; -use futures::{try_ready, Async, Future, Poll}; +use std::future::Future; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use std::{io, mem, path::Path}; use tokio_io; +use tokio_io::AsyncRead; /// Creates a future which will open a file for reading and read the entire /// contents into a buffer and return said buffer. @@ -25,7 +29,7 @@ use tokio_io; /// ``` pub fn read

(path: P) -> ReadFile

where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { ReadFile { state: State::Open(File::open(path)), @@ -34,41 +38,50 @@ where /// A future used to open a file and read its entire contents into a buffer. #[derive(Debug)] -pub struct ReadFile + Send + 'static> { +pub struct ReadFile + Send + Unpin + 'static> { state: State

, } #[derive(Debug)] -enum State + Send + 'static> { +enum State + Send + Unpin + 'static> { Open(file::OpenFuture

), Metadata(file::MetadataFuture), - Read(tokio_io::io::ReadToEnd), + Reading(Vec, usize, File), + Empty, } -impl + Send + 'static> Future for ReadFile

{ - type Item = Vec; - type Error = io::Error; +impl + Send + Unpin + 'static> Future for ReadFile

{ + type Output = io::Result>; - fn poll(&mut self) -> Poll { - let new_state = match &mut self.state { + fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + let inner = Pin::get_mut(self); + match &mut inner.state { State::Open(ref mut open_file) => { - let file = try_ready!(open_file.poll()); - State::Metadata(file.metadata()) + let file = ready!(Pin::new(open_file).poll(cx))?; + let new_state = State::Metadata(file.metadata()); + mem::replace(&mut inner.state, new_state); + Pin::new(inner).poll(cx) } State::Metadata(read_metadata) => { - let (file, metadata) = try_ready!(read_metadata.poll()); + let (file, metadata) = ready!(Pin::new(read_metadata).poll(cx))?; let buf = Vec::with_capacity(metadata.len() as usize + 1); - let read = tokio_io::io::read_to_end(file, buf); - State::Read(read) + let new_state = State::Reading(buf, 0, file); + mem::replace(&mut inner.state, new_state); + Pin::new(inner).poll(cx) } - State::Read(ref mut read) => { - let (_, buf) = try_ready!(read.poll()); - return Ok(Async::Ready(buf)); + State::Reading(buf, ref mut pos, file) => { + let n = ready!(Pin::new(file).poll_read_buf(cx, buf))?; + *pos += n; + if *pos >= buf.len() { + match mem::replace(&mut inner.state, State::Empty) { + State::Reading(buf, _, _) => Poll::Ready(Ok(buf)), + _ => panic!(), + } + } else { + Poll::Pending + } } - }; - - mem::replace(&mut self.state, new_state); - // Getting here means we transitionsed state. Must poll the new state. - self.poll() + State::Empty => panic!("poll a WriteFile after it's done"), + } } } diff --git a/tokio-fs/src/read_dir.rs b/tokio-fs/src/read_dir.rs index 855e23328..01e36d018 100644 --- a/tokio-fs/src/read_dir.rs +++ b/tokio-fs/src/read_dir.rs @@ -1,10 +1,14 @@ -use futures::{Future, Poll, Stream}; +use futures_core::stream::Stream; use std::ffi::OsString; use std::fs::{self, DirEntry as StdDirEntry, FileType, Metadata, ReadDir as StdReadDir}; +use std::future::Future; use std::io; #[cfg(unix)] use std::os::unix::fs::DirEntryExt; use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Returns a stream over the entries within a directory. /// @@ -40,10 +44,9 @@ impl

Future for ReadDirFuture

where P: AsRef + Send + 'static, { - type Item = ReadDir; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| Ok(ReadDir(fs::read_dir(&self.path)?))) } } @@ -68,15 +71,19 @@ where pub struct ReadDir(StdReadDir); impl Stream for ReadDir { - type Item = DirEntry; - type Error = io::Error; + type Item = io::Result; - fn poll(&mut self) -> Poll, Self::Error> { - crate::blocking_io(|| match self.0.next() { + fn poll_next(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + let inner = Pin::get_mut(self); + match crate::blocking_io(|| match inner.0.next() { Some(Err(err)) => Err(err), - Some(Ok(item)) => Ok(Some(DirEntry(item))), + Some(Ok(item)) => Ok(Some(Ok(DirEntry(item)))), None => Ok(None), - }) + }) { + Poll::Ready(Err(err)) => Poll::Ready(Some(Err(err))), + Poll::Ready(Ok(v)) => Poll::Ready(v), + Poll::Pending => Poll::Pending, + } } } @@ -181,7 +188,7 @@ impl DirEntry { /// /// tokio::run(fut); /// ``` - pub fn poll_metadata(&self) -> Poll { + pub fn poll_metadata(&self) -> Poll> { crate::blocking_io(|| self.0.metadata()) } @@ -213,7 +220,7 @@ impl DirEntry { /// /// tokio::run(fut); /// ``` - pub fn poll_file_type(&self) -> Poll { + pub fn poll_file_type(&self) -> Poll> { crate::blocking_io(|| self.0.file_type()) } } diff --git a/tokio-fs/src/read_link.rs b/tokio-fs/src/read_link.rs index a672c97d7..d5b9fe2d4 100644 --- a/tokio-fs/src/read_link.rs +++ b/tokio-fs/src/read_link.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::{Path, PathBuf}; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Reads a symbolic link, returning the file that the link points to. /// @@ -34,10 +37,9 @@ impl

Future for ReadLinkFuture

where P: AsRef, { - type Item = PathBuf; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::read_link(&self.path)) } } diff --git a/tokio-fs/src/remove_dir.rs b/tokio-fs/src/remove_dir.rs index 74acffe39..b0de35e16 100644 --- a/tokio-fs/src/remove_dir.rs +++ b/tokio-fs/src/remove_dir.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Removes an existing, empty directory. /// @@ -34,10 +37,9 @@ impl

Future for RemoveDirFuture

where P: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::remove_dir(&self.path)) } } diff --git a/tokio-fs/src/remove_file.rs b/tokio-fs/src/remove_file.rs index d2d4b86d3..ca8b80888 100644 --- a/tokio-fs/src/remove_file.rs +++ b/tokio-fs/src/remove_file.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Removes a file from the filesystem. /// @@ -38,10 +41,9 @@ impl

Future for RemoveFileFuture

where P: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::remove_file(&self.path)) } } diff --git a/tokio-fs/src/rename.rs b/tokio-fs/src/rename.rs index 1c2977ef5..742ced04b 100644 --- a/tokio-fs/src/rename.rs +++ b/tokio-fs/src/rename.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Rename a file or directory to a new name, replacing the original file if /// `to` already exists. @@ -41,10 +44,9 @@ where P: AsRef, Q: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::rename(&self.from, &self.to)) } } diff --git a/tokio-fs/src/set_permissions.rs b/tokio-fs/src/set_permissions.rs index dab062dc7..d5bdfd674 100644 --- a/tokio-fs/src/set_permissions.rs +++ b/tokio-fs/src/set_permissions.rs @@ -1,7 +1,10 @@ -use futures::{Future, Poll}; use std::fs; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Changes the permissions found on a file or a directory. /// @@ -38,10 +41,9 @@ impl

Future for SetPermissionsFuture

where P: AsRef, { - type Item = (); - type Error = io::Error; + type Output = io::Result<()>; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { crate::blocking_io(|| fs::set_permissions(&self.path, self.perm.clone())) } } diff --git a/tokio-fs/src/stderr.rs b/tokio-fs/src/stderr.rs index 97ee2343c..273e8e221 100644 --- a/tokio-fs/src/stderr.rs +++ b/tokio-fs/src/stderr.rs @@ -1,5 +1,7 @@ -use futures::Poll; use std::io::{self, Stderr as StdStderr, Write}; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use tokio_io::AsyncWrite; /// A handle to the standard error stream of a process. @@ -36,7 +38,25 @@ impl Write for Stderr { } impl AsyncWrite for Stderr { - fn shutdown(&mut self) -> Poll<(), io::Error> { - Ok(().into()) + fn poll_write( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + match Pin::get_mut(self).write(buf) { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + match Pin::get_mut(self).flush() { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) } } diff --git a/tokio-fs/src/stdin.rs b/tokio-fs/src/stdin.rs index f4339047d..3c5a387b1 100644 --- a/tokio-fs/src/stdin.rs +++ b/tokio-fs/src/stdin.rs @@ -1,4 +1,7 @@ use std::io::{self, Read, Stdin as StdStdin}; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use tokio_io::AsyncRead; /// A handle to the standard input stream of a process. @@ -37,7 +40,14 @@ impl Read for Stdin { } impl AsyncRead for Stdin { - unsafe fn prepare_uninitialized_buffer(&self, _: &mut [u8]) -> bool { - false + fn poll_read( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &mut [u8], + ) -> Poll> { + match Pin::get_mut(self).read(buf) { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } } } diff --git a/tokio-fs/src/stdout.rs b/tokio-fs/src/stdout.rs index 6141d7dd4..a661dbf34 100644 --- a/tokio-fs/src/stdout.rs +++ b/tokio-fs/src/stdout.rs @@ -1,5 +1,7 @@ -use futures::Poll; use std::io::{self, Stdout as StdStdout, Write}; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use tokio_io::AsyncWrite; /// A handle to the standard output stream of a process. @@ -36,7 +38,25 @@ impl Write for Stdout { } impl AsyncWrite for Stdout { - fn shutdown(&mut self) -> Poll<(), io::Error> { - Ok(().into()) + fn poll_write( + self: Pin<&mut Self>, + _cx: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + match Pin::get_mut(self).write(buf) { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_flush(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + match Pin::get_mut(self).flush() { + Err(ref e) if e.kind() == io::ErrorKind::WouldBlock => Poll::Pending, + other => Poll::Ready(other), + } + } + + fn poll_shutdown(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) } } diff --git a/tokio-fs/src/symlink_metadata.rs b/tokio-fs/src/symlink_metadata.rs index f3b256aa3..cc4509590 100644 --- a/tokio-fs/src/symlink_metadata.rs +++ b/tokio-fs/src/symlink_metadata.rs @@ -1,8 +1,11 @@ use super::blocking_io; -use futures::{Future, Poll}; use std::fs::{self, Metadata}; +use std::future::Future; use std::io; use std::path::Path; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; /// Queries the file system metadata for a path. /// @@ -38,10 +41,9 @@ impl

Future for SymlinkMetadataFuture

where P: AsRef + Send + 'static, { - type Item = Metadata; - type Error = io::Error; + type Output = io::Result; - fn poll(&mut self) -> Poll { + fn poll(self: Pin<&mut Self>, _cx: &mut Context<'_>) -> Poll { blocking_io(|| fs::symlink_metadata(&self.path)) } } diff --git a/tokio-fs/src/write.rs b/tokio-fs/src/write.rs index b6f7511d0..4b63984c8 100644 --- a/tokio-fs/src/write.rs +++ b/tokio-fs/src/write.rs @@ -1,7 +1,11 @@ use crate::{file, File}; -use futures::{try_ready, Async, Future, Poll}; +use std::future::Future; +use std::pin::Pin; +use std::task::Context; +use std::task::Poll; use std::{fmt, io, mem, path::Path}; use tokio_io; +use tokio_io::AsyncWrite; /// Creates a future that will open a file for writing and write the entire /// contents of `contents` to it. @@ -25,9 +29,9 @@ use tokio_io; /// /// tokio::run(task); /// ``` -pub fn write>(path: P, contents: C) -> WriteFile +pub fn write + Unpin>(path: P, contents: C) -> WriteFile where - P: AsRef + Send + 'static, + P: AsRef + Send + Unpin + 'static, { WriteFile { state: State::Create(File::create(path), Some(contents)), @@ -37,35 +41,62 @@ where /// A future used to open a file for writing and write the entire contents /// of some data to it. #[derive(Debug)] -pub struct WriteFile + Send + 'static, C: AsRef<[u8]>> { +pub struct WriteFile + Send + Unpin + 'static, C: AsRef<[u8]> + Unpin> { state: State, } #[derive(Debug)] -enum State + Send + 'static, C: AsRef<[u8]>> { +enum State + Send + Unpin + 'static, C: AsRef<[u8]> + Unpin> { Create(file::CreateFuture

, Option), - Write(tokio_io::io::WriteAll), + Writing { f: File, buf: C, pos: usize }, + Empty, } -impl + Send + 'static, C: AsRef<[u8]> + fmt::Debug> Future for WriteFile { - type Item = C; - type Error = io::Error; +fn zero_write() -> io::Error { + io::Error::new(io::ErrorKind::WriteZero, "zero-length write") +} - fn poll(&mut self) -> Poll { - let new_state = match &mut self.state { - State::Create(ref mut create_file, contents) => { - let file = try_ready!(create_file.poll()); - let write = tokio_io::io::write_all(file, contents.take().unwrap()); - State::Write(write) - } - State::Write(ref mut write) => { - let (_, contents) = try_ready!(write.poll()); - return Ok(Async::Ready(contents)); - } - }; +impl + Send + Unpin + 'static, C: AsRef<[u8]> + Unpin + fmt::Debug> Future + for WriteFile +{ + type Output = io::Result; - mem::replace(&mut self.state, new_state); - // We just entered the Write state, need to poll it before returning. - self.poll() + fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + let inner = Pin::get_mut(self); + match &mut inner.state { + State::Create(create_file, contents) => { + let file = ready!(Pin::new(create_file).poll(cx))?; + let contents = contents.take().unwrap(); + let new_state = State::Writing { + f: file, + buf: contents, + pos: 0, + }; + mem::replace(&mut inner.state, new_state); + // We just entered the Write state, need to poll it before returning. + return Pin::new(inner).poll(cx); + } + State::Empty => panic!("poll a WriteFile after it's done"), + _ => {} + } + + match mem::replace(&mut inner.state, State::Empty) { + State::Writing { + mut f, + buf, + mut pos, + } => { + let buf_ref = buf.as_ref(); + while pos < buf_ref.len() { + let n = ready!(Pin::new(&mut f).poll_write(cx, &buf_ref[pos..]))?; + pos += n; + if n == 0 { + return Poll::Ready(Err(zero_write())); + } + } + Poll::Ready(Ok(buf)) + } + _ => panic!(), + } } } diff --git a/tokio-fs/tests/dir.rs b/tokio-fs/tests/dir.rs index a22c048ea..74bb4d814 100644 --- a/tokio-fs/tests/dir.rs +++ b/tokio-fs/tests/dir.rs @@ -1,6 +1,8 @@ #![deny(warnings, rust_2018_idioms)] +#![feature(async_await)] -use futures::{Future, Stream}; +use futures_util::future; +use futures_util::try_stream::TryStreamExt; use std::fs; use std::sync::{Arc, Mutex}; use tempdir::TempDir; @@ -12,32 +14,44 @@ mod pool; fn create() { let base_dir = TempDir::new("base").unwrap(); let new_dir = base_dir.path().join("foo"); + let new_dir_2 = new_dir.clone(); - pool::run({ create_dir(new_dir.clone()) }); + pool::run(async move { + create_dir(new_dir).await?; + Ok(()) + }); - assert!(new_dir.is_dir()); + assert!(new_dir_2.is_dir()); } #[test] fn create_all() { let base_dir = TempDir::new("base").unwrap(); let new_dir = base_dir.path().join("foo").join("bar"); + let new_dir_2 = new_dir.clone(); - pool::run({ create_dir_all(new_dir.clone()) }); + pool::run(async move { + create_dir_all(new_dir).await?; + Ok(()) + }); - assert!(new_dir.is_dir()); + assert!(new_dir_2.is_dir()); } #[test] fn remove() { let base_dir = TempDir::new("base").unwrap(); let new_dir = base_dir.path().join("foo"); + let new_dir_2 = new_dir.clone(); fs::create_dir(new_dir.clone()).unwrap(); - pool::run({ remove_dir(new_dir.clone()) }); + pool::run(async move { + remove_dir(new_dir).await?; + Ok(()) + }); - assert!(!new_dir.exists()); + assert!(!new_dir_2.exists()); } #[test] @@ -53,12 +67,17 @@ fn read() { let f = files.clone(); let p = p.to_path_buf(); - pool::run({ - read_dir(p).flatten_stream().for_each(move |e| { - let s = e.file_name().to_str().unwrap().to_string(); - f.lock().unwrap().push(s); - Ok(()) - }) + + pool::run(async move { + let read_dir_fut = read_dir(p).await?; + read_dir_fut + .try_for_each(move |e| { + let s = e.file_name().to_str().unwrap().to_string(); + f.lock().unwrap().push(s); + future::ok(()) + }) + .await?; + Ok(()) }); let mut files = files.lock().unwrap(); diff --git a/tokio-fs/tests/file.rs b/tokio-fs/tests/file.rs index 421c0b36b..e2e028654 100644 --- a/tokio-fs/tests/file.rs +++ b/tokio-fs/tests/file.rs @@ -1,13 +1,13 @@ #![deny(warnings, rust_2018_idioms)] +#![feature(async_await)] -use futures::future::poll_fn; -use futures::Future; +use futures_util::future::poll_fn; use rand::{distributions, thread_rng, Rng}; use std::fs; use std::io::SeekFrom; use tempfile::Builder as TmpBuilder; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio_fs::*; -use tokio_io::io; mod pool; @@ -27,32 +27,25 @@ fn read_write() { .collect::() .into(); - pool::run({ - let file_path = file_path.clone(); - let contents = contents.clone(); + let file_path_2 = file_path.clone(); + let contents_2 = contents.clone(); - File::create(file_path) - .and_then(|file| file.metadata()) - .inspect(|&(_, ref metadata)| assert!(metadata.is_file())) - .and_then(move |(file, _)| io::write_all(file, contents)) - .and_then(|(mut file, _)| poll_fn(move || file.poll_sync_all())) - .then(|res| { - let _ = res.unwrap(); - Ok(()) - }) + pool::run(async move { + let file = File::create(file_path).await?; + let (mut file, metadata) = file.metadata().await?; + assert!(metadata.is_file()); + file.write(&contents).await?; + poll_fn(move |_cx| file.poll_sync_all()).await?; + Ok(()) }); - let dst = fs::read(&file_path).unwrap(); - assert_eq!(dst, contents); + let dst = fs::read(&file_path_2).unwrap(); + assert_eq!(dst, contents_2); - pool::run({ - File::open(file_path) - .and_then(|file| io::read_to_end(file, vec![])) - .then(move |res| { - let (_, buf) = res.unwrap(); - assert_eq!(buf, contents); - Ok(()) - }) + pool::run(async move { + let buf = read(file_path_2).await?; + assert_eq!(buf, contents_2); + Ok(()) }); } @@ -72,20 +65,21 @@ fn read_write_helpers() { .collect::() .into(); - pool::run(write(file_path.clone(), contents.clone()).then(|res| { - let _ = res.unwrap(); + let file_path_2 = file_path.clone(); + let contents_2 = contents.clone(); + + pool::run(async move { + write(file_path, contents).await?; Ok(()) - })); + }); - let dst = fs::read(&file_path).unwrap(); - assert_eq!(dst, contents); + let dst = fs::read(&file_path_2).unwrap(); + assert_eq!(dst, contents_2); - pool::run({ - read(file_path).then(move |res| { - let buf = res.unwrap(); - assert_eq!(buf, contents); - Ok(()) - }) + pool::run(async move { + let buf = read(file_path_2).await?; + assert_eq!(buf, contents_2); + Ok(()) }); } @@ -97,22 +91,12 @@ fn metadata() { .unwrap(); let file_path = dir.path().join("metadata.txt"); - pool::run({ - let file_path = file_path.clone(); - let file_path2 = file_path.clone(); - let file_path3 = file_path.clone(); - - tokio_fs::metadata(file_path) - .then(|r| { - let _ = r.err().unwrap(); - Ok(()) - }) - .and_then(|_| File::create(file_path2)) - .and_then(|_| tokio_fs::metadata(file_path3)) - .then(|r| { - assert!(r.unwrap().is_file()); - Ok(()) - }) + pool::run(async move { + assert!(tokio_fs::metadata(file_path.clone()).await.is_err()); + File::create(file_path.clone()).await?; + let metadata = tokio_fs::metadata(file_path.clone()).await?; + assert!(metadata.is_file()); + Ok(()) }); } @@ -124,28 +108,24 @@ fn seek() { .unwrap(); let file_path = dir.path().join("seek.txt"); - pool::run({ - OpenOptions::new() + pool::run(async move { + let mut file = OpenOptions::new() .create(true) .read(true) .write(true) .open(file_path) - .and_then(|file| io::write_all(file, "Hello, world!")) - .and_then(|(file, _)| file.seek(SeekFrom::End(-6))) - .and_then(|(file, _)| io::read_exact(file, vec![0; 5])) - .and_then(|(file, buf)| { - assert_eq!(buf, b"world"); - file.seek(SeekFrom::Start(0)) - }) - .and_then(|(file, _)| io::read_exact(file, vec![0; 5])) - .and_then(|(_, buf)| { - assert_eq!(buf, b"Hello"); - Ok(()) - }) - .then(|r| { - let _ = r.unwrap(); - Ok(()) - }) + .await + .unwrap(); + assert!(file.write(b"Hello, world!").await.is_ok()); + let mut file = file.seek(SeekFrom::End(-6)).await.unwrap().0; + let mut buf = vec![0; 5]; + assert!(file.read(buf.as_mut()).await.is_ok()); + assert_eq!(buf, b"world"); + let mut file = file.seek(SeekFrom::Start(0)).await.unwrap().0; + let mut buf = vec![0; 5]; + assert!(file.read(buf.as_mut()).await.is_ok()); + assert_eq!(buf, b"Hello"); + Ok(()) }); } @@ -158,24 +138,19 @@ fn clone() { .tempdir() .unwrap(); let file_path = dir.path().join("clone.txt"); + let file_path_2 = file_path.clone(); - pool::run( - File::create(file_path.clone()) - .and_then(|file| { - file.try_clone() - .map_err(|(_file, err)| err) - .and_then(|(file, clone)| { - io::write_all(file, "clone ") - .and_then(|_| io::write_all(clone, "successful")) - }) - }) - .then(|res| { - let _ = res.unwrap(); - Ok(()) - }), - ); + pool::run(async move { + let file = File::create(file_path.clone()).await.unwrap(); + let (mut file, mut clone) = file.try_clone().await.unwrap(); + assert!(AsyncWriteExt::write(&mut file, b"clone ").await.is_ok()); + assert!(AsyncWriteExt::write(&mut clone, b"successful") + .await + .is_ok()); + Ok(()) + }); - let mut file = std::fs::File::open(&file_path).unwrap(); + let mut file = std::fs::File::open(&file_path_2).unwrap(); let mut dst = vec![]; file.read_to_end(&mut dst).unwrap(); diff --git a/tokio-fs/tests/link.rs b/tokio-fs/tests/link.rs index 33364d985..b853085ad 100644 --- a/tokio-fs/tests/link.rs +++ b/tokio-fs/tests/link.rs @@ -1,4 +1,5 @@ #![deny(warnings, rust_2018_idioms)] +#![feature(async_await)] use std::fs; use std::io::prelude::*; @@ -19,7 +20,12 @@ fn test_hard_link() { file.write_all(b"hello").unwrap(); } - pool::run({ hard_link(src, dst.clone()) }); + let dst_2 = dst.clone(); + + pool::run(async move { + assert!(hard_link(src, dst_2.clone()).await.is_ok()); + Ok(()) + }); let mut content = String::new(); @@ -35,8 +41,6 @@ fn test_hard_link() { #[cfg(unix)] #[test] fn test_symlink() { - use futures::Future; - let dir = TempDir::new("base").unwrap(); let src = dir.path().join("src.txt"); let dst = dir.path().join("dst.txt"); @@ -46,7 +50,15 @@ fn test_symlink() { file.write_all(b"hello").unwrap(); } - pool::run({ os::unix::symlink(src.clone(), dst.clone()) }); + let src_2 = src.clone(); + let dst_2 = dst.clone(); + + pool::run(async move { + assert!(os::unix::symlink(src_2.clone(), dst_2.clone()) + .await + .is_ok()); + Ok(()) + }); let mut content = String::new(); @@ -58,6 +70,12 @@ fn test_symlink() { assert!(content == "hello"); - pool::run({ read_link(dst.clone()).map(move |x| assert!(x == src)) }); - pool::run({ symlink_metadata(dst.clone()).map(move |x| assert!(x.file_type().is_symlink())) }); + pool::run(async move { + let read = read_link(dst.clone()).await.unwrap(); + assert!(read == src); + + let symlink_meta = symlink_metadata(dst.clone()).await.unwrap(); + assert!(symlink_meta.file_type().is_symlink()); + Ok(()) + }); } diff --git a/tokio-fs/tests/pool/mod.rs b/tokio-fs/tests/pool/mod.rs index eeadd2123..9cce3b036 100644 --- a/tokio-fs/tests/pool/mod.rs +++ b/tokio-fs/tests/pool/mod.rs @@ -1,17 +1,19 @@ -use futures; use tokio_threadpool; use self::tokio_threadpool::Builder; -use futures::sync::oneshot; -use futures::Future; +use std::future::Future; use std::io; +use std::sync::mpsc; pub fn run(f: F) where - F: Future + Send + 'static, + F: Future> + Send + 'static, { let pool = Builder::new().pool_size(1).build(); - let (tx, rx) = oneshot::channel::<()>(); - pool.spawn(f.then(|_| tx.send(()))); - rx.wait().unwrap() + let (tx, rx) = mpsc::channel(); + pool.spawn(async move { + f.await.unwrap(); + tx.send(()).unwrap(); + }); + rx.recv().unwrap() } diff --git a/tokio/Cargo.toml b/tokio/Cargo.toml index 1963c16c2..62a4b3b80 100644 --- a/tokio/Cargo.toml +++ b/tokio/Cargo.toml @@ -68,7 +68,7 @@ bytes = { version = "0.4", optional = true } num_cpus = { version = "1.8.0", optional = true } tokio-codec = { version = "0.2.0", optional = true, path = "../tokio-codec" } tokio-current-thread = { version = "0.2.0", optional = true, path = "../tokio-current-thread" } -#tokio-fs = { version = "0.2.0", optional = true, path = "../tokio-fs" } +tokio-fs = { version = "0.2.0", optional = true, path = "../tokio-fs" } tokio-io = { version = "0.2.0", optional = true, path = "../tokio-io" } tokio-executor = { version = "0.2.0", optional = true, path = "../tokio-executor" } tokio-macros = { version = "0.2.0", optional = true, path = "../tokio-macros" }