Skip to main content

air_sys_syscall/io_uring/
linked.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 liées** (Temps 3c) : exprimer qu'une opération ne démarre
6//! qu'après l'achèvement de la précédente, **sans aller-retour userspace**.
7//! Sous-module `air-sys-syscall::io_uring::linked`. Repose sur les drapeaux SQE
8//! `IOSQE_IO_LINK` / `IOSQE_IO_HARDLINK` (axe D) et intègre l'op `LINK_TIMEOUT`
9//! (15). **Aucun register opcode.**
10//!
11//! Référence normative : `docs/specs/layer-0/io-uring-3c-linked.md`.
12//!
13//! ## Sémantique soft vs hard (point critique)
14//!
15//! L'ordre est garanti **dans les deux cas** ; ils diffèrent par la propagation
16//! des erreurs :
17//!
18//! - **soft** ([`LinkedChainBuilder::then`], `IOSQE_IO_LINK`) : si un maillon se
19//!   termine **en erreur**, la chaîne est **rompue** — les maillons restants sont
20//!   complétés avec **`-ECANCELED`**. **Définition large de l'« erreur »** : un
21//!   **`read`/`recv` court** (short read — moins d'octets que demandé) **compte
22//!   comme erreur** et rompt la chaîne. *Piège classique* : `read(N) → process`
23//!   casse si le `read` rend `< N` octets.
24//! - **hard** ([`LinkedChainBuilder::then_hard`], `IOSQE_IO_HARDLINK`) : une
25//!   erreur de **complétion** d'un maillon **n'interrompt pas** (suivants
26//!   exécutés) ; implique l'ordre du soft. Ne protège **pas** d'un **échec de
27//!   soumission** du parent (qui rompt quand même).
28//!
29//! ## Staging + réservation atomique + publication unique
30//!
31//! Pendant la durée du builder, le `&mut IoUring` est en **mode staging** : un
32//! `submit_*` appelé dans la closure réserve un slot S1 et **écrit** un SQE (avec
33//! le drapeau de lien posé sur le **prédécesseur**) **sans publier la queue**. La
34//! publication n'a lieu qu'au [`LinkedChainBuilder::submit`] final (**un seul
35//! store-release**, Temps 1 §3.2 → le kernel voit la chaîne complète d'un coup).
36//! On réutilise ainsi les `submit_*` existants **sans duplication** de `prepare_*`.
37//! Si un maillon ne peut être réservé (`EBUSY`), le builder **rembobine** (libère
38//! les slots déjà réservés, rembobine la `tail`) avant de retourner — **jamais de
39//! chaîne partielle publiée**.
40
41use super::owned::OwnedOp;
42use super::{IoUring, SubmissionToken, raw};
43use air_sys_types::Errno;
44use air_sys_types::io_uring::{TimeoutFlags, TimeoutSpec};
45use alloc::boxed::Box;
46use alloc::vec::Vec;
47
48impl IoUring {
49    /// Démarre la construction d'une **chaîne liée** (mode staging).
50    ///
51    /// Les `submit_*` appelés via le builder sont **mis en attente** (slot S1
52    /// réservé + SQE écrit, lien posé sur le prédécesseur) sans publier la queue ;
53    /// [`LinkedChainBuilder::submit`] publie la chaîne entière en une fois.
54    pub fn link_chain(&mut self) -> LinkedChainBuilder<'_> {
55        let start_tail = self.sq.local_tail();
56        LinkedChainBuilder {
57            ring: self,
58            start_tail,
59            tokens: Vec::new(),
60            prev_index: None,
61        }
62    }
63
64    /// `IORING_OP_LINK_TIMEOUT` (15) **staged** : borne le maillon précédent. Le
65    /// `__kernel_timespec` est gardé en vie dans le slot S1 (comme `submit_timeout`).
66    /// Réservé au [`LinkedChainBuilder`].
67    fn submit_link_timeout_staged(
68        &mut self,
69        spec: TimeoutSpec,
70        flags: TimeoutFlags,
71    ) -> Result<SubmissionToken, Errno> {
72        let ts = Box::new(super::async_ops::timespec_of(spec));
73        let ts_ptr = core::ptr::from_ref(&*ts) as u64;
74        let op_flags = flags.bits();
75        self.submit_op(Some(OwnedOp::Timeout(ts)), |sqe| {
76            sqe.opcode = raw::IORING_OP_LINK_TIMEOUT;
77            sqe.fd = -1;
78            sqe.addr_or_splice_off_in = ts_ptr;
79            sqe.len = 1;
80            sqe.op_flags = op_flags;
81        })
82    }
83}
84
85/// Constructeur d'une **chaîne d'opérations liées** sur un [`IoUring`] en mode
86/// staging (cf. module). Emprunte le ring `&mut` pour toute sa durée (aucune
87/// soumission directe ne peut s'intercaler).
88pub struct LinkedChainBuilder<'ring> {
89    /// Ring emprunté exclusivement (staging).
90    ring: &'ring mut IoUring,
91    /// `tail` locale avant le premier maillon (borne de rollback).
92    start_tail: u32,
93    /// Jetons des maillons stagés, dans l'ordre.
94    tokens: Vec<SubmissionToken>,
95    /// Index (dans l'anneau SQ) du **dernier** maillon stagé, pour y poser le
96    /// drapeau de lien au maillon suivant. `None` avant le premier.
97    prev_index: Option<u32>,
98}
99
100impl LinkedChainBuilder<'_> {
101    /// Cœur commun : stage un maillon (`op` = un `submit_*` normal). Si
102    /// `link_flag` est `Some`, le pose sur le **prédécesseur**. En cas d'échec de
103    /// réservation, **rembobine** la chaîne et propage l'erreur.
104    fn stage<F>(mut self, link_flag: Option<u8>, op: F) -> Result<Self, Errno>
105    where
106        F: FnOnce(&mut IoUring) -> Result<SubmissionToken, Errno>,
107    {
108        let index = self.ring.sq.local_tail();
109        match op(self.ring) {
110            Ok(token) => {
111                if let (Some(flag), Some(prev)) = (link_flag, self.prev_index) {
112                    // SAFETY: `prev` désigne le SQE du maillon précédent, stagé
113                    // pendant ce builder et **non publié** ; accès exclusif (`&mut`).
114                    unsafe {
115                        self.ring.sq.or_sqe_flags(prev, flag);
116                    }
117                }
118                self.tokens.push(token);
119                self.prev_index = Some(index);
120                Ok(self)
121            }
122            Err(error) => {
123                self.rollback();
124                Err(error)
125            }
126        }
127    }
128
129    /// **Premier maillon** : `op` appelle un `submit_*` normal, mis en attente
130    /// (staging). Ne pose aucun drapeau de lien (pas de prédécesseur).
131    ///
132    /// # Errors
133    ///
134    /// L'erreur du `submit_*` (ex. [`Errno::EBUSY`] si la SQ/le slab sont pleins) ;
135    /// la chaîne est alors rembobinée.
136    pub fn first<F>(self, op: F) -> Result<Self, Errno>
137    where
138        F: FnOnce(&mut IoUring) -> Result<SubmissionToken, Errno>,
139    {
140        self.stage(None, op)
141    }
142
143    /// Maillon **soft-lié** au précédent (`IOSQE_IO_LINK` posé sur le précédent) :
144    /// une **erreur de complétion** du précédent (`-errno`, **short read inclus**)
145    /// rompt la chaîne (suivants `-ECANCELED`).
146    ///
147    /// # Errors
148    ///
149    /// Voir [`LinkedChainBuilder::first`].
150    pub fn then<F>(self, op: F) -> Result<Self, Errno>
151    where
152        F: FnOnce(&mut IoUring) -> Result<SubmissionToken, Errno>,
153    {
154        self.stage(Some(raw::IOSQE_IO_LINK), op)
155    }
156
157    /// Maillon **hard-lié** au précédent (`IOSQE_IO_HARDLINK` posé sur le
158    /// précédent) : une erreur de complétion du précédent **n'interrompt pas** la
159    /// chaîne (l'ordre reste garanti).
160    ///
161    /// # Errors
162    ///
163    /// Voir [`LinkedChainBuilder::first`].
164    pub fn then_hard<F>(self, op: F) -> Result<Self, Errno>
165    where
166        F: FnOnce(&mut IoUring) -> Result<SubmissionToken, Errno>,
167    {
168        self.stage(Some(raw::IOSQE_IO_HARDLINK), op)
169    }
170
171    /// Borne dans le temps le **maillon précédent** via un `LINK_TIMEOUT` (op 15) :
172    /// si le timeout expire avant la fin du maillon, le maillon est annulé
173    /// (`-ECANCELED`) et la chaîne rompue (le timeout complète) ; si le maillon
174    /// finit avant, le timeout est annulé. **Unique** manière sûre d'émettre un
175    /// `LINK_TIMEOUT` (émission isolée refusée, Temps 2c).
176    ///
177    /// # Errors
178    ///
179    /// [`Errno::EINVAL`] s'il n'y a **aucun** maillon à borner ; sinon l'erreur du
180    /// staging (chaîne rembobinée).
181    pub fn with_link_timeout(
182        mut self,
183        spec: TimeoutSpec,
184        flags: TimeoutFlags,
185    ) -> Result<Self, Errno> {
186        let Some(prev) = self.prev_index else {
187            // Pas de maillon précédent à borner.
188            self.rollback();
189            return Err(Errno::EINVAL);
190        };
191        let index = self.ring.sq.local_tail();
192        match self.ring.submit_link_timeout_staged(spec, flags) {
193            Ok(token) => {
194                // SAFETY: `prev` (maillon borné) est stagé non publié ; on le lie
195                // (`IO_LINK`) au `LINK_TIMEOUT` qui le suit immédiatement. Accès
196                // exclusif (`&mut`).
197                unsafe {
198                    self.ring.sq.or_sqe_flags(prev, raw::IOSQE_IO_LINK);
199                }
200                self.tokens.push(token);
201                self.prev_index = Some(index);
202                Ok(self)
203            }
204            Err(error) => {
205                self.rollback();
206                Err(error)
207            }
208        }
209    }
210
211    /// Publie la chaîne **entière, atomiquement** (un seul store-release + un
212    /// `io_uring_enter`). Rend les jetons dans l'ordre de soumission.
213    ///
214    /// # Errors
215    ///
216    /// Les erreurs de [`IoUring::submit`] ([`Errno::EINTR`] remonté tel quel, etc.).
217    /// La réservation des slots ayant déjà réussi au staging, l'atomicité de la
218    /// chaîne est garantie en amont (jamais de chaîne partielle).
219    pub fn submit(mut self) -> Result<ChainTokens, Errno> {
220        self.ring.submit()?;
221        Ok(ChainTokens {
222            tokens: core::mem::take(&mut self.tokens),
223        })
224    }
225
226    /// Rembobine le staging : libère les slots déjà réservés (idempotent côté
227    /// slab) et rembobine la `tail` locale → aucun SQE de la chaîne ne sera publié.
228    fn rollback(&mut self) {
229        let tokens = core::mem::take(&mut self.tokens);
230        for token in tokens {
231            self.ring.slab.release(token);
232        }
233        self.ring.sq.rewind_to(self.start_tail);
234        self.prev_index = None;
235    }
236}
237
238/// Jetons des maillons d'une chaîne, **dans l'ordre de soumission** (corrèle
239/// chaque complétion à son maillon).
240#[derive(Debug, Clone)]
241pub struct ChainTokens {
242    tokens: Vec<SubmissionToken>,
243}
244
245impl ChainTokens {
246    /// Jetons des maillons, dans l'ordre de soumission.
247    #[must_use]
248    pub fn tokens(&self) -> &[SubmissionToken] {
249        &self.tokens
250    }
251
252    /// Jeton du maillon `index` (ordre de soumission), le cas échéant.
253    #[must_use]
254    pub fn get(&self, index: usize) -> Option<SubmissionToken> {
255        self.tokens.get(index).copied()
256    }
257
258    /// Nombre de maillons de la chaîne.
259    #[must_use]
260    pub fn len(&self) -> usize {
261        self.tokens.len()
262    }
263
264    /// `true` si la chaîne est vide.
265    #[must_use]
266    pub fn is_empty(&self) -> bool {
267        self.tokens.is_empty()
268    }
269
270    /// Consomme et rend les jetons possédés.
271    #[must_use]
272    pub fn into_vec(self) -> Vec<SubmissionToken> {
273        self.tokens
274    }
275}
276
277// ───────────────────────────────────────────────────────────────────────────
278// Tests
279// ───────────────────────────────────────────────────────────────────────────
280//
281// Les chaînes touchent `io_uring_enter` (non modélisé par Miri) →
282// `#[cfg_attr(miri, ignore)]`. Les accesseurs de `ChainTokens` (logique pure)
283// tournent **sous Miri**. Les erreurs déterministes sont injectées via
284// `submit_nop_with_result` (`IORING_NOP_INJECT_RESULT`) ; le piège du **short
285// read** est testé sur un pipe réel.
286
287#[cfg(test)]
288mod tests;