lowlet 0.1.2

Low-latency IPC library using shared memory and lock-free structures
Documentation
use std::thread;

use lowlet::channel;
use lowlet::lf::SpscQueue;
use lowlet::sync::{Barrier, Spinlock};

#[test]
fn test_channel_basic() {
    let (tx, rx) = channel::<u64, 64>();

    tx.send(42).unwrap();
    tx.send(43).unwrap();

    assert_eq!(rx.recv().unwrap(), 42);
    assert_eq!(rx.recv().unwrap(), 43);
}

#[test]
fn test_channel_full() {
    let (tx, rx) = channel::<u64, 4>();

    for i in 0..4 {
        tx.send(i).unwrap();
    }

    assert!(tx.send(100).is_err());

    rx.recv().unwrap();
    tx.send(100).unwrap();
}

#[test]
fn test_channel_threaded() {
    let (tx, rx) = channel::<u64, 1024>();

    let producer = thread::spawn(move || {
        for i in 0..1000u64 {
            loop {
                if tx.send(i).is_ok() {
                    break;
                }
                thread::yield_now();
            }
        }
    });

    let consumer = thread::spawn(move || {
        let mut sum = 0u64;
        for _ in 0..1000 {
            sum += rx.recv_spin().unwrap();
        }
        sum
    });

    producer.join().unwrap();
    let sum = consumer.join().unwrap();

    assert_eq!(sum, (0..1000u64).sum());
}

#[test]
fn test_spsc_queue() {
    let queue = SpscQueue::<u64, 64>::new();

    queue.push(1).unwrap();
    queue.push(2).unwrap();
    queue.push(3).unwrap();

    assert_eq!(queue.pop().unwrap(), 1);
    assert_eq!(queue.pop().unwrap(), 2);
    assert_eq!(queue.pop().unwrap(), 3);
    assert!(queue.pop().is_err());
}

#[test]
fn test_spinlock() {
    let lock = Spinlock::new(0u64);

    {
        let mut guard = lock.lock();
        *guard = 42;
    }

    assert_eq!(*lock.lock(), 42);
}

#[test]
fn test_spinlock_contention() {
    use std::sync::Arc;

    let lock = Arc::new(Spinlock::new(0u64));
    let mut handles = vec![];

    for _ in 0..4 {
        let lock = Arc::clone(&lock);
        handles.push(thread::spawn(move || {
            for _ in 0..1000 {
                let mut guard = lock.lock();
                *guard += 1;
            }
        }));
    }

    for h in handles {
        h.join().unwrap();
    }

    assert_eq!(*lock.lock(), 4000);
}

#[test]
fn test_barrier() {
    use std::sync::atomic::{AtomicU32, Ordering};
    use std::sync::Arc;

    let barrier = Arc::new(Barrier::new(4));
    let counter = Arc::new(AtomicU32::new(0));
    let mut handles = vec![];

    for _ in 0..4 {
        let barrier = Arc::clone(&barrier);
        let counter = Arc::clone(&counter);
        handles.push(thread::spawn(move || {
            counter.fetch_add(1, Ordering::SeqCst);
            barrier.wait();
            assert_eq!(counter.load(Ordering::SeqCst), 4);
        }));
    }

    for h in handles {
        h.join().unwrap();
    }
}