Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
33 changes: 21 additions & 12 deletions src/http/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -261,9 +261,10 @@ pub static OVERRIDDEN_DEFAULT_USER_AGENT: std::sync::OnceLock<&'static [u8]> =
std::sync::OnceLock::new();

/// Idle timeout for HTTP client sockets, in seconds. The timer is armed in
/// `on_open` (so it covers the TLS handshake) and re-armed on every read/write;
/// if no bytes move in either direction for this long the request fails with
/// `error.Timeout`. 0 disables the timer (matching `disable_timeout = true`).
/// `on_open` (so it covers the TLS handshake) and re-armed on writes and on
/// body-phase reads; response-header reads do not re-arm it, so it is an
/// absolute deadline for the header block to complete (undici `headersTimeout`
/// semantics). 0 disables the timer (matching `disable_timeout = true`).
Comment thread
robobun marked this conversation as resolved.
/// Overridable via `BUN_CONFIG_HTTP_IDLE_TIMEOUT`. Default is 5 minutes — the
/// previous hard-coded value — so unchanged environments see identical
/// behaviour except that the handshake phase is now also covered. Values
Expand Down Expand Up @@ -3618,11 +3619,12 @@ impl<'a> HTTPClient<'a> {
buffer.list.as_slice()
};

// Persist the unparsed tail for the next `on_data` and re-arm the
// receive timeout. When `needs_move`, `to_read` is a suffix of
// `incoming_data` and is copied into the (currently empty) accumulation
// buffer; otherwise `to_read` is a suffix of `buffer`, so the consumed
// prefix is drained and `buffer` is moved back into state.
// Persist the unparsed tail for the next `on_data`. When `needs_move`,
// `to_read` is a suffix of `incoming_data` and is copied into the
// (currently empty) accumulation buffer; otherwise `to_read` is a suffix
// of `buffer`, so the consumed prefix is drained and `buffer` is moved
// back into state. Does not re-arm the idle timer (header phase is an
// absolute deadline; see [`IDLE_TIMEOUT_SECONDS`]).
Comment thread
robobun marked this conversation as resolved.
macro_rules! short_read {
() => {{
bun_core::scoped_log!(fetch, "handleShortRead");
Expand All @@ -3638,7 +3640,6 @@ impl<'a> HTTPClient<'a> {
.drain_front(buffer.list.len().saturating_sub(keep));
self.state.response_message_buffer = buffer;
}
self.set_timeout(&socket);
return;
}};
}
Expand Down Expand Up @@ -3726,6 +3727,8 @@ impl<'a> HTTPClient<'a> {
return;
}
};
// Headers complete: start the body-idle window fresh (see [`IDLE_TIMEOUT_SECONDS`]).
self.set_timeout(&socket);
Comment thread
robobun marked this conversation as resolved.

