Skip to main content

wireshift_fallback/ops/
mod.rs

1use std::collections::HashMap;
2use std::sync::atomic::{AtomicBool, Ordering};
3use std::sync::Arc;
4use std::time::Duration;
5
6use wireshift_core::op::{CompletionPayload, OpDescriptor};
7use wireshift_core::{Error, Result};
8
9/// File open / openat operation handlers.
10pub mod open;
11/// Miscellaneous operations (statx, fsync, splice, madvise, linked, cancel).
12pub mod other;
13/// Read and vectored-read operation handlers.
14pub mod read;
15/// Write and vectored-write operation handlers.
16pub mod write;
17
18macro_rules! retry_eintr {
19    ($expr:expr) => {
20        loop {
21            match $expr {
22                Err(ref e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
23                result => break result,
24            }
25        }
26    };
27}
28pub(crate) use retry_eintr;
29
30pub(crate) fn execute_descriptor(
31    descriptor: OpDescriptor,
32    canceled: Arc<AtomicBool>,
33    timeout: Option<Duration>,
34) -> Result<CompletionPayload> {
35    execute_descriptor_with_fixed(descriptor, canceled, &mut HashMap::new(), timeout)
36}
37
38pub(crate) fn execute_descriptor_with_fixed(
39    descriptor: OpDescriptor,
40    canceled: Arc<AtomicBool>,
41    fixed_files: &mut HashMap<u32, std::fs::File>,
42    timeout: Option<Duration>,
43) -> Result<CompletionPayload> {
44    if canceled.load(Ordering::Relaxed) {
45        return Err(Error::canceled(
46            "operation canceled before dispatch",
47            "avoid canceling the request before the backend starts it",
48        ));
49    }
50    match descriptor {
51        OpDescriptor::Read {
52            file,
53            offset,
54            buffer,
55            len,
56        } => read::execute_read(file, offset, buffer, canceled, len),
57        OpDescriptor::Write {
58            file,
59            offset,
60            buffer,
61        } => write::execute_write(file, offset, buffer),
62        OpDescriptor::ReadVectored {
63            file,
64            offset,
65            buffers,
66        } => read::execute_read_vectored(file, offset, buffers),
67        OpDescriptor::ReadGpu {
68            file,
69            offset,
70            buffer,
71        } => read::execute_read_gpu(file, offset, buffer, canceled),
72        OpDescriptor::WriteVectored {
73            file,
74            offset,
75            buffers,
76        } => write::execute_write_vectored(file, offset, buffers),
77        OpDescriptor::Connect {
78            addr,
79            timeout: op_timeout,
80        } => crate::net_ops::execute_connect(addr, op_timeout.or(timeout), canceled),
81        OpDescriptor::Accept { listener } => {
82            crate::net_ops::execute_accept(listener, canceled, timeout)
83        }
84        OpDescriptor::Send { stream, buffer } => crate::net_ops::execute_send(stream, buffer),
85        OpDescriptor::Recv { stream, buffer } => {
86            crate::net_ops::execute_recv(stream, buffer, canceled, timeout)
87        }
88        OpDescriptor::OpenAt {
89            dir,
90            path,
91            flags,
92            read,
93            write,
94            create,
95            truncate,
96        } => open::execute_openat(dir, path, flags, read, write, create, truncate),
97        OpDescriptor::OpenAtDirect {
98            dir,
99            path,
100            flags,
101            read,
102            write,
103            create,
104            truncate,
105            slot,
106        } => open::execute_openat_direct(
107            dir,
108            path,
109            flags,
110            read,
111            write,
112            create,
113            truncate,
114            slot,
115            fixed_files,
116        ),
117        OpDescriptor::ReadFixed {
118            slot,
119            offset,
120            buffer,
121        } => read::execute_read_fixed(slot, offset, buffer, canceled, fixed_files),
122        OpDescriptor::CloseFixed { slot } => other::execute_close_fixed(slot, fixed_files),
123        OpDescriptor::Statx { path } => other::execute_statx(path),
124        OpDescriptor::Fsync { file } => other::execute_fsync(file),
125        OpDescriptor::Cancel { target } => Err(Error::Unsupported {
126            message: format!("direct cancel operation for request {target} is internal only")
127                .into(),
128            fix: "call Ring::cancel(request_id) instead".into(),
129        }),
130        OpDescriptor::Linked { descriptors } => {
131            other::execute_linked(descriptors, canceled, timeout)
132        }
133        OpDescriptor::Splice {
134            fd_in,
135            off_in,
136            fd_out,
137            off_out,
138            len,
139        } => other::execute_splice(fd_in, off_in, fd_out, off_out, len),
140        OpDescriptor::Madvise { file, advice } => other::execute_madvise(file, advice),
141        OpDescriptor::Nop => Ok(CompletionPayload::Unit),
142        _ => Err(Error::Unsupported {
143            message: "unsupported operation descriptor".into(),
144            fix: "upgrade wireshift-fallback to handle this descriptor".into(),
145        }),
146    }
147}