use crate::standard_pool::*;
use saddle_admission::*;
use std::{future::Future, pin::Pin, sync::{Arc, atomic::{AtomicUsize, Ordering}}, task::{Context, Poll}, time::Duration};
fn process(capacity: u32) -> ProfuseGwLightweightProcessOwner {
let pending = freeze_deployment_resource_budget(
1, 32, 5000, 1, capacity, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
)
.unwrap();
let (app, listener) = saddle_core::BootstrapRendezvousIssuer::issue()
.freeze_application(saddle_core::GeneratedApplicationFreezeSource::new(
"app",
b"descriptor",
&["route"],
))
.unwrap();
let listener = listener
.freeze_listener(saddle_core::ListenerStartupFreezeSource::new(
"app",
"127.0.0.1:8000".parse().unwrap(),
"127.0.0.1:9000".parse().unwrap(),
Duration::from_secs(5),
))
.ok()
.unwrap();
let (whole, receipt) = saddle_core::pair_bootstrap_rendezvous(app, listener)
.ok()
.unwrap();
prepare_profusegw_lightweight_profile(
bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
.ok()
.unwrap(),
)
.ok()
.unwrap()
}
fn execution(process: &ProfuseGwLightweightProcessOwner) -> ProfuseGwLightweightExecutionOwner {
let mut ingress = process.verified_profile().try_ingress().unwrap();
match process
.verified_profile()
.try_promote_ingress(&mut ingress, 0)
{
ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, _) => permit.into_execution(),
_ => panic!("pre-RPC requests must coexist"),
}
}
fn isolated(name: &str) -> bool {
if std::env::var("SADDLE_R1_CHILD").ok().as_deref() == Some(name) { return false; }
let result = std::process::Command::new(std::env::current_exe().unwrap())
.args(["--exact", &format!("standard_pool_tests::{name}"), "--nocapture"])
.env("SADDLE_R1_CHILD", name).output().unwrap();
assert!(result.status.success(), "isolated default consumer failed: {}\n{}",
String::from_utf8_lossy(&result.stdout), String::from_utf8_lossy(&result.stderr));
true
}
fn pool(process: &ProfuseGwLightweightProcessOwner) -> StandardConnectionPool {
let root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
StandardConnectionPool::new(root, process.rpc_stage_snapshot().capacity, vec![
FrozenConnectionTarget::new("alpha", "http://127.0.0.1:8001").unwrap(),
FrozenConnectionTarget::new("beta", "http://127.0.0.1:8002").unwrap(),
]).unwrap()
}
#[tokio::test]
async fn capacity_rejects_overlap_and_recovers_after_retirement() {
if isolated("capacity_rejects_overlap_and_recovers_after_retirement") { return; }
let process = process(1);
let pool = pool(&process);
let mut service = pool.start().unwrap();
let lease = pool.try_acquire(0).unwrap();
assert!(matches!(pool.try_acquire(0), Err(AdmissionError::CapacityRejected)));
assert!(matches!(pool.try_acquire(99), Err(AdmissionError::InvalidConfiguration)));
drop(lease);
pool.wait_retired().await;
assert!(pool.try_acquire(1).is_ok());
pool.stop_admission();
service.shutdown().await.unwrap();
}
#[tokio::test]
async fn old_generation_executor_and_waker_cannot_register_in_new_lease() {
if isolated("old_generation_executor_and_waker_cannot_register_in_new_lease") { return; }
let process = process(1);
let pool = pool(&process);
let mut service = pool.start().unwrap();
let old = pool.try_acquire(0).unwrap();
let old_generation = old.generation();
let executor = old.executor();
let wake = old.test_driver_waker();
drop(old);
pool.wait_retired().await;
let new = pool.try_acquire(0).unwrap();
assert_ne!(new.generation(), old_generation);
use hyper::rt::Executor;
executor.execute(Box::pin(async {}));
wake.wake();
assert_eq!(new.test_driver_count(), 0);
drop(new);
service.shutdown().await.unwrap();
}
struct PhysicalProbe {
dropped: Arc<AtomicUsize>,
process: Arc<ProfuseGwLightweightProcessOwner>,
}
impl Future for PhysicalProbe {
type Output = ();
fn poll(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll<()> { Poll::Pending }
}
impl Drop for PhysicalProbe {
fn drop(&mut self) {
assert_eq!(self.process.rpc_stage_snapshot().operations, 1,
"request credit must survive concrete driver destruction");
self.dropped.fetch_add(1, Ordering::SeqCst);
}
}
#[tokio::test]
async fn cancellation_destroys_driver_before_stage_credit_and_never_publishes_idle() {
if isolated("cancellation_destroys_driver_before_stage_credit_and_never_publishes_idle") { return; }
let process = Arc::new(process(1));
let execution = execution(&process);
let pool = pool(&process);
let baseline = process.resource_snapshot().framework_charged;
let mut service = pool.start().unwrap();
let mut lease = pool.try_acquire(0).unwrap();
lease.attach_stage(execution.try_begin_rpc_stage().unwrap()).unwrap();
let dropped = Arc::new(AtomicUsize::new(0));
use hyper::rt::Executor;
lease.executor().execute(Box::pin(PhysicalProbe { dropped: dropped.clone(), process: process.clone() }));
drop(lease);
pool.wait_retired().await;
assert_eq!(dropped.load(Ordering::SeqCst), 1);
assert_eq!(process.rpc_stage_snapshot().operations, 0);
assert_eq!(pool.test_idle_count(), 0);
assert!(process.resource_snapshot().framework_charged >= baseline);
service.shutdown().await.unwrap();
}
#[tokio::test]
async fn successful_response_is_required_for_idle_reuse() {
if isolated("successful_response_is_required_for_idle_reuse") { return; }
let process = process(1);
let pool = pool(&process);
let mut service = pool.start().unwrap();
let lease = pool.try_acquire(0).unwrap();
lease.finish(false);
pool.wait_retired().await;
assert_eq!(pool.test_idle_count(), 0);
pool.try_acquire(0).unwrap().finish(true);
pool.wait_retired().await;
assert_eq!(pool.test_idle_count(), 0);
service.shutdown().await.unwrap();
}
#[tokio::test]
async fn close_stops_admission_joins_and_releases_backing_after_last_handle() {
if isolated("close_stops_admission_joins_and_releases_backing_after_last_handle") { return; }
let process = process(1);
let baseline = process.resource_snapshot().framework_charged;
let pool = pool(&process);
let mut service = pool.start().unwrap();
let lease = pool.try_acquire(0).unwrap();
pool.stop_admission();
assert_eq!(process.rpc_stage_snapshot().standard_pool_owners, 1);
assert!(process.rpc_stage_snapshot().standard_pool_storage_bytes > 0);
assert!(matches!(pool.try_acquire(0), Err(AdmissionError::AccountClosed)));
service.shutdown().await.unwrap();
drop((lease, service, pool));
assert_eq!(process.resource_snapshot().framework_charged, baseline);
assert_eq!(process.rpc_stage_snapshot().standard_pool_owners, 0);
assert_eq!(process.rpc_stage_snapshot().standard_pool_storage_bytes, 0);
}
#[test]
fn startup_without_runtime_and_storage_rejection_roll_back() {
if isolated("startup_without_runtime_and_storage_rejection_roll_back") { return; }
let process = process(1);
let baseline = process.resource_snapshot().framework_charged;
let pool = pool(&process);
assert!(pool.start().is_err());
drop(pool);
assert_eq!(process.resource_snapshot().framework_charged, baseline);
let root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
let before = process.resource_snapshot().framework_charged;
assert!(StandardConnectionPool::new(root, usize::MAX, vec![]).is_err());
assert!(process.resource_snapshot().framework_charged < before);
assert_eq!(process.resource_snapshot().framework_charged, baseline);
}
#[tokio::test]
async fn real_transport_completes_reuses_and_closes_one_peer_connection() {
if isolated("real_transport_completes_reuses_and_closes_one_peer_connection") { return; }
real_transport_case().await;
}
#[tokio::test(flavor = "multi_thread", worker_threads = 2)]
async fn real_transport_multithread_reuses_after_physical_drain() {
if isolated("real_transport_multithread_reuses_after_physical_drain") { return; }
real_transport_case().await;
}
async fn real_transport_case() {
use http_body_util::{BodyExt, Full};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let authority = format!("http://{}", listener.local_addr().unwrap());
let accepted = Arc::new(AtomicUsize::new(0));
let connections = accepted.clone();
let stop = Arc::new(tokio::sync::Notify::new());
let stop_peer = stop.clone();
let peer = tokio::spawn(async move {
let mut tasks = tokio::task::JoinSet::new();
loop {
tokio::select! {
() = stop_peer.notified() => break,
accepted = listener.accept() => {
let (stream, _) = accepted.unwrap();
connections.fetch_add(1, Ordering::SeqCst);
tasks.spawn(async move {
let service = hyper::service::service_fn(|request: hyper::Request<hyper::body::Incoming>| async {
assert_eq!(request.uri().path(), "/synthetic.Unit/Echo");
let sequence = request.headers()["x-sequence"].clone();
let bytes = request.into_body().collect().await.unwrap().to_bytes();
assert_eq!(bytes.len(), 6);
assert_eq!(bytes[5].to_string(), sequence.to_str().unwrap());
let body = Full::new(bytes).with_trailers(async {
let mut trailers = hyper::HeaderMap::new();
trailers.insert("grpc-status", "0".parse().unwrap());
Some(Ok::<_, std::convert::Infallible>(trailers))
});
Ok::<_, std::convert::Infallible>(hyper::Response::builder()
.header("content-type", "application/grpc").body(body).unwrap())
});
hyper::server::conn::http2::Builder::new(hyper_util::rt::TokioExecutor::new())
.serve_connection(hyper_util::rt::TokioIo::new(stream), service).await
});
}
}
}
while let Some(join) = tasks.join_next().await { join.unwrap().unwrap(); }
});
let process = process(1);
let root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
let pool = StandardConnectionPool::new(root, process.rpc_stage_snapshot().capacity,
vec![FrozenConnectionTarget::new("alpha", &authority).unwrap()]).unwrap();
let mut service = pool.start().unwrap();
for sequence in 1..=3 {
let execution = execution(&process);
let memory = execution.request_memory();
let stage = execution.try_begin_rpc_stage().unwrap();
let wire = stage.framework_input(|builder| builder.write_bytes(1, |target| { target[0] = sequence; Ok(()) })).unwrap();
let mut lease = pool.try_acquire(0).unwrap();
lease.attach_stage(stage).unwrap();
let mut request = tonic::Request::new(wire);
request.metadata_mut().insert("x-sequence", sequence.to_string().parse().unwrap());
let response = tokio::time::timeout(Duration::from_secs(5), lease.unary_transport(request, memory,
"/synthetic.Unit/Echo", std::time::Instant::now() + Duration::from_secs(4))).await.unwrap()
.unwrap_or_else(|e| panic!("sequence={sequence} accepted={} idle={} failure={e:?}", accepted.load(Ordering::SeqCst), pool.test_idle_count()));
assert_eq!(response.get(), &[sequence]);
drop(response);
lease.drain_response(std::time::Instant::now() + Duration::from_secs(4)).await.unwrap();
lease.finish(true);
assert_eq!(pool.test_idle_count(), 1, "complete healthy response restores an idle connection");
assert_eq!(process.rpc_stage_snapshot().operations, 0, "idle connection retains no request operation");
drop(execution);
}
assert_eq!(accepted.load(Ordering::SeqCst), 1, "peer-observed TCP identity proves reuse");
service.shutdown().await.unwrap();
stop.notify_one();
tokio::time::timeout(Duration::from_secs(5), peer).await.unwrap().unwrap();
}
#[tokio::test]
async fn cancelled_shutdown_keeps_join_owner() {
if isolated("cancelled_shutdown_keeps_join_owner") { return; }
let process = process(1);
let baseline = process.resource_snapshot().framework_charged;
let pool = pool(&process);
let mut service = pool.start().unwrap();
{
let mut shutdown = Box::pin(service.shutdown());
assert!(std::future::Future::poll(shutdown.as_mut(), &mut std::task::Context::from_waker(std::task::Waker::noop())).is_pending());
}
service.shutdown().await.unwrap();
drop((service, pool));
assert_eq!(process.resource_snapshot().framework_charged, baseline);
}
#[tokio::test]
async fn registration_capacity_rejects_and_destroys_sixth_future() {
if isolated("registration_capacity_rejects_and_destroys_sixth_future") { return; }
use hyper::rt::Executor;
let process = Arc::new(process(1));
let pool = pool(&process);
let mut service = pool.start().unwrap();
let execution = execution(&process);
let mut lease = pool.try_acquire(0).unwrap();
lease.attach_stage(execution.try_begin_rpc_stage().unwrap()).unwrap();
let dropped = Arc::new(AtomicUsize::new(0));
for _ in 0..6 {
lease.executor().execute(Box::pin(PhysicalProbe { dropped: dropped.clone(), process: process.clone() }));
}
assert_eq!(dropped.load(Ordering::SeqCst), 1);
assert_eq!(lease.test_driver_count(), 5);
assert!(matches!(lease.failure(), Some(AdmissionError::CapacityRejected)));
drop(lease);
pool.wait_retired().await;
assert_eq!(dropped.load(Ordering::SeqCst), 6);
assert_eq!(process.rpc_stage_snapshot().operations, 0);
service.shutdown().await.unwrap();
}
#[test]
fn late_request_waker_keeps_stage_until_physical_drop() {
if isolated("late_request_waker_keeps_stage_until_physical_drop") { return; }
let process = process(1);
let execution = execution(&process);
let stage = execution.try_begin_rpc_stage().unwrap();
let bridge = StandardRpcRequestWake::prepare(&stage).unwrap();
let late = bridge.waker(std::task::Waker::noop());
bridge.close();
drop((bridge, stage));
assert_eq!(process.rpc_stage_snapshot().operations, 1);
late.wake_by_ref();
assert_eq!(process.rpc_stage_snapshot().operations, 1);
drop(late);
process.reclaim_retired_rpc();
assert_eq!(process.rpc_stage_snapshot().operations, 0);
}
#[test]
fn request_wake_callback_keeps_original_allocation_audit() {
if isolated("request_wake_callback_keeps_original_allocation_audit") { return; }
struct BusinessWake;
impl std::task::Wake for BusinessWake {
fn wake(self: Arc<Self>) { std::hint::black_box(Box::new([7u8; 64])); }
}
let process = process(1);
let execution = execution(&process);
let stage = execution.try_begin_rpc_stage().unwrap();
let bridge = StandardRpcRequestWake::prepare(&stage).unwrap();
let original = std::task::Waker::from(Arc::new(BusinessWake));
bridge.waker(&original).wake();
assert_eq!(process.resource_snapshot().health_failure, HealthFailure::RequestAllocationEscape);
bridge.close();
drop((bridge, stage));
process.reclaim_retired_rpc();
}
#[tokio::test]
async fn root_task_prepayment_shortage_rolls_back_before_spawn() {
if isolated("root_task_prepayment_shortage_rolls_back_before_spawn") { return; }
let process = process(1);
let baseline = process.resource_snapshot().framework_charged;
let root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
let occupation_root = root.try_reserve(StorageDemand::separate(&[]).unwrap()).unwrap();
let pool = StandardConnectionPool::new(root, 1, vec![FrozenConnectionTarget::new("alpha", "http://127.0.0.1:8001").unwrap()]).unwrap();
let snapshot = process.resource_snapshot();
let demand = StorageDemand::separate(&[(std::alloc::Layout::from_size_align(
snapshot.framework_capacity - snapshot.framework_charged - std::mem::size_of::<StoragePermit>(), 1).unwrap(), 1)]).unwrap();
let occupied = occupation_root.try_reserve(demand).unwrap();
assert!(pool.start().is_err());
drop((occupied, occupation_root));
let mut service = pool.start().unwrap();
service.shutdown().await.unwrap();
drop((service, pool));
assert_eq!(process.resource_snapshot().framework_charged, baseline);
}
#[tokio::test]
async fn frozen_alias_and_authority_keys_remain_distinct() {
if isolated("frozen_alias_and_authority_keys_remain_distinct") { return; }
let process = process(2);
let root = process.try_process_storage(StorageDemand::separate(&[]).unwrap()).unwrap();
let pool = StandardConnectionPool::new(root, 2, vec![
FrozenConnectionTarget::new("alpha", "http://{zone}.example:8001").unwrap(),
FrozenConnectionTarget::new("beta", "http://{zone}.example:8001").unwrap(),
]).unwrap();
let mut service = pool.start().unwrap();
assert!(pool.try_acquire_authority(0, "http://outside.example:8002").is_err());
let first = pool.try_acquire_authority(0, "http://east.example:8001").unwrap();
let second = pool.try_acquire_authority(1, "http://east.example:8001").unwrap();
assert!(matches!(pool.try_acquire_authority(0, "http://west.example:8001"), Err(AdmissionError::CapacityRejected)));
drop((first, second));
service.shutdown().await.unwrap();
}
#[tokio::test]
async fn shutdown_with_live_lease_rejects_late_transport_and_drain() {
if isolated("shutdown_with_live_lease_rejects_late_transport_and_drain") { return; }
let process = process(1);
let pool = pool(&process);
let mut service = pool.start().unwrap();
let execution = execution(&process);
let mut lease = pool.try_acquire(0).unwrap();
lease.attach_stage(execution.try_begin_rpc_stage().unwrap()).unwrap();
service.shutdown().await.unwrap();
assert!(matches!(lease.failure(), Some(AdmissionError::AccountClosed)));
use hyper::rt::Executor;
lease.executor().execute(Box::pin(async {}));
assert!(matches!(lease.drain_response(std::time::Instant::now() + Duration::from_secs(1)).await,
Err(TransportFailure::Resource(AdmissionError::AccountClosed))));
assert!(matches!(lease.connect(std::time::Instant::now() + Duration::from_secs(1)).await,
Err(TransportFailure::Resource(AdmissionError::AccountClosed))));
assert_eq!(process.rpc_stage_snapshot().operations, 0);
drop(lease);
}
#[tokio::test]
async fn native_pool_invoice_waits_for_the_last_physical_waker() {
if isolated("native_pool_invoice_waits_for_the_last_physical_waker") { return; }
let process = process(1);
let pool = pool(&process);
let mut service = pool.start().unwrap();
let lease = pool.try_acquire(0).unwrap();
let late = lease.test_driver_waker();
drop(lease);
service.shutdown().await.unwrap();
drop((service, pool));
process.reclaim_retired_rpc();
let held = process.rpc_stage_snapshot();
assert_eq!(held.standard_pool_owners, 1, "join cannot substitute for physical release");
assert!(held.standard_pool_storage_bytes > 0, "the original invoice remains charged");
drop(late);
process.reclaim_retired_rpc();
let closed = process.rpc_stage_snapshot();
assert_eq!(closed.standard_pool_owners, 0);
assert_eq!(closed.standard_pool_storage_bytes, 0);
}