Skip to main content

dope_core/driver/
completion.rs

1use 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(&timespec);
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}