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 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}