Skip to main content

dope_core/io/
mod.rs

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}