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::{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    /// # Safety
315    /// `cqe` was produced by the paired driver and has not been decoded before.
316    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;