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};
pub mod open;
pub mod other;
pub mod read;
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(),
}),
}
}