dope_core/driver/
completion.rs1use std::io;
2use std::time::Duration;
3
4use super::DriverContext;
5use crate::io::Cqe;
6
7pub trait Completion {
8 fn drain(&mut self, buf: &mut [Cqe]) -> usize;
9 fn wait(&mut self, timeout: Option<Duration>) -> io::Result<()>;
10}
11
12cfg_select! {
13 target_os = "linux" => {
14 use io_uring::types::{SubmitArgs, Timespec};
15
16 use crate::backend::uring::driver::{Disposition, Uring};
17
18 impl Completion for DriverContext<'_, '_> {
19 fn drain(&mut self, buf: &mut [Cqe]) -> usize {
20 self.flush_returned_buffers();
21 let state = self.backend();
22 let mut n = 0;
23 {
24 let Uring {
25 uring,
26 setsockopt,
27 files,
28 provided,
29 routes,
30 ..
31 } = state;
32 let mut cq = uring.completion();
33 while n < buf.len() {
34 let Some(item) = cq.next() else { break };
35 let result = item.result();
36 let user_data = match Uring::complete_cqe(
37 setsockopt,
38 files,
39 routes,
40 item.user_data(),
41 result,
42 item.flags(),
43 ) {
44 Disposition::Drop | Disposition::Internal => continue,
45 Disposition::DropBuffer(bid) => {
46 provided.defer(bid);
47 continue;
48 }
49 Disposition::Public(user_data) => user_data,
50 };
51 buf[n] = Cqe {
52 user_data,
53 result,
54 flags: item.flags(),
55 };
56 n += 1;
57 }
58 cq.sync();
59 }
60 state.flush_deferred_close();
61 state.flush_ready_create();
62 state.provided.flush();
63 n
64 }
65
66 fn wait(&mut self, timeout: Option<Duration>) -> io::Result<()> {
67 self.flush_returned_buffers();
68 let state = self.backend();
69 state.flush_deferred_close();
70 state.flush_ready_create();
71 state.provided.flush();
72 match timeout {
73 Some(timeout) => {
74 let timespec = Timespec::from(timeout);
75 let args = SubmitArgs::new().timespec(×pec);
76 match state.uring.submitter().submit_with_args(1, &args) {
77 Ok(_) => Ok(()),
78 Err(error) if error.raw_os_error() == Some(libc::ETIME) => Ok(()),
79 Err(error) => Err(error),
80 }
81 }
82 None => state.uring.submitter().submit_and_wait(1).map(|_| ()),
83 }
84 }
85 }
86 }
87 _ => {
88 use std::mem::MaybeUninit;
89 use std::slice;
90
91 use crate::backend::kqueue::driver::MAX_DRAIN_PER_FD;
92 use crate::backend::kqueue::driver::pending::PendingCompletion;
93 use crate::backend::kqueue::driver::read::dispatch::Dispatch;
94 use crate::driver::token::SHUTDOWN;
95 use crate::io::{BUFFER, BUFFER_SHIFT, MORE};
96
97 impl Completion for DriverContext<'_, '_> {
98 fn drain(&mut self, buf: &mut [Cqe]) -> usize {
99 self.flush_returned_buffers();
100 if self.backend_ref().pending.is_empty() {
101 let _ = Completion::wait(self, Some(Duration::ZERO));
102 }
103 let state = self.backend();
104 let mut n = 0;
105 while n < buf.len() {
106 let Some(pending) = state.pending.pop_front() else {
107 break;
108 };
109 buf[n] = match pending {
110 PendingCompletion::Accept { ud, result, more } => Cqe {
111 user_data: ud.raw(),
112 result,
113 flags: if more { MORE } else { 0 },
114 },
115 PendingCompletion::Recv {
116 ud,
117 result,
118 more,
119 bid,
120 } => {
121 let mut flags = if more { MORE } else { 0 };
122 if let Some(bid) = bid {
123 flags |= BUFFER | ((bid as u32) << BUFFER_SHIFT);
124 }
125 Cqe {
126 user_data: ud.raw(),
127 result,
128 flags,
129 }
130 }
131 PendingCompletion::Write { ud, result } => Cqe {
132 user_data: ud.raw(),
133 result,
134 flags: 0,
135 },
136 PendingCompletion::Create { ud, result, .. } => Cqe {
137 user_data: ud.raw(),
138 result,
139 flags: 0,
140 },
141 PendingCompletion::Timer { ud } => Cqe {
142 user_data: ud.raw(),
143 result: 0,
144 flags: 0,
145 },
146 PendingCompletion::Shutdown => Cqe {
147 user_data: SHUTDOWN.raw(),
148 result: 0,
149 flags: 0,
150 },
151 };
152 n += 1;
153 }
154 n
155 }
156
157 fn wait(&mut self, timeout: Option<Duration>) -> io::Result<()> {
158 self.flush_returned_buffers();
159 let state = self.backend();
160 state.resume_pending();
161 let mut events: [MaybeUninit<libc::kevent>; 64] = [const { MaybeUninit::uninit() }; 64];
162 if state.pending.remaining_capacity() < events.len() * MAX_DRAIN_PER_FD {
163 return Ok(());
164 }
165 let n = state.kevent_call(&mut events, timeout)?;
166 let ready = unsafe { slice::from_raw_parts(events.as_ptr().cast::<libc::kevent>(), n) };
167 for event in ready {
168 state.dispatch_event(event);
169 }
170 Ok(())
171 }
172 }
173 }
174}