use std::net::SocketAddr;
use std::time::Duration;
use arcature::prelude::*;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
const HANG_GUARD: Duration = Duration::from_secs(10);
async fn http_get(addr: SocketAddr, path: &str) -> String {
let mut stream = tokio::time::timeout(HANG_GUARD, TcpStream::connect(addr))
.await
.expect("connect did not hang")
.expect("connect succeeds");
let request = format!("GET {path} HTTP/1.1\r\nHost: localhost\r\nConnection: close\r\n\r\n");
stream
.write_all(request.as_bytes())
.await
.expect("write request");
let mut buffer = Vec::new();
tokio::time::timeout(HANG_GUARD, stream.read_to_end(&mut buffer))
.await
.expect("read did not hang")
.expect("read succeeds");
String::from_utf8_lossy(&buffer).into_owned()
}
#[tokio::test]
async fn run_with_lifecycle_no_subsystems_serves_and_shuts_down() {
let app = Application::new()
.routes(Routes::new().route("/", get(|| async { "lifecycle" })))
.bind("127.0.0.1")
.port(0)
.build();
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let (shutdown_tx, shutdown_rx) = tokio::sync::oneshot::channel::<()>();
let join = tokio::spawn(async move {
app.serve_with_lifecycle(listener, |_| (), async {
let _ = shutdown_rx.await;
})
.await
.expect("serve_with_lifecycle");
});
let response = http_get(addr, "/").await;
assert!(response.starts_with("HTTP/1.1 200 OK"), "{}", response);
let (_, body) = response.split_once("\r\n\r\n").unwrap();
assert_eq!(body, "lifecycle");
shutdown_tx.send(()).expect("shutdown");
tokio::time::timeout(HANG_GUARD, join)
.await
.expect("no hang")
.expect("no panic");
}
macro_rules! require_db {
() => {
if std::env::var("ARCATURE_TEST_DB_URL").is_err() {
eprintln!("skip: ARCATURE_TEST_DB_URL not set");
return;
}
};
}
#[cfg(all(feature = "db", feature = "jobs"))]
#[tokio::test]
async fn single_pgpool_shared_between_db_and_jobs() {
require_db!();
let db_url = std::env::var("ARCATURE_TEST_DB_URL").expect("checked by require_db");
let db = arcature_db::Db::connect(
arcature_db::DbConfig::new(&db_url)
.expect("valid db config")
.application_name("lifecycle-test"),
)
.await
.expect("db connect");
let jobs = arcature_jobs::Jobs::new(db.sqlx().clone());
assert!(
!db.sqlx().is_closed(),
"db pool should be open before close"
);
db.close().await;
assert!(
jobs.pool().is_closed(),
"jobs pool must be closed after db.close() — single-pool invariant"
);
}
#[cfg(feature = "db")]
#[tokio::test]
async fn db_ping_after_close_returns_closed() {
require_db!();
let db_url = std::env::var("ARCATURE_TEST_DB_URL").expect("checked by require_db");
let db = arcature_db::Db::connect(
arcature_db::DbConfig::new(&db_url)
.expect("valid db config")
.application_name("lifecycle-test-ping"),
)
.await
.expect("db connect");
db.ping().await.expect("ping succeeds before close");
db.close().await;
}
#[cfg(all(feature = "db", feature = "jobs"))]
#[tokio::test]
async fn worker_drains_before_pool_closes() {
require_db!();
let db_url = std::env::var("ARCATURE_TEST_DB_URL").expect("checked by require_db");
let db = arcature_db::Db::connect(
arcature_db::DbConfig::new(&db_url)
.expect("valid db config")
.application_name("lifecycle-test-drain"),
)
.await
.expect("db connect");
let jobs = arcature_jobs::Jobs::new(db.sqlx().clone());
jobs.migrate().await.expect("jobs migrate");
let registry = arcature_jobs::Registry::new();
let worker = arcature_jobs::Worker::builder(db.sqlx().clone(), registry)
.config(arcature_jobs::WorkerConfig::default())
.build();
let shutdown = tokio_util::sync::CancellationToken::new();
let worker_shutdown = shutdown.clone();
let join = tokio::spawn(async move { worker.run(worker_shutdown).await });
shutdown.cancel();
let worker_result = tokio::time::timeout(HANG_GUARD, join)
.await
.expect("worker did not hang on shutdown");
assert!(
worker_result.is_ok(),
"worker join must succeed after cancel"
);
let inner = worker_result.expect("join ok");
assert!(
inner.is_ok(),
"worker run must return Ok(()) on graceful shutdown, got: {:?}",
inner.err()
);
db.close().await;
}