wireshift-fallback 0.1.1

Blocking worker-pool fallback backend for wireshift
Documentation
use std::collections::HashMap;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;

use wireshift_core::op::{CompletionPayload, OpDescriptor};
use wireshift_core::{Error, Result};

/// File open / openat operation handlers.
pub mod open;
/// Miscellaneous operations (statx, fsync, splice, madvise, linked, cancel).
pub mod other;
/// Read and vectored-read operation handlers.
pub mod read;
/// Write and vectored-write operation handlers.
pub mod write;

macro_rules! retry_eintr {
    ($expr:expr) => {
        loop {
            match $expr {
                Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
                result => break result,
            }
        }
    };
}
pub(crate) use retry_eintr;

pub(crate) fn execute_descriptor(
    descriptor: OpDescriptor,
    canceled: Arc<AtomicBool>,
    timeout: Option<Duration>,
) -> Result<CompletionPayload> {
    execute_descriptor_with_fixed(descriptor, canceled, &mut HashMap::new(), timeout)
}

pub(crate) fn execute_descriptor_with_fixed(
    descriptor: OpDescriptor,
    canceled: Arc<AtomicBool>,
    fixed_files: &mut HashMap<u32, std::fs::File>,
    timeout: Option<Duration>,
) -> Result<CompletionPayload> {
    if canceled.load(Ordering::Relaxed) {
        return Err(Error::canceled(
            "operation canceled before dispatch",
            "avoid canceling the request before the backend starts it",
        ));
    }
    match descriptor {
        OpDescriptor::Read {
            file,
            offset,
            buffer,
            len,
        } => read::execute_read(file, offset, buffer, canceled, len),
        OpDescriptor::Write {
            file,
            offset,
            buffer,
        } => write::execute_write(file, offset, buffer),
        OpDescriptor::ReadVectored {
            file,
            offset,
            buffers,
        } => read::execute_read_vectored(file, offset, buffers),
        OpDescriptor::ReadGpu {
            file,
            offset,
            buffer,
        } => read::execute_read_gpu(file, offset, buffer, canceled),
        OpDescriptor::WriteVectored {
            file,
            offset,
            buffers,
        } => write::execute_write_vectored(file, offset, buffers),
        OpDescriptor::Connect {
            addr,
            timeout: op_timeout,
        } => crate::net_ops::execute_connect(addr, op_timeout.or(timeout), canceled),
        OpDescriptor::Accept { listener } => {
            crate::net_ops::execute_accept(listener, canceled, timeout)
        }
        OpDescriptor::Send { stream, buffer } => crate::net_ops::execute_send(stream, buffer),
        OpDescriptor::Recv { stream, buffer } => {
            crate::net_ops::execute_recv(stream, buffer, canceled, timeout)
        }
        OpDescriptor::OpenAt {
            dir,
            path,
            flags,
            read,
            write,
            create,
            truncate,
        } => open::execute_openat(dir, path, flags, read, write, create, truncate),
        OpDescriptor::OpenAtDirect {
            dir,
            path,
            flags,
            read,
            write,
            create,
            truncate,
            slot,
        } => open::execute_openat_direct(
            dir,
            path,
            flags,
            read,
            write,
            create,
            truncate,
            slot,
            fixed_files,
        ),
        OpDescriptor::ReadFixed {
            slot,
            offset,
            buffer,
        } => read::execute_read_fixed(slot, offset, buffer, canceled, fixed_files),
        OpDescriptor::CloseFixed { slot } => other::execute_close_fixed(slot, fixed_files),
        OpDescriptor::Statx { path } => other::execute_statx(path),
        OpDescriptor::Fsync { file } => other::execute_fsync(file),
        OpDescriptor::Cancel { target } => Err(Error::Unsupported {
            message: format!("direct cancel operation for request {target} is internal only")
                .into(),
            fix: "call Ring::cancel(request_id) instead".into(),
        }),
        OpDescriptor::Linked { descriptors } => {
            other::execute_linked(descriptors, canceled, timeout)
        }
        OpDescriptor::Splice {
            fd_in,
            off_in,
            fd_out,
            off_out,
            len,
        } => other::execute_splice(fd_in, off_in, fd_out, off_out, len),
        OpDescriptor::Madvise { file, advice } => other::execute_madvise(file, advice),
        OpDescriptor::Nop => Ok(CompletionPayload::Unit),
        _ => Err(Error::Unsupported {
            message: "unsupported operation descriptor".into(),
            fix: "upgrade wireshift-fallback to handle this descriptor".into(),
        }),
    }
}