1use 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: SharedBuffer,
23 next: AtomicUsize,
26 prev: AtomicUsize,
27 write_offset: AtomicUsize,
29 active_writers: AtomicUsize,
31}
32
33const FLAG_MASK: usize = 0b11 << 62;
34const FLAG_NORMAL: usize = 0b00 << 62;
36const FLAG_HEADER: usize = 0b01 << 62;
38const FLAG_UPDATE: usize = 0b10 << 62;
40
41const PAGE_ACTIVE_BIT: usize = 1 << 63;
44const PAGE_FINALIZED_BIT: usize = 1 << 62;
46const FLAGS_MASK: usize = PAGE_ACTIVE_BIT | PAGE_FINALIZED_BIT;
47const ACTIVE_WRITERS_MASK: usize = !FLAGS_MASK;
48
49const SPIN_SLEEP_DURATION: std::time::Duration = std::time::Duration::from_micros(50);
54
55const 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 fn finalize(&self) {
73 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 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 self.active_writers.fetch_and(!PAGE_ACTIVE_BIT, Ordering::Release);
93 }
94
95 fn release_writer(&self) {
97 let prev_writers = self.active_writers.fetch_sub(1, Ordering::Release);
100
101 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 self.finalize();
113 }
114 }
115 }
116}
117
118static_assertions::const_assert!(std::mem::size_of::<usize>() == 8);
120
121const RING_ENABLED_BIT: usize = 1 << 63;
124
125#[derive(RcuDroppable)]
126pub struct LocklessRingBuffer {
127 vmo: MemoryObject,
128 mapping: SharedBuffer,
129 nodes: Vec<Node>,
131 head_page: AtomicUsize,
133 tail_page: AtomicUsize,
135
136 reader_page: AtomicUsize,
138
139 ref_count: AtomicUsize,
142 overwrite: bool,
145 dropped_pages: std::sync::atomic::AtomicU64,
147 prev_timestamp: std::sync::atomic::AtomicU64,
150 write_event_async_id: fuchsia_trace::Id,
152 reader_active: std::sync::atomic::AtomicBool,
156 state_mutex: LockDepMutex<(), PerfRingBufferStateLock>,
159}
160impl LocklessRingBuffer {
161 pub const PAGE_HEADER_SIZE: usize = 16;
163
164 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 let pages = std::cmp::max(3, requested_pages);
176 let total_nodes = pages;
177 let capacity = total_nodes * (*PAGE_SIZE) as usize;
178
179 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 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 let mapping = unsafe { SharedBuffer::new(addr as *mut u8, capacity) };
197
198 let mut nodes = Vec::with_capacity(total_nodes);
200 let base_ptr = addr as *mut u8;
201 for i in 0..total_nodes {
202 let page_ptr = unsafe { base_ptr.add(i * (*PAGE_SIZE) as usize) };
206
207 nodes.push(Node {
208 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 active_writers: AtomicUsize::new(PAGE_ACTIVE_BIT),
215 });
216 }
217 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 nodes[circle_size - 1].next.store(Node::make_val(0, FLAG_HEADER), Ordering::Relaxed);
227 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 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#[derive(Default, Debug)]
261struct YieldTracker {
262 update_flag: u64,
264 node_match: u64,
266 head_lock: u64,
268}
269impl YieldTracker {
270 fn total(&self) -> u64 {
271 self.update_flag + self.node_match + self.head_lock
272 }
273
274 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
287pub 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 Advanced,
340 Yielded,
342 Error(Errno),
344}
345
346impl LocklessRingBuffer {
347 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 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 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 fn advance_to_next_page(
404 &self,
405 tail_node: &Node,
406 tail_val: usize,
407 yield_tracker: &mut YieldTracker,
408 ) -> AdvanceResult {
409 if self.tail_page.load(Ordering::Acquire) != tail_val {
411 return AdvanceResult::Advanced;
412 }
413
414 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 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 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 let expected_next = Node::make_val(next_idx, FLAG_HEADER);
454 let locked_next = Node::make_val(next_idx, FLAG_UPDATE);
455 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 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 let new_head_next = Node::make_val(head_next_idx, FLAG_HEADER);
476 head_node.next.store(new_head_next, Ordering::Release);
477
478 self.head_page.store(head_next_idx, Ordering::Release);
480
481 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 yield_tracker.head_lock += 1;
494 yield_tracker.yield_or_sleep();
495 return AdvanceResult::Yielded;
496 }
497 } else {
498 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 let (should_advance, won_reset) = {
509 let mut current = next_node.write_offset.load(Ordering::Acquire);
510 loop {
511 if self.tail_page.load(Ordering::Acquire) != tail_val {
513 break (false, false);
514 }
515 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 next_node.active_writers.store(PAGE_ACTIVE_BIT, Ordering::Release);
538 next_node.page_data.write_at(8, &0u64.to_le_bytes());
540 }
541
542 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 if size == 0 || size > (*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE {
562 return error!(EINVAL);
563 }
564 let mut val = self.ref_count.load(Ordering::Acquire);
566 loop {
567 if val & RING_ENABLED_BIT == 0 {
568 return error!(ENOMEM);
569 }
570 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 let ref_count_guard = scopeguard::guard(&self.ref_count, |ref_count| {
593 ref_count.fetch_sub(1, Ordering::Release);
594 });
595
596 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 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 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 if active_incremented {
660 self.nodes[tail_idx].release_writer();
661 }
662 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 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 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 scopeguard::ScopeGuard::into_inner(ref_count_guard);
691 Ok((res, now, delta))
692 }
693 Err(e) => {
694 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 return None;
713 }
714 let head_node = &self.nodes[head_idx];
715 let active_writers = head_node.active_writers.load(Ordering::Acquire);
718 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 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 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 let reader_idx = self.reader_page.load(Ordering::Acquire);
754 let next_node = &self.nodes[next_idx];
755 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 self.nodes[reader_idx]
765 .write_offset
766 .store(LocklessRingBuffer::PAGE_HEADER_SIZE, Ordering::Relaxed);
767 self.nodes[reader_idx].active_writers.store(PAGE_ACTIVE_BIT, Ordering::Relaxed);
769 self.nodes[reader_idx].page_data.write_at(8, &0u64.to_le_bytes());
771 let unlocked_val = Node::make_val(reader_idx, FLAG_NORMAL);
773 prev_node.next.store(unlocked_val, Ordering::Release);
774 self.head_page.store(next_idx, Ordering::Release);
776 self.reader_page.store(head_idx, Ordering::Release);
777 return Some(head_idx);
779 }
780 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 std::hint::spin_loop();
791 }
792 }
793 pub fn read(&self, buf: &mut dyn OutputBuffer) -> Result<usize, Errno> {
794 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 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 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 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 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 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 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 self.nodes[initial_pages].page_data.write_at(8, &0u64.to_le_bytes());
942 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 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 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 let res1 =
1044 buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1045 buffer.commit(res1.0);
1046 let res2 = buffer.reserve(100).unwrap();
1048 assert_eq!(res2.0.node_idx, 1);
1049 buffer.commit(res2.0);
1050 let old_head = buffer.swap_reader_page();
1053 assert_eq!(old_head, Some(0));
1054 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 let res1 = buffer
1090 .reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE - 50)
1091 .unwrap();
1092 buffer.commit(res1.0);
1093 let res2 = buffer.reserve(100).unwrap();
1095 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 let res1 =
1106 buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1107 buffer.commit(res1.0);
1108 let _res2 =
1110 buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1111 let res3 = buffer.reserve(100).unwrap();
1115 assert_eq!(res3.0.node_idx, 0);
1117 assert_eq!(buffer.dropped_pages(), 1);
1118 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 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 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 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 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 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 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 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 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 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_all_message_data(&all_messages, num_threads);
1310 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 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 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 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 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 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 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 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 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_all_message_data(&all_messages, num_threads);
1434 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 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 let handle = std::thread::spawn(move || {
1460 buffer_clone.disable().unwrap();
1461 disable_finished_clone.store(true, Ordering::Relaxed);
1462 });
1463 while buffer.ref_count.load(Ordering::Acquire) & RING_ENABLED_BIT != 0 {
1465 std::hint::spin_loop();
1466 }
1467 assert!(!disable_finished.load(Ordering::Relaxed));
1469 let res2 = buffer.reserve(100);
1471 assert_eq!(res2.unwrap_err(), starnix_uapi::errno!(ENOMEM));
1472 let data = vec![42u8; 100];
1475 res.0.write_at(0, &data);
1476 buffer.commit(res.0);
1477 handle.join().unwrap();
1479 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 let res = buffer.reserve(100).unwrap();
1489 assert_eq!(buffer.ref_count.load(Ordering::Relaxed), RING_ENABLED_BIT | 1);
1490 std::mem::drop(res.0);
1492 assert_eq!(buffer.ref_count.load(Ordering::Relaxed), RING_ENABLED_BIT);
1494 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 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 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 let res1 =
1543 buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1544 buffer.commit(res1.0);
1545 let res2 =
1547 buffer.reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE).unwrap();
1548 buffer.commit(res2.0);
1549 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 let buffer = Arc::new(
1569 LocklessRingBuffer::new(3 * (*PAGE_SIZE) as usize, true, fuchsia_trace::Id::new())
1570 .unwrap(),
1571 );
1572 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 let reader_handle = start_gated_reader_thread(buffer_reader, writers_done_reader, true);
1597
1598 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 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_all_message_data(&all_messages, num_threads);
1626 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 let (res1, _, _) = buffer.reserve(size1).unwrap();
1645 let (res2, _, _) = buffer.reserve(size2).unwrap();
1647
1648 buffer.commit(res2);
1650
1651 assert!(buffer.swap_reader_page().is_none());
1653
1654 buffer.commit(res1);
1656
1657 let _ = buffer.reserve(10).unwrap();
1659
1660 let swapped = buffer.swap_reader_page();
1662 assert!(swapped.is_some());
1663 assert_eq!(swapped.unwrap(), 0); }
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 self.barrier_in_available.wait();
1695
1696 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 barrier_in_available.wait();
1733
1734 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 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 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 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 let res1 = buffer.reserve(100).unwrap();
1795 assert_eq!(res1.0.node_idx, 0);
1796
1797 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 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 let res3 = buffer_clone
1812 .reserve((*PAGE_SIZE) as usize - LocklessRingBuffer::PAGE_HEADER_SIZE)
1813 .unwrap();
1814 buffer_clone.commit(res3.0);
1815
1816 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 std::thread::sleep(std::time::Duration::from_millis(50));
1826 assert!(!writer_finished.load(Ordering::Acquire));
1827
1828 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 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 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 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 let _dropped = buffer_clone.disable().unwrap();
1882
1883 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 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 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 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 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 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 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 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 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 let (res, _, _) = buffer.reserve(msg_size).unwrap();
1984 assert_eq!(res.node_idx, 1);
1985 buffer.commit(res);
1986
1987 assert_eq!(buffer.nodes[0].write_offset.load(Ordering::Acquire), expected_offset);
1991 }
1992}