#![expect(
clippy::unwrap_used,
reason = "test scope: failed precondition = test failure"
)]
use moirai_sync::{Barrier, ConcurrentHashMap, FutexMutex, LockFreeStack, WaitGroup};
use moirai_utils::LockFreeQueue;
use std::sync::Arc;
use std::thread;
use std::time::{Duration, Instant};
fn main() {
println!("Moirai Synchronization Primitives");
println!("=================================");
println!("\n1. Futex Mutex Example:");
let counter = Arc::new(FutexMutex::new(0));
let mut handles = vec![];
for i in 0..10 {
let counter_clone = Arc::clone(&counter);
let handle = thread::spawn(move || {
let mut guard = counter_clone.lock();
*guard += 1;
println!(" Thread {}: Counter = {}", i, *guard);
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!(" Final counter value: {}", *counter.lock());
println!("\n2. Concurrent HashMap Example:");
let map = Arc::new(ConcurrentHashMap::<String, i32>::new());
let mut handles = vec![];
for i in 0..5 {
let map_clone = Arc::clone(&map);
let handle = thread::spawn(move || {
for j in 0..10 {
let key = format!("thread_{}_item_{}", i, j);
map_clone
.insert(key.clone(), i * 10 + j)
.expect("concurrent map insert should succeed");
println!(" Thread {}: Inserted {}", i, key);
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!(" Map operations completed");
if let Ok(Some(value)) = map.get(&"thread_0_item_0".to_string()) {
println!(" Sample value: thread_0_item_0 = {}", value);
}
println!("\n3. Lock-Free Stack Example:");
let stack = Arc::new(LockFreeStack::new());
let mut handles = vec![];
for i in 0..5 {
let stack_clone = Arc::clone(&stack);
let handle = thread::spawn(move || {
for j in 0..5 {
let value = i * 10 + j;
match stack_clone.push(value) {
Ok(()) => println!(" Thread {}: Pushed {}", i, value),
Err(v) => println!(" Thread {}: stack full, dropped {}", i, v),
}
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!(" Popping values from stack:");
for _ in 0..5 {
if let Some(value) = stack.pop() {
println!(" Popped: {}", value);
}
}
println!("\n4. Lock-Free Queue Example:");
let queue: Arc<LockFreeQueue<i32>> = Arc::new(LockFreeQueue::new());
let queue_producer = Arc::clone(&queue);
let producer = thread::spawn(move || {
for i in 1..=10 {
queue_producer.enqueue(i);
println!(" Producer: Enqueued {}", i);
thread::sleep(Duration::from_millis(10));
}
});
let queue_consumer = Arc::clone(&queue);
let consumer = thread::spawn(move || {
thread::sleep(Duration::from_millis(50)); let mut consumed = 0;
while consumed < 10 {
if let Some(value) = queue_consumer.try_dequeue() {
println!(" Consumer: Dequeued {}", value);
consumed += 1;
} else {
thread::sleep(Duration::from_millis(20));
}
}
});
producer.join().unwrap();
consumer.join().unwrap();
println!("\n5. Barrier Example:");
let barrier = Arc::new(Barrier::new(3));
let mut handles = vec![];
for i in 0..3 {
let barrier_clone = Arc::clone(&barrier);
let handle = thread::spawn(move || {
println!(" Thread {}: Working...", i);
thread::sleep(Duration::from_millis((i + 1) as u64 * 100));
println!(" Thread {}: Waiting at barrier", i);
barrier_clone.wait();
println!(" Thread {}: Passed barrier!", i);
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
println!("\n6. WaitGroup Example:");
let wg = Arc::new(WaitGroup::new());
for i in 0..5 {
wg.add(1);
let wg_clone = Arc::clone(&wg);
thread::spawn(move || {
println!(" Worker {}: Starting task", i);
thread::sleep(Duration::from_millis((5 - i) as u64 * 50));
println!(" Worker {}: Task complete", i);
wg_clone.done();
});
}
println!(" Main: Waiting for all workers...");
wg.wait();
println!(" Main: All workers completed!");
println!("\n7. Performance Comparison (FutexMutex vs std::sync::Mutex):");
const BENCHMARK_ITERATIONS: usize = 100_000;
const BENCHMARK_WORKER_THREADS: usize = 4;
let iterations = BENCHMARK_ITERATIONS;
let futex_mutex = Arc::new(FutexMutex::new(0));
let start = Instant::now();
let mut handles = vec![];
for _ in 0..BENCHMARK_WORKER_THREADS {
let mutex_clone = Arc::clone(&futex_mutex);
let handle = thread::spawn(move || {
for _ in 0..iterations / BENCHMARK_WORKER_THREADS {
let mut guard = mutex_clone.lock();
*guard += 1;
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
let futex_time = start.elapsed();
let std_mutex = Arc::new(std::sync::Mutex::new(0));
let start = Instant::now();
let mut handles = vec![];
for _ in 0..BENCHMARK_WORKER_THREADS {
let mutex_clone = Arc::clone(&std_mutex);
let handle = thread::spawn(move || {
for _ in 0..iterations / BENCHMARK_WORKER_THREADS {
let mut guard = mutex_clone.lock().unwrap();
*guard += 1;
}
});
handles.push(handle);
}
for handle in handles {
handle.join().unwrap();
}
let std_time = start.elapsed();
println!(
" FutexMutex: {:?} ({} ops)",
futex_time,
*futex_mutex.lock()
);
println!(
" std::sync::Mutex: {:?} ({} ops)",
std_time,
*std_mutex.lock().unwrap()
);
println!(
" Speedup: {:.2}x",
std_time.as_secs_f64() / futex_time.as_secs_f64()
);
}