Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
fe5c433
Bun.serve: hold a pipelined request until the response ahead of it co…
robobun Aug 13, 2026
5575a57
Bun.serve: resume reads when an upgrade drops the requests held behin…
robobun Aug 13, 2026
1958e9a
uws: say what the resume in upgrade() costs the adopted WebSocket
robobun Aug 13, 2026
e809be2
test: run the upgrade pipelining cases over tls as well
robobun Aug 13, 2026
ac9a470
uws: keep a request-body resume from reopening reads over parked requ…
robobun Aug 14, 2026
5d2aab1
test: a request pipelined behind a streaming response
robobun Aug 14, 2026
b779b5f
uws: a connection with parked requests is not idle
robobun Aug 14, 2026
0369ed6
verify-baseline-static: allowlist llint_op_jmp_wide32 decode false po…
robobun Aug 27, 2026
cb7c07d
Merge branch 'main' into farm/3aa1ef0f/serve-pipelined-behind-pending…
robobun Sep 21, 2026
4fabf8d
Bun.serve: decide at dispatch, too, whether to hold the next request
robobun Sep 21, 2026
ca7f810
uws: answer the held requests of a peer that has already sent its FIN
robobun Sep 21, 2026
a8613df
test: a throw after an await, a HEAD and Expect: 100-continue ahead o…
robobun Sep 21, 2026
ce32ef9
uws: replay parked requests from a cursor instead of copying the rest…
robobun Sep 21, 2026
c5a70a7
Bun.serve: give the frames held behind a later server.upgrade() to th…
robobun Sep 21, 2026
348b1c9
uws: survive a us_socket_resume() that closes the socket
robobun Sep 21, 2026
b9acb38
uws: cut the pipelining comments down to what the code cannot say
robobun Sep 21, 2026
4a66fc4
Bun.serve: drop what arrives behind a closing request before the hold…
robobun Sep 21, 2026
55af365
Bun.serve: read what is unread behind parked requests before a close
robobun Sep 21, 2026
480dfe8
uws: the close in cork() reads what is unread behind parked requests,…
robobun Sep 21, 2026
8e68a68
Bun.serve: linger on a close that the peer is still writing behind pa…
robobun Sep 22, 2026
7bfa1ea
Merge branch 'main' into farm/3aa1ef0f/serve-pipelined-behind-pending…
robobun Sep 22, 2026
29164bb
uws: take the lingering close off the request path
robobun Sep 22, 2026
e589bad
Bun.serve: a lingering close needs a complete response and is never idle
robobun Sep 22, 2026
f987e1b
uws: keep the lingering-close bit across a dispatch, and fail a reset…
robobun Sep 22, 2026
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
4 changes: 3 additions & 1 deletion packages/bun-usockets/src/libusockets.h
Original file line number Diff line number Diff line change
Expand Up @@ -709,7 +709,9 @@ void us_socket_local_address(us_socket_r s, char *nonnull_arg buf, int *nonnull_

struct us_socket_t *us_socket_detach(us_socket_r s) nonnull_fn_decl;
int us_socket_ipc_write_fd(us_socket_r s, const char *data, int length, int fd) nonnull_fn_decl;
void us_socket_sendfile_needs_more(us_socket_r s) nonnull_fn_decl;
/* Dispatches on_writable once the socket is writable, as if a us_socket_write()
* had come up short. Safe inside on_writable. Leaves a paused read side alone. */
void us_socket_request_writable(us_socket_r s) nonnull_fn_decl;
void *us_listen_socket_ext(struct us_listen_socket_t *ls) nonnull_fn_decl;
LIBUS_SOCKET_DESCRIPTOR us_listen_socket_get_fd(struct us_listen_socket_t *ls) nonnull_fn_decl;
int us_listen_socket_port(struct us_listen_socket_t *ls) nonnull_fn_decl;
Expand Down
8 changes: 8 additions & 0 deletions packages/bun-usockets/src/socket.c
Original file line number Diff line number Diff line change
Expand Up @@ -415,6 +415,14 @@ static void us_internal_rearm_writable(struct us_socket_t *s) {
LIBUS_SOCKET_WRITABLE | ((s->flags.is_paused || s->read_eof) ? 0 : LIBUS_SOCKET_READABLE));
}

