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;