air_sys_syscall/io_uring/mod.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//! Module `air-sys-syscall::io_uring` — façade typée d'io_uring (cible 6.12).
6//!
7//! Niveau d'abstraction 2 (ADR-022, Décision 1) : soumission/complétion typée.
8//! Le niveau 1 (anneaux bruts) vivra dans [`raw`]. Les buffers suivent le
9//! modèle de transfert d'ownership (ADR-022, Décision 3), garés dans un slab
10//! pré-alloué (S1) ; le téardown est sûr par `Drop` quiescent +
11//! [`IoUring::shutdown`] (S2).
12//!
13//! **Périmètre — Temps 1 (cœur).** Cycle de vie du ring (setup flags, mmaps SQ/
14//! CQ/SQE, négociation des features), slab d'opérations en vol (S1), protocole
15//! d'anneau acquire/release (§3.2 de la spec), soumission/complétion,
16//! [`Completion`], téardown sûr (S2), introspection (probe + capabilities). Les
17//! opérations métier `submit_*` (Temps 2a–2d) et les sous-modules
18//! ([`shared`], [`multishot`], [`sandbox`], registration/provided/linked/cmd)
19//! relèvent de PRs ultérieures.
20//!
21//! Référence normative : `docs/specs/layer-0/io-uring-1-core.md` (contrat) et
22//! `docs/specs/layer-0/io-uring-0-inventaire.md` §7 (squelette).
23//!
24//! Les trois syscalls sous-jacents — `io_uring_setup` (425), `io_uring_enter`
25//! (426), `io_uring_register` (427), numéros identiques x86_64/ARM64 — sont
26//! appelés directement via `core::arch::asm!` (aucune dépendance externe).
27
28use crate::mem::{Mapping, madvise, mmap_file};
29use air_sys_types::Errno;
30use air_sys_types::fd::{AsFd, AsRawFd, BorrowedFd, FromRawFd, OwnedFd};
31use air_sys_types::mem::{MadviseAdvice, MapFlags, ProtectionFlags};
32use alloc::boxed::Box;
33use alloc::vec::Vec;
34use core::marker::PhantomData;
35use core::num::NonZeroU32;
36use core::time::Duration;
37
38pub mod multishot;
39pub mod raw;
40pub mod sandbox;
41pub mod shared;
42
43mod async_ops;
44mod cmd;
45mod fs_ops;
46mod linked;
47mod net_ops;
48mod owned;
49mod provided;
50mod registration;
51mod ring;
52mod slab;
53mod syscall;
54
55use owned::OwnedOp;
56use ring::{CompletionRing, SubmissionRing};
57use slab::InflightSlab;
58
59// Types de données purs (bitflags, opcodes) : définis dans `air-sys-types`,
60// re-exportés ici pour offrir une façade `air_sys_syscall::io_uring::*` unique.
61pub use air_sys_types::io_uring::{CompletionFlags, IoUringOpcode, SetupFlags, UringCmdFlags};
62
63// Type d'attente `futex_waitv` (porte une `MmapRegion`) : défini avec ses
64// façades dans `async_ops`, re-exporté sur la façade `io_uring`.
65pub use async_ops::FutexWaiter;
66
67// Trait des commandes passthrough `URING_CMD` (Temps 2d), défini avec ses
68// façades dans `cmd`, re-exporté sur la façade `io_uring`.
69pub use cmd::UringCommand;
70
71// Surface du Temps 3a (registration : ressources fixes), définie dans
72// `registration`, re-exportée sur la façade `io_uring`.
73pub use registration::{
74 ClockSource, FixedFdTable, FixedSlot, FixedSlotTarget, NapiConfig, Personality,
75 RegisteredBufferSlice, RegisteredBuffers, WorkQueueWorkerLimits,
76};
77
78// Surface du Temps 3b (buffers fournis ring-mapped), définie dans `provided`,
79// re-exportée sur la façade `io_uring`.
80pub use provided::{
81 ProvidedBuffer, ProvidedBufferRing, ProvidedBufferRingOptions, ProvidedBufferRingStatus,
82};
83
84// Surface du Temps 3c (opérations liées), définie dans `linked`, re-exportée sur
85// la façade `io_uring`.
86pub use linked::{ChainTokens, LinkedChainBuilder};
87
88// Surface du Temps 3e (usage multi-thread), définie dans `shared`, re-exportée
89// sur la façade `io_uring`.
90pub use shared::{LockedIoUring, RingHandle, RingPool, SqpollIoUring};
91
92// Surface du Temps 3f (confinement / sandbox), définie dans `sandbox`,
93// re-exportée sur la façade `io_uring`.
94pub use sandbox::{RegisterOp, RestrictionSet, SqeFlagSet};
95
96// Surface du Temps 4 (accès brut niveau 1), définie dans `raw`, re-exportée sur
97// la façade `io_uring`. `raw` reste par ailleurs le module des miroirs ABI.
98pub use raw::{RAW_USER_DATA_TAG, RawCompletionQueueEntry, RawOpcode, RawSubmissionQueueEntry};
99
100// ---------------------------------------------------------------------------
101// Construction & cycle de vie
102// ---------------------------------------------------------------------------
103
104/// Anneau io_uring : FD + mmaps (SQ/CQ/SQE) + slab d'opérations en vol (S1) +
105/// capabilities négociées.
106///
107/// `Send` mais **pas** `Sync` (ADR-022, Décision 6) : un reactor par thread.
108/// Pour le multi-thread, voir [`shared`]. L'invariant de sûreté (S1/S2) : tant
109/// qu'une opération est en vol, son buffer vit dans le slot du slab et le ring
110/// ne peut être détruit sans quiescence — aucune écriture kernel ne tombe sur
111/// de la mémoire libérée.
112pub struct IoUring {
113 /// FD du ring. `Option` pour la **discipline de téardown** : libéré (close)
114 /// après quiescence, ou *fuité* (`forget`) si la quiescence échoue.
115 fd: Option<OwnedFd>,
116 /// Mmap de l'anneau SQ (et de la CQ si `SINGLE_MMAP`). `Option` : voir `fd`.
117 sq_ring_map: Option<Mapping>,
118 /// Mmap du tableau de SQE. `Option` : voir `fd`.
119 sqes_map: Option<Mapping>,
120 /// Mmap de l'anneau CQ, ou `None` si partagée avec `sq_ring_map`
121 /// (`IORING_FEAT_SINGLE_MMAP`, toujours présent en 6.12).
122 cq_ring_map: Option<Mapping>,
123 /// Pointeurs/cache de l'anneau de soumission (dans `sq_ring_map`).
124 sq: SubmissionRing,
125 /// Pointeurs/cache de l'anneau de complétion.
126 cq: CompletionRing,
127 /// Slab d'opérations en vol (S1).
128 slab: InflightSlab,
129 /// Features négociées au setup (axe G).
130 capabilities: IoUringCapabilities,
131 /// Opcodes supportés par le kernel courant, indexés par numéro `IORING_OP_*`
132 /// (mis en cache à la construction via `IORING_REGISTER_PROBE`, §9).
133 supported_ops: [bool; 256],
134 /// Options appliquées à la prochaine soumission ([`IoUring::with`]).
135 pending_options: SubmitOptions,
136 /// Personality appliquée à la prochaine soumission ([`IoUring::with_personality`]) :
137 /// `0` = credentials du process (aucune personality). Posée dans
138 /// `sqe.personality` puis remise à `0` (Temps 3a, registration §6).
139 pending_personality: u16,
140 /// Index du ring fd enregistré ([`IoUring::register_ring_fd`], Temps 3a §4),
141 /// ou `None`. Quand présent, les `io_uring_enter` suivants l'utilisent avec
142 /// `IORING_ENTER_REGISTERED_RING` (pas de résolution de FD).
143 enter_ring_index: Option<u32>,
144 /// `true` une fois le téardown effectué (par `shutdown` ou un `Drop`
145 /// précédent) : empêche toute double libération (reminder S2).
146 disposed: bool,
147 /// Indices (`tail`) des SQE préparés via [`IoUring::raw_get_submission_queue_entry`]
148 /// (Temps 4) depuis la dernière soumission, pour la **vérification debug** du
149 /// tag `user_data` (§5). Vide tant que l'accès brut n'est pas utilisé (aucune
150 /// allocation sur le happy path niveau 2) ; purgé à chaque `submit`.
151 raw_pending: Vec<u32>,
152 /// Marqueur `Send + !Sync` : `Cell<()>` est `Send` mais pas `Sync`, ce qui
153 /// impose le « un reactor par thread » de la Décision 6 sans verrou.
154 _not_sync: PhantomData<core::cell::Cell<()>>,
155}
156
157// SAFETY: `IoUring` est **`Send`** (ADR-022 Décision 6) : il peut être *déplacé*
158// d'un thread à l'autre. Ses pointeurs bruts (`SubmissionRing`/`CompletionRing`)
159// désignent des zones **mmappées** valides pour toute la durée de vie du ring
160// (process-global, indépendantes du thread) ; le FD est process-global ; le slab
161// S1 et les payloads possédés (`OwnedFd`/`Vec`/`Box`) sont tous `Send`. Déplacer
162// le ring **transfère l'ownership exclusif** — `Send` concerne le transfert, pas
163// le partage. Le partage concurrent (`Sync`) reste **interdit** : `_not_sync`
164// (`PhantomData<Cell<()>>`) garde `IoUring: !Sync`, car le protocole d'ordering
165// du Temps 1 ne synchronise pas userspace↔userspace (cf. `shared` §1). Un seul
166// thread accède donc au ring à un instant donné (Send + !Sync = thread-per-core).
167unsafe impl Send for IoUring {}
168
169/// Construit un [`IoUring`] : applique setup flags et restrictions (S3) avant
170/// l'activation, puis `io_uring_setup(2)`.
171pub struct IoUringBuilder {
172 entries: NonZeroU32,
173 cq_entries: Option<NonZeroU32>,
174 max_inflight: Option<NonZeroU32>,
175 flags: SetupFlags,
176 sqpoll_idle: Option<Duration>,
177 sqpoll_cpu: Option<u32>,
178 attach_wq_fd: Option<i32>,
179 restrictions: Vec<Restriction>,
180}
181
182impl IoUringBuilder {
183 /// Démarre un builder pour `entries` SQE (arrondi par le kernel à la
184 /// puissance de 2 supérieure).
185 #[must_use]
186 pub fn new(entries: NonZeroU32) -> Self {
187 Self {
188 entries,
189 cq_entries: None,
190 max_inflight: None,
191 flags: SetupFlags::empty(),
192 sqpoll_idle: None,
193 sqpoll_cpu: None,
194 attach_wq_fd: None,
195 restrictions: Vec::new(),
196 }
197 }
198
199 /// Fixe explicitement la profondeur de la CQ (`IORING_SETUP_CQSIZE`).
200 #[must_use]
201 pub fn with_completion_queue_entries(mut self, entries: NonZeroU32) -> Self {
202 self.cq_entries = Some(entries);
203 self
204 }
205
206 /// Capacité du slab d'opérations en vol (S1). Défaut : `cq_entries`.
207 #[must_use]
208 pub fn max_inflight(mut self, n: NonZeroU32) -> Self {
209 self.max_inflight = Some(n);
210 self
211 }
212
213 /// Active des flags de setup (cf. [`SetupFlags`]).
214 #[must_use]
215 pub fn with_flags(mut self, flags: SetupFlags) -> Self {
216 self.flags |= flags;
217 self
218 }
219
220 /// Durée d'inactivité avant que le thread `SQPOLL` ne s'endorme
221 /// (`SETUP_SQPOLL`, Temps 3e).
222 #[must_use]
223 pub fn with_sqpoll_idle(mut self, idle: Duration) -> Self {
224 self.sqpoll_idle = Some(idle);
225 self
226 }
227
228 /// Épingle le thread `SQPOLL` sur un CPU (`SETUP_SQ_AFF`).
229 #[must_use]
230 pub fn with_sqpoll_cpu(mut self, cpu: u32) -> Self {
231 self.sqpoll_cpu = Some(cpu);
232 self
233 }
234
235 /// Partage le pool io-wq d'un ring existant (`SETUP_ATTACH_WQ`).
236 #[must_use]
237 pub fn attach_work_queue(mut self, other: &IoUring) -> Self {
238 self.attach_wq_fd = Some(other.fd_raw());
239 self
240 }
241
242 /// Crée le ring désactivé (`R_DISABLED`) et applique des restrictions (S3,
243 /// Temps 3f). Doit être suivi de [`IoUring::enable`].
244 #[must_use]
245 pub fn restrict(mut self, restrictions: &[Restriction]) -> Self {
246 self.restrictions = restrictions.to_vec();
247 self
248 }
249
250 /// Finalise : traduit la config en `io_uring_params`, appelle
251 /// `io_uring_setup(2)`, mmappe les anneaux (§3.1), alloue le slab (§4) et
252 /// lit les features. Si `restrict` a été utilisé, le ring est créé
253 /// désactivé et doit être activé via [`IoUring::enable`].
254 ///
255 /// # Errors
256 ///
257 /// - [`Errno::EINVAL`] : paramètres ou combinaison de flags invalides
258 /// (remontés tels quels, sans masquage).
259 /// - [`Errno::ENOMEM`] : mémoire insuffisante pour les anneaux/le slab.
260 /// - [`Errno::EPERM`] : `SQPOLL`/affinité sans privilège selon la config.
261 /// - [`Errno::EFAULT`] : pointeur de paramètres invalide.
262 /// - [`Errno::ENOSYS`] : io_uring indisponible (kernel ancien, sandbox,
263 /// container durci) — **pas** de fallback caché (cf. ADR-022 D10).
264 pub fn build(self) -> Result<IoUring, Errno> {
265 let mut params = raw::IoUringParams {
266 flags: self.flags.bits(),
267 ..raw::IoUringParams::default()
268 };
269 if let Some(cq) = self.cq_entries {
270 params.flags |= raw::IORING_SETUP_CQSIZE;
271 params.cq_entries = cq.get();
272 }
273 if let Some(idle) = self.sqpoll_idle {
274 params.sq_thread_idle = u32::try_from(idle.as_millis()).unwrap_or(u32::MAX);
275 }
276 if let Some(cpu) = self.sqpoll_cpu {
277 params.sq_thread_cpu = cpu;
278 }
279 if let Some(wq_fd) = self.attach_wq_fd {
280 params.flags |= raw::IORING_SETUP_ATTACH_WQ;
281 params.wq_fd = u32::from_ne_bytes(wq_fd.to_ne_bytes());
282 }
283 if !self.restrictions.is_empty() {
284 params.flags |= SetupFlags::R_DISABLED.bits();
285 }
286
287 let params_ptr = core::ptr::from_mut(&mut params) as u64;
288 // SAFETY: `params_ptr` pointe une `IoUringParams` valide en lecture/
289 // écriture pour la durée de l'appel ; `entries` est un scalaire.
290 let ret = unsafe { syscall::setup(self.entries.get(), params_ptr) };
291 if ret < 0 {
292 return Err(raw::errno_from_negative_syscall_ret(ret));
293 }
294 let raw_fd = i32::try_from(ret).map_err(|_| Errno::EINVAL)?;
295 // SAFETY: `io_uring_setup` a retourné un FD frais possédé.
296 let fd = unsafe { OwnedFd::from_raw_fd(raw_fd) };
297
298 if !self.restrictions.is_empty() {
299 // Ring créé `R_DISABLED` ci-dessus : applique la liste blanche S3
300 // (`REGISTER_RESTRICTIONS`, Temps 3f, sous-module `sandbox`) **tant
301 // qu'il est désactivé** (le kernel refuse sinon). En cas d'échec, `fd`
302 // (OwnedFd local) se ferme proprement au retour — aucune fuite. Le
303 // ring reste désactivé : l'appelant l'active via `IoUring::enable`,
304 // après quoi les restrictions sont **immuables** (imposé kernel).
305 apply_restrictions(&fd, &self.restrictions)?;
306 }
307
308 let mappings = map_rings(&fd, ¶ms, self.flags)?;
309 let slab_capacity = self.max_inflight.map_or(params.cq_entries, NonZeroU32::get);
310 let slab_capacity = NonZeroU32::new(slab_capacity).ok_or(Errno::EINVAL)?;
311 let supported_ops = probe_supported_ops(fd.as_raw_fd());
312
313 Ok(IoUring {
314 fd: Some(fd),
315 sq_ring_map: Some(mappings.sq_ring),
316 sqes_map: Some(mappings.sqes),
317 cq_ring_map: mappings.cq_ring,
318 sq: mappings.sq,
319 cq: mappings.cq,
320 slab: InflightSlab::with_capacity(slab_capacity),
321 capabilities: IoUringCapabilities {
322 features: params.features,
323 },
324 supported_ops,
325 pending_options: SubmitOptions::default(),
326 pending_personality: 0,
327 enter_ring_index: None,
328 disposed: false,
329 raw_pending: Vec::new(),
330 _not_sync: PhantomData,
331 })
332 }
333}
334
335/// Mappings et holders d'anneaux résultant de la construction.
336struct RingMappings {
337 sq_ring: Mapping,
338 sqes: Mapping,
339 cq_ring: Option<Mapping>,
340 sq: SubmissionRing,
341 cq: CompletionRing,
342}
343
344/// mmappe les trois zones (§3.1) et construit les holders de pointeurs.
345/// Tailles de mmap des trois zones, dérivées des `params` retournés par le
346/// kernel. **Décode PUR** (aucun syscall) — frontière de décode de données
347/// externes (Principe 3), fuzzée via [`fuzz_api`]. Arithmétique **checked** :
348/// toute donnée kernel incohérente (offsets/entrées énormes) ⇒ `EINVAL`, jamais
349/// d'overflow ni de panic.
350struct RingSizes {
351 single_mmap: bool,
352 sq_map: usize,
353 cq_map: usize,
354 sqes: usize,
355 /// Taille d'un SQE en octets (64, ou 128 si `SETUP_SQE128`).
356 sqe_size: usize,
357}
358
359fn ring_sizes(params: &raw::IoUringParams, flags: SetupFlags) -> Result<RingSizes, Errno> {
360 let single_mmap = params.features & raw::IORING_FEAT_SINGLE_MMAP != 0;
361 let sqe_size: usize = if flags.contains(SetupFlags::SQE128) {
362 128
363 } else {
364 64
365 };
366 let cqe_size: usize = if flags.contains(SetupFlags::CQE32) {
367 32
368 } else {
369 16
370 };
371
372 let sq_entries = usize_of(params.sq_entries);
373 let cq_entries = usize_of(params.cq_entries);
374 let sq_ring_sz = usize_of(params.sq_off.array)
375 .checked_add(sq_entries.checked_mul(4).ok_or(Errno::EINVAL)?)
376 .ok_or(Errno::EINVAL)?;
377 let cq_ring_sz = usize_of(params.cq_off.cqes)
378 .checked_add(cq_entries.checked_mul(cqe_size).ok_or(Errno::EINVAL)?)
379 .ok_or(Errno::EINVAL)?;
380 let sqes = sq_entries.checked_mul(sqe_size).ok_or(Errno::EINVAL)?;
381
382 let (sq_map, cq_map) = if single_mmap {
383 let m = sq_ring_sz.max(cq_ring_sz);
384 (m, m)
385 } else {
386 (sq_ring_sz, cq_ring_sz)
387 };
388 Ok(RingSizes {
389 single_mmap,
390 sq_map,
391 cq_map,
392 sqes,
393 sqe_size,
394 })
395}
396
397/// Retire un mapping d'anneau de l'héritage au `fork` (`MADV_DONTFORK`).
398///
399/// **Mesure de sécurité, pas d'optimisation** (ADR-132). Les zones d'un ring
400/// (SQ, CQ, tableau de SQE, et l'anneau de descripteurs des buffers fournis)
401/// sont mappées `MAP_SHARED` sur le ring fd : sans ce conseil, un enfant
402/// forké garde une fenêtre **en écriture** sur les *mêmes pages physiques* que
403/// le parent — fermer le descripteur (`close_range`) ne démappe rien, et les
404/// étages privsep sortent par `exit_group` sans destructeur. Un enfant pré-auth
405/// compromis pourrait alors réécrire un SQE que le parent (root) a préparé, ou
406/// avancer le `tail` de la SQ, et faire exécuter au parent l'opération de son
407/// choix — contournement complet de la séparation de privilèges.
408///
409/// `MADV_DONTFORK` pose `VM_DONTCOPY` sur le VMA : le processus courant garde
410/// son mapping intact (le réacteur continue de fonctionner normalement), seuls
411/// les **descendants** créés par `fork`/`clone` **sans** `CLONE_VM` ne le
412/// reçoivent pas. Aucun usage légitime n'hérite d'un anneau à travers un
413/// `fork` — ADR-122 §5 pose exactement l'inverse comme doctrine.
414///
415/// **Fail-closed** : toute erreur remonte et fait échouer la construction du
416/// ring. Un anneau sans cette protection ne doit pas exister.
417///
418/// `base`/`length` désignent un mapping vivant obtenu par `mmap` (base alignée
419/// sur la page ; `madvise(2)` arrondit `length` à la page supérieure, donc le
420/// VMA entier est couvert même si la taille utile ne l'est pas).
421fn deny_fork_inheritance(base: *const u8, length: usize) -> Result<(), Errno> {
422 let address =
423 core::ptr::NonNull::new(base.cast_mut()).expect("mmap rend une adresse non nulle");
424 madvise(address, length, MadviseAdvice::DontFork)
425}
426
427fn map_rings(
428 fd: &OwnedFd,
429 params: &raw::IoUringParams,
430 flags: SetupFlags,
431) -> Result<RingMappings, Errno> {
432 let sizes = ring_sizes(params, flags)?;
433 let single_mmap = sizes.single_mmap;
434
435 let prot = ProtectionFlags::READ | ProtectionFlags::WRITE;
436 let map = MapFlags::SHARED | MapFlags::POPULATE;
437
438 let mut sq_ring = mmap_file(fd.as_fd(), sizes.sq_map, raw::IORING_OFF_SQ_RING, prot, map)?;
439 deny_fork_inheritance(sq_ring.as_ptr(), sq_ring.len())?;
440 let mut cq_ring = if single_mmap {
441 None
442 } else {
443 let cq = mmap_file(fd.as_fd(), sizes.cq_map, raw::IORING_OFF_CQ_RING, prot, map)?;
444 deny_fork_inheritance(cq.as_ptr(), cq.len())?;
445 Some(cq)
446 };
447 let mut sqes = mmap_file(fd.as_fd(), sizes.sqes, raw::IORING_OFF_SQES, prot, map)?;
448 deny_fork_inheritance(sqes.as_ptr(), sqes.len())?;
449
450 let sq_base = sq_ring.as_mut_ptr();
451 let cq_base = cq_ring.as_mut().map_or(sq_base, Mapping::as_mut_ptr);
452 let sqes_base = sqes.as_mut_ptr();
453 let sq_mask = params.sq_entries.wrapping_sub(1);
454 let cq_mask = params.cq_entries.wrapping_sub(1);
455
456 // SAFETY: les trois bases pointent des mmaps fraîches dimensionnées selon
457 // les offsets kernel de `params` ; elles vivent dans l'`IoUring` (RAII) au
458 // moins aussi longtemps que les holders qui en dérivent les pointeurs.
459 let sq = unsafe {
460 SubmissionRing::new(
461 sq_base,
462 sqes_base,
463 sizes.sqe_size,
464 ¶ms.sq_off,
465 sq_mask,
466 params.sq_entries,
467 )
468 };
469 // SAFETY: idem pour l'anneau CQ (zone propre ou partagée si SINGLE_MMAP).
470 let cq = unsafe { CompletionRing::new(cq_base, ¶ms.cq_off, cq_mask) };
471
472 Ok(RingMappings {
473 sq_ring,
474 sqes,
475 cq_ring,
476 sq,
477 cq,
478 })
479}
480
481/// `u32` → `usize` (toujours valide sur cible LP64).
482fn usize_of(v: u32) -> usize {
483 usize::try_from(v).expect("usize ≥ u32 sur cible LP64")
484}
485
486/// Encode les [`Restriction`] en `struct io_uring_restriction` et les applique
487/// via `IORING_REGISTER_RESTRICTIONS` (S3, Temps 3f, sous-module [`sandbox`]).
488/// À n'appeler **que** sur un ring `R_DISABLED` (avant [`IoUring::enable`]) — le
489/// kernel refuse autrement. **Default-deny** : dès qu'un opcode est mis en liste
490/// blanche, le kernel refuse tout le reste (`-EACCES`).
491fn apply_restrictions(fd: &OwnedFd, restrictions: &[Restriction]) -> Result<(), Errno> {
492 let encoded: Vec<raw::IoUringRestriction> = restrictions
493 .iter()
494 .copied()
495 .map(encode_restriction)
496 .collect();
497 let nr_args = u32::try_from(encoded.len()).map_err(|_| Errno::EINVAL)?;
498 let arg = encoded.as_ptr() as u64;
499 // SAFETY: `arg` pointe un tableau de `nr_args` `io_uring_restriction` valides,
500 // vivants pour la durée de l'appel ; `fd` est un ring fd `R_DISABLED` possédé.
501 // `REGISTER_RESTRICTIONS` lit l'argument sans le réécrire.
502 let ret = unsafe {
503 syscall::register(
504 fd.as_raw_fd(),
505 raw::IORING_REGISTER_RESTRICTIONS,
506 arg,
507 nr_args,
508 )
509 };
510 if ret < 0 {
511 return Err(raw::errno_from_negative_syscall_ret(ret));
512 }
513 Ok(())
514}
515
516/// Traduit une [`Restriction`] en son `struct io_uring_restriction` : un opcode
517/// de **type** de restriction + un `u8` selon ce type (numéro d'opcode de
518/// soumission via [`opcode_number`], numéro de register-op, ou octet `IOSQE_*`).
519fn encode_restriction(restriction: Restriction) -> raw::IoUringRestriction {
520 let mut out = raw::IoUringRestriction::default();
521 match restriction {
522 Restriction::AllowOp(op) => {
523 out.opcode = raw::IORING_RESTRICTION_SQE_OP;
524 out.op_or_flags = opcode_number(op);
525 }
526 Restriction::AllowRegister(register_op) => {
527 out.opcode = raw::IORING_RESTRICTION_REGISTER_OP;
528 out.op_or_flags = register_op;
529 }
530 Restriction::SqeFlagsAllowed(options) => {
531 out.opcode = raw::IORING_RESTRICTION_SQE_FLAGS_ALLOWED;
532 out.op_or_flags = options.iosqe_flags();
533 }
534 Restriction::SqeFlagsRequired(options) => {
535 out.opcode = raw::IORING_RESTRICTION_SQE_FLAGS_REQUIRED;
536 out.op_or_flags = options.iosqe_flags();
537 }
538 }
539 out
540}
541
542impl IoUring {
543 /// Raccourci : `IoUringBuilder::new(entries).build()`.
544 ///
545 /// # Errors
546 ///
547 /// Voir [`IoUringBuilder::build`].
548 pub fn new(entries: NonZeroU32) -> Result<Self, Errno> {
549 IoUringBuilder::new(entries).build()
550 }
551
552 /// Active un ring créé avec `R_DISABLED` (`REGISTER_ENABLE_RINGS`).
553 ///
554 /// Inutile sur un ring déjà actif : retourne alors l'`EINVAL` du kernel.
555 ///
556 /// # Errors
557 ///
558 /// - [`Errno::EINVAL`] : ring déjà actif (remonté tel quel).
559 /// - [`Errno::EBADF`] : FD de ring invalide.
560 pub fn enable(&mut self) -> Result<(), Errno> {
561 // SAFETY: `fd` est un ring fd valide possédé ; `ENABLE_RINGS` ne lit
562 // aucune mémoire utilisateur (arg/nr_args nuls).
563 let ret =
564 unsafe { syscall::register(self.fd_raw(), raw::IORING_REGISTER_ENABLE_RINGS, 0, 0) };
565 if ret < 0 {
566 return Err(raw::errno_from_negative_syscall_ret(ret));
567 }
568 Ok(())
569 }
570
571 /// Vrai si la CQ a débordé (`IORING_SQ_CQ_OVERFLOW`). Avec `FEAT_NODROP`
572 /// (présent en 6.12) le kernel retient les complétions et les re-livre.
573 #[must_use]
574 pub fn completion_queue_overflowed(&self) -> bool {
575 self.sq.flags() & raw::IORING_SQ_CQ_OVERFLOW != 0
576 }
577
578 /// Nombre de places libres dans la SQ (calculé par `tail - head` masqué).
579 #[must_use]
580 pub fn submission_queue_space_left(&self) -> u32 {
581 self.sq.space_left()
582 }
583
584 /// Nombre d'opérations actuellement en vol (slots occupés du slab S1).
585 #[must_use]
586 pub fn in_flight(&self) -> u32 {
587 self.slab.in_flight()
588 }
589
590 /// Nombre de CQE bruts prêts dans la CQ (introspection de test — sert à
591 /// observer un CQE injecté par `msg_ring` côté cible, que `harvest_ready`
592 /// filtrerait comme périmé faute de slot correspondant).
593 #[cfg(test)]
594 pub(crate) fn cq_ready(&self) -> u32 {
595 self.cq.ready()
596 }
597
598 /// FD numérique du ring (toujours présent avant le téardown). `pub(crate)` :
599 /// le sous-module `registration` (Temps 3a) appelle `io_uring_register(2)`
600 /// directement sur ce FD.
601 pub(crate) fn fd_raw(&self) -> i32 {
602 self.fd.as_ref().map_or(-1, AsRawFd::as_raw_fd)
603 }
604
605 /// Cible d'un `io_uring_enter` : `(fd, flags additionnels)`. Avec un ring fd
606 /// enregistré (Temps 3a §4), l'`fd` est l'**index** enregistré et le flag
607 /// `IORING_ENTER_REGISTERED_RING` est OR-é ; sinon le FD ordinaire, sans flag.
608 fn enter_target(&self) -> (i32, u32) {
609 match self.enter_ring_index {
610 // L'index est petit (table de rings enregistrés) ⇒ tient dans un i32.
611 Some(index) => (index.cast_signed(), raw::IORING_ENTER_REGISTERED_RING),
612 None => (self.fd_raw(), 0),
613 }
614 }
615
616 /// Téardown propre (S2) : annule les ops en vol ([`IoUring::sync_cancel`]
617 /// `Any`), draine la CQ jusqu'à `in_flight() == 0`, `munmap` les anneaux et
618 /// ferme le FD. À préférer au `drop` implicite sur chemin chaud.
619 ///
620 /// # Errors
621 ///
622 /// Propage les erreurs de drainage/annulation (l'appelant décide de la
623 /// suite) ; voir [`IoUring::sync_cancel`].
624 pub fn shutdown(mut self) -> Result<(), Errno> {
625 let result = self.do_teardown();
626 // `do_teardown` a posé `disposed = true` et libéré/fuité les ressources.
627 // À la sortie, `self` est dropé → `Drop` voit `disposed` et ne refait
628 // **rien** (pas de double munmap/close/free — reminder S2 #3).
629 result
630 }
631
632 /// Quiescence + libération **ou fuite contrôlée** des ressources (idempotent
633 /// via `disposed`). Cœur de la décision S2 : on ne libère la mémoire que si
634 /// le kernel ne peut plus y écrire (`in_flight == 0`).
635 fn do_teardown(&mut self) -> Result<(), Errno> {
636 if self.disposed {
637 return Ok(());
638 }
639 self.disposed = true;
640 if self.quiesce() {
641 // RELEASE (sound : in_flight == 0, le kernel n'écrit plus) : munmap
642 // des 3 zones, close du fd, par RAII. Ordre : mmaps puis fd.
643 drop(self.sq_ring_map.take());
644 drop(self.sqes_map.take());
645 drop(self.cq_ring_map.take());
646 drop(self.fd.take());
647 Ok(())
648 } else {
649 // FUITE CONTRÔLÉE (Principe 5, reminder S2 #1) : la quiescence a
650 // échoué (op non annulable / timeout) — le kernel peut encore écrire
651 // dans les anneaux mmap et les buffers en vol. On ne libère RIEN
652 // (forget) : une fuite bornée à ce ring est sound, un UAF jamais.
653 leak_forget(self.sq_ring_map.take());
654 leak_forget(self.sqes_map.take());
655 leak_forget(self.cq_ring_map.take());
656 leak_forget(self.fd.take());
657 self.slab.leak_inflight_buffers();
658 Err(Errno::EBUSY)
659 }
660 }
661
662 /// Amène le ring à `in_flight == 0` : annule (best-effort) puis draine les
663 /// complétions (y compris `-ECANCELED`, déroulement **nominal**, reminder
664 /// S2 #2). Borné en temps. Retourne `true` si quiescent.
665 fn quiesce(&mut self) -> bool {
666 if self.slab.in_flight() == 0 {
667 return true;
668 }
669 // (a) **SOUMETTRE D'ABORD.** Le slab compte une op dès sa MISE EN FILE, mais le
670 // noyau ne la connaît qu'après l'`enter`. Une op mise en file et jamais soumise
671 // — le cas d'un anneau construit puis détruit sans avoir tourné une seule fois —
672 // était donc invisible aux DEUX étapes suivantes : `sync_cancel` ne trouvait rien
673 // à annuler, et le drainage attendait une complétion que rien ne pouvait produire.
674 // Résultat mesuré : 64 tours de 100 ms, soit 6,4 s, puis « fuite contrôlée » du fd
675 // d'anneau et de ses trois mappings, à CHAQUE destruction.
676 //
677 // Ce n'était pas un défaut de l'appelant : la soumission paresseuse est le motif
678 // normal (on met en file, on soumet au prochain tour). C'est la quiescence qui
679 // supposait, à tort, que tout ce qu'elle compte a été vu par le noyau.
680 //
681 // Best-effort, comme le reste du téardown : si l'`enter` échoue, on tente quand
682 // même l'annulation et le drainage — on ne rend pas la situation pire.
683 let _ = self.submit();
684 // (b) Annulation globale (best-effort : ENOENT/ETIME ignorés).
685 let _ = self.sync_cancel(CancelTarget::Any);
686 // (c) Drainage borné jusqu'à in_flight == 0.
687 const MAX_ROUNDS: u32 = 64;
688 for _ in 0..MAX_ROUNDS {
689 // Consomme tout ce qui est prêt (les CQE `-ECANCELED` libèrent leur
690 // slot via `harvest_ready`, sans propager d'erreur).
691 while self.harvest_ready().is_some() {}
692 if self.slab.in_flight() == 0 {
693 return true;
694 }
695 // Bloque (timeout borné) pour au moins une complétion supplémentaire.
696 if self.wait_bounded().is_err() {
697 // EINTR/ETIME : on retente quelques rounds avant d'abandonner.
698 }
699 }
700 self.slab.in_flight() == 0
701 }
702
703 /// `io_uring_enter(GETEVENTS, min_complete = 1)` avec un timeout court borné
704 /// (100 ms via `EXT_ARG`), pour le drainage de quiescence.
705 fn wait_bounded(&self) -> Result<(), Errno> {
706 let ts = raw::KernelTimespec {
707 tv_sec: 0,
708 tv_nsec: 100_000_000,
709 };
710 let arg = raw::IoUringGeteventsArg {
711 sigmask: 0,
712 sigmask_sz: 0,
713 min_wait_usec: 0,
714 ts: core::ptr::from_ref(&ts) as u64,
715 };
716 let arg_sz = u64::try_from(core::mem::size_of::<raw::IoUringGeteventsArg>())
717 .expect("taille getevents_arg ≤ u64");
718 let (efd, eflag) = self.enter_target();
719 // SAFETY: `fd`/index valide ; `arg`/`ts` vivants pour la durée de l'appel.
720 let ret = unsafe {
721 syscall::enter(
722 efd,
723 0,
724 1,
725 raw::IORING_ENTER_GETEVENTS | raw::IORING_ENTER_EXT_ARG | eflag,
726 core::ptr::from_ref(&arg) as u64,
727 arg_sz,
728 )
729 };
730 if ret < 0 {
731 return Err(raw::errno_from_negative_syscall_ret(ret));
732 }
733 Ok(())
734 }
735}
736
737/// Interroge `IORING_REGISTER_PROBE` et retourne, par numéro d'opcode, le
738/// support du kernel courant (mis en cache par [`IoUringBuilder::build`], §9).
739///
740/// Best-effort : si le probe échoue (kernel ancien, sandbox/`R_DISABLED`), tous
741/// les opcodes sont marqués **non supportés** — conservateur, force le fallback
742/// synchrone (ADR-022 D8).
743fn probe_supported_ops(fd: i32) -> [bool; 256] {
744 let mut probe = Box::new(raw::IoUringProbe {
745 last_op: 0,
746 ops_len: 0,
747 resv: 0,
748 resv2: [0; 3],
749 ops: [raw::IoUringProbeOp::default(); 256],
750 });
751 let probe_ptr = core::ptr::from_mut(&mut *probe) as u64;
752 // SAFETY: `fd` ring valide ; `probe` pointe une `io_uring_probe` vivante,
753 // dimensionnée pour `IO_URING_PROBE_OPS` entrées, écrite par le kernel.
754 let ret = unsafe {
755 syscall::register(
756 fd,
757 raw::IORING_REGISTER_PROBE,
758 probe_ptr,
759 raw::IO_URING_PROBE_OPS,
760 )
761 };
762 let mut supported = [false; 256];
763 if ret < 0 {
764 return supported;
765 }
766 for (op, entry) in probe.ops.iter().enumerate() {
767 supported[op] = entry.flags & raw::IO_URING_OP_SUPPORTED != 0;
768 }
769 supported
770}
771
772/// Fuite contrôlée d'une ressource RAII optionnelle (`forget` si présente).
773fn leak_forget<T>(resource: Option<T>) {
774 if let Some(value) = resource {
775 core::mem::forget(value);
776 }
777}
778
779impl Drop for IoUring {
780 /// Filet de sécurité (S2) : si des ops sont en vol, quiesce (annule +
781 /// draine, **bloquant et best-effort**, boucle bornée) avant de libérer la
782 /// mémoire. Coût (blocage potentiel) **assumé et documenté** : sur chemin
783 /// chaud, préférer [`IoUring::shutdown`]. Un `Drop` qui laisserait le kernel
784 /// écrire dans de la mémoire libérée serait un défaut de soundness
785 /// inacceptable en couche 0 (Principe 5 : sur-sécuriser puis dégraisser).
786 fn drop(&mut self) {
787 // Filet de sécurité S2 : quiesce (annule + draine, **bloquant**,
788 // best-effort borné) puis libère — ou **fuite** de façon contrôlée si la
789 // quiescence échoue (jamais d'UAF). No-op si `shutdown` a déjà disposé.
790 let _ = self.do_teardown();
791 }
792}
793
794// ---------------------------------------------------------------------------
795// Soumission
796// ---------------------------------------------------------------------------
797
798impl IoUring {
799 /// Soumet les SQE en attente (`io_uring_enter`, sans attente). Retourne le
800 /// nombre de SQE consommés par le kernel.
801 ///
802 /// Publie la queue SQ par un **store release** (§3.2) avant l'`enter`. En
803 /// mode `SQPOLL` chaud, peut n'avoir **aucun** syscall à faire (juste la
804 /// publication), ne réveillant le thread que si `SQ_NEED_WAKEUP`.
805 ///
806 /// # Errors
807 ///
808 /// - [`Errno::EINTR`] : remonté tel quel (ADR-021 conv. 2 ; l'appelant écrit
809 /// sa propre boucle de retry).
810 /// - [`Errno::EAGAIN`] : ressources momentanément indisponibles.
811 /// - [`Errno::EBUSY`] : CQ pleine non drainée (selon contexte).
812 /// - [`Errno::EINVAL`], [`Errno::EFAULT`], [`Errno::EBADF`].
813 pub fn submit(&mut self) -> Result<u32, Errno> {
814 self.debug_assert_raw_tags();
815 let to_submit = self.sq.publish_and_pending();
816 if to_submit == 0 {
817 return Ok(0);
818 }
819 let (efd, eflag) = self.enter_target();
820 // SAFETY: `fd`/index est valide ; pas d'argument étendu.
821 let ret = unsafe { syscall::enter(efd, to_submit, 0, eflag, 0, 0) };
822 if ret < 0 {
823 // EINTR/EAGAIN remontés tels quels (ADR-021 conv. 2 — aucun retry).
824 return Err(raw::errno_from_negative_syscall_ret(ret));
825 }
826 Ok(u32::try_from(ret).unwrap_or(0))
827 }
828
829 /// Soumet puis attend au moins `want` complétions
830 /// (`IORING_ENTER_GETEVENTS`, `min_complete = want`).
831 ///
832 /// # Errors
833 ///
834 /// Mêmes erreurs que [`IoUring::submit`] ; [`Errno::EINTR`] remonté tel quel.
835 pub fn submit_and_wait(&mut self, want: u32) -> Result<u32, Errno> {
836 self.debug_assert_raw_tags();
837 let to_submit = self.sq.publish_and_pending();
838 let (efd, eflag) = self.enter_target();
839 // SAFETY: `fd`/index est valide ; pas d'argument étendu.
840 let ret = unsafe {
841 syscall::enter(
842 efd,
843 to_submit,
844 want,
845 raw::IORING_ENTER_GETEVENTS | eflag,
846 0,
847 0,
848 )
849 };
850 if ret < 0 {
851 // EINTR remonté tel quel (ADR-021 conv. 2).
852 return Err(raw::errno_from_negative_syscall_ret(ret));
853 }
854 Ok(u32::try_from(ret).unwrap_or(0))
855 }
856
857 /// Applique des options par-opération ([`SubmitOptions`]) à la **prochaine**
858 /// soumission `submit_*`. Style builder (chaînable).
859 pub fn with(&mut self, opts: SubmitOptions) -> &mut Self {
860 self.pending_options = opts;
861 self
862 }
863
864 /// Exécute la **prochaine** soumission `submit_*` avec les credentials d'une
865 /// [`Personality`] enregistrée (`sqe.personality`, Temps 3a §6). Consommée
866 /// pour cette seule opération. Style builder (chaînable).
867 pub fn with_personality(&mut self, personality: Personality) -> &mut Self {
868 self.pending_personality = personality.raw();
869 self
870 }
871
872 /// Prépare et met en file un `IORING_OP_NOP` — la première opération de
873 /// validation du cœur. Réserve un slot S1, écrit le SQE, et retourne le
874 /// jeton ; la soumission effective se fait par [`IoUring::submit`] /
875 /// [`IoUring::submit_and_wait`].
876 ///
877 /// # Errors
878 ///
879 /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins (back-pressure
880 /// structurelle, **aucun** syscall — §4.3).
881 pub fn submit_nop(&mut self) -> Result<SubmissionToken, Errno> {
882 if self.sq.space_left() == 0 {
883 return Err(Errno::EBUSY);
884 }
885 let token = self.slab.reserve(None).map_err(|_| Errno::EBUSY)?;
886 // SAFETY: place SQ vérifiée > 0 ; le SQE est rempli intégralement
887 // ci-dessous avant toute publication de la `tail`.
888 let sqe = unsafe { self.sq.prepare() }.expect("place SQ vérifiée");
889 let opts = self.pending_options;
890 // SAFETY: `sqe` pointe un slot SQE valide fraîchement zéro-initialisé,
891 // exclusivement détenu jusqu'à la publication.
892 unsafe {
893 (*sqe).opcode = raw::IORING_OP_NOP;
894 (*sqe).flags = opts.iosqe_flags();
895 (*sqe).user_data = token.to_user_data();
896 }
897 self.pending_options = SubmitOptions::default();
898 if opts.skips_cqe_on_success() {
899 // Aucun CQE attendu en cas de succès → libère le slot S1 dès la
900 // soumission (§6.2). Un échec produirait un CQE périmé, filtré.
901 self.slab.release(token);
902 }
903 Ok(token)
904 }
905
906 /// Moissonne une complétion **non périmée** déjà disponible (sans syscall),
907 /// en filtrant les complétions périmées (§4.2). `None` si la CQ est vide.
908 fn harvest_ready(&mut self) -> Option<Completion> {
909 loop {
910 let cqe = self.cq.peek()?;
911 if cqe.user_data & raw::RAW_USER_DATA_TAG != 0 {
912 // Complétion **brute** (Temps 4 §5) : ne PAS toucher au slab ni
913 // consommer le CQE — laissé à `raw_peek_completion_queue_entry` /
914 // `raw_advance_completion_queue`. La moisson gérée s'arrête ici
915 // (FIFO : aucune complétion gérée disponible tant que cette brute,
916 // en tête, n'est pas consommée par l'appelant).
917 return None;
918 }
919 self.cq.advance();
920 let flags = CompletionFlags::from_bits_truncate(cqe.flags);
921 // Marquage multishot lu **avant** `complete` (qui peut libérer le slot
922 // à la terminale) : permet à `Completion::multishot_token` de
923 // distinguer un flux multishot d'une op mono-coup (Temps 3d).
924 let multishot = self.slab.is_multishot(cqe.user_data);
925 match self
926 .slab
927 .complete(cqe.user_data, flags.contains(CompletionFlags::MORE))
928 {
929 slab::SlotOutcome::Stale => {} // périmée : passer à la suivante
930 slab::SlotOutcome::More => {
931 return Some(Completion {
932 token: SubmissionToken::from_user_data(cqe.user_data),
933 res: cqe.res,
934 raw_flags: cqe.flags,
935 payload: None,
936 multishot,
937 });
938 }
939 slab::SlotOutcome::Final { payload } => {
940 return Some(Completion {
941 token: SubmissionToken::from_user_data(cqe.user_data),
942 res: cqe.res,
943 raw_flags: cqe.flags,
944 payload,
945 multishot,
946 });
947 }
948 }
949 }
950 }
951}
952
953// ───────────────────────────────────────────────────────────────────────────
954// Temps 4 — accès brut niveau 1 (soupape ADR-022 Décision 1)
955// ───────────────────────────────────────────────────────────────────────────
956impl IoUring {
957 /// Capacité totale de la SQ (entrées effectives, arrondi kernel). **Sûr.**
958 #[must_use]
959 pub fn submission_queue_capacity(&self) -> u32 {
960 self.sq.capacity()
961 }
962
963 /// Emplacements SQ actuellement disponibles. **Sûr.**
964 #[must_use]
965 pub fn submission_queue_available(&self) -> u32 {
966 self.sq.space_left()
967 }
968
969 /// Capacité totale de la CQ. **Sûr.**
970 #[must_use]
971 pub fn completion_queue_capacity(&self) -> u32 {
972 self.cq.capacity()
973 }
974
975 /// Réserve un emplacement SQE libre et le rend pour **remplissage manuel**
976 /// (accès brut niveau 1, Temps 4). `None` si la SQ est pleine. La
977 /// **publication reste `submit()`/`submit_and_wait()`** (sûrs) : la façade
978 /// fait le store release et garde la main sur l'ordering d'anneau ; l'`unsafe`
979 /// ne porte donc QUE sur le **contenu** du SQE et la **validité des buffers**.
980 ///
981 /// # Safety
982 ///
983 /// L'appelant garantit (contrat §6 de la spec) :
984 /// 1. **opcode valide** et supporté par le kernel courant ([`IoUring::supports_op`]/probe) ;
985 /// 2. **cohérence des champs** du SQE (`fd`/`addr`/`len`/flags) pour cet opcode ;
986 /// 3. **validité mémoire** : tout buffer pointé par `addr`/`addr2`/`addr3` reste
987 /// valide et **non déplacé jusqu'à la complétion** (le slab S1 ne le protège
988 /// **pas** en accès brut) ;
989 /// 4. **tag `user_data`** : poser [`RAW_USER_DATA_TAG`] (§5) pour coexister avec
990 /// le niveau 2 sans collision (vérifié en build debug à la soumission) ;
991 /// 5. **pas de double consommation** de la complétion correspondante
992 /// ([`IoUring::raw_advance_completion_queue`] cohérent avec les peeks).
993 pub unsafe fn raw_get_submission_queue_entry(
994 &mut self,
995 ) -> Option<&mut RawSubmissionQueueEntry> {
996 let tail = self.sq.local_tail();
997 // SAFETY: préconditions déléguées à l'appelant (contrat ci-dessus) ;
998 // `prepare` rend un slot fraîchement zéro-initialisé, exclusivement détenu
999 // jusqu'à la publication par `submit`.
1000 let sqe = unsafe { self.sq.prepare() }?;
1001 self.raw_pending.push(tail);
1002 // SAFETY: `RawSubmissionQueueEntry` est `#[repr(transparent)]` sur
1003 // `IoUringSqe` (layout identique, asserté statiquement) ; `sqe` désigne un
1004 // slot SQE valide de la mmap, emprunté `&mut` exclusivement le temps du
1005 // remplissage (avant publication).
1006 Some(unsafe { &mut *sqe.cast::<RawSubmissionQueueEntry>() })
1007 }
1008
1009 /// Inspecte la prochaine complétion **sans la consommer** (`None` si la CQ est
1010 /// vide). **Sûr** : lecture seule via le load acquire interne. L'interprétation
1011 /// de `res`/`flags` et la lecture du tag `user_data` (§5) sont à la charge de
1012 /// l'appelant (selon l'opcode soumis).
1013 #[must_use]
1014 pub fn raw_peek_completion_queue_entry(&self) -> Option<&RawCompletionQueueEntry> {
1015 let cqe = self.cq.peek_entry()?;
1016 // SAFETY: `cqe` pointe un slot CQE vivant de la mmap (stable tant qu'aucune
1017 // avance ; l'emprunt `&self` interdit `advance`/`&mut` concurrent) ;
1018 // `RawCompletionQueueEntry` est un miroir `#[repr(C)]` d'`IoUringCqe`
1019 // (layout identique, asserté).
1020 Some(unsafe { &*cqe.cast::<RawCompletionQueueEntry>() })
1021 }
1022
1023 /// Consomme `n` complétions de la CQ (avance la tête + store release par la
1024 /// façade). **Sûr** : borné au nombre réellement prêt (Principe 4) ; l'appelant
1025 /// garantit la cohérence avec ses [`IoUring::raw_peek_completion_queue_entry`]
1026 /// (pas de double consommation, §6 pt 5).
1027 pub fn raw_advance_completion_queue(&mut self, n: u32) {
1028 let n = n.min(self.cq.ready());
1029 for _ in 0..n {
1030 self.cq.advance();
1031 }
1032 }
1033
1034 /// Vérifie (build **debug**) que tout SQE préparé via
1035 /// [`IoUring::raw_get_submission_queue_entry`] porte [`RAW_USER_DATA_TAG`]
1036 /// (§5) — sinon une op brute collisionnerait avec l'encodage du slab S1. En
1037 /// release, `debug_assert!` est inerte ; la purge de `raw_pending` a lieu dans
1038 /// tous les cas (idempotente, sans allocation sur liste vide).
1039 fn debug_assert_raw_tags(&mut self) {
1040 // COUVERTURE (ADR-035, STRUCTURAL — cf. docs/COVERAGE-EXCEPTIONS.md,
1041 // section `io_uring`). Le `}` fermant ce bloc `if cfg!(debug_assertions)`
1042 // porte une région LLVM dégénérée à `count == 0` alors que le **corps**
1043 // de la boucle est prouvé exécuté (mesure : 308 itérations via les tests
1044 // `raw_nop_*`). Artefact d'instrumentation du `cfg!`-constant analogue aux
1045 // ombres `[True:0, False:0]` documentées en tête de registre — non
1046 // couvrable par aucun test.
1047 if cfg!(debug_assertions) {
1048 for &tail in &self.raw_pending {
1049 // SAFETY: `tail` désigne un SQE préparé pendant ce staging, non
1050 // encore publié ; lecture seule de son `user_data`.
1051 let user_data = unsafe { self.sq.sqe_user_data_at(tail) };
1052 debug_assert!(
1053 user_data & raw::RAW_USER_DATA_TAG != 0,
1054 "opération brute (Temps 4) soumise sans RAW_USER_DATA_TAG : \
1055 collision possible avec le slab S1 (§5)"
1056 );
1057 }
1058 }
1059 self.raw_pending.clear();
1060 }
1061}
1062
1063/// Jeton opaque reliant une soumission à sa complétion (décision S1).
1064///
1065/// Encapsule un index de slot du slab d'opérations en vol **et** une génération
1066/// (compteur anti-réutilisation, protège contre l'ABA et les complétions
1067/// périmées). Le `user_data` kernel est dérivé en interne — l'API ne laisse
1068/// **jamais** l'appelant manipuler un `user_data` brut.
1069#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1070pub struct SubmissionToken {
1071 slot: u32,
1072 generation: u32,
1073}
1074
1075impl SubmissionToken {
1076 /// Construit un jeton à partir d'un index de slot et d'une génération.
1077 pub(crate) const fn new(slot: u32, generation: u32) -> Self {
1078 Self { slot, generation }
1079 }
1080
1081 /// Index de slot dans le slab d'opérations en vol.
1082 pub(crate) const fn slot(self) -> u32 {
1083 self.slot
1084 }
1085
1086 /// Génération du slot au moment de la réservation (anti-réutilisation/ABA).
1087 pub(crate) const fn generation(self) -> u32 {
1088 self.generation
1089 }
1090
1091 /// Encode le jeton en `user_data` kernel :
1092 /// `((generation as u64) << 32) | (slot as u64)` (§4.2).
1093 pub(crate) fn to_user_data(self) -> u64 {
1094 // `wrapping_shl(32)` : la valeur (u32 élargie) ne déborde jamais sur 64
1095 // bits ; choix wrapping explicite pour rester lint-clean (Principe 2).
1096 u64::from(self.generation).wrapping_shl(32) | u64::from(self.slot)
1097 }
1098
1099 /// Décode un `user_data` kernel en `(generation, slot)`.
1100 ///
1101 /// # Panics
1102 ///
1103 /// Jamais : les masques 32 bits garantissent que chaque moitié tient dans
1104 /// un `u32` (les `expect` sont structurellement inatteignables).
1105 pub(crate) fn from_user_data(user_data: u64) -> Self {
1106 let slot = u32::try_from(user_data & 0xFFFF_FFFF).expect("masque 32 bits bas");
1107 let generation = u32::try_from(user_data >> 32).expect("décalage 32 bits haut");
1108 Self { slot, generation }
1109 }
1110}
1111
1112/// Jeton d'une opération multishot (un slot, plusieurs complétions). Voir
1113/// [`multishot`] (Temps 3d).
1114#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1115pub struct MultishotToken {
1116 slot: u32,
1117 generation: u32,
1118}
1119
1120impl MultishotToken {
1121 /// Construit un jeton multishot à partir d'un index de slot et d'une
1122 /// génération (sémantique Temps 3d).
1123 pub(crate) const fn new(slot: u32, generation: u32) -> Self {
1124 Self { slot, generation }
1125 }
1126
1127 /// `SubmissionToken` équivalent (même `(slot, generation)`) — pour
1128 /// l'annulation par jeton et la corrélation aux complétions.
1129 pub(crate) const fn as_submission(self) -> SubmissionToken {
1130 SubmissionToken::new(self.slot, self.generation)
1131 }
1132}
1133
1134/// Options de soumission par-opération (flags `IOSQE_*` exposés sûrement).
1135///
1136/// Construit par combinateurs ; appliqué via [`IoUring::with`]. Les liens
1137/// (`IO_LINK`/`IO_HARDLINK`) passent par le `LinkedChainBuilder` (Temps 3c) ; la
1138/// sélection de buffer (`BUFFER_SELECT`) par `ProvidedBufferRing` (Temps 3b).
1139#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1140pub struct SubmitOptions {
1141 // Représentation interne (Pass B) : bitmask IOSQE_*.
1142 flags: u32,
1143}
1144
1145impl SubmitOptions {
1146 /// Active `IOSQE_IO_DRAIN` : draine les opérations déjà soumises avant
1147 /// celle-ci.
1148 #[must_use]
1149 pub fn drain(mut self) -> Self {
1150 self.flags |= u32::from(raw::IOSQE_IO_DRAIN);
1151 self
1152 }
1153
1154 /// Active `IOSQE_ASYNC` : force l'exécution sur le pool io-wq.
1155 #[must_use]
1156 pub fn force_async(mut self) -> Self {
1157 self.flags |= u32::from(raw::IOSQE_ASYNC);
1158 self
1159 }
1160
1161 /// Active `IOSQE_CQE_SKIP_SUCCESS` : pas de CQE en cas de succès (libère le
1162 /// slot S1 **à la soumission**). Nécessite la feature `FEAT_CQE_SKIP`.
1163 #[must_use]
1164 pub fn skip_cqe_on_success(mut self) -> Self {
1165 self.flags |= u32::from(raw::IOSQE_CQE_SKIP_SUCCESS);
1166 self
1167 }
1168
1169 /// Active `IOSQE_FIXED_FILE` : le `fd` du SQE est un **index de slot fixe**
1170 /// (Temps 3a — `FixedSlot`/`fixed_fd_install`). Interne : posé par les
1171 /// façades fixed-file, jamais par l'appelant.
1172 pub(crate) fn fixed_file(mut self) -> Self {
1173 self.flags |= u32::from(raw::IOSQE_FIXED_FILE);
1174 self
1175 }
1176
1177 /// Active `IOSQE_BUFFER_SELECT` : l'op sélectionne automatiquement un buffer
1178 /// dans un groupe fourni (Temps 3b). Interne : posé par les façades
1179 /// `submit_*_provided`, jamais par l'appelant.
1180 pub(crate) fn buffer_select(mut self) -> Self {
1181 self.flags |= u32::from(raw::IOSQE_BUFFER_SELECT);
1182 self
1183 }
1184
1185 /// Octet de flags `IOSQE_*` à poser dans le SQE (tous < 256).
1186 fn iosqe_flags(self) -> u8 {
1187 u8::try_from(self.flags).unwrap_or(0)
1188 }
1189
1190 /// Vrai si `skip_cqe_on_success` est demandé (libération anticipée du slot).
1191 fn skips_cqe_on_success(self) -> bool {
1192 self.flags & u32::from(raw::IOSQE_CQE_SKIP_SUCCESS) != 0
1193 }
1194}
1195
1196// ---------------------------------------------------------------------------
1197// Complétion
1198// ---------------------------------------------------------------------------
1199
1200impl IoUring {
1201 /// Attend (bloquant) la prochaine complétion et la consomme
1202 /// (`io_uring_enter(.., GETEVENTS, min_complete = 1)` ; lecture de la CQ par
1203 /// **load acquire**, §3.2). Les complétions **périmées** (génération de slot
1204 /// non concordante, §4.2) sont filtrées en silence.
1205 ///
1206 /// # Errors
1207 ///
1208 /// - [`Errno::EINTR`] : remonté tel quel.
1209 /// - [`Errno::EBADF`].
1210 pub fn wait_completion(&mut self) -> Result<Completion, Errno> {
1211 loop {
1212 if let Some(completion) = self.harvest_ready() {
1213 return Ok(completion);
1214 }
1215 // Aucune complétion non périmée disponible : bloque jusqu'à au moins
1216 // une via `io_uring_enter(GETEVENTS, min_complete = 1)`.
1217 let (efd, eflag) = self.enter_target();
1218 // SAFETY: `fd`/index est valide ; pas d'argument étendu.
1219 let ret =
1220 unsafe { syscall::enter(efd, 0, 1, raw::IORING_ENTER_GETEVENTS | eflag, 0, 0) };
1221 if ret < 0 {
1222 // EINTR remonté tel quel (ADR-021 conv. 2). La boucle ne
1223 // re-`enter` que sur succès renvoyant des CQE périmés à filtrer.
1224 return Err(raw::errno_from_negative_syscall_ret(ret));
1225 }
1226 }
1227 }
1228
1229 /// Attend une complétion avec délai maximal (`IORING_ENTER_EXT_ARG` +
1230 /// `io_uring_getevents_arg`, ou `ABS_TIMER` si `FEAT_MIN_TIMEOUT`).
1231 /// Retourne `Ok(None)` à l'expiration.
1232 ///
1233 /// # Errors
1234 ///
1235 /// - [`Errno::EINTR`] : remonté tel quel.
1236 /// - [`Errno::EBADF`], [`Errno::EINVAL`].
1237 pub fn wait_completion_timeout(
1238 &mut self,
1239 timeout: Duration,
1240 ) -> Result<Option<Completion>, Errno> {
1241 if let Some(completion) = self.harvest_ready() {
1242 return Ok(Some(completion));
1243 }
1244 let ts = raw::KernelTimespec {
1245 tv_sec: i64::try_from(timeout.as_secs()).unwrap_or(i64::MAX),
1246 tv_nsec: i64::from(timeout.subsec_nanos()),
1247 };
1248 let arg = raw::IoUringGeteventsArg {
1249 sigmask: 0,
1250 sigmask_sz: 0,
1251 min_wait_usec: 0,
1252 ts: core::ptr::from_ref(&ts) as u64,
1253 };
1254 let arg_sz = u64::try_from(core::mem::size_of::<raw::IoUringGeteventsArg>())
1255 .expect("taille getevents_arg ≤ u64");
1256 let (efd, eflag) = self.enter_target();
1257 // SAFETY: `fd`/index valide ; `arg` pointe une `io_uring_getevents_arg`
1258 // vivante de taille `arg_sz` cohérente avec le flag `EXT_ARG`, et `ts`
1259 // reste en vie pendant l'appel (variables locales).
1260 let ret = unsafe {
1261 syscall::enter(
1262 efd,
1263 0,
1264 1,
1265 raw::IORING_ENTER_GETEVENTS | raw::IORING_ENTER_EXT_ARG | eflag,
1266 core::ptr::from_ref(&arg) as u64,
1267 arg_sz,
1268 )
1269 };
1270 if ret < 0 {
1271 let err = raw::errno_from_negative_syscall_ret(ret);
1272 if err == Errno::ETIME {
1273 // Expiration : aucune complétion dans le délai imparti.
1274 return Ok(None);
1275 }
1276 // EINTR et autres errno remontés tels quels (ADR-021 conv. 2).
1277 return Err(err);
1278 }
1279 Ok(self.harvest_ready())
1280 }
1281
1282 /// Récupère une complétion si disponible, **sans bloquer** ni syscall
1283 /// (load acquire de la CQ). Filtre les complétions périmées.
1284 pub fn try_completion(&mut self) -> Option<Completion> {
1285 self.harvest_ready()
1286 }
1287
1288 /// Itère les complétions actuellement disponibles dans la CQ, avançant
1289 /// `head` (store release) au fil de l'itération.
1290 pub fn completions(&mut self) -> CompletionIter<'_> {
1291 CompletionIter { ring: self }
1292 }
1293
1294 /// Vrai si le kernel courant supporte `op` (`IORING_REGISTER_PROBE`, mis en
1295 /// cache à la construction ; ADR-022 Décision 8). Permet le fallback vers le
1296 /// syscall synchrone correspondant.
1297 #[must_use]
1298 pub fn supports_op(&self, op: IoUringOpcode) -> bool {
1299 self.supported_ops
1300 .get(usize::from(opcode_number(op)))
1301 .copied()
1302 .unwrap_or(false)
1303 }
1304
1305 /// Features io_uring négociées au setup (axe G), exposées par prédicats
1306 /// stables.
1307 #[must_use]
1308 pub fn capabilities(&self) -> IoUringCapabilities {
1309 self.capabilities
1310 }
1311
1312 /// Annulation synchrone globale ou ciblée (`IORING_REGISTER_SYNC_CANCEL`,
1313 /// avec timeout). Brique de [`IoUring::shutdown`] (S2). Retourne le nombre
1314 /// d'opérations annulées.
1315 ///
1316 /// # Errors
1317 ///
1318 /// - [`Errno::EINTR`] : remonté tel quel.
1319 /// - [`Errno::EALREADY`] : annulation déjà en cours pour la cible.
1320 /// - [`Errno::ENOENT`] : aucune opération ne correspond à la cible.
1321 /// - [`Errno::EINVAL`].
1322 pub fn sync_cancel(&mut self, target: CancelTarget<'_>) -> Result<u32, Errno> {
1323 let mut reg = raw::IoUringSyncCancelReg {
1324 // Timeout interne borné : le kernel attend jusqu'à 2 s que les
1325 // annulations postent leurs CQE, puis rend la main.
1326 timeout_sec: 2,
1327 ..raw::IoUringSyncCancelReg::default()
1328 };
1329 match target {
1330 CancelTarget::Token(token) => reg.addr = token.to_user_data(),
1331 CancelTarget::Fd(fd) => {
1332 reg.fd = fd.as_raw_fd();
1333 reg.flags = raw::IORING_ASYNC_CANCEL_FD;
1334 }
1335 CancelTarget::Op(op) => {
1336 reg.opcode = opcode_number(op);
1337 reg.flags = raw::IORING_ASYNC_CANCEL_OP;
1338 }
1339 CancelTarget::Any => reg.flags = raw::IORING_ASYNC_CANCEL_ANY,
1340 }
1341 let reg_ptr = core::ptr::from_ref(®) as u64;
1342 // SAFETY: `fd` ring valide ; `reg` pointe une `io_uring_sync_cancel_reg`
1343 // vivante (variable locale) pour la durée de l'appel (`nr_args = 1`).
1344 let ret = unsafe {
1345 syscall::register(self.fd_raw(), raw::IORING_REGISTER_SYNC_CANCEL, reg_ptr, 1)
1346 };
1347 if ret < 0 {
1348 return Err(raw::errno_from_negative_syscall_ret(ret));
1349 }
1350 Ok(u32::try_from(ret).unwrap_or(0))
1351 }
1352}
1353
1354/// Numéro d'opcode kernel `IORING_OP_*` correspondant à un [`IoUringOpcode`].
1355///
1356/// Les 3 opcodes obsolètes (18 `OPENAT`, 31 `PROVIDE_BUFFERS`, 32
1357/// `REMOVE_BUFFERS`) ne figurent pas dans l'énumération : leurs numéros sont
1358/// donc absents (trous 18/31/32).
1359fn opcode_number(op: IoUringOpcode) -> u8 {
1360 match op {
1361 IoUringOpcode::Nop => 0,
1362 IoUringOpcode::Readv => 1,
1363 IoUringOpcode::Writev => 2,
1364 IoUringOpcode::Fsync => 3,
1365 IoUringOpcode::ReadFixed => 4,
1366 IoUringOpcode::WriteFixed => 5,
1367 IoUringOpcode::PollAdd => 6,
1368 IoUringOpcode::PollRemove => 7,
1369 IoUringOpcode::SyncFileRange => 8,
1370 IoUringOpcode::Sendmsg => 9,
1371 IoUringOpcode::Recvmsg => 10,
1372 IoUringOpcode::Timeout => 11,
1373 IoUringOpcode::TimeoutRemove => 12,
1374 IoUringOpcode::Accept => 13,
1375 IoUringOpcode::AsyncCancel => 14,
1376 IoUringOpcode::LinkTimeout => 15,
1377 IoUringOpcode::Connect => 16,
1378 IoUringOpcode::Fallocate => 17,
1379 IoUringOpcode::Close => 19,
1380 IoUringOpcode::FilesUpdate => 20,
1381 IoUringOpcode::Statx => 21,
1382 IoUringOpcode::Read => 22,
1383 IoUringOpcode::Write => 23,
1384 IoUringOpcode::Fadvise => 24,
1385 IoUringOpcode::Madvise => 25,
1386 IoUringOpcode::Send => 26,
1387 IoUringOpcode::Recv => 27,
1388 IoUringOpcode::Openat2 => 28,
1389 IoUringOpcode::EpollCtl => 29,
1390 IoUringOpcode::Splice => 30,
1391 IoUringOpcode::Tee => 33,
1392 IoUringOpcode::Shutdown => 34,
1393 IoUringOpcode::Renameat => 35,
1394 IoUringOpcode::Unlinkat => 36,
1395 IoUringOpcode::Mkdirat => 37,
1396 IoUringOpcode::Symlinkat => 38,
1397 IoUringOpcode::Linkat => 39,
1398 IoUringOpcode::MsgRing => 40,
1399 IoUringOpcode::Fsetxattr => 41,
1400 IoUringOpcode::Setxattr => 42,
1401 IoUringOpcode::Fgetxattr => 43,
1402 IoUringOpcode::Getxattr => 44,
1403 IoUringOpcode::Socket => 45,
1404 IoUringOpcode::UringCmd => 46,
1405 IoUringOpcode::SendZc => 47,
1406 IoUringOpcode::SendmsgZc => 48,
1407 IoUringOpcode::ReadMultishot => 49,
1408 IoUringOpcode::Waitid => 50,
1409 IoUringOpcode::FutexWait => 51,
1410 IoUringOpcode::FutexWake => 52,
1411 IoUringOpcode::FutexWaitv => 53,
1412 IoUringOpcode::FixedFdInstall => 54,
1413 IoUringOpcode::Ftruncate => 55,
1414 IoUringOpcode::Bind => 56,
1415 IoUringOpcode::Listen => 57,
1416 // `IoUringOpcode` est `#[non_exhaustive]` : un opcode d'un kernel futur
1417 // non encore mappé ⇒ numéro invalide (`255`), qui ne matchera aucune
1418 // opération réelle côté `sync_cancel`/restriction (plutôt qu'un faux
1419 // match sur l'opcode 0).
1420 _ => u8::MAX,
1421 }
1422}
1423
1424/// Une complétion (CQE décodé). Au Temps 1, les interprétations sont
1425/// **génériques** ; les variantes typées riches (fd accepté/ouvert, etc.)
1426/// arrivent aux Temps 2a–2d (ADR-022 Décision 9). Convention de résultat : un
1427/// `res` négatif est un `-errno`, converti en `Err(Errno)` par les méthodes
1428/// `into_*` ; un `res ≥ 0` est la valeur utile.
1429pub struct Completion {
1430 /// Jeton décodé depuis `user_data`.
1431 token: SubmissionToken,
1432 /// Résultat brut du kernel (`res` du CQE).
1433 res: i32,
1434 /// Flags bruts du CQE (`IORING_CQE_F_*` + id de buffer dans les bits hauts).
1435 raw_flags: u32,
1436 /// Payload possédé récupéré du slot S1 (`None` si l'op n'en portait pas).
1437 /// Extrait par les accesseurs typés (`into_buffer_result`, `into_statx`…).
1438 payload: Option<OwnedOp>,
1439 /// `true` si cette complétion provient d'une opération **multishot** (Temps
1440 /// 3d) : le slot a été réservé multishot (lu avant `complete`). Sert à
1441 /// [`Completion::multishot_token`].
1442 multishot: bool,
1443}
1444
1445impl Completion {
1446 /// Jeton de l'opération qui a produit cette complétion.
1447 #[must_use]
1448 pub fn token(&self) -> SubmissionToken {
1449 self.token
1450 }
1451
1452 /// Jeton **multishot** si cette complétion provient d'une opération multishot
1453 /// (Temps 3d), `None` sinon. Toutes les complétions d'un même multishot
1454 /// (intermédiaires `CQE_F_MORE` **et** terminale) portent le **même** jeton ;
1455 /// [`Completion::has_more`] distingue les intermédiaires de la terminale.
1456 #[must_use]
1457 pub fn multishot_token(&self) -> Option<MultishotToken> {
1458 if self.multishot {
1459 Some(MultishotToken::new(
1460 self.token.slot(),
1461 self.token.generation(),
1462 ))
1463 } else {
1464 None
1465 }
1466 }
1467
1468 /// Résultat brut du kernel (sémantique dépendante de l'op soumise).
1469 #[must_use]
1470 pub fn raw_result(&self) -> i32 {
1471 self.res
1472 }
1473
1474 /// Flags de complétion (cf. [`CompletionFlags`]).
1475 #[must_use]
1476 pub fn flags(&self) -> CompletionFlags {
1477 CompletionFlags::from_bits_truncate(self.raw_flags)
1478 }
1479
1480 /// `IORING_CQE_F_MORE` : d'autres complétions suivront (multishot).
1481 #[must_use]
1482 pub fn has_more(&self) -> bool {
1483 self.flags().contains(CompletionFlags::MORE)
1484 }
1485
1486 /// `IORING_CQE_F_NOTIF` : complétion de notification zero-copy.
1487 #[must_use]
1488 pub fn is_notif(&self) -> bool {
1489 self.flags().contains(CompletionFlags::NOTIF)
1490 }
1491
1492 /// ID du buffer fourni consommé (`IORING_CQE_F_BUFFER`), le cas échéant.
1493 /// L'identifiant occupe les 16 bits de poids fort des flags du CQE.
1494 #[must_use]
1495 pub fn buffer_id(&self) -> Option<u16> {
1496 if self.flags().contains(CompletionFlags::BUFFER) {
1497 Some(u16::try_from(self.raw_flags >> 16).unwrap_or(0))
1498 } else {
1499 None
1500 }
1501 }
1502
1503 /// `IORING_CQE_F_SOCK_NONEMPTY` : données restantes lisibles après un `recv`.
1504 #[must_use]
1505 pub fn socket_has_pending_data(&self) -> bool {
1506 self.flags().contains(CompletionFlags::SOCK_NONEMPTY)
1507 }
1508
1509 /// Interprétation générique : `res ≥ 0` → valeur utile, `res < 0` →
1510 /// `Err(Errno)` (= `-res`).
1511 ///
1512 /// # Errors
1513 ///
1514 /// [`Errno`] correspondant à `-res` si le kernel a signalé un échec.
1515 pub fn into_result(self) -> Result<i32, Errno> {
1516 res_to_result(self.res)
1517 }
1518
1519 /// Récupère le buffer **déplacé** hors du slot (S1, zéro copie) avec le
1520 /// nombre d'octets traités (`submit_read`/`submit_write`).
1521 ///
1522 /// # Errors
1523 ///
1524 /// [`Errno`] correspondant à `-res` si l'opération a échoué (le buffer est
1525 /// alors consommé sans être restitué).
1526 pub fn into_buffer_result(self) -> Result<(Vec<u8>, usize), Errno> {
1527 let bytes = res_to_result(self.res)?;
1528 let length = usize::try_from(bytes).unwrap_or(0);
1529 let buffer = match self.payload {
1530 Some(OwnedOp::Bytes(buffer)) => buffer,
1531 // Op sans buffer `Bytes` (NOP, ou mauvais accesseur) : buffer vide.
1532 _ => Vec::new(),
1533 };
1534 Ok((buffer, length))
1535 }
1536
1537 /// Récupère les buffers vectorisés **déplacés** hors du slot avec le nombre
1538 /// total d'octets traités (`submit_readv`/`submit_writev`).
1539 ///
1540 /// # Errors
1541 ///
1542 /// [`Errno`] correspondant à `-res` si l'opération a échoué (les buffers
1543 /// sont alors consommés sans être restitués).
1544 pub fn into_vectored_result(self) -> Result<(Vec<Vec<u8>>, usize), Errno> {
1545 let bytes = res_to_result(self.res)?;
1546 let length = usize::try_from(bytes).unwrap_or(0);
1547 let buffers = match self.payload {
1548 Some(OwnedOp::Vectored { buffers, .. }) => buffers,
1549 _ => Vec::new(),
1550 };
1551 Ok((buffers, length))
1552 }
1553
1554 /// Récupère le tampon `statx` **déplacé** hors du slot, rempli par le kernel
1555 /// (`submit_statx`).
1556 ///
1557 /// # Errors
1558 ///
1559 /// [`Errno`] correspondant à `-res` si l'opération a échoué (le tampon est
1560 /// alors consommé sans être restitué).
1561 pub fn into_statx(self) -> Result<Box<air_sys_types::fs::Statx>, Errno> {
1562 res_to_result(self.res)?;
1563 match self.payload {
1564 Some(OwnedOp::Statx { out, .. }) => Ok(out),
1565 // Accesseur appelé sur une complétion sans tampon statx : tampon
1566 // zéro-initialisé (défensif ; jamais sur le chemin nominal `statx`).
1567 _ => Ok(Box::new(air_sys_types::fs::Statx::default())),
1568 }
1569 }
1570
1571 /// Récupère le buffer de valeur xattr **déplacé** hors du slot avec la
1572 /// taille rendue par le kernel (`submit_getxattr`/`submit_fgetxattr`). La
1573 /// taille permet de retailler (un buffer trop court ⇒ `Err(ERANGE)`).
1574 ///
1575 /// # Errors
1576 ///
1577 /// [`Errno`] correspondant à `-res` si l'opération a échoué (le buffer est
1578 /// alors consommé sans être restitué).
1579 pub fn into_xattr_result(self) -> Result<(Vec<u8>, usize), Errno> {
1580 let bytes = res_to_result(self.res)?;
1581 let length = usize::try_from(bytes).unwrap_or(0);
1582 let value = match self.payload {
1583 Some(OwnedOp::Xattr { value, .. }) => value,
1584 _ => Vec::new(),
1585 };
1586 Ok((value, length))
1587 }
1588
1589 /// Récupère le descripteur ouvert par `submit_openat2`.
1590 ///
1591 /// `res` (≥ 0) est le numéro de FD ; il devient un [`OwnedFd`] possédé
1592 /// (CLOEXEC selon `how`). Le chemin source, gardé en vie dans le slot, est
1593 /// libéré ici.
1594 ///
1595 /// # Errors
1596 ///
1597 /// [`Errno`] correspondant à `-res` si l'ouverture a échoué.
1598 pub fn opened_fd(self) -> Result<OwnedFd, Errno> {
1599 let fd = res_to_result(self.res)?;
1600 // SAFETY: `res ≥ 0` est un FD frais possédé, retourné par `openat2` via
1601 // io_uring ; `from_raw_fd` en prend l'ownership exclusif (pas de double
1602 // close : le chemin parké est libéré, aucun autre détenteur).
1603 Ok(unsafe { OwnedFd::from_raw_fd(fd) })
1604 }
1605
1606 /// Récupère le socket créé par `submit_socket` (CLOEXEC).
1607 ///
1608 /// # Errors
1609 ///
1610 /// [`Errno`] correspondant à `-res` si la création a échoué.
1611 pub fn into_socket_fd(self) -> Result<OwnedFd, Errno> {
1612 let fd = res_to_result(self.res)?;
1613 // SAFETY: `res ≥ 0` est un FD socket frais possédé (`IORING_OP_SOCKET`) ;
1614 // `from_raw_fd` en prend l'ownership exclusif.
1615 Ok(unsafe { OwnedFd::from_raw_fd(fd) })
1616 }
1617
1618 /// Récupère le socket accepté par `submit_accept` (CLOEXEC), sans adresse
1619 /// pair.
1620 ///
1621 /// # Errors
1622 ///
1623 /// [`Errno`] correspondant à `-res` si l'`accept` a échoué.
1624 pub fn accepted_fd(self) -> Result<OwnedFd, Errno> {
1625 let fd = res_to_result(self.res)?;
1626 // SAFETY: `res ≥ 0` est un FD socket connecté frais possédé
1627 // (`IORING_OP_ACCEPT`) ; ownership exclusif via `from_raw_fd`.
1628 Ok(unsafe { OwnedFd::from_raw_fd(fd) })
1629 }
1630
1631 /// Récupère le socket accepté **et** l'adresse du pair
1632 /// (`submit_accept_with_peer`). L'adresse est décodée du stockage `sockaddr`
1633 /// possédé écrit par le kernel.
1634 ///
1635 /// # Errors
1636 ///
1637 /// [`Errno`] correspondant à `-res` si l'`accept` a échoué.
1638 pub fn into_accept_result(self) -> Result<(OwnedFd, air_sys_types::net::SocketAddr), Errno> {
1639 let fd = res_to_result(self.res)?;
1640 let address = match self.payload {
1641 Some(OwnedOp::Accept { addr, addrlen }) => {
1642 crate::net::raw_to_socket_addr(&addr.bytes, *addrlen)
1643 // Décodage défensif : un kernel renvoyant une famille inconnue ou une
1644 // longueur nulle ⇒ adresse « non nommée » plutôt que panique.
1645 .unwrap_or(air_sys_types::net::SocketAddr::Unix(
1646 air_sys_types::net::UnixSocketAddr::Unnamed,
1647 ))
1648 }
1649 _ => air_sys_types::net::SocketAddr::Unix(air_sys_types::net::UnixSocketAddr::Unnamed),
1650 };
1651 // SAFETY: `res ≥ 0` est un FD socket connecté frais possédé ; ownership
1652 // exclusif via `from_raw_fd`.
1653 Ok((unsafe { OwnedFd::from_raw_fd(fd) }, address))
1654 }
1655
1656 /// Décode le résultat d'un `submit_receive_message` : buffers reçus +
1657 /// métadonnées (`ReceiveMessageMeta`, adossé à `io_uring_recvmsg_out`). Les
1658 /// **FD reçus** via `SCM_RIGHTS` sont matérialisés en [`OwnedFd`] **CLOEXEC**
1659 /// (drapeau `MSG_CMSG_CLOEXEC` posé par la façade). Un cmsg tronqué
1660 /// (`MSG_CTRUNC`) est signalé via `meta.control_truncated()` ; les FD
1661 /// incomplets ne sont **jamais** matérialisés (le kernel les ferme) ⇒ aucune
1662 /// fuite de FD.
1663 ///
1664 /// # Errors
1665 ///
1666 /// [`Errno`] correspondant à `-res` si le `recvmsg` a échoué.
1667 pub fn into_receive_message_result(
1668 self,
1669 ) -> Result<
1670 (
1671 air_sys_types::net::OwnedReceiveMessage,
1672 air_sys_types::net::ReceiveMessageMeta,
1673 ),
1674 Errno,
1675 > {
1676 use air_sys_types::net::{MessageFlags, OwnedReceiveMessage, ReceiveMessageMeta};
1677 let bytes = res_to_result(self.res)?;
1678 let payloadlen = u32::try_from(bytes).unwrap_or(0);
1679 let Some(OwnedOp::RecvMsg(state)) = self.payload else {
1680 // Accesseur appelé sur une complétion sans état recvmsg : résultat
1681 // vide (défensif ; jamais sur le chemin nominal `recvmsg`).
1682 return Ok((
1683 OwnedReceiveMessage::default(),
1684 ReceiveMessageMeta {
1685 namelen: 0,
1686 controllen: 0,
1687 payloadlen,
1688 flags: MessageFlags::empty(),
1689 address: None,
1690 fds: Vec::new(),
1691 },
1692 ));
1693 };
1694 let owned::RecvMsgState {
1695 msghdr,
1696 buffers,
1697 name,
1698 control,
1699 ..
1700 } = *state;
1701 let namelen = msghdr.msg_namelen;
1702 let controllen_usize = usize::try_from(msghdr.msg_controllen).unwrap_or(0);
1703 let controllen = u32::try_from(msghdr.msg_controllen).unwrap_or(0);
1704 let flags = MessageFlags::from_bits_truncate(msghdr.msg_flags);
1705 // Extrait les FD complets du buffer de contrôle (jamais de FD partiel).
1706 let fds = crate::net::parse_scm_rights(&control, controllen_usize.min(control.len()));
1707 let address = if namelen > 0 {
1708 crate::net::raw_to_socket_addr(&name.bytes, namelen)
1709 } else {
1710 None
1711 };
1712 let control_capacity = control.len();
1713 let meta = ReceiveMessageMeta {
1714 namelen,
1715 controllen,
1716 payloadlen,
1717 flags,
1718 address,
1719 fds,
1720 };
1721 let message = OwnedReceiveMessage {
1722 buffers,
1723 control_capacity,
1724 flags,
1725 };
1726 Ok((message, meta))
1727 }
1728
1729 /// Restitue le buffer zero-copy de `submit_send_zero_copy` à la complétion
1730 /// **NOTIF** (`is_notif`) — le kernel ne référence plus la mémoire. En cas
1731 /// d'**échec précoce** (CQE de résultat `res < 0` **sans** `F_MORE`, donc
1732 /// aucune NOTIF à venir), renvoie l'erreur ; le buffer est alors libéré (sûr,
1733 /// aucune fuite).
1734 ///
1735 /// Pour un `submit_send_message_zero_copy` **multi-buffers**, utiliser
1736 /// [`Completion::into_zero_copy_buffers`] (restitution intégrale, ADR-032) ;
1737 /// appelé ici, ses buffers sont **concaténés** sans perte d'octet (les
1738 /// frontières sont perdues, mais aucune donnée n'est *discardée*).
1739 ///
1740 /// # Errors
1741 ///
1742 /// [`Errno`] correspondant à `-res` sur le chemin d'échec précoce.
1743 pub fn into_zero_copy_buffer(self) -> Result<Vec<u8>, Errno> {
1744 let is_notif = self.flags().contains(CompletionFlags::NOTIF);
1745 let res = self.res;
1746 let buffer = match self.payload {
1747 Some(OwnedOp::Bytes(buffer)) => buffer,
1748 // ADR-032 : ne *discarde* jamais — concatène tous les buffers (aucun
1749 // octet perdu) plutôt que de ne rendre que le premier.
1750 Some(OwnedOp::SendMsg(state)) => state.buffers.into_iter().flatten().collect(),
1751 _ => Vec::new(),
1752 };
1753 if is_notif {
1754 // NOTIF : succès ; le `res` peut porter le bit « copié » (REPORT_USAGE),
1755 // ce n'est pas une erreur — la restitution est inconditionnelle.
1756 Ok(buffer)
1757 } else {
1758 // Échec précoce (résultat sans F_MORE) : propage l'erreur.
1759 res_to_result(res).map(|_| buffer)
1760 }
1761 }
1762
1763 /// Restitue **l'intégralité** des buffers zero-copy d'un
1764 /// `submit_send_message_zero_copy` à la complétion **NOTIF** — **tous** les
1765 /// buffers du message, **intacts et dans l'ordre** (ADR-032, zéro perte de
1766 /// donnée). Le slot S1 les a tous retenus jusqu'au NOTIF. Pour un
1767 /// `submit_send_zero_copy` (mono-buffer), renvoie un `Vec` à un seul élément.
1768 ///
1769 /// Échec précoce (CQE résultat `res < 0` sans `F_MORE`) : renvoie l'erreur ;
1770 /// les buffers sont alors libérés (sûr, aucune fuite).
1771 ///
1772 /// # Errors
1773 ///
1774 /// [`Errno`] correspondant à `-res` sur le chemin d'échec précoce.
1775 pub fn into_zero_copy_buffers(self) -> Result<Vec<Vec<u8>>, Errno> {
1776 let is_notif = self.flags().contains(CompletionFlags::NOTIF);
1777 let res = self.res;
1778 let buffers = match self.payload {
1779 // `sendmsg_zc` : tous les buffers, dans l'ordre — jamais un sous-ensemble.
1780 Some(OwnedOp::SendMsg(state)) => state.buffers,
1781 // `send_zc` mono-buffer : enveloppé dans un `Vec` à un élément.
1782 Some(OwnedOp::Bytes(buffer)) => vec![buffer],
1783 _ => Vec::new(),
1784 };
1785 if is_notif {
1786 Ok(buffers)
1787 } else {
1788 res_to_result(res).map(|_| buffers)
1789 }
1790 }
1791
1792 /// `IORING_NOTIF_USAGE_ZC_COPIED` : sur une complétion NOTIF avec
1793 /// `ZeroCopyFlags::REPORT_USAGE`, indique que le kernel a dû **copier** (pas
1794 /// de vrai zero-copy). `false` hors NOTIF ou sans `REPORT_USAGE`.
1795 #[must_use]
1796 pub fn zero_copy_copied(&self) -> bool {
1797 self.flags().contains(CompletionFlags::NOTIF)
1798 && (self.res & raw::IORING_NOTIF_USAGE_ZC_COPIED) != 0
1799 }
1800
1801 /// Décode les événements prêts d'un `submit_poll_add` (`res ≥ 0` = bitset
1802 /// `PollEvents`).
1803 ///
1804 /// # Errors
1805 ///
1806 /// [`Errno`] correspondant à `-res` si le poll a échoué.
1807 pub fn into_poll_result(self) -> Result<air_sys_types::io_uring::PollEvents, Errno> {
1808 let bits = res_to_result(self.res)?;
1809 // `res` ≥ 0 est un bitset d'événements (`u32`) — jamais négatif ici.
1810 let bits = u32::try_from(bits).unwrap_or(0);
1811 Ok(air_sys_types::io_uring::PollEvents::from_bits_truncate(
1812 bits,
1813 ))
1814 }
1815
1816 /// Récupère le tampon `siginfo_t` **déplacé** hors du slot, rempli par le
1817 /// kernel (`submit_waitid`). À lire via
1818 /// [`SignalInfo::as_bytes`](air_sys_types::signal::SignalInfo::as_bytes).
1819 ///
1820 /// # Errors
1821 ///
1822 /// [`Errno`] correspondant à `-res` si le `waitid` a échoué (le tampon est
1823 /// alors consommé sans être restitué).
1824 pub fn into_waitid_result(self) -> Result<Box<air_sys_types::signal::SignalInfo>, Errno> {
1825 res_to_result(self.res)?;
1826 match self.payload {
1827 Some(OwnedOp::Waitid(info)) => Ok(info),
1828 // Accesseur appelé hors chemin `waitid` : `siginfo` zéro-initialisé.
1829 _ => Ok(Box::new(air_sys_types::signal::SignalInfo::zeroed())),
1830 }
1831 }
1832
1833 /// Construit une [`Completion`] de test (sans kernel) pour exercer les
1834 /// accesseurs typés sur des `(res, flags, payload)` arbitraires.
1835 #[cfg(test)]
1836 pub(crate) fn for_test(
1837 token: SubmissionToken,
1838 res: i32,
1839 raw_flags: u32,
1840 payload: Option<OwnedOp>,
1841 ) -> Self {
1842 Self {
1843 token,
1844 res,
1845 raw_flags,
1846 payload,
1847 multishot: false,
1848 }
1849 }
1850
1851 /// Succès sans valeur de retour (close, fsync, nop…).
1852 ///
1853 /// # Errors
1854 ///
1855 /// [`Errno`] correspondant à `-res` si l'opération a échoué.
1856 pub fn completed(&self) -> Result<(), Errno> {
1857 res_to_result(self.res).map(|_| ())
1858 }
1859}
1860
1861/// Convertit le `res` d'un CQE en `Result` : `res < 0` ⇒ `Err(-res)`, sinon
1862/// `Ok(res)`. Convention io_uring (ADR-022 Décision 9).
1863fn res_to_result(res: i32) -> Result<i32, Errno> {
1864 if res < 0 {
1865 // `res` négatif ⇒ `-res` est un errno kernel (1..=4095) non nul.
1866 // `checked_neg` ne renvoie `None` que pour `i32::MIN` (défensif, jamais
1867 // en pratique pour un CQE kernel).
1868 match res.checked_neg().and_then(core::num::NonZeroI32::new) {
1869 Some(nz) => Err(Errno::from_nonzero(nz)),
1870 None => Err(Errno::EINVAL),
1871 }
1872 } else {
1873 Ok(res)
1874 }
1875}
1876
1877/// Surface de fuzzing (frontière de décode des données écrites par le kernel,
1878/// Principe 3). N'existe QUE sous `--cfg fuzzing` (posé par `cargo-fuzz`) :
1879/// zéro surface en build normal. Le harnais est **pur** (aucun syscall, aucun
1880/// ring réel) — on modélise un kernel buggé/hostile remplissant la mémoire
1881/// partagée. **Invariant : TOTALITÉ** — toute entrée donne une valeur typée ou
1882/// une erreur typée, jamais de panic/UB/accès hors-bornes.
1883#[cfg(fuzzing)]
1884pub mod fuzz_api {
1885 use super::owned::OwnedOp;
1886 use super::{Completion, SetupFlags, SubmissionToken, raw, ring_sizes};
1887 use crate::io_uring::slab::InflightSlab;
1888 use air_sys_types::fs::{Statx, StatxMask};
1889 // Crate `#![no_std]` (sauf `cfg(test)`) : sous `--cfg fuzzing` on n'est pas en
1890 // test, donc les types possédés viennent d'`alloc`, pas du prélude `std`. La
1891 // macro `vec!` reste fournie par `#[macro_use] extern crate alloc` (crate root).
1892 use alloc::boxed::Box;
1893 use alloc::ffi::CString;
1894 use alloc::vec::Vec;
1895 use core::num::NonZeroU32;
1896
1897 /// Décode les **résultats kernel externes** des opérations 2a (`statx`,
1898 /// `getxattr`, `readv`) via les accesseurs typés de [`Completion`], pour
1899 /// tout `(res, raw_flags)` et tout contenu de tampon. Invariant de
1900 /// **totalité** : aucune entrée ne produit panic/UB — la valeur typée ou
1901 /// l'erreur typée, jamais autre chose. (Le kernel écrit `Statx`/la valeur
1902 /// xattr ; ce sont des données externes au sens du Principe 3.)
1903 pub fn decode_2a_results(res: i32, raw_flags: u32, mask: u32, value: Vec<u8>) {
1904 let token = SubmissionToken::new(0, 0);
1905
1906 // statx : tampon écrit par le kernel, masque arbitraire.
1907 let mut out = Box::new(Statx::default());
1908 out.mask = mask;
1909 out.size = u64::from(raw_flags);
1910 let statx_completion = Completion {
1911 token,
1912 res,
1913 raw_flags,
1914 payload: Some(OwnedOp::Statx {
1915 out,
1916 path: CString::new("p").expect("pas de NUL"),
1917 }),
1918 multishot: false,
1919 };
1920 if let Ok(decoded) = statx_completion.into_statx() {
1921 let _ = decoded.has(StatxMask::from_bits_truncate(mask));
1922 let _ = decoded.size;
1923 let _ = decoded.mtime;
1924 }
1925
1926 // getxattr : valeur de sortie de taille arbitraire.
1927 let xattr_completion = Completion {
1928 token,
1929 res,
1930 raw_flags,
1931 payload: Some(OwnedOp::Xattr {
1932 value: value.clone(),
1933 name: CString::new("n").expect("pas de NUL"),
1934 path: None,
1935 }),
1936 multishot: false,
1937 };
1938 let _ = xattr_completion.into_xattr_result();
1939
1940 // readv : restitution vectorisée (un buffer + son iovec dérivé).
1941 let mut buffers = vec![value];
1942 let iovecs: Box<[raw::Iovec]> = buffers
1943 .iter_mut()
1944 .map(|buffer| raw::Iovec {
1945 iov_base: buffer.as_mut_ptr(),
1946 iov_len: buffer.len(),
1947 })
1948 .collect();
1949 let vectored_completion = Completion {
1950 token,
1951 res,
1952 raw_flags,
1953 payload: Some(OwnedOp::Vectored { buffers, iovecs }),
1954 multishot: false,
1955 };
1956 let _ = vectored_completion.into_vectored_result();
1957 }
1958
1959 /// Décode les **complétions `URING_CMD` socket** (Temps 2d) : `getsockopt`
1960 /// restitue son **buffer de sortie** (rempli par le kernel) via
1961 /// `into_buffer_result` (avec `res` = longueur effective annoncée par le
1962 /// kernel — donnée externe, Principe 3) ; `inq`/`outq` via `into_result`.
1963 /// Invariant de **totalité** : pour tout `(res, value)`, ni panic ni UB —
1964 /// la longueur retournée peut dépasser `value.len()` si le kernel ment, mais
1965 /// l'accesseur ne déréférence pas au-delà du buffer (il rend `(Vec, usize)`).
1966 pub fn decode_2d_results(res: i32, value: Vec<u8>) {
1967 let token = SubmissionToken::new(0, 0);
1968 // getsockopt : buffer de sortie + longueur effective (`res`).
1969 let getsockopt_completion = Completion {
1970 token,
1971 res,
1972 raw_flags: 0,
1973 payload: Some(OwnedOp::Bytes(value)),
1974 multishot: false,
1975 };
1976 let _ = getsockopt_completion.into_buffer_result();
1977 // inq/outq : pas de buffer, `res` = compteur d'octets.
1978 let count_completion = Completion {
1979 token,
1980 res,
1981 raw_flags: 0,
1982 payload: None,
1983 multishot: false,
1984 };
1985 let _ = count_completion.completed();
1986 let _ = count_completion.into_result();
1987 }
1988
1989 /// Décode les **données réseau externes** d'un `recvmsg` io_uring (Temps 2b)
1990 /// via `into_receive_message_result` : le **buffer de contrôle** (cmsg
1991 /// `SCM_RIGHTS`) et le **stockage d'adresse** (`sockaddr`) sont remplis par
1992 /// le kernel ⇒ entrées externes (Principe 3). Invariant de **totalité** :
1993 /// pour tout `(res, controllen, namelen, flags, control, name)`, le parsing
1994 /// (`parse_scm_rights` + `raw_to_socket_addr`) ne panique ni ne déborde. Les
1995 /// `i32` extraits du cmsg ne sont **pas** matérialisés en FD ici (on ne
1996 /// fournit pas de vrai descripteur), évitant toute fermeture parasite.
1997 pub fn decode_2b_recvmsg(
1998 res: i32,
1999 controllen: u32,
2000 namelen: u32,
2001 msg_flags: i32,
2002 control: Vec<u8>,
2003 name: Vec<u8>,
2004 ) {
2005 let mut name_storage = Box::new(crate::net::RawSockaddrStorage {
2006 bytes: [0u8; crate::net::MAX_SOCKADDR_LEN],
2007 });
2008 let take = name.len().min(crate::net::MAX_SOCKADDR_LEN);
2009 name_storage.bytes[..take].copy_from_slice(&name[..take]);
2010 let state = crate::io_uring::owned::RecvMsgState {
2011 msghdr: Box::new(raw::Msghdr {
2012 msg_namelen: namelen,
2013 msg_controllen: u64::from(controllen),
2014 msg_flags,
2015 ..raw::Msghdr::default()
2016 }),
2017 iovecs: Vec::new().into_boxed_slice(),
2018 buffers: Vec::new(),
2019 name: name_storage,
2020 control: control.into_boxed_slice(),
2021 };
2022 let completion = Completion {
2023 token: SubmissionToken::new(0, 0),
2024 res,
2025 raw_flags: 0,
2026 payload: Some(OwnedOp::RecvMsg(Box::new(state))),
2027 multishot: false,
2028 };
2029 if let Ok((_msg, meta)) = completion.into_receive_message_result() {
2030 let _ = meta.payload_truncated();
2031 let _ = meta.control_truncated();
2032 // Ne pas `drop` les FD ici : ce sont des `i32` arbitraires enveloppés
2033 // en `OwnedFd` ; on les **oublie** pour éviter de fermer des FD réels
2034 // du processus de fuzz (totalité sans effet de bord).
2035 for fd in meta.fds {
2036 core::mem::forget(fd);
2037 }
2038 }
2039 }
2040
2041 /// Décode les **résultats kernel externes** des accesseurs 2c (`poll`,
2042 /// `waitid`) pour tout `res`. Invariant de **totalité** : `into_poll_result`
2043 /// (bitset d'événements écrit par le kernel) et `into_waitid_result` (tampon
2044 /// `siginfo` rempli par le kernel) ne paniquent ni ne débordent. Le tampon
2045 /// `siginfo` est de la **mémoire pure** (aucun FD) ⇒ droppé normalement.
2046 pub fn decode_2c_results(res: i32) {
2047 let token = SubmissionToken::new(0, 0);
2048 let poll = Completion {
2049 token,
2050 res,
2051 raw_flags: 0,
2052 payload: None,
2053 multishot: false,
2054 };
2055 let _ = poll.into_poll_result();
2056 let waitid = Completion {
2057 token,
2058 res,
2059 raw_flags: 0,
2060 payload: Some(OwnedOp::Waitid(Box::new(
2061 air_sys_types::signal::SignalInfo::zeroed(),
2062 ))),
2063 multishot: false,
2064 };
2065 if let Ok(info) = waitid.into_waitid_result() {
2066 let _ = info.as_bytes();
2067 }
2068 }
2069
2070 /// Décode un CQE (champs bruts kernel) et exerce toutes les interprétations
2071 /// de [`Completion`] : aucune ne doit paniquer pour quelque entrée que ce
2072 /// soit (`res` incluant `i32::MIN`, flags arbitraires).
2073 pub fn decode_cqe(user_data: u64, res: i32, raw_flags: u32) {
2074 let completion = Completion {
2075 token: SubmissionToken::from_user_data(user_data),
2076 res,
2077 raw_flags,
2078 payload: None,
2079 multishot: false,
2080 };
2081 let _ = completion.token();
2082 let _ = completion.raw_result();
2083 let _ = completion.flags();
2084 let _ = completion.has_more();
2085 let _ = completion.is_notif();
2086 let _ = completion.buffer_id();
2087 let _ = completion.socket_has_pending_data();
2088 let _ = completion.completed();
2089 let _ = completion.into_result();
2090 }
2091
2092 /// Décode une **complétion brute** (Temps 4) à partir de **données kernel
2093 /// externes** arbitraires (`user_data`/`res`/`flags`, Principe 3). Invariant de
2094 /// **totalité** : aucune entrée ne panique ; le routage du tag (§5) est total
2095 /// (brut ⇔ bit 63 posé) ; un `user_data` **non tagué** hostile passé au
2096 /// décodage géré (`SubmissionToken::from_user_data`) reste borné (jamais d'OOB,
2097 /// jamais de panic). Couvre la frontière de décode brut/géré coexistant.
2098 pub fn decode_raw_cqe(user_data: u64, res: i32, flags: u32) {
2099 let raw_cqe = raw::RawCompletionQueueEntry {
2100 user_data,
2101 res,
2102 flags,
2103 };
2104 // Lecture totale des champs (POD) — jamais de panic.
2105 let _ = raw_cqe.user_data;
2106 let _ = raw_cqe.res;
2107 let _ = raw_cqe.flags;
2108 // Routage du tag : exactement deux classes, exhaustives.
2109 if raw_cqe.user_data & raw::RAW_USER_DATA_TAG == 0 {
2110 // Classe « gérée » : le décodage slab doit rester borné même sur un
2111 // `user_data` hostile non tagué (S1).
2112 let _ = SubmissionToken::from_user_data(raw_cqe.user_data);
2113 }
2114 }
2115
2116 /// Invariant de soundness S1 : un `user_data` **hostile** (slot arbitraire,
2117 /// potentiellement ≥ capacité) passé à `slab.complete` doit être **borné**
2118 /// (rejeté en périmé via `get`), jamais indexer le slab hors-bornes.
2119 pub fn slab_complete_is_bounded(capacity: u8, hostile_user_data: &[u64]) {
2120 let cap = NonZeroU32::new(u32::from(capacity).max(1)).expect("≥ 1");
2121 let mut slab = InflightSlab::with_capacity(cap);
2122 // Quelques réservations réelles pour avoir des slots en vol.
2123 let mut tokens = Vec::new();
2124 for _ in 0..u32::from(capacity).min(4) {
2125 if let Ok(t) = slab.reserve(None) {
2126 tokens.push(t);
2127 }
2128 }
2129 // Complétions hostiles entrelacées : aucune ne doit paniquer/OOB.
2130 for &user_data in hostile_user_data {
2131 let _ = slab.complete(user_data, user_data & 1 == 0);
2132 }
2133 for token in tokens {
2134 let _ = slab.complete(token.to_user_data(), false);
2135 }
2136 }
2137
2138 /// Décode les tailles d'anneau depuis des `io_uring_params` bruts : toute
2139 /// combinaison d'entrées/offsets/features ⇒ `Ok`/`Err(EINVAL)`, jamais
2140 /// d'overflow ni de panic (arithmétique checked).
2141 pub fn ring_sizes_decode(
2142 sq_entries: u32,
2143 cq_entries: u32,
2144 sq_array_off: u32,
2145 cq_cqes_off: u32,
2146 features: u32,
2147 setup_bits: u32,
2148 ) {
2149 let mut params = raw::IoUringParams {
2150 sq_entries,
2151 cq_entries,
2152 features,
2153 ..raw::IoUringParams::default()
2154 };
2155 params.sq_off.array = sq_array_off;
2156 params.cq_off.cqes = cq_cqes_off;
2157 let _ = ring_sizes(¶ms, SetupFlags::from_bits_truncate(setup_bits));
2158 }
2159}
2160
2161/// Itérateur sur les complétions disponibles d'un [`IoUring`] (emprunte le ring).
2162pub struct CompletionIter<'ring> {
2163 ring: &'ring mut IoUring,
2164}
2165
2166impl Iterator for CompletionIter<'_> {
2167 type Item = Completion;
2168
2169 fn next(&mut self) -> Option<Completion> {
2170 self.ring.harvest_ready()
2171 }
2172}
2173
2174// ---------------------------------------------------------------------------
2175// Capacités, annulation, sandbox (types couplés au wrapper)
2176// ---------------------------------------------------------------------------
2177
2178/// Miroir typé de `struct io_uring_params` (entrées/sorties de
2179/// `io_uring_setup(2)`).
2180///
2181/// Porte les setup flags demandés, les profondeurs effectives
2182/// (`sq_entries`/`cq_entries`, arrondies par le kernel), les `features`
2183/// négociées et les offsets d'anneau. Le layout `#[repr(C)]` réel est posé en
2184/// Pass B (les champs ABI conservent leurs noms kernel, ADR-029 zone
2185/// d'interface). Opaque au Temps 1.
2186#[derive(Debug, Clone, Copy)]
2187pub struct IoUringParams {
2188 // Layout #[repr(C)] complet posé en Pass B.
2189 _opaque: [u8; 0],
2190}
2191
2192/// Features io_uring du kernel courant (axe G, 16 features en 6.12). Lues dans
2193/// `io_uring_params.features` après setup, exposées par **prédicats stables**.
2194#[derive(Debug, Clone, Copy)]
2195pub struct IoUringCapabilities {
2196 // Champ interne (Pass B) : bitmask IORING_FEAT_*.
2197 features: u32,
2198}
2199
2200impl IoUringCapabilities {
2201 /// Vrai si le bit de feature `bit` (un `IORING_FEAT_*`) est négocié.
2202 fn has(&self, bit: u32) -> bool {
2203 self.features & bit != 0
2204 }
2205
2206 /// `IORING_FEAT_SINGLE_MMAP` : SQ et CQ partagent une seule mmap.
2207 #[must_use]
2208 pub fn single_mmap(&self) -> bool {
2209 self.has(raw::IORING_FEAT_SINGLE_MMAP)
2210 }
2211 /// `IORING_FEAT_NODROP` : la CQ ne perd pas de complétions (overflow retenu).
2212 #[must_use]
2213 pub fn nodrop(&self) -> bool {
2214 self.has(raw::IORING_FEAT_NODROP)
2215 }
2216 /// `IORING_FEAT_SUBMIT_STABLE` : données de soumission stables après `enter`.
2217 #[must_use]
2218 pub fn submit_stable(&self) -> bool {
2219 self.has(raw::IORING_FEAT_SUBMIT_STABLE)
2220 }
2221 /// `IORING_FEAT_RW_CUR_POS` : offset `-1` = position courante du fichier.
2222 #[must_use]
2223 pub fn rw_cur_pos(&self) -> bool {
2224 self.has(raw::IORING_FEAT_RW_CUR_POS)
2225 }
2226 /// `IORING_FEAT_CUR_PERSONALITY` : applique la personnalité courante.
2227 #[must_use]
2228 pub fn cur_personality(&self) -> bool {
2229 self.has(raw::IORING_FEAT_CUR_PERSONALITY)
2230 }
2231 /// `IORING_FEAT_FAST_POLL` : poll interne rapide.
2232 #[must_use]
2233 pub fn fast_poll(&self) -> bool {
2234 self.has(raw::IORING_FEAT_FAST_POLL)
2235 }
2236 /// `IORING_FEAT_POLL_32BITS` : masques de poll 32 bits.
2237 #[must_use]
2238 pub fn poll_32bits(&self) -> bool {
2239 self.has(raw::IORING_FEAT_POLL_32BITS)
2240 }
2241 /// `IORING_FEAT_SQPOLL_NONFIXED` : SQPOLL sans FD fixes obligatoires.
2242 #[must_use]
2243 pub fn sqpoll_nonfixed(&self) -> bool {
2244 self.has(raw::IORING_FEAT_SQPOLL_NONFIXED)
2245 }
2246 /// `IORING_FEAT_EXT_ARG` : argument étendu d'`enter` (timeout d'attente).
2247 #[must_use]
2248 pub fn ext_arg(&self) -> bool {
2249 self.has(raw::IORING_FEAT_EXT_ARG)
2250 }
2251 /// `IORING_FEAT_NATIVE_WORKERS` : workers io-wq natifs.
2252 #[must_use]
2253 pub fn native_workers(&self) -> bool {
2254 self.has(raw::IORING_FEAT_NATIVE_WORKERS)
2255 }
2256 /// `IORING_FEAT_RSRC_TAGS` : étiquetage des ressources enregistrées.
2257 #[must_use]
2258 pub fn rsrc_tags(&self) -> bool {
2259 self.has(raw::IORING_FEAT_RSRC_TAGS)
2260 }
2261 /// `IORING_FEAT_CQE_SKIP` : support de `IOSQE_CQE_SKIP_SUCCESS`.
2262 #[must_use]
2263 pub fn cqe_skip(&self) -> bool {
2264 self.has(raw::IORING_FEAT_CQE_SKIP)
2265 }
2266 /// `IORING_FEAT_LINKED_FILE` : résolution de FD dans les chaînes liées.
2267 #[must_use]
2268 pub fn linked_file(&self) -> bool {
2269 self.has(raw::IORING_FEAT_LINKED_FILE)
2270 }
2271 /// `IORING_FEAT_REG_REG_RING` : enregistrement du ring fd via ring enregistré.
2272 #[must_use]
2273 pub fn reg_reg_ring(&self) -> bool {
2274 self.has(raw::IORING_FEAT_REG_REG_RING)
2275 }
2276 /// `IORING_FEAT_RECVSEND_BUNDLE` : recv/send groupés (bundle).
2277 #[must_use]
2278 pub fn recvsend_bundle(&self) -> bool {
2279 self.has(raw::IORING_FEAT_RECVSEND_BUNDLE)
2280 }
2281 /// `IORING_FEAT_MIN_TIMEOUT` : timeout minimal d'attente (`ABS_TIMER`).
2282 #[must_use]
2283 pub fn min_timeout(&self) -> bool {
2284 self.has(raw::IORING_FEAT_MIN_TIMEOUT)
2285 }
2286}
2287
2288/// Cible d'une annulation synchrone ([`IoUring::sync_cancel`]).
2289///
2290/// Mappée sur les flags `IORING_ASYNC_CANCEL_*`. Le paramètre de durée de vie
2291/// `'fd` borne l'emprunt du descripteur ciblé par [`CancelTarget::Fd`].
2292#[derive(Debug, Clone, Copy)]
2293pub enum CancelTarget<'fd> {
2294 /// Annule l'opération identifiée par ce jeton (flag `USERDATA`).
2295 Token(SubmissionToken),
2296 /// Annule toutes les opérations portant sur ce descripteur (flag `FD`).
2297 Fd(BorrowedFd<'fd>),
2298 /// Annule toutes les opérations de cet opcode (flag `OP`).
2299 Op(IoUringOpcode),
2300 /// Annule **toutes** les opérations en vol (flag `ANY`).
2301 Any,
2302}
2303
2304/// Restriction appliquée à un ring **avant** son activation (décision S3 ;
2305/// `struct io_uring_restriction`).
2306///
2307/// Posée via [`IoUringBuilder::restrict`], elle réduit la surface d'attaque du
2308/// ring (capability, Temps 3f / [`sandbox`]). Le ring est créé désactivé
2309/// (`R_DISABLED`) puis activé par [`IoUring::enable`].
2310#[derive(Debug, Clone, Copy)]
2311pub enum Restriction {
2312 /// N'autoriser que cet opcode de soumission.
2313 AllowOp(IoUringOpcode),
2314 /// N'autoriser que ce register opcode (numéro `IORING_REGISTER_*`).
2315 AllowRegister(u8),
2316 /// Flags SQE (`IOSQE_*`) autorisés.
2317 SqeFlagsAllowed(SubmitOptions),
2318 /// Flags SQE requis sur chaque soumission.
2319 SqeFlagsRequired(SubmitOptions),
2320}
2321
2322#[cfg(test)]
2323mod token_tests;
2324
2325#[cfg(test)]
2326mod construction_tests;
2327
2328#[cfg(test)]
2329mod submission_tests;
2330
2331#[cfg(test)]
2332mod coverage_tests;
2333
2334#[cfg(test)]
2335mod simulator_tests;
2336
2337#[cfg(test)]
2338mod teardown_tests;
2339
2340#[cfg(test)]
2341mod fork_isolation_tests;