Skip to main content

air_sys_syscall/io_uring/
multishot.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 **multishot** (Temps 3d) : une **unique soumission** produit un
6//! **flux** de complétions. Sous-module `air-sys-syscall::io_uring::multishot`.
7//! Repose sur le drapeau de complétion `CQE_F_MORE` (axe E), sur les buffers
8//! fournis du **Temps 3b** (recv/read) et les descripteurs directs du **Temps
9//! 3a** (accept direct). **Aucun register opcode** ; un opcode dédié
10//! (`READ_MULTISHOT`, 49) et des drapeaux d'op.
11//!
12//! Référence normative : `docs/specs/layer-0/io-uring-3d-multishot.md`.
13//!
14//! ## Cycle de vie (slab S1)
15//!
16//! Chaque complétion intermédiaire porte **`CQE_F_MORE`** (« d'autres suivront
17//! pour ce SQE ») ; la **terminale** ne le porte pas (erreur, pénurie de
18//! buffers, ou annulation). Le slot S1 reste **vivant** tant que `CQE_F_MORE`
19//! et n'est **libéré qu'à la terminale** — second cas (après le NOTIF zero-copy
20//! du Temps 2b) où un slot survit à sa première complétion. La **génération** du
21//! jeton filtre les CQE tardifs après annulation (Temps 1 §4.2).
22//!
23//! Une op multishot rend un [`MultishotToken`](super::MultishotToken) (distinct
24//! du `SubmissionToken` mono-coup) ; toutes ses complétions le portent
25//! ([`Completion::multishot_token`](super::Completion::multishot_token)).
26//! [`Completion::has_more`](super::Completion::has_more) distingue les
27//! intermédiaires de la terminale. Le **réarmement** (resoumission) est explicite.
28
29use super::owned::OwnedOp;
30use super::raw::{self, IoUringSqe};
31use super::{CancelTarget, IoUring, MultishotToken, ProvidedBufferRing, SubmitOptions};
32use air_sys_types::Errno;
33use air_sys_types::fd::{AsRawFd, BorrowedFd};
34use air_sys_types::io_uring::{PollEvents, TimeoutFlags, TimeoutSpec};
35use air_sys_types::net::{AcceptFlags, MessageFlags};
36use alloc::boxed::Box;
37use core::time::Duration;
38
39impl IoUring {
40    /// Réserve un slot **multishot** (S1, survit aux `CQE_F_MORE`), prépare le SQE
41    /// (rempli par `fill`), pose `user_data`/`personality`/flags `IOSQE_*`, et rend
42    /// un [`MultishotToken`]. Sans publication (comme `submit_op`). Pas de gestion
43    /// `skip_cqe` : un multishot **consomme chaque CQE** (le slot vit jusqu'à la
44    /// terminale).
45    fn submit_op_multishot(
46        &mut self,
47        payload: Option<OwnedOp>,
48        fill: impl FnOnce(&mut IoUringSqe),
49    ) -> Result<MultishotToken, Errno> {
50        if self.sq.space_left() == 0 {
51            return Err(Errno::EBUSY);
52        }
53        let token = self
54            .slab
55            .reserve_multishot(payload)
56            .map_err(|_| Errno::EBUSY)?;
57        let opts = self.pending_options;
58        self.pending_options = SubmitOptions::default();
59        let personality = self.pending_personality;
60        self.pending_personality = 0;
61        // SAFETY: place SQ vérifiée > 0 juste au-dessus ; le slot (zéro-initialisé
62        // par `prepare`) est rempli intégralement ci-dessous avant toute publication.
63        let sqe_ptr = unsafe { self.sq.prepare() }.expect("place SQ vérifiée");
64        // SAFETY: `sqe_ptr` pointe le slot SQE, exclusivement détenu jusqu'à la
65        // publication ; le `&mut` ne vit que cette portée.
66        unsafe {
67            let sqe = &mut *sqe_ptr;
68            fill(sqe);
69            sqe.user_data = token.to_user_data();
70            sqe.personality = personality;
71            sqe.flags = opts.iosqe_flags();
72        }
73        Ok(MultishotToken::new(token.slot(), token.generation()))
74    }
75
76    // ── Accept multishot (§2) ─────────────────────────────────────────────
77
78    /// `ACCEPT` (13) + `IORING_ACCEPT_MULTISHOT` : **un seul SQE accepte en
79    /// continu** — chaque connexion produit une complétion portant le FD accepté
80    /// ([`Completion::accepted_fd`](super::Completion::accepted_fd)), `CQE_F_MORE`
81    /// maintenu jusqu'à la terminale (ex. listener fermé). `O_CLOEXEC` posé d'office.
82    ///
83    /// # Errors
84    ///
85    /// [`Errno::EBUSY`] si la SQ ou le slab sont pleins.
86    pub fn submit_accept_multishot(
87        &mut self,
88        listener: BorrowedFd<'_>,
89        flags: AcceptFlags,
90    ) -> Result<MultishotToken, Errno> {
91        let fd = listener.as_raw_fd();
92        let accept_flags = (flags | AcceptFlags::CLOEXEC).bits().cast_unsigned();
93        self.submit_op_multishot(None, |sqe| {
94            sqe.opcode = raw::IORING_OP_ACCEPT;
95            sqe.fd = fd;
96            sqe.op_flags = accept_flags;
97            sqe.ioprio |= raw::IORING_ACCEPT_MULTISHOT;
98        })
99    }
100
101    /// Variante **descripteurs directs** : chaque connexion atterrit dans un slot
102    /// **auto-alloué** de la `FixedFdTable` (Temps 3a, `FixedSlotTarget::Alloc`).
103    /// Le slot choisi est rendu dans `cqe->res`
104    /// ([`Completion::allocated_slot`](super::Completion::allocated_slot)). Pour
105    /// les serveurs à fort taux de connexions (pas de FD ordinaire).
106    ///
107    /// Le kernel **rejette** `SOCK_CLOEXEC` sur un descripteur direct : non posé
108    /// (le `O_CLOEXEC` est appliqué à `fixed_fd_install`).
109    ///
110    /// # Errors
111    ///
112    /// Voir [`IoUring::submit_accept_multishot`].
113    pub fn submit_accept_multishot_direct(
114        &mut self,
115        listener: BorrowedFd<'_>,
116        flags: AcceptFlags,
117    ) -> Result<MultishotToken, Errno> {
118        let fd = listener.as_raw_fd();
119        let accept_flags = flags.bits().cast_unsigned();
120        self.submit_op_multishot(None, |sqe| {
121            sqe.opcode = raw::IORING_OP_ACCEPT;
122            sqe.fd = fd;
123            sqe.op_flags = accept_flags;
124            sqe.ioprio |= raw::IORING_ACCEPT_MULTISHOT;
125            // Auto-allocation dans la table de FD fixes (chaque connexion → slot).
126            sqe.splice_fd_in_or_file_index = raw::IORING_FILE_INDEX_ALLOC;
127        })
128    }
129
130    // ── Recv / Read multishot (§3, buffers fournis 3b) ────────────────────
131
132    /// `RECV` (27) + `IORING_RECV_MULTISHOT` sur les buffers de `group`
133    /// (`IOSQE_BUFFER_SELECT`) : chaque arrivée de données prend un buffer du
134    /// groupe et produit une complétion. Complétion via
135    /// [`Completion::into_provided_buffer`](super::Completion::into_provided_buffer).
136    ///
137    /// **Pénurie de buffers** : si le groupe est vide à l'arrivée de données, la
138    /// complétion porte **`-ENOBUFS` et TERMINE** le multishot (pas de `F_MORE`) —
139    /// à distinguer d'une erreur réseau. L'application réapprovisionne puis
140    /// **resoumet**.
141    ///
142    /// # Errors
143    ///
144    /// [`Errno::EBUSY`] si la SQ/le slab sont pleins.
145    pub fn submit_receive_multishot(
146        &mut self,
147        sock: BorrowedFd<'_>,
148        group: &ProvidedBufferRing,
149        flags: MessageFlags,
150    ) -> Result<MultishotToken, Errno> {
151        let fd = sock.as_raw_fd();
152        let bgid = group.group_id();
153        let msg_flags = flags.bits().cast_unsigned();
154        self.pending_options = self.pending_options.buffer_select();
155        self.submit_op_multishot(None, |sqe| {
156            sqe.opcode = raw::IORING_OP_RECV;
157            sqe.fd = fd;
158            sqe.op_flags = msg_flags;
159            sqe.buf_index_or_group = bgid;
160            sqe.ioprio |= raw::IORING_RECV_MULTISHOT;
161        })
162    }
163
164    /// `READ_MULTISHOT` (49) sur les buffers de `group` : flux de lectures, chaque
165    /// arrivée prenant un buffer du groupe. Même sémantique de pénurie `-ENOBUFS`
166    /// terminante que [`IoUring::submit_receive_multishot`]. L'opcode dédié implique
167    /// le multishot (pas de drapeau d'op séparé).
168    ///
169    /// # Errors
170    ///
171    /// Voir [`IoUring::submit_receive_multishot`].
172    pub fn submit_read_multishot(
173        &mut self,
174        fd: BorrowedFd<'_>,
175        group: &ProvidedBufferRing,
176        offset: Option<u64>,
177    ) -> Result<MultishotToken, Errno> {
178        let fd = fd.as_raw_fd();
179        let bgid = group.group_id();
180        let off = offset.unwrap_or(u64::MAX);
181        self.pending_options = self.pending_options.buffer_select();
182        self.submit_op_multishot(None, |sqe| {
183            sqe.opcode = raw::IORING_OP_READ_MULTISHOT;
184            sqe.fd = fd;
185            sqe.off_or_addr2 = off;
186            sqe.buf_index_or_group = bgid;
187        })
188    }
189
190    // ── Poll multishot (§4) ───────────────────────────────────────────────
191
192    /// `POLL_ADD` (6) + `IORING_POLL_ADD_MULTI` : chaque transition d'état du FD
193    /// vers les `events` surveillés produit une complétion
194    /// ([`Completion::into_poll_result`](super::Completion::into_poll_result)),
195    /// `F_MORE` maintenu. **Edge-triggered** (la surface validée ne porte pas
196    /// d'option level-triggered).
197    ///
198    /// # Errors
199    ///
200    /// [`Errno::EBUSY`] si la SQ/le slab sont pleins.
201    pub fn submit_poll_multishot(
202        &mut self,
203        fd: BorrowedFd<'_>,
204        events: PollEvents,
205    ) -> Result<MultishotToken, Errno> {
206        let fd = fd.as_raw_fd();
207        let poll_events = events.bits();
208        self.submit_op_multishot(None, |sqe| {
209            sqe.opcode = raw::IORING_OP_POLL_ADD;
210            sqe.fd = fd;
211            sqe.op_flags = poll_events; // poll32_events
212            sqe.len = raw::IORING_POLL_ADD_MULTI; // drapeaux poll (multishot)
213        })
214    }
215
216    // ── Timeout multishot (§5) ────────────────────────────────────────────
217
218    /// `TIMEOUT` (11) + `IORING_TIMEOUT_MULTISHOT` : émet une complétion à
219    /// **intervalle régulier** (timer répétitif) jusqu'à annulation — battement
220    /// périodique dans le reactor sans resoumettre à chaque tick. Le
221    /// `__kernel_timespec` est gardé en vie dans le slot (vivant jusqu'à la
222    /// terminale).
223    ///
224    /// # Errors
225    ///
226    /// [`Errno::EBUSY`] si la SQ/le slab sont pleins.
227    pub fn submit_timeout_multishot(
228        &mut self,
229        interval: Duration,
230        flags: TimeoutFlags,
231    ) -> Result<MultishotToken, Errno> {
232        let ts = Box::new(super::async_ops::timespec_of(TimeoutSpec::after(interval)));
233        let ts_ptr = core::ptr::from_ref(&*ts) as u64;
234        let op_flags = flags.bits() | raw::IORING_TIMEOUT_MULTISHOT;
235        self.submit_op_multishot(Some(OwnedOp::Timeout(ts)), |sqe| {
236            sqe.opcode = raw::IORING_OP_TIMEOUT;
237            sqe.fd = -1;
238            sqe.addr_or_splice_off_in = ts_ptr;
239            sqe.len = 1;
240            sqe.op_flags = op_flags;
241        })
242    }
243
244    // ── Annulation (§6) ───────────────────────────────────────────────────
245
246    /// Annule un multishot en vol, en ciblant son `token`. La complétion terminale
247    /// (sans `F_MORE`, souvent `-ECANCELED`) **libère le slot** ; les CQE tardifs
248    /// éventuels sont filtrés par la **génération** (S1).
249    ///
250    /// Utilise l'annulation **synchrone** par jeton (`IORING_REGISTER_SYNC_CANCEL`,
251    /// `CancelTarget::Token`) — cohérente avec un retour `Result<()>` (pas de CQE
252    /// d'annulation supplémentaire à drainer ; la terminale du multishot suffit).
253    /// `submit_cancel` (Temps 2c) reste disponible pour les annulations groupées.
254    ///
255    /// # Errors
256    ///
257    /// [`Errno::ENOENT`] si le multishot s'est déjà terminé (rien à annuler) ;
258    /// autres erreurs de `io_uring_register`.
259    pub fn cancel_multishot(&mut self, token: MultishotToken) -> Result<(), Errno> {
260        self.sync_cancel(CancelTarget::Token(token.as_submission()))
261            .map(|_count| ())
262    }
263}
264
265// ───────────────────────────────────────────────────────────────────────────
266// Tests
267// ───────────────────────────────────────────────────────────────────────────
268//
269// Intégration kernel (multishot non modélisé par Miri) → `#[cfg_attr(miri, ignore)]`.
270// Le cycle de vie S1 multishot (slot vivant tant que F_MORE, libéré à la
271// terminale, génération anti-CQE-tardif) est prouvé **sous Miri** par les tests
272// purs du slab (`slab::proptests::multishot_slot_survives_until_final_…`).
273
274#[cfg(test)]
275mod tests;