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
1 change: 1 addition & 0 deletions src/js/internal/http.ts
Original file line number Diff line number Diff line change
Expand Up @@ -77,6 +77,7 @@ export const enum NodeHTTPResponseFlags {
request_has_completed = 1 << 1,
ended = 1 << 2,
upgraded = 1 << 3,
dispatch_threw_while_queued = 1 << 9,

closed_or_completed = socket_closed | request_has_completed,
}
Expand Down
56 changes: 45 additions & 11 deletions src/js/node/_http_server.ts
Original file line number Diff line number Diff line change
Expand Up @@ -690,6 +690,13 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
socket = new (getNodeHTTPServerSocket())(server, socketHandle, !!tls);
}

if (isPipelinedDispatch) {
// The native handle holds this request's turn until its response exists, also when a throw prevents that.
(socket[kPipelinedResponses] ??= []).push(handle);
// For a throw before the kick below. drainMicrotasks() further down runs this one when nothing throws.
kickPipelineIfIdle(server, socket);
}

// Like Node.js's resetSocketTimeout (parserOnIncoming): a new request
// arriving on a kept-alive connection replaces the keep-alive idle
// timeout with the server's regular per-socket timeout.
Expand Down Expand Up @@ -890,14 +897,12 @@ Server.prototype[kRealListen] = function (tls, port, host, socketPath, reusePort
// Node.js, this response is queued (res.socket === null) and its
// writes are buffered until the in-flight response finishes and the
// pipeline assigns it the socket (advanceResponsePipeline).
socket[kPipelinedResponses]?.pop(); // the turn that the native handle held
queuePipelinedResponse(socket, http_res, !!isAncientHTTP);
// A pipelined dispatch can arrive after the previous response finished and detached
// (bytes still flushing keep it pending), leaving nothing in flight to advance the
// queue. Kick the pipeline once this dispatch settles.
if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) {
socket[kPipelineKickScheduled] = true;
process.nextTick(advancePipelineIfIdleNT, server, socket);
}
kickPipelineIfIdle(server, socket);
// Node's parserOnIncoming stops reading the connection once the bytes
// queued on responses that do not own the socket yet reach the
// socket's high water mark, so pipelined requests cannot flood it.
Expand Down Expand Up @@ -2498,6 +2503,12 @@ function releasePipelineOutgoingData(socket, bytes) {
// connection's current response, is assigned the socket, and its buffered
// output is flushed.
const kPipelineKickScheduled = Symbol("kPipelineKickScheduled");
function kickPipelineIfIdle(server, socket) {
if (socket._httpMessage == null && !socket[kPipelineKickScheduled]) {
socket[kPipelineKickScheduled] = true;
process.nextTick(advancePipelineIfIdleNT, server, socket);
}
}
function advancePipelineIfIdleNT(server, socket) {
socket[kPipelineKickScheduled] = false;
if (socket._httpMessage == null && socket[kPipelinedResponses]?.length) {
Expand Down Expand Up @@ -2533,6 +2544,8 @@ function abortQueuedPipelinedResponses(socket) {
socket[kPipelinedResponses] = undefined;
for (let i = 0; i < pipelinedLength; i++) {
const queuedRes = pipelined[i];
// A turn that the native handle still holds: the native close path notifies that one.
if (queuedRes[kPipelinedQueuedState] === undefined) continue;
const queuedReq = queuedRes.req;
if (queuedReq && !queuedReq.destroyed) {
queuedReq[kHandle] = undefined;
Expand All @@ -2551,6 +2564,15 @@ function abortQueuedPipelinedResponses(socket) {
}
}

// The head of the queue can never be sent: close once the bytes of the responses ahead of it have left.
function closeAfterLastSendableResponse(socket) {
if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) {
socket[kHandle]?.closeWhenDrained();
} else if (!socket.destroyed) {
socket.destroy();
}
}
Comment thread
claude[bot] marked this conversation as resolved.

function advanceResponsePipeline(server, socket) {
// The previous response on this connection closed it (Connection: close,
// HTTP/1.0, maxRequestsPerSocket): like Node.js's resOnFinish, advancing
Expand All @@ -2564,24 +2586,36 @@ function advanceResponsePipeline(server, socket) {
if (!queue || queue.length === 0) {
return;
}
const res = queue.shift();
const res = queue[0];
const queued = res[kPipelinedQueuedState];
res[kPipelinedQueuedState] = undefined;
releasePipelineOutgoingData(socket, queued.bytes);
if (queued === undefined) {
// The native handle still holds this turn: its dispatch queued no response, and after a throw none can come.
if ((res.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0) {
closeAfterLastSendableResponse(socket);
}
return;
}
const handle = res[kHandle];

if (res.destroyed || !handle) {
if (
res.destroyed ||
!handle ||
// Its dispatch threw and nothing ended it since: natively only the current response is answered for a throw.
(!queued.ended && (handle.flags & NodeHTTPResponseFlags.dispatch_threw_while_queued) !== 0)
) {
// The queued response was destroyed before it could be sent; the
// connection cannot produce a response for this slot, so it is unusable.
// Deliberate divergence from Node v26, which assigns the destroyed
// message and wedges the connection until requestTimeout: an HTTP/1.1
// connection cannot skip a response slot, so reset it instead.
if (!socket.destroyed) {
socket.destroy();
}
closeAfterLastSendableResponse(socket); // it stays queued for the close path
return;
}

queue.shift();
res[kPipelinedQueuedState] = undefined;
releasePipelineOutgoingData(socket, queued.bytes);

if (NodeHTTPServerSocket && socket instanceof NodeHTTPServerSocket) {
const socketHandle = socket[kHandle];
if (
Expand Down
23 changes: 23 additions & 0 deletions src/jsc/bindings/node/JSNodeHTTPServerSocket.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -280,6 +280,29 @@ bool JSNodeHTTPServerSocket::shutdownAfterResponseDrains()
return deferShutdownUntilResponseDrains<false>(socket);
}

template<bool SSL>
static void closeWhenDrainedImpl(us_socket_t* socket)
{
auto* httpResponseData = reinterpret_cast<uWS::HttpResponseData<SSL>*>(us_socket_ext(socket));
/* uWS's close gates (below, or onWritable after the flush) close a connection marked like this. */
httpResponseData->state |= uWS::HttpResponseData<SSL>::HTTP_CONNECTION_CLOSE;
/* A response that ended inside the read being parsed is still in the cork buffer. */
reinterpret_cast<uWS::AsyncSocket<SSL>*>(socket)->uncork();
reinterpret_cast<uWS::HttpResponse<SSL>*>(socket)->closeIfDoneAndMarked(httpResponseData);
}

void JSNodeHTTPServerSocket::closeWhenDrained()
{
if (!socket || upgraded || us_socket_is_closed(socket)) {
return;
}
if (is_ssl) {
closeWhenDrainedImpl<true>(socket);
} else {
closeWhenDrainedImpl<false>(socket);
}
}

template<bool SSL>
static bool isRequestTimedOutImpl(us_socket_t* socket, uint64_t headersTimeoutMs, uint64_t requestTimeoutMs)
{
Expand Down
3 changes: 3 additions & 0 deletions src/jsc/bindings/node/JSNodeHTTPServerSocket.h
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,9 @@ class JSNodeHTTPServerSocket : public JSC::JSDestructibleObject {
* truncate the response. Returns true after handing the close to uWS. */
bool shutdownAfterResponseDrains();

/* Close once the bytes of the responses that ended have left. close() discards them, end() waits for the peer. */
void closeWhenDrained();

/* Switch the connection into CONNECT-style tunnel mode after an accepted
* Upgrade: subsequent bytes bypass the HTTP parser and stream to the
* ondata callback as opaque data. With afterBody, the switch is deferred
Expand Down
11 changes: 11 additions & 0 deletions src/jsc/bindings/node/JSNodeHTTPServerSocketPrototype.cpp
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketSetResponseTrailers);
JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketIsRequestTimedOut);
JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStartPipelinedResponse);
JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStopParsing);
JSC_DECLARE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketCloseWhenDrained);
JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterResponse);
JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterRemoteAddress);
JSC_DECLARE_CUSTOM_GETTER(jsNodeHttpServerSocketGetterLocalAddress);
Expand Down Expand Up @@ -83,6 +84,7 @@ static const JSC::HashTableValue JSNodeHTTPServerSocketPrototypeTableValues[] =
{ "isRequestTimedOut"_s, static_cast<unsigned>(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketIsRequestTimedOut, 2 } },
{ "startPipelinedResponse"_s, static_cast<unsigned>(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketStartPipelinedResponse, 3 } },
{ "stopParsing"_s, static_cast<unsigned>(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketStopParsing, 0 } },
{ "closeWhenDrained"_s, static_cast<unsigned>(JSC::PropertyAttribute::Function | JSC::PropertyAttribute::DontEnum), JSC::NoIntrinsic, { JSC::HashTableValue::NativeFunctionType, jsFunctionNodeHTTPServerSocketCloseWhenDrained, 0 } },
{ "secureEstablished"_s, static_cast<unsigned>(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterIsSecureEstablished, noOpSetter } },
{ "servername"_s, static_cast<unsigned>(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterServername, noOpSetter } },
{ "authorizationError"_s, static_cast<unsigned>(JSC::PropertyAttribute::CustomAccessor | JSC::PropertyAttribute::ReadOnly), JSC::NoIntrinsic, { JSC::HashTableValue::GetterSetterType, jsNodeHttpServerSocketGetterAuthorizationError, noOpSetter } },
Expand Down Expand Up @@ -203,6 +205,15 @@ JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketStopParsing, (JSC::JSGlob
return JSValue::encode(JSC::jsUndefined());
}

JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketCloseWhenDrained, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame))
{
auto* thisObject = dynamicDowncast<JSNodeHTTPServerSocket>(callFrame->thisValue());
if (thisObject) [[likely]] {
thisObject->closeWhenDrained();
}
return JSValue::encode(JSC::jsUndefined());
}

JSC_DEFINE_HOST_FUNCTION(jsFunctionNodeHTTPServerSocketWrite, (JSC::JSGlobalObject * globalObject, JSC::CallFrame* callFrame))
{
auto* thisObject = dynamicDowncast<JSNodeHTTPServerSocket>(callFrame->thisValue());
Expand Down
14 changes: 14 additions & 0 deletions src/runtime/server/NodeHTTPResponse.rs
Original file line number Diff line number Diff line change
Expand Up @@ -95,6 +95,8 @@ bitflags! {
/// node:http handed this connection to a raw 'upgrade'/'connect'
/// tunnel (JSNodeHTTPServerSocket::upgradeToTunnelMode).
const TUNNELED = 1 << 8;
/// Its dispatch threw while it was queued (pipelining): advanceResponsePipeline decides at its turn.
const DISPATCH_THREW_WHILE_QUEUED = 1 << 9;
}
}

Expand Down Expand Up @@ -442,6 +444,18 @@ impl NodeHTTPResponse {
Bun__getNodeHTTPResponseThisValue(any_response_is_ssl(&raw), raw.socket().cast())
}

/// Flags this response when another one is the connection's current response, and says so.
pub(crate) fn mark_dispatch_threw_if_queued(&self) -> bool {
let queued = self
.get_this_value()
.as_class_ref::<Self>()
.is_some_and(|current| !ptr::eq(current, self));
if queued {
self.update_flags(|f| f.insert(Flags::DISPATCH_THREW_WHILE_QUEUED));
}
queued
}

fn get_server_socket_value(&self) -> JSValue {
let flags = self.flags.get();
let Some(raw) = self.raw_response.get() else {
Expand Down
7 changes: 6 additions & 1 deletion src/runtime/server/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1460,7 +1460,12 @@ impl<const SSL: bool, const DEBUG: bool> NewServer<SSL, DEBUG> {
)
};

if !node_http_response.is_null() {
// A pipelined response stays queued: `raw_response` describes the one ahead of it.
let threw_while_queued = !node_http_response.is_null()
// SAFETY: see `nhr` above.
&& unsafe { &*node_http_response }.mark_dispatch_threw_if_queued();

if !node_http_response.is_null() && !threw_while_queued {
// SAFETY: see `nhr` above.
let nhr = unsafe { &*node_http_response };
let nhr_flags = nhr.flags.get();
Expand Down
Loading
Loading