#[cfg(test)]
mod tests {
use crate::proxy::ring_buffer::RingBuffer;
use crate::proxy::types::*;
use std::sync::Arc;
use std::thread;
use std::time::Instant;
#[test]
fn test_single_threaded_performance() {
let config = RingBufferConfig {
buffer_size: BufferSize::try_new(1024 * 1024).expect("valid size"), slot_size: SlotSize::try_new(1024).expect("valid size"), };
let request_id = RequestId::new();
let data = vec![0u8; 512]; let iterations = 10000;
let ring_buffer = RingBuffer::new(&config);
let start = Instant::now();
for _ in 0..iterations {
let _ = ring_buffer.write(request_id, &data);
}
let duration = start.elapsed();
let ns_per_write = duration.as_nanos() as f64 / iterations as f64;
println!("Single-threaded performance ({iterations} writes):");
println!(" Total duration: {duration:?}");
println!(" Per write: {ns_per_write:.2}ns");
assert!(
ns_per_write < 1000.0,
"Ring buffer should achieve <1μs per write, got {ns_per_write:.2}ns"
);
if ns_per_write > 100.0 {
println!(" WARNING: Performance degraded above expected 100ns threshold");
}
}
#[test]
fn test_concurrent_performance() {
let config = RingBufferConfig {
buffer_size: BufferSize::try_new(4 * 1024 * 1024).expect("valid size"), slot_size: SlotSize::try_new(1024).expect("valid size"), };
let thread_count = 4;
let writes_per_thread = 1000;
println!(
"Concurrent performance test ({thread_count} threads, {writes_per_thread} writes each):"
);
let ring_buffer = Arc::new(RingBuffer::new(&config));
let start = Instant::now();
let handles: Vec<_> = (0..thread_count)
.map(|_| {
let buffer = Arc::clone(&ring_buffer);
thread::spawn(move || {
let request_id = RequestId::new();
let data = vec![0u8; 256];
for _ in 0..writes_per_thread {
let _ = buffer.write(request_id, &data);
}
})
})
.collect();
for handle in handles {
handle.join().unwrap();
}
let duration = start.elapsed();
let total_operations = thread_count * writes_per_thread;
let ops_per_sec = total_operations as f64 / duration.as_secs_f64();
let ns_per_op = duration.as_nanos() as f64 / total_operations as f64;
println!(" Total duration: {duration:?}");
println!(" Operations per second: {ops_per_sec:.0}");
println!(" Per operation: {ns_per_op:.2}ns");
assert!(
ns_per_op < 100000.0,
"Concurrent operations should be <100μs, got {ns_per_op:.2}ns"
);
}
#[test]
fn test_ring_buffer_correctness() {
let config = RingBufferConfig {
buffer_size: BufferSize::try_new(1024).expect("valid size"),
slot_size: SlotSize::try_new(256).expect("valid size"),
};
let request_id = RequestId::new();
let data = b"test data for correctness";
let ring_buffer = RingBuffer::new(&config);
assert!(ring_buffer.write(request_id, data).is_ok());
let (read_id, read_data) = ring_buffer.read().expect("Should read data");
assert_eq!(read_id, request_id);
assert_eq!(&read_data[..], data);
println!("✅ Ring buffer correctness verified!");
}
#[test]
fn test_performance_under_pressure() {
let config = RingBufferConfig {
buffer_size: BufferSize::try_new(4 * 1024 * 1024).expect("valid size"), slot_size: SlotSize::try_new(1024).expect("valid size"), };
let ring_buffer = Arc::new(RingBuffer::new(&config));
let thread_count = 8;
let duration_secs = 1;
println!("Stress test: {thread_count} threads for {duration_secs} second(s)");
let start = Instant::now();
let stop_flag = Arc::new(AtomicBool::new(false));
let handles: Vec<_> = (0..thread_count)
.map(|_i| {
let buffer = Arc::clone(&ring_buffer);
let stop = Arc::clone(&stop_flag);
thread::spawn(move || {
let request_id = RequestId::new();
let data = vec![0u8; 512];
let mut count = 0u64;
while !stop.load(Ordering::Relaxed) {
if buffer.write(request_id, &data).is_ok() {
count += 1;
}
}
count
})
})
.collect();
thread::sleep(Duration::from_secs(duration_secs));
stop_flag.store(true, Ordering::Relaxed);
let mut total_writes = 0u64;
for (i, handle) in handles.into_iter().enumerate() {
let count = handle.join().unwrap();
total_writes += count;
println!(" Thread {i}: {count} writes");
}
let actual_duration = start.elapsed();
let writes_per_sec = total_writes as f64 / actual_duration.as_secs_f64();
let ns_per_write = actual_duration.as_nanos() as f64 / total_writes as f64;
println!(" Total writes: {total_writes}");
println!(" Writes per second: {writes_per_sec:.0}");
println!(" Average latency: {ns_per_write:.2}ns");
assert!(
writes_per_sec > 1_000.0,
"Should handle >1K ops/sec under stress, got {writes_per_sec:.0}"
);
}
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
}