extern crate aio_bindings;
extern crate futures;
extern crate futures_cpupool;
extern crate libc;
extern crate memmap;
extern crate mio;
extern crate rand;
extern crate tokio;
use std::convert;
use std::error;
use std::fmt;
use std::io;
use std::mem;
use std::ops;
use std::ptr;
use std::os::unix::io::RawFd;
use libc::{c_long, c_void, mlock};
use futures::Future;
use ops::Deref;
mod aio;
mod eventfd;
mod sync;
struct IocbInfo {
opcode: u32,
fd: RawFd,
offset: u64,
buf: u64,
len: u64,
flags: u32,
}
#[derive(Debug)]
struct RequestState {
request: aio::iocb,
completed_receiver: futures::sync::oneshot::Receiver<c_long>,
completed_sender: Option<futures::sync::oneshot::Sender<c_long>>,
}
struct AioBaseFuture {
context: std::sync::Arc<AioContextInner>,
iocb_info: IocbInfo,
state: Option<Box<RequestState>>,
acquire_state: Option<sync::SemaphoreHandle>,
}
impl AioBaseFuture {
fn submit_request(&mut self) -> Result<futures::Async<()>, io::Error> {
if self.state.is_none() {
if self.acquire_state.is_none() {
self.acquire_state = Some(self.context.have_capacity.acquire());
}
match self.acquire_state.as_mut().unwrap().poll() {
Err(err) => return Err(err),
Ok(futures::Async::NotReady) => return Ok(futures::Async::NotReady),
Ok(futures::Async::Ready(_)) => {
let mut guard = self.context.capacity.write();
match guard {
Ok(ref mut guard) => {
self.state = guard.state.pop();
}
Err(_) => panic!("TODO: Figure out how to handle this kind of error"),
}
}
}
assert!(self.state.is_some());
let state = self.state.as_mut().unwrap();
let state_addr = state.deref().deref() as *const RequestState;
state.request.aio_data = unsafe { mem::transmute(state_addr) };
state.request.aio_resfd = self.context.completed_fd as u32;
state.request.aio_flags = aio::IOCB_FLAG_RESFD | self.iocb_info.flags;
state.request.aio_fildes = self.iocb_info.fd as u32;
state.request.aio_offset = self.iocb_info.offset as i64;
state.request.aio_buf = self.iocb_info.buf;
state.request.aio_nbytes = self.iocb_info.len;
state.request.aio_lio_opcode = self.iocb_info.opcode as u16;
let (sender, receiver) = futures::sync::oneshot::channel();
state.completed_receiver = receiver;
state.completed_sender = Some(sender);
let mut request_ptr_array: [*mut aio::iocb; 1] =
[&mut state.request as *mut aio::iocb; 1];
let result = unsafe {
aio::io_submit(
self.context.context,
1,
&mut request_ptr_array[0] as *mut *mut aio::iocb,
)
};
if result != 1 {
return Err(io::Error::last_os_error());
}
}
Ok(futures::Async::Ready(()))
}
fn retrieve_result(&mut self) -> Result<futures::Async<()>, io::Error> {
let result_code = match self.state.as_mut().unwrap().completed_receiver.poll() {
Err(err) => return Err(io::Error::new(io::ErrorKind::Other, err)),
Ok(futures::Async::NotReady) => return Ok(futures::Async::NotReady),
Ok(futures::Async::Ready(n)) => n,
};
match self.context.capacity.write() {
Ok(ref mut guard) => {
guard.state.push(self.state.take().unwrap());
}
Err(_) => panic!("TODO: Figure out how to handle this kind of error"),
}
self.context.have_capacity.release();
if result_code < 0 {
Err(io::Error::from_raw_os_error(result_code as i32))
} else {
Ok(futures::Async::Ready(()))
}
}
}
impl futures::Future for AioBaseFuture {
type Item = ();
type Error = io::Error;
fn poll(&mut self) -> Result<futures::Async<()>, io::Error> {
let result = self.submit_request();
match result {
Ok(futures::Async::Ready(())) => self.retrieve_result(),
Ok(futures::Async::NotReady) => Ok(futures::Async::NotReady),
Err(err) => Err(err),
}
}
}
pub struct AioError<Handle> {
pub buffer: Handle,
pub error: io::Error,
}
impl<Handle> fmt::Debug for AioError<Handle> {
fn fmt(&self, f: &mut fmt::Formatter) -> Result<(), fmt::Error> {
self.error.fmt(f)
}
}
impl<Handle> fmt::Display for AioError<Handle> {
fn fmt(&self, f: &mut fmt::Formatter) -> Result<(), fmt::Error> {
self.error.fmt(f)
}
}
impl<Handle> error::Error for AioError<Handle> {
fn description(&self) -> &str {
self.error.description()
}
fn cause(&self) -> Option<&error::Error> {
self.error.cause()
}
}
pub struct AioReadResultFuture<ReadWriteHandle>
where
ReadWriteHandle: convert::AsMut<[u8]>,
{
base: AioBaseFuture,
buffer: Option<ReadWriteHandle>,
}
impl<ReadWriteHandle> futures::Future for AioReadResultFuture<ReadWriteHandle>
where
ReadWriteHandle: convert::AsMut<[u8]>,
{
type Item = ReadWriteHandle;
type Error = AioError<ReadWriteHandle>;
fn poll(&mut self) -> Result<futures::Async<Self::Item>, Self::Error> {
self.base
.poll()
.map(|val| val.map(|_| self.buffer.take().unwrap()))
.map_err(|err| AioError {
buffer: self.buffer.take().unwrap(),
error: err,
})
}
}
pub struct AioWriteResultFuture<ReadOnlyHandle>
where
ReadOnlyHandle: convert::AsRef<[u8]>,
{
base: AioBaseFuture,
buffer: Option<ReadOnlyHandle>,
}
impl<ReadOnlyHandle> futures::Future for AioWriteResultFuture<ReadOnlyHandle>
where
ReadOnlyHandle: convert::AsRef<[u8]>,
{
type Item = ReadOnlyHandle;
type Error = AioError<ReadOnlyHandle>;
fn poll(&mut self) -> Result<futures::Async<Self::Item>, Self::Error> {
self.base
.poll()
.map(|val| val.map(|_| self.buffer.take().unwrap()))
.map_err(|err| AioError {
buffer: self.buffer.take().unwrap(),
error: err,
})
}
}
pub struct AioSyncResultFuture
{
base: AioBaseFuture,
}
impl futures::Future for AioSyncResultFuture
{
type Item = ();
type Error = io::Error;
fn poll(&mut self) -> Result<futures::Async<Self::Item>, Self::Error> {
self.base.poll()
}
}
pub struct AioPollFuture {
context: aio::aio_context_t,
eventfd: eventfd::EventFd,
events: Vec<aio::io_event>,
}
impl futures::Future for AioPollFuture {
type Item = ();
type Error = io::Error;
fn poll(&mut self) -> Result<futures::Async<Self::Item>, Self::Error> {
loop {
let available = match self.eventfd.read() {
Err(err) => return Err(err),
Ok(futures::Async::NotReady) => return Ok(futures::Async::NotReady),
Ok(futures::Async::Ready(value)) => value as usize,
};
assert!(available > 0);
self.events.clear();
unsafe {
let result = aio::io_getevents(
self.context,
available as c_long,
available as c_long,
self.events.as_mut_ptr(),
ptr::null_mut::<aio::timespec>(),
);
if result < 0 {
return Err(io::Error::last_os_error());
}
assert!(result as usize == available);
self.events.set_len(available);
};
for ref event in &self.events {
let request_state: &mut RequestState = unsafe { mem::transmute(event.data) };
request_state
.completed_sender
.take()
.unwrap()
.send(event.res)
.unwrap();
}
}
}
}
#[derive(Debug)]
struct Capacity {
state: Vec<Box<RequestState>>,
}
impl Capacity {
fn new(nr: usize) -> Result<Capacity, io::Error> {
let mut state = Vec::with_capacity(nr);
for _ in 0..nr {
let (_, receiver) = futures::sync::oneshot::channel();
state.push(Box::new(RequestState {
request: unsafe { mem::zeroed() },
completed_receiver: receiver,
completed_sender: None,
}));
}
Ok(Capacity { state })
}
}
#[derive(Debug)]
struct AioContextInner {
context: aio::aio_context_t,
completed_fd: RawFd,
have_capacity: sync::Semaphore,
capacity: std::sync::RwLock<Capacity>,
poll_task_handle: Option<futures::sync::oneshot::SpawnHandle<(), io::Error>>,
}
impl AioContextInner {
fn new(fd: RawFd, nr: usize) -> Result<AioContextInner, io::Error> {
let mut context: aio::aio_context_t = 0;
unsafe {
if aio::io_setup(nr as c_long, &mut context) != 0 {
return Err(io::Error::last_os_error());
}
};
Ok(AioContextInner {
context,
capacity: std::sync::RwLock::new(Capacity::new(nr)?),
have_capacity: sync::Semaphore::new(nr),
completed_fd: fd,
poll_task_handle: None,
})
}
}
impl Drop for AioContextInner {
fn drop(&mut self) {
let result = unsafe { aio::io_destroy(self.context) };
assert!(result == 0);
}
}
#[derive(Clone, Debug)]
pub struct AioContext {
inner: std::sync::Arc<AioContextInner>,
}
#[derive(Copy, Clone, Debug)]
pub enum SyncLevel {
None = 0,
Data = aio::RWF_DSYNC as isize,
Full = aio::RWF_SYNC as isize,
}
impl AioContext {
pub fn new<E>(executor: &E, nr: usize) -> Result<AioContext, io::Error>
where
E: futures::future::Executor<futures::sync::oneshot::Execute<AioPollFuture>>,
{
let eventfd = eventfd::EventFd::create(0, false)?;
let fd = eventfd.evented.get_ref().fd;
let mut inner = AioContextInner::new(fd, nr)?;
let context = inner.context;
let poll_future = AioPollFuture {
context,
eventfd,
events: Vec::with_capacity(nr),
};
inner.poll_task_handle = Some(futures::sync::oneshot::spawn(poll_future, executor));
Ok(AioContext {
inner: std::sync::Arc::new(inner),
})
}
pub fn read<ReadWriteHandle>(
&self,
fd: RawFd,
offset: u64,
mut buffer_obj: ReadWriteHandle,
) -> AioReadResultFuture<ReadWriteHandle>
where
ReadWriteHandle: convert::AsMut<[u8]>,
{
let (ptr, len) = {
let buffer = buffer_obj.as_mut();
let len = buffer.len() as u64;
let ptr = unsafe { mem::transmute(buffer.as_ptr()) };
(ptr, len)
};
AioReadResultFuture {
base: AioBaseFuture {
context: self.inner.clone(),
iocb_info: IocbInfo {
opcode: aio::IOCB_CMD_PREAD,
fd,
offset,
len,
buf: ptr,
flags: 0,
},
state: None,
acquire_state: None,
},
buffer: Some(buffer_obj),
}
}
pub fn write<ReadOnlyHandle>(
&self,
fd: RawFd,
offset: u64,
buffer: ReadOnlyHandle,
) -> AioWriteResultFuture<ReadOnlyHandle>
where
ReadOnlyHandle: convert::AsRef<[u8]>,
{
self.write_sync(fd, offset, buffer, SyncLevel::None)
}
pub fn write_sync<ReadOnlyHandle>(
&self,
fd: RawFd,
offset: u64,
buffer_obj: ReadOnlyHandle,
sync_level: SyncLevel
) -> AioWriteResultFuture<ReadOnlyHandle>
where
ReadOnlyHandle: convert::AsRef<[u8]>,
{
let (ptr, len) = {
let buffer = buffer_obj.as_ref();
let len = buffer.len() as u64;
let ptr = unsafe { mem::transmute(buffer.as_ptr()) };
(ptr, len)
};
AioWriteResultFuture {
base: AioBaseFuture {
context: self.inner.clone(),
iocb_info: IocbInfo {
opcode: aio::IOCB_CMD_PWRITE,
fd,
offset,
len,
buf: ptr,
flags: sync_level as u32,
},
state: None,
acquire_state: None,
},
buffer: Some(buffer_obj),
}
}
pub fn sync(
&self,
fd: RawFd,
) -> AioSyncResultFuture
{
AioSyncResultFuture {
base: AioBaseFuture {
context: self.inner.clone(),
iocb_info: IocbInfo {
opcode: aio::IOCB_CMD_FSYNC,
fd,
buf: 0,
len: 0,
offset: 0,
flags: 0,
},
state: None,
acquire_state: None,
},
}
}
pub fn data_sync(
&self,
fd: RawFd,
) -> AioSyncResultFuture
{
AioSyncResultFuture {
base: AioBaseFuture {
context: self.inner.clone(),
iocb_info: IocbInfo {
opcode: aio::IOCB_CMD_FDSYNC,
fd,
buf: 0,
len: 0,
offset: 0,
flags: 0,
},
state: None,
acquire_state: None,
},
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::borrow::{Borrow, BorrowMut};
use std::env;
use std::fs;
use std::io::Write;
use std::os::unix::ffi::OsStrExt;
use std::path;
use std::sync;
use rand::Rng;
use tokio::executor::current_thread;
use memmap;
use futures_cpupool;
use libc::{close, open, O_DIRECT, O_RDWR};
const FILE_SIZE: u64 = 1024 * 512;
fn temp_file_name() -> path::PathBuf {
let mut rng = rand::thread_rng();
let mut result = env::temp_dir();
let filename = format!("test-aio-{}.dat", rng.gen::<u64>());
result.push(filename);
result
}
fn create_temp_file(path: &path::Path) {
let mut file = fs::File::create(path).unwrap();
let mut data: [u8; FILE_SIZE as usize] = [0; FILE_SIZE as usize];
for index in 0..data.len() {
data[index] = index as u8;
}
let result = file.write(&data).and_then(|_| file.sync_all());
assert!(result.is_ok());
}
fn remove_file(path: &path::Path) {
let _ = fs::remove_file(path);
}
#[test]
fn create_and_drop() {
let pool = futures_cpupool::CpuPool::new(3);
let _context = AioContext::new(&pool, 10).unwrap();
}
struct MemoryBlock {
bytes: sync::RwLock<memmap::MmapMut>,
}
impl MemoryBlock {
fn new() -> MemoryBlock {
let map = memmap::MmapMut::map_anon(8192).unwrap();
unsafe { mlock(map.as_ref().as_ptr() as *const c_void, map.len()) };
MemoryBlock {
bytes: sync::RwLock::new(map),
}
}
}
struct MemoryHandle {
block: sync::Arc<MemoryBlock>,
}
impl MemoryHandle {
fn new() -> MemoryHandle {
MemoryHandle {
block: sync::Arc::new(MemoryBlock::new()),
}
}
}
impl Clone for MemoryHandle {
fn clone(&self) -> MemoryHandle {
MemoryHandle {
block: self.block.clone(),
}
}
}
impl convert::AsRef<[u8]> for MemoryHandle {
fn as_ref(&self) -> &[u8] {
unsafe { mem::transmute(&(*self.block.bytes.read().unwrap())[..]) }
}
}
impl convert::AsMut<[u8]> for MemoryHandle {
fn as_mut(&mut self) -> &mut [u8] {
unsafe { mem::transmute(&mut (*self.block.bytes.write().unwrap())[..]) }
}
}
#[test]
fn read_block_mt() {
let file_name = temp_file_name();
create_temp_file(&file_name);
{
let owned_fd = OwnedFd::new_from_raw_fd(unsafe {
open(
mem::transmute(file_name.as_os_str().as_bytes().as_ptr()),
O_DIRECT | O_RDWR,
)
});
let fd = owned_fd.fd;
let pool = futures_cpupool::CpuPool::new(5);
let buffer = MemoryHandle::new();
{
let context = AioContext::new(&pool, 10).unwrap();
let read_future = context
.read(fd, 0, buffer)
.map(move |result_buffer| {
assert!(validate_block(result_buffer.as_ref()));
})
.map_err(|err| {
panic!("{:?}", err);
});
let cpu_future = pool.spawn(read_future);
let result = cpu_future.wait();
assert!(result.is_ok());
}
}
remove_file(&file_name);
}
#[test]
fn write_block_mt() {
use io::{Read, Seek};
let file_name = temp_file_name();
create_temp_file(&file_name);
{
let owned_fd = OwnedFd::new_from_raw_fd(unsafe {
open(
mem::transmute(file_name.as_os_str().as_bytes().as_ptr()),
O_DIRECT | O_RDWR,
)
});
let fd = owned_fd.fd;
let pool = futures_cpupool::CpuPool::new(5);
let mut buffer = MemoryHandle::new();
fill_pattern(65u8, buffer.as_mut());
{
let context = AioContext::new(&pool, 2).unwrap();
let write_future = context.write(fd, 16384, buffer).map_err(|err| {
panic!("{:?}", err);
});
let cpu_future = pool.spawn(write_future);
let result = cpu_future.wait();
assert!(result.is_ok());
}
}
let mut file = fs::File::open(&file_name).unwrap();
file.seek(io::SeekFrom::Start(16384)).unwrap();
let mut read_buffer: [u8; 8192] = [0u8; 8192];
file.read(&mut read_buffer).unwrap();
assert!(validate_pattern(65u8, &read_buffer));
}
#[test]
fn write_block_sync_mt() {
use io::{Read, Seek};
let file_name = temp_file_name();
create_temp_file(&file_name);
{
let owned_fd = OwnedFd::new_from_raw_fd(unsafe {
open(
mem::transmute(file_name.as_os_str().as_bytes().as_ptr()),
O_DIRECT | O_RDWR,
)
});
let fd = owned_fd.fd;
let pool = futures_cpupool::CpuPool::new(5);
let context = AioContext::new(&pool, 2).unwrap();
{
let mut buffer = MemoryHandle::new();
fill_pattern(65u8, buffer.as_mut());
let write_future = context.write(fd, 16384, buffer).map_err(|err| {
panic!("{:?}", err);
});
let cpu_future = pool.spawn(write_future);
let result = cpu_future.wait();
assert!(result.is_ok());
}
{
let mut buffer = MemoryHandle::new();
fill_pattern(66u8, buffer.as_mut());
let write_future = context.write(fd, 32768, buffer).map_err(|err| {
panic!("{:?}", err);
});
let cpu_future = pool.spawn(write_future);
let result = cpu_future.wait();
assert!(result.is_ok());
}
{
let mut buffer = MemoryHandle::new();
fill_pattern(67u8, buffer.as_mut());
let write_future = context.write(fd, 49152, buffer).map_err(|err| {
panic!("{:?}", err);
});
let cpu_future = pool.spawn(write_future);
let result = cpu_future.wait();
assert!(result.is_ok());
}
}
let mut file = fs::File::open(&file_name).unwrap();
let mut read_buffer: [u8; 8192] = [0u8; 8192];
file.seek(io::SeekFrom::Start(16384)).unwrap();
file.read(&mut read_buffer).unwrap();
assert!(validate_pattern(65u8, &read_buffer));
file.seek(io::SeekFrom::Start(32768)).unwrap();
file.read(&mut read_buffer).unwrap();
assert!(validate_pattern(66u8, &read_buffer));
file.seek(io::SeekFrom::Start(49152)).unwrap();
file.read(&mut read_buffer).unwrap();
assert!(validate_pattern(67u8, &read_buffer));
}
#[test]
fn read_invalid_fd() {
let fd = 2431;
let pool = futures_cpupool::CpuPool::new(5);
let buffer = MemoryHandle::new();
{
let context = AioContext::new(&pool, 10).unwrap();
let read_future = context
.read(fd, 0, buffer)
.map(move |_| {
assert!(false);
})
.map_err(|err| {
assert!(err.error.kind() == io::ErrorKind::Other);
err
});
let cpu_future = pool.spawn(read_future);
let result = cpu_future.wait();
assert!(result.is_err());
}
}
#[test]
fn read_many_blocks_mt() {
let file_name = temp_file_name();
create_temp_file(&file_name);
{
let owned_fd = OwnedFd::new_from_raw_fd(unsafe {
open(
mem::transmute(file_name.as_os_str().as_bytes().as_ptr()),
O_DIRECT | O_RDWR,
)
});
let fd = owned_fd.fd;
let pool = futures_cpupool::CpuPool::new(5);
{
let num_slots = 7;
let context = AioContext::new(&pool, num_slots).unwrap();
for _wave in 0..50 {
let mut futures = Vec::new();
for index in 0..100 {
let buffer = MemoryHandle::new();
let read_future = context
.read(fd, (index * 8192) % FILE_SIZE, buffer)
.map(move |result_buffer| {
assert!(validate_block(result_buffer.as_ref()));
})
.map_err(|err| {
panic!("{:?}", err);
});
futures.push(pool.spawn(read_future));
}
let result = futures::future::join_all(futures).wait();
assert!(result.is_ok());
assert!(context.inner.have_capacity.current_capacity() == num_slots);
}
}
}
remove_file(&file_name);
}
#[test]
fn mixed_read_write() {
let file_name = temp_file_name();
create_temp_file(&file_name);
let owned_fd = OwnedFd::new_from_raw_fd(unsafe {
open(
mem::transmute(file_name.as_os_str().as_bytes().as_ptr()),
O_DIRECT | O_RDWR,
)
});
let fd = owned_fd.fd;
let mut futures = Vec::new();
let pool = futures_cpupool::CpuPool::new(5);
let context = AioContext::new(&pool, 7).unwrap();
let buffer1 = MemoryHandle::new();
let sequence1 = {
let context1 = context.clone();
let context2 = context.clone();
let context3 = context.clone();
let context4 = context.clone();
let context5 = context.clone();
let context6 = context.clone();
context1
.read(fd, 8192, buffer1)
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(0u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context2.write(fd, 8192, buffer))
.and_then(move |buffer| context3.read(fd, 0, buffer))
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(1u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context4.write(fd, 0, buffer))
.and_then(move |buffer| context5.read(fd, 8192, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(0u8, buffer.as_ref()));
buffer
})
.and_then(move |buffer| context6.read(fd, 0, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(1u8, buffer.as_ref()));
buffer
})
.map_err(|err| {
panic!("{:?}", err);
})
};
let buffer2 = MemoryHandle::new();
let sequence2 = {
let context1 = context.clone();
let context2 = context.clone();
let context3 = context.clone();
let context4 = context.clone();
let context5 = context.clone();
let context6 = context.clone();
context1
.read(fd, 16384, buffer2)
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(2u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context2.write(fd, 16384, buffer))
.and_then(move |buffer| context3.read(fd, 24576, buffer))
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(3u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context4.write(fd, 24576, buffer))
.and_then(move |buffer| context5.read(fd, 16384, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(2u8, buffer.as_ref()));
buffer
})
.and_then(move |buffer| context6.read(fd, 24576, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(3u8, buffer.as_ref()));
buffer
})
.map_err(|err| {
panic!("{:?}", err);
})
};
let buffer3 = MemoryHandle::new();
let sequence3 = {
let context1 = context.clone();
let context2 = context.clone();
let context3 = context.clone();
let context4 = context.clone();
let context5 = context.clone();
let context6 = context.clone();
context1
.read(fd, 40960, buffer3)
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(5u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context2.write(fd, 40960, buffer))
.and_then(move |buffer| context3.read(fd, 32768, buffer))
.map(|mut buffer| -> MemoryHandle {
assert!(validate_block(buffer.as_ref()));
fill_pattern(6u8, buffer.as_mut());
buffer
})
.and_then(move |buffer| context4.write(fd, 32768, buffer))
.and_then(move |buffer| context5.read(fd, 40960, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(5u8, buffer.as_ref()));
buffer
})
.and_then(move |buffer| context6.read(fd, 32768, buffer))
.map(|buffer| -> MemoryHandle {
assert!(validate_pattern(6u8, buffer.as_ref()));
buffer
})
.map_err(|err| {
panic!("{:?}", err);
})
};
futures.push(pool.spawn(sequence1));
futures.push(pool.spawn(sequence2));
futures.push(pool.spawn(sequence3));
let result = futures::future::join_all(futures).wait();
assert!(result.is_ok());
}
fn fill_pattern(key: u8, buffer: &mut [u8]) {
assert!(buffer.len() % 2 == 0);
for index in 0..buffer.len() / 2 {
buffer[index * 2] = key;
buffer[index * 2 + 1] = index as u8;
}
}
fn validate_pattern(key: u8, buffer: &[u8]) -> bool {
assert!(buffer.len() % 2 == 0);
for index in 0..buffer.len() / 2 {
if (buffer[index * 2] != key) || (buffer[index * 2 + 1] != (index as u8)) {
return false;
}
}
return true;
}
fn validate_block(data: &[u8]) -> bool {
for index in 0..data.len() {
if data[index] != index as u8 {
return false;
}
}
true
}
struct OwnedFd {
fd: RawFd,
}
impl OwnedFd {
fn new_from_raw_fd(fd: RawFd) -> OwnedFd {
OwnedFd { fd }
}
}
impl Drop for OwnedFd {
fn drop(&mut self) {
let result = unsafe { close(self.fd) };
assert!(result == 0);
}
}
}