Skip to content
Closed
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
57 changes: 45 additions & 12 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -340,12 +340,37 @@ struct HttpContext {
auto result = httpResponseData->template consumePostPadded<IsNodeHttp>(httpContextData->maxHeaderSize, httpResponseData->isConnectRequest, httpContextData->flags.requireHostHeader,httpContextData->flags.useStrictMethodValidation, httpContextData->flags.useInsecureHTTPParser, nodeHttpRequestTrailers, &httpResponseData->chunkedExtensionsByteCount, data, (unsigned int) length, s, proxyParser, [httpContextData](void *s, HttpRequest *httpRequest) -> void * {


HttpResponseData<SSL> *httpResponseData = (HttpResponseData<SSL> *) us_socket_ext((us_socket_t *) s);

/* Are we not ready for another request yet? Bun.serve has one
* HttpResponseData per socket (no queue): drop this request and
* mark the in-flight response for connection-close so the socket
* closes once it drains. Closing here would lose the in-flight
* response before it reaches the wire. RFC 9112 9.3.2: a
* pipelining client must be prepared to retry unanswered requests
* on a new connection.
* Runs before the timeout clear below so a dropped request cannot
* disarm the in-flight request's idleTimeout; pause() stops
* further segments so the connection cannot be held open by a
* pipelined flood or a dropped request's body bytes that extend
* past this segment. getHeaders() writes isConnectRequest via
* bool& for every parsed request; reset it so a dropped CONNECT
* cannot switch the in-flight response into tunnel handling. */
if constexpr (!IsNodeHttp) {
if (httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) {
httpResponseData->state |= HttpResponseData<SSL>::HTTP_PIPELINED_DROP;
httpResponseData->isConnectRequest = false;
/* HttpResponse::pause() also does timeout(0); re-arm after. */
((HttpResponse<SSL> *) s)->pause();
((HttpResponse<SSL> *) s)->resetTimeout();
return s;
}
Comment thread
robobun marked this conversation as resolved.
}

/* For every request we reset the timeout and hang until user makes action */
/* Warning: if we are in shutdown state, resetting the timer is a security issue! */
us_socket_timeout((us_socket_t *) s, 0);

HttpResponseData<SSL> *httpResponseData = (HttpResponseData<SSL> *) us_socket_ext((us_socket_t *) s);

/* node:http compat: the JS layer stopped HTTP processing on this
* connection (the user emitted 'close' on the socket - Node frees
* the parser there); abandon the rest of the buffer. */
Expand All @@ -367,17 +392,13 @@ struct HttpContext {
nodeHttpResponseData->headersCompleted = true;
}

/* Are we not ready for another request yet? Terminate the connection.
* Important for denying async pipelining until, if ever, we want to support it.
* Otherwise requests can get mixed up on the same connection. We still support sync pipelining. */
/* node:http compat: the request arrived while an earlier response on
* this connection is still in flight.
* (For !IsNodeHttp, HTTP_RESPONSE_PENDING was handled above and
* hasQueuedPipelinedResponses stays false, so the else runs.) */
bool hasQueuedPipelinedResponses = false;
if constexpr (IsNodeHttp) hasQueuedPipelinedResponses = httpResponseData->nodeHttpQueuedPipelinedCount > 0;
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) || hasQueuedPipelinedResponses) {
if constexpr (!IsNodeHttp) {
us_socket_close((us_socket_t *) s, 0, nullptr);
return nullptr;
} else {

if (IsNodeHttp && ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) || hasQueuedPipelinedResponses)) {
/* node:http supports async pipelining: the request is dispatched
* while the previous response is still in flight and the JS layer
* queues its response (res.socket === null until it becomes the
Expand All @@ -397,7 +418,6 @@ struct HttpContext {
httpResponseData->state |= HttpResponseData<SSL>::HTTP_NODE_READS_PAUSED;
((HttpResponse<SSL> *) s)->pause();
}
}
} else {
/* Reset httpResponse */
httpResponseData->offset = 0;
Expand Down Expand Up @@ -576,6 +596,19 @@ struct HttpContext {
if(httpContextData->onClientError) {
httpContextData->onClientError(SSL, s, result.parserError, data, length);
}
/* A parse error in bytes trailing a dropped pipelined request must
* not overwrite the in-flight response: its bytes have not reached
* the wire yet, so a 4xx here would be read as the answer to the
* first (valid) request. The drop path already marked the
* connection for close-after-drain and paused reads; leave the
* in-flight response to finish and close. */
if constexpr (!IsNodeHttp) {
if (httpResponseData->state & HttpResponseData<SSL>::HTTP_PIPELINED_DROP) {
us_socket_unref(s);
((AsyncSocket<SSL> *) s)->uncork();
return s;
}
}
/* For errors, we only deliver them "at most once". We don't care if they get halfways delivered or not. */
us_socket_write(s, httpErrorResponses[httpErrorStatusCode].data(), (int) httpErrorResponses[httpErrorStatusCode].length());
us_socket_shutdown(s);
Expand Down
37 changes: 33 additions & 4 deletions packages/bun-uws/src/HttpResponse.h
Original file line number Diff line number Diff line change
Expand Up @@ -187,8 +187,10 @@ struct HttpResponse : public AsyncSocket<SSL> {
}
httpResponseData->markDone(this);

/* We need to check if we should close this socket here now */
if (!Super::isCorked()) {
/* keepCorked=true (only upgrade()) means the caller is adopting the
* socket right after this returns; closing here would destruct the
* ext block out from under it. */
if (!keepCorked && !Super::isCorked()) {
if (httpResponseData->shouldCloseConnection()) {
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) == 0) {
if (((AsyncSocket<SSL> *) this)->getBufferedAmount() == 0) {
Expand Down Expand Up @@ -253,8 +255,10 @@ struct HttpResponse : public AsyncSocket<SSL> {
if (httpResponseData->offset == totalSize) {
httpResponseData->markDone(this);

/* We need to check if we should close this socket here now */
if (!Super::isCorked()) {
/* keepCorked=true (only upgrade()) means the caller is adopting
* the socket right after this returns; closing here would
* destruct the ext block out from under it. */
if (!keepCorked && !Super::isCorked()) {
if (httpResponseData->shouldCloseConnection()) {
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) == 0) {
if (((AsyncSocket<SSL> *) this)->getBufferedAmount() == 0) {
Expand Down Expand Up @@ -367,6 +371,10 @@ struct HttpResponse : public AsyncSocket<SSL> {
/* Grab the httpContext from res */
HttpContext<SSL> *httpContext = HttpContext<SSL>::fromSocket((struct us_socket_t *) this);

/* A pipelined request dropped behind this handshake may have paused
* reads; the adopted WebSocket needs them. No-op if not paused. */
Super::resume();

/* Move any backpressure out of HttpResponse */
auto* responseData = getHttpResponseData();
BackPressure backpressure(std::move(((AsyncSocketData<SSL> *) responseData)->buffer));
Expand Down Expand Up @@ -834,6 +842,10 @@ struct HttpResponse : public AsyncSocket<SSL> {
HttpResponse *cork(MoveOnlyFunction<void()> &&handler) {
if (!Super::isCorked()) {
LoopData *loopData = Super::getLoopData();
/* Captured before handler(): an upgrade inside the handler moves
* the socket to the WebSocket group (in-place) or marks the old
* allocation closed (reallocated). */
us_socket_group_t *httpGroup = us_socket_group((us_socket_t *) this);
Super::cork();
handler();

Expand All @@ -845,6 +857,23 @@ struct HttpResponse : public AsyncSocket<SSL> {
* The upgrade case is handled by HttpContext's uncork or the drain
* loop. */
if (loopData->findCorkSlot(this) == LoopData::INVALID_CORK_SLOT) {
/* (a) from an async handler completing outside this socket's
* onData (which corks it, so that socket takes the else branch
* below) lands here via internalEnd()'s uncork with a
* connection-close mark nothing else acts on. For (b) the ext
* is still HttpResponseData and PENDING is still set; for (c)
* the group changed (in-place adopt) or this allocation was
* marked closed (reallocated adopt). */
if (!us_socket_is_closed((us_socket_t *) this)
&& us_socket_group((us_socket_t *) this) == httpGroup) {
HttpResponseData<SSL> *httpResponseData = getHttpResponseData();
if (httpResponseData->shouldCloseConnection()
&& (httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) == 0
&& ((AsyncSocket<SSL> *) this)->getBufferedAmount() == 0) {
((AsyncSocket<SSL> *) this)->shutdown();
((AsyncSocket<SSL> *) this)->close();
}
}
return this;
}

Expand Down
10 changes: 9 additions & 1 deletion packages/bun-uws/src/HttpResponseData.h
Original file line number Diff line number Diff line change
Expand Up @@ -140,6 +140,14 @@ struct HttpResponseData : AsyncSocketData<SSL>, HttpParser {
* into the shared word so the shared response-end path (internalEnd) never
* has to touch the node-only field. */
HTTP_NODE_HAS_RESPONSE_TRAILERS = 1 << 16,
/* A pipelined request arrived while the previous Bun.serve response was
* still pending and was dropped: the connection closes after the
* in-flight response drains. Distinct from HTTP_CONNECTION_CLOSE so the
* Connection: close header-write guards (internalEnd,
* uws_res_end_without_body) still fire; those guards key off
* HTTP_CONNECTION_CLOSE meaning "the client already knows / a caller
* already wrote it". */
HTTP_PIPELINED_DROP = 1 << 17,

/* Bits that describe the connection rather than the response in flight.
* There is one HttpResponseData per socket, reused by every request on a
Expand Down Expand Up @@ -210,7 +218,7 @@ struct HttpResponseData : AsyncSocketData<SSL>, HttpParser {
/* Whether the connection should be torn down once the in-flight response (if
* any) has completed and all buffered outgoing data has been flushed. */
bool shouldCloseConnection() const {
return (state & HTTP_CONNECTION_CLOSE)
return (state & (HTTP_CONNECTION_CLOSE | HTTP_PIPELINED_DROP))
|| ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0);
}

Expand Down
3 changes: 2 additions & 1 deletion src/uws_sys/Response.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1005,6 +1005,7 @@ bitflags::bitflags! {
const HTTP_RESPONSE_PENDING = 8;
const HTTP_CONNECTION_CLOSE = 16;
const HTTP_WROTE_CONTENT_LENGTH_HEADER = 32;
const HTTP_PIPELINED_DROP = 1 << 17;
}
}

Expand Down Expand Up @@ -1036,7 +1037,7 @@ impl State {

#[inline]
pub fn is_http_connection_close(self) -> bool {
self.bits() & State::HTTP_CONNECTION_CLOSE.bits() != 0
self.bits() & (State::HTTP_CONNECTION_CLOSE.bits() | State::HTTP_PIPELINED_DROP.bits()) != 0
}
}

Expand Down
64 changes: 64 additions & 0 deletions test/js/bun/http/bun-serve-file.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1187,3 +1187,67 @@ test("file route serves a burst of concurrent requests after reloads", async ()
const a = await fetch(`${server.url}a`).then(r => r.text());
expect(a).toBe("a-new");
});

// On POSIX, FileResponseStream reads a regular file with a blocking read()
// inside the request handler, so the response is complete before the parser
// moves on to the second pipelined request. On Windows the read goes to the
// libuv threadpool and completes on a later loop tick, so the second request
// is parsed while HTTP_RESPONSE_PENDING is still set. The in-flight response
// must still reach the wire; the pipelined request may be dropped with the
// connection closing afterwards (RFC 9112 9.3.2) or, like on POSIX, served.
test("Bun.file route answers pipelined HTTP/1.1 requests without dropping the in-flight response", async () => {
using dir = tempDir("serve-file-pipelined", {
"hello.txt": "hello",
"fixture.ts": /* ts */ `
import { connect } from "node:net";
import { join } from "node:path";

const server = Bun.serve({
port: 0,
routes: { "/": new Response(Bun.file(join(import.meta.dir, "hello.txt"))) },
fetch: () => new Response("fallback"),
});

const socket = connect({ port: server.port, host: "127.0.0.1" });
socket.on("error", () => {});
await new Promise((resolve, reject) => { socket.once("connect", resolve); socket.once("error", reject); });
// Two pipelined requests in one TCP segment.
socket.write("GET / HTTP/1.1\\r\\nHost: x\\r\\n\\r\\nGET / HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n");
let data = "";
socket.on("data", chunk => {
data += chunk;
// POSIX keeps the connection open (both served synchronously); end from the
// client once the first response body is on the wire so both platforms
// converge.
if (data.includes("hello")) socket.end();
});
await new Promise(r => socket.once("close", r));
Comment thread
robobun marked this conversation as resolved.

const firstLine = data.slice(0, data.indexOf("\\r\\n"));
const responses = data.split("HTTP/1.1 ").length - 1;
console.log(JSON.stringify({ firstLine, body: data.includes("hello"), responses }));
server.stop(true);
`,
});

await using proc = Bun.spawn({
cmd: [bunExe(), "fixture.ts"],
env: bunEnv,
cwd: String(dir),
stdout: "pipe",
stderr: "pipe",
});

const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);

const result = JSON.parse(stdout.trim() || "{}");
expect({ firstLine: result.firstLine, body: result.body, stderr, exitCode }).toEqual({
firstLine: "HTTP/1.1 200 OK",
body: true,
stderr: "",
exitCode: 0,
});
// At least one full response. POSIX serves both (sync file read); Windows
// serves the first and closes.
expect(result.responses).toBeGreaterThanOrEqual(1);
});
80 changes: 80 additions & 0 deletions test/js/bun/http/serve.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3053,6 +3053,86 @@ server.listen(0, "127.0.0.1", () => {
expect(exitCode).toBe(0);
});

// uWS has one HttpResponseData per socket, so a request that arrives while the
// previous fetch handler is still producing its response cannot be dispatched
// (async pipelining). The in-flight response must still reach the wire; the
// socket closes once it drains and the client retries the dropped request on a
// new connection (RFC 9112 9.3.2). Previously the socket was hard-closed with
// the cork buffer discarded, so zero bytes were written.
describe("pipelined request behind an async fetch handler", () => {
// idleTimeout: 0 so a regression in the close-after-drain path cannot be
// masked by the idle timer closing the socket.
async function pipelined(segment: string): Promise<{ wire: string; calls: number }> {
let calls = 0;
await using server = Bun.serve({
port: 0,
idleTimeout: 0,
async fetch(req) {
calls++;
await new Promise(r => setImmediate(r));
return new Response(new URL(req.url).pathname === "/a" ? "FIRST" : "SECOND");
},
});
const wire = await new Promise<string>((resolve, reject) => {
const socket = net.connect(server.port!, "127.0.0.1");
let data = "";
socket.on("connect", () => socket.write(segment));
socket.on("data", chunk => (data += chunk));
socket.on("close", () => resolve(data));
socket.on("error", reject);
socket.setTimeout(3000, () => {
socket.destroy();
reject(new Error("server did not close the connection after draining"));
});
});
return { wire, calls };
}

it("delivers the in-flight response and closes for pipelined GETs", async () => {
const { wire, calls } = await pipelined(
"GET /a HTTP/1.1\r\nHost: x\r\nConnection: keep-alive\r\n\r\n" +
"GET /b HTTP/1.1\r\nHost: x\r\nConnection: keep-alive\r\n\r\n",
);
expect(wire).toStartWith("HTTP/1.1 200 OK\r\n");
// RFC 9112 9.6: the final response on a connection the server is closing
// SHOULD carry Connection: close so the client knows not to wait for more.
expect(wire).toMatch(/^connection:\s*close\r$/im);
expect(wire).toContain("FIRST");
// The pipelined request is dropped without dispatch; only the first handler runs.
expect(wire).not.toContain("SECOND");
expect(calls).toBe(1);
});

// The dropped request may carry a body; its Content-Length must not let a
// third pipelined request reach the handler once the connection is marked
// for close.
it("drops a pipelined POST body and a third pipelined request", async () => {
const { wire, calls } = await pipelined(
"GET /a HTTP/1.1\r\nHost: x\r\n\r\n" +
"POST /b HTTP/1.1\r\nHost: x\r\nContent-Length: 5\r\n\r\nhello" +
"GET /c HTTP/1.1\r\nHost: x\r\n\r\n",
);
expect(wire).toStartWith("HTTP/1.1 200 OK\r\n");
expect(wire).toContain("FIRST");
expect(wire).not.toContain("SECOND");
expect(calls).toBe(1);
});

// A parse error in the dropped request's bytes (here, an invalid chunked
// size) must not write a 4xx ahead of the in-flight response: its bytes have
// not reached the wire yet, so a 4xx would be read as the first request's
// answer and onClose would abort it.
it("still delivers the in-flight response when the dropped request's bytes are malformed", async () => {
const { wire, calls } = await pipelined(
"GET /a HTTP/1.1\r\nHost: x\r\n\r\n" + "POST /b HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked\r\n\r\nZZ\r\n",
);
expect(wire).toStartWith("HTTP/1.1 200 OK\r\n");
expect(wire).toContain("FIRST");
expect(wire).not.toMatch(/HTTP\/1\.1 4\d\d/);
expect(calls).toBe(1);
});
});

it("only serves /bun:info to loopback clients in development mode", async () => {
using server = Bun.serve({
port: 0,
Expand Down
Loading