Skip to main content

starnix_core/perf/
lockless_ring_buffer.rs

1// Copyright 2026 The Fuchsia Authors. All rights reserved.
2// Use of this source code is governed by a BSD-style license that can be
3// found in the LICENSE file.
4
5//! This is a lockless ring buffer modeled after https://docs.kernel.org/trace/ring-buffer-design.html
6use crate::mm::memory::MemoryObject;
7use crate::vfs::OutputBuffer;
8use fuchsia_rcu::RcuDroppable;
9use fuchsia_runtime::vmar_root_self;
10use fuchsia_trace;
11use shared_buffer::SharedBuffer;
12use starnix_logging::{log_error, log_info, log_warn};
13use starnix_sync::{LockDepMutex, PerfRingBufferStateLock};
14use starnix_types::PAGE_SIZE;
15use starnix_uapi::errors::Errno;
16use starnix_uapi::{errno, error, from_status_like_fdio};
17use std::sync::atomic::{AtomicUsize, Ordering};
18
19#[derive(RcuDroppable)]
20struct Node {
21    // page_data is the memory view for this page.
22    page_data: SharedBuffer,
23    // next and prev are pointers to the next and previous nodes.
24    // creating a linked list of pages.
25    next: AtomicUsize,
26    prev: AtomicUsize,
27    // offset to write new data in this page.
28    write_offset: AtomicUsize,
29    // Number of active writers currently reserving or writing to this page.
30    active_writers: AtomicUsize,
31}
32
33const FLAG_MASK: usize = 0b11 << 62;
34// Normal link between pages in the circular list.
35const FLAG_NORMAL: usize = 0b00 << 62;
36// Marks the link pointing to the head page (the oldest page with data).
37const FLAG_HEADER: usize = 0b01 << 62;
38// Set on the link pointing to the head page when it is being swapped by the reader.
39const FLAG_UPDATE: usize = 0b10 << 62;
40
41// The highest bit (bit 63) of `active_writers` is used as a flag indicating if the page is active.
42// When the page is finalized, this bit is cleared to 0.
43const PAGE_ACTIVE_BIT: usize = 1 << 63;
44// The second highest bit (bit 62) of `active_writers` is used to coordinate/claim finalization.
45const PAGE_FINALIZED_BIT: usize = 1 << 62;
46const FLAGS_MASK: usize = PAGE_ACTIVE_BIT | PAGE_FINALIZED_BIT;
47const ACTIVE_WRITERS_MASK: usize = !FLAGS_MASK;
48
49// When yielding under high contention, if we busy-spin we risk priority inversion and CPU starvation
50// of preempted writers. Sleeping for 50 microseconds gives the OS scheduler sufficient window (accounting
51// for 2-5us context switch overhead and 15-30us thread execution time) to reschedule and execute the
52// preempted writer so it can release its reservation.
53const SPIN_SLEEP_DURATION: std::time::Duration = std::time::Duration::from_micros(50);
54
55// When progressive sleep is triggered (exceeding 1000 yields), we back off with a 10 milliseconds sleep.
56// This is a robust scheduling window to guarantee that Zircon schedules the preempted writer thread
57// to complete its commit/release logic.
58const PROGRESSIVE_SLEEP_DURATION: std::time::Duration = std::time::Duration::from_millis(10);
59
60impl Node {
61    fn get_index(val: usize) -> usize {
62        val & !FLAG_MASK
63    }
64    fn get_flags(val: usize) -> usize {
65        val & FLAG_MASK
66    }
67    fn make_val(index: usize, flags: usize) -> usize {
68        index | flags
69    }
70
71    /// Finalizes the page by writing the data size to the header if tracing is enabled.
72    fn finalize(&self) {
73        // Attempt to claim the finalization role for this page.
74        // If another thread already claimed it, return early to avoid a data race on `page_data`.
75        let old_val = self.active_writers.fetch_or(PAGE_FINALIZED_BIT, Ordering::AcqRel);
76        if old_val & PAGE_FINALIZED_BIT != 0 {
77            return;
78        }
79
80        // - Acquire: We load `write_offset` with Acquire ordering to synchronize with all concurrent
81        //   writers that updated it via CAS in `try_reserve_on_page`. This ensures that
82        //   all non-atomic writes to `page_data` performed by the writers before their reservation was completed
83        //   are fully visible to this finalizing thread.
84        let write_offset = self.write_offset.load(Ordering::Acquire);
85        let data_size = std::cmp::min(write_offset, (*PAGE_SIZE) as usize)
86            - LocklessRingBuffer::PAGE_HEADER_SIZE;
87        self.page_data.write_at(8, &(data_size as u64).to_le_bytes());
88
89        // Clear the `PAGE_ACTIVE_BIT` from `active_writers` with Release ordering to publish the completed page.
90        // Release semantics here guarantees that if a reader observes the PAGE_ACTIVE_BIT being cleared,
91        // it is safe to read the page as it will also observe all the completed writes from the writers for that page.
92        self.active_writers.fetch_and(!PAGE_ACTIVE_BIT, Ordering::Release);
93    }
94
95    /// Releases a writer reservation on this node, and finalizes the node if it was the last writer.
96    fn release_writer(&self) {
97        // Using Release ordering with fetch_sub make the loading part use Relaxed ordering, and the
98        // store part Release. This means the value will be updated before the Acquire loading.
99        let prev_writers = self.active_writers.fetch_sub(1, Ordering::Release);
100
101        // If this was the last writer finalizing an active page (`PAGE_ACTIVE_BIT | 1`), we ensure the page is finalized.
102        // Note that if the page was already finalized by the advancing thread, `PAGE_FINALIZED_BIT` will be set,
103        // meaning `prev_writers` will have it set and we won't enter this branch, which is correct since it's already finalized.
104        if prev_writers == PAGE_ACTIVE_BIT | 1 {
105            let write_offset = self.write_offset.load(Ordering::Acquire);
106            if write_offset >= (*PAGE_SIZE) as usize {
107                // Since this thread holds a valid `Reservation`, `ref_count` remains > 0, which blocks
108                // `disable()` from completing and shrinking the VMO. This guarantees that the ring VMO
109                // remains enabled and valid for memory access.
110                // Without this check and finalization on failure, there is a liveness bug where a full page
111                // could never be finalized because the successful writer already skipped finalization.
112                self.finalize();
113            }
114        }
115    }
116}
117
118// We are only 64 bit today, but make it easy to find this assumption if we go to a smaller arch.
119static_assertions::const_assert!(std::mem::size_of::<usize>() == 8);
120
121// Use the high bit to indicate the ring is enabled. This makes the default, 0,
122// which is disabled, an easy value to reason about. And the ref count is the lower 63 bits.
123const RING_ENABLED_BIT: usize = 1 << 63;
124
125#[derive(RcuDroppable)]
126pub struct LocklessRingBuffer {
127    vmo: MemoryObject,
128    mapping: SharedBuffer,
129    // linked list of pages that represent the ring, backed by the mapping (and the vmo).
130    nodes: Vec<Node>,
131    // The page where to read the data from the ring.
132    head_page: AtomicUsize,
133    // The page where writes add to the ring.
134    tail_page: AtomicUsize,
135
136    // A page used by the reader of the ring.
137    reader_page: AtomicUsize,
138
139    // Tracks the global number of active readers and writers of the ring in the lower 63 bits.
140    // The highest bit (RING_ENABLED_BIT) is set when the ring is enabled.
141    ref_count: AtomicUsize,
142    // True indicates drop old pages to write new data when the ring is full. If false, writes fail until
143    // the data is read, effectively dropping new data and keeping the old data.
144    overwrite: bool,
145    // Number of dropped pages when overwrite is true.
146    dropped_pages: std::sync::atomic::AtomicU64,
147    // Used to calculate the delta from the last event. Since this is written to by multiple threads,
148    // we keep an atomic value for this.
149    prev_timestamp: std::sync::atomic::AtomicU64,
150    // The async trace event ID for write events.
151    write_event_async_id: fuchsia_trace::Id,
152    // Tracks whether a reader is currently active, used to validate that concurrent reads
153    // on a single ring buffer are not supported and do not happen. This can be removed if we
154    // are convinced it is not necessary.
155    reader_active: std::sync::atomic::AtomicBool,
156    // Mutex to serialize enable() and disable() calls to prevent racing. This should
157    // never be accessed by other methods to avoid locking during reading and writing.
158    state_mutex: LockDepMutex<(), PerfRingBufferStateLock>,
159}
160impl LocklessRingBuffer {
161    // This is an Ftrace page header consisting of a u64 timestamp at offset 0 and a u64 data size at offset 8.
162    pub const PAGE_HEADER_SIZE: usize = 16;
163
164    /// Creates a new LocklessRingBuffer.
165    /// size_bytes: Size of the ring buffer. Must be at least PAGE_SIZE * 3 and will be rounded up to a multiple of PAGE_SIZE.
166    /// overwrite: If true, old pages will be dropped when the ring is full. If false, writes will fail until the data is read.
167    /// write_event_async_id: The async trace event ID for write events.
168    pub fn new(
169        size_bytes: usize,
170        overwrite: bool,
171        write_event_async_id: fuchsia_trace::Id,
172    ) -> Result<Self, Errno> {
173        let requested_pages = (size_bytes + (*PAGE_SIZE) as usize - 1) / (*PAGE_SIZE) as usize;
174        // 3 pages are needed, 1 for the read page, 1 for head, and one to swap.
175        let pages = std::cmp::max(3, requested_pages);
176        let total_nodes = pages;
177        let capacity = total_nodes * (*PAGE_SIZE) as usize;
178
179        // Create VMO
180        let vmo: MemoryObject =
181            zx::Vmo::create_with_opts(zx::VmoOptions::RESIZABLE, capacity as u64)
182                .map_err(|_| errno!(ENOMEM))?
183                .into();
184        let vmo = vmo.with_zx_name(b"starnix:tracefs");
185        // Map VMO
186        let addr = vmar_root_self()
187            .map(
188                0,
189                vmo.as_vmo().expect("vmo must exist"),
190                0,
191                capacity,
192                zx::VmarFlags::PERM_READ | zx::VmarFlags::PERM_WRITE,
193            )
194            .map_err(|e| from_status_like_fdio!(e))?;
195        // SAFETY: The address returned by `vmar.map` is valid for `capacity` bytes.
196        let mapping = unsafe { SharedBuffer::new(addr as *mut u8, capacity) };
197
198        // Create nodes
199        let mut nodes = Vec::with_capacity(total_nodes);
200        let base_ptr = addr as *mut u8;
201        for i in 0..total_nodes {
202            // SAFETY: `base_ptr` points to a valid mapping of size `capacity`.
203            // `i * PAGE_SIZE` is within the bounds of this mapping since `i < total_nodes`.
204            // The memory range `[page_ptr, page_ptr + PAGE_SIZE)` is valid and mapped.
205            let page_ptr = unsafe { base_ptr.add(i * (*PAGE_SIZE) as usize) };
206
207            nodes.push(Node {
208                // SAFETY: `page_ptr` is in bounds of the mapped region and valid for `PAGE_SIZE` bytes.
209                page_data: unsafe { SharedBuffer::new(page_ptr, (*PAGE_SIZE) as usize) },
210                next: AtomicUsize::new(0),
211                prev: AtomicUsize::new(0),
212                write_offset: AtomicUsize::new(LocklessRingBuffer::PAGE_HEADER_SIZE),
213                // Initialize active_writers to PAGE_ACTIVE_BIT (1) since the page is active.
214                active_writers: AtomicUsize::new(PAGE_ACTIVE_BIT),
215            });
216        }
217        // Link first `pages - 1` nodes in a circle. The last page is reserved for the reader initialization.
218        let circle_size = pages - 1;
219        for i in 0..circle_size {
220            let next_idx = (i + 1) % circle_size;
221            let prev_idx = (i + circle_size - 1) % circle_size;
222            nodes[i].next.store(Node::make_val(next_idx, FLAG_NORMAL), Ordering::Relaxed);
223            nodes[i].prev.store(Node::make_val(prev_idx, FLAG_NORMAL), Ordering::Relaxed);
224        }
225        // Initialize the reader page, mark head page (node 0).
226        nodes[circle_size - 1].next.store(Node::make_val(0, FLAG_HEADER), Ordering::Relaxed);
227        // Initialize reader page pointers to point to the head page and its predecessor
228        nodes[pages - 1].next.store(Node::make_val(0, FLAG_NORMAL), Ordering::Relaxed);
229        nodes[pages - 1]
230            .prev
231            .store(Node::make_val(circle_size - 1, FLAG_NORMAL), Ordering::Relaxed);
232        // Initialize reader page size to 0
233        nodes[pages - 1].page_data.write_at(8, &0u64.to_le_bytes());
234        let buffer = Self {
235            vmo,
236            mapping,
237            nodes,
238            head_page: AtomicUsize::new(0),
239            tail_page: AtomicUsize::new(0),
240            reader_page: AtomicUsize::new(pages - 1),
241
242            ref_count: AtomicUsize::new(RING_ENABLED_BIT),
243            overwrite,
244            dropped_pages: std::sync::atomic::AtomicU64::new(0),
245            prev_timestamp: std::sync::atomic::AtomicU64::new(
246                zx::BootInstant::get().into_nanos() as u64
247            ),
248            write_event_async_id,
249            reader_active: std::sync::atomic::AtomicBool::new(false),
250            state_mutex: Default::default(),
251        };
252        Ok(buffer)
253    }
254    pub fn dropped_pages(&self) -> u64 {
255        self.dropped_pages.load(Ordering::Relaxed)
256    }
257}
258// Debugging support for detecting live locks. Consider removing once
259// the ring has baked for some time.
260#[derive(Default, Debug)]
261struct YieldTracker {
262    /// Number of times we yielded because the next page was being updated by the reader.
263    update_flag: u64,
264    /// Number of times we yielded because the page to overwrite still had active writers.
265    node_match: u64,
266    /// Number of times we yielded because we failed to lock the head page for overwrite.
267    head_lock: u64,
268}
269impl YieldTracker {
270    fn total(&self) -> u64 {
271        self.update_flag + self.node_match + self.head_lock
272    }
273
274    /// Yields or sleeps progressively based on the total retry count.
275    fn yield_or_sleep(&self) {
276        let total = self.total();
277        if total > 1000 {
278            std::thread::sleep(PROGRESSIVE_SLEEP_DURATION);
279        } else if total > 100 {
280            std::thread::sleep(SPIN_SLEEP_DURATION);
281        } else {
282            std::thread::yield_now();
283        }
284    }
285}
286
287/// Represents a reserved region of the ring buffer that is allocated, but not committed yet.
288pub struct Reservation<'a> {
289    pub offset: usize,
290    pub node_idx: usize,
291    pub size: usize,
292    buffer: &'a LocklessRingBuffer,
293    committed: bool,
294}
295
296impl<'a> Reservation<'a> {
297    pub fn write_at(&self, rel_offset: usize, data: &[u8]) {
298        assert!(rel_offset + data.len() <= self.size, "Write exceeds reservation size");
299        self.buffer.mapping.write_at(self.offset + rel_offset, data);
300    }
301
302    fn release(&mut self) {
303        if self.committed {
304            return;
305        }
306
307        let node_idx = self.node_idx;
308        let node = &self.buffer.nodes[node_idx];
309
310        node.release_writer();
311
312        self.buffer.ref_count.fetch_sub(1, Ordering::Release);
313        self.committed = true;
314    }
315}
316
317impl<'a> Drop for Reservation<'a> {
318    fn drop(&mut self) {
319        if !self.committed {
320            starnix_logging::log_warn!("LocklessRingBuffer: Reservation dropped without commit");
321        }
322        self.release();
323    }
324}
325
326impl<'a> std::fmt::Debug for Reservation<'a> {
327    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
328        f.debug_struct("Reservation")
329            .field("offset", &self.offset)
330            .field("node_idx", &self.node_idx)
331            .field("size", &self.size)
332            .field("committed", &self.committed)
333            .finish()
334    }
335}
336#[derive(Debug, PartialEq, Eq)]
337enum AdvanceResult {
338    /// The tail page has been successfully advanced or transitioned.
339    Advanced,
340    /// Failed to advance due to transient contention. The caller should yield and retry.
341    Yielded,
342    /// A hard error occurred (e.g. buffer is full and overwrite is disabled).
343    Error(Errno),
344}
345
346impl LocklessRingBuffer {
347    /// Attempts to reserve space on the current tail page.
348    ///
349    /// Returns `Ok((offset, now, delta))` if successful, or `Err(())` if the page is full.
350    fn try_reserve_on_page(
351        &self,
352        tail_node: &Node,
353        size: usize,
354    ) -> Result<(usize, zx::BootInstant, zx::Duration<zx::BootTimeline>), ()> {
355        if !self.is_enabled() {
356            log_warn!(
357                "LocklessRingBuffer: canceling try_reserve_on_page because ring is disabled."
358            );
359            return Err(());
360        }
361
362        // Atomically reserve space on the page using a bounded CAS loop.
363        // By verifying that `current_offset + size <= PAGE_SIZE` before executing the CAS,
364        // we guarantee that failed reservations never advance `write_offset`, keeping it exactly
365        // clamped to valid data boundaries.
366        let mut current_offset = tail_node.write_offset.load(Ordering::Acquire);
367        loop {
368            if current_offset + size > (*PAGE_SIZE) as usize {
369                return Err(());
370            }
371            match tail_node.write_offset.compare_exchange_weak(
372                current_offset,
373                current_offset + size,
374                Ordering::AcqRel,
375                Ordering::Acquire,
376            ) {
377                Ok(_) => break,
378                Err(actual) => current_offset = actual,
379            }
380        }
381
382        let now_candidate = zx::BootInstant::get().into_nanos() as u64;
383        let actual_prev = self.prev_timestamp.fetch_max(now_candidate, Ordering::AcqRel);
384        let final_now_nanos = std::cmp::max(now_candidate, actual_prev);
385        let now = zx::BootInstant::from_nanos(final_now_nanos as i64);
386
387        // First write on this page. Set timestamp.
388        let delta_nanos = if current_offset == LocklessRingBuffer::PAGE_HEADER_SIZE {
389            tail_node.page_data.write_at(0, &final_now_nanos.to_le_bytes());
390            0
391        } else {
392            final_now_nanos.saturating_sub(actual_prev)
393        };
394
395        let delta = zx::Duration::from_nanos(delta_nanos as i64);
396
397        Ok((current_offset, now, delta))
398    }
399    /// Advances the tail page to the next page in the ring.
400    ///
401    /// Handles marking the current page as full, moving the head page if overwrite is enabled,
402    /// and updating the global `tail_page` pointer.
403    fn advance_to_next_page(
404        &self,
405        tail_node: &Node,
406        tail_val: usize,
407        yield_tracker: &mut YieldTracker,
408    ) -> AdvanceResult {
409        // Short-circuit if another thread has already advanced the tail page past tail_val.
410        if self.tail_page.load(Ordering::Acquire) != tail_val {
411            return AdvanceResult::Advanced;
412        }
413
414        // If there are no active writers, finalize the page.
415        // We need to check both here and in commit() (and Drop) to avoid race conditions.
416        // Specifically, a writer might finish and see the page is not full yet, while the
417        // thread that failed to fit its data and is moving to the next page hasn't marked it full yet.
418        // By checking in both places, we ensure that at least one thread sees both conditions.
419        if self.is_enabled() {
420            tail_node.finalize();
421        }
422
423        let next_val = tail_node.next.load(Ordering::Acquire);
424        let next_idx = Node::get_index(next_val);
425        let next_flags = Node::get_flags(next_val);
426        if next_flags == FLAG_UPDATE {
427            // Reader is swapping this page. Wait.
428            starnix_logging::log_warn!(
429                "Reservation yielding due to FLAG_UPDATE on node {}",
430                next_idx
431            );
432            yield_tracker.update_flag += 1;
433            yield_tracker.yield_or_sleep();
434            return AdvanceResult::Yielded;
435        } else if next_flags == FLAG_HEADER {
436            if self.overwrite {
437                // Check if any active writer is on the page we want to overwrite.
438                if (self.nodes[next_idx].active_writers.load(Ordering::Acquire)
439                    & ACTIVE_WRITERS_MASK)
440                    > 0
441                {
442                    yield_tracker.node_match += 1;
443                    yield_tracker.yield_or_sleep();
444                    return AdvanceResult::Yielded;
445                }
446
447                starnix_logging::log_warn!(
448                    "LocklessRingBuffer Overwriting page {} (overwriting={})",
449                    next_idx,
450                    self.overwrite
451                );
452                // Try to lock the head page
453                let expected_next = Node::make_val(next_idx, FLAG_HEADER);
454                let locked_next = Node::make_val(next_idx, FLAG_UPDATE);
455                // AcqRel on success:
456                // - Acquire: Ensures we see the up-to-date links of the head page we are about to move.
457                // - Release: Ensures our transition to FLAG_UPDATE (locking) is visible to others.
458                // Relaxed on failure: We fail to acquire the lock and just retry.
459                if tail_node
460                    .next
461                    .compare_exchange_weak(
462                        expected_next,
463                        locked_next,
464                        Ordering::AcqRel,
465                        Ordering::Relaxed,
466                    )
467                    .is_ok()
468                {
469                    // Move the head page.
470                    let head_node = &self.nodes[next_idx];
471                    let head_next_val = head_node.next.load(Ordering::Acquire);
472                    let head_next_idx = Node::get_index(head_next_val);
473
474                    // Set FLAG_HEADER on the new head pointer
475                    let new_head_next = Node::make_val(head_next_idx, FLAG_HEADER);
476                    head_node.next.store(new_head_next, Ordering::Release);
477
478                    // Update global head_page
479                    self.head_page.store(head_next_idx, Ordering::Release);
480
481                    // Unlock the old head pointer (tail_node.next) and make it FLAG_NORMAL
482                    let unlocked_next = Node::make_val(next_idx, FLAG_NORMAL);
483                    tail_node.next.store(unlocked_next, Ordering::Release);
484                    self.dropped_pages.fetch_add(1, Ordering::Relaxed);
485                    fuchsia_trace::async_instant!(
486                        self.write_event_async_id,
487                        "starnix:trace_meta",
488                        "page dropped"
489                    );
490                    return AdvanceResult::Advanced;
491                } else {
492                    // Failed to lock head_page. Retry.
493                    yield_tracker.head_lock += 1;
494                    yield_tracker.yield_or_sleep();
495                    return AdvanceResult::Yielded;
496                }
497            } else {
498                // Buffer is full and overwrite is false.
499                starnix_logging::log_error!("LocklessRingBuffer is full");
500                return AdvanceResult::Error(errno!(ENOSPC));
501            }
502        }
503
504        let next_node = &self.nodes[next_idx];
505        // Only the thread that successfully transitions write_offset from PAGE_SIZE (or full)
506        // to PAGE_HEADER_SIZE is allowed to reset the page and advance the tail pointer.
507        // This prevents the tail-advancing pointer race and slow-thread clobbering.
508        let (should_advance, won_reset) = {
509            let mut current = next_node.write_offset.load(Ordering::Acquire);
510            loop {
511                // Another thread advanced the tail page.
512                if self.tail_page.load(Ordering::Acquire) != tail_val {
513                    break (false, false);
514                }
515                // Edge Case: If `current` is already `PAGE_HEADER_SIZE`, another concurrent writer thread has
516                // already won the reset race on this target page. Since `self.tail_page` hasn't advanced yet,
517                // we must still return `should_advance = true` to cooperate in advancing the global `tail_page`
518                // pointer, while returning `won_reset = false` to avoid redundantly clobbering page metadata.
519                if current == LocklessRingBuffer::PAGE_HEADER_SIZE {
520                    break (true, false);
521                }
522                match next_node.write_offset.compare_exchange_weak(
523                    current,
524                    LocklessRingBuffer::PAGE_HEADER_SIZE,
525                    Ordering::AcqRel,
526                    Ordering::Relaxed,
527                ) {
528                    Ok(_) => break (true, true),
529                    Err(actual) => current = actual,
530                }
531            }
532        };
533
534        if should_advance {
535            if won_reset {
536                // Reset `active_writers` to PAGE_ACTIVE_BIT (1) since the page is active.
537                next_node.active_writers.store(PAGE_ACTIVE_BIT, Ordering::Release);
538                // Clear the size field in the page header to avoid reader seeing stale size.
539                next_node.page_data.write_at(8, &0u64.to_le_bytes());
540            }
541
542            // Try to move tail_page to next_val.
543            // If it fails, it means another thread has already successfully advanced tail_page.
544            if let Err(err) = self.tail_page.compare_exchange(
545                tail_val,
546                next_val,
547                Ordering::AcqRel,
548                Ordering::Relaxed,
549            ) {
550                starnix_logging::log_debug!("Tail page already advanced by another thread: {err}");
551            }
552        }
553        AdvanceResult::Advanced
554    }
555
556    pub fn reserve(
557        &self,
558        size: usize,
559    ) -> Result<(Reservation<'_>, zx::BootInstant, zx::Duration<zx::BootTimeline>), Errno> {
560        // Check that the reservation is non-zero and fits within a page.
561        if size == 0 || size > (*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE {
562            return error!(EINVAL);
563        }
564        // Increment ref_count if enabled.
565        let mut val = self.ref_count.load(Ordering::Acquire);
566        loop {
567            if val & RING_ENABLED_BIT == 0 {
568                return error!(ENOMEM);
569            }
570            // Acquire on success:
571            // - Acquire: Synchronizes with the Release store in `enable()`. This guarantees that if the writer
572            //   observes that the ring is enabled, it will also observe all prior state initialization and memory
573            //   setups made by `enable()` before starting any writes.
574            // - Release: Synchronizes with `disable()`'s Acquire load on `ref_count`, ensuring that if `disable()`
575            //   observes a zero active writer count, it is guaranteed that this writer has either not yet entered
576            //   or has fully completed its operations.
577            // Note: Only the decrement in `release()` needs to be Release to ensure that the data payload writes
578            // are fully visible. Since the writer has not yet written any payload data to the page, there are no
579            // memory writes that need to be published to other threads via Release ordering at this point.
580            match self.ref_count.compare_exchange_weak(
581                val,
582                val + 1,
583                Ordering::Acquire,
584                Ordering::Relaxed,
585            ) {
586                Ok(_) => break,
587                Err(actual) => val = actual,
588            }
589        }
590
591        // Create a scope guard to decrement ref_count on error, panic, or early exit.
592        let ref_count_guard = scopeguard::guard(&self.ref_count, |ref_count| {
593            ref_count.fetch_sub(1, Ordering::Release);
594        });
595
596        // Lock the reservation.
597        let mut yield_tracker = YieldTracker::default();
598        let result = loop {
599            if !self.is_enabled() {
600                log_info!("LocklessRingBuffer: canceling reserve because ring is disabled.");
601                break error!(ENOMEM);
602            }
603            let tail_val = self.tail_page.load(Ordering::Acquire);
604            let tail_idx = Node::get_index(tail_val);
605            let tail_node = &self.nodes[tail_idx];
606
607            // Increment the writer count before trying to reserve space, if it fails decrement it.
608            // We use compare_exchange_weak in a loop to ensure we only increment if the PAGE_ACTIVE_BIT is still set.
609            // Note: The success and failure orderings are Relaxed because the writer has not yet written any payload
610            // data to the page. Therefore, there are no memory writes that need to be published to other threads
611            // via Release ordering.
612            let mut active_incremented = false;
613            let mut old_writers = tail_node.active_writers.load(Ordering::Acquire);
614            loop {
615                if old_writers & PAGE_ACTIVE_BIT == 0 {
616                    break;
617                }
618                match tail_node.active_writers.compare_exchange_weak(
619                    old_writers,
620                    old_writers + 1,
621                    Ordering::Relaxed,
622                    Ordering::Relaxed,
623                ) {
624                    Ok(_) => {
625                        active_incremented = true;
626                        break;
627                    }
628                    Err(actual) => old_writers = actual,
629                }
630            }
631
632            let reservation_result = if active_incremented {
633                self.try_reserve_on_page(tail_node, size)
634            } else {
635                Err(())
636            };
637
638            // Try to claim space on the current page.
639            match reservation_result {
640                Ok((offset, now, delta)) => {
641                    break Ok((
642                        Reservation {
643                            offset: tail_idx * (*PAGE_SIZE) as usize + offset,
644                            node_idx: tail_idx,
645                            size,
646                            buffer: self,
647                            committed: false,
648                        },
649                        now,
650                        delta,
651                    ));
652                }
653                Err(()) => {
654                    // Handle the race where the last successful writer decrements active_writers but
655                    // sees prev_writers > 1 because another concurrent writer (this thread) has already
656                    // incremented it. This thread fails reservation because the page is full, decrements
657                    // active_writers to 0, and gets prev_writers == 1. We must check and finalize here to
658                    // guarantee that full pages are always finalized even if the finalizer thread failed its reserve.
659                    if active_incremented {
660                        self.nodes[tail_idx].release_writer();
661                    }
662                    // Page is full. Advance to the next page.
663                    match self.advance_to_next_page(tail_node, tail_val, &mut yield_tracker) {
664                        AdvanceResult::Advanced | AdvanceResult::Yielded => {}
665                        AdvanceResult::Error(e) => {
666                            break Err(e);
667                        }
668                    }
669                }
670            }
671            // Debug logging for high yield counts.
672            if yield_tracker.total() % 100_000_000 == 0 && yield_tracker.total() > 0 {
673                log_warn!(
674                    "LocklessRingBuffer: spinning in reservation loop.  details = {yield_tracker:?}",
675                );
676            }
677        };
678        // Cleanup and return
679        let total_yields = yield_tracker.total();
680        if total_yields > 500_000 {
681            starnix_logging::log_info!(
682                "Reservation completed with yields: total={}, details={:?}",
683                total_yields,
684                yield_tracker
685            );
686        }
687        match result {
688            Ok((res, now, delta)) => {
689                // Success: Disarm the scope guard so ref_count is NOT decremented.
690                scopeguard::ScopeGuard::into_inner(ref_count_guard);
691                Ok((res, now, delta))
692            }
693            Err(e) => {
694                // Error: The scope guard will automatically decrement ref_count when dropped.
695                Err(e)
696            }
697        }
698    }
699    pub fn commit(&self, mut reservation: Reservation<'_>) {
700        reservation.release();
701    }
702
703    pub fn swap_reader_page(&self) -> Option<usize> {
704        let mut retries = 0;
705        loop {
706            let head_val = self.head_page.load(Ordering::Acquire);
707            let head_idx = Node::get_index(head_val);
708            let tail_val = self.tail_page.load(Ordering::Acquire);
709            let tail_idx = Node::get_index(tail_val);
710            if head_idx == tail_idx {
711                // Ring is empty or we are writing to the only page.
712                return None;
713            }
714            let head_node = &self.nodes[head_idx];
715            // Perform a single Acquire load on `active_writers` to establish happens-before with the finalizing writer.
716            // We check concurrently that the page is finalized (PAGE_ACTIVE_BIT is cleared to 0) and that there are no active writers (count is 0).
717            let active_writers = head_node.active_writers.load(Ordering::Acquire);
718            // We can swap the page if it is no longer active and has no active writers.
719            // We ignore the PAGE_FINALIZED_BIT when determining if the page is ready.
720            if (active_writers & (PAGE_ACTIVE_BIT | ACTIVE_WRITERS_MASK)) != 0 {
721                return None;
722            }
723            let mut size_bytes = [0u8; 8];
724            head_node.page_data.read_at(8, &mut size_bytes);
725            let size = u64::from_le_bytes(size_bytes);
726            if size == 0 {
727                // This can happen even if `head_page != tail_page` (the ring is not empty).
728                // Specifically, if a writer advanced the tail to the next page, the global
729                // `head_page` pointer may have advanced, but the active writer on this new
730                // `head_page` has not yet called `commit()` or finalized the page.
731                // Thus, the data size field at offset 8 in the page header remains 0.
732                return None;
733            }
734            let next_val = head_node.next.load(Ordering::Acquire);
735            let next_idx = Node::get_index(next_val);
736            let prev_val = head_node.prev.load(Ordering::Acquire);
737            let prev_idx = Node::get_index(prev_val);
738            let prev_node = &self.nodes[prev_idx];
739            // Try to lock the head page by setting UPDATE flag on prev.next
740            let expected_next = Node::make_val(head_idx, FLAG_HEADER);
741            let locked_next = Node::make_val(head_idx, FLAG_UPDATE);
742            if prev_node
743                .next
744                .compare_exchange_weak(
745                    expected_next,
746                    locked_next,
747                    Ordering::AcqRel,
748                    Ordering::Relaxed,
749                )
750                .is_ok()
751            {
752                // Locked! Now perform the swap.
753                let reader_idx = self.reader_page.load(Ordering::Acquire);
754                let next_node = &self.nodes[next_idx];
755                // Update pointers to insert reader node.
756                self.nodes[reader_idx]
757                    .prev
758                    .store(Node::make_val(prev_idx, FLAG_NORMAL), Ordering::Relaxed);
759                self.nodes[reader_idx]
760                    .next
761                    .store(Node::make_val(next_idx, FLAG_HEADER), Ordering::Relaxed);
762                next_node.prev.store(Node::make_val(reader_idx, FLAG_NORMAL), Ordering::Relaxed);
763                // Reset write offset of the new head page (old reader page)
764                self.nodes[reader_idx]
765                    .write_offset
766                    .store(LocklessRingBuffer::PAGE_HEADER_SIZE, Ordering::Relaxed);
767                // Reset `active_writers` to PAGE_ACTIVE_BIT before recycling.
768                self.nodes[reader_idx].active_writers.store(PAGE_ACTIVE_BIT, Ordering::Relaxed);
769                // Clear the size field in the page header to avoid reader seeing stale size.
770                self.nodes[reader_idx].page_data.write_at(8, &0u64.to_le_bytes());
771                // Unlock prev_node and set its next to reader_idx as NORMAL
772                let unlocked_val = Node::make_val(reader_idx, FLAG_NORMAL);
773                prev_node.next.store(unlocked_val, Ordering::Release);
774                // Update global pointers
775                self.head_page.store(next_idx, Ordering::Release);
776                self.reader_page.store(head_idx, Ordering::Release);
777                // Return the old head page index (now reader page)
778                return Some(head_idx);
779            }
780            // Failed to lock, retry.
781            starnix_logging::log_warn!("swap_reader_page failed to lock, yielding");
782            retries += 1;
783            if retries >= 100_000 {
784                starnix_logging::log_error!("LocklessRingBuffer: HUNG in swap_reader_page loop");
785                return None;
786            }
787            // Failed to acquire the lock for swapping. Spin and retry.
788            // Note: The log message above says "yielding" but we actually spin here
789            // to avoid context switch overhead for this short wait.
790            std::hint::spin_loop();
791        }
792    }
793    pub fn read(&self, buf: &mut dyn OutputBuffer) -> Result<usize, Errno> {
794        // Validate that concurrent reads on a single ring buffer do not happen.
795        // By design in both ftrace/tracefs and Starnix VFS, only a single reader is active at a time.
796        if self.reader_active.swap(true, Ordering::AcqRel) {
797            starnix_logging::log_error!(
798                "LocklessRingBuffer: concurrent reads detected! Concurrent reads are not supported by design."
799            );
800            return error!(EBUSY);
801        }
802
803        let _lock_guard = scopeguard::guard(&self.reader_active, |reader_active| {
804            reader_active.store(false, Ordering::Release);
805        });
806
807        // Increment ref_count if enabled to prevent disable() from shrinking the VMO
808        // while we are reading from it. We use a compare_exchange loop instead of
809        // fetch_add to avoid modifying the counter if the ring is already disabled,
810        // which prevents live-locking disable().
811        let mut val = self.ref_count.load(Ordering::Acquire);
812        loop {
813            if val & RING_ENABLED_BIT == 0 {
814                return error!(EAGAIN);
815            }
816            match self.ref_count.compare_exchange_weak(
817                val,
818                val + 1,
819                Ordering::Acquire,
820                Ordering::Relaxed,
821            ) {
822                Ok(_) => break,
823                Err(actual) => val = actual,
824            }
825        }
826
827        let _guard = scopeguard::guard(&self.ref_count, |ref_count| {
828            ref_count.fetch_sub(1, Ordering::Release);
829        });
830
831        if buf.available() < (*PAGE_SIZE) as usize {
832            starnix_logging::log_error!(
833                "Buffer is too small: {} bytes (needs {} bytes)",
834                buf.available(),
835                (*PAGE_SIZE) as usize
836            );
837            return error!(EINVAL);
838        }
839        // TODO(https://fxbug.dev/505532201): Consider handling buffers larger than a page.
840        // this currently only works for a single page.
841        if let Some(idx) = self.swap_reader_page() {
842            let node = &self.nodes[idx];
843            let mut offset = 0;
844            let page_size = (*PAGE_SIZE) as usize;
845            let bytes_written = buf.write_each(&mut move |segment| {
846                let available = page_size - offset;
847                if available == 0 {
848                    return Ok(0);
849                }
850                let size = std::cmp::min(segment.len(), available);
851                // SAFETY: We are writing to this segment, so it's safe to treat it as initialized after we write.
852                // We cast it to `&mut [u8]` to pass to `read_at`.
853                let segment_mut_u8 = unsafe {
854                    std::slice::from_raw_parts_mut(segment.as_mut_ptr() as *mut u8, size)
855                };
856                node.page_data.read_at(offset, segment_mut_u8);
857                offset += size;
858                Ok(size)
859            })?;
860            return Ok(bytes_written);
861        }
862        error!(EAGAIN)
863    }
864
865    #[cfg(test)]
866    pub(crate) fn size_bytes(&self) -> usize {
867        self.nodes.len() * (*PAGE_SIZE) as usize
868    }
869
870    pub fn is_enabled(&self) -> bool {
871        self.ref_count.load(Ordering::Acquire) & RING_ENABLED_BIT != 0
872    }
873    pub fn disable(&self) -> Result<u64, Errno> {
874        let _lock = self.state_mutex.lock();
875        // Clear the enabled bit.
876        self.ref_count.fetch_and(!RING_ENABLED_BIT, Ordering::AcqRel);
877        let mut yield_count: u64 = 0;
878        loop {
879            let active_writers = self.ref_count.load(Ordering::Acquire);
880            if active_writers == 0 {
881                break;
882            }
883            yield_count += 1;
884            if yield_count % 1_000_000 == 0 {
885                log_error!(
886                    "LocklessRingBuffer: disable is waiting for {active_writers} active writers",
887                );
888            }
889            // We use yield_now() here as a compromise:
890            // - spin_wait would waste too much CPU if the writer is descheduled.
891            // - sleep adds too much latency for what should be a very short wait (copying data).
892            // Yielding gives other threads a chance to run while keeping latency low.
893            std::thread::yield_now();
894        }
895        if yield_count > 500_000 {
896            log_warn!("LocklessRingBuffer disable took {} yields", yield_count);
897        }
898        if let Err(e) = self.vmo.set_size(0) {
899            starnix_logging::log_error!(
900                "LocklessRingBuffer disable failed to set_size(0): {:?}",
901                e
902            );
903            return Err(from_status_like_fdio!(e));
904        }
905        let dropped = self.dropped_pages.load(Ordering::Relaxed);
906        Ok(dropped)
907    }
908    pub fn enable(&self) -> Result<zx::BootInstant, Errno> {
909        let _lock = self.state_mutex.lock();
910        let initial_pages = self.nodes.len() - 1;
911        let capacity = self.nodes.len() * (*PAGE_SIZE) as usize;
912        if let Err(e) = self.vmo.set_size(capacity as u64) {
913            starnix_logging::log_error!("LocklessRingBuffer enable failed to set_size: {:?}", e);
914            return Err(from_status_like_fdio!(e));
915        }
916        let now = zx::BootInstant::get();
917        // Reset state
918        self.head_page.store(0, Ordering::Release);
919        self.tail_page.store(0, Ordering::Release);
920
921        self.reader_page.store(initial_pages, Ordering::Release);
922
923        self.dropped_pages.store(0, Ordering::Release);
924        self.prev_timestamp.store(now.into_nanos() as u64, Ordering::Release);
925
926        for i in 0..self.nodes.len() {
927            self.nodes[i]
928                .write_offset
929                .store(LocklessRingBuffer::PAGE_HEADER_SIZE, Ordering::Release);
930            self.nodes[i].active_writers.store(PAGE_ACTIVE_BIT, Ordering::Release);
931        }
932        for i in 0..initial_pages {
933            let next_idx = (i + 1) % initial_pages;
934            let prev_idx = (i + initial_pages - 1) % initial_pages;
935            self.nodes[i].next.store(Node::make_val(next_idx, FLAG_NORMAL), Ordering::Relaxed);
936            self.nodes[i].prev.store(Node::make_val(prev_idx, FLAG_NORMAL), Ordering::Relaxed);
937        }
938        self.nodes[initial_pages - 1].next.store(Node::make_val(0, FLAG_HEADER), Ordering::Relaxed);
939
940        // Initialize reader page size to 0
941        self.nodes[initial_pages].page_data.write_at(8, &0u64.to_le_bytes());
942        // Initialize first page.
943        let _ = self.vmo.as_vmo().expect("vmo must exist").op_range(zx::VmoOp::ZERO, 0, *PAGE_SIZE);
944        let nanos = now.into_nanos() as u64;
945        self.nodes[0].page_data.write_at(0, &nanos.to_le_bytes());
946
947        // Set enabled bit LAST to ensure state is fully reset before writes start.
948        self.ref_count.store(RING_ENABLED_BIT, Ordering::Release);
949        Ok(now)
950    }
951}
952impl Drop for LocklessRingBuffer {
953    fn drop(&mut self) {
954        let (ptr, len) = self.mapping.as_ptr_len();
955        // SAFETY: The mapping was created in `new` and is valid for this lifetime.
956        // We are freeing it now because the object is being destroyed.
957        unsafe {
958            let _ = vmar_root_self().unmap(ptr as usize, len);
959        }
960    }
961}
962
963#[cfg(test)]
964mod tests {
965    use super::*;
966    use crate::vfs::buffers::VecOutputBuffer;
967    use crate::vfs::{Buffer, OutputBufferCallback, PeekBufferSegmentsCallback};
968    use std::sync::Arc;
969    #[repr(C)]
970    #[derive(Clone, Copy, Debug, PartialEq, Eq)]
971    struct TestMessage {
972        thread_index: u32,
973        timestamp_nanos: u64,
974        delta: u64,
975        data: [u8; 12],
976    }
977    impl TestMessage {
978        pub const SIZE: usize = 32;
979        fn to_bytes(&self) -> [u8; TestMessage::SIZE] {
980            let mut bytes = [0u8; TestMessage::SIZE];
981            bytes[0..4].copy_from_slice(&self.thread_index.to_le_bytes());
982            bytes[4..12].copy_from_slice(&self.timestamp_nanos.to_le_bytes());
983            bytes[12..20].copy_from_slice(&self.delta.to_le_bytes());
984            bytes[20..32].copy_from_slice(&self.data);
985            bytes
986        }
987        fn from_bytes(bytes: &[u8]) -> Self {
988            let mut thread_index = [0u8; 4];
989            thread_index.copy_from_slice(&bytes[0..4]);
990            let mut timestamp_nanos = [0u8; 8];
991            timestamp_nanos.copy_from_slice(&bytes[4..12]);
992            let mut delta = [0u8; 8];
993            delta.copy_from_slice(&bytes[12..20]);
994            let mut data = [0u8; 12];
995            data.copy_from_slice(&bytes[20..32]);
996            Self {
997                thread_index: u32::from_le_bytes(thread_index),
998                timestamp_nanos: u64::from_le_bytes(timestamp_nanos),
999                delta: u64::from_le_bytes(delta),
1000                data,
1001            }
1002        }
1003    }
1004    #[test]
1005    fn test_init() {
1006        let buffer =
1007            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1008                .unwrap();
1009        assert_eq!(buffer.size_bytes(), 3 * (*PAGE_SIZE) as usize);
1010    }
1011    #[test]
1012    fn test_reserve() {
1013        let buffer =
1014            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1015                .unwrap();
1016        let res = buffer.reserve(100).unwrap();
1017        assert_eq!(res.0.size, 100);
1018        assert_eq!(res.0.offset, LocklessRingBuffer::PAGE_HEADER_SIZE);
1019        assert_eq!(res.0.node_idx, 0);
1020    }
1021    #[test]
1022    fn test_commit() {
1023        let buffer =
1024            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1025                .unwrap();
1026        let res = buffer.reserve(100).unwrap();
1027        buffer.commit(res.0);
1028        assert_eq!(buffer.nodes[0].active_writers.load(Ordering::Relaxed), PAGE_ACTIVE_BIT);
1029    }
1030    #[test]
1031    fn test_swap_reader_page_empty() {
1032        let buffer =
1033            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1034                .unwrap();
1035        assert_eq!(buffer.swap_reader_page(), None);
1036    }
1037    #[test]
1038    fn test_swap_reader_page_success() {
1039        let buffer =
1040            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1041                .unwrap();
1042        // Fill node 0
1043        let res1 =
1044            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1045        buffer.commit(res1.0);
1046        // Next reserve should go to node 1
1047        let res2 = buffer.reserve(100).unwrap();
1048        assert_eq!(res2.0.node_idx, 1);
1049        buffer.commit(res2.0);
1050        // Page 0 has size > 0 and tail advanced to Page 1.
1051        // So swap_reader_page should succeed!
1052        let old_head = buffer.swap_reader_page();
1053        assert_eq!(old_head, Some(0));
1054        // And head_page should now have advanced to index 1.
1055        assert_eq!(buffer.head_page.load(Ordering::Relaxed), 1);
1056    }
1057    #[test]
1058    fn test_concurrent_reserve() {
1059        use std::sync::Arc;
1060        let buffer = Arc::new(
1061            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1062                .unwrap(),
1063        );
1064        let mut handles = vec![];
1065        for _ in 0..5 {
1066            let buffer_clone = Arc::clone(&buffer);
1067            let handle = std::thread::spawn(move || {
1068                for _ in 0..20 {
1069                    if let Ok(res) = buffer_clone.reserve(10) {
1070                        buffer_clone.commit(res.0);
1071                    }
1072                }
1073            });
1074            handles.push(handle);
1075        }
1076        for handle in handles {
1077            handle.join().unwrap();
1078        }
1079        for i in 0..buffer.nodes.len() {
1080            assert_eq!(buffer.nodes[i].active_writers.load(Ordering::Relaxed), PAGE_ACTIVE_BIT);
1081        }
1082    }
1083    #[test]
1084    fn test_reserve_moves_to_next_page() {
1085        let buffer =
1086            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1087                .unwrap();
1088        // Reserve most of the page, leaving only 50 bytes
1089        let res1 = buffer
1090            .reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE - 50)
1091            .unwrap();
1092        buffer.commit(res1.0);
1093        // Try to reserve 100 bytes. It shouldn't fit on node 0.
1094        let res2 = buffer.reserve(100).unwrap();
1095        // It should be on node 1.
1096        assert_eq!(res2.0.node_idx, 1);
1097        assert_eq!(res2.0.offset, (*PAGE_SIZE) as usize + LocklessRingBuffer::PAGE_HEADER_SIZE);
1098    }
1099    #[test]
1100    fn test_reserve_overwrite() {
1101        let buffer =
1102            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1103                .unwrap();
1104        // Fill node 0 and commit
1105        let res1 =
1106            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1107        buffer.commit(res1.0);
1108        // Fill node 1 but do NOT commit
1109        let _res2 =
1110            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1111        // Now tail is at node 1. Next reserve wants to move to node 0.
1112        // But commit is still at node 0.
1113        // So it should trigger overwrite.
1114        let res3 = buffer.reserve(100).unwrap();
1115        // It should succeed and be on node 0.
1116        assert_eq!(res3.0.node_idx, 0);
1117        assert_eq!(buffer.dropped_pages(), 1);
1118        // And head_page should have advanced to 1.
1119        assert_eq!(buffer.head_page.load(Ordering::Relaxed), 1);
1120    }
1121    #[test]
1122    fn test_read() {
1123        let buffer =
1124            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1125                .unwrap();
1126        let res1 =
1127            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1128        let data = vec![1u8; (*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE];
1129        res1.0.write_at(0, &data);
1130        buffer.commit(res1.0);
1131        // Advance commit page to node 1
1132        let res2 = buffer.reserve(100).unwrap();
1133        buffer.commit(res2.0);
1134        let mut dest = VecOutputBuffer::new((*PAGE_SIZE) as usize);
1135        let result = buffer.read(&mut dest);
1136        assert!(result.is_ok());
1137        assert_eq!(result.unwrap(), (*PAGE_SIZE) as usize);
1138        assert_eq!(&dest.data()[LocklessRingBuffer::PAGE_HEADER_SIZE..], &data[..]);
1139    }
1140    #[test]
1141    fn test_enable_disable() {
1142        let buffer =
1143            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1144                .unwrap();
1145        assert!(buffer.disable().is_ok());
1146        assert_eq!(buffer.ref_count.load(Ordering::Relaxed) & RING_ENABLED_BIT, 0);
1147        let res = buffer.reserve(100);
1148        assert_eq!(res.unwrap_err(), starnix_uapi::errno!(ENOMEM));
1149        assert!(buffer.enable().is_ok());
1150        assert_eq!(buffer.size_bytes(), 3 * (*PAGE_SIZE) as usize);
1151        let res = buffer.reserve(100);
1152        assert!(res.is_ok());
1153    }
1154    // Spawns a reader thread that periodically drains the lockless ring buffer.
1155    fn start_reader_thread(
1156        buffer_reader: Arc<LocklessRingBuffer>,
1157        writers_done_reader: Arc<std::sync::atomic::AtomicBool>,
1158    ) -> std::thread::JoinHandle<Vec<TestMessage>> {
1159        start_gated_reader_thread(buffer_reader, writers_done_reader, false)
1160    }
1161
1162    // Spawns a reader thread that optionally yields until at least one page has been
1163    // dropped (overwritten) by writers before draining.
1164    fn start_gated_reader_thread(
1165        buffer_reader: Arc<LocklessRingBuffer>,
1166        writers_done_reader: Arc<std::sync::atomic::AtomicBool>,
1167        wait_for_dropped_pages: bool,
1168    ) -> std::thread::JoinHandle<Vec<TestMessage>> {
1169        std::thread::spawn(move || {
1170            if wait_for_dropped_pages {
1171                while buffer_reader.dropped_pages() == 0
1172                    && !writers_done_reader.load(Ordering::Acquire)
1173                {
1174                    std::thread::yield_now();
1175                }
1176            }
1177            let mut all_messages = Vec::new();
1178            let mut dest = VecOutputBuffer::new((*PAGE_SIZE) as usize);
1179            let mut consecutive_eagain = 0;
1180            loop {
1181                dest.reset();
1182                match buffer_reader.read(&mut dest) {
1183                    Ok(bytes_read) => {
1184                        consecutive_eagain = 0;
1185                        assert_eq!(bytes_read, (*PAGE_SIZE) as usize);
1186                        let mut header_ts_bytes = [0u8; 8];
1187                        header_ts_bytes.copy_from_slice(&dest.data()[0..8]);
1188                        let header_timestamp = u64::from_le_bytes(header_ts_bytes);
1189                        let mut size_bytes = [0u8; 8];
1190                        size_bytes.copy_from_slice(&dest.data()[8..16]);
1191                        let data_size = u64::from_le_bytes(size_bytes) as usize;
1192                        let max_offset = LocklessRingBuffer::PAGE_HEADER_SIZE + data_size;
1193                        let mut offset = LocklessRingBuffer::PAGE_HEADER_SIZE;
1194                        let mut first_msg = true;
1195                        while offset + TestMessage::SIZE <= max_offset {
1196                            let msg_bytes = &dest.data()[offset..offset + TestMessage::SIZE];
1197                            let msg = TestMessage::from_bytes(msg_bytes);
1198                            if first_msg {
1199                                if header_timestamp != msg.timestamp_nanos {
1200                                    println!(
1201                                        "HEADER TIMESTAMP MISMATCH: header={}, msg={}",
1202                                        header_timestamp, msg.timestamp_nanos
1203                                    );
1204                                }
1205                                first_msg = false;
1206                            }
1207                            all_messages.push(msg);
1208                            offset += TestMessage::SIZE;
1209                        }
1210                    }
1211                    Err(e) if e == starnix_uapi::errno!(EAGAIN) => {
1212                        if writers_done_reader.load(Ordering::Acquire) {
1213                            let head_val = buffer_reader.head_page.load(Ordering::Relaxed);
1214                            let tail_val = buffer_reader.tail_page.load(Ordering::Relaxed);
1215                            // If the writers are done and the reader has caught up to the tail page
1216                            // index, all readable pages have been drained, so we break the loop.
1217                            if Node::get_index(head_val) == Node::get_index(tail_val) {
1218                                break;
1219                            }
1220                            consecutive_eagain += 1;
1221                            assert!(
1222                                consecutive_eagain < 500,
1223                                "LocklessRingBuffer reader stuck: consecutive EAGAIN limit exceeded (5s) after writers finished. head_page={}, tail_page={}, reader_page={}",
1224                                Node::get_index(head_val),
1225                                Node::get_index(tail_val),
1226                                buffer_reader.reader_page.load(Ordering::Relaxed)
1227                            );
1228                        }
1229                        std::thread::sleep(std::time::Duration::from_millis(10));
1230                    }
1231                    Err(e) => panic!("Unexpected error from read: {:?}", e),
1232                }
1233            }
1234            all_messages
1235        })
1236    }
1237    fn check_all_message_data(all_messages: &[TestMessage], num_threads: u32) {
1238        let mut prev_timestamp = 0;
1239        let mut thread_counts = vec![0; num_threads as usize];
1240        let mut out_of_order = 0;
1241        let mut corrupted = 0;
1242        for msg in all_messages {
1243            if msg.timestamp_nanos < prev_timestamp {
1244                println!("OUT OF ORDER: prev_timestamp={}, current={:?}", prev_timestamp, msg);
1245                out_of_order += 1;
1246            }
1247            prev_timestamp = msg.timestamp_nanos;
1248            if &msg.data != b"Event data\0\0" || msg.thread_index >= num_threads {
1249                println!("CORRUPTED: msg={:?}", msg);
1250                corrupted += 1;
1251            } else {
1252                thread_counts[msg.thread_index as usize] += 1;
1253            }
1254        }
1255        println!(
1256            "TEST_RESULT: Read {} messages, {} out of order, {} corrupted. Thread counts: {:?}",
1257            all_messages.len(),
1258            out_of_order,
1259            corrupted,
1260            thread_counts
1261        );
1262        // We do not assert out_of_order is 0 because wait-free reservation can cause slight reordering of timestamps.
1263        assert_eq!(corrupted, 0, "Found corrupted messages");
1264    }
1265    #[test]
1266    fn test_concurrent_read_write_4() {
1267        let num_threads = 4;
1268        let msgs_per_thread = 64;
1269        // 4 threads * 64 msgs * 32 bytes = 8192 bytes.
1270        // Each page can hold (4096 - 16) / 32 = 127 msgs (with 16 bytes unused at the end).
1271        // 256 msgs will take 2 full pages (254 msgs) + 2 msgs on a 3rd page.
1272        let buffer = Arc::new(
1273            LocklessRingBuffer::new(4 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1274                .unwrap(),
1275        );
1276        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1277        let mut handles = vec![];
1278        // Spawn reader
1279        let buffer_reader = Arc::clone(&buffer);
1280        let writers_done_reader = Arc::clone(&writers_done);
1281        let reader_handle = start_reader_thread(buffer_reader, writers_done_reader);
1282        // Spawn writers
1283        for thread_index in 0..num_threads {
1284            let buffer_clone = Arc::clone(&buffer);
1285            let handle = std::thread::spawn(move || {
1286                std::thread::sleep(std::time::Duration::from_millis(10 + thread_index as u64));
1287                for _ in 0..msgs_per_thread {
1288                    // Reserve exactly TestMessage::SIZE bytes
1289                    let (res, now, delta) = buffer_clone.reserve(TestMessage::SIZE).unwrap();
1290                    let msg = TestMessage {
1291                        thread_index,
1292                        timestamp_nanos: now.into_nanos() as u64,
1293                        delta: delta.into_nanos() as u64,
1294                        data: *b"Event data\0\0",
1295                    };
1296                    res.write_at(0, &msg.to_bytes());
1297                    buffer_clone.commit(res);
1298                    std::thread::sleep(std::time::Duration::from_nanos(thread_index as u64));
1299                }
1300            });
1301            handles.push(handle);
1302        }
1303        for handle in handles {
1304            handle.join().unwrap();
1305        }
1306        writers_done.store(true, Ordering::Release);
1307        let all_messages = reader_handle.join().unwrap();
1308        // Check the messages
1309        check_all_message_data(&all_messages, num_threads);
1310        // Trailing messages might remain on the final unfinalized tail page. We expect between 250 and 256 messages to be successfully read.
1311        assert!(
1312            all_messages.len() >= 250 && all_messages.len() <= 256,
1313            "Expected between 250 and 256 messages, got {}",
1314            all_messages.len()
1315        );
1316    }
1317    #[test]
1318    fn test_concurrent_read_write_1_thread() {
1319        let num_threads = 1;
1320        let msgs_per_thread = 256;
1321        // 1 thread * 256 msgs = 256 msgs.
1322        // Each page holds 127 msgs. 256 msgs requires 3 data pages.
1323        // We use 5 pages total to provide enough capacity.
1324        let buffer = Arc::new(
1325            LocklessRingBuffer::new(5 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1326                .unwrap(),
1327        );
1328        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1329        let mut handles = vec![];
1330        // Spawn reader
1331        let buffer_reader = Arc::clone(&buffer);
1332        let writers_done_reader = Arc::clone(&writers_done);
1333        let reader_handle = start_reader_thread(buffer_reader, writers_done_reader);
1334        // Spawn writers
1335        for thread_index in 0..num_threads {
1336            let buffer_clone = Arc::clone(&buffer);
1337            let handle = std::thread::spawn(move || {
1338                std::thread::sleep(std::time::Duration::from_millis(10 + thread_index as u64));
1339                for _ in 0..msgs_per_thread {
1340                    let (res, now, delta) = buffer_clone.reserve(TestMessage::SIZE).unwrap();
1341                    let msg = TestMessage {
1342                        thread_index,
1343                        timestamp_nanos: now.into_nanos() as u64,
1344                        delta: delta.into_nanos() as u64,
1345                        data: *b"Event data\0\0",
1346                    };
1347                    res.write_at(0, &msg.to_bytes());
1348                    buffer_clone.commit(res);
1349                    std::thread::sleep(std::time::Duration::from_nanos(thread_index as u64));
1350                }
1351            });
1352            handles.push(handle);
1353        }
1354        for handle in handles {
1355            handle.join().unwrap();
1356        }
1357        writers_done.store(true, Ordering::Release);
1358        let all_messages = reader_handle.join().unwrap();
1359        check_all_message_data(&all_messages, num_threads);
1360        assert_eq!(
1361            all_messages.len(),
1362            254,
1363            "Expected exactly 254 messages, got {}",
1364            all_messages.len()
1365        );
1366    }
1367    #[test]
1368    fn test_concurrent_read_write_8_threads() {
1369        let num_threads = 8;
1370        let msgs_per_thread = 128;
1371        // 4 threads * 64 msgs * 32 bytes = 8192 bytes.
1372        // 8 threads * 128 msgs = 1024 msgs.
1373        // Each page holds 127 msgs. 1024 msgs requires 9 data pages.
1374        // We use 12 pages total to provide enough capacity so no pages are overwritten.
1375        let buffer = Arc::new(
1376            LocklessRingBuffer::new(12 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1377                .unwrap(),
1378        );
1379        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1380        let _writers_done_guard = scopeguard::guard(Arc::clone(&writers_done), |done| {
1381            done.store(true, Ordering::Release);
1382        });
1383        let mut handles = vec![];
1384        // Spawn reader
1385        let buffer_reader = Arc::clone(&buffer);
1386        let writers_done_reader = Arc::clone(&writers_done);
1387        let reader_handle = start_reader_thread(buffer_reader, writers_done_reader);
1388        // Spawn writers
1389        for thread_index in 0..num_threads {
1390            let buffer_clone = Arc::clone(&buffer);
1391            let handle = std::thread::spawn(move || {
1392                std::thread::sleep(std::time::Duration::from_millis(10 + thread_index as u64));
1393                for _ in 0..msgs_per_thread {
1394                    // Reserve exactly TestMessage::SIZE bytes
1395                    let mut retries = 0;
1396                    let (res, now, delta) = loop {
1397                        match buffer_clone.reserve(TestMessage::SIZE) {
1398                            Ok(r) => break r,
1399                            Err(e) if e == starnix_uapi::errno!(ENOSPC) => {
1400                                retries += 1;
1401                                assert!(
1402                                    retries < 100,
1403                                    "LocklessRingBuffer: ENOSPC transient limit exceeded (100ms). Reader may have hung."
1404                                );
1405                                // Under high contention, writers can transiently catch up to the head page.
1406                                // This is expected and will clear up once the reader thread gets scheduled
1407                                // and swaps a page.
1408                                std::thread::sleep(std::time::Duration::from_millis(1));
1409                            }
1410                            Err(e) => panic!("Unexpected error: {:?}", e),
1411                        }
1412                    };
1413                    let msg = TestMessage {
1414                        thread_index,
1415                        timestamp_nanos: now.into_nanos() as u64,
1416                        delta: delta.into_nanos() as u64,
1417                        data: *b"Event data\0\0",
1418                    };
1419                    res.write_at(0, &msg.to_bytes());
1420                    buffer_clone.commit(res);
1421                    std::thread::sleep(std::time::Duration::from_nanos(thread_index as u64));
1422                }
1423            });
1424            handles.push(handle);
1425        }
1426        for handle in handles {
1427            handle.join().unwrap();
1428        }
1429
1430        writers_done.store(true, Ordering::Release);
1431        let all_messages = reader_handle.join().unwrap();
1432        // Check the messages
1433        check_all_message_data(&all_messages, num_threads);
1434        // The last page in the circle (current tail page) cannot be swapped by the reader
1435        // (to prevent reading uncommitted/partial data). The number of messages left sitting on
1436        // the tail page when the writers finish depends entirely on concurrent thread scheduling
1437        // variance under high contention. We expect between 1008 and 1024 messages to be read
1438        // (leaving 0 to 16 messages unread on the final tail page).
1439        assert!(
1440            all_messages.len() >= 1008 && all_messages.len() <= 1024,
1441            "Expected between 1008 and 1024 messages, got {}",
1442            all_messages.len()
1443        );
1444    }
1445    #[test]
1446    fn test_disable_waits_for_ref_count() {
1447        use std::sync::Arc;
1448        use std::sync::atomic::{AtomicBool, Ordering};
1449        let buffer = Arc::new(
1450            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1451                .unwrap(),
1452        );
1453        // Thread 1: Reserves space, keeping ref_count = 1 (plus enabled bit).
1454        let res = buffer.reserve(100).unwrap();
1455        let disable_finished = Arc::new(AtomicBool::new(false));
1456        let disable_finished_clone = Arc::clone(&disable_finished);
1457        let buffer_clone = Arc::clone(&buffer);
1458        // Thread 2: Calls disable. This should block because ref_count has active writers.
1459        let handle = std::thread::spawn(move || {
1460            buffer_clone.disable().unwrap();
1461            disable_finished_clone.store(true, Ordering::Relaxed);
1462        });
1463        // Wait for Thread 2 to start and set size_bytes to 0
1464        while buffer.ref_count.load(Ordering::Acquire) & RING_ENABLED_BIT != 0 {
1465            std::hint::spin_loop();
1466        }
1467        // Verify that disable has NOT finished
1468        assert!(!disable_finished.load(Ordering::Relaxed));
1469        // Verify that new reserves fail with ENOMEM because disable() cleared RING_ENABLED_BIT
1470        let res2 = buffer.reserve(100);
1471        assert_eq!(res2.unwrap_err(), starnix_uapi::errno!(ENOMEM));
1472        // Now write data and commit the first reservation.
1473        // If the VMO had been shrunk, write_data would panic/fault here.
1474        let data = vec![42u8; 100];
1475        res.0.write_at(0, &data);
1476        buffer.commit(res.0);
1477        // Wait for Thread 2 to finish
1478        handle.join().unwrap();
1479        // Verify disable finished
1480        assert!(disable_finished.load(Ordering::Relaxed));
1481    }
1482    #[test]
1483    fn test_reservation_drop_cancels_writer_count() {
1484        let buffer =
1485            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1486                .unwrap();
1487        // 1. Reserve space. ref_count becomes 1 (plus enabled bit).
1488        let res = buffer.reserve(100).unwrap();
1489        assert_eq!(buffer.ref_count.load(Ordering::Relaxed), RING_ENABLED_BIT | 1);
1490        // 2. Drop the reservation without commit.
1491        std::mem::drop(res.0);
1492        // 3. Verify ref_count becomes 0 (plus enabled bit).
1493        assert_eq!(buffer.ref_count.load(Ordering::Relaxed), RING_ENABLED_BIT);
1494        // 4. Verify that a new reservation can be committed (proving the dropped one didn't block).
1495        let res2 = buffer.reserve(100).unwrap();
1496        buffer.commit(res2.0);
1497    }
1498    #[test]
1499    fn test_overwrite_and_read_4_pages() {
1500        let num_threads = 1;
1501        // Use 6 pages total (5 data pages) to allow exactly 4 readable pages
1502        // (1 page is the active commit page).
1503        let buffer = Arc::new(
1504            LocklessRingBuffer::new(6 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1505                .unwrap(),
1506        );
1507        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1508        // 5 data pages capacity + 2 overwritten pages = 7 pages to write.
1509        // Each page holds exactly 127 messages.
1510        let msgs_per_thread = 7 * 127;
1511        let buffer_clone = Arc::clone(&buffer);
1512        let thread_index = 0;
1513        let writer_handle = std::thread::spawn(move || {
1514            for _ in 0..msgs_per_thread {
1515                let (res, now, delta) = buffer_clone.reserve(TestMessage::SIZE).unwrap();
1516                let msg = TestMessage {
1517                    thread_index,
1518                    timestamp_nanos: now.into_nanos() as u64,
1519                    delta: delta.into_nanos() as u64,
1520                    data: *b"Event data\0\0",
1521                };
1522                res.write_at(0, &msg.to_bytes());
1523                buffer_clone.commit(res);
1524            }
1525        });
1526        writer_handle.join().unwrap();
1527        assert_eq!(buffer.dropped_pages(), 2, "Expected 2 pages to be dropped");
1528        writers_done.store(true, Ordering::Release);
1529        let buffer_reader = Arc::clone(&buffer);
1530        let writers_done_reader = Arc::clone(&writers_done);
1531        let reader_handle = start_reader_thread(buffer_reader, writers_done_reader);
1532        let all_messages = reader_handle.join().unwrap();
1533        check_all_message_data(&all_messages, num_threads);
1534        assert_eq!(all_messages.len(), 4 * 127, "Expected exactly 4 pages of messages");
1535    }
1536    #[test]
1537    fn test_reserve_full_producer_consumer() {
1538        let buffer =
1539            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1540                .unwrap();
1541        // Fill node 0
1542        let res1 =
1543            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1544        buffer.commit(res1.0);
1545        // Fill node 1
1546        let res2 =
1547            buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1548        buffer.commit(res2.0);
1549        // Buffer is full (2 data pages used, head at 0, commit at 1 but next is head).
1550        // Next reserve should fail with ENOSPC because overwrite is false.
1551        let res3 = buffer.reserve(100);
1552        assert_eq!(res3.unwrap_err(), starnix_uapi::errno!(ENOSPC));
1553    }
1554
1555    #[test]
1556    fn test_reserve_zero_size() {
1557        let buffer =
1558            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1559                .unwrap();
1560        let res = buffer.reserve(0);
1561        assert_eq!(res.unwrap_err(), starnix_uapi::errno!(EINVAL));
1562    }
1563    #[test]
1564    fn test_concurrent_overwrite_stability() {
1565        let num_threads = 4;
1566        let msgs_per_thread = 256;
1567        // Small buffer to force overwrites
1568        let buffer = Arc::new(
1569            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1570                .unwrap(),
1571        );
1572        // Synchronously fill the active circular buffer (2 data pages) to capacity.
1573        // The reader thread is gated below so it does not drain the buffer until
1574        // the writer threads wrap around to Node 0 and trigger a dropped page.
1575        let max_payload = (*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE;
1576        let num_msgs = max_payload / TestMessage::SIZE;
1577        for _ in 0..(2 * num_msgs) {
1578            let (res, now, delta) = buffer.reserve(TestMessage::SIZE).unwrap();
1579            let msg = TestMessage {
1580                thread_index: 0,
1581                timestamp_nanos: now.into_nanos() as u64,
1582                delta: delta.into_nanos() as u64,
1583                data: *b"Event data\0\0",
1584            };
1585            res.write_at(0, &msg.to_bytes());
1586            buffer.commit(res);
1587        }
1588
1589        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1590        let mut handles = vec![];
1591
1592        let buffer_reader = Arc::clone(&buffer);
1593        let writers_done_reader = Arc::clone(&writers_done);
1594        // Gate reader until writers wrap around and trigger an overwrite (dropped page),
1595        // guaranteeing overwrite stability is exercised without relying on timing or sleeps.
1596        let reader_handle = start_gated_reader_thread(buffer_reader, writers_done_reader, true);
1597
1598        // Spawn writers
1599        for thread_index in 0..num_threads {
1600            let buffer_clone = Arc::clone(&buffer);
1601            let handle = std::thread::spawn(move || {
1602                for _ in 0..msgs_per_thread {
1603                    if let Ok((res, now, delta)) = buffer_clone.reserve(TestMessage::SIZE) {
1604                        let msg = TestMessage {
1605                            thread_index,
1606                            timestamp_nanos: now.into_nanos() as u64,
1607                            delta: delta.into_nanos() as u64,
1608                            data: *b"Event data\0\0",
1609                        };
1610                        res.write_at(0, &msg.to_bytes());
1611                        buffer_clone.commit(res);
1612                    }
1613                    // Yield to allow other threads to run and cause contention
1614                    std::thread::yield_now();
1615                }
1616            });
1617            handles.push(handle);
1618        }
1619        for handle in handles {
1620            handle.join().unwrap();
1621        }
1622        writers_done.store(true, Ordering::Release);
1623        let all_messages = reader_handle.join().unwrap();
1624        // Check that data is not corrupted
1625        check_all_message_data(&all_messages, num_threads);
1626        // Verify that some pages were dropped (overwritten)
1627        assert!(buffer.dropped_pages() > 0, "Expected at least some dropped pages");
1628    }
1629
1630    #[test]
1631    fn test_out_of_order_completion_same_page() {
1632        let buffer =
1633            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1634                .unwrap();
1635
1636        let header_size = LocklessRingBuffer::PAGE_HEADER_SIZE;
1637        let page_size = (*PAGE_SIZE) as usize;
1638        let available_space = page_size - header_size;
1639
1640        let size1 = available_space / 2;
1641        let size2 = available_space - size1;
1642
1643        // Reserve first chunk
1644        let (res1, _, _) = buffer.reserve(size1).unwrap();
1645        // Reserve second chunk (fills the page)
1646        let (res2, _, _) = buffer.reserve(size2).unwrap();
1647
1648        // Commit second chunk first (out of order)
1649        buffer.commit(res2);
1650
1651        // Page should not be finalized yet because res1 is still active
1652        assert!(buffer.swap_reader_page().is_none());
1653
1654        // Commit first chunk
1655        buffer.commit(res1);
1656
1657        // Force advance tail by making a reservation that won't fit on Page 0.
1658        let _ = buffer.reserve(10).unwrap();
1659
1660        // Now page should be finalized
1661        let swapped = buffer.swap_reader_page();
1662        assert!(swapped.is_some());
1663        assert_eq!(swapped.unwrap(), 0); // First data page is node 0
1664    }
1665
1666    #[test]
1667    fn test_concurrent_readers_rejection() {
1668        #[derive(Debug)]
1669        struct SleepingOutputBuffer {
1670            barrier_in_available: Arc<std::sync::Barrier>,
1671            barrier_test_done: Arc<std::sync::Barrier>,
1672        }
1673        impl Buffer for SleepingOutputBuffer {
1674            fn segments_count(&self) -> Result<usize, Errno> {
1675                Ok(1)
1676            }
1677            fn peek_each_segment(
1678                &mut self,
1679                _callback: &mut PeekBufferSegmentsCallback<'_>,
1680            ) -> Result<(), Errno> {
1681                Ok(())
1682            }
1683        }
1684        impl OutputBuffer for SleepingOutputBuffer {
1685            fn write_each(
1686                &mut self,
1687                _callback: &mut OutputBufferCallback<'_>,
1688            ) -> Result<usize, Errno> {
1689                Ok(0)
1690            }
1691            fn available(&self) -> usize {
1692                // Signal the main thread that we are active and holding the read lock, and block
1693                // until the main thread reaches its wait point.
1694                self.barrier_in_available.wait();
1695
1696                // Block here until the main thread has completed its concurrent read assertion.
1697                self.barrier_test_done.wait();
1698
1699                (*PAGE_SIZE) as usize
1700            }
1701            fn bytes_written(&self) -> usize {
1702                0
1703            }
1704            fn zero(&mut self) -> Result<usize, Errno> {
1705                Ok(0)
1706            }
1707            unsafe fn advance(&mut self, _length: usize) -> Result<(), Errno> {
1708                Ok(())
1709            }
1710        }
1711
1712        let buffer = Arc::new(
1713            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, false, fuchsia_trace::Id::new())
1714                .unwrap(),
1715        );
1716
1717        let barrier_in_available = Arc::new(std::sync::Barrier::new(2));
1718        let barrier_test_done = Arc::new(std::sync::Barrier::new(2));
1719
1720        let buffer_clone = Arc::clone(&buffer);
1721        let barrier_in_available_clone = Arc::clone(&barrier_in_available);
1722        let barrier_test_done_clone = Arc::clone(&barrier_test_done);
1723        let handle = std::thread::spawn(move || {
1724            let mut dest = SleepingOutputBuffer {
1725                barrier_in_available: barrier_in_available_clone,
1726                barrier_test_done: barrier_test_done_clone,
1727            };
1728            let _ = buffer_clone.read(&mut dest);
1729        });
1730
1731        // Wait until the background thread enters `available()` and is holding the read lock.
1732        barrier_in_available.wait();
1733
1734        // Perform concurrent read while background thread holds `reader_active` lock.
1735        let mut dest = VecOutputBuffer::new((*PAGE_SIZE) as usize);
1736        let res = buffer.read(&mut dest);
1737        assert_eq!(res.unwrap_err(), starnix_uapi::errno!(EBUSY));
1738
1739        // Signal the background thread that the assertion is done and it can proceed.
1740        barrier_test_done.wait();
1741
1742        handle.join().unwrap();
1743    }
1744
1745    #[test]
1746    fn test_stale_offset_livelock() {
1747        let num_threads = 8;
1748        let msgs_per_thread = 200;
1749        // Small 3-page buffer to force wraparounds and high contention
1750        let buffer = Arc::new(
1751            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1752                .unwrap(),
1753        );
1754        let mut handles = vec![];
1755
1756        for thread_index in 0..num_threads {
1757            let buffer_clone = Arc::clone(&buffer);
1758            let handle = std::thread::spawn(move || {
1759                for _ in 0..msgs_per_thread {
1760                    if let Ok((res, now, delta)) = buffer_clone.reserve(TestMessage::SIZE) {
1761                        let msg = TestMessage {
1762                            thread_index,
1763                            timestamp_nanos: now.into_nanos() as u64,
1764                            delta: delta.into_nanos() as u64,
1765                            data: *b"Event data\0\0",
1766                        };
1767                        res.write_at(0, &msg.to_bytes());
1768                        buffer_clone.commit(res);
1769                    }
1770                    std::thread::yield_now();
1771                }
1772            });
1773            handles.push(handle);
1774        }
1775
1776        for handle in handles {
1777            handle.join().unwrap();
1778        }
1779
1780        let dropped = buffer.dropped_pages();
1781        println!("High contention completed. Total dropped pages = {}", dropped);
1782        // Dropped pages should be reasonable and not spike to infinity or cause deadlock
1783        assert!(dropped > 0);
1784    }
1785
1786    #[test]
1787    fn test_writer_preemption_and_overwrite_prevention() {
1788        let buffer = Arc::new(
1789            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1790                .unwrap(),
1791        );
1792
1793        // 1. Thread 1 reserves space on Node 0, but delays committing (preempted)
1794        let res1 = buffer.reserve(100).unwrap();
1795        assert_eq!(res1.0.node_idx, 0);
1796
1797        // 2. Thread 2 fills the rest of Node 0, Node 1, and wraps around to Node 0.
1798        // Thread 2 should block or spin since Node 0 has active writers (active_writers > 0)
1799        let buffer_clone = Arc::clone(&buffer);
1800        let writer_finished = Arc::new(std::sync::atomic::AtomicBool::new(false));
1801        let writer_finished_clone = Arc::clone(&writer_finished);
1802
1803        let handle = std::thread::spawn(move || {
1804            // Fill Node 0
1805            let res2 = buffer_clone
1806                .reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE - 100)
1807                .unwrap();
1808            buffer_clone.commit(res2.0);
1809
1810            // Fill Node 1
1811            let res3 = buffer_clone
1812                .reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE)
1813                .unwrap();
1814            buffer_clone.commit(res3.0);
1815
1816            // Try to reserve again - wraps around to Node 0 (which still has res1 active)
1817            // This should block/spin until we commit res1.
1818            let res4 = buffer_clone.reserve(100).unwrap();
1819            buffer_clone.commit(res4.0);
1820
1821            writer_finished_clone.store(true, Ordering::Release);
1822        });
1823
1824        // Give the thread some time to reach the block point
1825        std::thread::sleep(std::time::Duration::from_millis(50));
1826        assert!(!writer_finished.load(Ordering::Acquire));
1827
1828        // Now commit res1, which unblocks the wraparound writer
1829        buffer.commit(res1.0);
1830
1831        handle.join().unwrap();
1832        assert!(writer_finished.load(Ordering::Acquire));
1833    }
1834
1835    #[test]
1836    fn test_extreme_disable_enable_stress() {
1837        let num_threads = 8;
1838        let msgs_per_thread = 100;
1839        let buffer = Arc::new(
1840            LocklessRingBuffer::new(5 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1841                .unwrap(),
1842        );
1843        let mut handles = vec![];
1844
1845        // 1. Spawn 8 writer threads writing continuously until they successfully write msgs_per_thread.
1846        for thread_index in 0..num_threads {
1847            let buffer_clone = Arc::clone(&buffer);
1848            let handle = std::thread::spawn(move || {
1849                let mut count = 0;
1850                while count < msgs_per_thread {
1851                    match buffer_clone.reserve(TestMessage::SIZE) {
1852                        Ok((res, now, delta)) => {
1853                            let msg = TestMessage {
1854                                thread_index,
1855                                timestamp_nanos: now.into_nanos() as u64,
1856                                delta: delta.into_nanos() as u64,
1857                                data: *b"Event data\0\0",
1858                            };
1859                            res.write_at(0, &msg.to_bytes());
1860                            buffer_clone.commit(res);
1861                            count += 1;
1862                        }
1863                        Err(e) if e == starnix_uapi::errno!(ENOMEM) => {
1864                            // Expected when disabled. Yield and retry.
1865                            std::thread::yield_now();
1866                        }
1867                        Err(e) => panic!("Unexpected error during reserve: {:?}", e),
1868                    }
1869                }
1870                count
1871            });
1872            handles.push(handle);
1873        }
1874
1875        // 2. Coordinator thread repeatedly disables and enables the ring buffer.
1876        let buffer_clone = Arc::clone(&buffer);
1877        let coordinator = std::thread::spawn(move || {
1878            for _ in 0..20 {
1879                std::thread::sleep(std::time::Duration::from_millis(5));
1880                // Disable the ring buffer. This must successfully wait for all active writers.
1881                let _dropped = buffer_clone.disable().unwrap();
1882
1883                // Ensure subsequent reservations fail while disabled.
1884                let res = buffer_clone.reserve(TestMessage::SIZE);
1885                assert_eq!(res.unwrap_err(), starnix_uapi::errno!(ENOMEM));
1886
1887                std::thread::sleep(std::time::Duration::from_millis(2));
1888                // Re-enable the ring buffer.
1889                let _now = buffer_clone.enable().unwrap();
1890            }
1891        });
1892
1893        coordinator.join().unwrap();
1894        let mut total_writes = 0;
1895        for handle in handles {
1896            total_writes += handle.join().unwrap();
1897        }
1898        println!("Disable/enable stress completed. Total writes = {}", total_writes);
1899        assert_eq!(total_writes, num_threads * msgs_per_thread);
1900    }
1901
1902    #[test]
1903    fn test_reader_loop_swapping_high_contention() {
1904        let num_threads = 8usize;
1905        let msgs_per_thread = 100;
1906        // 8 threads * 100 msgs * 32 bytes = 25600 bytes.
1907        // Each page holds 127 messages.
1908        // Let's make a 10-page buffer to have ample space.
1909        let buffer = Arc::new(
1910            LocklessRingBuffer::new(10 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1911                .unwrap(),
1912        );
1913        let writers_done = Arc::new(std::sync::atomic::AtomicBool::new(false));
1914        let mut handles = vec![];
1915
1916        // Spawn unified reader thread that loops read() draining all pages.
1917        let buffer_reader = Arc::clone(&buffer);
1918        let writers_done_reader = Arc::clone(&writers_done);
1919        let reader_handle = start_reader_thread(buffer_reader, writers_done_reader);
1920
1921        // Spawn 8 writer threads.
1922        for thread_index in 0..num_threads {
1923            let buffer_clone = Arc::clone(&buffer);
1924            let handle = std::thread::spawn(move || {
1925                for _ in 0..msgs_per_thread {
1926                    let (res, now, delta) = buffer_clone.reserve(TestMessage::SIZE).unwrap();
1927                    let msg = TestMessage {
1928                        thread_index: thread_index as u32,
1929                        timestamp_nanos: now.into_nanos() as u64,
1930                        delta: delta.into_nanos() as u64,
1931                        data: *b"Event data\0\0",
1932                    };
1933                    res.write_at(0, &msg.to_bytes());
1934                    buffer_clone.commit(res);
1935                    // Yield to cause maximum interleaving
1936                    std::thread::yield_now();
1937                }
1938            });
1939            handles.push(handle);
1940        }
1941
1942        for handle in handles {
1943            handle.join().unwrap();
1944        }
1945        writers_done.store(true, Ordering::Release);
1946        let all_messages = reader_handle.join().unwrap();
1947        let messages_read = all_messages.len();
1948        check_all_message_data(&all_messages, num_threads as u32);
1949        let dropped = buffer.dropped_pages();
1950        println!(
1951            "Reader high contention completed. Messages read = {}, dropped pages = {}",
1952            messages_read, dropped
1953        );
1954        // Verify that the total number of read messages + dropped pages' messages covers all writes.
1955        let messages_per_page =
1956            ((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE) / TestMessage::SIZE;
1957        let total_written = num_threads * msgs_per_thread;
1958        let total_accounted = messages_read + (dropped as usize * messages_per_page);
1959        assert!(total_accounted >= total_written - messages_per_page);
1960    }
1961
1962    #[test]
1963    fn test_failed_reservation_offset_boundary() {
1964        let buffer =
1965            LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1966                .unwrap();
1967        let page_size = (*PAGE_SIZE) as usize;
1968        let max_payload = page_size - LocklessRingBuffer::PAGE_HEADER_SIZE;
1969
1970        // Reserve enough space to almost fill Node 0
1971        let msg_size = 100;
1972        let num_msgs = max_payload / msg_size;
1973        for _ in 0..num_msgs {
1974            let (res, _, _) = buffer.reserve(msg_size).unwrap();
1975            buffer.commit(res);
1976        }
1977
1978        // Check current write_offset of Node 0
1979        let expected_offset = LocklessRingBuffer::PAGE_HEADER_SIZE + num_msgs * msg_size;
1980        assert_eq!(buffer.nodes[0].write_offset.load(Ordering::Acquire), expected_offset);
1981
1982        // Now try a reservation that exceeds remaining space on Node 0.
1983        let (res, _, _) = buffer.reserve(msg_size).unwrap();
1984        assert_eq!(res.node_idx, 1);
1985        buffer.commit(res);
1986
1987        // Check write_offset of Node 0.
1988        // Before CAS fix, write_offset was 4116 (> page_size).
1989        // After CAS fix, write_offset is exactly expected_offset (4016).
1990        assert_eq!(buffer.nodes[0].write_offset.load(Ordering::Acquire), expected_offset);
1991    }
1992}