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::{SHUTDOWN, Token, kind};
15use fd::FdSlot;
16use provided::ProvidedLease;
17
18#[derive(Clone, Copy, Debug, PartialEq, Eq)]
19pub struct DecodeError;
20
21impl fmt::Display for DecodeError {
22 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
23 f.write_str("invalid completion")
24 }
25}
26
27impl StdError for DecodeError {}
28
29pub(crate) const BUFFER: u32 = 1 << 0;
30pub(crate) const MORE: u32 = 1 << 1;
31pub(crate) const BUFFER_SHIFT: u32 = 16;
32
33#[derive(Clone, Copy)]
34pub(crate) struct Cqe {
35 user_data: u64,
36 result: i32,
37 flags: u32,
38}
39
40impl Cqe {
41 pub(crate) const fn new(user_data: u64, result: i32, flags: u32) -> Self {
42 Self {
43 user_data,
44 result,
45 flags,
46 }
47 }
48
49 pub(crate) fn kind(self) -> u8 {
50 (self.user_data >> token::KIND_SHIFT) as u8
51 }
52
53 fn more(self) -> bool {
54 self.flags & MORE != 0
55 }
56
57 fn bid_raw(self) -> u16 {
58 (self.flags >> BUFFER_SHIFT) as u16
59 }
60
61 fn has_buffer(self) -> bool {
62 self.flags & BUFFER != 0
63 }
64}
65
66pub enum RecvEvent<'d> {
67 Data(ProvidedLease<'d>),
68 Discarded { len: u32 },
69 Eof,
70 Cancelled,
71 Starved,
72 Failed(i32),
73}
74
75#[derive(Clone, Copy)]
76pub enum SendEvent {
77 Sent(u32),
78 Failed(i32),
79}
80
81#[derive(Clone, Copy)]
82pub enum WriteEvent {
83 Wrote(u32),
84 Failed(i32),
85}
86
87#[derive(Clone, Copy)]
88pub enum SyncEvent {
89 Synced,
90 Failed(i32),
91}
92
93#[derive(Clone, Copy)]
94pub enum OpenEvent {
95 Opened(i32),
96 Failed(i32),
97}
98
99#[derive(Clone, Copy)]
100pub enum ReadEvent {
101 Read(u32),
102 Eof,
103 Failed(i32),
104}
105
106impl ReadEvent {
107 fn from_result(result: i32) -> Self {
108 match result {
109 n if n > 0 => Self::Read(n as u32),
110 0 => Self::Eof,
111 n => Self::Failed(-n),
112 }
113 }
114}
115
116#[derive(Clone, Copy)]
117pub enum SpliceEvent {
118 Moved(u32),
119 Eof,
120 Failed(i32),
121}
122
123#[derive(Clone, Copy)]
124pub enum StatEvent {
125 Done,
126 Failed(i32),
127}
128
129#[derive(Clone, Copy)]
130pub enum AcceptEvent {
131 Accepted(FdSlot),
132 Failed,
133}
134
135pub enum SocketEvent {
136 Created,
137 Failed(io::Error),
138}
139
140pub enum ConnectEvent {
141 Connected,
142 Failed(io::Error),
143}
144
145pub struct Event<'d> {
146 kind: EventKind<'d>,
147 result: i32,
148 operation: u8,
149}
150
151pub enum EventKind<'d> {
152 Accept(Token, bool, AcceptEvent),
153 Recv(Token, bool, RecvEvent<'d>),
154 Send(Token, SendEvent),
155 Timer(Token),
156 Socket(Token, SocketEvent),
157 Connect(Token, ConnectEvent),
158 Write(Token, WriteEvent),
159 Sync(Token, SyncEvent),
160 Open(Token, OpenEvent),
161 Read(Token, ReadEvent),
162 ReadBlock(Token, ReadEvent),
163 Splice(Token, SpliceEvent),
164 Stat(Token, StatEvent),
165 Shutdown,
166}
167
168#[derive(Clone, Copy)]
169enum DecodedRecv {
170 Data { len: u32, bid: u16 },
171 Discarded { len: u32 },
172 Eof,
173 Cancelled,
174 Starved,
175 Failed(i32),
176}
177
178impl DecodedRecv {
179 fn from_errno(result: i32) -> Self {
180 match -result {
181 libc::ECANCELED => Self::Cancelled,
182 libc::ENOBUFS | libc::EAGAIN | libc::EINTR => Self::Starved,
183 errno => Self::Failed(errno),
184 }
185 }
186}
187
188#[derive(Clone, Copy)]
189enum DecodedEvent {
190 Accept(Token, bool, AcceptEvent),
191 Recv(Token, bool, DecodedRecv),
192 Send(Token, SendEvent),
193 Timer(Token),
194 Socket(Token, i32),
195 Connect(Token, i32),
196 Write(Token, WriteEvent),
197 Sync(Token, SyncEvent),
198 Open(Token, OpenEvent),
199 Read(Token, ReadEvent),
200 ReadBlock(Token, ReadEvent),
201 Splice(Token, SpliceEvent),
202 Stat(Token, StatEvent),
203 Shutdown,
204}
205
206impl DecodedEvent {
207 fn decode(c: Cqe) -> Result<Self, DecodeError> {
208 let token = Token::try_from_raw(c.user_data).ok_or(DecodeError)?;
209 if token == SHUTDOWN {
210 return Ok(Self::Shutdown);
211 }
212 match c.kind() {
213 kind::ACCEPT => {
214 let e = match c.result {
215 n if n >= 0 => AcceptEvent::Accepted(FdSlot::new(n as u32)),
216 _ => AcceptEvent::Failed,
217 };
218 Ok(Self::Accept(token, c.more(), e))
219 }
220 kind::RECV => {
221 let e = match c.result {
222 n if n > 0 => {
223 if !c.has_buffer() {
224 debug_assert!(false, "RECV data cqe without buffer flag");
225 return Err(DecodeError);
226 }
227 DecodedRecv::Data {
228 len: n as u32,
229 bid: c.bid_raw(),
230 }
231 }
232 0 => DecodedRecv::Eof,
233 n => DecodedRecv::from_errno(n),
234 };
235 Ok(Self::Recv(token, c.more(), e))
236 }
237 kind::RECV_DISCARD => {
238 let e = match c.result {
239 n if n > 0 => DecodedRecv::Discarded { len: n as u32 },
240 0 => DecodedRecv::Eof,
241 n => DecodedRecv::from_errno(n),
242 };
243 Ok(Self::Recv(token, c.more(), e))
244 }
245 kind::SEND => {
246 let e = if c.result >= 0 {
247 SendEvent::Sent(c.result as u32)
248 } else {
249 SendEvent::Failed(-c.result)
250 };
251 Ok(Self::Send(token, e))
252 }
253 kind::WRITE => {
254 let e = if c.result >= 0 {
255 WriteEvent::Wrote(c.result as u32)
256 } else {
257 WriteEvent::Failed(-c.result)
258 };
259 Ok(Self::Write(token, e))
260 }
261 kind::SYNC => {
262 let e = if c.result >= 0 {
263 SyncEvent::Synced
264 } else {
265 SyncEvent::Failed(-c.result)
266 };
267 Ok(Self::Sync(token, e))
268 }
269 kind::OPEN => {
270 let e = if c.result >= 0 {
271 OpenEvent::Opened(c.result)
272 } else {
273 OpenEvent::Failed(-c.result)
274 };
275 Ok(Self::Open(token, e))
276 }
277 kind::READ => Ok(Self::Read(token, ReadEvent::from_result(c.result))),
278 kind::READ_BLOCK => Ok(Self::ReadBlock(token, ReadEvent::from_result(c.result))),
279 kind::SPLICE => {
280 let e = match c.result {
281 n if n > 0 => SpliceEvent::Moved(n as u32),
282 0 => SpliceEvent::Eof,
283 n => SpliceEvent::Failed(-n),
284 };
285 Ok(Self::Splice(token, e))
286 }
287 kind::STAT => {
288 let e = if c.result >= 0 {
289 StatEvent::Done
290 } else {
291 StatEvent::Failed(-c.result)
292 };
293 Ok(Self::Stat(token, e))
294 }
295 kind::TIMER => Ok(Self::Timer(token)),
296 kind::SOCKET => Ok(Self::Socket(token, c.result)),
297 kind::CONNECT => Ok(Self::Connect(token, c.result)),
298 _ => Err(DecodeError),
299 }
300 }
301}
302
303pub enum EventRef<'a, 'd> {
304 Accept(Token, bool, &'a AcceptEvent),
305 Recv(Token, bool, &'a RecvEvent<'d>),
306 Send(Token, &'a SendEvent),
307 Timer(Token),
308 Socket(Token, &'a SocketEvent),
309 Connect(Token, &'a ConnectEvent),
310 Write(Token, &'a WriteEvent),
311 Sync(Token, &'a SyncEvent),
312 Open(Token, &'a OpenEvent),
313 Read(Token, &'a ReadEvent),
314 Splice(Token, &'a SpliceEvent),
315 Stat(Token, &'a StatEvent),
316 Shutdown,
317}
318
319impl<'d> Event<'d> {
320 pub(crate) fn from_cqe(
321 cqe: Cqe,
322 provided: impl FnOnce(u32, u16) -> ProvidedLease<'d>,
323 ) -> Result<Self, DecodeError> {
324 let result = cqe.result;
325 let operation = cqe.kind();
326 let mut provided = cqe
327 .has_buffer()
328 .then(|| provided(result.max(0) as u32, cqe.bid_raw()));
329 let kind = match DecodedEvent::decode(cqe)? {
330 DecodedEvent::Accept(token, more, event) => EventKind::Accept(token, more, event),
331 DecodedEvent::Recv(token, more, event) => {
332 let event = match event {
333 DecodedRecv::Data { len, bid } => {
334 let lease = provided.take().ok_or(DecodeError)?;
335 debug_assert_eq!(lease.as_slice().len(), len as usize);
336 debug_assert_eq!(bid, cqe.bid_raw());
337 RecvEvent::Data(lease)
338 }
339 DecodedRecv::Discarded { len } => RecvEvent::Discarded { len },
340 DecodedRecv::Eof => RecvEvent::Eof,
341 DecodedRecv::Cancelled => RecvEvent::Cancelled,
342 DecodedRecv::Starved => RecvEvent::Starved,
343 DecodedRecv::Failed(errno) => RecvEvent::Failed(errno),
344 };
345 EventKind::Recv(token, more, event)
346 }
347 DecodedEvent::Send(token, event) => EventKind::Send(token, event),
348 DecodedEvent::Timer(token) => EventKind::Timer(token),
349 DecodedEvent::Socket(token, result) => EventKind::Socket(
350 token,
351 if result >= 0 {
352 SocketEvent::Created
353 } else {
354 SocketEvent::Failed(Error::from_raw_os_error(-result))
355 },
356 ),
357 DecodedEvent::Connect(token, result) => EventKind::Connect(
358 token,
359 if result >= 0 {
360 ConnectEvent::Connected
361 } else {
362 ConnectEvent::Failed(Error::from_raw_os_error(-result))
363 },
364 ),
365 DecodedEvent::Write(token, event) => EventKind::Write(token, event),
366 DecodedEvent::Sync(token, event) => EventKind::Sync(token, event),
367 DecodedEvent::Open(token, event) => EventKind::Open(token, event),
368 DecodedEvent::Read(token, event) => EventKind::Read(token, event),
369 DecodedEvent::ReadBlock(token, event) => EventKind::ReadBlock(token, event),
370 DecodedEvent::Splice(token, event) => EventKind::Splice(token, event),
371 DecodedEvent::Stat(token, event) => EventKind::Stat(token, event),
372 DecodedEvent::Shutdown => EventKind::Shutdown,
373 };
374 Ok(Self {
375 kind,
376 result,
377 operation,
378 })
379 }
380
381 pub fn into_kind(self) -> EventKind<'d> {
382 self.kind
383 }
384
385 pub fn as_ref(&self) -> EventRef<'_, 'd> {
386 match &self.kind {
387 EventKind::Accept(t, more, e) => EventRef::Accept(*t, *more, e),
388 EventKind::Recv(t, more, e) => EventRef::Recv(*t, *more, e),
389 EventKind::Send(t, e) => EventRef::Send(*t, e),
390 EventKind::Timer(t) => EventRef::Timer(*t),
391 EventKind::Socket(t, e) => EventRef::Socket(*t, e),
392 EventKind::Connect(t, e) => EventRef::Connect(*t, e),
393 EventKind::Write(t, e) => EventRef::Write(*t, e),
394 EventKind::Sync(t, e) => EventRef::Sync(*t, e),
395 EventKind::Open(t, e) => EventRef::Open(*t, e),
396 EventKind::Read(t, e) => EventRef::Read(*t, e),
397 EventKind::ReadBlock(t, e) => EventRef::Read(*t, e),
398 EventKind::Splice(t, e) => EventRef::Splice(*t, e),
399 EventKind::Stat(t, e) => EventRef::Stat(*t, e),
400 EventKind::Shutdown => EventRef::Shutdown,
401 }
402 }
403
404 pub const fn result(&self) -> i32 {
405 self.result
406 }
407
408 pub const fn operation(&self) -> u8 {
409 self.operation
410 }
411
412 pub const fn is_shutdown(&self) -> bool {
413 matches!(self.kind, EventKind::Shutdown)
414 }
415
416 pub const fn token(&self) -> Option<Token> {
417 match &self.kind {
418 EventKind::Accept(token, ..)
419 | EventKind::Recv(token, ..)
420 | EventKind::Send(token, _)
421 | EventKind::Timer(token)
422 | EventKind::Socket(token, _)
423 | EventKind::Connect(token, _)
424 | EventKind::Write(token, _)
425 | EventKind::Sync(token, _)
426 | EventKind::Open(token, _)
427 | EventKind::Read(token, _)
428 | EventKind::ReadBlock(token, _)
429 | EventKind::Splice(token, _)
430 | EventKind::Stat(token, _) => Some(*token),
431 EventKind::Shutdown => None,
432 }
433 }
434
435 pub fn route(&self) -> u8 {
436 match &self.kind {
437 EventKind::Accept(t, ..) => t.route(),
438 EventKind::Recv(t, ..) => t.route(),
439 EventKind::Send(t, _) => t.route(),
440 EventKind::Timer(t) => t.route(),
441 EventKind::Socket(t, _) => t.route(),
442 EventKind::Connect(t, _) => t.route(),
443 EventKind::Write(t, _) => t.route(),
444 EventKind::Sync(t, _) => t.route(),
445 EventKind::Open(t, _) => t.route(),
446 EventKind::Read(t, _) => t.route(),
447 EventKind::ReadBlock(t, _) => t.route(),
448 EventKind::Splice(t, _) => t.route(),
449 EventKind::Stat(t, _) => t.route(),
450 EventKind::Shutdown => SHUTDOWN.route(),
451 }
452 }
453}