Skip to main content

air_sys_syscall/io_uring/
net_ops.rs

1// This Source Code Form is subject to the terms of the Mozilla Public
2// License, v. 2.0. If a copy of the MPL was not distributed with this
3// file, You can obtain one at https://mozilla.org/MPL/2.0/.
4
5//! Opérations **réseau** d'io_uring (Temps 2b), par-dessus le cœur Temps 1 et
6//! les conventions Temps 2a (ownership S1, accesseurs typés). Stratégique pour
7//! AirCom (ADR-001) : **zero-copy** du data plane et **passage de FD**
8//! (capabilities) via `sendmsg`/`recvmsg`.
9//!
10//! Référence normative : `docs/specs/layer-0/io-uring-2b-network.md`.
11//!
12//! **Invariants Air** : `MSG_NOSIGNAL` par défaut sur tout envoi (EPIPE en
13//! complétion, jamais de `SIGPIPE`), `SOCK_CLOEXEC`/`CLOEXEC` par défaut sur les
14//! FD créés (`socket`/`accept`) et les FD reçus (`MSG_CMSG_CLOEXEC`).
15//!
16//! **Hors périmètre** : variantes direct-descriptor (`accept_direct`,
17//! `socket_direct`) → Temps 3a ; multishot → Temps 3d ; buffers fournis /
18//! bundle → Temps 3b.
19//!
20//! **Cycle zero-copy à deux complétions** (§4.1) : `send_zc`/`sendmsg_zc`
21//! produisent une CQE **résultat** (`res ≥ 0` + `F_MORE`, sans restitution) puis
22//! une CQE **NOTIF** (`F_NOTIF`) qui restitue le buffer et libère le slot. C'est
23//! l'unique cas où un slot survit à sa première complétion — porté par le
24//! mécanisme `has_more` du slab (Temps 1), sans extension.
25
26use super::owned::{OwnedOp, RecvMsgState, SendMsgState};
27use super::raw::{self, Iovec, Msghdr};
28use super::{IoUring, SubmissionToken};
29use crate::net::{MAX_SOCKADDR_LEN, RawSockaddrStorage, build_scm_rights_cmsg, socket_addr_to_raw};
30use air_sys_types::Errno;
31use air_sys_types::fd::{AsFd, AsRawFd, BorrowedFd};
32use air_sys_types::net::{
33    AcceptFlags, MessageFlags, OwnedReceiveMessage, OwnedSendMessage, ShutdownMode, SocketAddr,
34    SocketDomain, SocketType, ZeroCopyFlags,
35};
36use alloc::boxed::Box;
37use alloc::vec::Vec;
38
39/// `usize` → `u64` saturant (longueurs ; jamais déclenché sur cible LP64 où
40/// `usize == u64` — `unwrap_or` n'introduit pas de branche d'erreur morte).
41fn u64_of(n: usize) -> u64 {
42    u64::try_from(n).unwrap_or(u64::MAX)
43}
44
45/// `usize` → `u32` du champ `len` du SQE, `EINVAL` si débordement (validation
46/// amont, Principe 4 — un buffer réseau > 4 GiB).
47fn len_u32(n: usize) -> Result<u32, Errno> {
48    u32::try_from(n).map_err(|_| Errno::EINVAL)
49}
50
51/// Sérialise une [`SocketAddr`] en stockage `sockaddr` **possédé** (boxé, à
52/// adresse stable) + sa longueur effective, réutilisant la sérialisation de la
53/// famille `net` (DRY).
54fn owned_sockaddr(address: &SocketAddr) -> (Box<RawSockaddrStorage>, u32) {
55    let (storage, len) = socket_addr_to_raw(address);
56    (Box::new(storage), len)
57}
58
59/// Construit l'état `sendmsg` possédé : `msghdr` + `iovec` + adresse + cmsg
60/// `SCM_RIGHTS`, chacun dans une allocation indépendante à adresse stable. Pose
61/// `msg_flags` dans le `msghdr`. Retourne l'état boxé **et** le pointeur du
62/// `msghdr` (capturé avant déplacement dans le slot).
63fn build_send_state(request: OwnedSendMessage, msg_flags: i32) -> (Box<SendMsgState>, u64) {
64    let OwnedSendMessage {
65        mut buffers,
66        address,
67        fds,
68        flags: _,
69    } = request;
70    // iovec depuis les tas (stables) des buffers possédés.
71    let mut iovec_vec: Vec<Iovec> = Vec::with_capacity(buffers.len());
72    for buf in &mut buffers {
73        iovec_vec.push(Iovec {
74            iov_base: buf.as_mut_ptr(),
75            iov_len: buf.len(),
76        });
77    }
78    let iovecs = iovec_vec.into_boxed_slice();
79    let iov_ptr = iovecs.as_ptr() as u64;
80    let iov_len = u64_of(iovecs.len());
81    // Adresse de destination possédée (optionnelle).
82    let addr = address.as_ref().map(owned_sockaddr);
83    let (name_ptr, name_len) = addr
84        .as_ref()
85        .map_or((0u64, 0u32), |(s, l)| (s.bytes.as_ptr() as u64, *l));
86    // cmsg SCM_RIGHTS depuis les FD (le buffer d'octets est possédé/boxé).
87    let fd_borrows: Vec<BorrowedFd<'_>> = fds.iter().map(AsFd::as_fd).collect();
88    let control = build_scm_rights_cmsg(&fd_borrows).into_boxed_slice();
89    let (control_ptr, control_len) = if control.is_empty() {
90        (0u64, 0u64)
91    } else {
92        (control.as_ptr() as u64, u64_of(control.len()))
93    };
94    let msghdr = Box::new(Msghdr {
95        msg_name: name_ptr,
96        msg_namelen: name_len,
97        pad1: 0,
98        msg_iov: iov_ptr,
99        msg_iovlen: iov_len,
100        msg_control: control_ptr,
101        msg_controllen: control_len,
102        msg_flags,
103        pad2: 0,
104    });
105    let hdr_ptr = core::ptr::from_ref(&*msghdr) as u64;
106    let state = Box::new(SendMsgState {
107        msghdr,
108        iovecs,
109        buffers,
110        addr: addr.map(|(s, _)| s),
111        control,
112        fds,
113    });
114    (state, hdr_ptr)
115}
116
117impl IoUring {
118    // ── Établissement de connexion ────────────────────────────────────────
119
120    /// Soumet un `IORING_OP_ACCEPT` (13) sans capturer l'adresse pair.
121    /// `CLOEXEC` est ajouté d'office. Complétion via
122    /// [`Completion::accepted_fd`](super::Completion::accepted_fd).
123    ///
124    /// # Errors
125    ///
126    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
127    pub fn submit_accept(
128        &mut self,
129        listener: BorrowedFd<'_>,
130        flags: AcceptFlags,
131    ) -> Result<SubmissionToken, Errno> {
132        let fd = listener.as_raw_fd();
133        let accept_flags = (flags | AcceptFlags::CLOEXEC).bits().cast_unsigned();
134        self.submit_op(None, |sqe| {
135            sqe.opcode = raw::IORING_OP_ACCEPT;
136            sqe.fd = fd;
137            sqe.op_flags = accept_flags;
138        })
139    }
140
141    /// Soumet un `IORING_OP_ACCEPT` (13) en capturant l'adresse pair dans un
142    /// stockage possédé du slot. Complétion via
143    /// [`Completion::into_accept_result`](super::Completion::into_accept_result).
144    ///
145    /// # Errors
146    ///
147    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
148    pub fn submit_accept_with_peer(
149        &mut self,
150        listener: BorrowedFd<'_>,
151        flags: AcceptFlags,
152    ) -> Result<SubmissionToken, Errno> {
153        let fd = listener.as_raw_fd();
154        let accept_flags = (flags | AcceptFlags::CLOEXEC).bits().cast_unsigned();
155        let addr = Box::new(RawSockaddrStorage {
156            bytes: [0u8; MAX_SOCKADDR_LEN],
157        });
158        let addrlen = Box::new(u32::try_from(MAX_SOCKADDR_LEN).unwrap_or(0));
159        let addr_ptr = addr.bytes.as_ptr() as u64;
160        let addrlen_ptr = core::ptr::from_ref(&*addrlen) as u64;
161        self.submit_op(Some(OwnedOp::Accept { addr, addrlen }), |sqe| {
162            sqe.opcode = raw::IORING_OP_ACCEPT;
163            sqe.fd = fd;
164            sqe.addr_or_splice_off_in = addr_ptr;
165            sqe.off_or_addr2 = addrlen_ptr;
166            sqe.op_flags = accept_flags;
167        })
168    }
169
170    /// Soumet un `IORING_OP_CONNECT` (16). L'adresse est sérialisée et déplacée
171    /// dans le slot. Complétion via
172    /// [`Completion::completed`](super::Completion::completed).
173    ///
174    /// # Errors
175    ///
176    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
177    pub fn submit_connect(
178        &mut self,
179        sock: BorrowedFd<'_>,
180        address: SocketAddr,
181    ) -> Result<SubmissionToken, Errno> {
182        let fd = sock.as_raw_fd();
183        let (storage, len) = owned_sockaddr(&address);
184        let addr_ptr = storage.bytes.as_ptr() as u64;
185        self.submit_op(Some(OwnedOp::SockAddr(storage)), |sqe| {
186            sqe.opcode = raw::IORING_OP_CONNECT;
187            sqe.fd = fd;
188            sqe.addr_or_splice_off_in = addr_ptr;
189            sqe.off_or_addr2 = u64::from(len);
190        })
191    }
192
193    /// Soumet un `IORING_OP_SOCKET` (45). `SOCK_CLOEXEC` est ajouté d'office.
194    /// Complétion via
195    /// [`Completion::into_socket_fd`](super::Completion::into_socket_fd).
196    ///
197    /// # Errors
198    ///
199    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
200    pub fn submit_socket(
201        &mut self,
202        domain: SocketDomain,
203        ty: SocketType,
204        protocol: i32,
205    ) -> Result<SubmissionToken, Errno> {
206        let domain = domain as i32;
207        let ty = (ty as i32) | raw::SOCK_CLOEXEC;
208        let protocol = protocol.cast_unsigned();
209        self.submit_op(None, |sqe| {
210            sqe.opcode = raw::IORING_OP_SOCKET;
211            sqe.fd = domain;
212            sqe.off_or_addr2 = u64::from(ty.cast_unsigned());
213            sqe.len = protocol;
214        })
215    }
216
217    /// Soumet un `IORING_OP_BIND` (56, kernel ≥ 6.11). Complétion via `completed`.
218    ///
219    /// # Errors
220    ///
221    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
222    pub fn submit_bind(
223        &mut self,
224        sock: BorrowedFd<'_>,
225        address: SocketAddr,
226    ) -> Result<SubmissionToken, Errno> {
227        let fd = sock.as_raw_fd();
228        let (storage, len) = owned_sockaddr(&address);
229        let addr_ptr = storage.bytes.as_ptr() as u64;
230        self.submit_op(Some(OwnedOp::SockAddr(storage)), |sqe| {
231            sqe.opcode = raw::IORING_OP_BIND;
232            sqe.fd = fd;
233            sqe.addr_or_splice_off_in = addr_ptr;
234            sqe.off_or_addr2 = u64::from(len);
235        })
236    }
237
238    /// Soumet un `IORING_OP_LISTEN` (57, kernel ≥ 6.11). Complétion via
239    /// `completed`.
240    ///
241    /// # Errors
242    ///
243    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
244    pub fn submit_listen(
245        &mut self,
246        sock: BorrowedFd<'_>,
247        backlog: u32,
248    ) -> Result<SubmissionToken, Errno> {
249        let fd = sock.as_raw_fd();
250        self.submit_op(None, |sqe| {
251            sqe.opcode = raw::IORING_OP_LISTEN;
252            sqe.fd = fd;
253            sqe.len = backlog;
254        })
255    }
256
257    // ── Transfert one-shot ────────────────────────────────────────────────
258
259    /// Soumet un `IORING_OP_SEND` (26). `MSG_NOSIGNAL` ajouté par défaut. Buffer
260    /// déplacé dans le slot ; complétion via
261    /// [`Completion::into_buffer_result`](super::Completion::into_buffer_result).
262    ///
263    /// # Errors
264    ///
265    /// [`Errno::EBUSY`] (SQ/slab plein) ; [`Errno::EINVAL`] si `buffer.len()`
266    /// déborde un `u32`.
267    pub fn submit_send(
268        &mut self,
269        sock: BorrowedFd<'_>,
270        mut buffer: Vec<u8>,
271        flags: MessageFlags,
272    ) -> Result<SubmissionToken, Errno> {
273        let fd = sock.as_raw_fd();
274        let len = len_u32(buffer.len())?;
275        let addr = buffer.as_mut_ptr() as u64;
276        let msg_flags = (flags | MessageFlags::NOSIGNAL).bits().cast_unsigned();
277        self.submit_op(Some(OwnedOp::Bytes(buffer)), |sqe| {
278            sqe.opcode = raw::IORING_OP_SEND;
279            sqe.fd = fd;
280            sqe.addr_or_splice_off_in = addr;
281            sqe.len = len;
282            sqe.op_flags = msg_flags;
283        })
284    }
285
286    /// Soumet un `IORING_OP_RECV` (27). Buffer déplacé dans le slot ; complétion
287    /// via `into_buffer_result`. `socket_has_pending_data()` (CQE_F_SOCK_NONEMPTY)
288    /// indique s'il reste des données à lire.
289    ///
290    /// # Errors
291    ///
292    /// [`Errno::EBUSY`] (SQ/slab plein) ; [`Errno::EINVAL`] si `buffer.len()`
293    /// déborde un `u32`.
294    pub fn submit_receive(
295        &mut self,
296        sock: BorrowedFd<'_>,
297        mut buffer: Vec<u8>,
298        flags: MessageFlags,
299    ) -> Result<SubmissionToken, Errno> {
300        let fd = sock.as_raw_fd();
301        let len = len_u32(buffer.len())?;
302        let addr = buffer.as_mut_ptr() as u64;
303        let msg_flags = flags.bits().cast_unsigned();
304        self.submit_op(Some(OwnedOp::Bytes(buffer)), |sqe| {
305            sqe.opcode = raw::IORING_OP_RECV;
306            sqe.fd = fd;
307            sqe.addr_or_splice_off_in = addr;
308            sqe.len = len;
309            sqe.op_flags = msg_flags;
310        })
311    }
312
313    /// Soumet un `IORING_OP_SENDMSG` (9) : envoi vectorisé + **passage de FD**
314    /// (`SCM_RIGHTS`). `MSG_NOSIGNAL` par défaut. Tout l'état est déplacé dans le
315    /// slot. Complétion via
316    /// [`Completion::into_result`](super::Completion::into_result) (octets).
317    ///
318    /// # Errors
319    ///
320    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
321    pub fn submit_send_message(
322        &mut self,
323        sock: BorrowedFd<'_>,
324        request: OwnedSendMessage,
325    ) -> Result<SubmissionToken, Errno> {
326        let fd = sock.as_raw_fd();
327        let msg_flags = (request.flags | MessageFlags::NOSIGNAL).bits();
328        let (state, hdr_ptr) = build_send_state(request, msg_flags);
329        self.submit_op(Some(OwnedOp::SendMsg(state)), |sqe| {
330            sqe.opcode = raw::IORING_OP_SENDMSG;
331            sqe.fd = fd;
332            sqe.addr_or_splice_off_in = hdr_ptr;
333            sqe.len = 1;
334        })
335    }
336
337    /// Soumet un `IORING_OP_RECVMSG` (10) : réception vectorisée + **réception de
338    /// FD** (`SCM_RIGHTS`). `MSG_CMSG_CLOEXEC` par défaut (FD reçus CLOEXEC).
339    /// Complétion via [`Completion::into_receive_message_result`](super::Completion::into_receive_message_result).
340    ///
341    /// # Errors
342    ///
343    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
344    pub fn submit_receive_message(
345        &mut self,
346        sock: BorrowedFd<'_>,
347        request: OwnedReceiveMessage,
348    ) -> Result<SubmissionToken, Errno> {
349        let fd = sock.as_raw_fd();
350        let msg_flags = (request.flags | MessageFlags::CMSG_CLOEXEC).bits();
351        let OwnedReceiveMessage {
352            mut buffers,
353            control_capacity,
354            flags: _,
355        } = request;
356        let mut iovec_vec: Vec<Iovec> = Vec::with_capacity(buffers.len());
357        for buf in &mut buffers {
358            iovec_vec.push(Iovec {
359                iov_base: buf.as_mut_ptr(),
360                iov_len: buf.len(),
361            });
362        }
363        let iovecs = iovec_vec.into_boxed_slice();
364        let iov_ptr = iovecs.as_ptr() as u64;
365        let iov_len = u64_of(iovecs.len());
366        let name = Box::new(RawSockaddrStorage {
367            bytes: [0u8; MAX_SOCKADDR_LEN],
368        });
369        let name_ptr = name.bytes.as_ptr() as u64;
370        let control = vec![0u8; control_capacity].into_boxed_slice();
371        let (control_ptr, control_len) = if control.is_empty() {
372            (0u64, 0u64)
373        } else {
374            (control.as_ptr() as u64, u64_of(control.len()))
375        };
376        let name_capacity = u32::try_from(MAX_SOCKADDR_LEN).unwrap_or(0);
377        let msghdr = Box::new(Msghdr {
378            msg_name: name_ptr,
379            msg_namelen: name_capacity,
380            pad1: 0,
381            msg_iov: iov_ptr,
382            msg_iovlen: iov_len,
383            msg_control: control_ptr,
384            msg_controllen: control_len,
385            msg_flags,
386            pad2: 0,
387        });
388        let hdr_ptr = core::ptr::from_ref(&*msghdr) as u64;
389        let state = Box::new(RecvMsgState {
390            msghdr,
391            iovecs,
392            buffers,
393            name,
394            control,
395        });
396        self.submit_op(Some(OwnedOp::RecvMsg(state)), |sqe| {
397            sqe.opcode = raw::IORING_OP_RECVMSG;
398            sqe.fd = fd;
399            sqe.addr_or_splice_off_in = hdr_ptr;
400            sqe.len = 1;
401        })
402    }
403
404    // ── Zero-copy (data plane AirCom) ─────────────────────────────────────
405
406    /// Soumet un `IORING_OP_SEND_ZC` (47). **Deux complétions** (§4.1) : résultat
407    /// (`F_MORE`, `into_result`) puis NOTIF (`into_zero_copy_buffer`). Le buffer
408    /// reste vivant **jusqu'au NOTIF**. `MSG_NOSIGNAL` par défaut.
409    ///
410    /// # Errors
411    ///
412    /// [`Errno::EBUSY`] (SQ/slab plein) ; [`Errno::EINVAL`] si `buffer.len()`
413    /// déborde un `u32`.
414    pub fn submit_send_zero_copy(
415        &mut self,
416        sock: BorrowedFd<'_>,
417        mut buffer: Vec<u8>,
418        flags: MessageFlags,
419        zero_copy: ZeroCopyFlags,
420    ) -> Result<SubmissionToken, Errno> {
421        let fd = sock.as_raw_fd();
422        let len = len_u32(buffer.len())?;
423        let addr = buffer.as_mut_ptr() as u64;
424        let msg_flags = (flags | MessageFlags::NOSIGNAL).bits().cast_unsigned();
425        let zc = u16::try_from(zero_copy.bits()).unwrap_or(raw::IORING_SEND_ZC_REPORT_USAGE);
426        self.submit_op(Some(OwnedOp::Bytes(buffer)), |sqe| {
427            sqe.opcode = raw::IORING_OP_SEND_ZC;
428            sqe.fd = fd;
429            sqe.addr_or_splice_off_in = addr;
430            sqe.len = len;
431            sqe.op_flags = msg_flags;
432            sqe.ioprio = zc;
433        })
434    }
435
436    /// Soumet un `IORING_OP_SENDMSG_ZC` (48) : `sendmsg` zero-copy (vectorisé +
437    /// FD passing). Deux complétions (§4.1). `MSG_NOSIGNAL` par défaut.
438    ///
439    /// # Errors
440    ///
441    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
442    pub fn submit_send_message_zero_copy(
443        &mut self,
444        sock: BorrowedFd<'_>,
445        request: OwnedSendMessage,
446        zero_copy: ZeroCopyFlags,
447    ) -> Result<SubmissionToken, Errno> {
448        let fd = sock.as_raw_fd();
449        let msg_flags = (request.flags | MessageFlags::NOSIGNAL).bits();
450        let (state, hdr_ptr) = build_send_state(request, msg_flags);
451        let zc = u16::try_from(zero_copy.bits()).unwrap_or(raw::IORING_SEND_ZC_REPORT_USAGE);
452        self.submit_op(Some(OwnedOp::SendMsg(state)), |sqe| {
453            sqe.opcode = raw::IORING_OP_SENDMSG_ZC;
454            sqe.fd = fd;
455            sqe.addr_or_splice_off_in = hdr_ptr;
456            sqe.len = 1;
457            sqe.ioprio = zc;
458        })
459    }
460
461    // ── Fermeture ─────────────────────────────────────────────────────────
462
463    /// Soumet un `IORING_OP_SHUTDOWN` (34). Complétion via `completed`.
464    ///
465    /// # Errors
466    ///
467    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
468    pub fn submit_shutdown(
469        &mut self,
470        sock: BorrowedFd<'_>,
471        how: ShutdownMode,
472    ) -> Result<SubmissionToken, Errno> {
473        let fd = sock.as_raw_fd();
474        let how = (how as i32).cast_unsigned();
475        self.submit_op(None, |sqe| {
476            sqe.opcode = raw::IORING_OP_SHUTDOWN;
477            sqe.fd = fd;
478            sqe.len = how;
479        })
480    }
481}
482
483// ───────────────────────────────────────────────────────────────────────────
484// Tests
485// ───────────────────────────────────────────────────────────────────────────
486//
487// Intégration (kernel 6.12 réel : AF_UNIX socketpair + AF_UNIX path) :
488// `#[cfg_attr(miri, ignore)]` (Miri ne modélise pas `io_uring_enter`). Les tests
489// purs (accesseurs, ownership zero-copy/recvmsg) tournent SOUS Miri.
490
491#[cfg(test)]
492mod tests;