if (self.state.content_encoding_i as usize) < response.headers.list.len()
&& !self.state.flags.did_set_content_encoding
Expand Down Expand Up @@ -3782,7 +3785,6 @@ impl<'a> HTTPClient<'a> {
return;
}
} else if self.state.response_stage == ResponseStage::BodyChunk {
self.set_timeout(&socket);
let report_progress = match self.handle_response_body_chunked_encoding(to_read) {
Ok(b) => b,
Err(err) => {
Expand Down Expand Up @@ -3817,8 +3819,15 @@ impl<'a> HTTPClient<'a> {
}

if self.proxy_tunnel.is_some() {
// if we have a tunnel we dont care about the other stages, we will just tunnel the data
self.set_timeout(&socket);
// Body phase only, mirroring the non-proxy dispatch below (header
// phase is an absolute deadline; see [`IDLE_TIMEOUT_SECONDS`]).
Comment thread
robobun marked this conversation as resolved.
debug_assert!(!self.state.flags.receive_paused); // maybe_pause_receive bails on proxy_tunnel
if matches!(
self.state.response_stage,
ResponseStage::Body | ResponseStage::BodyChunk
) {
self.set_timeout(&socket);
}
Comment thread
robobun marked this conversation as resolved.
Comment thread
coderabbitai[bot] marked this conversation as resolved.
self.proxy_tunnel_mut().unwrap().receive(incoming_data);
return;
}
Expand Down
74 changes: 74 additions & 0 deletions test/js/web/fetch/fetch.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3068,3 +3068,77 @@ it("an explicit numeric `timeout` extends the socket idle deadline past the defa
expect(out.withDefault).toStartWith("ERR:");
expect(exitCode).toBe(0);
}, 60_000);

it("the idle timer is an absolute deadline for the response header block (not re-armed by a byte drip)", async () => {
// A server that trickles one response-header byte at a time, each interval
// shorter than the request's idle timeout, must not be able to keep the
// request alive indefinitely. The idle timer is armed when the request is
// written and is not re-armed on partial header reads, so it bounds how long
// the header block may take to arrive in total (undici `headersTimeout`
// semantics). Once the header block completes the body path re-arms per
// chunk, so a slow-but-steady body is still accepted.
const BODY = "abc";
const HEAD = `HTTP/1.1 200 OK\r\nContent-Length: ${BODY.length}\r\n\r\n`;
const DRIP_MS = 2_000;
const DRIP_N = 10; // header drip sends this many single bytes, then the rest at once
const IDLE_MS = 5_000;

const sockets = new Set<net.Socket>();
const intervals = new Set<ReturnType<typeof setInterval>>();
const server = net.createServer(sock => {
sockets.add(sock);
sock.on("close", () => sockets.delete(sock));
sock.on("error", () => {});
sock.once("data", chunk => {
// /h drips DRIP_N header bytes then bursts the rest + body.
// /b bursts the header block then drips the body byte-by-byte.
const headerDrip = chunk.includes("/h ");
if (!headerDrip) sock.write(HEAD);
const dripped = headerDrip ? HEAD.slice(0, DRIP_N) : BODY;
const tail = headerDrip ? HEAD.slice(DRIP_N) + BODY : "";
Comment thread
coderabbitai[bot] marked this conversation as resolved.
let i = 0;
const iv = setInterval(() => {
if (sock.destroyed) {
clearInterval(iv);
intervals.delete(iv);
return;
}
if (i < dripped.length) {
sock.write(dripped[i++]);
} else {
clearInterval(iv);
intervals.delete(iv);
sock.end(tail);
}
}, DRIP_MS);
intervals.add(iv);
});
});
await new Promise<void>(r => server.listen(0, "127.0.0.1", () => r()));
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const port = (server.address() as AddressInfo).port;

try {
const settle = (path: string) =>
fetch(`http://127.0.0.1:${port}${path}`, { timeout: IDLE_MS }).then(
async r => ({ ok: true as const, status: r.status, body: await r.text() }),
e => ({ ok: false as const, name: e?.name as string, message: String(e?.message ?? e) }),
);

// /h: DRIP_N bytes * DRIP_MS = ~20s of drip before the response would
// complete; the 5s idle deadline (uSockets 4s-tick sweep, so ~5-9s) must
// fire first. A build that re-arms on every partial header read resolves
// 200 after the full drip instead.
// /b: headers arrive in one write, then the 3-byte body trickles at
// DRIP_MS/byte (~8s). Each body chunk re-arms the idle timer, so this
// resolves despite taking longer than IDLE_MS overall.
const [hdr, bod] = await Promise.all([settle("/h"), settle("/b")]);
expect({ hdr, bod }).toEqual({
hdr: { ok: false, name: "TimeoutError", message: "The operation timed out." },
bod: { ok: true, status: 200, body: BODY },
});
} finally {
for (const iv of intervals) clearInterval(iv);
for (const s of sockets) s.destroy();
await new Promise<void>(r => server.close(() => r()));
}
}, 60_000);
Loading