use futures::Future;
use futures::future::{Either, ok};
use futures_cpupool::CpuPool;
use std::io::{self, Write};
use std::mem;
use std::ops::Deref;
use common::*;
pub struct BufWriter<W> {
inner: W,
buf: Box<[u8]>,
pos: usize,
w_start: usize,
pool: Option<CpuPool>,
}
type OkWrite<W, B> = (W, B);
type ErrWrite<W, B> = (W, B, io::Error);
impl<W: Write + Send + 'static> BufWriter<W> {
pub fn with_pool_and_capacity(pool: CpuPool, cap: usize, inner: W) -> BufWriter<W> {
let mut buf = Vec::with_capacity(cap);
unsafe {
buf.set_len(cap);
}
BufWriter::with_pool_and_buf(pool, buf.into_boxed_slice(), inner)
}
pub fn with_pool_and_buf(pool: CpuPool, buf: Box<[u8]>, inner: W) -> BufWriter<W> {
BufWriter {
inner: inner,
buf: buf,
pos: 0,
w_start: 0,
pool: Some(pool),
}
}
pub fn get_ref(&self) -> &W {
&self.inner
}
pub fn get_mut(&mut self) -> &W {
&mut self.inner
}
pub unsafe fn components(mut self) -> (W, Box<[u8]>, CpuPool) {
let w = mem::replace(&mut self.inner, mem::uninitialized());
let buf = mem::replace(&mut self.buf, mem::uninitialized());
let mut pool = mem::replace(&mut self.pool, mem::uninitialized());
let pool = pool.take().expect(EXP_POOL);
mem::forget(self);
(w, buf, pool)
}
pub unsafe fn set_pos(&mut self, pos: usize) {
self.pos = pos;
}
pub fn write_all<B>(
mut self,
buf: B,
) -> impl Future<Item = OkWrite<Self, B>, Error = ErrWrite<Self, B>>
where
B: Deref<Target = [u8]> + Send + 'static,
{
let mut rem = buf.len();
let mut at = 0;
let mut write_buf = false;
if self.pos == 0 {
if buf.len() < self.buf.len() {
self.pos = copy(&mut self.buf, &*buf);
return Either::A(ok::<OkWrite<Self, B>, ErrWrite<Self, B>>((self, buf)));
}
} else {
at = copy(&mut self.buf[self.pos..], &*buf);
self.pos += at;
rem -= at;
if self.pos != self.buf.len() {
return Either::A(ok::<OkWrite<Self, B>, ErrWrite<Self, B>>((self, buf)));
}
write_buf = true;
}
let pool = self.pool.take().expect(EXP_POOL);
let fut = pool.spawn_fn(move || {
if write_buf {
if let Err(e) = self.inner.write_all(&self.buf[self.w_start..]) {
return Err((self, buf, e));
}
self.w_start = 0;
}
if rem >= self.buf.len() {
let n_write = rem -
if self.buf.len() != 0 {
rem % self.buf.len()
} else {
0
};
if let Err(e) = self.inner.write_all(&buf[at..at + n_write]) {
return Err((self, buf, e));
}
at += n_write;
rem -= n_write;
}
self.pos = copy(&mut self.buf, &buf[at..]);
Ok((self, buf))
});
Either::B(fut.then(|res| match res {
Ok(mut x) => {
x.0.pool = Some(pool);
Ok(x)
}
Err(mut x) => {
x.0.pool = Some(pool);
Err(x)
}
}))
}
pub fn flush_buf(mut self) -> impl Future<Item = Self, Error = (Self, io::Error)> {
if self.w_start == self.pos {
return Either::A(ok::<Self, (Self, io::Error)>(self));
}
let pool = self.pool.take().expect(EXP_POOL);
let fut = pool.spawn_fn(move || {
if let Err(e) = self.inner.write_all(&self.buf[self.w_start..self.pos]) {
return Err((self, e));
}
self.w_start = self.pos;
Ok(self)
});
Either::B(fut.then(|res| match res {
Ok(mut me) => {
me.pool = Some(pool);
Ok(me)
}
Err(mut x) => {
x.0.pool = Some(pool);
Err(x)
}
}))
}
pub fn flush_inner(mut self) -> impl Future<Item = Self, Error = (Self, io::Error)> {
let pool = self.pool.take().expect(EXP_POOL);
let fut = pool.spawn_fn(move || {
if let Err(e) = self.inner.flush() {
return Err((self, e));
}
Ok(self)
});
fut.then(|res| match res {
Ok(mut me) => {
me.pool = Some(pool);
Ok(me)
}
Err(mut x) => {
x.0.pool = Some(pool);
Err(x)
}
})
}
}
#[test]
fn test_write() {
use std::fs;
use std::io::Read;
fn assert_foo(exp: &'static str) {
let mut foo = fs::File::open("foo.txt").expect("re-open");
let mut contents = String::new();
foo.read_to_string(&mut contents).expect(
"unable to read file",
);
assert_eq!(contents, exp);
}
let f = BufWriter::with_pool_and_capacity(
CpuPool::new(1),
10,
fs::OpenOptions::new()
.write(true)
.read(true) .create_new(true)
.open("foo.txt")
.expect("foo not exclusively created?"),
);
let (f, buf) = f.write_all(b"hello".to_vec()).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to write to file: {}", e)
},
);
assert_eq!(f.pos, 5);
assert_eq!(f.w_start, 0);
assert_eq!(&*buf, b"hello");
assert_foo("");
let f = f.flush_buf().wait().unwrap_or_else(|(_, e)| {
panic!("unable to flush buf: {}", e)
});
assert_eq!(f.pos, 5);
assert_eq!(f.w_start, 5);
assert_foo("hello");
let f = f.flush_inner().wait().unwrap_or_else(|(_, e)| {
panic!("unable to flush file: {}", e)
});
assert_eq!(f.pos, 5);
assert_eq!(f.w_start, 5);
assert_foo("hello");
let (f, buf) = f.write_all(b"tw".to_vec()).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to write to file: {}", e)
},
);
assert_eq!(f.pos, 7);
assert_eq!(f.w_start, 5);
assert_eq!(&*buf, b"tw");
assert_foo("hello");
let f = f.flush_buf().wait().unwrap_or_else(|(_, e)| {
panic!("unable to flush buf: {}", e)
});
assert_eq!(f.pos, 7);
assert_eq!(f.w_start, 7);
assert_foo("hellotw");
let (f, buf) = f.write_all(b"goodbye".to_vec()).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to write to file: {}", e)
},
);
assert_eq!(f.pos, 4);
assert_eq!(f.w_start, 0);
assert_eq!(&*buf, b"goodbye");
assert_foo("hellotwgoo");
let (f, buf) = f.write_all(b"more++andthenten".to_vec())
.wait()
.unwrap_or_else(|(_, _, e)| panic!("unable to write to file: {}", e));
assert_eq!(f.pos, 0);
assert_eq!(f.w_start, 0);
assert_eq!(&*buf, b"more++andthenten");
assert_foo("hellotwgoodbyemore++andthenten");
let (f, buf) = f.write_all(b"andtenmore".to_vec()).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to write to file: {}", e)
},
);
assert_eq!(f.pos, 0);
assert_eq!(f.w_start, 0);
assert_eq!(&*buf, b"andtenmore");
assert_foo("hellotwgoodbyemore++andthentenandtenmore");
let (f, buf) = f.write_all(b"this is rly old".to_vec())
.wait()
.unwrap_or_else(|(_, _, e)| panic!("unable to write to file: {}", e));
assert_eq!(f.pos, 5);
assert_eq!(f.w_start, 0);
assert_eq!(&*buf, b"this is rly old");
assert_foo("hellotwgoodbyemore++andthentenandtenmorethis is rl");
let f = f.flush_buf().wait().unwrap_or_else(|(_, e)| {
panic!("unable to flush buf: {}", e)
});
assert_eq!(f.pos, 5);
assert_eq!(f.w_start, 5);
assert_foo("hellotwgoodbyemore++andthentenandtenmorethis is rly old");
fs::remove_file("foo.txt").expect("expected file to be removed");
}