diff --git a/src/js/node/net.ts b/src/js/node/net.ts index 2046ed40ae41..c93acb6f7bd6 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -883,7 +883,9 @@ const ServerHandlers: SocketHandler = { self.emit("secureConnect", verifyError); if (server?.pauseOnConnect) { self.pause(); - } else { + } else if (self.readableFlowing === null) { + // See onconnection: honor a pause()/'data'/'readable' touched inside + // the secureConnection/secureConnect handlers. self.resume(); } }, @@ -1063,8 +1065,11 @@ function onconnection(err, clientHandle) { } self.emit("connection", _socket); - // the duplex implementation start paused, so we resume when pauseOnConnect is falsy - if (!pauseOnConnect && !isTLS) { + // Honor a pause()/'data'/'readable' touched inside the handler. null (the + // handler left flowing untouched) still resumes so a write-only handler's + // peer-close tears down; Node would leave it null and release the loop via + // UV_EOF readStop, which needs accepted sockets to hold the loop themselves. + if (!pauseOnConnect && !isTLS && _socket.readableFlowing === null) { _socket.resume(); } } diff --git a/test/js/node/net/net-server-accepted-socket-pause-fixture.js b/test/js/node/net/net-server-accepted-socket-pause-fixture.js new file mode 100644 index 000000000000..9a7ec9d0bc7a --- /dev/null +++ b/test/js/node/net/net-server-accepted-socket-pause-fixture.js @@ -0,0 +1,69 @@ +"use strict"; +// onconnection's post-emit resume() previously stomped a pause() made inside +// the 'connection' handler: flowing went back to true with no listener, bytes +// were discarded, and the writer saw 'drain' against a paused peer. +const net = require("node:net"); +const { once } = require("node:events"); + +function waitFor(cond) { + return new Promise(resolve => { + const check = () => (cond() ? resolve() : setImmediate(check)); + check(); + }); +} + +(async () => { + let acc; + const server = net.createServer(); + server.on("connection", s => { + acc = s; + s.on("error", () => {}); + s.pause(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const cli = net.connect(server.address().port, "127.0.0.1"); + cli.on("error", () => {}); + await once(cli, "connect"); + await waitFor(() => acc); + // readableFlowing must be false (the handler's pause()) now that the + // handler has returned; before the fix it was flipped back to true. + console.log("flowing", acc.readableFlowing); + + const chunk = Buffer.alloc(64 * 1024, 0x61); + let writes = 0; + let drains = 0; + cli.on("drain", () => drains++); + // Write until backpressure or an upper bound: with a paused peer the kernel + // buffers should fill long before 4000 chunks. + while (writes < 4000) { + writes++; + if (!cli.write(chunk)) break; + } + let turns = 0; + await waitFor(() => drains > 0 || ++turns > 200); + // 'drain' can fire while the peer is paused (the Writable buffer flushes + // into the kernel send buffer, whose capacity is platform-dependent); the + // distinguishing observables are readableFlowing and whether bytes are + // delivered after resume(). + console.log("backpressured", writes < 4000); + + let got = 0; + acc.on("data", d => (got += d.length)); + acc.resume(); + const want = writes * chunk.length; + // If the bytes were already discarded (the bug) nothing more is coming; + // bound the wait so the broken build reports false instead of hanging. + let dturns = 0; + await waitFor(() => got >= want || acc.destroyed || ++dturns > 2000); + console.log("delivered", got === want); + + cli.destroy(); + acc.destroy(); + server.close(); +})().then( + () => process.exit(0), + err => { + console.error(err && err.stack ? err.stack : String(err)); + process.exit(1); + }, +); diff --git a/test/js/node/net/node-net.test.ts b/test/js/node/net/node-net.test.ts index b82972d3691d..f48cd5c4097d 100644 --- a/test/js/node/net/node-net.test.ts +++ b/test/js/node/net/node-net.test.ts @@ -1132,3 +1132,24 @@ it.skipIf(isWindows)("connect({ localPort }) succeeds when the local port has TI target.close(); } }); + +// onconnection / ServerHandlers.handshake previously resume()d after emit, +// stomping a pause() made inside the handler. Subprocess-isolated so nothing +// in the test runner touches readableFlowing. +describe.each([ + ["net.createServer 'connection'", "net-server-accepted-socket-pause-fixture.js"], + ["tls.createServer 'secureConnection'", "tls-server-accepted-socket-pause-fixture.js"], +])("accepted socket honors pause() made inside the %s handler", (_, fixture) => { + it("leaves readableFlowing false and delivers every byte after resume()", async () => { + await using proc = Bun.spawn({ + cmd: [bunExe(), join(import.meta.dir, fixture)], + 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().split("\n")).toEqual(["flowing false", "backpressured true", "delivered true"]); + expect(exitCode).toBe(0); + }); +}); diff --git a/test/js/node/net/tls-server-accepted-socket-pause-fixture.js b/test/js/node/net/tls-server-accepted-socket-pause-fixture.js new file mode 100644 index 000000000000..fa7cc5ea7ead --- /dev/null +++ b/test/js/node/net/tls-server-accepted-socket-pause-fixture.js @@ -0,0 +1,69 @@ +"use strict"; +// TLS sibling of net-server-accepted-socket-pause-fixture.js: pause() inside +// 'secureConnection' must be honored by ServerHandlers.handshake's post-emit +// resume() gate. +const tls = require("node:tls"); +const fs = require("node:fs"); +const path = require("node:path"); +const { once } = require("node:events"); + +const keys = path.join(__dirname, "..", "test", "fixtures", "keys"); + +function waitFor(cond) { + return new Promise(resolve => { + const check = () => (cond() ? resolve() : setImmediate(check)); + check(); + }); +} + +(async () => { + let acc; + const server = tls.createServer({ + key: fs.readFileSync(path.join(keys, "agent1-key.pem")), + cert: fs.readFileSync(path.join(keys, "agent1-cert.pem")), + }); + server.on("secureConnection", s => { + acc = s; + s.on("error", () => {}); + s.pause(); + }); + await once(server.listen(0, "127.0.0.1"), "listening"); + const cli = tls.connect({ port: server.address().port, host: "127.0.0.1", rejectUnauthorized: false }); + cli.on("error", () => {}); + await once(cli, "secureConnect"); + await waitFor(() => acc); + console.log("flowing", acc.readableFlowing); + + const chunk = Buffer.alloc(64 * 1024, 0x61); + let writes = 0; + let drains = 0; + cli.on("drain", () => drains++); + while (writes < 4000) { + writes++; + if (!cli.write(chunk)) break; + } + let turns = 0; + await waitFor(() => drains > 0 || ++turns > 200); + // drains can be >0 under TLS (engine-internal buffering flushes to the + // kernel independently of the peer's read state); the distinguishing + // observables are readableFlowing and whether the bytes are delivered. + console.log("backpressured", writes < 4000); + + let got = 0; + acc.on("data", d => (got += d.length)); + acc.resume(); + const want = writes * chunk.length; + let dturns = 0; + await waitFor(() => got >= want || acc.destroyed || ++dturns > 2000); + console.log("delivered", got === want); + + cli.destroy(); + acc.destroy(); + server.close(); +})().then( + () => process.exit(0), + err => { + console.error(err && err.stack ? err.stack : String(err)); + process.exit(1); + }, +);