1use crate::backend::Sqe;
2
3use super::{DriverContext, PushError};
4
5pub trait Submission {
6 fn push(&mut self, sqe: Sqe) -> Result<(), PushError>;
7 fn flush_submissions(&mut self) -> bool;
8}
9
10cfg_select! {
11 target_os = "linux" => {
12 use crate::backend::uring::driver::Uring;
13 use crate::backend::uring::driver::files::Admission;
14
15 impl Submission for DriverContext<'_, '_> {
16 fn push(&mut self, sqe: Sqe) -> Result<(), PushError> {
17 let state = self.backend();
18 let Some(create) = sqe.create_meta() else {
19 return Uring::entry_push(&mut state.uring, sqe.entry());
20 };
21 match state.files.admission(create.slot) {
22 Admission::Start => {
23 Uring::entry_push(&mut state.uring, sqe.entry())?;
24 state.files.begin_create(create);
25 Ok(())
26 }
27 Admission::Defer => {
28 state.files.defer_create(create, sqe);
29 Ok(())
30 }
31 Admission::Reject => Err(PushError),
32 }
33 }
34
35 fn flush_submissions(&mut self) -> bool {
36 let state = self.backend();
37 state.flush_deferred_close();
38 state.flush_ready_create();
39 state.uring.submit().is_ok()
40 }
41 }
42 }
43 _ => {
44 use crate::backend::kqueue::driver::pending::PendingCompletion;
45 use crate::backend::kqueue::driver::read::arm::Arm;
46 use crate::backend::kqueue::driver::submit::Submit;
47 use crate::backend::kqueue::sqe::SqeInner;
48
49 impl Submission for DriverContext<'_, '_> {
50 fn push(&mut self, sqe: Sqe) -> Result<(), PushError> {
51 let state = self.backend();
52 if state.pending.is_full()
53 && !matches!(
54 &sqe.0,
55 SqeInner::Quickack
56 | SqeInner::Shutdown { .. }
57 | SqeInner::Cancel { .. }
58 | SqeInner::CancelCreate { .. }
59 )
60 {
61 return Err(PushError);
62 }
63 let accepted = match sqe.0 {
64 SqeInner::AcceptOneshot {
65 listener,
66 addr_ptr,
67 addrlen_ptr,
68 ud,
69 } => {
70 let Some(raw) = state.raw_fd(listener) else {
71 state.push_pending(PendingCompletion::Accept {
72 ud,
73 result: -libc::EBADF,
74 more: false,
75 });
76 return Ok(());
77 };
78 state.arm_accept_oneshot_inner(ud, raw, addr_ptr, addrlen_ptr)
79 }
80 SqeInner::RecvMulti { slot, ud } => state.arm_recv_multi_inner(ud, slot),
81 SqeInner::RecvMsgMulti { slot, msghdr, ud } => unsafe {
82 state.arm_recv_msg_multi_inner(ud, slot, msghdr)
83 },
84 SqeInner::Send { slot, ptr, len, ud } => {
85 state.submit_send_tagged_inner(ud, slot, ptr, len)
86 }
87 SqeInner::WriteFd {
88 fd,
89 ptr,
90 len,
91 offset,
92 ud,
93 } => state.submit_write_fd_inner(ud, fd, ptr, len, offset),
94 SqeInner::OpenAt {
95 dir,
96 path,
97 flags,
98 mode,
99 ud,
100 } => state.submit_openat_inner(ud, dir, path, flags, mode, None),
101 SqeInner::OpenAtFixed {
102 dir,
103 path,
104 flags,
105 mode,
106 slot,
107 ud,
108 } => state.submit_openat_inner(ud, dir, path, flags, mode, Some(slot)),
109 SqeInner::Read {
110 fd,
111 ptr,
112 len,
113 offset,
114 ud,
115 } => state.submit_read_inner(ud, fd, ptr, len, offset),
116 SqeInner::ReadFixed {
117 slot,
118 ptr,
119 len,
120 offset,
121 ud,
122 } => match state.raw_fd(slot) {
123 Some(fd) => state.submit_read_inner(ud, fd, ptr, len, offset),
124 None => {
125 state.push_pending(PendingCompletion::Write {
126 ud,
127 result: -libc::EBADF,
128 });
129 true
130 }
131 },
132 SqeInner::StatPath { path, stat, ud } => {
133 let rc = unsafe { libc::stat(path, stat) };
134 state.complete_io(ud, rc as isize)
135 }
136 SqeInner::StatFd { fd, stat, ud } => {
137 let rc = unsafe { libc::fstat(fd, stat) };
138 state.complete_io(ud, rc as isize)
139 }
140 SqeInner::Splice {
141 fd_in,
142 off_in,
143 fd_out,
144 off_out,
145 len,
146 ud,
147 } => state.submit_splice_inner(ud, fd_in, off_in, fd_out, off_out, len),
148 SqeInner::SendMsg { slot, msg, ud } => unsafe {
149 state.submit_send_msg_tagged_inner(ud, slot, msg)
150 },
151 SqeInner::Quickack => true,
152 SqeInner::Shutdown { slot, how } => {
153 if let Some(raw) = state.raw_fd(slot) {
154 unsafe { libc::shutdown(raw, how) };
155 }
156 true
157 }
158 SqeInner::Cancel { target } => state.cancel_inner(target),
159 SqeInner::Interval { sec, nsec, ud } => {
160 let micros = (i128::from(sec) * 1_000_000 + i128::from(nsec) / 1_000)
161 .clamp(1, libc::intptr_t::MAX as i128)
162 as libc::intptr_t;
163 state.changes.push(libc::kevent {
164 ident: ud.raw() as libc::uintptr_t,
165 filter: libc::EVFILT_TIMER,
166 flags: libc::EV_ADD,
167 fflags: libc::NOTE_USECONDS,
168 data: micros,
169 udata: ud.raw() as usize as *mut libc::c_void,
170 });
171 state.flush_changes_if_full();
172 true
173 }
174 SqeInner::CancelCreate { slot } => {
175 state.pending.cancel_create(slot);
176 state.close_fd(slot);
177 true
178 }
179 SqeInner::SocketAt {
180 domain,
181 socket_type,
182 protocol,
183 slot,
184 ud,
185 } => state.submit_socket_at(domain, socket_type, protocol, slot, ud),
186 SqeInner::Connect {
187 slot,
188 addr_ptr,
189 addr_len,
190 ud,
191 } => state.submit_connect(slot, addr_ptr, addr_len, ud),
192 };
193 if accepted {
194 Ok(())
195 } else {
196 Err(PushError)
197 }
198 }
199
200 fn flush_submissions(&mut self) -> bool {
201 false
202 }
203 }
204 }
205}