mirror of
https://github.com/tokio-rs/tokio.git
synced 2026-09-09 00:00:08 +02:00
fs: move into tokio (#1672)
A step towards collapsing Tokio sub crates into a single `tokio` crate (#1318). The `fs` implementation is now provided by the main `tokio` crate. The `fs` functionality may still be excluded from the build by skipping the `fs` feature flag.
This commit is contained in:
@@ -0,0 +1,74 @@
|
||||
#![warn(rust_2018_idioms)]
|
||||
|
||||
use tokio::fs;
|
||||
use tokio_test::assert_ok;
|
||||
|
||||
use futures_util::future;
|
||||
use futures_util::try_stream::TryStreamExt;
|
||||
use std::sync::{Arc, Mutex};
|
||||
use tempfile::tempdir;
|
||||
|
||||
#[tokio::test]
|
||||
async fn create_dir() {
|
||||
let base_dir = tempdir().unwrap();
|
||||
let new_dir = base_dir.path().join("foo");
|
||||
let new_dir_2 = new_dir.clone();
|
||||
|
||||
assert_ok!(fs::create_dir(new_dir).await);
|
||||
|
||||
assert!(new_dir_2.is_dir());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn create_all() {
|
||||
let base_dir = tempdir().unwrap();
|
||||
let new_dir = base_dir.path().join("foo").join("bar");
|
||||
let new_dir_2 = new_dir.clone();
|
||||
|
||||
assert_ok!(fs::create_dir_all(new_dir).await);
|
||||
assert!(new_dir_2.is_dir());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn remove() {
|
||||
let base_dir = tempdir().unwrap();
|
||||
let new_dir = base_dir.path().join("foo");
|
||||
let new_dir_2 = new_dir.clone();
|
||||
|
||||
std::fs::create_dir(new_dir.clone()).unwrap();
|
||||
|
||||
assert_ok!(fs::remove_dir(new_dir).await);
|
||||
assert!(!new_dir_2.exists());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn read() {
|
||||
let base_dir = tempdir().unwrap();
|
||||
|
||||
let p = base_dir.path();
|
||||
std::fs::create_dir(p.join("aa")).unwrap();
|
||||
std::fs::create_dir(p.join("bb")).unwrap();
|
||||
std::fs::create_dir(p.join("cc")).unwrap();
|
||||
|
||||
let files = Arc::new(Mutex::new(Vec::new()));
|
||||
|
||||
let f = files.clone();
|
||||
let p = p.to_path_buf();
|
||||
|
||||
let read_dir_fut = fs::read_dir(p).await.unwrap();
|
||||
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
|
||||
.unwrap();
|
||||
|
||||
let mut files = files.lock().unwrap();
|
||||
files.sort(); // because the order is not guaranteed
|
||||
assert_eq!(
|
||||
*files,
|
||||
vec!["aa".to_string(), "bb".to_string(), "cc".to_string()]
|
||||
);
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
#![warn(rust_2018_idioms)]
|
||||
|
||||
use tokio::fs::File;
|
||||
use tokio::prelude::*;
|
||||
|
||||
use std::io::prelude::*;
|
||||
use tempfile::NamedTempFile;
|
||||
|
||||
const HELLO: &[u8] = b"hello world...";
|
||||
|
||||
#[tokio::test]
|
||||
async fn basic_read() {
|
||||
let mut tempfile = tempfile();
|
||||
tempfile.write_all(HELLO).unwrap();
|
||||
|
||||
let mut file = File::open(tempfile.path()).await.unwrap();
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let n = file.read(&mut buf).await.unwrap();
|
||||
|
||||
assert_eq!(n, HELLO.len());
|
||||
assert_eq!(&buf[..n], HELLO);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn basic_write() {
|
||||
let tempfile = tempfile();
|
||||
|
||||
let mut file = File::create(tempfile.path()).await.unwrap();
|
||||
|
||||
file.write_all(HELLO).await.unwrap();
|
||||
file.flush().await.unwrap();
|
||||
|
||||
let file = std::fs::read(tempfile.path()).unwrap();
|
||||
assert_eq!(file, HELLO);
|
||||
}
|
||||
|
||||
fn tempfile() -> NamedTempFile {
|
||||
NamedTempFile::new().unwrap()
|
||||
}
|
||||
@@ -0,0 +1,748 @@
|
||||
#![warn(rust_2018_idioms)]
|
||||
|
||||
// Load source
|
||||
#[allow(warnings)]
|
||||
#[path = "../src/fs/file.rs"]
|
||||
mod file;
|
||||
use file::File;
|
||||
|
||||
#[allow(warnings)]
|
||||
#[path = "../src/fs/blocking.rs"]
|
||||
mod blocking;
|
||||
|
||||
// Load mocked types
|
||||
mod support {
|
||||
pub(crate) mod mock_file;
|
||||
pub(crate) mod mock_pool;
|
||||
}
|
||||
pub(crate) use support::mock_pool as pool;
|
||||
|
||||
// Place them where the source expects them
|
||||
pub(crate) mod fs {
|
||||
pub(crate) use crate::blocking;
|
||||
|
||||
pub(crate) mod sys {
|
||||
pub(crate) use crate::support::mock_file::File;
|
||||
pub(crate) use crate::support::mock_pool::{run, Blocking};
|
||||
}
|
||||
|
||||
pub(crate) use crate::support::mock_pool::asyncify;
|
||||
}
|
||||
use fs::sys;
|
||||
|
||||
use tokio::prelude::*;
|
||||
use tokio_test::{assert_pending, assert_ready, assert_ready_err, assert_ready_ok, task};
|
||||
|
||||
use std::io::SeekFrom;
|
||||
|
||||
const HELLO: &[u8] = b"hello world...";
|
||||
const FOO: &[u8] = b"foo bar baz...";
|
||||
|
||||
#[test]
|
||||
fn open_read() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_eq!(0, pool::len());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, mock.remaining());
|
||||
assert_eq!(1, pool::len());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
assert_eq!(0, mock.remaining());
|
||||
assert!(t.is_woken());
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, HELLO.len());
|
||||
assert_eq!(&buf[..n], HELLO);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_twice_before_dispatch() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(pool::len(), 1);
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(&buf[..n], HELLO);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_with_smaller_buf() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
{
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
{
|
||||
let mut buf = [0; 4];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, 4);
|
||||
assert_eq!(&buf[..], &HELLO[..n]);
|
||||
}
|
||||
|
||||
// Calling again immediately succeeds with the rest of the buffer
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, 10);
|
||||
assert_eq!(&buf[..n], &HELLO[4..]);
|
||||
|
||||
assert_eq!(0, pool::len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_with_bigger_buf() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(&HELLO[..4]).read(&HELLO[4..]);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
{
|
||||
let mut buf = [0; 4];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
{
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, 4);
|
||||
assert_eq!(&buf[..n], &HELLO[..n]);
|
||||
}
|
||||
|
||||
// Calling again immediately succeeds with the rest of the buffer
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, 10);
|
||||
assert_eq!(&buf[..n], &HELLO[4..]);
|
||||
|
||||
assert_eq!(0, pool::len());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_err_then_read_success() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read_err().read(&HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
{
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
{
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
|
||||
assert_eq!(n, HELLO.len());
|
||||
assert_eq!(&buf[..n], HELLO);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
|
||||
assert_eq!(0, pool::len());
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
assert_eq!(1, mock.remaining());
|
||||
assert_eq!(1, pool::len());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
assert_eq!(0, mock.remaining());
|
||||
assert!(!t.is_woken());
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn flush_while_idle() {
|
||||
let (_mock, file) = sys::File::mock();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_with_buffer_larger_than_max() {
|
||||
// Chunks
|
||||
let a = 16 * 1024;
|
||||
let b = a * 2;
|
||||
let c = a * 3;
|
||||
let d = a * 4;
|
||||
|
||||
assert_eq!(d / 1024, 64);
|
||||
|
||||
let mut data = vec![];
|
||||
for i in 0..(d - 1) {
|
||||
data.push((i % 151) as u8);
|
||||
}
|
||||
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(&data[0..a])
|
||||
.read(&data[a..b])
|
||||
.read(&data[b..c])
|
||||
.read(&data[c..]);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut actual = vec![0; d];
|
||||
let mut pos = 0;
|
||||
|
||||
while pos < data.len() {
|
||||
let mut t = task::spawn(file.read(&mut actual[pos..]));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
pool::run_one();
|
||||
assert!(t.is_woken());
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert!(n <= a);
|
||||
|
||||
pos += n;
|
||||
}
|
||||
|
||||
assert_eq!(mock.remaining(), 0);
|
||||
assert_eq!(data, &actual[..data.len()]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_with_buffer_larger_than_max() {
|
||||
// Chunks
|
||||
let a = 16 * 1024;
|
||||
let b = a * 2;
|
||||
let c = a * 3;
|
||||
let d = a * 4;
|
||||
|
||||
assert_eq!(d / 1024, 64);
|
||||
|
||||
let mut data = vec![];
|
||||
for i in 0..(d - 1) {
|
||||
data.push((i % 151) as u8);
|
||||
}
|
||||
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(&data[0..a])
|
||||
.write(&data[a..b])
|
||||
.write(&data[b..c])
|
||||
.write(&data[c..]);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut rem = &data[..];
|
||||
|
||||
let mut first = true;
|
||||
|
||||
while !rem.is_empty() {
|
||||
let mut t = task::spawn(file.write(rem));
|
||||
|
||||
if !first {
|
||||
assert_pending!(t.poll());
|
||||
pool::run_one();
|
||||
assert!(t.is_woken());
|
||||
}
|
||||
|
||||
first = false;
|
||||
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
|
||||
rem = &rem[n..];
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
assert_eq!(mock.remaining(), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_twice_before_dispatch() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).write(FOO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(pool::len(), 1);
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(pool::len(), 1);
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn incomplete_read_followed_by_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO)
|
||||
.seek_current_ok(-(HELLO.len() as i64), 0)
|
||||
.write(FOO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 32];
|
||||
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
assert_eq!(pool::len(), 1);
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn incomplete_partial_read_followed_by_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO).seek_current_ok(-10, 0).write(FOO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 32];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut buf = [0; 4];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
assert_eq!(pool::len(), 1);
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn incomplete_read_followed_by_flush() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO)
|
||||
.seek_current_ok(-(HELLO.len() as i64), 0)
|
||||
.write(FOO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 32];
|
||||
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn incomplete_flush_followed_by_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).write(FOO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
assert_eq!(n, HELLO.len());
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
// TODO: Move under write
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn read_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read_err();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
assert!(t.is_woken());
|
||||
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_write_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write_err();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_read_write_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write_err().read(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_read_flush_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write_err().read(HELLO);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut buf = [0; 1024];
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_seek_write_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write_err().seek_start_ok(0);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
{
|
||||
let mut t = task::spawn(file.seek(SeekFrom::Start(0)));
|
||||
assert_pending!(t.poll());
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.write(FOO));
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn write_seek_flush_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write_err().seek_start_ok(0);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
pool::run_one();
|
||||
|
||||
{
|
||||
let mut t = task::spawn(file.seek(SeekFrom::Start(0)));
|
||||
assert_pending!(t.poll());
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
let mut t = task::spawn(file.flush());
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sync_all_ordered_after_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).sync_all();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.sync_all());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sync_all_err_ordered_after_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).sync_all_err();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.sync_all());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sync_data_ordered_after_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).sync_data();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.sync_data());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn sync_data_err_ordered_after_write() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.write(HELLO).sync_data_err();
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.write(HELLO));
|
||||
assert_ready_ok!(t.poll());
|
||||
|
||||
let mut t = task::spawn(file.sync_data());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_pending!(t.poll());
|
||||
|
||||
assert_eq!(1, pool::len());
|
||||
pool::run_one();
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_set_len_ok() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.set_len(123);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.set_len(123));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
assert_eq!(1, mock.remaining());
|
||||
|
||||
pool::run_one();
|
||||
assert_eq!(0, mock.remaining());
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn open_set_len_err() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.set_len_err(123);
|
||||
|
||||
let mut file = File::from_std(file);
|
||||
let mut t = task::spawn(file.set_len(123));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
assert_eq!(1, mock.remaining());
|
||||
|
||||
pool::run_one();
|
||||
assert_eq!(0, mock.remaining());
|
||||
|
||||
assert!(t.is_woken());
|
||||
assert_ready_err!(t.poll());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn partial_read_set_len_ok() {
|
||||
let (mock, file) = sys::File::mock();
|
||||
mock.read(HELLO)
|
||||
.seek_current_ok(-14, 0)
|
||||
.set_len(123)
|
||||
.read(FOO);
|
||||
|
||||
let mut buf = [0; 32];
|
||||
let mut file = File::from_std(file);
|
||||
|
||||
{
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
}
|
||||
|
||||
pool::run_one();
|
||||
|
||||
{
|
||||
let mut t = task::spawn(file.set_len(123));
|
||||
|
||||
assert_pending!(t.poll());
|
||||
pool::run_one();
|
||||
assert_ready_ok!(t.poll());
|
||||
}
|
||||
|
||||
let mut t = task::spawn(file.read(&mut buf));
|
||||
assert_pending!(t.poll());
|
||||
pool::run_one();
|
||||
let n = assert_ready_ok!(t.poll());
|
||||
|
||||
assert_eq!(n, FOO.len());
|
||||
assert_eq!(&buf[..n], FOO);
|
||||
}
|
||||
@@ -0,0 +1,69 @@
|
||||
#![warn(rust_2018_idioms)]
|
||||
|
||||
use tokio::fs;
|
||||
|
||||
use std::io::prelude::*;
|
||||
use std::io::BufReader;
|
||||
use tempfile::tempdir;
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_hard_link() {
|
||||
let dir = tempdir().unwrap();
|
||||
let src = dir.path().join("src.txt");
|
||||
let dst = dir.path().join("dst.txt");
|
||||
|
||||
{
|
||||
let mut file = std::fs::File::create(&src).unwrap();
|
||||
file.write_all(b"hello").unwrap();
|
||||
}
|
||||
|
||||
let dst_2 = dst.clone();
|
||||
|
||||
assert!(fs::hard_link(src, dst_2.clone()).await.is_ok());
|
||||
|
||||
let mut content = String::new();
|
||||
|
||||
{
|
||||
let file = std::fs::File::open(dst).unwrap();
|
||||
let mut reader = BufReader::new(file);
|
||||
reader.read_to_string(&mut content).unwrap();
|
||||
}
|
||||
|
||||
assert!(content == "hello");
|
||||
}
|
||||
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn test_symlink() {
|
||||
let dir = tempdir().unwrap();
|
||||
let src = dir.path().join("src.txt");
|
||||
let dst = dir.path().join("dst.txt");
|
||||
|
||||
{
|
||||
let mut file = std::fs::File::create(&src).unwrap();
|
||||
file.write_all(b"hello").unwrap();
|
||||
}
|
||||
|
||||
let src_2 = src.clone();
|
||||
let dst_2 = dst.clone();
|
||||
|
||||
assert!(fs::os::unix::symlink(src_2.clone(), dst_2.clone())
|
||||
.await
|
||||
.is_ok());
|
||||
|
||||
let mut content = String::new();
|
||||
|
||||
{
|
||||
let file = std::fs::File::open(dst.clone()).unwrap();
|
||||
let mut reader = BufReader::new(file);
|
||||
reader.read_to_string(&mut content).unwrap();
|
||||
}
|
||||
|
||||
assert!(content == "hello");
|
||||
|
||||
let read = fs::read_link(dst.clone()).await.unwrap();
|
||||
assert!(read == src);
|
||||
|
||||
let symlink_meta = fs::symlink_metadata(dst.clone()).await.unwrap();
|
||||
assert!(symlink_meta.file_type().is_symlink());
|
||||
}
|
||||
@@ -0,0 +1,265 @@
|
||||
use std::collections::VecDeque;
|
||||
use std::fmt;
|
||||
use std::fs::{Metadata, Permissions};
|
||||
use std::io;
|
||||
use std::io::prelude::*;
|
||||
use std::io::SeekFrom;
|
||||
use std::path::PathBuf;
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
pub struct File {
|
||||
shared: Arc<Mutex<Shared>>,
|
||||
}
|
||||
|
||||
pub struct Handle {
|
||||
shared: Arc<Mutex<Shared>>,
|
||||
}
|
||||
|
||||
struct Shared {
|
||||
calls: VecDeque<Call>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
enum Call {
|
||||
Read(io::Result<Vec<u8>>),
|
||||
Write(io::Result<Vec<u8>>),
|
||||
Seek(SeekFrom, io::Result<u64>),
|
||||
SyncAll(io::Result<()>),
|
||||
SyncData(io::Result<()>),
|
||||
SetLen(u64, io::Result<()>),
|
||||
}
|
||||
|
||||
impl Handle {
|
||||
pub fn read(&self, data: &[u8]) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls.push_back(Call::Read(Ok(data.to_owned())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn read_err(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::Read(Err(io::ErrorKind::Other.into())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn write(&self, data: &[u8]) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls.push_back(Call::Write(Ok(data.to_owned())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn write_err(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::Write(Err(io::ErrorKind::Other.into())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn seek_start_ok(&self, offset: u64) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::Seek(SeekFrom::Start(offset), Ok(offset)));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn seek_current_ok(&self, offset: i64, ret: u64) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::Seek(SeekFrom::Current(offset), Ok(ret)));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn sync_all(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls.push_back(Call::SyncAll(Ok(())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn sync_all_err(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::SyncAll(Err(io::ErrorKind::Other.into())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn sync_data(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls.push_back(Call::SyncData(Ok(())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn sync_data_err(&self) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::SyncData(Err(io::ErrorKind::Other.into())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn set_len(&self, size: u64) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls.push_back(Call::SetLen(size, Ok(())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn set_len_err(&self, size: u64) -> &Self {
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
s.calls
|
||||
.push_back(Call::SetLen(size, Err(io::ErrorKind::Other.into())));
|
||||
self
|
||||
}
|
||||
|
||||
pub fn remaining(&self) -> usize {
|
||||
let s = self.shared.lock().unwrap();
|
||||
s.calls.len()
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for Handle {
|
||||
fn drop(&mut self) {
|
||||
if !std::thread::panicking() {
|
||||
let s = self.shared.lock().unwrap();
|
||||
assert_eq!(0, s.calls.len());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl File {
|
||||
pub fn open(_: PathBuf) -> io::Result<File> {
|
||||
unimplemented!();
|
||||
}
|
||||
|
||||
pub fn create(_: PathBuf) -> io::Result<File> {
|
||||
unimplemented!();
|
||||
}
|
||||
|
||||
pub fn mock() -> (Handle, File) {
|
||||
let shared = Arc::new(Mutex::new(Shared {
|
||||
calls: VecDeque::new(),
|
||||
}));
|
||||
|
||||
let handle = Handle {
|
||||
shared: shared.clone(),
|
||||
};
|
||||
let file = File { shared };
|
||||
|
||||
(handle, file)
|
||||
}
|
||||
|
||||
pub fn sync_all(&self) -> io::Result<()> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(SyncAll(ret)) => ret,
|
||||
Some(op) => panic!("expected next call to be {:?}; was sync_all", op),
|
||||
None => panic!("did not expect call"),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn sync_data(&self) -> io::Result<()> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(SyncData(ret)) => ret,
|
||||
Some(op) => panic!("expected next call to be {:?}; was sync_all", op),
|
||||
None => panic!("did not expect call"),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn set_len(&self, size: u64) -> io::Result<()> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(SetLen(arg, ret)) => {
|
||||
assert_eq!(arg, size);
|
||||
ret
|
||||
}
|
||||
Some(op) => panic!("expected next call to be {:?}; was sync_all", op),
|
||||
None => panic!("did not expect call"),
|
||||
}
|
||||
}
|
||||
|
||||
pub fn metadata(&self) -> io::Result<Metadata> {
|
||||
unimplemented!();
|
||||
}
|
||||
|
||||
pub fn set_permissions(&self, _perm: Permissions) -> io::Result<()> {
|
||||
unimplemented!();
|
||||
}
|
||||
|
||||
pub fn try_clone(&self) -> io::Result<Self> {
|
||||
unimplemented!();
|
||||
}
|
||||
}
|
||||
|
||||
impl Read for &'_ File {
|
||||
fn read(&mut self, dst: &mut [u8]) -> io::Result<usize> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(Read(Ok(data))) => {
|
||||
assert!(dst.len() >= data.len());
|
||||
assert!(dst.len() <= 16 * 1024, "actual = {}", dst.len()); // max buffer
|
||||
|
||||
&mut dst[..data.len()].copy_from_slice(&data);
|
||||
Ok(data.len())
|
||||
}
|
||||
Some(Read(Err(e))) => Err(e),
|
||||
Some(op) => panic!("expected next call to be {:?}; was a read", op),
|
||||
None => panic!("did not expect call"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Write for &'_ File {
|
||||
fn write(&mut self, src: &[u8]) -> io::Result<usize> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(Write(Ok(data))) => {
|
||||
assert_eq!(src, &data[..]);
|
||||
Ok(src.len())
|
||||
}
|
||||
Some(Write(Err(e))) => Err(e),
|
||||
Some(op) => panic!("expected next call to be {:?}; was write", op),
|
||||
None => panic!("did not expect call"),
|
||||
}
|
||||
}
|
||||
|
||||
fn flush(&mut self) -> io::Result<()> {
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
impl Seek for &'_ File {
|
||||
fn seek(&mut self, pos: SeekFrom) -> io::Result<u64> {
|
||||
use self::Call::*;
|
||||
|
||||
let mut s = self.shared.lock().unwrap();
|
||||
|
||||
match s.calls.pop_front() {
|
||||
Some(Seek(expect, res)) => {
|
||||
assert_eq!(expect, pos);
|
||||
res
|
||||
}
|
||||
Some(op) => panic!("expected call {:?}; was `seek`", op),
|
||||
None => panic!("did not expect call; was `seek`"),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Debug for File {
|
||||
fn fmt(&self, fmt: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
fmt.debug_struct("mock::File").finish()
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,66 @@
|
||||
use tokio_sync::oneshot;
|
||||
|
||||
use std::cell::RefCell;
|
||||
use std::collections::VecDeque;
|
||||
use std::future::Future;
|
||||
use std::io;
|
||||
use std::pin::Pin;
|
||||
use std::task::{Context, Poll};
|
||||
|
||||
thread_local! {
|
||||
static QUEUE: RefCell<VecDeque<Box<dyn FnOnce() + Send>>> = RefCell::new(VecDeque::new())
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct Blocking<T> {
|
||||
rx: oneshot::Receiver<T>,
|
||||
}
|
||||
|
||||
pub(crate) fn run<F, R>(f: F) -> Blocking<R>
|
||||
where
|
||||
F: FnOnce() -> R + Send + 'static,
|
||||
R: Send + 'static,
|
||||
{
|
||||
let (tx, rx) = oneshot::channel();
|
||||
let task = Box::new(move || {
|
||||
let _ = tx.send(f());
|
||||
});
|
||||
|
||||
QUEUE.with(|cell| cell.borrow_mut().push_back(task));
|
||||
|
||||
Blocking { rx }
|
||||
}
|
||||
|
||||
impl<T> Future for Blocking<T> {
|
||||
type Output = T;
|
||||
|
||||
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
|
||||
use std::task::Poll::*;
|
||||
|
||||
match Pin::new(&mut self.rx).poll(cx) {
|
||||
Ready(Ok(v)) => Ready(v),
|
||||
Ready(Err(e)) => panic!("error = {:?}", e),
|
||||
Pending => Pending,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn asyncify<F, T>(f: F) -> io::Result<T>
|
||||
where
|
||||
F: FnOnce() -> io::Result<T> + Send + 'static,
|
||||
T: Send + 'static,
|
||||
{
|
||||
run(f).await
|
||||
}
|
||||
|
||||
pub(crate) fn len() -> usize {
|
||||
QUEUE.with(|cell| cell.borrow().len())
|
||||
}
|
||||
|
||||
pub(crate) fn run_one() {
|
||||
let task = QUEUE
|
||||
.with(|cell| cell.borrow_mut().pop_front())
|
||||
.expect("expected task to run, but none ready");
|
||||
|
||||
task();
|
||||
}
|
||||
Reference in New Issue
Block a user