use crate::common;
use camber::RuntimeError;
use camber::http::{self, Request, Response, Router};
use camber::runtime;
use std::sync::Arc;
use std::time::{Duration, Instant};
const EVENT_TIMEOUT: Duration = Duration::from_secs(2);
const CONCURRENT_REQUESTS: usize = 3;
const DRAIN_TIMEOUT: Duration = Duration::from_millis(200);
#[test]
fn builder_configures_concurrent_requests() {
runtime::builder()
.worker_threads(2)
.keepalive_timeout(Duration::from_millis(100))
.shutdown_timeout(Duration::from_secs(2))
.run(|| {
let barrier = Arc::new(tokio::sync::Barrier::new(CONCURRENT_REQUESTS));
let mut router = Router::new();
router.get("/slow", move |_req: &Request| {
let handler_barrier = Arc::clone(&barrier);
async move {
handler_barrier.wait().await;
Response::text(200, "done")
}
});
let server = common::spawn_server_ready(router, EVENT_TIMEOUT).unwrap();
let addr = server.local_addr();
let url = format!("http://{addr}/slow");
let responses = common::block_on(async {
tokio::time::timeout(
EVENT_TIMEOUT,
futures_util::future::join_all(
(0..CONCURRENT_REQUESTS).map(|_| http::get(&url)),
),
)
.await
})
.expect("the server never held all three requests in flight at once");
assert_eq!(
responses.len(),
CONCURRENT_REQUESTS,
"the barrier released without one response per request"
);
responses.into_iter().for_each(|response| {
assert_eq!(response.unwrap().status(), 200);
});
server.shutdown_bounded(EVENT_TIMEOUT).unwrap();
runtime::request_shutdown();
})
.unwrap();
}
#[test]
fn builder_configures_shutdown_timeout() {
let wedged = common::WedgedHandle::new();
let closure_wedged = wedged.clone();
let outcome = runtime::builder()
.shutdown_timeout(DRAIN_TIMEOUT)
.run(move || {
let (entered_tx, entered_rx) = tokio::sync::oneshot::channel();
let task = camber::spawn_async(async move {
let _ = entered_tx.send(());
std::future::pending::<()>().await;
});
common::block_on(common::join_bounded(entered_rx, EVENT_TIMEOUT)).unwrap();
runtime::request_shutdown();
closure_wedged.record((task, Instant::now()));
});
assert!(
matches!(outcome, Err(RuntimeError::ScopeDrainTimeout(1))),
"the safety-net timeout did not report the wedged child: {outcome:?}"
);
let (task, drain_started) = wedged.take();
let elapsed = drain_started.elapsed();
assert!(
elapsed >= DRAIN_TIMEOUT && elapsed < Duration::from_secs(1),
"expected the configured {DRAIN_TIMEOUT:?} shutdown timeout to drive escalation, got {elapsed:?}"
);
let joined = common::block_on_detached(common::join_bounded(task, EVENT_TIMEOUT));
common::assert_forced_abort(&joined);
}