diff --git a/src/js/node/_http_server.ts b/src/js/node/_http_server.ts index 3c7f33a417ed..d543cf3fe149 100644 --- a/src/js/node/_http_server.ts +++ b/src/js/node/_http_server.ts @@ -2446,10 +2446,17 @@ function stopServerResponsePerf(this: any) { // arm keep-alive) runs first because onResponseFinishHandleSocket's guards // read pre-detach state, then detach the socket and advance the pipeline. function emitResponseFinish() { + // If the user never called req.read(), and didn't pipe() or + // .resume() or .on('data'), then we call req._dump() so that the + // bytes will be pulled off the wire. + const req = this.req; + if (req && !req._consuming && !req._readableState?.resumeScheduled) { + req._dump(); + } // req.socket is nulled by the stream destroyer (pipeline/compose cleanup); // the response's own socket (set by assignSocket, cleared only by // detachSocket) still references the connection then. - const socket = this.req?.socket ?? this.socket; + const socket = req?.socket ?? this.socket; onResponseFinishHandleSocket(socket?.server, socket, this); // The dispatcher detached a synchronously-finished response itself; // advancing the pipeline again here would skip a queued response. @@ -3218,10 +3225,6 @@ ServerResponse.prototype.end = function (chunk, encoding, callback) { } } this._header = " "; - const req = this.req; - if (!req._consuming && !req?._readableState?.resumeScheduled) { - req._dump(); - } // The socket is NOT detached here: like Node.js, res.socket stays assigned // until the response 'finish' machinery runs (the dispatcher detaches it // right after a synchronously-finished handler returns, or via its 'finish' diff --git a/src/runtime/server/NodeHTTPResponse.rs b/src/runtime/server/NodeHTTPResponse.rs index 86e0db34a795..2582d2220e59 100644 --- a/src/runtime/server/NodeHTTPResponse.rs +++ b/src/runtime/server/NodeHTTPResponse.rs @@ -610,7 +610,7 @@ impl NodeHTTPResponse { || flags.contains(Flags::ENDED)) && (self.body_read_ref.get().has || self.body_read_state.get() == BodyReadState::Pending) - && (!flags.contains(Flags::HAS_CUSTOM_ON_DATA) + && (this_value.is_empty_or_undefined_or_null() || js::on_data_get_cached(this_value).is_none()) { let had_ref = self.body_read_ref.get().has; @@ -632,6 +632,17 @@ impl NodeHTTPResponse { fn should_request_be_pending(&self) -> bool { let flags = self.flags.get(); + // The body's fin was parsed onto the paused shim and the tail is still + // buffered: mark_request_as_done() would free it before the consumer + // drains, so stay pending regardless of which terminal flag is set. + if flags.contains(Flags::IS_DATA_BUFFERED_DURING_PAUSE_LAST) + && !self + .buffered_request_body_data_during_pause + .get() + .is_empty() + { + return true; + } // Once the socket is closed or has been adopted by the WebSocket // layer, the HTTP request/response cycle is over — no further uws // callbacks will arrive on `raw_response` to balance the @@ -644,12 +655,12 @@ impl NodeHTTPResponse { // A raw 'upgrade'/'connect' tunnel handoff ends the HTTP exchange the // same way, except an Upgrade carrying a body keeps parsing as HTTP // until the body's fin chunk (the actual tunnel start). - if flags.contains(Flags::TUNNELED) { - return self.body_read_state.get() == BodyReadState::Pending; - } - - if flags.contains(Flags::ENDED) { - return self.body_read_state.get() == BodyReadState::Pending; + if flags.contains(Flags::TUNNELED) || flags.contains(Flags::ENDED) { + return self.body_read_state.get() == BodyReadState::Pending + || !self + .buffered_request_body_data_during_pause + .get() + .is_empty(); } true @@ -701,6 +712,7 @@ impl NodeHTTPResponse { let vm = vm_get(); self.clear_on_data_callback(self.get_this_value(), vm.global()); + self.body_read_ref.with_mut(|r| r.unref(vm)); self.clear_pending_pinned_write(vm.global(), JSValue::ZERO); self.upgrade_context.with_mut(|c| c.reset()); @@ -1347,10 +1359,11 @@ impl NodeHTTPResponse { let Some(raw) = self.raw_response.get() else { return Ok(JSValue::FALSE); }; - if flags.contains(Flags::REQUEST_HAS_COMPLETED) - || flags.contains(Flags::SOCKET_CLOSED) - || flags.contains(Flags::ENDED) - || flags.contains(Flags::UPGRADED) + if flags.contains(Flags::SOCKET_CLOSED) || flags.contains(Flags::UPGRADED) { + return Ok(JSValue::FALSE); + } + if (flags.contains(Flags::REQUEST_HAS_COMPLETED) || flags.contains(Flags::ENDED)) + && self.body_read_state.get() != BodyReadState::Pending { return Ok(JSValue::FALSE); } @@ -1400,10 +1413,15 @@ impl NodeHTTPResponse { let bytes = self .buffered_request_body_data_during_pause .replace(Vec::new()); - return Some(JSValue::create_buffer_from_box( - global_object, - bytes.into_boxed_slice(), - )); + let buf = JSValue::create_buffer_from_box(global_object, bytes.into_boxed_slice()); + if self + .flags + .get() + .contains(Flags::IS_DATA_BUFFERED_DURING_PAUSE_LAST) + { + self.mark_request_as_done_if_necessary(); + } + return Some(buf); } None } @@ -1419,11 +1437,11 @@ impl NodeHTTPResponse { let Some(raw) = self.raw_response.get() else { return JSValue::FALSE; }; - if flags.contains(Flags::REQUEST_HAS_COMPLETED) - || flags.contains(Flags::SOCKET_CLOSED) - || flags.contains(Flags::ENDED) - || flags.contains(Flags::UPGRADED) - { + if flags.contains(Flags::SOCKET_CLOSED) || flags.contains(Flags::UPGRADED) { + return JSValue::FALSE; + } + let ended = flags.contains(Flags::REQUEST_HAS_COMPLETED) || flags.contains(Flags::ENDED); + if ended && self.body_read_state.get() != BodyReadState::Pending { return JSValue::FALSE; } // Body already delivered: re-arming onData/onTimeout would overwrite a @@ -1435,6 +1453,13 @@ impl NodeHTTPResponse { self.set_on_aborted_handler(); raw.on_data(on_data_shim, self.as_ctx_ptr()); } + // After the response has ended detachSocket() cleared + // socket._httpMessage, so #resumeSocket() can no longer route the + // drain to the IncomingMessage; leave it for _read()'s + // handle.drainRequestBody() which pushes to `this` directly. + if ended { + return JSValue::TRUE; + } self.update_flags(|f| f.remove(Flags::IS_DATA_BUFFERED_DURING_PAUSE)); let mut result: JSValue = JSValue::TRUE; @@ -1599,10 +1624,11 @@ impl NodeHTTPResponse { if last { self.capture_request_trailers(); self.update_flags(|f| f.insert(Flags::IS_DATA_BUFFERED_DURING_PAUSE_LAST)); + self.body_read_state.set(BodyReadState::Done); if self.body_read_ref.get().has { self.body_read_ref.with_mut(|r| r.unref(vm_get())); - self.mark_request_as_done_if_necessary(); } + self.mark_request_as_done_if_necessary(); } } @@ -1688,8 +1714,8 @@ impl NodeHTTPResponse { if last { if self.body_read_ref.get().has { self.body_read_ref.with_mut(|r| r.unref(vm_get())); - self.mark_request_as_done_if_necessary(); } + self.mark_request_as_done_if_necessary(); self.deref(); } } @@ -2013,12 +2039,11 @@ impl NodeHTTPResponse { self.spill_pending_pinned_write(global_object); if IS_END { - // Discard the body read ref if it's pending and no onData callback is set at this point. - // This is the equivalent of req._dump(). + // Connection: close / HTTP/1.0: end(close=true) shuts the socket + // down before any later body segment could reach the parser. if self.body_read_ref.get().has && self.body_read_state.get() == BodyReadState::Pending - && (!self.flags.get().contains(Flags::HAS_CUSTOM_ON_DATA) - || js::on_data_get_cached(this_value).is_none()) + && state.is_http_connection_close() { self.body_read_ref.with_mut(|r| r.unref(vm_get())); self.body_read_state.set(BodyReadState::None); @@ -2039,6 +2064,26 @@ impl NodeHTTPResponse { } else { raw_response.end_stream(state.is_http_connection_close()); } + // markDone() nulled inStream; re-arm so the request body keeps + // flowing into the IncomingMessage after the response is sent. + // SOCKET_CLOSED means the HttpResponseData ext was just + // destructed by the context onClose, so it cannot be touched. + let post_end_flags = self.flags.get(); + if self.body_read_state.get() == BodyReadState::Pending + && !post_end_flags.contains(Flags::SOCKET_CLOSED) + { + if let Some(raw_response) = self.raw_response.get() { + if post_end_flags.contains(Flags::IS_DATA_BUFFERED_DURING_PAUSE) + && !post_end_flags.contains(Flags::IS_DATA_BUFFERED_DURING_PAUSE_LAST) + { + raw_response.on_data(on_buffer_paused_shim, self.as_ctx_ptr()); + #[cfg(not(windows))] + self.pause_socket(); + } else { + raw_response.on_data(on_data_shim, self.as_ctx_ptr()); + } + } + } self.on_request_complete(); Ok(JSValue::js_number_from_uint64(bytes.len() as u64)) @@ -2212,11 +2257,17 @@ impl NodeHTTPResponse { fn clear_on_data_callback(&self, this_value: JSValue, global_object: &JSGlobalObject) { scoped_log!(NodeHTTPResponse, "clearOnDataCallback"); if self.body_read_state.get() != BodyReadState::None { - if !this_value.is_empty() { - js::on_data_set_cached(this_value, global_object, JSValue::UNDEFINED); - } let flags = self.flags.get(); - if !flags.contains(Flags::SOCKET_CLOSED) && !flags.contains(Flags::UPGRADED) { + // Once REQUEST_HAS_COMPLETED is set the per-socket HttpResponseData + // (and get_this_value()'s currentResponseObject) may belong to the + // next keep-alive request; do not touch either. + if !flags.contains(Flags::SOCKET_CLOSED) + && !flags.contains(Flags::UPGRADED) + && !flags.contains(Flags::REQUEST_HAS_COMPLETED) + { + if !this_value.is_empty() { + js::on_data_set_cached(this_value, global_object, JSValue::UNDEFINED); + } scoped_log!(NodeHTTPResponse, "clearOnData"); if let Some(raw_response) = self.raw_response.get() { raw_response.clear_on_data(); @@ -2246,10 +2297,11 @@ impl NodeHTTPResponse { || flags.contains(Flags::UPGRADED) { js::on_data_set_cached(this_value, global_object, JSValue::UNDEFINED); + let was_pending = self.body_read_state.get() == BodyReadState::Pending; // defer { if body_read_ref.has { unref } } — moved to tail of this branch. match self.body_read_state.get() { BodyReadState::Pending | BodyReadState::Done => { - if !flags.contains(Flags::REQUEST_HAS_COMPLETED) + if (was_pending || !flags.contains(Flags::REQUEST_HAS_COMPLETED)) && !flags.contains(Flags::SOCKET_CLOSED) && !flags.contains(Flags::UPGRADED) { @@ -2266,6 +2318,9 @@ impl NodeHTTPResponse { self.body_read_ref .with_mut(|r| r.unref(bun_vm_mut(global_object))); } + self.buffered_request_body_data_during_pause + .with_mut(|b| b.clear_and_free()); + self.mark_request_as_done_if_necessary(); return; } diff --git a/test/js/node/http/node-http-proxy.js b/test/js/node/http/node-http-proxy.js index 8b82678ae9e4..51b9dca9ac26 100644 --- a/test/js/node/http/node-http-proxy.js +++ b/test/js/node/http/node-http-proxy.js @@ -37,7 +37,7 @@ export async function run() { const options = { protocol: "http:", - hostname: "localhost", + hostname: address.address, port: address.port, path: "/", // Change path to / headers: { diff --git a/test/js/node/http/node-http.test.ts b/test/js/node/http/node-http.test.ts index a204e6d37259..cc57a5860dfe 100644 --- a/test/js/node/http/node-http.test.ts +++ b/test/js/node/http/node-http.test.ts @@ -4054,3 +4054,285 @@ it("OutgoingMessage outputData is per-instance and _flushOutput is defined", () c.outputData.push({ data: "y", encoding: "utf8", callback: null }); expect(d.outputData.length).toBe(0); }); + +// https://github.com/oven-sh/bun/issues/4733 +// Ending the response inside the request handler must not drop the request +// body: like Node.js, the body keeps flowing to req's 'data' listeners and +// piped destinations until it has been fully received, and resOnFinish (the +// 'finish' listener) decides whether to _dump() based on the state *at* +// 'finish', so a consumer attached after res.end() in the same tick still +// receives the body. +describe("request body still flows after res.end() was called in the handler", () => { + async function run(handler: (req: IncomingMessage, res: ServerResponse, out: Writable) => void) { + const events: string[] = []; + const chunks: Buffer[] = []; + const { promise: finished, resolve: finish, reject } = Promise.withResolvers(); + + let reqRef: IncomingMessage | undefined; + await using server = createServer((req, res) => { + reqRef = req; + const out = new Writable({ + write(chunk, _enc, cb) { + chunks.push(Buffer.from(chunk)); + cb(); + }, + }); + req.once("end", () => events.push("req-end")); + req.once("close", () => events.push("req-close")); + req.once("error", reject); + out.once("error", reject); + out.once("finish", () => { + events.push("out-finish"); + finish(); + }); + handler(req, res, out); + }); + + await once(server.listen(0, "127.0.0.1"), "listening"); + const { port } = server.address() as AddressInfo; + const resp = await fetch(`http://127.0.0.1:${port}/`, { method: "POST", body: "testing-body" }); + expect(await resp.text()).toBe("ok"); + await finished; + + expect({ + body: Buffer.concat(chunks).toString(), + events, + dumped: reqRef!._dumped, + }).toEqual({ + body: "testing-body", + events: ["req-end", "out-finish", "req-close"], + dumped: false, + }); + } + + it("req.pipe(out) before res.end()", async () => { + await run((req, res, out) => { + req.pipe(out); + res.end("ok"); + }); + }); + + it("req.on('data') before res.end()", async () => { + await run((req, res, out) => { + req.on("data", c => out.write(c)); + req.on("end", () => out.end()); + res.end("ok"); + }); + }); + + it("req.pipe(out) after res.end() in the same tick", async () => { + await run((req, res, out) => { + res.end("ok"); + req.pipe(out); + }); + }); + + it("req.on('data') after res.end() in the same tick", async () => { + await run((req, res, out) => { + res.end("ok"); + req.on("data", c => out.write(c)); + req.on("end", () => out.end()); + }); + }); + + it("req.pipe(out) with res.write() + res.end()", async () => { + await run((req, res, out) => { + req.pipe(out); + res.write("o"); + res.end("k"); + }); + }); + + it("req.pipe(out) with res.end() on nextTick", async () => { + await run((req, res, out) => { + req.pipe(out); + process.nextTick(() => res.end("ok")); + }); + }); + + it("with no consumer, req is _dumped on 'finish' like Node's resOnFinish", async () => { + let dumpedAtFinish: boolean | undefined; + let reqRef: IncomingMessage | undefined; + const { promise: closed, resolve, reject } = Promise.withResolvers(); + await using server = createServer((req, res) => { + reqRef = req; + req.once("error", reject); + req.once("close", resolve); + res.end("ok"); + // emitResponseFinish is registered before the 'request' event, so by the + // time this listener runs req._dump() has already been called. + res.on("finish", () => (dumpedAtFinish = req._dumped)); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const { port } = server.address() as AddressInfo; + const resp = await fetch(`http://127.0.0.1:${port}/`, { method: "POST", body: "testing-body" }); + expect(await resp.text()).toBe("ok"); + await closed; + expect({ dumpedAtFinish, dumped: reqRef!._dumped }).toEqual({ dumpedAtFinish: true, dumped: true }); + }); + + it("req.resume() after res.end() in the same tick prevents _dump()", async () => { + let dumped: boolean | undefined; + const { promise: ended, resolve, reject } = Promise.withResolvers(); + await using server = createServer((req, res) => { + req.once("error", reject); + res.end("ok"); + req.resume(); + req.once("end", () => { + dumped = req._dumped; + resolve(); + }); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const { port } = server.address() as AddressInfo; + const resp = await fetch(`http://127.0.0.1:${port}/`, { method: "POST", body: "testing-body" }); + expect(await resp.text()).toBe("ok"); + await ended; + expect(dumped).toBe(false); + }); + + it("keep-alive connection reused after a consumed body releases the request", async () => { + // The body's fin arriving on the re-armed inStream after res.end() must + // release the pending-request ref so server.close() resolves and the + // process exits. + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + `const http = require("http"); + const { once } = require("events"); + (async () => { + const bodies = []; + const server = http.createServer((req, res) => { + let body = ""; + req.on("data", c => body += c); + req.once("end", () => bodies.push(body)); + res.end("ok"); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const agent = new http.Agent({ keepAlive: true, maxSockets: 1 }); + const port = server.address().port; + for (let i = 0; i < 3; i++) { + await new Promise((resolve, reject) => { + const r = http.request({ agent, method: "POST", port }, res => { + res.resume(); + res.on("end", resolve); + res.on("error", reject); + }); + r.on("error", reject); + r.end("body" + i); + }); + } + agent.destroy(); + await new Promise(r => server.close(r)); + console.log(JSON.stringify(bodies)); + })();`, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + expect(stdout.trim()).toBe('["body0","body1","body2"]'); + expect(exitCode).toBe(0); + }, 20_000); + + it("chunked request body split across writes", async () => { + let body = ""; + const { promise: ended, resolve, reject } = Promise.withResolvers(); + await using server = createServer((req, res) => { + req.on("data", c => (body += c)); + req.once("end", resolve); + req.once("error", reject); + req.once("close", () => reject(new Error("closed before end"))); + res.end("ok"); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const { port } = server.address() as AddressInfo; + + const sock = connect(port, "127.0.0.1"); + await once(sock, "connect"); + sock.write("POST / HTTP/1.1\r\nHost: x\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhello\r\n"); + // A second TCP segment carrying the remainder (via setNoDelay + event-loop + // bounce) so res.end() has already run when it arrives. + sock.setNoDelay(true); + await new Promise(r => setImmediate(r)); + sock.write("6\r\n world\r\n0\r\n\r\n"); + + await ended; + sock.end(); + expect(body).toBe("hello world"); + }); + + it("socket closed mid-upload does not strand the event loop", async () => { + // Covers the case where the body's fin never arrives after the response + // ended: the body-read ref must be released on teardown. + await using proc = Bun.spawn({ + cmd: [ + bunExe(), + "-e", + `const http = require("http"); + const net = require("net"); + const { once } = require("events"); + (async () => { + const events = []; + const server = http.createServer((req, res) => { + req.on("data", c => events.push("data(" + c.length + ")")); + req.on("end", () => events.push("end")); + res.end("ok"); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const s = net.connect(server.address().port, "127.0.0.1"); + await once(s, "connect"); + s.write("POST / HTTP/1.1\\r\\nHost: x\\r\\nContent-Length: 100\\r\\n\\r\\nabc"); + s.resume(); + await once(s, "data"); + s.destroy(); + await once(s, "close"); + server.close(); + process.on("beforeExit", () => console.log(JSON.stringify(events))); + })();`, + ], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect(stderr).toBe(""); + // Like Node.js: the partial body is delivered, 'end' is not (incomplete), + // and the process reaches beforeExit. + expect(stdout.trim()).toBe('["data(3)"]'); + expect(exitCode).toBe(0); + }, 20_000); + + it("Connection: close does not hang when the body consumer stays attached", async () => { + // The socket is shut down right after the response is written, so later + // body segments are not delivered; 'end' must still fire so the consumer + // is not left waiting. + let ended = false; + const { promise: closed, resolve, reject } = Promise.withResolvers(); + await using server = createServer((req, res) => { + req.on("data", () => {}); + req.once("end", () => (ended = true)); + req.once("error", reject); + req.once("close", resolve); + res.end("ok"); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const { port } = server.address() as AddressInfo; + + const sock = connect(port, "127.0.0.1"); + await once(sock, "connect"); + sock.on("error", () => {}); + sock.write("POST / HTTP/1.1\r\nHost: x\r\nConnection: close\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhello\r\n"); + sock.setNoDelay(true); + await new Promise(r => setImmediate(r)); + sock.write("6\r\n world\r\n0\r\n\r\n"); + sock.resume(); + + await closed; + sock.end(); + expect(ended).toBe(true); + }); +});