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;