use std::io::{Read, Write};
use std::net::TcpListener;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};
use bun_core::MutableString;
use bun_http::{AsyncHTTP, FetchRedirect, HTTPClientResult, HTTPClientResultCallback, Method,
async_http};
#[unsafe(no_mangle)]
extern "Rust" fn __bun_run_file_poll(_poll: *mut bun_io::FilePoll, _size_or_offset: i64) {}
#[unsafe(no_mangle)]
extern "Rust" fn __bun_crash_handler_out_of_memory() -> ! {
eprintln!("bun: out of memory");
std::process::abort()
}
#[derive(Debug, Clone)]
struct Outcome {
status: Option<u32>,
failed: bool,
timed_out: bool,
has_more: bool,
}
struct Recorder {
tx: std::sync::mpsc::Sender<Outcome>,
}
fn recorder_callback(
this: *mut Recorder,
async_http: *mut AsyncHTTP<'static>,
result: HTTPClientResult<'_>,
) {
let rec: &Recorder = unsafe { &*this };
let _ = rec.tx.send(Outcome {
status: result.metadata.as_ref().map(|m| m.response.status_code),
failed: result.fail.is_some(),
timed_out: result.is_timeout(),
has_more: result.has_more,
});
if !result.has_more {
let real = unsafe { (*async_http).real };
if let Some(r) = real {
drop(unsafe { Box::from_raw(r.as_ptr()) });
}
let buf = unsafe { (*async_http).response_buffer };
if !buf.is_null() {
drop(unsafe { Box::from_raw(buf) });
}
}
}
fn spawn_request(url: String) -> std::sync::mpsc::Receiver<Outcome> {
let (tx, rx) = std::sync::mpsc::channel();
let recorder = Box::into_raw(Box::new(Recorder { tx }));
let url_bytes: &'static [u8] = Box::leak(url.into_bytes().into_boxed_slice());
let parsed_url = bun_url::URL::parse(url_bytes);
let response_buffer = Box::into_raw(Box::new(MutableString::default()));
let ah = AsyncHTTP::init(
Method::GET,
parsed_url,
Default::default(),
b"",
response_buffer,
b"",
HTTPClientResultCallback::new(recorder, recorder_callback),
FetchRedirect::Follow,
async_http::Options::default(),
);
let ah_ptr = bun_core::heap::into_raw(Box::new(ah));
let batch = bun_threading::thread_pool::Batch::from(unsafe {
core::ptr::addr_of_mut!((*ah_ptr).task)
});
bun_http::HTTPThread::schedule(batch);
rx
}
fn collect_bounded(rx: &std::sync::mpsc::Receiver<Outcome>, bound: Duration) -> Vec<Outcome> {
let deadline = Instant::now() + bound;
let mut out = Vec::new();
loop {
let Some(remaining) = deadline.checked_duration_since(Instant::now()) else {
break;
};
let Ok(d) = rx.recv_timeout(remaining) else {
break;
};
let terminal = !d.has_more;
out.push(d);
if terminal {
break;
}
}
out
}
struct CpuHogs {
stop: Arc<AtomicBool>,
handles: Vec<std::thread::JoinHandle<()>>,
}
impl CpuHogs {
fn spawn(multiplier: usize) -> Self {
let stop = Arc::new(AtomicBool::new(false));
let cores = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(8);
let handles = (0..cores.saturating_mul(multiplier))
.map(|_| {
let stop = Arc::clone(&stop);
std::thread::spawn(move || {
while !stop.load(Ordering::Relaxed) {
std::hint::spin_loop();
}
})
})
.collect();
Self { stop, handles }
}
}
impl Drop for CpuHogs {
fn drop(&mut self) {
self.stop.store(true, Ordering::Relaxed);
for h in self.handles.drain(..) {
let _ = h.join();
}
}
}
fn spawn_accept_server<F>(deadline_secs: u64, serve: F) -> u16
where
F: Fn(std::net::TcpStream) + Send + Sync + 'static,
{
let listener = TcpListener::bind("127.0.0.1:0").expect("bind");
let port = listener.local_addr().unwrap().port();
let serve = Arc::new(serve);
listener.set_nonblocking(true).ok();
std::thread::spawn(move || {
let deadline = Instant::now() + Duration::from_secs(deadline_secs);
while Instant::now() < deadline {
match listener.accept() {
Ok((stream, _)) => {
let serve = Arc::clone(&serve);
std::thread::spawn(move || serve(stream));
}
Err(_) => std::thread::sleep(Duration::from_millis(2)),
}
}
});
port
}
fn read_request_head(stream: &mut std::net::TcpStream) {
stream.set_read_timeout(Some(Duration::from_secs(5))).ok();
let mut buf = [0u8; 4096];
let mut seen = Vec::new();
loop {
match stream.read(&mut buf) {
Ok(0) => return,
Ok(n) => {
seen.extend_from_slice(&buf[..n]);
if seen.windows(4).any(|w| w == b"\r\n\r\n") {
return;
}
}
Err(_) => return,
}
}
}
#[test]
fn connection_deadline_both_directions_under_oversubscription() {
bao_native_stubs::force_link();
bun_core::Output::init_test();
bun_http::http_thread::init(&Default::default());
let prev_idle = bun_http::IDLE_TIMEOUT_SECONDS.load(Ordering::Relaxed);
bun_http::IDLE_TIMEOUT_SECONDS.store(300, Ordering::Relaxed);
let accepted = Arc::new(AtomicUsize::new(0));
let accepted_sink = Arc::clone(&accepted);
let port = spawn_accept_server(60, move |mut stream| {
accepted_sink.fetch_add(1, Ordering::Relaxed);
read_request_head(&mut stream);
let resp = "HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nok";
let _ = stream.write_all(resp.as_bytes());
let _ = stream.flush();
});
std::thread::sleep(Duration::from_millis(20));
const CONCURRENT: usize = 8;
let rxs: Vec<_> = (0..CONCURRENT)
.map(|i| spawn_request(format!("http://127.0.0.1:{port}/load-{i}")))
.collect();
let hogs = CpuHogs::spawn(2);
let start = Instant::now();
let results: Vec<Vec<Outcome>> = rxs
.iter()
.map(|rx| collect_bounded(rx, Duration::from_secs(30)))
.collect();
let elapsed = start.elapsed();
drop(hogs);
assert_eq!(
accepted.load(Ordering::Relaxed),
CONCURRENT,
"healthy server must have served every request"
);
for (i, deliveries) in results.iter().enumerate() {
let terminal = deliveries
.last()
.unwrap_or_else(|| panic!("request {i}: no delivery at all (hung?)"));
assert!(
!terminal.has_more,
"request {i}: last delivery must be terminal, got {terminal:?}"
);
assert!(
!terminal.failed,
"request {i} mis-killed under load: {terminal:?} (all: {deliveries:?})"
);
assert_eq!(
terminal.status,
Some(200),
"request {i} status (all: {deliveries:?})"
);
assert!(
!terminal.timed_out,
"request {i} was misjudged error.Timeout under load — deadline clock domain broken"
);
}
assert!(
elapsed < Duration::from_secs(25),
"phase 1 took {elapsed:?} — oversubscription degraded into starvation"
);
bun_http::IDLE_TIMEOUT_SECONDS.store(4, Ordering::Relaxed);
let stall_port = spawn_accept_server(60, move |mut stream| {
read_request_head(&mut stream);
let park = Instant::now() + Duration::from_secs(45);
while Instant::now() < park {
std::thread::sleep(Duration::from_millis(100));
let mut probe = [0u8; 16];
match stream.read(&mut probe) {
Ok(0) | Err(_) => return,
Ok(_) => {}
}
}
});
std::thread::sleep(Duration::from_millis(20));
let stall_rx = spawn_request(format!("http://127.0.0.1:{stall_port}/stall"));
let hogs = CpuHogs::spawn(2);
let start = Instant::now();
let deliveries = collect_bounded(&stall_rx, Duration::from_secs(30));
let elapsed = start.elapsed();
drop(hogs);
bun_http::IDLE_TIMEOUT_SECONDS.store(prev_idle, Ordering::Relaxed);
let terminal = deliveries
.last()
.expect("stalled request must terminate (idle deadline), got no delivery");
assert!(
terminal.timed_out,
"stalled server must fail with error.Timeout, got {terminal:?} (all: {deliveries:?})"
);
assert!(
elapsed < Duration::from_secs(30),
"timeout firing delayed to {elapsed:?} — tick lag grew unbounded"
);
bun_http::http_thread::shutdown_for_exit();
std::process::exit(0);
}