1pub mod datagram;
2pub mod fd;
3pub(crate) mod ffi;
4pub mod file;
5pub mod pipe;
6pub mod provided;
7pub mod socket;
8
9use std::error::Error as StdError;
10use std::fmt;
11use std::io::{self, Error};
12
13use crate::driver::token;
14use crate::driver::token::{Token, kind};
15use fd::FdSlot;
16
17#[derive(Clone, Copy, Debug, PartialEq, Eq)]
18pub struct DecodeError;
19
20impl fmt::Display for DecodeError {
21 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
22 f.write_str("invalid completion")
23 }
24}
25
26impl StdError for DecodeError {}
27
28pub(crate) const BUFFER: u32 = 1 << 0;
29pub(crate) const MORE: u32 = 1 << 1;
30pub(crate) const BUFFER_SHIFT: u32 = 16;
31
32#[derive(Clone, Copy)]
33pub struct Cqe {
34 pub user_data: u64,
35 pub result: i32,
36 pub flags: u32,
37}
38
39impl Cqe {
40 pub const ZERO: Self = Self {
41 user_data: 0,
42 result: 0,
43 flags: 0,
44 };
45
46 pub fn route(self) -> u8 {
47 (self.user_data >> token::ROUTE_SHIFT) as u8
48 }
49
50 pub fn kind(self) -> u8 {
51 (self.user_data >> token::KIND_SHIFT) as u8
52 }
53
54 fn more(self) -> bool {
55 self.flags & MORE != 0
56 }
57
58 fn bid_raw(self) -> u16 {
59 (self.flags >> BUFFER_SHIFT) as u16
60 }
61
62 fn has_buffer(self) -> bool {
63 self.flags & BUFFER != 0
64 }
65}
66
67#[derive(Clone, Copy)]
68pub enum RecvEvent {
69 Data { len: u32, bid: u16 },
70 Discarded { len: u32 },
71 Eof,
72 Cancelled,
73 Starved,
74 Failed(i32),
75}
76
77impl RecvEvent {
78 fn from_errno(result: i32) -> Self {
79 match -result {
80 libc::ECANCELED => Self::Cancelled,
81 libc::ENOBUFS => Self::Starved,
82 libc::EAGAIN | libc::EINTR => Self::Starved,
83 errno => Self::Failed(errno),
84 }
85 }
86}
87
88#[derive(Clone, Copy)]
89pub enum SendEvent {
90 Sent(u32),
91 Failed(i32),
92}
93
94#[derive(Clone, Copy)]
95pub enum WriteEvent {
96 Wrote(u32),
97 Failed(i32),
98}
99
100#[derive(Clone, Copy)]
101pub enum SyncEvent {
102 Synced,
103 Failed(i32),
104}
105
106#[derive(Clone, Copy)]
107pub enum OpenEvent {
108 Opened(i32),
109 Failed(i32),
110}
111
112#[derive(Clone, Copy)]
113pub enum ReadEvent {
114 Read(u32),
115 Eof,
116 Failed(i32),
117}
118
119impl ReadEvent {
120 fn from_result(result: i32) -> Self {
121 match result {
122 n if n > 0 => Self::Read(n as u32),
123 0 => Self::Eof,
124 n => Self::Failed(-n),
125 }
126 }
127}
128
129#[derive(Clone, Copy)]
130pub enum SpliceEvent {
131 Moved(u32),
132 Eof,
133 Failed(i32),
134}
135
136#[derive(Clone, Copy)]
137pub enum StatEvent {
138 Done,
139 Failed(i32),
140}
141
142#[derive(Clone, Copy)]
143pub enum AcceptEvent {
144 Accepted(FdSlot),
145 Failed,
146}
147
148pub enum SocketEvent {
149 Created,
150 Failed(io::Error),
151}
152
153pub enum ConnectEvent {
154 Connected,
155 Failed(io::Error),
156}
157
158pub struct Event(EventKind);
159
160pub enum EventKind {
161 Accept(Token, bool, AcceptEvent),
162 Recv(Token, bool, RecvEvent),
163 Send(Token, SendEvent),
164 Timer(Token),
165 Socket(Token, SocketEvent),
166 Connect(Token, ConnectEvent),
167 Write(Token, WriteEvent),
168 Sync(Token, SyncEvent),
169 Open(Token, OpenEvent),
170 Read(Token, ReadEvent),
171 ReadBlock(Token, ReadEvent),
172 Splice(Token, SpliceEvent),
173 Stat(Token, StatEvent),
174}
175
176impl Event {
177 pub fn decode(c: Cqe) -> Result<Self, DecodeError> {
178 let token = Token::try_from_raw(c.user_data).ok_or(DecodeError)?;
179 match c.kind() {
180 kind::ACCEPT => {
181 let e = match c.result {
182 n if n >= 0 => AcceptEvent::Accepted(FdSlot::new(n as u32)),
183 _ => AcceptEvent::Failed,
184 };
185 Ok(Self(EventKind::Accept(token, c.more(), e)))
186 }
187 kind::RECV => {
188 let e = match c.result {
189 n if n > 0 => {
190 if !c.has_buffer() {
191 debug_assert!(false, "RECV data cqe without buffer flag");
192 return Err(DecodeError);
193 }
194 RecvEvent::Data {
195 len: n as u32,
196 bid: c.bid_raw(),
197 }
198 }
199 0 => RecvEvent::Eof,
200 n => RecvEvent::from_errno(n),
201 };
202 Ok(Self(EventKind::Recv(token, c.more(), e)))
203 }
204 kind::RECV_DISCARD => {
205 let e = match c.result {
206 n if n > 0 => RecvEvent::Discarded { len: n as u32 },
207 0 => RecvEvent::Eof,
208 n => RecvEvent::from_errno(n),
209 };
210 Ok(Self(EventKind::Recv(token, c.more(), e)))
211 }
212 kind::SEND => {
213 let e = if c.result >= 0 {
214 SendEvent::Sent(c.result as u32)
215 } else {
216 SendEvent::Failed(-c.result)
217 };
218 Ok(Self(EventKind::Send(token, e)))
219 }
220 kind::WRITE => {
221 let e = if c.result >= 0 {
222 WriteEvent::Wrote(c.result as u32)
223 } else {
224 WriteEvent::Failed(-c.result)
225 };
226 Ok(Self(EventKind::Write(token, e)))
227 }
228 kind::SYNC => {
229 let e = if c.result >= 0 {
230 SyncEvent::Synced
231 } else {
232 SyncEvent::Failed(-c.result)
233 };
234 Ok(Self(EventKind::Sync(token, e)))
235 }
236 kind::OPEN => {
237 let e = if c.result >= 0 {
238 OpenEvent::Opened(c.result)
239 } else {
240 OpenEvent::Failed(-c.result)
241 };
242 Ok(Self(EventKind::Open(token, e)))
243 }
244 kind::READ => Ok(Self(EventKind::Read(
245 token,
246 ReadEvent::from_result(c.result),
247 ))),
248 kind::READ_BLOCK => Ok(Self(EventKind::ReadBlock(
249 token,
250 ReadEvent::from_result(c.result),
251 ))),
252 kind::SPLICE => {
253 let e = match c.result {
254 n if n > 0 => SpliceEvent::Moved(n as u32),
255 0 => SpliceEvent::Eof,
256 n => SpliceEvent::Failed(-n),
257 };
258 Ok(Self(EventKind::Splice(token, e)))
259 }
260 kind::STAT => {
261 let e = if c.result >= 0 {
262 StatEvent::Done
263 } else {
264 StatEvent::Failed(-c.result)
265 };
266 Ok(Self(EventKind::Stat(token, e)))
267 }
268 kind::TIMER => Ok(Self(EventKind::Timer(token))),
269 kind::SOCKET => {
270 let e = if c.result >= 0 {
271 SocketEvent::Created
272 } else {
273 SocketEvent::Failed(Error::from_raw_os_error(-c.result))
274 };
275 Ok(Self(EventKind::Socket(token, e)))
276 }
277 kind::CONNECT => {
278 let e = if c.result >= 0 {
279 ConnectEvent::Connected
280 } else {
281 ConnectEvent::Failed(Error::from_raw_os_error(-c.result))
282 };
283 Ok(Self(EventKind::Connect(token, e)))
284 }
285 _ => Err(DecodeError),
286 }
287 }
288}
289
290impl TryFrom<Cqe> for EventKind {
291 type Error = DecodeError;
292
293 fn try_from(cqe: Cqe) -> Result<Self, DecodeError> {
294 Event::decode(cqe).map(Event::into_kind)
295 }
296}
297
298pub enum EventRef<'a> {
299 Accept(Token, bool, &'a AcceptEvent),
300 Recv(Token, bool, &'a RecvEvent),
301 Send(Token, &'a SendEvent),
302 Timer(Token),
303 Socket(Token, &'a SocketEvent),
304 Connect(Token, &'a ConnectEvent),
305 Write(Token, &'a WriteEvent),
306 Sync(Token, &'a SyncEvent),
307 Open(Token, &'a OpenEvent),
308 Read(Token, &'a ReadEvent),
309 Splice(Token, &'a SpliceEvent),
310 Stat(Token, &'a StatEvent),
311}
312
313impl Event {
314 pub unsafe fn from_cqe(cqe: Cqe) -> Result<Self, DecodeError> {
317 Self::decode(cqe)
318 }
319
320 pub fn into_kind(self) -> EventKind {
321 self.0
322 }
323
324 pub fn as_ref(&self) -> EventRef<'_> {
325 match &self.0 {
326 EventKind::Accept(t, more, e) => EventRef::Accept(*t, *more, e),
327 EventKind::Recv(t, more, e) => EventRef::Recv(*t, *more, e),
328 EventKind::Send(t, e) => EventRef::Send(*t, e),
329 EventKind::Timer(t) => EventRef::Timer(*t),
330 EventKind::Socket(t, e) => EventRef::Socket(*t, e),
331 EventKind::Connect(t, e) => EventRef::Connect(*t, e),
332 EventKind::Write(t, e) => EventRef::Write(*t, e),
333 EventKind::Sync(t, e) => EventRef::Sync(*t, e),
334 EventKind::Open(t, e) => EventRef::Open(*t, e),
335 EventKind::Read(t, e) => EventRef::Read(*t, e),
336 EventKind::ReadBlock(t, e) => EventRef::Read(*t, e),
337 EventKind::Splice(t, e) => EventRef::Splice(*t, e),
338 EventKind::Stat(t, e) => EventRef::Stat(*t, e),
339 }
340 }
341
342 pub fn route(&self) -> u8 {
343 match &self.0 {
344 EventKind::Accept(t, ..) => t.route(),
345 EventKind::Recv(t, ..) => t.route(),
346 EventKind::Send(t, _) => t.route(),
347 EventKind::Timer(t) => t.route(),
348 EventKind::Socket(t, _) => t.route(),
349 EventKind::Connect(t, _) => t.route(),
350 EventKind::Write(t, _) => t.route(),
351 EventKind::Sync(t, _) => t.route(),
352 EventKind::Open(t, _) => t.route(),
353 EventKind::Read(t, _) => t.route(),
354 EventKind::ReadBlock(t, _) => t.route(),
355 EventKind::Splice(t, _) => t.route(),
356 EventKind::Stat(t, _) => t.route(),
357 }
358 }
359}
360
361#[cfg(test)]
362mod tests;