rafts 0.1.0

Minimal experimental async block device I/O abstraction
Documentation
//! Example: File-backed block device using the tokio runtime with blocking I/O in spawn_blocking.
//!
//! Run:
//!   cargo run --example file_block_device --features tokio
//!
//! Demonstrates implementing the callback based `BlockAsyncIo` trait on top of
//! a regular `std::fs::File` by off-loading blocking reads / writes to
//! `tokio::task::spawn_blocking`, then using the async extension helpers.

use rafts::error::Errno;
use rafts::io::BlockAsyncIo;
use std::fs::OpenOptions;
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::PathBuf;
use std::sync::{Arc, Mutex};

struct FileBlockDevice {
    file: Arc<Mutex<std::fs::File>>,
    block_size: usize,
}

impl FileBlockDevice {
    fn open(path: PathBuf, block_size: usize) -> std::io::Result<Self> {
        let file = OpenOptions::new()
            .read(true)
            .write(true)
            .create(true)
            .open(path)?;
        Ok(Self {
            file: Arc::new(Mutex::new(file)),
            block_size,
        })
    }
}

impl BlockAsyncIo for FileBlockDevice {
    fn block_size(&self) -> usize {
        self.block_size
    }

    fn read(
        &self,
        block: u64,
        buf: &mut [u8],
        completion: Box<dyn FnOnce(Result<(), Errno>) + Send + 'static>,
    ) {
        let file = self.file.clone();
        let bs = self.block_size;
        let len = buf.len();
        let offset = block * bs as u64;
        let ptr = buf.as_mut_ptr() as usize; // make Send
        tokio::task::spawn_blocking(move || {
            let res = (|| {
                if len % bs != 0 {
                    return Err(Errno::Unknown);
                }
                let mut f = file.lock().unwrap();
                f.seek(SeekFrom::Start(offset))
                    .map_err(|_| Errno::Unknown)?;
                unsafe {
                    let slice = std::slice::from_raw_parts_mut(ptr as *mut u8, len);
                    f.read_exact(slice).map_err(|_| Errno::Unknown)?;
                }
                Ok(())
            })();
            completion(res);
        });
    }

    fn write(
        &self,
        block: u64,
        buf: &[u8],
        completion: Box<dyn FnOnce(Result<(), Errno>) + Send + 'static>,
    ) {
        let file = self.file.clone();
        let bs = self.block_size;
        let len = buf.len();
        let offset = block * bs as u64;
        let owned = buf.to_vec();
        tokio::task::spawn_blocking(move || {
            let res = (|| {
                if len % bs != 0 {
                    return Err(Errno::Unknown);
                }
                let mut f = file.lock().unwrap();
                f.seek(SeekFrom::Start(offset))
                    .map_err(|_| Errno::Unknown)?;
                f.write_all(&owned).map_err(|_| Errno::Unknown)?;
                Ok(())
            })();
            completion(res);
        });
    }
}

fn main() {
    use rafts::io::BlockAsyncIoExt; // async extension trait
    let rt = tokio::runtime::Runtime::new().expect("runtime");
    rt.block_on(async {
        let dev = Arc::new(FileBlockDevice::open(PathBuf::from("block.img"), 4096).expect("open"));
        let mut write_buf = vec![0u8; 4096];
        write_buf[0..11].copy_from_slice(b"hello world");
        dev.write_async(0, &write_buf).await.expect("write");

        let mut read_buf = vec![0u8; 4096];
        dev.read_async(0, &mut read_buf).await.expect("read");
        println!("Read back: {}", String::from_utf8_lossy(&read_buf[..11]));
    });
}