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
23 changes: 23 additions & 0 deletions packages/bun-usockets/src/bsd.c
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
3 changes: 3 additions & 0 deletions packages/bun-usockets/src/internal/networking/bsd.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
13 changes: 13 additions & 0 deletions packages/bun-usockets/src/libusockets.h
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
7 changes: 7 additions & 0 deletions packages/bun-usockets/src/socket.c
Original file line number Diff line number Diff line change
Expand Up @@ -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. */
Expand Down
21 changes: 21 additions & 0 deletions src/http/HTTPContext.rs
Original file line number Diff line number Diff line change
Expand Up @@ -876,6 +876,27 @@ impl<const SSL: bool> HTTPContext<SSL> {
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.
Comment thread
robobun marked this conversation as resolved.
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<RefPtr<ProxyTunnel>> = socket.proxy_tunnel.take();
socket.target_hostname = Box::default();
Expand Down
1 change: 1 addition & 0 deletions src/uws/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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).
Expand Down
2 changes: 1 addition & 1 deletion src/uws_sys/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<SSL>` / `Response<SSL>`.
Expand Down
13 changes: 12 additions & 1 deletion src/uws_sys/socket.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
};

Expand Down Expand Up @@ -301,6 +301,17 @@ impl<const IS_SSL: bool> NewSocketHandler<IS_SSL> {
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.
Comment thread
robobun marked this conversation as resolved.
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(),
Expand Down
31 changes: 31 additions & 0 deletions src/uws_sys/us_socket_t.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Comment thread
robobun marked this conversation as resolved.
#[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)]
Expand Down Expand Up @@ -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,
}
Comment thread
robobun marked this conversation as resolved.
}
}

/// Raw externs. Private — every operation has a typed method on `us_socket_t`.
Expand Down Expand Up @@ -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(
Expand Down
134 changes: 133 additions & 1 deletion test/js/web/fetch/fetch-keepalive.test.ts
Original file line number Diff line number Diff line change
@@ -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({
Expand Down Expand Up @@ -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;
Expand Down
Loading