#[cfg(target_os = "linux")]
use safer_ring::Ring;
use safer_ring::{BufferPool, PinnedBuffer};
use std::sync::Arc;
use std::time::{Duration, Instant};
use tokio::time::timeout;
#[tokio::test]
async fn stress_test_buffer_pool_high_frequency() {
const POOL_SIZE: usize = 100;
const BUFFER_SIZE: usize = 4096;
const OPERATIONS: usize = 10000;
const CONCURRENT_TASKS: usize = 10;
let pool = Arc::new(BufferPool::new(POOL_SIZE, BUFFER_SIZE));
let start_time = Instant::now();
let mut handles = Vec::new();
for task_id in 0..CONCURRENT_TASKS {
let pool_clone = Arc::clone(&pool);
let handle = tokio::spawn(async move {
let mut local_operations = 0;
let operations_per_task = OPERATIONS / CONCURRENT_TASKS;
for i in 0..operations_per_task {
let mut buffer = loop {
if let Some(buf) = pool_clone.get() {
break buf;
}
tokio::task::yield_now().await;
};
let mut slice = buffer.as_mut_slice();
for (idx, byte) in slice.iter_mut().enumerate() {
*byte = ((task_id + i + idx) % 256) as u8;
}
for (idx, &byte) in slice.iter().enumerate() {
assert_eq!(byte, ((task_id + i + idx) % 256) as u8);
}
local_operations += 1;
if i % 100 == 0 {
tokio::task::yield_now().await;
}
}
local_operations
});
handles.push(handle);
}
let results = timeout(Duration::from_secs(30), futures::future::join_all(handles))
.await
.expect("Stress test timed out");
let total_operations: usize = results.into_iter().map(|r| r.unwrap()).sum();
let elapsed = start_time.elapsed();
println!("Completed {total_operations} operations in {elapsed:?}");
println!(
"Operations per second: {:.2}",
total_operations as f64 / elapsed.as_secs_f64()
);
let stats = pool.stats();
assert_eq!(stats.total_buffers, POOL_SIZE);
assert_eq!(stats.available_buffers, POOL_SIZE);
assert_eq!(stats.in_use_buffers, 0);
}
#[cfg(target_os = "linux")]
#[tokio::test]
async fn stress_test_ring_operations() {
const RING_SIZE: usize = 256;
const OPERATIONS: usize = 5000;
const BUFFER_SIZE: usize = 1024;
let ring = Ring::new(RING_SIZE as u32).unwrap();
let pool = Arc::new(BufferPool::new(RING_SIZE * 2, BUFFER_SIZE));
let start_time = Instant::now();
let mut handles = Vec::new();
let temp_dir = tempfile::tempdir().unwrap();
let mut temp_files = Vec::new();
for i in 0..10 {
let file_path = temp_dir.path().join(format!("test_file_{i}.txt"));
let file = std::fs::File::create(&file_path).unwrap();
temp_files.push(file);
}
for task_id in 0..5 {
let pool_clone = Arc::clone(&pool);
let handle = tokio::spawn(async move {
let mut completed_ops = 0;
for i in 0..(OPERATIONS / 5) {
let mut buffer = loop {
if let Some(buf) = pool_clone.get() {
break buf;
}
tokio::task::yield_now().await;
};
let mut slice = buffer.as_mut_slice();
for (idx, byte) in slice.iter_mut().enumerate() {
*byte = ((task_id + i + idx) % 256) as u8;
}
tokio::task::yield_now().await;
for (idx, &byte) in slice.iter().enumerate() {
assert_eq!(byte, ((task_id + i + idx) % 256) as u8);
}
completed_ops += 1;
if i % 50 == 0 {
tokio::task::yield_now().await;
}
}
completed_ops
});
handles.push(handle);
}
let results = timeout(Duration::from_secs(60), futures::future::join_all(handles))
.await
.expect("Ring stress test timed out");
let total_operations: usize = results.into_iter().map(|r| r.unwrap()).sum();
let elapsed = start_time.elapsed();
println!("Ring completed {total_operations} operations in {elapsed:?}");
println!(
"Ring operations per second: {:.2}",
total_operations as f64 / elapsed.as_secs_f64()
);
assert_eq!(ring.operations_in_flight(), 0);
}
#[cfg(not(target_os = "linux"))]
#[tokio::test]
async fn stress_test_ring_operations() {
use safer_ring::Ring;
match Ring::new(64) {
Ok(_) => panic!("Ring creation should fail on non-Linux platforms"),
Err(e) => {
println!("Expected error on non-Linux platform: {}", e);
}
}
}
#[tokio::test]
async fn stress_test_memory_pressure() {
const ITERATIONS: usize = 1000;
const BUFFERS_PER_ITERATION: usize = 50;
const BUFFER_SIZE: usize = 8192;
for iteration in 0..ITERATIONS {
let mut buffers = Vec::new();
for _ in 0..BUFFERS_PER_ITERATION {
let buffer = PinnedBuffer::with_capacity(BUFFER_SIZE);
buffers.push(buffer);
}
for (i, buffer) in buffers.iter_mut().enumerate() {
let mut slice = buffer.as_mut_slice();
for (idx, byte) in slice.iter_mut().enumerate() {
*byte = ((iteration + i + idx) % 256) as u8;
}
}
for (i, buffer) in buffers.iter().enumerate() {
let slice = buffer.as_slice();
for (idx, &byte) in slice.iter().enumerate() {
assert_eq!(byte, ((iteration + i + idx) % 256) as u8);
}
}
if iteration % 100 == 0 {
println!("Memory pressure test: iteration {iteration}/{ITERATIONS}");
}
if iteration % 10 == 0 {
tokio::task::yield_now().await;
}
}
println!("Memory pressure test completed successfully");
}
#[tokio::test]
async fn stress_test_error_recovery() {
const POOL_SIZE: usize = 20;
const BUFFER_SIZE: usize = 1024;
const ERROR_RATE: usize = 10;
let pool = Arc::new(BufferPool::new(POOL_SIZE, BUFFER_SIZE));
let mut handles = Vec::new();
for task_id in 0..5 {
let pool_clone = Arc::clone(&pool);
let handle = tokio::spawn(async move {
let mut successful_ops = 0;
let mut failed_ops = 0;
for i in 0..200 {
let mut buffer = loop {
if let Some(buf) = pool_clone.get() {
break buf;
}
tokio::task::yield_now().await;
};
if (task_id + i) % ERROR_RATE == 0 {
failed_ops += 1;
drop(buffer); } else {
let mut slice = buffer.as_mut_slice();
slice[0] = (i % 256) as u8;
successful_ops += 1;
}
tokio::task::yield_now().await;
}
(successful_ops, failed_ops)
});
handles.push(handle);
}
let results = futures::future::join_all(handles).await;
let (total_success, total_failures): (usize, usize) = results
.into_iter()
.map(|r| r.unwrap())
.fold((0, 0), |(s1, f1), (s2, f2)| (s1 + s2, f1 + f2));
println!("Error recovery test: {total_success} successful, {total_failures} failed operations");
let stats = pool.stats();
assert_eq!(stats.total_buffers, POOL_SIZE);
assert_eq!(stats.available_buffers, POOL_SIZE);
assert_eq!(stats.in_use_buffers, 0);
}