/* See libusockets.h. loop.c drops writable interest after an on_writable that
* left last_write_failed clear, so this sets it. */
void us_socket_request_writable(struct us_socket_t *s) {
if (us_socket_is_closed(s)) return;
s->flags.last_write_failed = 1;
us_internal_rearm_writable(s);
}

/* See libusockets.h: whether a zero-progress write on a writable event proves
* the peer is gone. Only the libuv backend has to ask the kernel. */
int us_socket_stalled_write_means_peer_gone(struct us_socket_t *s) {
Expand Down
130 changes: 108 additions & 22 deletions packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
#include <span>
#include <array>
#include <mutex>
#include <utility>


extern "C" void Bun__NodeHTTP__onReadsResumable(int ssl, struct us_socket_t *s);
Expand Down Expand Up @@ -117,11 +118,23 @@ struct HttpContext {
static unsigned char socketKind() { return SSL ? US_SOCKET_KIND_UWS_HTTP_TLS : US_SOCKET_KIND_UWS_HTTP; }

public:
/* node:http flood prevention: re-feed parked request bytes through the same
* parse path fresh socket data takes. The caller guarantees the buffer has
* LIBUS_RECV_BUFFER_PADDING of writable slack past `length`. */
static us_socket_t *feedNodeHttpData(us_socket_t *s, char *data, int length) {
return onData<true>(s, data, length);
/* Re-feeds parkedRequestBytes through onData and returns what it returns. The
* buffer belongs to this call while it is parsed: a dispatch can close or
* upgrade the socket, which destructs the HTTP state. */
template <bool IsNodeHttp>
static us_socket_t *replayParkedRequestBytes(us_socket_t *s) {
auto *httpResponseData = reinterpret_cast<HttpResponseData<SSL> *>(us_socket_ext(s));
WTF::Vector<char> replayed = std::exchange(httpResponseData->parkedRequestBytes, {});
size_t start = std::exchange(httpResponseData->parkedRequestBytesStart, 0);
size_t length = replayed.size() - start;
/* The parser fences the buffer by writing past its logical end. */
replayed.grow(replayed.size() + LIBUS_RECV_BUFFER_PADDING);
httpResponseData->replayedRequestBytes = &replayed;
us_socket_t *returned = onData<IsNodeHttp>(s, replayed.mutableSpan().data() + start, static_cast<int>(length));
if (!us_socket_is_closed(s) && us_socket_kind(s) == socketKind()) {
httpResponseData->replayedRequestBytes = nullptr;
}
return returned;
}

us_socket_group_t *getSocketGroup() {
Expand Down Expand Up @@ -307,6 +320,12 @@ struct HttpContext {
return us_socket_close(s, 0, nullptr);
}

/* Bun.serve: a connection has one response slot. While it is taken, the next
* request head is parked instead of parsed (HttpParser::parkAtNextBoundary). */
static bool cannotDispatchAnotherRequest(HttpResponseData<SSL> *httpResponseData) {
return (httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) != 0;
}

template <bool IsNodeHttp>
static us_socket_t *onData(us_socket_t *s, char *data, int length) {
// ref the socket to make sure we process it entirely before it is closed
Expand All @@ -325,6 +344,17 @@ struct HttpContext {
/* Balance the us_socket_ref above — every other return path
* reaches the unref via returnedData. */
us_socket_unref(s);
/* Bun.serve: a lingering close drops what the peer still sends, up to
* a limit (HttpResponse::shutdownAndClose). */
if constexpr (!IsNodeHttp) {
auto *lingering = (HttpResponseData<SSL> *) us_socket_ext(s);
if (lingering->state & HttpResponseData<SSL>::HTTP_LINGERING_CLOSE) {
lingering->received_bytes_per_timeout += (unsigned int) length;
if (lingering->received_bytes_per_timeout > HttpResponse<SSL>::LINGERING_CLOSE_MAX_BYTES) {
return ((AsyncSocket<SSL> *) s)->close();
}
}
}
return s;
}

Expand Down Expand Up @@ -400,6 +430,12 @@ struct HttpContext {
httpContextData->parsingSocket = s;
httpResponseData->isIdle = false;

/* Bun.serve: derived on every read, so a request that arrives after the
* response completed takes the ordinary path. */
if constexpr (!IsNodeHttp) {
httpResponseData->parkAtNextBoundary = cannotDispatchAnotherRequest(httpResponseData);
}

/* node:http compat: maintain the headers/request timeout window (see
* the requestHandler/dataHandler hooks and the post-parse check). */
const bool trackNodeHttpTimings = IsNodeHttp && !httpResponseData->isConnectRequest;
Expand Down Expand Up @@ -461,15 +497,16 @@ 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. */
/* Is the previous response on this connection still in flight? */
bool hasQueuedPipelinedResponses = false;
if constexpr (IsNodeHttp) hasQueuedPipelinedResponses = httpResponseData->nodeHttpQueuedPipelinedCount > 0;
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) || hasQueuedPipelinedResponses) {
if constexpr (!IsNodeHttp) {
/* Responses that completed earlier in this read can still
* sit in the cork buffer. close() sends them first. */
/* The parser parks a head that arrives while the response slot is
* taken, so this is a backstop against interleaved responses.
* Responses that completed earlier in this read can still sit in
* the cork buffer. close() sends them first. */
ASSERT_NOT_REACHED();
((AsyncSocket<SSL> *) s)->close();
return nullptr;
} else {
Expand All @@ -493,7 +530,7 @@ struct HttpContext {
httpResponseData->state |= HttpResponseData<SSL>::HTTP_NODE_READS_PAUSED;
/* Also stop the request loop over the buffer being parsed
* right now — pausing the socket alone cannot bound it. */
httpResponseData->nodeHttpParkAtNextBoundary = true;
httpResponseData->parkAtNextBoundary = true;
((HttpResponse<SSL> *) s)->pause();
}
}
Expand Down Expand Up @@ -526,7 +563,7 @@ struct HttpContext {
* and park already-received requests. No already-paused guard (replay clears the park flag only). */
if (((AsyncSocket<SSL> *) s)->getBufferedAmount() > 0) {
httpResponseData->state |= HttpResponseData<SSL>::HTTP_NODE_READS_PAUSED;
httpResponseData->nodeHttpParkAtNextBoundary = true;
httpResponseData->parkAtNextBoundary = true;
((HttpResponse<SSL> *) s)->pause();
}
}
Expand Down Expand Up @@ -593,6 +630,12 @@ struct HttpContext {
((HttpResponse<SSL> *) s)->resetTimeout();
}

/* Bun.serve: derived here too, because a Content-Length: 0 head completed
* from the fallback buffer gets no end-of-message callback. */
if constexpr (!IsNodeHttp) {
httpResponseData->parkAtNextBoundary = cannotDispatchAnotherRequest(httpResponseData);
}

/* Continue parsing */
return s;

Expand Down Expand Up @@ -665,6 +708,14 @@ struct HttpContext {
httpResponseData->inStream = nullptr;
}
}

/* Bun.serve: the handler may have completed the response since the
* dispatch, also from inside this body callback. */
if constexpr (!IsNodeHttp) {
if (fin) {
httpResponseData->parkAtNextBoundary = cannotDispatchAnotherRequest(httpResponseData);
}
}
return user;
});

Expand Down Expand Up @@ -702,9 +753,13 @@ struct HttpContext {
((AsyncSocket<SSL> *) s)->uncork();
/* 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);
/* Close any socket on HTTP errors */
us_socket_close(s, 0, nullptr);
if constexpr (!IsNodeHttp) {
((HttpResponse<SSL> *) s)->shutdownAndClose(httpResponseData);
Comment thread
robobun marked this conversation as resolved.
} else {
us_socket_shutdown(s);
us_socket_close(s, 0, nullptr);
}
}

auto returnedData = result.returnedData;
Expand Down Expand Up @@ -733,14 +788,24 @@ struct HttpContext {
((HttpResponse<SSL> *) s)->resetTimeout();
}

/* Bun.serve: reads stay paused while requests are parked, which bounds
* them to one recv. markDone() arms the replay, unless the response was
* complete before anything was parked. AsyncSocket::pause and not
* HttpResponse::pause: the pending response's timeout must stay armed. */
if constexpr (!IsNodeHttp) {
if (!httpResponseData->parkedRequestBytes.isEmpty()) [[unlikely]] {
((AsyncSocket<SSL> *) s)->pause();
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) == 0) {
us_socket_request_writable(s);
}
}
}
Comment thread
robobun marked this conversation as resolved.

