Skip to main content

strat9_kernel/async_io/
ring.rs

1//! Ring buffer for async I/O submission and completion.
2//!
3//! Each ring consists of two pages mapped into the owning process's address
4//! space: a Submission Queue (SQ) page and a Completion Queue (CQ) page.
5//! The kernel tracks metadata (head/tail pointers, ownership) internally;
6//! only the raw SQE/CQE arrays are visible to userspace.
7
8use crate::{
9    memory::{
10        self,
11        address_space::{VmaFlags, VmaPageSize, VmaType},
12        phys_to_virt, PhysFrame,
13    },
14    sync::SpinLock,
15};
16use alloc::{sync::Arc, vec::Vec};
17use core::sync::atomic::{AtomicU32, AtomicU64, Ordering};
18
19use super::ops::{AsyncCqe, AsyncSqe, DEFAULT_RING_ENTRIES};
20
21// =============================================================================
22// Ring
23// =============================================================================
24
25/// A shared submission/completion ring for a single process.
26pub struct Ring {
27    /// Opaque handle returned to userspace.
28    pub id: u64,
29    /// Process that owns this ring (by PID).
30    pub owner_pid: u32,
31    /// Stable mapping capability for the SQ pages.
32    pub sq_mapping_cap_id: crate::capability::CapId,
33    /// Physical frame backing the SQ page.
34    pub sq_frame: PhysFrame,
35    /// Kernel virtual address of the SQ metadata header.
36    pub sq_virt: u64,
37    /// Stable mapping capability for the CQ pages.
38    pub cq_mapping_cap_id: crate::capability::CapId,
39    /// Physical frame backing the CQ page.
40    pub cq_frame: PhysFrame,
41    /// Kernel virtual address of the CQ metadata header.
42    pub cq_virt: u64,
43    /// Number of entries (power of two, min 2, max MAX_IN_FLIGHT).
44    pub entries: u32,
45    /// Number of in-flight operations (not yet completed).
46    pub in_flight: AtomicU32,
47    /// Whether this ring has been destroyed (no new ops accepted).
48    pub destroyed: AtomicU32,
49    /// Lock for pushing completions.
50    pub cq_lock: SpinLock<()>,
51    /// Completions retained temporarily when the visible CQ is full.
52    pub completion_backlog: SpinLock<Vec<AsyncCqe>>,
53    /// Order used for SQ allocation.
54    pub sq_order: u8,
55    /// Order used for CQ allocation.
56    pub cq_order: u8,
57    /// Wait queue for tasks blocked on this ring.
58    pub wq: crate::sync::WaitQueue,
59}
60
61impl Ring {
62    /// Create a new ring and map it into the calling process's address space.
63    pub fn create(pid: u32, entries: u32) -> Result<u64, RingError> {
64        if RING_REGISTRY.lock().len() >= MAX_RINGS {
65            return Err(RingError::TooManyRings);
66        }
67
68        let entries = entries.max(2).min(DEFAULT_RING_ENTRIES).next_power_of_two();
69
70        let sq_bytes = RING_META_SIZE + (entries as usize * core::mem::size_of::<AsyncSqe>());
71        let cq_bytes = RING_META_SIZE + (entries as usize * core::mem::size_of::<AsyncCqe>());
72
73        let sq_order = ((sq_bytes + 4095) / 4096)
74            .next_power_of_two()
75            .trailing_zeros() as u8;
76        let cq_order = ((cq_bytes + 4095) / 4096)
77            .next_power_of_two()
78            .trailing_zeros() as u8;
79
80        // Allocate SQ frames
81        let sq_frame =
82            crate::sync::with_irqs_disabled(|t| memory::allocate_phys_contiguous(t, sq_order))
83                .map_err(|_| RingError::Alloc)?;
84        let sq_virt = phys_to_virt(sq_frame.start_address.as_u64());
85
86        // Allocate CQ frames
87        let cq_frame = match crate::sync::with_irqs_disabled(|t| {
88            memory::allocate_phys_contiguous(t, cq_order)
89        }) {
90            Ok(f) => f,
91            Err(_) => {
92                crate::sync::with_irqs_disabled(|t| {
93                    memory::free_phys_contiguous(t, sq_frame, sq_order);
94                });
95                return Err(RingError::Alloc);
96            }
97        };
98        let cq_virt = phys_to_virt(cq_frame.start_address.as_u64());
99
100        // Initialize metadata
101        unsafe {
102            let sq = sq_virt as *mut RingMeta;
103            (*sq).entries.store(entries, Ordering::Release);
104            (*sq).mask.store(entries - 1, Ordering::Release);
105            (*sq).head.store(0, Ordering::Release);
106            (*sq).tail.store(0, Ordering::Release);
107            (*sq).flags.store(0, Ordering::Release);
108
109            let cq = cq_virt as *mut RingMeta;
110            (*cq).entries.store(entries, Ordering::Release);
111            (*cq).mask.store(entries - 1, Ordering::Release);
112            (*cq).head.store(0, Ordering::Release);
113            (*cq).tail.store(0, Ordering::Release);
114            (*cq).flags.store(0, Ordering::Release);
115        }
116
117        let ring = Arc::new(Ring {
118            id: next_ring_id(),
119            owner_pid: pid,
120            sq_mapping_cap_id: memory::allocate_mapping_cap_id(),
121            sq_frame,
122            sq_virt,
123            cq_mapping_cap_id: memory::allocate_mapping_cap_id(),
124            cq_frame,
125            cq_virt,
126            entries,
127            in_flight: AtomicU32::new(0),
128            destroyed: AtomicU32::new(0),
129            cq_lock: SpinLock::new(()),
130            completion_backlog: SpinLock::new(Vec::with_capacity(16)),
131            sq_order,
132            cq_order,
133            wq: crate::sync::WaitQueue::new(),
134        });
135
136        let id = ring.id;
137        let mut registry = RING_REGISTRY.lock();
138        if registry.len() >= MAX_RINGS {
139            return Err(RingError::TooManyRings);
140        }
141        registry.push(ring);
142        Ok(id)
143    }
144
145    /// Map the SQ and CQ buffers into the provided address space.
146    pub fn map_into_process(
147        &self,
148        addr_space: &crate::memory::AddressSpace,
149    ) -> Result<super::syscall::AsyncRingMapping, &'static str> {
150        let sq_page_count = 1usize << self.sq_order;
151        let cq_page_count = 1usize << self.cq_order;
152        let sq_phys_addrs = self.contiguous_phys_addrs(self.sq_frame, sq_page_count);
153        let cq_phys_addrs = self.contiguous_phys_addrs(self.cq_frame, cq_page_count);
154        let sq_mapping_cap_ids = alloc::vec![self.sq_mapping_cap_id; sq_page_count];
155        let cq_mapping_cap_ids = alloc::vec![self.cq_mapping_cap_id; cq_page_count];
156
157        let sq_base = addr_space
158            .find_free_vma_range(crate::kaslr::mmap_base(), sq_page_count, VmaPageSize::Small)
159            .ok_or("async ring: no free range for SQ")?;
160        addr_space.map_shared_frames_with_cap_ids(
161            sq_base,
162            &sq_phys_addrs,
163            Some(&sq_mapping_cap_ids),
164            VmaFlags {
165                readable: true,
166                writable: true,
167                executable: false,
168                user_accessible: true,
169            },
170            VmaType::Anonymous,
171        )?;
172
173        let cq_base = match addr_space.find_free_vma_range(
174            crate::syscall::mmap::MMAP_BASE,
175            cq_page_count,
176            VmaPageSize::Small,
177        ) {
178            Some(base) => base,
179            None => {
180                let _ = addr_space.unmap_region(sq_base, sq_page_count, VmaPageSize::Small);
181                return Err("async ring: no free range for CQ");
182            }
183        };
184
185        if let Err(err) = addr_space.map_shared_frames_with_cap_ids(
186            cq_base,
187            &cq_phys_addrs,
188            Some(&cq_mapping_cap_ids),
189            VmaFlags {
190                readable: true,
191                writable: true,
192                executable: false,
193                user_accessible: true,
194            },
195            VmaType::Anonymous,
196        ) {
197            let _ = addr_space.unmap_region(sq_base, sq_page_count, VmaPageSize::Small);
198            return Err(err);
199        }
200
201        Ok(super::syscall::AsyncRingMapping {
202            sq_base,
203            cq_base,
204            sq_size: (sq_page_count * 4096) as u64,
205            cq_size: (cq_page_count * 4096) as u64,
206            entries: self.entries,
207        })
208    }
209
210    fn contiguous_phys_addrs(&self, frame: PhysFrame, page_count: usize) -> Vec<u64> {
211        let base = frame.start_address.as_u64();
212        (0..page_count)
213            .map(|index| base + (index as u64) * 4096)
214            .collect()
215    }
216
217    /// Read the SQ tail (written by userspace, read by kernel).
218    #[inline]
219    pub fn sq_tail(&self) -> u32 {
220        self.sq_meta().tail.load(Ordering::Acquire)
221    }
222
223    /// Read the SQ head (written by kernel, read by userspace).
224    #[inline]
225    pub fn sq_head(&self) -> u32 {
226        self.sq_meta().head.load(Ordering::Acquire)
227    }
228
229    /// Advance the SQ head after processing submissions.
230    #[inline]
231    pub fn sq_advance_head(&self, count: u32) {
232        self.sq_meta().head.fetch_add(count, Ordering::Release);
233    }
234
235    /// Advance the CQ tail after pushing completions.
236    #[inline]
237    pub fn cq_advance_tail(&self, count: u32) {
238        self.cq_meta().tail.fetch_add(count, Ordering::Release);
239    }
240
241    /// Read the CQ head (written by userspace after consuming).
242    #[inline]
243    pub fn cq_head(&self) -> u32 {
244        self.cq_meta().head.load(Ordering::Acquire)
245    }
246
247    /// Get a pointer to the SQE at the given index.
248    #[inline]
249    pub unsafe fn sqe_at(&self, index: u32) -> *const AsyncSqe {
250        let base = (self.sq_virt + RING_META_SIZE as u64) as *const AsyncSqe;
251        base.add(index as usize)
252    }
253
254    /// Get a mutable pointer to the CQE at the given index.
255    #[inline]
256    pub unsafe fn cqe_at(&self, index: u32) -> *mut AsyncCqe {
257        let base = (self.cq_virt + RING_META_SIZE as u64) as *mut AsyncCqe;
258        base.add(index as usize)
259    }
260
261    /// Mark the ring as destroyed (no new submissions accepted).
262    #[inline]
263    pub fn destroy(&self) {
264        self.destroyed.store(1, Ordering::Release);
265    }
266
267    pub(crate) fn sq_meta(&self) -> &RingMeta {
268        unsafe { &*(self.sq_virt as *const RingMeta) }
269    }
270    pub(crate) fn cq_meta(&self) -> &RingMeta {
271        unsafe { &*(self.cq_virt as *const RingMeta) }
272    }
273}
274
275impl Drop for Ring {
276    fn drop(&mut self) {
277        let _ = memory::revoke_mapping_cap_id(self.sq_mapping_cap_id);
278        let _ = memory::revoke_mapping_cap_id(self.cq_mapping_cap_id);
279        crate::sync::with_irqs_disabled(|t| {
280            memory::free_phys_contiguous(t, self.sq_frame, self.sq_order);
281            memory::free_phys_contiguous(t, self.cq_frame, self.cq_order);
282        });
283    }
284}
285
286// =============================================================================
287// RingMeta : header at the start of each ring page
288// =============================================================================
289
290/// Metadata header embedded at offset 0 of each ring page.
291/// The SQE/CQE array starts at `RING_META_SIZE` bytes from the page base.
292#[repr(C)]
293pub(crate) struct RingMeta {
294    pub(crate) head: AtomicU32,    // Producer offset
295    pub(crate) tail: AtomicU32,    // Consumer offset
296    pub(crate) mask: AtomicU32,    // entries - 1 (for index wrapping)
297    pub(crate) entries: AtomicU32, // total number of slots
298    pub(crate) flags: AtomicU32,   // ring flags (NEED_WAKEUP, etc.)
299}
300
301const RING_META_SIZE: usize = 64; // leave room for future flags
302
303// =============================================================================
304// RingRegistry
305// =============================================================================
306
307const MAX_RINGS: usize = 128;
308
309static RING_REGISTRY: SpinLock<Vec<Arc<Ring>>> = SpinLock::new(Vec::new());
310
311static NEXT_RING_ID: AtomicU64 = AtomicU64::new(1);
312
313fn next_ring_id() -> u64 {
314    NEXT_RING_ID.fetch_add(1, Ordering::Relaxed)
315}
316
317/// Find a ring by its opaque id.
318pub fn find_ring(id: u64) -> Option<Arc<Ring>> {
319    let guard = RING_REGISTRY.lock();
320    guard.iter().find(|r| r.id == id).cloned()
321}
322
323/// Remove a ring from the registry and free its pages.
324pub fn destroy_ring(id: u64) -> Result<(), RingError> {
325    let mut guard = RING_REGISTRY.lock();
326    if let Some(pos) = guard.iter().position(|r| r.id == id) {
327        let ring = guard.remove(pos);
328        ring.destroy();
329        ring.wq.wake_all();
330        crate::hardware::storage::ahci::discard_deferred_async_read_completions(id);
331        let _ = memory::revoke_mapping_cap_id(ring.sq_mapping_cap_id);
332        let _ = memory::revoke_mapping_cap_id(ring.cq_mapping_cap_id);
333        crate::sync::with_irqs_disabled(|t| {
334            memory::free_phys_contiguous(t, ring.sq_frame, ring.sq_order);
335            memory::free_phys_contiguous(t, ring.cq_frame, ring.cq_order);
336        });
337        Ok(())
338    } else {
339        Err(RingError::NotFound)
340    }
341}
342
343// =============================================================================
344// RingError
345// =============================================================================
346
347#[derive(Debug, Clone, Copy, PartialEq, Eq)]
348pub enum RingError {
349    Alloc,
350    NotFound,
351    TooManyRings,
352}