ringfire 0.5.2

Zero-copy lock-free inter-process communication (IPC) ring buffer and shared memory bus in Rust, with ring mirroring across hosts
Documentation
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();

    // Write 150 items into a capacity 64 buffer (lapping consumer by > 2x capacity)
    for i in 1..=150 {
        producer.push(&i);
    }

    // Consumer reads: was lapped by full buffer -> jumps to oldest surviving item (150 - 64 + 1 = 87)
    let status = consumer.recv_status();
    match status {
        RecvStatus::Lapped { skipped, item } => {
            assert_eq!(skipped, 86); // 87 - 1 = 86 skipped
            assert_eq!(item, 87);
        }
        _ => panic!("Expected Lapped status"),
    }

    assert_eq!(consumer.lapped_count(), 86);

    // Subsequent reads drain from 88 up to 150 in sequence (all 64 surviving messages in the ring!)
    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();

    // Producer writes 1,000,000 items while consumer is completely asleep
    for i in 1..=1_000_000 {
        producer.push(&i);
    }

    // Reader checks lag
    assert_eq!(consumer.lag(), 1_000_000);

    // Reader can jump directly to latest
    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();

    // Spawn a background cooperative counter task on the Tokio runtime
    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;
        }
    });

    // Spawn async reader task
    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
    });

    // Push messages with small delay
    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]);

    // Verify background task made continuous progress (not starved by the reader!)
    bg_task.await.unwrap();
    assert_eq!(counter.load(std::sync::atomic::Ordering::SeqCst), 100);
}