/* We need to check if we should close this socket here now */
if (httpResponseData->shouldCloseConnection()) {
if ((httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING) == 0) {
if (((AsyncSocket<SSL> *) s)->hasFullyDrained()) {
((AsyncSocket<SSL> *) s)->shutdown();
/* We need to force close after sending FIN since we want to hinder
* clients from keeping to send their huge data */
((AsyncSocket<SSL> *) s)->close();
((HttpResponse<SSL> *) s)->shutdownAndClose(httpResponseData);
}
}
}
Expand Down Expand Up @@ -899,19 +964,40 @@ struct HttpContext {
}
}
if (responseDone && asyncSocket->hasFullyDrained()) {
asyncSocket->shutdown();
/* We need to force close after sending FIN since we want to hinder
* clients from keeping to send their huge data */
asyncSocket->close();
reinterpret_cast<HttpResponse<SSL> *>(s)->shutdownAndClose(httpResponseData);
}
}

/* Expect another writable event, or another request within the timeout */
reinterpret_cast<HttpResponse<SSL> *>(s)->resetTimeout();

if constexpr (!IsNodeHttp) {
return replayParkedRequestsIfResponseComplete(s);
}
return s;
}

/* Bun.serve: the tail of every writable dispatch. Nothing of the completed
* response is on the stack here, and onWritable's close gate has run, so a
* connection that is closing replays nothing. */
static us_socket_t *replayParkedRequestsIfResponseComplete(us_socket_t *s) {
if (us_socket_is_closed(s) || us_socket_is_shut_down(s)) {
return s;
}
auto *httpResponseData = reinterpret_cast<HttpResponseData<SSL> *>(us_socket_ext(s));
if (httpResponseData->parkedRequestBytes.isEmpty()
|| (httpResponseData->state & HttpResponseData<SSL>::HTTP_RESPONSE_PENDING)
|| !reinterpret_cast<AsyncSocket<SSL> *>(s)->hasFullyDrained()) {
return s;
}
reinterpret_cast<AsyncSocket<SSL> *>(s)->resume();
/* us_socket_resume closes a socket that the kernel does not take back. */
if (us_socket_is_closed(s)) {
return s;
}
return replayParkedRequestBytes<false>(s);
Comment thread
robobun marked this conversation as resolved.
}

