use ranked_semaphore::{PriorityConfig, QueueStrategy, RankedSemaphore, TryAcquireError};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::Duration;
use tokio::time::{sleep, timeout};
#[tokio::test]
async fn test_priority_basic_ordering() {
let sem = Arc::new(RankedSemaphore::new_fifo(1));
let _hold_permit = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
let order1 = Arc::clone(&execution_order);
let sem1 = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem1.acquire_with_priority(-10).await.unwrap();
order1.store(1, Ordering::Relaxed);
}));
let order2 = Arc::clone(&execution_order);
let sem2 = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem2.acquire_with_priority(10).await.unwrap();
order2.store(2, Ordering::Relaxed);
}));
let order3 = Arc::clone(&execution_order);
let sem3 = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem3.acquire_with_priority(5).await.unwrap();
order3.store(3, Ordering::Relaxed);
}));
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
let final_order = execution_order.load(Ordering::Relaxed);
assert!(final_order > 0, "At least one task should have completed");
}
#[tokio::test]
async fn test_priority_try_acquire_basic() {
let sem = RankedSemaphore::new_fifo(3);
let permit1 = sem.try_acquire().unwrap();
let permit2 = sem.try_acquire().unwrap();
let permit3 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
let result = sem.try_acquire();
assert!(matches!(result, Err(TryAcquireError::NoPermits)));
drop(permit1);
drop(permit2);
drop(permit3);
}
#[tokio::test]
async fn test_priority_acquire_many_with_priority() {
let sem = RankedSemaphore::new_fifo(5);
let permit = sem.try_acquire_many(3).unwrap();
assert_eq!(permit.num_permits(), 3);
assert_eq!(sem.available_permits(), 2);
let permit2 = timeout(
Duration::from_millis(50),
sem.acquire_many_with_priority(-1, 2),
)
.await
.unwrap()
.unwrap();
assert_eq!(permit2.num_permits(), 2);
assert_eq!(sem.available_permits(), 0);
drop(permit);
drop(permit2);
}
#[tokio::test]
async fn test_queue_strategy_fifo() {
let sem = Arc::new(RankedSemaphore::new_fifo(1));
let _hold_permit = sem.try_acquire().unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for i in 1..=3 {
let order = Arc::clone(&execution_order);
let sem_clone = Arc::clone(&sem);
let task_id = i;
handles.push(tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(0).await.unwrap();
let _ = order.compare_exchange(0, task_id, Ordering::Relaxed, Ordering::Relaxed);
}));
}
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
sleep(Duration::from_millis(10)).await;
assert_eq!(execution_order.load(Ordering::Relaxed), 1);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
}
#[tokio::test]
async fn test_queue_strategy_lifo() {
let sem = Arc::new(RankedSemaphore::new_lifo(1));
let _hold_permit = sem.try_acquire().unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for i in 1..=3 {
let order = Arc::clone(&execution_order);
let sem_clone = Arc::clone(&sem);
let task_id = i;
handles.push(tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(0).await.unwrap();
let _ = order.compare_exchange(0, task_id, Ordering::Relaxed, Ordering::Relaxed);
}));
sleep(Duration::from_millis(1)).await;
}
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
sleep(Duration::from_millis(10)).await;
assert_eq!(execution_order.load(Ordering::Relaxed), 3);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
}
#[tokio::test]
async fn test_priority_config_exact_priority() {
let config = PriorityConfig::new()
.default_strategy(QueueStrategy::Fifo)
.exact(10, QueueStrategy::Lifo);
let sem = Arc::new(RankedSemaphore::new_with_config(1, config));
let _hold_permit = sem.try_acquire().unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for i in 1..=2 {
let order = Arc::clone(&execution_order);
let sem_clone = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(10).await.unwrap();
let _ = order.compare_exchange(0, i, Ordering::Relaxed, Ordering::Relaxed);
}));
sleep(Duration::from_millis(1)).await;
}
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
sleep(Duration::from_millis(10)).await;
assert_eq!(execution_order.load(Ordering::Relaxed), 2);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
}
#[tokio::test]
async fn test_priority_config_range() {
let config = PriorityConfig::new()
.default_strategy(QueueStrategy::Fifo)
.range(5, 15, QueueStrategy::Lifo);
let sem = RankedSemaphore::new_with_config(5, config);
let permit1 = sem.try_acquire().unwrap(); let permit2 = sem.try_acquire().unwrap(); let permit3 = sem.try_acquire().unwrap(); let permit4 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 1);
drop(permit1);
drop(permit2);
drop(permit3);
drop(permit4);
}
#[tokio::test]
async fn test_priority_config_greater_or_equal() {
let config = PriorityConfig::new()
.default_strategy(QueueStrategy::Fifo)
.greater_or_equal(10, QueueStrategy::Lifo);
let sem = RankedSemaphore::new_with_config(3, config);
let permit1 = sem.try_acquire().unwrap(); let permit2 = sem.try_acquire().unwrap(); let permit3 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
drop(permit1);
drop(permit2);
drop(permit3);
}
#[tokio::test]
async fn test_priority_config_less_or_equal() {
let config = PriorityConfig::new()
.default_strategy(QueueStrategy::Fifo)
.less_or_equal(0, QueueStrategy::Lifo);
let sem = RankedSemaphore::new_with_config(3, config);
let permit1 = sem.try_acquire().unwrap(); let permit2 = sem.try_acquire().unwrap(); let permit3 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
drop(permit1);
drop(permit2);
drop(permit3);
}
#[tokio::test]
async fn test_priority_config_multiple_rules() {
let config = PriorityConfig::new()
.default_strategy(QueueStrategy::Fifo)
.exact(0, QueueStrategy::Lifo)
.range(10, 20, QueueStrategy::Lifo)
.greater_or_equal(100, QueueStrategy::Lifo);
let sem = RankedSemaphore::new_with_config(5, config);
let permit1 = sem.try_acquire().unwrap(); let permit2 = sem.try_acquire().unwrap(); let permit3 = sem.try_acquire().unwrap(); let permit4 = sem.try_acquire().unwrap(); let permit5 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
drop(permit1);
drop(permit2);
drop(permit3);
drop(permit4);
drop(permit5);
}
#[tokio::test]
async fn test_extreme_priority_values() {
let sem = RankedSemaphore::new_fifo(3);
let permit1 = sem.try_acquire().unwrap();
let permit2 = sem.try_acquire().unwrap();
let permit3 = sem.try_acquire().unwrap();
assert_eq!(sem.available_permits(), 0);
drop(permit1);
drop(permit2);
drop(permit3);
}
#[tokio::test]
async fn test_priority_with_many_permits() {
let sem = Arc::new(RankedSemaphore::new_fifo(1));
let _hold_permit = sem.try_acquire().unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
let order1 = Arc::clone(&execution_order);
let sem1 = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
sem1.add_permits(2); let _permit = sem1.acquire_many_with_priority(2, 10).await.unwrap();
order1.store(1, Ordering::Relaxed);
}));
let order2 = Arc::clone(&execution_order);
let sem2 = Arc::clone(&sem);
handles.push(tokio::spawn(async move {
let _permit = sem2.acquire_with_priority(5).await.unwrap();
order2.store(2, Ordering::Relaxed);
}));
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
let final_order = execution_order.load(Ordering::Relaxed);
assert!(final_order > 0, "At least one task should have completed");
}
#[tokio::test]
async fn test_priority_ordering_stress() {
let sem = Arc::new(RankedSemaphore::new_fifo(1));
let _hold_permit = sem.try_acquire().unwrap();
let mut handles = Vec::new();
let results = Arc::new(std::sync::Mutex::new(Vec::new()));
let priorities = vec![1, 5, 3, 10, 2, 8, 4, 9, 6, 7];
for (index, priority) in priorities.into_iter().enumerate() {
let sem_clone = Arc::clone(&sem);
let results_clone = Arc::clone(&results);
handles.push(tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(priority).await.unwrap();
results_clone.lock().unwrap().push((index, priority));
}));
sleep(Duration::from_millis(1)).await;
}
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
for _ in 0..3 {
sleep(Duration::from_millis(10)).await;
sem.add_permits(1);
}
sleep(Duration::from_millis(20)).await;
let priority_completed = {
let final_results = results.lock().unwrap();
if !final_results.is_empty() {
final_results[0].1
} else {
0
}
};
if priority_completed > 0 {
println!("First completed priority: {priority_completed}");
}
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
}
#[tokio::test]
async fn test_mixed_acquire_types_with_priority() {
let sem = Arc::new(RankedSemaphore::new_fifo(5));
let permit1 = sem.try_acquire().unwrap();
let permit2 = timeout(Duration::from_millis(50), sem.acquire_with_priority(3))
.await
.unwrap()
.unwrap();
let permit3 = sem.try_acquire_many(2).unwrap();
assert_eq!(sem.available_permits(), 1);
assert_eq!(permit1.num_permits(), 1);
assert_eq!(permit2.num_permits(), 1);
assert_eq!(permit3.num_permits(), 2);
drop(permit1);
drop(permit2);
drop(permit3);
assert_eq!(sem.available_permits(), 5);
}
#[tokio::test]
async fn test_priority_with_closed_semaphore() {
let sem = RankedSemaphore::new_fifo(1);
sem.close();
let result1 = sem.try_acquire();
assert!(matches!(result1, Err(TryAcquireError::Closed)));
let result2 = sem.try_acquire_many(1);
assert!(matches!(result2, Err(TryAcquireError::Closed)));
}
#[tokio::test]
async fn test_priority_zero_and_negative() {
let sem = Arc::new(RankedSemaphore::new_fifo(1));
let _hold_permit = sem.try_acquire().unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = Vec::new();
for (i, priority) in [0, -1, -10, -100].iter().enumerate() {
let order = Arc::clone(&execution_order);
let sem_clone = Arc::clone(&sem);
let task_id = i + 1;
let prio = *priority;
handles.push(tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(prio).await.unwrap();
let _ = order.compare_exchange(0, task_id, Ordering::Relaxed, Ordering::Relaxed);
}));
}
sleep(Duration::from_millis(10)).await;
drop(_hold_permit);
sleep(Duration::from_millis(10)).await;
assert_eq!(execution_order.load(Ordering::Relaxed), 1);
for handle in handles {
let _ = timeout(Duration::from_millis(100), handle).await;
}
}
#[tokio::test]
async fn test_greater_than_convenience_method() {
let config = PriorityConfig::new()
.greater_than(0, QueueStrategy::Fifo)
.default_strategy(QueueStrategy::Lifo);
let sem = Arc::new(RankedSemaphore::new_with_config(1, config));
let _permit = sem.acquire_with_priority(10).await.unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = vec![];
for i in 0..3 {
let sem_clone = Arc::clone(&sem);
let order_clone = Arc::clone(&execution_order);
let handle = tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(1).await.unwrap();
let order = order_clone.fetch_add(1, Ordering::Relaxed);
(i, order)
});
handles.push(handle);
tokio::time::sleep(Duration::from_millis(1)).await; }
for i in 3..6 {
let sem_clone = Arc::clone(&sem);
let order_clone = Arc::clone(&execution_order);
let handle = tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(0).await.unwrap();
let order = order_clone.fetch_add(1, Ordering::Relaxed);
(i, order)
});
handles.push(handle);
tokio::time::sleep(Duration::from_millis(1)).await; }
tokio::time::sleep(Duration::from_millis(10)).await;
drop(_permit);
let mut results = vec![];
for handle in handles {
if let Ok(Ok(result)) = timeout(Duration::from_millis(100), handle).await {
results.push(result);
}
}
results.sort_by_key(|&(_, order)| order);
let task_ids: Vec<usize> = results.into_iter().map(|(id, _)| id).collect();
assert_eq!(task_ids, vec![0, 1, 2, 5, 4, 3]);
}
#[tokio::test]
async fn test_less_than_convenience_method() {
let config = PriorityConfig::new()
.less_than(5, QueueStrategy::Fifo)
.default_strategy(QueueStrategy::Lifo);
let sem = Arc::new(RankedSemaphore::new_with_config(1, config));
let _permit = sem.acquire_with_priority(10).await.unwrap();
let execution_order = Arc::new(AtomicUsize::new(0));
let mut handles = vec![];
for i in 0..3 {
let sem_clone = Arc::clone(&sem);
let order_clone = Arc::clone(&execution_order);
let handle = tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(4).await.unwrap();
let order = order_clone.fetch_add(1, Ordering::Relaxed);
(i, order)
});
handles.push(handle);
tokio::time::sleep(Duration::from_millis(1)).await; }
for i in 3..6 {
let sem_clone = Arc::clone(&sem);
let order_clone = Arc::clone(&execution_order);
let handle = tokio::spawn(async move {
let _permit = sem_clone.acquire_with_priority(5).await.unwrap();
let order = order_clone.fetch_add(1, Ordering::Relaxed);
(i, order)
});
handles.push(handle);
tokio::time::sleep(Duration::from_millis(1)).await; }
tokio::time::sleep(Duration::from_millis(10)).await;
drop(_permit);
let mut results = vec![];
for handle in handles {
if let Ok(Ok(result)) = timeout(Duration::from_millis(100), handle).await {
results.push(result);
}
}
results.sort_by_key(|&(_, order)| order);
let task_ids: Vec<usize> = results.into_iter().map(|(id, _)| id).collect();
assert_eq!(task_ids, vec![5, 4, 3, 0, 1, 2]);
}
#[tokio::test]
async fn test_edge_cases_for_greater_and_less_than() {
let config = PriorityConfig::new()
.greater_than(isize::MAX - 1, QueueStrategy::Fifo)
.less_than(isize::MIN + 1, QueueStrategy::Lifo)
.default_strategy(QueueStrategy::Fifo);
let sem = Arc::new(RankedSemaphore::new_with_config(3, config));
let _permit1 = sem.acquire_with_priority(isize::MAX).await.unwrap();
let _permit2 = sem.acquire_with_priority(isize::MIN).await.unwrap();
let _permit3 = sem.acquire_with_priority(0).await.unwrap();
assert_eq!(sem.available_permits(), 0);
}