use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
async fn await_waiters(lane: &EndpointLane, n: usize) {
while lane.waiter_count() != n {
tokio::task::yield_now().await;
}
}
#[test]
fn client_id_accepts_bounded_identity() {
assert_eq!(ClientId::parse("agent-7").as_str(), "agent-7");
assert_eq!(ClientId::parse(" tenant.a:1 ").as_str(), "tenant.a:1");
}
#[test]
fn client_id_rejects_invalid_and_oversized() {
assert_eq!(ClientId::parse("").as_str(), ClientId::DEFAULT);
assert_eq!(ClientId::parse("has space").as_str(), ClientId::DEFAULT);
assert_eq!(ClientId::parse("bad/slash").as_str(), ClientId::DEFAULT);
let oversized = "x".repeat(ClientId::MAX_LEN + 1);
assert_eq!(ClientId::parse(&oversized).as_str(), ClientId::DEFAULT);
assert_eq!(ClientId::from_header(None).as_str(), ClientId::DEFAULT);
}
#[tokio::test]
async fn unlimited_lane_admits_immediately() {
let lane = EndpointLane::unlimited();
let permit = lane.admit("default").await.unwrap();
assert!(permit.limited.is_none());
}
#[tokio::test]
async fn concurrency_one_serializes_two_admits() {
let lane = EndpointLane::new(1, &QueueConfig::default());
let phase = Arc::new(AtomicUsize::new(0));
let first = lane.admit("a").await.unwrap();
phase.store(1, Ordering::SeqCst);
let lane2 = lane.clone();
let phase2 = Arc::clone(&phase);
let second = tokio::spawn(async move {
phase2.store(2, Ordering::SeqCst);
let permit = lane2.admit("a").await.unwrap();
phase2.store(3, Ordering::SeqCst);
permit
});
await_waiters(&lane, 1).await;
assert_eq!(
phase.load(Ordering::SeqCst),
2,
"second admit still waiting"
);
drop(first);
let _second_permit = second.await.unwrap();
assert_eq!(phase.load(Ordering::SeqCst), 3);
}
#[tokio::test]
async fn queue_full_rejects_when_max_depth_waiters_present() {
let lane = EndpointLane::new(
1,
&QueueConfig {
max_depth: 1,
fair_scheduling: true,
},
);
let inflight = lane.admit("a").await.unwrap();
let lane_wait = lane.clone();
let waiting = tokio::spawn(async move { lane_wait.admit("b").await });
await_waiters(&lane, 1).await;
let err = lane.admit("c").await.unwrap_err();
assert_eq!(err, AdmitError::QueueFull);
drop(inflight);
let _waiting_permit = waiting.await.unwrap().unwrap();
}
#[tokio::test]
async fn fair_scheduling_interleaves_clients() {
let lane = EndpointLane::new(
1,
&QueueConfig {
max_depth: 10,
fair_scheduling: true,
},
);
let held = lane.admit("A").await.unwrap();
let order = Arc::new(Mutex::new(Vec::new()));
let mut handles = Vec::new();
for (registered, key) in ["A", "A", "B"].into_iter().enumerate() {
let lane_task = lane.clone();
let order = Arc::clone(&order);
let key = key.to_string();
handles.push(tokio::spawn(async move {
let permit = lane_task.admit(&key).await.unwrap();
order
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(key);
drop(permit);
}));
await_waiters(&lane, registered + 1).await;
}
drop(held);
for handle in handles {
handle.await.unwrap();
}
let got = order
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.clone();
assert_eq!(got, vec!["A", "B", "A"]);
}
#[tokio::test]
async fn cancel_while_queued_frees_the_waiter_slot() {
use std::future::{Future, poll_fn};
use std::task::Poll;
let lane = EndpointLane::new(
1,
&QueueConfig {
max_depth: 10,
fair_scheduling: true,
},
);
let held = lane.admit("a").await.unwrap();
let mut waiting = Box::pin(lane.admit("b"));
poll_fn(|cx| match waiting.as_mut().poll(cx) {
Poll::Ready(_) => panic!("waiter should be queued, not ready"),
Poll::Pending => Poll::Ready(()),
})
.await;
await_waiters(&lane, 1).await;
drop(waiting);
await_waiters(&lane, 0).await;
let lane_next = lane.clone();
let next = tokio::spawn(async move { lane_next.admit("c").await });
await_waiters(&lane, 1).await;
drop(held);
let _permit = next.await.unwrap().unwrap();
}
#[test]
fn queue_full_and_unavailable_are_distinct() {
assert_ne!(AdmitError::QueueFull, AdmitError::Unavailable);
assert_ne!(
AdmitError::QueueFull.to_string(),
AdmitError::Unavailable.to_string()
);
}
#[tokio::test]
async fn cancel_after_wake_does_not_leak_slot() {
use std::future::{Future, poll_fn};
use std::task::Poll;
let lane = EndpointLane::new(
1,
&QueueConfig {
max_depth: 10,
fair_scheduling: true,
},
);
let held = lane.admit("a").await.unwrap();
let mut waiting = Box::pin(lane.admit("b"));
poll_fn(|cx| match waiting.as_mut().poll(cx) {
Poll::Ready(_) => panic!("waiter should be queued, not ready"),
Poll::Pending => Poll::Ready(()),
})
.await;
drop(held);
drop(waiting);
let permit = tokio::time::timeout(Duration::from_millis(500), lane.admit("c"))
.await
.expect("slot leaked: subsequent admit timed out")
.unwrap();
drop(permit);
}
#[test]
fn new_is_total_and_never_panics_on_zero_settings() {
let unlimited = EndpointLane::new(0, &QueueConfig::default());
assert!(matches!(unlimited.inner, LaneInner::Unlimited));
let clamped = EndpointLane::new(
1,
&QueueConfig {
max_depth: 0,
fair_scheduling: true,
},
);
assert!(matches!(clamped.inner, LaneInner::Limited(_)));
}
#[tokio::test]
async fn distinct_client_labels_are_capped_to_bound_fair_scheduling() {
let lane = EndpointLane::new(
1,
&QueueConfig {
max_depth: 1000,
fair_scheduling: true,
},
);
let _held = lane.admit("holder").await.unwrap();
let total = MAX_DISTINCT_CLIENTS + 10;
let mut handles = Vec::new();
for i in 0..total {
let lane = lane.clone();
handles.push(tokio::spawn(async move {
let _ = lane.admit(&format!("client-{i}")).await;
}));
}
await_waiters(&lane, total).await;
assert!(
lane.distinct_clients() <= MAX_DISTINCT_CLIENTS + 1,
"distinct client buckets {} exceeded bound {}",
lane.distinct_clients(),
MAX_DISTINCT_CLIENTS + 1
);
for handle in handles {
handle.abort();
}
}