diff --git a/docs/runtime/http/server.mdx b/docs/runtime/http/server.mdx index 2692a9f700b4..fc935f9b93a6 100644 --- a/docs/runtime/http/server.mdx +++ b/docs/runtime/http/server.mdx @@ -298,7 +298,18 @@ await server.stop(); await server.stop(true); ``` -By default, `stop()` allows in-flight requests and WebSocket connections to complete. Pass `true` to immediately terminate all connections. +By default, `stop()` allows in-flight requests and WebSocket connections to complete. Idle keep-alive connections are closed immediately, and connections with a request in flight close once their response has been sent. Pass `true` to immediately terminate all connections instead. The returned promise resolves once every connection has closed. + +### `server.closeIdleConnections()` + +To close keep-alive connections that are not currently serving a request, without stopping the server: + +```ts +const closed = server.closeIdleConnections(); +console.log(`closed ${closed} idle connections`); +``` + +It returns the number of connections it closed. Connections with a request in flight and open WebSockets are untouched, and the server keeps accepting new connections. This mirrors `node:http`'s `server.closeIdleConnections()`, which returns nothing. ### `server.ref()` and `server.unref()` @@ -558,6 +569,12 @@ interface Server extends Disposable { */ stop(closeActiveConnections?: boolean): Promise; + /** + * Close idle keep-alive connections without stopping the server. + * @returns The number of connections closed + */ + closeIdleConnections(): number; + /** * Update handlers without restarting the server. * Only fetch and error handlers can be updated. diff --git a/packages/bun-types/serve.d.ts b/packages/bun-types/serve.d.ts index fee1ee6e0f2a..4dbe707dba76 100644 --- a/packages/bun-types/serve.d.ts +++ b/packages/bun-types/serve.d.ts @@ -910,13 +910,29 @@ declare module "bun" { /** * Stop listening to prevent new connections from being accepted. * - * By default, it does not cancel in-flight requests or websockets. That means it may take some time before all network activity stops. + * By default, it does not cancel in-flight requests or websockets. Idle + * keep-alive connections are closed right away, and connections with a + * request in flight close as soon as their response completes. That means + * it may take some time before all network activity stops. + * + * The returned promise resolves once every connection is closed. * * @param closeActiveConnections Immediately terminate in-flight requests, websockets, and stop accepting new connections. * @default false */ stop(closeActiveConnections?: boolean): Promise; + /** + * Close every connection that is not currently sending a request or + * waiting for a response, without stopping the server. + * + * In-flight requests and open WebSockets are untouched, and the server + * keeps accepting new connections. + * + * @returns The number of connections that were closed. + */ + closeIdleConnections(): number; + /** * Update the `fetch` and `error` handlers without restarting the server. * diff --git a/packages/bun-uws/src/App.h b/packages/bun-uws/src/App.h index 03fc56452570..eb5ee2ba132d 100644 --- a/packages/bun-uws/src/App.h +++ b/packages/bun-uws/src/App.h @@ -397,20 +397,28 @@ struct TemplatedApp { return std::move(*this); } - /** Closes all connections connected to this server which are not sending a request or waiting for a response. Does not close the listen socket. */ - TemplatedApp &&closeIdle() { + /** Closes all connections connected to this server which are not sending a request or waiting for a response. Does not close the listen socket. + * With closeWhenIdle set, connections that are busy right now are marked to close as soon as their in-flight work completes (graceful shutdown); + * upgraded WebSockets and CONNECT/Upgrade tunnels never become idle, so they are left alone either way. + * Returns the number of connections closed. */ + size_t closeIdle(bool closeWhenIdle = false) { auto *group = httpContext->getSocketGroup(); struct us_socket_t *s = group->head_sockets; + size_t closed = 0; while (s) { - // no matter the type of socket will always contain the AsyncSocketData - auto *data = ((AsyncSocket *) s)->getAsyncSocketData(); + /* The HTTP group only holds HTTP sockets (an upgraded WebSocket is + * adopted into its own group), so the ext block is an HttpResponseData. */ + auto *data = (HttpResponseData *) ((AsyncSocket *) s)->getAsyncSocketData(); struct us_socket_t *next = s->next; if (data->isIdle) { us_socket_close(s, LIBUS_SOCKET_CLOSE_CODE_CLEAN_SHUTDOWN, 0); + closed++; + } else if (closeWhenIdle) { + data->state |= HttpResponseData::HTTP_CLOSE_WHEN_IDLE; } s = next; } - return std::move(*this); + return closed; } template diff --git a/packages/bun-uws/src/HttpContext.h b/packages/bun-uws/src/HttpContext.h index d0870e94fe44..2a88ad687ad8 100644 --- a/packages/bun-uws/src/HttpContext.h +++ b/packages/bun-uws/src/HttpContext.h @@ -328,8 +328,12 @@ struct HttpContext { /* Cork this socket */ ((AsyncSocket *) s)->cork(); - /* Mark that we are inside the parser now */ + /* Mark that we are inside the parser now. Save/restore the parsed + * socket: node:http's read replay can nest a parse inside another + * socket's dispatch. */ httpContextData->flags.isParsingHttp = true; + struct us_socket_t *prevParsingSocket = httpContextData->parsingSocket; + httpContextData->parsingSocket = s; httpResponseData->isIdle = false; /* node:http compat: maintain the headers/request timeout window (see @@ -581,6 +585,7 @@ struct HttpContext { /* Mark that we are no longer parsing Http */ httpContextData->flags.isParsingHttp = false; + httpContextData->parsingSocket = prevParsingSocket; /* If we got fullptr that means the parser wants us to close the socket from error (same as calling the errorHandler) */ if (httpErrorStatusCode) { /* node:http compat: parse errors surface as the server's 'clientError' diff --git a/packages/bun-uws/src/HttpContextData.h b/packages/bun-uws/src/HttpContextData.h index 6767a8a481aa..d2a0a47267af 100644 --- a/packages/bun-uws/src/HttpContextData.h +++ b/packages/bun-uws/src/HttpContextData.h @@ -68,6 +68,14 @@ struct alignas(16) HttpContextData { /* This is the currently browsed-to router when using SNI */ HttpRouter *currentRouter = &router; + /* The socket onData is currently parsing, nullptr outside a parse. The + * close gates in internalEnd need the per-socket identity: a DIFFERENT + * socket's response can complete inside this window (a microtask drained + * during a request dispatch), and the context-wide isParsingHttp bit + * alone would wrongly defer its close to a post-parse gate that only + * checks the parsed socket. */ + struct us_socket_t *parsingSocket = nullptr; + /* This is the default router for default SNI or non-SSL */ HttpRouter router; void *upgradedWebSocket = nullptr; diff --git a/packages/bun-uws/src/HttpResponse.h b/packages/bun-uws/src/HttpResponse.h index 07934d33df1e..ac140fa59a73 100644 --- a/packages/bun-uws/src/HttpResponse.h +++ b/packages/bun-uws/src/HttpResponse.h @@ -96,6 +96,24 @@ struct HttpResponse : public AsyncSocket { getHttpResponseData()->state |= HttpResponseData::HTTP_WROTE_DATE_HEADER; } + /* Shutdown+close when the connection is marked to close (Connection: + * close, peer FIN, close-when-idle), the response is complete and every + * outgoing byte has been flushed. Returns true when the socket was closed. */ + bool closeIfDoneAndMarked(HttpResponseData *httpResponseData) { + if (httpResponseData->shouldCloseConnection()) { + if ((httpResponseData->state & HttpResponseData::HTTP_RESPONSE_PENDING) == 0) { + if (((AsyncSocket *) this)->hasFullyDrained()) { + ((AsyncSocket *) this)->shutdown(); + /* We need to force close after sending FIN since we want to hinder + * clients from keeping to send their huge data */ + ((AsyncSocket *) this)->close(); + return true; + } + } + } + return false; + } + /* Returns true on success, indicating that it might be feasible to write more data. * Will start timeout if stream reaches totalSize or write failure. * keepCorked: if true, skip the trailing uncork so the caller can batch @@ -189,19 +207,22 @@ struct HttpResponse : public AsyncSocket { /* We need to check if we should close this socket here now */ if (!Super::isCorked()) { - if (httpResponseData->shouldCloseConnection()) { - if ((httpResponseData->state & HttpResponseData::HTTP_RESPONSE_PENDING) == 0) { - if (((AsyncSocket *) this)->hasFullyDrained()) { - ((AsyncSocket *) this)->shutdown(); - /* We need to force close after sending FIN since we want to hinder - * clients from keeping to send their huge data */ - ((AsyncSocket *) this)->close(); - return true; - } - } + if (closeIfDoneAndMarked(httpResponseData)) { + return true; } } else if (!keepCorked) { this->uncork(); + /* That uncork released our cork slot, so the cork() wrapper's + * post-uncork close gate will not run. When THIS socket is the + * one being parsed, onData's post-parse gate closes it once + * the buffer is fully consumed; any other socket (an async + * handler completing, possibly inside another socket's parse + * window via a drained microtask) gets no later gate, so close + * here. */ + if (HttpContext::fromSocket((us_socket_t *) this)->getSocketContextData()->parsingSocket != (us_socket_t *) this + && closeIfDoneAndMarked(httpResponseData)) { + return true; + } } /* tryEnd can never fail when in chunked mode, since we do not have tryWrite (yet), only write */ @@ -255,18 +276,15 @@ struct HttpResponse : public AsyncSocket { /* We need to check if we should close this socket here now */ if (!Super::isCorked()) { - if (httpResponseData->shouldCloseConnection()) { - if ((httpResponseData->state & HttpResponseData::HTTP_RESPONSE_PENDING) == 0) { - if (((AsyncSocket *) this)->hasFullyDrained()) { - ((AsyncSocket *) this)->shutdown(); - /* We need to force close after sending FIN since we want to hinder - * clients from keeping to send their huge data */ - ((AsyncSocket *) this)->close(); - } - } - } + closeIfDoneAndMarked(httpResponseData); } else if (!keepCorked) { this->uncork(); + /* Same as the chunked arm above: the cork slot is gone, so + * run the close gate here unless THIS socket is the one + * being parsed (then onData's post-parse gate handles it). */ + if (HttpContext::fromSocket((us_socket_t *) this)->getSocketContextData()->parsingSocket != (us_socket_t *) this) { + closeIfDoneAndMarked(httpResponseData); + } } } @@ -591,6 +609,12 @@ struct HttpResponse : public AsyncSocket { * Starts a timeout in some cases. Returns [ok, hasResponded] */ std::pair tryEnd(std::string_view data, uintmax_t totalSize = 0, bool closeConnection = false) { bool ok = internalEnd(data, totalSize, true, true, closeConnection); + /* internalEnd's close gate may have closed the socket (destructing the + * ext hasResponded() reads); that only happens once the response has + * completed, so report responded without touching it. */ + if (us_socket_is_closed((us_socket_t *) this)) { + return {ok, true}; + } return {ok, hasResponded()}; } diff --git a/packages/bun-uws/src/HttpResponseData.h b/packages/bun-uws/src/HttpResponseData.h index c27b94add094..c64ff4493f79 100644 --- a/packages/bun-uws/src/HttpResponseData.h +++ b/packages/bun-uws/src/HttpResponseData.h @@ -59,7 +59,9 @@ struct HttpResponseData : AsyncSocketData, HttpParser { this->state &= ~HttpResponseData::HTTP_RESPONSE_PENDING; HttpResponseData *httpResponseData = uwsRes->getHttpResponseData(); - httpResponseData->isIdle = true; + /* A queued pipelined response (node:http) still owes output on this + * connection, so it is not idle between the responses. */ + httpResponseData->isIdle = httpResponseData->nodeHttpQueuedPipelinedCount == 0; } /* Caller of onWritable. It is possible onWritable calls markDone so we need to borrow it. */ @@ -143,13 +145,19 @@ struct HttpResponseData : AsyncSocketData, 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, + /* Close this connection the next time it is idle (no request being + * received, no response in flight or queued). Set by + * App::closeIdle(true) on connections that were busy during a graceful + * shutdown sweep; the shouldCloseConnection() gates act on it once the + * in-flight work completes. */ + HTTP_CLOSE_WHEN_IDLE = 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 * keep-alive connection, so starting a new response clears the rest of the * word (resetResponseState) - these have to survive that. */ HTTP_CONNECTION_SCOPED = HTTP_NODE_PARSING_STOPPED | HTTP_NODE_READS_PAUSED - | HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN, + | HTTP_NODE_TUNNEL_AFTER_BODY | HTTP_NODE_RECEIVED_FIN | HTTP_CLOSE_WHEN_IDLE, }; /* Begin a new response on this connection. Clearing the word in one go is @@ -158,6 +166,9 @@ struct HttpResponseData : AsyncSocketData, HttpParser { * keep-alive socket; only the connection-scoped bits are carried over. */ void resetResponseState() { state = (state & HTTP_CONNECTION_SCOPED) | HTTP_RESPONSE_PENDING; + /* A response is in flight again (a new request dispatched, or a queued + * pipelined response activated), so the connection is not idle. */ + this->isIdle = false; } /* Set or clear a flag from a runtime bool. */ @@ -214,7 +225,8 @@ struct HttpResponseData : AsyncSocketData, HttpParser { * any) has completed and all buffered outgoing data has been flushed. */ bool shouldCloseConnection() const { return (state & HTTP_CONNECTION_CLOSE) - || ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0); + || ((state & HTTP_NODE_RECEIVED_FIN) && nodeHttpQueuedPipelinedCount == 0) + || ((state & HTTP_CLOSE_WHEN_IDLE) && this->isIdle); } #ifdef UWS_WITH_PROXY diff --git a/src/runtime/server/FileResponseStream.rs b/src/runtime/server/FileResponseStream.rs index a5ff5adc3737..2778b660b0db 100644 --- a/src/runtime/server/FileResponseStream.rs +++ b/src/runtime/server/FileResponseStream.rs @@ -479,6 +479,24 @@ impl FileResponseStream { let resp = self.resp.get(); resp.end_send_file(self.sendfile.get().offset, resp.should_close_connection()); (self.on_complete.get())(self.ctx.get(), resp); + // `end_send_file` bypasses every shouldCloseConnection() gate: it does + // not go through internalEnd, and the onWritable gate is skipped + // because this frame returns `false` to it. Run the gate here — after + // `on_complete`, which must see a live socket — so Connection: close + // and the graceful-stop close-when-idle mark actually close. + // + // `resp` is still valid here: usockets never frees a socket + // synchronously — us_socket_close only links it onto the loop's + // closed list, freed by us_internal_free_closed_sockets at the end of + // the loop iteration — so the allocation outlives this frame no + // matter what `on_complete` did (the same invariant that makes + // passing `resp` to `on_complete` after the end sound). It is still + // *this* HTTP socket: an upgrade (us_socket_adopt) is only reachable + // from a live in-flight request, and this one just completed. And if + // anything in the frame closed it, the shim's leading + // us_socket_is_closed check returns before touching the destructed + // ext. Only `finish()` runs after this, and it never touches `resp`. + resp.close_if_done_and_marked(); self.finish(); } @@ -543,6 +561,10 @@ impl FileResponseStream { let resp = self.resp.get(); resp.end_without_body(resp.should_close_connection()); (self.on_complete.get())(self.ctx.get(), resp); + // This end runs uncorked (reader callbacks), so no cork or parser + // gate will run the close check; do it here, after `on_complete` + // like `end_sendfile`, so the callbacks see a live socket. + resp.close_if_done_and_marked(); } // Release the owner ref from `heap::into_raw` in `start()`. Every entry diff --git a/src/runtime/server/RequestContext.rs b/src/runtime/server/RequestContext.rs index 9a7e2f089d16..7797c93dceb6 100644 --- a/src/runtime/server/RequestContext.rs +++ b/src/runtime/server/RequestContext.rs @@ -1274,6 +1274,12 @@ where self.detach_response(); // SAFETY: FFI handle resp.end_without_body(close_connection); + // This end can run uncorked (e.g. render_production_error from a + // rejection microtask), where no cork or parser gate runs the + // close check for Connection: close or a graceful-stop mark. The + // shim no-ops when the socket is corked (the cork wrapper's own + // gate runs later) or already closed. + resp.close_if_done_and_marked(); // end_request_streaming_and_drain() must run after the last // `resp` access: its drain_microtasks() can re-enter lsquic (H3) // and free the stream out from under the local `resp` copy. diff --git a/src/runtime/server/mod.rs b/src/runtime/server/mod.rs index 1f5baac3e14a..0fa188e5010d 100644 --- a/src/runtime/server/mod.rs +++ b/src/runtime/server/mod.rs @@ -1700,6 +1700,26 @@ impl NewServer { if !abrupt { // S012: `app::ListenSocket` is a ZST opaque — safe deref. bun_opaque::opaque_deref_mut(listener).close(); + // Close idle keep-alive connections now and mark busy ones to + // close once their in-flight work completes; open websockets are + // untouched and drain on their own. Each close reaches + // `on_connection_filter(-1)` synchronously, so hold the guard so + // that path cannot form a second `&mut self` under this frame — + // `stop()` runs `deinit_if_we_can` right after this returns. + // + // node:http servers are exempt: Node's `close()` sweeps idle + // connections exactly once (the JS layer already called + // `closeIdleConnections()`), and a connection whose response + // completes after `close()` stays keep-alive until its timeout + // reaps it — verified against Node v26. + if self.config.on_node_http_request.is_empty() { + if let Some(app) = self.app { + self.deinit_running.set(true); + // S012: `NewApp` is a ZST opaque — safe `*mut → &mut` deref. + let _closed = bun_opaque::opaque_deref_mut(app).close_idle_connections(true); + self.deinit_running.set(false); + } + } } else if !self.flags.contains(ServerFlags::TERMINATED) { if let Some(ws) = self.config.websocket.as_mut() { ws.handler.app = None; diff --git a/src/runtime/server/server_body.rs b/src/runtime/server/server_body.rs index 2867551551ae..200ea2acbeee 100644 --- a/src/runtime/server/server_body.rs +++ b/src/runtime/server/server_body.rs @@ -2607,16 +2607,17 @@ where _callframe: &CallFrame, ) -> JsResult { if self.app.is_none() || self.deinit_running.get() { - return Ok(JSValue::UNDEFINED); + return Ok(JSValue::js_number(0.0)); } // On a Bun.serve server each close reaches `on_connection_filter(-1)` // synchronously; hold the guard so it cannot re-derive `&mut self` - // while this frame owns it. + // while this frame owns it. One-shot sweep (Node semantics): busy + // connections are spared and are NOT marked to close later. self.deinit_running.set(true); - self.app_mut().close_idle_connections(); + let closed = self.app_mut().close_idle_connections(false); self.deinit_running.set(false); self.deinit_if_we_can(); - Ok(JSValue::UNDEFINED) + Ok(JSValue::js_number(closed as f64)) } pub(crate) fn stop_from_js(&mut self, abruptly: Option) -> JSValue { diff --git a/src/uws_sys/App.rs b/src/uws_sys/App.rs index ddef6e230663..c051b5aeb846 100644 --- a/src/uws_sys/App.rs +++ b/src/uws_sys/App.rs @@ -100,8 +100,13 @@ impl App { c::uws_app_close(Self::SSL_FLAG, self.as_raw()) } - pub fn close_idle_connections(&mut self) { - c::uws_app_close_idle(Self::SSL_FLAG, self.as_raw()) + /// Close every HTTP connection that is idle (no request being received, no + /// response in flight). With `close_when_idle`, connections that are busy + /// right now are additionally marked to close as soon as their in-flight + /// work completes. Never touches WebSockets or the listen socket. + /// Returns the number of connections closed. + pub fn close_idle_connections(&mut self, close_when_idle: bool) -> usize { + c::uws_app_close_idle(Self::SSL_FLAG, self.as_raw(), i32::from(close_when_idle)) } pub fn create(opts: &BunSocketContextOptions) -> Option<*mut Self> { @@ -467,7 +472,11 @@ pub mod c { unsafe extern "C" { pub(crate) safe fn uws_app_close(ssl: i32, app: &mut uws_app_s); - pub(crate) safe fn uws_app_close_idle(ssl: i32, app: &mut uws_app_s); + pub(crate) safe fn uws_app_close_idle( + ssl: i32, + app: &mut uws_app_s, + close_when_idle: i32, + ) -> usize; // safe: `&mut uws_app_s` is ABI-identical to a non-null `*mut`; // `handler`/`user_data` are stored opaquely (never dereferenced by the // C++ shim itself) — no preconditions on this call. diff --git a/src/uws_sys/Response.rs b/src/uws_sys/Response.rs index 1cc6fe6c301f..d0650880c162 100644 --- a/src/uws_sys/Response.rs +++ b/src/uws_sys/Response.rs @@ -243,6 +243,13 @@ impl Response { ) } + /// Completion gate for end paths that bypass `internalEnd` (the sendfile + /// path): closes the socket when the connection is marked to close, the + /// response is complete, and every outgoing byte has been flushed. + pub(crate) fn close_if_done_and_marked(&mut self) { + c::uws_res_close_if_done_and_marked(Self::ssl_flag(), self.as_raw()) + } + pub fn timeout(&mut self, seconds: u8) { c::uws_res_timeout(Self::ssl_flag(), self.as_raw(), seconds) } @@ -726,6 +733,12 @@ impl AnyResponse { any_dispatch!(self, |r| r.end_send_file(write_offset, close_connection)) } + /// Completion gate for end paths that bypass `internalEnd` (the sendfile + /// path); see `Response::close_if_done_and_marked`. + pub fn close_if_done_and_marked(self) { + any_dispatch!(self, |r| r.close_if_done_and_marked()) + } + pub fn socket(self) -> *mut c::uws_res { match self { AnyResponse::H3(_) => panic!("socket() is not available for HTTP/3 responses"), @@ -1138,6 +1151,7 @@ pub mod c { ); pub(crate) safe fn uws_res_timeout(ssl: i32, res: &mut uws_res, timeout: u8); pub(crate) safe fn uws_res_reset_timeout(ssl: i32, res: &mut uws_res); + pub(crate) safe fn uws_res_close_if_done_and_marked(ssl: i32, res: &mut uws_res); pub(crate) safe fn uws_res_get_buffered_amount(ssl: i32, res: &mut uws_res) -> u64; pub(crate) fn uws_res_write( ssl: i32, diff --git a/src/uws_sys/h3.rs b/src/uws_sys/h3.rs index 655757c98eb7..da8962fa9ffa 100644 --- a/src/uws_sys/h3.rs +++ b/src/uws_sys/h3.rs @@ -104,6 +104,9 @@ impl Response { pub(crate) fn end_send_file(&mut self, write_offset: u64, close_connection: bool) { c::uws_h3_res_end_sendfile(self, write_offset, close_connection) } + /// H3 streams tear down through the QUIC engine; the TCP close-when-idle + /// gate has no equivalent here. + pub(crate) fn close_if_done_and_marked(&mut self) {} pub(crate) fn write(&mut self, data: &[u8]) -> WriteResult { let mut len: usize = data.len(); // SAFETY: self is a live FFI handle; data ptr valid for read; len out-ptr is a valid local diff --git a/src/uws_sys/libuwsockets.cpp b/src/uws_sys/libuwsockets.cpp index 81ae20f23cd0..989380a317b1 100644 --- a/src/uws_sys/libuwsockets.cpp +++ b/src/uws_sys/libuwsockets.cpp @@ -387,17 +387,17 @@ extern "C" } } - void uws_app_close_idle(int ssl, uws_app_t *app) + size_t uws_app_close_idle(int ssl, uws_app_t *app, int close_when_idle) { if (ssl) { uWS::SSLApp *uwsApp = (uWS::SSLApp *)app; - uwsApp->closeIdle(); + return uwsApp->closeIdle(close_when_idle != 0); } else { uWS::App *uwsApp = (uWS::App *)app; - uwsApp->closeIdle(); + return uwsApp->closeIdle(close_when_idle != 0); } } @@ -1186,6 +1186,11 @@ extern "C" void uws_res_pause(int ssl, uws_res_r res) { + /* No-op on a closed socket; see uws_res_on_aborted. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1200,6 +1205,12 @@ extern "C" void uws_res_resume(int ssl, uws_res_r res) { + /* No-op on a closed socket (resume's resetTimeout reads the destructed + * ext); see uws_res_on_aborted. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1339,6 +1350,36 @@ extern "C" uwsRes->resetTimeout(); } } + /* Completion gate for response-end paths that bypass internalEnd (the + * sendfile path): closes the socket when the connection is marked to close + * (Connection: close, peer FIN, close-when-idle), the response is complete, + * and every outgoing byte has been flushed. Corked responses are left to the + * cork() wrapper's own post-uncork gate. */ + void uws_res_close_if_done_and_marked(int ssl, uws_res_r res) + { + /* A callback upstream of this gate may already have closed the socket; + * onClose destructs the ext block, so bail before touching it. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } + if (ssl) + { + uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; + if (!uwsRes->AsyncSocket::isCorked()) + { + uwsRes->closeIfDoneAndMarked(uwsRes->getHttpResponseData()); + } + } + else + { + uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; + if (!uwsRes->AsyncSocket::isCorked()) + { + uwsRes->closeIfDoneAndMarked(uwsRes->getHttpResponseData()); + } + } + } void uws_res_reset_timeout(int ssl, uws_res_r res) { if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1379,6 +1420,12 @@ extern "C" data->state |= uWS::HttpResponseData::HTTP_END_CALLED; data->markDone(uwsRes); uwsRes->resetTimeout(); + /* No close gate here: callers (FileResponseStream::finish, + * DevServer/HTMLBundle error paths) keep using the response after this + * returns, so closing inside this call would destruct the ext under + * them. Corked callers get the cork() wrapper's post-uncork gate; + * uncorked ones run uws_res_close_if_done_and_marked themselves once + * they are done with the response. */ } else { @@ -1401,6 +1448,7 @@ extern "C" data->state |= uWS::HttpResponseData::HTTP_END_CALLED; data->markDone(uwsRes); uwsRes->resetTimeout(); + /* No close gate here; see the SSL arm above. */ } } @@ -1499,6 +1547,10 @@ extern "C" } void uws_res_clear_on_writable(int ssl, uws_res_r res) { + /* No-op on a closed socket; see uws_res_on_aborted. */ + if (us_socket_is_closed((struct us_socket_t *)res)) { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; uwsRes->clearOnWritable(); @@ -1512,6 +1564,14 @@ extern "C" void (*handler)(uws_res_r res, void *optional_data), void *optional_data) { + /* A closed socket's ext block is already destructed (HttpContext::onClose) + * and its callbacks can never fire again; registering or clearing one is a + * no-op. Completion bookkeeping runs after ends that may have closed the + * socket via a shouldCloseConnection() gate, so this must not touch ext. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1544,6 +1604,11 @@ extern "C" void (*handler)(uws_res_r res, void *optional_data), void *optional_data) { + /* No-op on a closed socket; see uws_res_on_aborted. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1578,6 +1643,11 @@ extern "C" void *optional_data), void *optional_data) { + /* No-op on a closed socket; see uws_res_on_aborted. */ + if (us_socket_is_closed((struct us_socket_t *)res)) + { + return; + } if (ssl) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; @@ -1835,7 +1905,10 @@ __attribute__((callback (corker, ctx))) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; auto pair = uwsRes->tryEnd(stringViewFromC(bytes, len), total_len, close); - if (pair.first) { + /* A completed tryEnd may have closed the socket through a + * shouldCloseConnection() gate, destructing the ext; markDone already + * cleared the callbacks in that case. */ + if (pair.first && !us_socket_is_closed((struct us_socket_t *)res)) { uwsRes->clearOnWritableAndAborted(); } @@ -1845,7 +1918,8 @@ __attribute__((callback (corker, ctx))) { uWS::HttpResponse *uwsRes = (uWS::HttpResponse *)res; auto pair = uwsRes->tryEnd(stringViewFromC(bytes, len), total_len, close); - if (pair.first) { + /* See the SSL arm above. */ + if (pair.first && !us_socket_is_closed((struct us_socket_t *)res)) { uwsRes->clearOnWritableAndAborted(); } diff --git a/test/js/bun/http/bun-server.test.ts b/test/js/bun/http/bun-server.test.ts index 13c302bc00bf..d45f7f3913e5 100644 --- a/test/js/bun/http/bun-server.test.ts +++ b/test/js/bun/http/bun-server.test.ts @@ -598,11 +598,13 @@ test("should be able to await server.stop()", async () => { }); describe.concurrent("server.stop() drain promise counts open connections", () => { - // The drain promise must not resolve while a keep-alive HTTP connection is - // still open. Previously it counted only in-flight requests (plus the - // listener and websockets), so an idle keep-alive socket let the promise - // resolve in 0 ms and the socket kept answering requests. - async function runDrainFixture(mode: "idle" | "inflight" | "force" | "closeIdle") { + // The drain promise must not resolve while a connection is still open, and + // a graceful stop() must actively drain: idle keep-alive connections close + // right away, busy ones close as soon as their in-flight work completes. + // The client never hangs up first (except in "partial" mode), so every + // observed close below is server-initiated; idleTimeout is long enough that + // a timeout-driven close would flake the runtime budget long before firing. + async function runDrainFixture(mode: "idle" | "inflight" | "inflightHead" | "force" | "partial") { await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -615,6 +617,7 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const server = Bun.serve({ port: 0, hostname: "127.0.0.1", + idleTimeout: 255, async fetch(req) { if (new URL(req.url).pathname === "/slow") { inflight.resolve(); @@ -634,26 +637,51 @@ describe.concurrent("server.stop() drain promise counts open connections", () => c.on("connect", resolve); c.on("error", reject); }); - c.write("GET /" + (mode === "inflight" ? "slow" : "fast") + " HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); - if (mode === "inflight") { - await inflight.promise; + if (mode === "partial") { + // Half a request head: enough to be accepted and parsed as an + // in-progress request (a zero-byte connection would sit in the + // TCP_DEFER_ACCEPT queue on Linux and never reach the server), + // never completed. + c.write("GET /slow HTTP/1.1\\r\\nHost: x\\r\\n"); + // Wait until the server has accepted it (it shows up in + // pendingRequests only once parsed, so poll a tick batch). + for (let i = 0; i < 20; i++) await new Promise(r => setImmediate(r)); } else { - while (!buf.includes("\\r\\nok")) await new Promise(r => setImmediate(r)); + const method = mode === "inflightHead" ? "HEAD" : "GET"; + c.write(method + " /" + (mode === "idle" ? "fast" : "slow") + " HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + if (mode === "idle") { + while (!buf.includes("\\r\\nok")) await new Promise(r => setImmediate(r)); + } else { + await inflight.promise; + } } - // One keep-alive connection exists (idle or with a request in flight). let resolved = false; const stopped = server.stop(false).then(() => { resolved = true; }); await new Promise(r => setImmediate(r)); const resolvedEarly = resolved; - if (mode === "inflight") release.resolve(); - // Connection is (now) idle; promise must still be pending. - while (!buf.includes("\\r\\nok")) await new Promise(r => setImmediate(r)); - await new Promise(r => setImmediate(r)); - const resolvedWhileOpen = resolved; - // Drain it: the client hangs up or the caller escalates or sweeps. - if (mode === "force") server.stop(true); - else if (mode === "closeIdle") server.closeIdleConnections(); - else c.destroy(); + let responseAfterStop = false; + if (mode === "inflight" || mode === "inflightHead") { + // The held request completes; its full response must reach the + // client before the server closes the now-idle connection. A HEAD + // response has no body, so its completion goes through the + // no-body end path rather than internalEnd. + release.resolve(); + const doneMark = mode === "inflightHead" ? "\\r\\n\\r\\n" : "\\r\\nok"; + while (!buf.includes(doneMark) && !events.includes("close")) await new Promise(r => setImmediate(r)); + responseAfterStop = buf.includes(doneMark); + } else if (mode === "force") { + // Escalation cuts the still-held request; release the handler so + // the aborted request can settle. + server.stop(true); + release.resolve(); + } else if (mode === "partial") { + // A connection mid-request ("sending a request", Node's idle + // definition excludes it) is spared; the promise stays pending + // until the client hangs up. + for (let i = 0; i < 20; i++) await new Promise(r => setImmediate(r)); + if (resolved) throw new Error("stop() resolved while a mid-request connection was open"); + c.destroy(); + } await stopped; // The client socket's 'close' and the server-side filter → promise // resolution race; poll so the assertion is order-independent. @@ -663,7 +691,7 @@ describe.concurrent("server.stop() drain promise counts open connections", () => await new Promise(r => setImmediate(r)); closed = events.includes("close"); } - console.log(JSON.stringify({ resolvedEarly, resolvedWhileOpen, resolved, closed })); + console.log(JSON.stringify({ resolvedEarly, responseAfterStop, resolved, closed })); c.destroy(); `, ], @@ -675,37 +703,354 @@ describe.concurrent("server.stop() drain promise counts open connections", () => return { stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }; } - test("idle keep-alive connection holds the promise until the client closes", async () => { + test("stop() closes an idle keep-alive connection and the promise resolves", async () => { + // The sweep closes the socket inside stop() itself, so the promise may + // already be resolved one tick later; resolvedEarly is not meaningful. expect(await runDrainFixture("idle")).toEqual({ stderr: "", - out: { resolvedEarly: false, resolvedWhileOpen: false, resolved: true, closed: true }, + out: { resolvedEarly: expect.any(Boolean), responseAfterStop: false, resolved: true, closed: true }, exitCode: 0, }); }); - test("in-flight request's connection holds the promise past response end", async () => { + test("in-flight request completes across stop(), then its connection is closed", async () => { expect(await runDrainFixture("inflight")).toEqual({ stderr: "", - out: { resolvedEarly: false, resolvedWhileOpen: false, resolved: true, closed: true }, + out: { resolvedEarly: false, responseAfterStop: true, resolved: true, closed: true }, + exitCode: 0, + }); + }); + + test("in-flight HEAD request completes across stop(), then its connection is closed", async () => { + expect(await runDrainFixture("inflightHead")).toEqual({ + stderr: "", + out: { resolvedEarly: false, responseAfterStop: true, resolved: true, closed: true }, + exitCode: 0, + }); + }); + + test("a handler rejection on a HEAD request still closes its drained connection", async () => { + // The production 500 for a rejected handler renders from the rejection + // microtask, uncorked, and a HEAD response ends without a body: no cork + // or parser gate runs, so RequestContext::end_without_body has to run the + // close gate itself for the stop() mark to take effect. + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const net = require("net"); + const inflight = Promise.withResolvers(); + const release = Promise.withResolvers(); + const server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + idleTimeout: 255, + development: false, + async fetch() { + inflight.resolve(); + await release.promise; + throw new Error("boom"); + }, + }); + const c = net.connect(server.port, "127.0.0.1"); + let buf = ""; + const closed = Promise.withResolvers(); + c.on("data", d => (buf += d)); + c.on("close", () => closed.resolve()); + c.on("error", () => {}); + await new Promise((resolve, reject) => { c.on("connect", resolve); c.on("error", reject); }); + c.write("HEAD /slow HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + await inflight.promise; + let resolved = false; + const stopped = server.stop(false).then(() => { resolved = true; }); + await new Promise(r => setImmediate(r)); + const resolvedEarly = resolved; + // The handler rejects; the 500 head renders from the microtask and + // the drained connection must close on its own. + release.resolve(); + await stopped; + await closed.promise; + console.log(JSON.stringify({ resolvedEarly, got500: buf.includes(" 500 "), resolved })); + // The rejected handler marks the process exit code; the drain + // assertions above are what this test is about. + process.exit(0); + `, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + // The rejected handler is logged to stderr by design; drain it but assert + // only the drain behavior. + const [stdout, , exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ + out: { resolvedEarly: false, got500: true, resolved: true }, exitCode: 0, }); }); - test("stop(true) after stop(false) force-closes the surviving connection", async () => { + test("stop(true) after stop(false) force-closes the still-busy connection", async () => { expect(await runDrainFixture("force")).toEqual({ stderr: "", - out: { resolvedEarly: false, resolvedWhileOpen: false, resolved: true, closed: true }, + out: { resolvedEarly: false, responseAfterStop: false, resolved: true, closed: true }, + exitCode: 0, + }); + }); + + test("a connection mid-request survives stop() until the client closes", async () => { + expect(await runDrainFixture("partial")).toEqual({ + stderr: "", + out: { resolvedEarly: false, responseAfterStop: false, resolved: true, closed: true }, + exitCode: 0, + }); + }); + + test("in-flight Bun.file (sendfile) response completes across stop(), then its connection is closed", async () => { + // The sendfile completion path (uws_res_end_sendfile) bypasses internalEnd + // and returns `false` to uWS's onWritable, so none of the parser-side + // shouldCloseConnection() gates run; the explicit gate after the stream's + // on_complete is what closes the drained connection here. + const dir = tempDirWithFiles("drain-sendfile", {}); + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const net = require("net"); + const SIZE = 16 * 1024 * 1024; + await Bun.write("big.bin", Buffer.alloc(SIZE, "x")); + const server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + idleTimeout: 255, + fetch: () => new Response(Bun.file("big.bin")), + }); + const c = net.connect(server.port, "127.0.0.1"); + // Parse the response framing so the assertion is on body bytes, not + // raw socket bytes (headers must not mask a truncated body). + let head = ""; + let headDone = false; + let contentLength = -1; + let bodyBytes = 0; + let sawClose = false; + const firstData = Promise.withResolvers(); + c.on("data", d => { + if (!headDone) { + head += d.toString("latin1"); + const he = head.indexOf("\\r\\n\\r\\n"); + if (he !== -1) { + headDone = true; + contentLength = +(/\\r\\ncontent-length: *(\\d+)/i.exec(head.slice(0, he))?.[1] ?? -1); + bodyBytes = head.length - he - 4; + } + } else { + bodyBytes += d.length; + } + firstData.resolve(); + }); + c.on("close", () => (sawClose = true)); + c.on("error", () => {}); + await new Promise((resolve, reject) => { c.on("connect", resolve); c.on("error", reject); }); + c.write("GET / HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + // First bytes of the response have arrived; 16 MB cannot fit in the + // socket buffers, so the transfer is still in flight server-side. + await firstData.promise; + c.pause(); + let resolved = false; + const stopped = server.stop(false).then(() => { resolved = true; }); + await new Promise(r => setImmediate(r)); + const resolvedEarly = resolved; + c.resume(); + // The full body must arrive, then the server closes the connection + // (the client never hangs up; idleTimeout is far beyond the test + // budget, so only the drain can close it). + while (!sawClose) await new Promise(r => setImmediate(r)); + await stopped; + console.log(JSON.stringify({ resolvedEarly, contentLength, bodyBytes, resolved })); + `, + ], + env: bunEnv, + cwd: dir, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ + stderr: "", + out: { resolvedEarly: false, contentLength: 16 * 1024 * 1024, bodyBytes: 16 * 1024 * 1024, resolved: true }, + exitCode: 0, + }); + }); + + test("a response completing inside another socket's parse window still closes its drained connection", async () => { + // internalEnd's post-uncork close gate must key on WHICH socket the + // parser is on, not the context-wide isParsingHttp bit: B's parked + // response below completes in the microtask drain inside A's onData + // dispatch (A resolves it), and B gets no later gate of its own. With the + // context-wide bit, B lingered until idleTimeout and stop() hung. + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const net = require("net"); + const releaseB = Promise.withResolvers(); + const bHeld = Promise.withResolvers(); + const server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + idleTimeout: 255, + async fetch(req) { + const path = new URL(req.url).pathname; + if (path === "/hold") { + bHeld.resolve(); + await releaseB.promise; + return new Response("held"); + } + if (path === "/poke") { + // B's completion runs in the microtask drain while A's + // socket is the one being parsed. + releaseB.resolve(); + return new Response("poked"); + } + return new Response("ok"); + }, + }); + function dial() { + const c = net.connect(server.port, "127.0.0.1"); + const state = { c, buf: "", closed: false }; + c.on("data", d => (state.buf += d)); + c.on("close", () => (state.closed = true)); + c.on("error", () => {}); + return new Promise((res, rej) => { + c.on("connect", () => res(state)); + c.on("error", rej); + }); + } + const b = await dial(); + b.c.write("GET /hold HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + await bHeld.promise; + const a = await dial(); + // Partial head: A is mid-request at the sweep, so it is spared and + // marked close-when-idle. Give the bytes a tick batch to arrive (a + // zero-byte connection would sit in the defer-accept queue and the + // sweep could not see it). + a.c.write("GET /poke HTTP/1.1\\r\\nHost: x\\r\\n"); + for (let i = 0; i < 20; i++) await new Promise(r => setImmediate(r)); + let resolved = false; + const stopped = server.stop(false).then(() => { resolved = true; }); + await new Promise(r => setImmediate(r)); + const resolvedEarly = resolved; + // Complete A's request; its dispatch resolves B inside A's parse + // window. Both responses must be delivered, then both connections + // close server-initiated and the drain promise resolves. + a.c.write("\\r\\n"); + while (!a.closed || !b.closed) await new Promise(r => setImmediate(r)); + await stopped; + console.log(JSON.stringify({ + resolvedEarly, + aGotResponse: a.buf.includes("poked"), + bGotResponse: b.buf.includes("held"), + resolved, + })); + `, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ + stderr: "", + out: { resolvedEarly: false, aGotResponse: true, bGotResponse: true, resolved: true }, exitCode: 0, }); }); - test("closeIdleConnections() after stop(false) drains the surviving connection", async () => { - // close_idle_connections suppresses the filter's own dispatch while the - // uWS call runs (re-entrance guard); its trailing deinit_if_we_can() is - // what resolves the promise on this path. - expect(await runDrainFixture("closeIdle")).toEqual({ + test("closeIdleConnections() is a one-shot sweep that spares busy connections and keeps listening", async () => { + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + ` + const net = require("net"); + const inflight = Promise.withResolvers(); + const release = Promise.withResolvers(); + const server = Bun.serve({ + port: 0, + hostname: "127.0.0.1", + idleTimeout: 255, + async fetch(req) { + if (new URL(req.url).pathname === "/slow") { + inflight.resolve(); + await release.promise; + } + return new Response("ok"); + }, + }); + const port = server.port; + function dial() { + const c = net.connect(port, "127.0.0.1"); + const state = { c, buf: "", closed: false }; + c.on("data", d => (state.buf += d)); + c.on("close", () => (state.closed = true)); + c.on("error", () => {}); + return new Promise((resolve, reject) => { + c.on("connect", () => resolve(state)); + c.on("error", reject); + }); + } + const countOks = s => (s.buf.match(/\\r\\nok/g) || []).length; + const idle = await dial(); + idle.c.write("GET /fast HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + while (countOks(idle) < 1) await new Promise(r => setImmediate(r)); + const busy = await dial(); + busy.c.write("GET /slow HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + await inflight.promise; + + const closedFirst = server.closeIdleConnections(); + // Idle connection closes; the busy one is spared. + while (!idle.closed) await new Promise(r => setImmediate(r)); + const busyClosedBySweep = busy.closed; + release.resolve(); + while (countOks(busy) < 1) await new Promise(r => setImmediate(r)); + // One-shot: the spared connection was not marked close-when-idle, so + // it keeps serving keep-alive requests after its response completed. + busy.c.write("GET /fast HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + while (countOks(busy) < 2) await new Promise(r => setImmediate(r)); + // The listener is untouched: a fresh connection still gets served. + const fresh = await dial(); + fresh.c.write("GET /fast HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); + while (countOks(fresh) < 1) await new Promise(r => setImmediate(r)); + // Both surviving connections are idle keep-alive now; a second + // sweep closes them both and reports the count. + const closedSecond = server.closeIdleConnections(); + while (!busy.closed || !fresh.closed) await new Promise(r => setImmediate(r)); + console.log(JSON.stringify({ + closedFirst, + busyClosedBySweep, + busyServedAfterSweep: countOks(busy) === 2, + freshServed: countOks(fresh) === 1, + closedSecond, + })); + server.stop(true); + `, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ stderr: "", - out: { resolvedEarly: false, resolvedWhileOpen: false, resolved: true, closed: true }, + out: { + closedFirst: 1, + busyClosedBySweep: false, + busyServedAfterSweep: true, + freshServed: true, + closedSecond: 2, + }, exitCode: 0, }); }); @@ -761,8 +1106,10 @@ describe.concurrent("server.stop() drain promise counts open connections", () => test("pre-handshake TLS close does not steal another connection's count", async () => { // For TLS, +1 fires in onHandshake, -1 in onClose. A socket that RSTs // before the handshake reaches onClose without a matching +1; without the - // per-socket filteredOpen gate that -1 would decrement the count for the - // live handshaken connection and stop(false) would resolve under it. + // per-socket filteredOpen gate that -1 would steal the live handshaken + // connection's count and leave it stuck after that connection closes, so + // stop(false) would never resolve. The live connection holds a request + // across stop() so the sweep cannot close it before the raw closes land. await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -771,22 +1118,30 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const net = require("net"); const tls = require("tls"); const { tls: serverTls } = require(${JSON.stringify(require.resolve("harness"))}); + const inflight = Promise.withResolvers(); + const release = Promise.withResolvers(); const server = Bun.serve({ port: 0, hostname: "127.0.0.1", tls: serverTls, - fetch: () => new Response("ok"), + async fetch() { + inflight.resolve(); + await release.promise; + return new Response("ok"); + }, }); const port = server.port; // One real TLS keep-alive connection: handshake completes, count=1. const c = tls.connect({ port, host: "127.0.0.1", ca: serverTls.cert, rejectUnauthorized: false }); let buf = ""; + const closed = Promise.withResolvers(); c.on("data", d => (buf += d)); + c.on("close", () => closed.resolve()); c.on("error", () => {}); await new Promise((resolve, reject) => { c.on("secureConnect", resolve); c.on("error", reject); }); c.write("GET / HTTP/1.1\\r\\nHost: x\\r\\n\\r\\n"); - while (!buf.includes("\\r\\nok")) await new Promise(r => setImmediate(r)); + await inflight.promise; // Three raw TCP connects that close before the handshake. onClose // fires for each; without filteredOpen, each -1 would steal c's // count (and the rest would be swallowed by the prev==0 guard). @@ -808,9 +1163,16 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const stopped = server.stop(false).then(() => { resolved = true; }); await new Promise(r => setImmediate(r)); const resolvedEarly = resolved; - c.destroy(); - await Promise.race([stopped, new Promise(r => setTimeout(r, 2000))]); - console.log(JSON.stringify({ resolvedEarly, resolved })); + // The held response completes, reaches the client, and the drain + // closes the connection (server-initiated; the client never hangs + // up); the promise must resolve on its own. The runner timeout is + // the stall bound for both awaits. + release.resolve(); + await stopped; + // Reaching the log below proves the server-initiated close arrived. + await closed.promise; + const gotResponse = buf.includes("\\r\\nok"); + console.log(JSON.stringify({ resolvedEarly, resolved, gotResponse })); process.exit(0); `, ], @@ -821,14 +1183,17 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ stderr: "", - out: { resolvedEarly: false, resolved: true }, + out: { resolvedEarly: false, resolved: true, gotResponse: true }, exitCode: 0, }); }); test("server.reload() keeps the connection count coherent", async () => { // clearRoutes() used to wipe filterHandlers, so a connection open across - // reload left active_connection_count stuck > 0 forever. + // reload left active_connection_count stuck > 0 forever. With the stuck + // count (or a wiped filter, whose close would no longer decrement), the + // sweep in stop(false) closes the idle connection but the drain promise + // never resolves. await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -837,12 +1202,15 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const net = require("net"); const server = Bun.serve({ port: 0, hostname: "127.0.0.1", + idleTimeout: 255, fetch: () => new Response("ok"), }); const port = server.port; const c = net.connect(port, "127.0.0.1"); let buf = ""; + const closed = Promise.withResolvers(); c.on("data", d => (buf += d)); + c.on("close", () => closed.resolve()); c.on("error", () => {}); await new Promise((resolve, reject) => { c.on("connect", resolve); @@ -855,11 +1223,12 @@ describe.concurrent("server.stop() drain promise counts open connections", () => server.reload({ fetch: () => new Response("ok") }); let resolved = false; const stopped = server.stop(false).then(() => { resolved = true; }); - await new Promise(r => setImmediate(r)); - const resolvedEarly = resolved; - c.destroy(); - await Promise.race([stopped, new Promise(r => setTimeout(r, 2000))]); - console.log(JSON.stringify({ resolvedEarly, resolved })); + // stop() closes the idle connection itself; the client never hangs + // up, so both awaits are server-driven and the runner timeout is + // the stall bound (a stuck count would hang right here). + await stopped; + await closed.promise; + console.log(JSON.stringify({ resolved })); process.exit(0); `, ], @@ -870,7 +1239,7 @@ describe.concurrent("server.stop() drain promise counts open connections", () => const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stderr, out: JSON.parse(stdout.trim() || "null"), exitCode }).toEqual({ stderr: "", - out: { resolvedEarly: false, resolved: true }, + out: { resolved: true }, exitCode: 0, }); }); @@ -912,24 +1281,26 @@ test("should be able to await server.stop(true) with keep alive", async () => { expect(async () => await fetch(server.url)).toThrow(); }); -// Shared rig for the two "late keep-alive" tests below: open a raw TCP -// socket, hold the first request in-flight across stop()/close(), pipeline a -// second request behind it, release, GC, and print the second response's -// status line. The subprocess runs the rig so a (former) panic in the dispatch -// trampoline surfaces as a non-zero exit instead of taking down the runner. +// Shared rig for the "late keep-alive" tests below: open a raw TCP socket, +// hold the first request in-flight across stop()/close(), pipeline a second +// request behind it, release, GC, and print the second response's status +// line (or "" if the connection closed instead). The subprocess runs the rig +// so a (former) panic in the dispatch trampoline surfaces as a non-zero exit +// instead of taking down the runner. // // `deinit_if_we_can` defers the wrapper downgrade while the connection is -// still open, so both requests dispatch against a live wrapper. The pipelined -// request carries `Connection: close`, so once it completes the connection -// closes, the wrapper downgrades to Weak, and the GC pass below must collect -// it cleanly. The `respond_stopped_503` guard in the trampolines is a safety -// net for the `Finalized` case (wrapper GC'd while `self` still lives between -// `finalize()` and the next-tick `schedule_deinit`); that window is not -// deterministically reachable from a test. +// still open, so the held request always dispatches and completes against a +// live wrapper. What happens to the pipelined request depends on the server: +// a Bun.serve graceful stop marks the busy connection close-when-idle, so the +// connection closes right after the held response and the pipelined request +// is dropped (secondOutcome: "closed"); node:http queues pipelined responses, +// so close() still delivers it before the connection closes +// (secondOutcome: "200"). Either way the wrapper then downgrades to Weak and +// the GC pass below must collect it cleanly. // // `serverSnippet` must define `port` (the listen port) and `stop()` in scope, // and may read `release`/`inflight`/`hits` for the hold protocol. -async function runLateKeepAlive(reqPath: string, serverSnippet: string) { +async function runLateKeepAlive(reqPath: string, serverSnippet: string, secondOutcome: "200" | "closed") { await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -994,10 +1365,12 @@ async function runLateKeepAlive(reqPath: string, serverSnippet: string) { })(); // The only server binding is now out of scope. - // First request completes; the connection is still open so the - // wrapper stays Strong, and the pipelined request dispatches against - // it → 200. Previously: panic (or 503 when the gate checked - // Strong-only). + // First request completes against a live wrapper (the connection is + // open, so it stays Strong). What happens to the pipelined request + // depends on the caller: Bun.serve's drain closes the connection at + // idle and drops it (second = ""), node:http delivers the queued + // response (second = 200). Previously the late dispatch could panic + // (or 503 when the gate checked Strong-only). release.resolve(); const first = await nextResponse(); if (!first.includes("200")) throw new Error("first request failed: " + first); @@ -1021,18 +1394,23 @@ async function runLateKeepAlive(reqPath: string, serverSnippet: string) { }); const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - // Must actually dispatch — empty would mean the socket was closed before the - // pipelined request reached the trampoline. expect({ stdout: stdout.trim(), stderr, exitCode }).toEqual({ - stdout: expect.stringMatching(/^HTTP\/1\.1 200\b/), + // "200": the pipelined request must actually dispatch and be answered. + // "closed": the connection must close cleanly after the held response + // without the pipelined request being answered (and without a panic). + stdout: secondOutcome === "200" ? expect.stringMatching(/^HTTP\/1\.1 200\b/) : "", stderr: "", exitCode: 0, }); } -test("late keep-alive request to a route after stop() still dispatches", async () => { - // Per-route handlers live in ServerRouteList, which is reachable from JS only - // through the Server wrapper — exercises on_user_route_request's gate. +test("stop() completes the in-flight request, then closes the connection instead of serving a late pipelined request", async () => { + // The route handler is held across stop(), so the connection is busy during + // the sweep and gets marked close-when-idle; the held response must still be + // delivered in full, after which the connection closes and the pipelined + // request behind it is dropped (a pipelining client retries it elsewhere, + // RFC 9112 9.3.2). Must not panic: the close path runs against a server + // whose only JS binding went out of scope before the drain. await runLateKeepAlive( "/r", ` @@ -1052,17 +1430,17 @@ test("late keep-alive request to a route after stop() still dispatches", async ( const port = server.port; const stop = () => server.stop(); `, + "closed", ); }); -test("late keep-alive WebSocket upgrade after stop() succeeds while the connection is still open", async () => { +test("stop() closes the drained connection before a late pipelined WebSocket upgrade dispatches", async () => { // Sibling of the HTTP late-keep-alive test for the WebSocket upgrade path. - // `deinit_if_we_can` defers the wrapper downgrade (and the - // `handler.server`/`handler.app` clear) while the keep-alive connection is - // still open, so a pipelined upgrade on that connection reaches a live - // handler and `server.upgrade()` succeeds. On upgrade the socket's count - // moves from the HTTP tally to the WebSocket tally, so the server still - // drains cleanly once the WebSocket closes. + // The connection is busy during stop()'s sweep (held response), so it is + // marked close-when-idle: the held response is delivered, then the + // connection closes, and the upgrade request pipelined behind it never + // reaches the fetch handler (the client reconnects elsewhere). Must exit + // cleanly: the close runs on a gracefully stopped server. await using proc = Bun.spawn({ cmd: [ bunExe(), @@ -1105,7 +1483,8 @@ test("late keep-alive WebSocket upgrade after stop() succeeds while the connecti sock.write("GET / HTTP/1.1\\r\\nHost: x\\r\\nUpgrade: websocket\\r\\nConnection: Upgrade\\r\\nSec-WebSocket-Key: " + key + "\\r\\nSec-WebSocket-Version: 13\\r\\n\\r\\n"); server.stop(); release.resolve(); - // Wait for both responses. + // Wait for the held response and the server-initiated close (the + // pipelined upgrade is dropped by the drain, so no second response). while (!sockClosed && (received.match(/\\r\\n\\r\\n/g) || []).length < 2) { await waiter.promise; waiter = Promise.withResolvers(); } @@ -1124,16 +1503,18 @@ test("late keep-alive WebSocket upgrade after stop() succeeds while the connecti const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); const out = JSON.parse(stdout.trim() || "{}"); expect({ stderr, exitCode }).toEqual({ stderr: "", exitCode: 0 }); - // First request 200 (held handler), pipelined upgrade accepted → 101. - expect(out.upgraded).toBe(true); - expect(out.statuses?.[0]).toMatch(/^HTTP\/1\.1 200\b/); - expect(out.statuses?.[1]).toMatch(/^HTTP\/1\.1 101\b/); + // Held request 200; the pipelined upgrade never dispatches (fetch is not + // called again, so `upgraded` is never assigned) and no 101 is written. + expect(out.upgraded).toBeUndefined(); + expect(out.statuses).toEqual([expect.stringMatching(/^HTTP\/1\.1 200\b/)]); }); test("late keep-alive request to a node:http server after close() still dispatches", async () => { // Same shape but through node:http so the request dispatches via // on_node_http_request_with_upgrade_ctx — the trampoline that would panic - // on a stale shadow without the `js_value_for_dispatch` gate. + // on a stale shadow without the `js_value_for_dispatch` gate. Unlike + // Bun.serve, node:http queues pipelined requests, so the queued response is + // still delivered after close() before the drain closes the connection. await runLateKeepAlive( "/", ` @@ -1152,6 +1533,7 @@ test("late keep-alive request to a node:http server after close() still dispatch // Also drops node:http's own reference to the Bun server. const stop = () => srv.close(); `, + "200", ); });