diff --git a/packages/bun-usockets/src/bsd.c b/packages/bun-usockets/src/bsd.c index 4fe48236ca30..6b2456b13f26 100644 --- a/packages/bun-usockets/src/bsd.c +++ b/packages/bun-usockets/src/bsd.c @@ -901,6 +901,29 @@ ssize_t bsd_recv(LIBUS_SOCKET_DESCRIPTOR fd, void *buf, int length, int flags) { } } +int bsd_queued_input(LIBUS_SOCKET_DESCRIPTOR fd) { + /* Windows has no MSG_DONTWAIT. Every socket uSockets owns there is already + * non-blocking, which is what bsd_recv relies on too. */ +#ifdef _WIN32 + const int peek_flags = MSG_PEEK; +#else + const int peek_flags = MSG_PEEK | MSG_DONTWAIT; +#endif + char byte; + ssize_t ret; + do { + ret = recv(fd, &byte, 1, peek_flags); + } while (UNLIKELY(IS_EINTR(ret))); + + if (ret > 0) { + return LIBUS_QUEUED_INPUT_DATA; + } + if (ret == 0) { + return LIBUS_QUEUED_INPUT_EOF; + } + return bsd_would_block() ? LIBUS_QUEUED_INPUT_NONE : LIBUS_QUEUED_INPUT_ERROR; +} + #if !defined(_WIN32) ssize_t bsd_recvmsg(LIBUS_SOCKET_DESCRIPTOR fd, struct msghdr *msg, int flags) { ssize_t injected = 0; int unused = 0; diff --git a/packages/bun-usockets/src/internal/networking/bsd.h b/packages/bun-usockets/src/internal/networking/bsd.h index 4d53a9fa3fa6..09d0236cd1e9 100644 --- a/packages/bun-usockets/src/internal/networking/bsd.h +++ b/packages/bun-usockets/src/internal/networking/bsd.h @@ -214,6 +214,9 @@ int bsd_addr_get_port(struct bsd_addr_t *addr); LIBUS_SOCKET_DESCRIPTOR bsd_accept_socket(LIBUS_SOCKET_DESCRIPTOR fd, struct bsd_addr_t *addr); ssize_t bsd_recv(LIBUS_SOCKET_DESCRIPTOR fd, void *buf, int length, int flags); +/* One of the LIBUS_QUEUED_INPUT_* codes. Peeks, so it consumes nothing and the + * next readable event still reports whatever it found. */ +int bsd_queued_input(LIBUS_SOCKET_DESCRIPTOR fd); #if !defined(_WIN32) ssize_t bsd_recvmsg(LIBUS_SOCKET_DESCRIPTOR fd, struct msghdr *msg, int flags); #endif diff --git a/packages/bun-usockets/src/libusockets.h b/packages/bun-usockets/src/libusockets.h index 0a5d82ae32c3..ad8974160762 100644 --- a/packages/bun-usockets/src/libusockets.h +++ b/packages/bun-usockets/src/libusockets.h @@ -677,6 +677,19 @@ void us_socket_shutdown(us_socket_r s) nonnull_fn_decl; void us_socket_shutdown_read(us_socket_r s) nonnull_fn_decl; int us_socket_is_shut_down(us_socket_r s) nonnull_fn_decl; int us_socket_is_closed(us_socket_r s) nonnull_fn_decl; + +/* Return codes of us_socket_queued_input. */ +#define LIBUS_QUEUED_INPUT_NONE 0 /* a read would block: nothing is queued */ +#define LIBUS_QUEUED_INPUT_DATA 1 /* at least one byte is readable */ +#define LIBUS_QUEUED_INPUT_EOF 2 /* the peer sent a FIN */ +#define LIBUS_QUEUED_INPUT_ERROR 3 /* the read side failed, e.g. a reset */ +/* What the read side of the socket holds right now, without a trip through + * the event loop. The peek consumes nothing, so the poll still reports the + * same input later and the normal read path still handles it. Callers that + * own a socket between loop iterations use this to tell an idle connection + * from one the peer has already written to or closed. */ +int us_socket_queued_input(us_socket_r s) nonnull_fn_decl; + int us_socket_is_ssl_handshake_finished(us_socket_r s) nonnull_fn_decl; int us_socket_ssl_handshake_callback_has_fired(us_socket_r s) nonnull_fn_decl; /* TLS ciphertext bytes already sealed for this socket and reported as diff --git a/packages/bun-usockets/src/socket.c b/packages/bun-usockets/src/socket.c index 20f351f74b6c..9c8f1fa82671 100644 --- a/packages/bun-usockets/src/socket.c +++ b/packages/bun-usockets/src/socket.c @@ -169,6 +169,13 @@ __attribute__((always_inline)) int us_socket_is_established(struct us_socket_t * return us_internal_poll_type((struct us_poll_t *) s) != POLL_TYPE_SEMI_SOCKET; } +int us_socket_queued_input(struct us_socket_t *s) { + if (s->flags.is_closed || !us_socket_is_established(s)) { + return LIBUS_QUEUED_INPUT_NONE; + } + return bsd_queued_input(us_poll_fd(&s->p)); +} + /* Detach c from its group + drop the borrowed SSL_CTX ref, but leave c * allocated. After this, c->group is NULL and the embedding owner may safely * deinit; the only remaining link is into a loop-owned list. */ diff --git a/src/http/HTTPContext.rs b/src/http/HTTPContext.rs index 01d1552f1f7e..b49758c85ff0 100644 --- a/src/http/HTTPContext.rs +++ b/src/http/HTTPContext.rs @@ -876,6 +876,27 @@ impl HTTPContext { continue; } + // `HTTPThread::drain_events` hands a socket out before the loop + // polls it, so input the origin already wrote is unread in the + // kernel and invisible to the checks above; reuse answers this + // request with it. Same verdicts the idle handlers reach after + // a poll: `Handler::on_data` terminates, `Handler::on_end` + // closes. HTTP/2 idle frames are healthy, so `on_idle_data` + // keeps deciding for those. + if socket.h2_session.is_none() { + match http_socket.queued_input() { + uws::QueuedInput::None => {} + uws::QueuedInput::Eof => { + Self::close_socket(http_socket); + continue; + } + uws::QueuedInput::Data | uws::QueuedInput::Error => { + Self::terminate_socket(http_socket); + continue; + } + } + } + // Transfer tunnel ownership (the parked strong ref) to the caller. let tunnel: Option> = socket.proxy_tunnel.take(); socket.target_hostname = Box::default(); diff --git a/src/uws/lib.rs b/src/uws/lib.rs index 627f9ddd52c8..d564a2265830 100644 --- a/src/uws/lib.rs +++ b/src/uws/lib.rs @@ -1367,6 +1367,7 @@ pub use bun_uws_sys::SocketKind; pub type DispatchKind = SocketKind; pub use bun_uws_sys::CloseCode; +pub use bun_uws_sys::QueuedInput; /// Legacy alias — `bun_uws_sys::CloseCode` is the one canonical `#[repr(i32)]` /// enum (`normal`/`failure`/`fast_shutdown`, with `Normal`/`Failure`/ /// `FastShutdown` associated-const aliases). diff --git a/src/uws_sys/lib.rs b/src/uws_sys/lib.rs index e1ac6fdd9e3b..ad9170bca813 100644 --- a/src/uws_sys/lib.rs +++ b/src/uws_sys/lib.rs @@ -492,7 +492,7 @@ pub use response::{AnyResponse, SocketAddress, WebSocketUpgradeContext}; pub use socket_context::BunSocketContextOptions; pub use socket_group::ConnectResult; pub use socket_group::SocketGroup; -pub use us_socket::{CloseCode, UsIoVec, us_socket_stream_buffer_t, us_socket_t}; +pub use us_socket::{CloseCode, QueuedInput, UsIoVec, us_socket_stream_buffer_t, us_socket_t}; pub use web_socket::{AnyWebSocket, RawWebSocket, WebSocketBehavior}; /// Legacy aliases for `App` / `Response`. diff --git a/src/uws_sys/socket.rs b/src/uws_sys/socket.rs index 6f5bb3e0484a..016e77932502 100644 --- a/src/uws_sys/socket.rs +++ b/src/uws_sys/socket.rs @@ -21,7 +21,7 @@ use bun_core::{Fd, ZStr}; use crate::WindowsNamedPipe; use crate::{ CloseCode, ConnectResult, ConnectingSocket, LIBUS_SOCKET_ALLOW_HALF_OPEN, - LIBUS_SOCKET_DESCRIPTOR, SocketGroup, SocketKind, SslCtx, UpgradedDuplex, + LIBUS_SOCKET_DESCRIPTOR, QueuedInput, SocketGroup, SocketKind, SslCtx, UpgradedDuplex, us_bun_verify_error_t, us_socket_t, }; @@ -301,6 +301,17 @@ impl NewSocketHandler { self.is_closed() || self.is_shutdown() || self.get_error() != 0 } + /// What the read side holds right now, without waiting for the loop to + /// poll. Transports the loop does not `recv()` report nothing queued. + pub fn queued_input(&self) -> QueuedInput { + on_socket!(self.socket; + connected s => s.queued_input(), + duplex _d => QueuedInput::None, + pipe _p => QueuedInput::None, + else => QueuedInput::None, + ) + } + pub fn get_verify_error(&self) -> us_bun_verify_error_t { on_socket!(self.socket; connected s => s.get_verify_error(), diff --git a/src/uws_sys/us_socket_t.rs b/src/uws_sys/us_socket_t.rs index 6a26fd207166..e4a1f1a96d6c 100644 --- a/src/uws_sys/us_socket_t.rs +++ b/src/uws_sys/us_socket_t.rs @@ -44,6 +44,27 @@ pub enum CloseCode { fast_shutdown = 2, } +/// `LIBUS_QUEUED_INPUT_*` in libusockets.h, mirrored by name. +pub const LIBUS_QUEUED_INPUT_NONE: c_int = 0; +pub const LIBUS_QUEUED_INPUT_DATA: c_int = 1; +pub const LIBUS_QUEUED_INPUT_EOF: c_int = 2; +pub const LIBUS_QUEUED_INPUT_ERROR: c_int = 3; + +/// What a socket's read side holds right now. The peek behind it consumes +/// nothing, so the normal read path still gets the same bytes. +#[repr(i32)] +#[derive(Copy, Clone, Eq, PartialEq, Debug)] +pub enum QueuedInput { + /// A read would block: the peer has written nothing since the last read. + None = LIBUS_QUEUED_INPUT_NONE, + /// At least one byte is readable. + Data = LIBUS_QUEUED_INPUT_DATA, + /// The peer sent a FIN. + Eof = LIBUS_QUEUED_INPUT_EOF, + /// The read side failed, for example a reset. + Error = LIBUS_QUEUED_INPUT_ERROR, +} + /// Layout-compatible with `struct us_iovec_t` in libusockets.h (== POSIX iovec). #[repr(C)] #[derive(Clone, Copy)] @@ -464,6 +485,15 @@ impl us_socket_t { pub(crate) fn is_established(&self) -> bool { c::us_socket_is_established(self) > 0 } + + pub(crate) fn queued_input(&self) -> QueuedInput { + match c::us_socket_queued_input(self) { + LIBUS_QUEUED_INPUT_DATA => QueuedInput::Data, + LIBUS_QUEUED_INPUT_EOF => QueuedInput::Eof, + LIBUS_QUEUED_INPUT_ERROR => QueuedInput::Error, + _ => QueuedInput::None, + } + } } /// Raw externs. Private — every operation has a typed method on `us_socket_t`. @@ -576,6 +606,7 @@ mod c { pub(super) safe fn us_socket_verify_error(s: &us_socket_t) -> us_bun_verify_error_t; pub(super) safe fn us_socket_get_error(s: &us_socket_t) -> c_int; pub(super) safe fn us_socket_is_established(s: &us_socket_t) -> i32; + pub(super) safe fn us_socket_queued_input(s: &us_socket_t) -> c_int; /// ssl_ctx is required (the whole point); sni may be null. pub(super) fn us_socket_adopt_tls( diff --git a/test/js/web/fetch/fetch-keepalive.test.ts b/test/js/web/fetch/fetch-keepalive.test.ts index efd2ac85a2b5..c08303c08ff6 100644 --- a/test/js/web/fetch/fetch-keepalive.test.ts +++ b/test/js/web/fetch/fetch-keepalive.test.ts @@ -1,7 +1,8 @@ import { expect, test } from "bun:test"; -import { bunEnv, bunExe, isWindows, tls } from "harness"; +import { bunEnv, bunExe, isMacOS, isWindows, tempDir, tls } from "harness"; import { once } from "node:events"; import { createServer } from "node:net"; +import { join } from "node:path"; test("keepalive", async () => { using server = Bun.serve({ @@ -700,6 +701,137 @@ for (const [label, earlyReply, body, first, onWindows] of earlyReplyCases) { }); } +// A connection parked in the keep-alive pool is handed to the next request by +// `HTTPThread::drain_events`, which runs before the event loop polls. So input +// the origin wrote after its last response can still be unread in the kernel at +// that moment, and `is_closed`/`is_shutdown`/`get_error` cannot see it. Writing +// the next request onto that connection makes bun answer the request with bytes +// that were already on the wire before it: an unsolicited response, or the +// `HTTP/1.1 408 Request Timeout` plus `Connection: close` that many servers and +// load balancers use to retire an idle keep-alive connection. The request was +// never processed, yet `fetch()` resolves with 408 and the origin's timeout +// body. +// +// Every round writes the injected event BEFORE it queues request 2, so the +// bytes are in bun's kernel buffer before bun can write that request anywhere. +// Two ballast requests to an origin that never answers are queued first, which +// takes the HTTP thread out of `poll()`: without them the loop reads the +// injected bytes on the idle connection and retires it through +// `Handler::on_data`, which is the behaviour this check extends to the checkout +// window. Whichever of the two gets there first, request 2 must be answered on +// a connection the origin accepted later. +const idleInjections: [label: string, bytes: string][] = [ + [ + "408 Request Timeout and Connection: close", + "HTTP/1.1 408 Request Timeout\r\nConnection: close\r\nContent-Length: 5\r\n\r\nT-OUT", + ], + ["a complete unsolicited response", "HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nPWNED"], + ["one stray CRLF", "\r\n"], + ["64 bytes of garbage", Buffer.alloc(64, "!").toString()], + ["a FIN", ""], +]; +// "Before" only holds when write() returns with the bytes already in the +// peer's receive buffer. An AF_UNIX stream does that everywhere, and TCP +// loopback does it on Linux and Windows. macOS hands a loopback segment to +// the dlil input thread first, so it can land after the checkout, which no +// client can tell from an answer: there only the unix-socket pool is checked. +// fetch() does not pool unix sockets on Windows. Both pools check a socket +// out through the same code. +const idleTransports = [...(isMacOS ? [] : ["tcp"]), ...(isWindows ? [] : ["unix"])]; + +test.concurrent.each( + idleTransports.flatMap(transport => idleInjections.map(([label, bytes]) => [transport, label, bytes])), +)("%s: a pooled connection the origin answered with %s is not reused", async (transport, _label, bytes) => { + const rounds = 12; + using dir = tempDir("fetch-ka-idle", {}); + // Subprocess so the keep-alive pool starts empty and so the origin shares a + // thread with the client: the injected write and the next fetch() are then + // ordered by the JS thread, not by a timer. + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const bytes = ${JSON.stringify(bytes)}; + const unix = ${JSON.stringify(transport === "unix" ? join(String(dir), "o.sock") : null)}; + const misattributed = []; + + // Accepts connections and never answers them. + using sink = Bun.listen({ + hostname: "127.0.0.1", + port: 0, + socket: { open() {}, data() {}, close() {}, error() {}, drain() {} }, + }); + + let accepted = 0; + let idle = null; + using server = Bun.listen({ + ...(unix ? { unix } : { hostname: "127.0.0.1", port: 0 }), + socket: { + open(socket) { + socket.data = { id: ++accepted, buf: "" }; + }, + data(socket, chunk) { + socket.data.buf += chunk.toString("latin1"); + let end; + while ((end = socket.data.buf.indexOf("\\r\\n\\r\\n")) >= 0) { + const path = socket.data.buf.slice(0, end).split(" ")[1]; + socket.data.buf = socket.data.buf.slice(end + 4); + const body = "c" + socket.data.id; + socket.write("HTTP/1.1 200 OK\\r\\nX-Conn: " + body + "\\r\\nContent-Length: " + body.length + "\\r\\n\\r\\n" + body); + if (path === "/warm") idle = socket; + } + }, + close() {}, error() {}, drain() {}, + }, + }); + const origin = unix ? "http://localhost" : "http://127.0.0.1:" + server.port; + const init = unix ? { unix } : {}; + + for (let round = 0; round < ${rounds}; round++) { + const warm = await fetch(origin + "/warm", init); + const warmConn = warm.headers.get("x-conn"); + if ((await warm.text()) !== warmConn) throw new Error("warm body " + warmConn); + + const ac = new AbortController(); + const ballast = [0, 1].map(b => + fetch("http://127.0.0.1:" + sink.port + "/" + round + "/" + b, { signal: ac.signal }).catch(() => {}), + ); + if (bytes.length === 0) idle.end(); + else idle.write(bytes); + const pending = fetch(origin + "/next", init); + + let got; + try { + const next = await pending; + got = next.status + ":" + (await next.text()); + } catch (e) { + got = "ERR:" + (e.code ?? e.name); + } + ac.abort(); + await Promise.all(ballast); + + // The parked connection was written to before this request existed, + // so the answer has to be an honest 200 from a connection the origin + // accepted later. Anything else means bun read the injected bytes as + // the answer. + const honest = /^200:c\\d+$/.test(got) && got !== "200:" + warmConn; + if (!honest) misattributed.push("round " + round + " -> " + got); + } + console.log(JSON.stringify({ misattributed })); + process.exit(0); + `, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + const result = stdout.startsWith("{") ? JSON.parse(stdout.trim()) : { stdout, stderr }; + expect({ result, exitCode }).toEqual({ result: { misattributed: [] }, exitCode: 0 }); +}); + test.skipIf(isWindows)("a full keep-alive pool evicts the longest-idle connection", async () => { function makeServer() { let connections = 0;