template <bool IsNodeHttp>
static us_socket_t *onEnd(us_socket_t *s) {
auto *asyncSocket = reinterpret_cast<AsyncSocket<SSL> *>(s);
Expand Down
60 changes: 44 additions & 16 deletions packages/bun-uws/src/HttpParser.h
Original file line number Diff line number Diff line change
Expand Up @@ -30,8 +30,10 @@
#include <algorithm>
#include <chrono>
#include <climits>
#include <functional>
#include <string_view>
#include <span>
#include <utility>
#include <wtf/Vector.h>
#include "MoveOnlyFunction.h"
#include "ChunkedEncoding.h"
Expand Down Expand Up @@ -633,15 +635,36 @@ struct HttpResponseData;
private:
std::string fallback;
public:
/* node:http flood prevention. HTTP_NODE_READS_PAUSED (state bit) = the socket's raw reads are
* paused and stays set through spill replay; this flag = "the parse loop running now must stop
* at the next request boundary and park the rest", cleared for replay so it can make progress. */
bool nodeHttpParkAtNextBoundary = false;
/* The parse loop must stop at the next request boundary and park the rest in
* parkedRequestBytes. Bun.serve: a response is pending (HttpContext::onData
* derives it). node:http: set on the flood-prevention pause edge, cleared for
* the replay so it can make progress (HTTP_NODE_READS_PAUSED stays set). */
bool parkAtNextBoundary = false;
bool nodeHttpSpillReplayScheduled = false;
/* A request on this connection had Connection: close or was HTTP/1.0, or a Bun.serve response closed it (RFC 9112 9.6). */
bool sawConnectionClose = false;
WTF::Vector<char> nodeHttpPausedSpill;
/* Where the parked bytes start in parkedRequestBytes: a replay that parks
* again gives its buffer back instead of copying what it did not reach. */
unsigned int parkedRequestBytesStart = 0;
WTF::Vector<char> parkedRequestBytes;
/* The buffer being replayed, during HttpContext::replayParkedRequestBytes. */
WTF::Vector<char> *replayedRequestBytes = nullptr;
private:
/* In a replay, [data, data + length) is the tail of the replayed buffer: it
* comes back whole, less the replay's fence, and only the start moves. */
void parkRequestBytes(char *data, unsigned int length) {
if (WTF::Vector<char> *replayed = std::exchange(replayedRequestBytes, nullptr)) {
char *begin = replayed->mutableSpan().data();
if (parkedRequestBytes.isEmpty() && std::greater_equal<char *>{}(data, begin)
&& std::less_equal<char *>{}(data + length, begin + replayed->size())) {
parkedRequestBytesStart = (unsigned int) (data - begin);
replayed->shrink(parkedRequestBytesStart + length);
parkedRequestBytes = std::exchange(*replayed, {});
return;
}
}
parkedRequestBytes.append(std::span<const char>(data, length));
}
/* This guy really has only 30 bits since we reserve two highest bits to chunked encoding parsing state */
uint64_t remainingStreamingBytes = 0;

Expand Down Expand Up @@ -1150,16 +1173,22 @@ struct HttpResponseData;
consumedTotal += length;
return HttpParserResult::success(consumedTotal, returnedUser);
}
/* node:http flood prevention: a dispatch earlier in this buffer paused reads.
* Stop at this request boundary, park the rest, report it as consumed so the
* caller does not spill it into the size-capped header fallback buffer. */
if constexpr (IsNodeHttp) {
if (nodeHttpParkAtNextBoundary) [[unlikely]] {
nodeHttpPausedSpill.append(std::span<const char>(data, length));
consumedTotal += length;
return HttpParserResult::success(consumedTotal, user);
/* Bun.serve: a closing connection takes nothing more (RFC 9112 9.6). Ahead
* of the park, which pauses reads: a close over bytes left unread resets
* the connection behind the complete response. */
if constexpr (!IsNodeHttp) {
if (sawConnectionClose) [[unlikely]] {
return HttpParserResult::success(consumedTotal + length, user);
}
}
/* Before getHeaders touches the next head. Reported as consumed so the
* caller does not spill it into the size-capped fallback buffer. Reads
* that arrive while bytes are parked go behind them, to keep wire order. */
if (parkAtNextBoundary || !parkedRequestBytes.isEmpty()) [[unlikely]] {
parkRequestBytes(data, length);
consumedTotal += length;
return HttpParserResult::success(consumedTotal, user);
}
/* RFC 9112 2.2: ignore empty lines (CRLF) received prior to the
* request-line, like Node/llhttp - e.g. a stray "\r\n" sent on an
* idle keep-alive connection must not be treated as a bad request.
Expand All @@ -1182,11 +1211,10 @@ struct HttpResponseData;
}
}
/* Must stay below the tunnel check, the park and the CR/LF skip, like llhttp's closed state. */
if (sawConnectionClose) {
if constexpr (IsNodeHttp) {
if constexpr (IsNodeHttp) {
if (sawConnectionClose) {
return HttpParserResult::error(HTTP_ERROR_400_BAD_REQUEST, HTTP_PARSER_ERROR_CLOSED_CONNECTION);
}
return HttpParserResult::success(consumedTotal + length, user);
}
auto result = getHeaders(data, data + length, req->headers, req->ancientHttp, isConnectRequest, useStrictMethodValidation, useInsecureHTTPParser, maxHeaderSize);
if(result.isError()) {
Expand Down
Loading
Loading