use futures::Future;
use futures::future::{Either, ok};
use futures_cpupool::CpuPool;
use std::ops::DerefMut;
use std::io::{self, Read};
use std::mem;
use common::EXP_POOL;
pub struct BufReader<R> {
inner: R,
buf: Box<[u8]>,
pos: usize,
cap: usize,
pool: Option<CpuPool>,
}
type OkRead<R, B> = (R, B, usize);
type ErrRead<R, B> = (R, B, io::Error);
impl<R: Read + Send + 'static> BufReader<R> {
pub fn with_pool_and_capacity(pool: CpuPool, cap: usize, inner: R) -> BufReader<R> {
let mut buf = Vec::with_capacity(cap);
unsafe {
buf.set_len(cap);
}
BufReader::with_pool_and_buf(pool, buf.into_boxed_slice(), inner)
}
pub fn with_pool_and_buf(pool: CpuPool, buf: Box<[u8]>, inner: R) -> BufReader<R> {
let cap = buf.len();
BufReader {
inner: inner,
buf: buf,
pos: cap,
cap: cap,
pool: Some(pool),
}
}
pub fn get_ref(&self) -> &R {
&self.inner
}
pub fn get_mut(&mut self) -> &R {
&mut self.inner
}
pub unsafe fn set_pos(&mut self, pos: usize) {
self.pos = pos;
}
pub unsafe fn components(mut self) -> (R, Box<[u8]>, CpuPool) {
let r = 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);
(r, buf, pool)
}
pub fn try_read_full<B>(
mut self,
mut buf: B,
) -> impl Future<Item = OkRead<Self, B>, Error = ErrRead<Self, B>>
where
B: DerefMut<Target = [u8]> + Send + 'static,
{
const U8READ: &str = "&[u8] reads never error";
let mut rem = buf.len();
let mut at = 0;
if self.pos != self.cap {
at = (&self.buf[self.pos..self.cap]).read(&mut buf).expect(
U8READ,
);
rem -= at;
self.pos += at;
if rem == 0 {
return Either::A(ok::<OkRead<Self, B>, ErrRead<Self, B>>((self, buf, at)));
}
}
let pool = self.pool.take().expect(EXP_POOL);
let block = if self.cap > 0 {
rem - rem % self.cap
} else {
rem
};
let fut = pool.spawn_fn(move || {
if block > 0 {
let (block_read, err) = try_read_full(&mut self.inner, &mut buf[at..at + block]);
if let Some(e) = err {
return Err((self, buf, e));
}
at += block_read;
rem -= block_read;
if rem == 0 {
return Ok((self, buf, at));
}
}
let (buf_read, err) = try_read_full(&mut self.inner, &mut self.buf);
match err {
Some(e) => Err((self, buf, e)),
None => {
self.cap = buf_read;
self.pos = (&self.buf[..self.cap]).read(&mut buf[at..]).expect(U8READ);
at += self.pos;
Ok((self, buf, at))
}
}
});
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)
}
}))
}
}
fn try_read_full<R: Read>(r: &mut R, mut buf: &mut [u8]) -> (usize, Option<io::Error>) {
let mut nn: usize = 0;
while !buf.is_empty() {
match r.read(buf) {
Ok(0) => break,
Ok(n) => {
let tmp = buf;
buf = &mut tmp[n..];
nn += n;
}
Err(ref e) if e.kind() == io::ErrorKind::Interrupted => {}
Err(e) => return (nn, Some(e)),
}
}
(nn, None)
}
#[test]
fn test_read() {
use std::fs;
use std::io::Write;
fs::OpenOptions::new()
.write(true)
.create_new(true)
.open("bar.txt")
.expect("unable to exclusively create foo")
.write_all(
"Strapped down to my bed, feet cold, eyes red. I'm out of my head. Am I alive? Am I \
dead?"
.as_bytes(),
)
.expect("unable to write all");
let f = BufReader::with_pool_and_capacity(
CpuPool::new(1),
10,
fs::OpenOptions::new().read(true).open("bar.txt").expect(
"foo does not exist?",
),
);
assert_eq!(f.pos, 10);
let (f, buf, n) = f.try_read_full(vec![0; 5]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 5);
assert_eq!(&*buf, b"Strap");
assert_eq!(f.pos, 5);
assert_eq!(f.cap, 10);
assert_eq!(&*f.buf, b"Strapped d");
let (f, buf, n) = f.try_read_full(vec![0; 2]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 2);
assert_eq!(&*buf, b"pe");
assert_eq!(f.pos, 7);
assert_eq!(f.cap, 10);
assert_eq!(&*f.buf, b"Strapped d");
let (f, buf, n) = f.try_read_full(vec![0; 25]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 25);
assert_eq!(&*buf, b"d down to my bed, feet co");
assert_eq!(f.pos, 2);
assert_eq!(f.cap, 10);
assert_eq!(&*f.buf, b"cold, eyes");
let (f, buf, n) = f.try_read_full(vec![0; 18]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 18);
assert_eq!(&*buf, b"ld, eyes red. I'm ");
assert_eq!(f.pos, 10);
assert_eq!(f.cap, 10);
assert_eq!(&*f.buf, b"cold, eyes");
let (f, buf, n) = f.try_read_full(vec![0; 10]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 10);
assert_eq!(&*buf, b"out of my ");
assert_eq!(f.pos, 10);
assert_eq!(f.cap, 10);
assert_eq!(&*f.buf, b"cold, eyes");
let (f, buf, n) = f.try_read_full(vec![0; 29]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 28);
assert_eq!(&*buf, b"head. Am I alive? Am I dead?\0");
assert_eq!(f.pos, 8);
assert_eq!(f.cap, 8);
assert_eq!(&*f.buf, b" I dead?es");
let (f, buf, n) = f.try_read_full(vec![0; 2]).wait().unwrap_or_else(
|(_, _, e)| {
panic!("unable to read: {}", e)
},
);
assert_eq!(n, 0);
assert_eq!(&*buf, b"\0\0");
assert_eq!(f.pos, 0);
assert_eq!(f.cap, 0);
assert_eq!(&*f.buf, b" I dead?es");
fs::remove_file("bar.txt").expect("expected file to be removed");
}