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
9pub mod open;
11pub mod other;
13pub mod read;
15pub 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}