use super::*;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(C)]
struct TestTrade {
time_ns: u64,
px: u64,
sz: u64,
side: u8,
}
#[test]
fn test_ring_buffer_push_and_recv() {
let tmp_path = std::env::temp_dir().join("test_ringfire_spmc.shm");
let _ = std::fs::remove_file(&tmp_path);
let mut producer = RingProducer::<TestTrade>::create(&tmp_path, 1024).unwrap();
let mut consumer = RingConsumer::<TestTrade>::attach(&tmp_path).unwrap();
assert_eq!(consumer.try_recv(), None);
let t1 = TestTrade {
time_ns: 100,
px: 81200,
sz: 15,
side: b'B',
};
let t2 = TestTrade {
time_ns: 200,
px: 81205,
sz: 20,
side: b'S',
};
producer.push(&t1);
producer.push(&t2);
assert_eq!(consumer.try_recv(), Some(t1));
assert_eq!(consumer.try_recv(), Some(t2));
assert_eq!(consumer.try_recv(), None);
}
#[test]
fn test_ring_buffer_batch_recv() {
let tmp_path = std::env::temp_dir().join("test_ringfire_batch.shm");
let _ = std::fs::remove_file(&tmp_path);
let mut producer = RingProducer::<u64>::create(&tmp_path, 1024).unwrap();
let mut consumer = RingConsumer::<u64>::attach(&tmp_path).unwrap();
let items: Vec<u64> = (1..=50).collect();
producer.push_batch(&items);
let mut buf = [0u64; 32];
let n1 = consumer.recv_batch(&mut buf);
assert_eq!(n1, 32);
assert_eq!(&buf[..], &(1..=32).collect::<Vec<u64>>()[..]);
let n2 = consumer.recv_batch(&mut buf);
assert_eq!(n2, 18);
assert_eq!(&buf[..18], &(33..=50).collect::<Vec<u64>>()[..]);
assert_eq!(consumer.recv_batch(&mut buf), 0);
}
#[test]
fn test_ring_buffer_latest_wins_lapping() {
let tmp_path = std::env::temp_dir().join("test_ringfire_lap.shm");
let _ = std::fs::remove_file(&tmp_path);
let capacity = 64;
let mut producer = RingProducer::<u64>::create(&tmp_path, capacity).unwrap();
let mut consumer = RingConsumer::<u64>::attach(&tmp_path).unwrap();
for i in 1..=150 {
producer.push(&i);
}
let status = consumer.recv_status();
match status {
RecvStatus::Lapped { skipped, item } => {
assert_eq!(skipped, 86); assert_eq!(item, 87);
}
_ => panic!("Expected Lapped status"),
}
assert_eq!(consumer.lapped_count(), 86);
for expected in 88..=150 {
assert_eq!(consumer.try_recv(), Some(expected));
}
assert_eq!(consumer.try_recv(), None);
}
#[test]
fn test_reader_extreme_lag_and_jump_to_latest() {
let tmp_path = std::env::temp_dir().join("test_ringfire_extreme_lag.shm");
let _ = std::fs::remove_file(&tmp_path);
let capacity = 1024;
let mut producer = RingProducer::<u64>::create(&tmp_path, capacity).unwrap();
let mut consumer = RingConsumer::<u64>::attach(&tmp_path).unwrap();
for i in 1..=1_000_000 {
producer.push(&i);
}
assert_eq!(consumer.lag(), 1_000_000);
let skipped = consumer.jump_to_latest();
assert_eq!(skipped, 999_999);
assert_eq!(consumer.try_recv(), Some(1_000_000));
assert_eq!(consumer.try_recv(), None);
}
#[test]
fn test_blackboard_read_write() {
let tmp_path = std::env::temp_dir().join("test_ringfire_bb.shm");
let _ = std::fs::remove_file(&tmp_path);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(C)]
struct BboSnapshot {
bid_px: u64,
ask_px: u64,
bid_sz: u64,
ask_sz: u64,
}
let mut producer = BlackboardProducer::<BboSnapshot>::create(&tmp_path, 256).unwrap();
let consumer = BlackboardConsumer::<BboSnapshot>::attach(&tmp_path).unwrap();
assert_eq!(consumer.read(0).unwrap(), None);
let bbo_btc = BboSnapshot {
bid_px: 8_210_000,
ask_px: 8_210_050,
bid_sz: 500,
ask_sz: 300,
};
producer.write(1, &bbo_btc).unwrap();
let read_val = consumer.read(1).unwrap().unwrap();
assert_eq!(read_val, bbo_btc);
}
#[test]
fn test_wait_strategy_yield_and_futex() {
let tmp_path = std::env::temp_dir().join("test_ringfire_wait.shm");
let _ = std::fs::remove_file(&tmp_path);
let mut producer = RingProducer::<u64>::create(&tmp_path, 1024).unwrap();
let mut consumer = RingConsumer::<u64>::attach(&tmp_path).unwrap();
let handle = std::thread::spawn(move || {
let mut wait = FutexWait::default();
consumer.recv_blocking(&mut wait)
});
std::thread::sleep(std::time::Duration::from_millis(10));
let val = 424242u64;
producer.push(&val);
let res = handle.join().unwrap();
assert_eq!(res, val);
}
#[test]
fn test_mpmc_producers() {
let tmp_path = std::env::temp_dir().join("test_ringfire_mpmc.shm");
let _ = std::fs::remove_file(&tmp_path);
let p1 = MpmcProducer::<u64>::create(&tmp_path, 1024).unwrap();
let p2 = MpmcProducer::<u64>::attach(&tmp_path).unwrap();
let mut consumer = RingConsumer::<u64>::attach(&tmp_path).unwrap();
p1.push(&100);
p2.push(&200);
let v1 = consumer.try_recv().unwrap();
let v2 = consumer.try_recv().unwrap();
assert_eq!(v1 + v2, 300);
}
#[cfg(feature = "tokio")]
#[tokio::test]
async fn test_tokio_async_consumer() {
let tmp_path = std::env::temp_dir().join("test_tokio_async.shm");
let _ = std::fs::remove_file(&tmp_path);
let mut producer = RingProducer::<u64>::create(&tmp_path, 1024).unwrap();
let mut async_consumer = AsyncRingConsumer::<u64>::attach(&tmp_path).unwrap();
let counter = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
let c_clone = counter.clone();
let bg_task = tokio::spawn(async move {
for _ in 0..100 {
c_clone.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
tokio::task::yield_now().await;
}
});
let reader_task = tokio::spawn(async move {
let mut received = Vec::new();
for _ in 0..5 {
let msg = async_consumer.recv().await;
received.push(msg);
}
received
});
for val in 1..=5 {
tokio::time::sleep(std::time::Duration::from_millis(5)).await;
producer.push(&val);
}
let msgs = reader_task.await.unwrap();
assert_eq!(msgs, vec![1, 2, 3, 4, 5]);
bg_task.await.unwrap();
assert_eq!(counter.load(std::sync::atomic::Ordering::SeqCst), 100);
}