use super::*;
use std::sync::atomic::Ordering;
#[test]
fn test_spsc_concurrency() {
let (mut prod, mut cons) = RingBuffer::<i32>::new(64);
let handle = std::thread::spawn(move || {
let mut count = 0;
while count < 1000 {
if prod.push(count).is_ok() {
count += 1;
}
std::thread::yield_now(); }
});
let mut count = 0;
while count < 1000 {
if let Ok(val) = cons.pop() {
assert_eq!(val, count);
count += 1;
}
std::thread::yield_now();
}
handle.join().unwrap();
}
#[test]
fn test_spsc_full_empty() {
let (mut prod, mut cons) = RingBuffer::<i32>::new(4);
assert!(cons.pop().is_err());
assert!(prod.push(1).is_ok());
assert!(prod.push(2).is_ok());
assert!(prod.push(3).is_ok());
assert!(prod.push(4).is_ok());
assert!(prod.push(5).is_err());
assert_eq!(cons.pop(), Ok(1));
assert!(prod.push(5).is_ok());
assert!(prod.push(6).is_err());
}
#[test]
fn test_gc_overflow_overwrite() {
use std::sync::Arc;
use std::sync::atomic::AtomicU32;
let rt_status = RtStatusFlags::new();
let overflow = GcOverflowBuffer::new(64);
let counter = Arc::new(AtomicU32::new(0));
for _ in 0..64 {
let item = GcItem::Test(Box::new(counter.clone()));
overflow.push(item);
}
let item_65 = GcItem::Test(Box::new(counter.clone()));
overflow.push(item_65);
let drained = overflow.drain(&rt_status);
assert_eq!(drained.len(), 64);
}
#[test]
fn test_gc_stress_no_leak() {
use std::sync::Arc;
use std::sync::atomic::AtomicU32;
let (mut gc_prod, mut gc_cons) = RingBuffer::<GcItem>::new(32);
let rt_status = RtStatusFlags::new();
let overflow = GcOverflowBuffer::new(32);
let counter = Arc::new(AtomicU32::new(0));
for _ in 0..1000 {
let item = GcItem::Test(Box::new(counter.clone()));
if let Err(rtrb::PushError::Full(returned_item)) = gc_prod.push(item) {
overflow.push(returned_item);
}
super::drain_gc_channels(&mut gc_cons, &overflow, &rt_status);
}
drop(gc_cons);
for _ in overflow.drain(&rt_status) {}
}
#[test]
fn test_gc_corrupted_slot_unknown_type() {
use std::sync::Arc;
use std::sync::atomic::AtomicU32;
let rt_status = RtStatusFlags::new();
let overflow = GcOverflowBuffer::new(4);
let counter = Arc::new(AtomicU32::new(0));
let valid_item = GcItem::Test(Box::new(counter.clone()));
let valid_packed = valid_item.into_packed();
let ptr = (valid_packed & 0x00FF_FFFF_FFFF_FFFF) as *mut std::ffi::c_void;
overflow.push(GcItem::Test(Box::new(counter.clone())));
overflow.push(GcItem::Test(Box::new(counter.clone())));
overflow.push(GcItem::Test(Box::new(counter.clone())));
overflow.slots[3].store(
(99u64 << 56) | (ptr as u64 & 0x00FF_FFFF_FFFF_FFFF),
Ordering::Release,
);
assert!(!rt_status.check_flag(RT_STATUS_GC_CORRUPTED));
let drained = overflow.drain(&rt_status);
assert_eq!(drained.len(), 3);
drop(drained);
assert!(rt_status.check_flag(RT_STATUS_GC_CORRUPTED));
let drained_again = overflow.drain(&rt_status);
assert_eq!(drained_again.len(), 0);
drop(drained_again);
}
#[test]
fn test_gc_corrupted_slot_null_type_non_null_ptr() {
use std::sync::Arc;
use std::sync::atomic::AtomicU32;
let rt_status = RtStatusFlags::new();
let overflow = GcOverflowBuffer::new(4);
let counter = Arc::new(AtomicU32::new(0));
let valid_item = GcItem::Test(Box::new(counter.clone()));
let packed = valid_item.into_packed();
let ptr = (packed & 0x00FF_FFFF_FFFF_FFFF) as *mut std::ffi::c_void;
overflow.push(GcItem::Test(Box::new(counter.clone())));
overflow.push(GcItem::Test(Box::new(counter.clone())));
overflow.slots[2].store(ptr as u64 & 0x00FF_FFFF_FFFF_FFFF, Ordering::Release);
assert!(!rt_status.check_flag(RT_STATUS_GC_CORRUPTED));
let drained = overflow.drain(&rt_status);
assert_eq!(drained.len(), 2);
drop(drained);
assert!(rt_status.check_flag(RT_STATUS_GC_CORRUPTED));
let drained_again = overflow.drain(&rt_status);
assert_eq!(drained_again.len(), 0);
drop(drained_again);
}
#[test]
fn test_gc_concurrent_push_drain() {
use std::sync::Arc;
use std::sync::atomic::AtomicBool;
let rt_status = RtStatusFlags::new();
let overflow = Arc::new(GcOverflowBuffer::new(32));
let overflow_prod = Arc::clone(&overflow);
let done = Arc::new(AtomicBool::new(false));
let done_prod = Arc::clone(&done);
let producer = std::thread::spawn(move || {
let counter = Arc::new(std::sync::atomic::AtomicU32::new(0));
while !done_prod.load(Ordering::Relaxed) {
let item = GcItem::Test(Box::new(counter.clone()));
overflow_prod.push(item);
}
});
let overflow_cons = Arc::clone(&overflow);
let done_cons = Arc::clone(&done);
let consumer = std::thread::spawn(move || {
let local_rt_status = RtStatusFlags::new();
for _ in 0..2000 {
let drained = overflow_cons.drain(&local_rt_status);
drop(drained);
std::thread::yield_now();
}
done_cons.store(true, Ordering::Relaxed);
});
producer.join().unwrap();
consumer.join().unwrap();
let final_drained = overflow.drain(&rt_status);
for item in final_drained {
match item {
GcItem::Test(_) => {} _ => panic!("Unexpected GcItem variant in concurrent test"),
}
}
}
#[test]
fn test_drain_gc_empty_is_noop() {
let (_, mut consumer) = RingBuffer::<GcItem>::new(16);
let overflow = GcOverflowBuffer::new(16);
let rt_status = RtStatusFlags::new();
let drained = drain_gc_channels(&mut consumer, &overflow, &rt_status);
assert_eq!(drained, 0);
}