Skip to main content

air_sys_syscall/io_uring/
async_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 **async-spécifiques** d'io_uring (Temps 2c) : contrôle du ring
6//! (`nop`), temporisation (`timeout`), annulation asynchrone (`cancel`),
7//! surveillance d'événements (`poll`, `epoll_ctl`), attente de processus
8//! (`waitid`), notification inter-ring (`msg_ring` data). Par-dessus le cœur
9//! Temps 1, conventions Temps 2a/2b (ownership S1, accesseurs typés).
10//!
11//! Référence normative : `docs/specs/layer-0/io-uring-2c-async.md`.
12//!
13//! **`futex_wait`/`wake`/`waitv` (51/52/53)** sont désormais **implémentés**
14//! (PR coordonnée `family-mem`) : le mot futex vit dans une [`MmapRegion`]
15//! partageable et le slot S1 retient sa garde de vivacité — la signature
16//! historique `&AtomicU32` (insoundable en async) est remplacée par
17//! `&MmapRegion` + offset (cf. `family-mem-mmap-region.md §3`).
18//!
19//! **Différé (signalé)** :
20//! - `link_timeout` (15) → **Temps 3c** (n'a de sens qu'attaché à une chaîne
21//!   liée via le `LinkedChainBuilder` ; jamais en op isolée).
22//! - `files_update` (20), `fixed_fd_install` (54), `msg_ring_fd` (40
23//!   `MSG_SEND_FD`) → **Temps 3a** (table de FD enregistrée / descripteurs
24//!   directs).
25
26use super::owned::OwnedOp;
27use super::raw::{self, FutexWaitvRaw, KernelTimespec};
28use super::{CancelTarget, IoUring, SubmissionToken, opcode_number};
29use crate::mem::MmapRegion;
30use air_sys_types::Errno;
31use air_sys_types::fd::{AsRawFd, BorrowedFd};
32use air_sys_types::io_uring::{
33    CancelFlags, EpollEvent, EpollOp, FutexFlags, MessageRingFlags, PollEvents, TimeoutFlags,
34    TimeoutSpec,
35};
36use air_sys_types::process::{WaitOptions, WaitTarget};
37use air_sys_types::signal::SignalInfo;
38use alloc::boxed::Box;
39use alloc::vec::Vec;
40use core::time::Duration;
41
42/// Construit un `__kernel_timespec` depuis la durée d'un [`TimeoutSpec`]
43/// (`None` ⇒ zéro, le timeout est alors purement par comptage).
44pub(crate) fn timespec_of(spec: TimeoutSpec) -> KernelTimespec {
45    let d = spec.duration.unwrap_or(Duration::ZERO);
46    KernelTimespec {
47        tv_sec: i64::try_from(d.as_secs()).unwrap_or(i64::MAX),
48        tv_nsec: i64::from(d.subsec_nanos()),
49    }
50}
51
52/// `(idtype, id)` kernel d'une [`WaitTarget`] (sentinelles typées, ADR-021).
53fn waitid_target(target: WaitTarget<'_>) -> (i32, i32) {
54    match target {
55        WaitTarget::AnyChild => (raw::P_ALL, 0),
56        WaitTarget::Pid(p) => (raw::P_PID, p.as_raw()),
57        WaitTarget::ProcessGroup(p) => (raw::P_PGID, p.as_raw()),
58        WaitTarget::AnyProcessGroup => (raw::P_PGID, 0),
59        WaitTarget::PidFd(fd) => (raw::P_PIDFD, fd.as_raw_fd()),
60    }
61}
62
63impl IoUring {
64    // ── Contrôle : nop ────────────────────────────────────────────────────
65
66    /// Soumet un `IORING_OP_NOP` (0) **injectant un résultat** (`res = injected`)
67    /// via `IORING_NOP_INJECT_RESULT`. Précieux pour exercer les chemins
68    /// d'erreur des accesseurs **sans** provoquer de vraie erreur kernel.
69    /// Complétion via [`Completion::into_result`](super::Completion::into_result).
70    ///
71    /// # Errors
72    ///
73    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
74    pub fn submit_nop_with_result(&mut self, injected: i32) -> Result<SubmissionToken, Errno> {
75        let len = injected.cast_unsigned();
76        self.submit_op(None, |sqe| {
77            sqe.opcode = raw::IORING_OP_NOP;
78            sqe.len = len;
79            sqe.op_flags = raw::IORING_NOP_INJECT_RESULT;
80        })
81    }
82
83    // ── Temporisation ─────────────────────────────────────────────────────
84
85    /// Soumet un `IORING_OP_TIMEOUT` (11) : échéance temporelle et/ou seuil de
86    /// complétions. `res == -ETIME` à l'expiration (succès si
87    /// `TimeoutFlags::ETIME_SUCCESS`). Complétion via
88    /// [`Completion::completed`](super::Completion::completed).
89    ///
90    /// # Errors
91    ///
92    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
93    pub fn submit_timeout(
94        &mut self,
95        spec: TimeoutSpec,
96        flags: TimeoutFlags,
97    ) -> Result<SubmissionToken, Errno> {
98        let ts = Box::new(timespec_of(spec));
99        let ts_ptr = core::ptr::from_ref(&*ts) as u64;
100        let count = u64::from(spec.count);
101        let op_flags = flags.bits();
102        self.submit_op(Some(OwnedOp::Timeout(ts)), |sqe| {
103            sqe.opcode = raw::IORING_OP_TIMEOUT;
104            sqe.fd = -1;
105            sqe.addr_or_splice_off_in = ts_ptr;
106            sqe.len = 1;
107            sqe.off_or_addr2 = count;
108            sqe.op_flags = op_flags;
109        })
110    }
111
112    /// Soumet un `IORING_OP_TIMEOUT_REMOVE` (12) : annule un timeout en vol
113    /// désigné par son jeton. Complétion via `completed` (`ENOENT` si la cible
114    /// est inconnue/déjà déclenchée).
115    ///
116    /// # Errors
117    ///
118    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
119    pub fn submit_timeout_remove(
120        &mut self,
121        target: SubmissionToken,
122    ) -> Result<SubmissionToken, Errno> {
123        let addr = target.to_user_data();
124        self.submit_op(None, |sqe| {
125            sqe.opcode = raw::IORING_OP_TIMEOUT_REMOVE;
126            sqe.fd = -1;
127            sqe.addr_or_splice_off_in = addr;
128        })
129    }
130
131    /// Soumet un `IORING_OP_TIMEOUT_REMOVE` + `IORING_TIMEOUT_UPDATE` (12) :
132    /// re-arme un timeout en vol avec un nouveau `spec`. Complétion via
133    /// `completed`.
134    ///
135    /// # Errors
136    ///
137    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
138    pub fn submit_timeout_update(
139        &mut self,
140        target: SubmissionToken,
141        spec: TimeoutSpec,
142        flags: TimeoutFlags,
143    ) -> Result<SubmissionToken, Errno> {
144        let ts = Box::new(timespec_of(spec));
145        let ts_ptr = core::ptr::from_ref(&*ts) as u64;
146        let addr = target.to_user_data();
147        let op_flags = flags.bits() | raw::IORING_TIMEOUT_UPDATE;
148        self.submit_op(Some(OwnedOp::Timeout(ts)), |sqe| {
149            sqe.opcode = raw::IORING_OP_TIMEOUT_REMOVE;
150            sqe.fd = -1;
151            sqe.addr_or_splice_off_in = addr;
152            sqe.off_or_addr2 = ts_ptr;
153            sqe.op_flags = op_flags;
154        })
155    }
156
157    // ── Annulation asynchrone ─────────────────────────────────────────────
158
159    /// Soumet un `IORING_OP_ASYNC_CANCEL` (14) : variante **asynchrone** (dans le
160    /// flux) du `sync_cancel` du Temps 1. La complétion de l'op annulée arrive en
161    /// `-ECANCELED`, son slot S1 (et son buffer) restitué normalement.
162    /// Complétion via `into_result` (= nombre d'opérations annulées ; `-ENOENT`
163    /// si rien ne correspond, `-EALREADY` si déjà en cours).
164    ///
165    /// # Errors
166    ///
167    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
168    pub fn submit_cancel(
169        &mut self,
170        target: CancelTarget<'_>,
171        flags: CancelFlags,
172    ) -> Result<SubmissionToken, Errno> {
173        let mut fd = -1;
174        let mut addr = 0u64;
175        let mut len = 0u32;
176        let mut cancel_flags = flags.bits();
177        match target {
178            CancelTarget::Token(token) => addr = token.to_user_data(),
179            CancelTarget::Fd(target_fd) => {
180                fd = target_fd.as_raw_fd();
181                cancel_flags |= raw::IORING_ASYNC_CANCEL_FD;
182            }
183            CancelTarget::Op(op) => {
184                len = u32::from(opcode_number(op));
185                cancel_flags |= raw::IORING_ASYNC_CANCEL_OP;
186            }
187            CancelTarget::Any => cancel_flags |= raw::IORING_ASYNC_CANCEL_ANY,
188        }
189        self.submit_op(None, |sqe| {
190            sqe.opcode = raw::IORING_OP_ASYNC_CANCEL;
191            sqe.fd = fd;
192            sqe.addr_or_splice_off_in = addr;
193            sqe.len = len;
194            sqe.op_flags = cancel_flags;
195        })
196    }
197
198    // ── Surveillance d'événements ─────────────────────────────────────────
199
200    /// Soumet un `IORING_OP_POLL_ADD` (6) : surveille `fd` pour `events` (mono-
201    /// coup ; multishot/level-triggered au Temps 3d). Complétion via
202    /// [`Completion::into_poll_result`](super::Completion::into_poll_result).
203    ///
204    /// # Errors
205    ///
206    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
207    pub fn submit_poll_add(
208        &mut self,
209        fd: BorrowedFd<'_>,
210        events: PollEvents,
211    ) -> Result<SubmissionToken, Errno> {
212        let fd = fd.as_raw_fd();
213        let events = events.bits();
214        self.submit_op(None, |sqe| {
215            sqe.opcode = raw::IORING_OP_POLL_ADD;
216            sqe.fd = fd;
217            sqe.op_flags = events;
218        })
219    }
220
221    /// Soumet un `IORING_OP_POLL_REMOVE` (7) : retire un poll en vol désigné par
222    /// son jeton. Complétion via `completed`.
223    ///
224    /// # Errors
225    ///
226    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
227    pub fn submit_poll_remove(
228        &mut self,
229        target: SubmissionToken,
230    ) -> Result<SubmissionToken, Errno> {
231        let addr = target.to_user_data();
232        self.submit_op(None, |sqe| {
233            sqe.opcode = raw::IORING_OP_POLL_REMOVE;
234            sqe.fd = -1;
235            sqe.addr_or_splice_off_in = addr;
236        })
237    }
238
239    /// Soumet un `IORING_OP_POLL_REMOVE` + `IORING_POLL_UPDATE_EVENTS` (7) :
240    /// modifie les événements surveillés d'un poll en vol. Complétion via
241    /// `completed`.
242    ///
243    /// # Errors
244    ///
245    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
246    pub fn submit_poll_update(
247        &mut self,
248        target: SubmissionToken,
249        events: PollEvents,
250    ) -> Result<SubmissionToken, Errno> {
251        let ud = target.to_user_data();
252        let events = events.bits();
253        self.submit_op(None, |sqe| {
254            sqe.opcode = raw::IORING_OP_POLL_REMOVE;
255            sqe.fd = -1;
256            sqe.addr_or_splice_off_in = ud;
257            sqe.off_or_addr2 = ud;
258            sqe.len = raw::IORING_POLL_UPDATE_EVENTS;
259            sqe.op_flags = events;
260        })
261    }
262
263    /// Soumet un `IORING_OP_EPOLL_CTL` (29) : add/mod/del sur un set epoll, en
264    /// asynchrone. L'`epoll_event` est gardé en vie dans le slot (lu par le
265    /// kernel). Complétion via `completed`.
266    ///
267    /// # Errors
268    ///
269    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
270    pub fn submit_epoll_ctl(
271        &mut self,
272        epfd: BorrowedFd<'_>,
273        op: EpollOp,
274        fd: BorrowedFd<'_>,
275        event: EpollEvent,
276    ) -> Result<SubmissionToken, Errno> {
277        let epfd = epfd.as_raw_fd();
278        let target_fd = u64::from(fd.as_raw_fd().cast_unsigned());
279        let op = op.as_raw().cast_unsigned();
280        let event = Box::new(event);
281        let event_ptr = core::ptr::from_ref(&*event) as u64;
282        self.submit_op(Some(OwnedOp::Epoll(event)), |sqe| {
283            sqe.opcode = raw::IORING_OP_EPOLL_CTL;
284            sqe.fd = epfd;
285            sqe.addr_or_splice_off_in = event_ptr;
286            sqe.len = op;
287            sqe.off_or_addr2 = target_fd;
288        })
289    }
290
291    // ── Attente de processus ──────────────────────────────────────────────
292
293    /// Soumet un `IORING_OP_WAITID` (50) : récolte la fin d'un enfant **sans
294    /// bloquer** le reactor. Le kernel remplit `info` (`siginfo_t`, déplacé dans
295    /// le slot). Complétion via
296    /// [`Completion::into_waitid_result`](super::Completion::into_waitid_result).
297    ///
298    /// # Errors
299    ///
300    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
301    pub fn submit_waitid(
302        &mut self,
303        target: WaitTarget<'_>,
304        options: WaitOptions,
305        info: Box<SignalInfo>,
306    ) -> Result<SubmissionToken, Errno> {
307        let (idtype, id) = waitid_target(target);
308        let info_ptr = info.as_bytes().as_ptr() as u64;
309        let options = options.bits().cast_unsigned();
310        self.submit_op(Some(OwnedOp::Waitid(info)), |sqe| {
311            sqe.opcode = raw::IORING_OP_WAITID;
312            sqe.fd = id;
313            sqe.len = idtype.cast_unsigned();
314            sqe.off_or_addr2 = info_ptr;
315            sqe.splice_fd_in_or_file_index = options;
316        })
317    }
318
319    // ── Inter-ring (notification) ─────────────────────────────────────────
320
321    /// Soumet un `IORING_OP_MSG_RING` (40, `MSG_DATA`) : poste un CQE (`res` +
322    /// `user_data` arbitraires) dans le **ring cible** — notification inter-
323    /// reactor très légère. Complétion (côté émetteur) via `completed` ; côté
324    /// cible, un CQE apparaît dans sa CQ.
325    ///
326    /// # Errors
327    ///
328    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
329    pub fn submit_message_ring_data(
330        &mut self,
331        target: &IoUring,
332        res: i32,
333        user_data: u64,
334        flags: MessageRingFlags,
335    ) -> Result<SubmissionToken, Errno> {
336        let target_fd = target.fd_raw();
337        let len = res.cast_unsigned();
338        let op_flags = flags.bits();
339        self.submit_op(None, |sqe| {
340            sqe.opcode = raw::IORING_OP_MSG_RING;
341            sqe.fd = target_fd;
342            sqe.addr_or_splice_off_in = raw::IORING_MSG_DATA;
343            sqe.len = len;
344            sqe.off_or_addr2 = user_data;
345            sqe.op_flags = op_flags;
346        })
347    }
348
349    // ── Futex (mémoire partagée) ──────────────────────────────────────────
350    //
351    // Le mot futex vit dans une `MmapRegion` partageable : le kernel le lit
352    // **de façon asynchrone** jusqu'à la complétion. Le slot S1 retient la
353    // garde de vivacité de la région (`OwnedOp::MemLiveness` /
354    // `OwnedOp::FutexWaitv`) ⇒ `munmap` impossible tant que l'op est en vol :
355    // ni use-after-unmap, ni fuite (cf. `family-mem-mmap-region.md`).
356
357    /// Encode les drapeaux `FUTEX2_*` posés dans `sqe->fd` (`WAIT`/`WAKE`). La
358    /// taille `FUTEX2_SIZE_U32` est **imposée** (le mot d'une `MmapRegion` est un
359    /// `AtomicU32`) ; seul `PRIVATE` provient de `flags`.
360    fn futex2_fd(flags: FutexFlags) -> Result<i32, Errno> {
361        i32::try_from(raw::FUTEX2_SIZE_U32 | flags.bits()).map_err(|_| Errno::EINVAL)
362    }
363
364    /// Soumet un `IORING_OP_FUTEX_WAIT` (51) : attend que le mot futex à `offset`
365    /// dans `region` soit réveillé (ou `-EAGAIN` si sa valeur ≠ `expected` à
366    /// l'armement). `mask` est le bitset futex (bits surveillés). La région est
367    /// gardée vivante par le slot S1. Complétion via
368    /// [`Completion::completed`](super::Completion::completed).
369    ///
370    /// # Errors
371    ///
372    /// [`Errno::EINVAL`] si `offset` est hors bornes, non aligné sur 4, ou si la
373    /// région n'est pas lisible (cf. [`MmapRegion::futex_word`]) ;
374    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
375    pub fn submit_futex_wait(
376        &mut self,
377        region: &MmapRegion,
378        offset: usize,
379        expected: u64,
380        mask: u64,
381        flags: FutexFlags,
382    ) -> Result<SubmissionToken, Errno> {
383        let uaddr = core::ptr::from_ref(region.futex_word(offset)?) as u64;
384        let futex2 = Self::futex2_fd(flags)?;
385        let guard = region.liveness_handle();
386        self.submit_op(Some(OwnedOp::MemLiveness(guard)), |sqe| {
387            sqe.opcode = raw::IORING_OP_FUTEX_WAIT;
388            sqe.fd = futex2;
389            sqe.addr_or_splice_off_in = uaddr;
390            sqe.off_or_addr2 = expected;
391            sqe.addr3 = mask;
392        })
393    }
394
395    /// Soumet un `IORING_OP_FUTEX_WAKE` (52) : réveille jusqu'à `nr` waiters sur
396    /// le mot futex à `offset` dans `region` (filtrés par `mask`). Complétion via
397    /// [`Completion::into_result`](super::Completion::into_result) (nombre de
398    /// waiters réveillés).
399    ///
400    /// # Errors
401    ///
402    /// [`Errno::EINVAL`] si `offset` est invalide (cf. [`MmapRegion::futex_word`]) ;
403    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
404    pub fn submit_futex_wake(
405        &mut self,
406        region: &MmapRegion,
407        offset: usize,
408        nr: u64,
409        mask: u64,
410        flags: FutexFlags,
411    ) -> Result<SubmissionToken, Errno> {
412        let uaddr = core::ptr::from_ref(region.futex_word(offset)?) as u64;
413        let futex2 = Self::futex2_fd(flags)?;
414        let guard = region.liveness_handle();
415        self.submit_op(Some(OwnedOp::MemLiveness(guard)), |sqe| {
416            sqe.opcode = raw::IORING_OP_FUTEX_WAKE;
417            sqe.fd = futex2;
418            sqe.addr_or_splice_off_in = uaddr;
419            sqe.off_or_addr2 = nr;
420            sqe.addr3 = mask;
421        })
422    }
423
424    /// Soumet un `IORING_OP_FUTEX_WAITV` (53) : attend sur **plusieurs** mots
425    /// futex à la fois (réveil dès qu'**un** se déclenche). Chaque [`FutexWaiter`]
426    /// désigne une région + offset, sa valeur attendue. Le tableau de
427    /// `futex_waitv` et toutes les gardes de vivacité sont garés dans le slot S1.
428    /// Complétion via [`Completion::completed`](super::Completion::completed).
429    ///
430    /// **Note ABI.** `struct futex_waitv` du kernel ne porte **pas** de masque
431    /// par-attendu (champs `{ val, uaddr, flags, __reserved }`) : [`FutexWaiter`]
432    /// n'expose donc **pas** de masque (un champ silencieusement ignoré serait un
433    /// footgun, cf. ADR-032). Pour une attente **masquée**, utiliser
434    /// [`IoUring::submit_futex_wait`] (mono-attente, qui garde `mask`).
435    ///
436    /// # Errors
437    ///
438    /// [`Errno::EINVAL`] si `waiters` est vide, dépasse `FUTEX_WAITV_MAX` (128),
439    /// ou si un offset est invalide ; [`Errno::EBUSY`] si la SQ ou le slab sont
440    /// pleins.
441    pub fn submit_futex_waitv(
442        &mut self,
443        waiters: Vec<FutexWaiter>,
444        flags: FutexFlags,
445    ) -> Result<SubmissionToken, Errno> {
446        // Validation amont : au moins un attendu, pas plus que le maximum kernel.
447        if waiters.is_empty() || waiters.len() > raw::FUTEX_WAITV_MAX {
448            return Err(Errno::EINVAL);
449        }
450        let futex2 = raw::FUTEX2_SIZE_U32 | flags.bits();
451        let mut raws: Vec<FutexWaitvRaw> = Vec::with_capacity(waiters.len());
452        let mut guards = Vec::with_capacity(waiters.len());
453        for waiter in &waiters {
454            // Valide chaque offset (bornes + alignement) AVANT toute soumission.
455            let uaddr = core::ptr::from_ref(waiter.region.futex_word(waiter.offset)?) as u64;
456            raws.push(FutexWaitvRaw {
457                val: waiter.expected,
458                uaddr,
459                flags: futex2,
460                reserved: 0,
461            });
462            guards.push(waiter.region.liveness_handle());
463        }
464        let raws = raws.into_boxed_slice();
465        let arr_ptr = raws.as_ptr() as u64;
466        // `len() ≤ FUTEX_WAITV_MAX = 128`, donc le `u32` ne déborde jamais.
467        let nr = u32::try_from(raws.len()).map_err(|_| Errno::EINVAL)?;
468        self.submit_op(
469            Some(OwnedOp::FutexWaitv {
470                waiters: raws,
471                guards,
472            }),
473            |sqe| {
474                sqe.opcode = raw::IORING_OP_FUTEX_WAITV;
475                sqe.fd = 0;
476                sqe.addr_or_splice_off_in = arr_ptr;
477                sqe.len = nr;
478            },
479        )
480    }
481}
482
483/// Un attendu de [`IoUring::submit_futex_waitv`] : la région partageable + le
484/// décalage du mot futex, et la valeur attendue.
485///
486/// La `region` est **possédée** (un `clone` de la `MmapRegion`) : elle *est* la
487/// garde de vivacité naturelle — tant que le `FutexWaiter` (et donc la garde
488/// garée dans le slot) vit, la région reste mappée.
489#[derive(Debug, Clone)]
490pub struct FutexWaiter {
491    /// Région partageable contenant le mot futex (gardée vivante par le slot).
492    pub region: MmapRegion,
493    /// Décalage du mot futex dans la région (borné, aligné sur 4).
494    pub offset: usize,
495    /// Valeur attendue au mot à l'armement.
496    pub expected: u64,
497}
498
499// ───────────────────────────────────────────────────────────────────────────
500// Tests
501// ───────────────────────────────────────────────────────────────────────────
502//
503// Intégration (kernel 6.12 réel) : `#[cfg_attr(miri, ignore)]`. Tests purs
504// (accesseurs, ownership des payloads possédés) sous Miri.
505
506#[cfg(test)]
507mod tests;