#![forbid(unsafe_code)]
use super::*;
use serial_test::serial;
use std::sync::atomic::AtomicUsize;
use std::time::Duration;
#[test]
fn auto_limit_is_clamped() {
let n = auto_limit();
assert!(n >= MIN_CONCURRENCY);
assert!(n <= HARD_CAP);
}
#[test]
fn resolve_cli_override_clamps() {
assert_eq!(resolve_limit(Some(0)), MIN_CONCURRENCY);
assert_eq!(resolve_limit(Some(1)), 1);
assert_eq!(resolve_limit(Some(9999)), HARD_CAP);
}
#[test]
fn worker_threads_sane() {
let w = worker_threads();
assert!((2..=16).contains(&w));
}
fn reset_signal_flags() {
crate::signals::cancellation_flag().store(false, Ordering::Release);
crate::signals::sigterm_flag().store(false, Ordering::Release);
crate::signals::force_exit_flag().store(false, Ordering::Release);
}
#[tokio::test(flavor = "current_thread")]
#[serial]
async fn map_bounded_respects_peak() {
reset_signal_flags();
let limit = 3usize;
let items: Vec<u32> = (0..20).collect();
let current = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let cur_c = Arc::clone(¤t);
let peak_c = Arc::clone(&peak);
let results = map_bounded(items, limit, move |n| {
let cur_c = Arc::clone(&cur_c);
let peak_c = Arc::clone(&peak_c);
async move {
let now = cur_c.fetch_add(1, Ordering::SeqCst) + 1;
peak_c.fetch_max(now, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(15)).await;
cur_c.fetch_sub(1, Ordering::SeqCst);
n * 2
}
})
.await;
assert_eq!(results.len(), 20);
let observed = peak.load(Ordering::SeqCst);
assert!(observed <= limit, "peak={observed} limit={limit}");
assert!(observed >= 1);
for (i, r) in results.iter().enumerate() {
assert_eq!(r.index, i);
assert_eq!(r.outcome.as_ref().unwrap(), &(i as u32 * 2));
}
}
#[tokio::test(flavor = "current_thread")]
#[serial]
async fn permit_released_after_panic_in_task() {
reset_signal_flags();
let limit = 2usize;
let progressed = Arc::new(AtomicUsize::new(0));
let current = Arc::new(AtomicUsize::new(0));
let peak = Arc::new(AtomicUsize::new(0));
let p = Arc::clone(&progressed);
let cur_c = Arc::clone(¤t);
let peak_c = Arc::clone(&peak);
let items: Vec<u32> = (0..6).collect();
let results = map_bounded(items, limit, move |n| {
let p = Arc::clone(&p);
let cur_c = Arc::clone(&cur_c);
let peak_c = Arc::clone(&peak_c);
async move {
struct Guard(Arc<AtomicUsize>);
impl Drop for Guard {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::SeqCst);
}
}
let now = cur_c.fetch_add(1, Ordering::SeqCst) + 1;
peak_c.fetch_max(now, Ordering::SeqCst);
let _g = Guard(Arc::clone(&cur_c));
p.fetch_add(1, Ordering::SeqCst);
if n == 1 {
panic!("boom");
}
tokio::time::sleep(Duration::from_millis(5)).await;
n
}
})
.await;
assert_eq!(results.len(), 6);
let panics = results.iter().filter(|r| r.outcome.is_err()).count();
assert_eq!(panics, 1);
let panic_row = results
.iter()
.find(|r| r.outcome.is_err())
.expect("one panic");
assert_eq!(panic_row.index, 1, "panic must preserve input index");
assert!(progressed.load(Ordering::SeqCst) >= 5);
let observed = peak.load(Ordering::SeqCst);
assert!(observed <= limit, "peak={observed} limit={limit}");
assert_eq!(current.load(Ordering::SeqCst), 0);
}
#[tokio::test(flavor = "current_thread")]
#[serial]
async fn panic_preserves_input_index_not_usize_max() {
reset_signal_flags();
let items: Vec<u32> = vec![10, 20, 30];
let results = map_bounded(items, 2, |n| async move {
if n == 20 {
panic!("mid");
}
n
})
.await;
assert_eq!(results.len(), 3);
let panic_row = results
.iter()
.find(|r| r.outcome.is_err())
.expect("panic");
assert_ne!(panic_row.index, usize::MAX);
assert_eq!(panic_row.index, 1);
assert!(results[0].outcome.is_ok());
assert!(results[2].outcome.is_ok());
}
#[tokio::test]
async fn semaphore_acquire_owned_raii() {
let sem = semaphore(1);
let p1 = acquire_owned(&sem).await;
assert_eq!(sem.available_permits(), 0);
drop(p1);
assert_eq!(sem.available_permits(), 1);
}
#[tokio::test(flavor = "current_thread")]
#[serial]
async fn map_bounded_stops_admission_on_should_stop() {
reset_signal_flags();
let started = Arc::new(AtomicUsize::new(0));
let started_c = Arc::clone(&started);
let items: Vec<u32> = (0..20).collect();
let limit = 2usize;
let arm = tokio::spawn(async {
tokio::time::sleep(Duration::from_millis(20)).await;
crate::signals::cancellation_flag().store(true, Ordering::Release);
});
let results = map_bounded(items, limit, move |_n| {
let started_c = Arc::clone(&started_c);
async move {
started_c.fetch_add(1, Ordering::SeqCst);
tokio::time::sleep(Duration::from_millis(40)).await;
1u32
}
})
.await;
let _ = arm.await;
reset_signal_flags();
let n_started = started.load(Ordering::SeqCst);
assert!(
n_started < 20,
"admission must stop after should_stop; started={n_started}"
);
assert!(n_started >= 1, "at least seed batch should start");
assert_eq!(results.len(), n_started);
assert!(results.iter().all(|r| r.outcome.is_ok()));
}
#[tokio::test(flavor = "current_thread")]
#[serial]
async fn map_bounded_force_aborts_inflight() {
reset_signal_flags();
let items: Vec<u32> = (0..8).collect();
let limit = 4usize;
let arm = tokio::spawn(async {
tokio::time::sleep(Duration::from_millis(15)).await;
crate::signals::cancellation_flag().store(true, Ordering::Release);
crate::signals::force_exit_flag().store(true, Ordering::Release);
});
let results = map_bounded(items, limit, move |_n| async move {
tokio::time::sleep(Duration::from_secs(30)).await;
1u32
})
.await;
let _ = arm.await;
reset_signal_flags();
assert!(!results.is_empty(), "some tasks should have been admitted");
let cancelled_or_done = results
.iter()
.filter(|r| match &r.outcome {
Ok(_) => true,
Err(e) => e.is_cancelled() || e.is_panic(),
})
.count();
assert_eq!(cancelled_or_done, results.len());
let any_cancel = results.iter().any(|r| {
r.outcome
.as_ref()
.err()
.is_some_and(|e| e.is_cancelled())
});
assert!(
any_cancel || results.len() < 8,
"force_exit should abort in-flight or stop admission early"
);
}