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 StatEvent {
118 Done,
119 Failed(i32),
120}
121
122#[derive(Clone, Copy)]
123pub enum AcceptEvent {
124 Accepted(FdSlot),
125 Failed,
126}
127
128pub enum SocketEvent {
129 Created,
130 Failed(io::Error),
131}
132
133pub enum ConnectEvent {
134 Connected,
135 Failed(io::Error),
136}
137
138pub struct Event<'d> {
139 kind: EventKind<'d>,
140 result: i32,
141 operation: u8,
142}
143
144pub enum EventKind<'d> {
145 Accept(Token, bool, AcceptEvent),
146 Recv(Token, bool, RecvEvent<'d>),
147 Send(Token, SendEvent),
148 Timer(Token),
149 Socket(Token, SocketEvent),
150 Connect(Token, ConnectEvent),
151 Write(Token, WriteEvent),
152 Sync(Token, SyncEvent),
153 Open(Token, OpenEvent),
154 Read(Token, ReadEvent),
155 Stat(Token, StatEvent),
156 Shutdown,
157}
158
159#[derive(Clone, Copy)]
160enum DecodedRecv {
161 Data { len: u32, bid: u16 },
162 Discarded { len: u32 },
163 Eof,
164 Cancelled,
165 Starved,
166 Failed(i32),
167}
168
169impl DecodedRecv {
170 fn from_errno(result: i32) -> Self {
171 match -result {
172 libc::ECANCELED => Self::Cancelled,
173 libc::ENOBUFS | libc::EAGAIN | libc::EINTR => Self::Starved,
174 errno => Self::Failed(errno),
175 }
176 }
177}
178
179#[derive(Clone, Copy)]
180enum DecodedEvent {
181 Accept(Token, bool, AcceptEvent),
182 Recv(Token, bool, DecodedRecv),
183 Send(Token, SendEvent),
184 Timer(Token),
185 Socket(Token, i32),
186 Connect(Token, i32),
187 Write(Token, WriteEvent),
188 Sync(Token, SyncEvent),
189 Open(Token, OpenEvent),
190 Read(Token, ReadEvent),
191 Stat(Token, StatEvent),
192 Shutdown,
193}
194
195impl DecodedEvent {
196 fn decode(c: Cqe) -> Result<Self, DecodeError> {
197 let token = Token::try_from_raw(c.user_data).ok_or(DecodeError)?;
198 if token == SHUTDOWN {
199 return Ok(Self::Shutdown);
200 }
201 match c.kind() {
202 kind::ACCEPT => {
203 let e = match c.result {
204 n if n >= 0 => AcceptEvent::Accepted(FdSlot::new(n as u32)),
205 _ => AcceptEvent::Failed,
206 };
207 Ok(Self::Accept(token, c.more(), e))
208 }
209 kind::RECV => {
210 let e = match c.result {
211 n if n > 0 => {
212 if !c.has_buffer() {
213 debug_assert!(false, "RECV data cqe without buffer flag");
214 return Err(DecodeError);
215 }
216 DecodedRecv::Data {
217 len: n as u32,
218 bid: c.bid_raw(),
219 }
220 }
221 0 => DecodedRecv::Eof,
222 n => DecodedRecv::from_errno(n),
223 };
224 Ok(Self::Recv(token, c.more(), e))
225 }
226 kind::RECV_DISCARD => {
227 let e = match c.result {
228 n if n > 0 => DecodedRecv::Discarded { len: n as u32 },
229 0 => DecodedRecv::Eof,
230 n => DecodedRecv::from_errno(n),
231 };
232 Ok(Self::Recv(token, c.more(), e))
233 }
234 kind::SEND => {
235 let e = if c.result >= 0 {
236 SendEvent::Sent(c.result as u32)
237 } else {
238 SendEvent::Failed(-c.result)
239 };
240 Ok(Self::Send(token, e))
241 }
242 kind::WRITE => {
243 let e = if c.result >= 0 {
244 WriteEvent::Wrote(c.result as u32)
245 } else {
246 WriteEvent::Failed(-c.result)
247 };
248 Ok(Self::Write(token, e))
249 }
250 kind::SYNC => {
251 let e = if c.result >= 0 {
252 SyncEvent::Synced
253 } else {
254 SyncEvent::Failed(-c.result)
255 };
256 Ok(Self::Sync(token, e))
257 }
258 kind::OPEN => {
259 let e = if c.result >= 0 {
260 OpenEvent::Opened(c.result)
261 } else {
262 OpenEvent::Failed(-c.result)
263 };
264 Ok(Self::Open(token, e))
265 }
266 kind::READ => Ok(Self::Read(token, ReadEvent::from_result(c.result))),
267 kind::STAT => {
268 let e = if c.result >= 0 {
269 StatEvent::Done
270 } else {
271 StatEvent::Failed(-c.result)
272 };
273 Ok(Self::Stat(token, e))
274 }
275 kind::TIMER => Ok(Self::Timer(token)),
276 kind::SOCKET => Ok(Self::Socket(token, c.result)),
277 kind::CONNECT => Ok(Self::Connect(token, c.result)),
278 _ => Err(DecodeError),
279 }
280 }
281}
282
283pub enum EventRef<'a, 'd> {
284 Accept(Token, bool, &'a AcceptEvent),
285 Recv(Token, bool, &'a RecvEvent<'d>),
286 Send(Token, &'a SendEvent),
287 Timer(Token),
288 Socket(Token, &'a SocketEvent),
289 Connect(Token, &'a ConnectEvent),
290 Write(Token, &'a WriteEvent),
291 Sync(Token, &'a SyncEvent),
292 Open(Token, &'a OpenEvent),
293 Read(Token, &'a ReadEvent),
294 Stat(Token, &'a StatEvent),
295 Shutdown,
296}
297
298impl<'d> Event<'d> {
299 pub(crate) fn from_cqe(
300 cqe: Cqe,
301 provided: impl FnOnce(u32, u16) -> ProvidedLease<'d>,
302 ) -> Result<Self, DecodeError> {
303 let result = cqe.result;
304 let operation = cqe.kind();
305 let mut provided = cqe
306 .has_buffer()
307 .then(|| provided(result.max(0) as u32, cqe.bid_raw()));
308 let kind = match DecodedEvent::decode(cqe)? {
309 DecodedEvent::Accept(token, more, event) => EventKind::Accept(token, more, event),
310 DecodedEvent::Recv(token, more, event) => {
311 let event = match event {
312 DecodedRecv::Data { len, bid } => {
313 let lease = provided.take().ok_or(DecodeError)?;
314 debug_assert_eq!(lease.as_slice().len(), len as usize);
315 debug_assert_eq!(bid, cqe.bid_raw());
316 RecvEvent::Data(lease)
317 }
318 DecodedRecv::Discarded { len } => RecvEvent::Discarded { len },
319 DecodedRecv::Eof => RecvEvent::Eof,
320 DecodedRecv::Cancelled => RecvEvent::Cancelled,
321 DecodedRecv::Starved => RecvEvent::Starved,
322 DecodedRecv::Failed(errno) => RecvEvent::Failed(errno),
323 };
324 EventKind::Recv(token, more, event)
325 }
326 DecodedEvent::Send(token, event) => EventKind::Send(token, event),
327 DecodedEvent::Timer(token) => EventKind::Timer(token),
328 DecodedEvent::Socket(token, result) => EventKind::Socket(
329 token,
330 if result >= 0 {
331 SocketEvent::Created
332 } else {
333 SocketEvent::Failed(Error::from_raw_os_error(-result))
334 },
335 ),
336 DecodedEvent::Connect(token, result) => EventKind::Connect(
337 token,
338 if result >= 0 {
339 ConnectEvent::Connected
340 } else {
341 ConnectEvent::Failed(Error::from_raw_os_error(-result))
342 },
343 ),
344 DecodedEvent::Write(token, event) => EventKind::Write(token, event),
345 DecodedEvent::Sync(token, event) => EventKind::Sync(token, event),
346 DecodedEvent::Open(token, event) => EventKind::Open(token, event),
347 DecodedEvent::Read(token, event) => EventKind::Read(token, event),
348 DecodedEvent::Stat(token, event) => EventKind::Stat(token, event),
349 DecodedEvent::Shutdown => EventKind::Shutdown,
350 };
351 Ok(Self {
352 kind,
353 result,
354 operation,
355 })
356 }
357
358 pub fn into_kind(self) -> EventKind<'d> {
359 self.kind
360 }
361
362 pub fn as_ref(&self) -> EventRef<'_, 'd> {
363 match &self.kind {
364 EventKind::Accept(t, more, e) => EventRef::Accept(*t, *more, e),
365 EventKind::Recv(t, more, e) => EventRef::Recv(*t, *more, e),
366 EventKind::Send(t, e) => EventRef::Send(*t, e),
367 EventKind::Timer(t) => EventRef::Timer(*t),
368 EventKind::Socket(t, e) => EventRef::Socket(*t, e),
369 EventKind::Connect(t, e) => EventRef::Connect(*t, e),
370 EventKind::Write(t, e) => EventRef::Write(*t, e),
371 EventKind::Sync(t, e) => EventRef::Sync(*t, e),
372 EventKind::Open(t, e) => EventRef::Open(*t, e),
373 EventKind::Read(t, e) => EventRef::Read(*t, e),
374 EventKind::Stat(t, e) => EventRef::Stat(*t, e),
375 EventKind::Shutdown => EventRef::Shutdown,
376 }
377 }
378
379 pub const fn result(&self) -> i32 {
380 self.result
381 }
382
383 pub const fn operation(&self) -> u8 {
384 self.operation
385 }
386
387 pub const fn is_shutdown(&self) -> bool {
388 matches!(self.kind, EventKind::Shutdown)
389 }
390
391 pub const fn token(&self) -> Option<Token> {
392 match &self.kind {
393 EventKind::Accept(token, ..)
394 | EventKind::Recv(token, ..)
395 | EventKind::Send(token, _)
396 | EventKind::Timer(token)
397 | EventKind::Socket(token, _)
398 | EventKind::Connect(token, _)
399 | EventKind::Write(token, _)
400 | EventKind::Sync(token, _)
401 | EventKind::Open(token, _)
402 | EventKind::Read(token, _)
403 | EventKind::Stat(token, _) => Some(*token),
404 EventKind::Shutdown => None,
405 }
406 }
407
408 pub fn route(&self) -> u8 {
409 match &self.kind {
410 EventKind::Accept(t, ..) => t.route(),
411 EventKind::Recv(t, ..) => t.route(),
412 EventKind::Send(t, _) => t.route(),
413 EventKind::Timer(t) => t.route(),
414 EventKind::Socket(t, _) => t.route(),
415 EventKind::Connect(t, _) => t.route(),
416 EventKind::Write(t, _) => t.route(),
417 EventKind::Sync(t, _) => t.route(),
418 EventKind::Open(t, _) => t.route(),
419 EventKind::Read(t, _) => t.route(),
420 EventKind::Stat(t, _) => t.route(),
421 EventKind::Shutdown => SHUTDOWN.route(),
422 }
423 }
424}