Skip to main content

dope_core/driver/
submission.rs

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}