From df3eb3849a28ed0d74f2642422ad6f5c18639097 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Wed, 7 Oct 2026 22:39:03 -0700 Subject: [PATCH 1/3] fix(tls): destroy wrapped transports without ending them Match Node 24 destruction for Duplex-backed TLS while keeping graceful shutdown separate. Retain adopted-fd close ownership and release HTTP/2 injected transports when their TLS proxy is destroyed. Adapt HTTP/2 lifecycle coverage from oven-sh/bun#38154. Co-authored-by: robobun <117481402+robobun@users.noreply.github.com> --- CHANGELOG.md | 3 +- docs/runtime/nodejs-compat.mdx | 2 + src/js/node/_http2_upgrade.ts | 5 +- src/js/node/net.ts | 8 +- src/runtime/socket/UpgradedDuplex.rs | 2 - src/runtime/socket/socket_body.rs | 6 +- .../js/node/http2/node-http2-upgrade.test.mts | 63 ++++ test/js/node/tls/node-tls-connect.test.ts | 75 ++-- .../node-tls-duplex-close-throw-uaf.test.ts | 39 +-- .../tls/tls-destroy-transport-fixture.cjs | 320 ++++++++++++++++++ 10 files changed, 461 insertions(+), 62 deletions(-) create mode 100644 test/js/node/tls/tls-destroy-transport-fixture.cjs diff --git a/CHANGELOG.md b/CHANGELOG.md index d10c8e3c7712..44dd7bf7249f 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -192,7 +192,6 @@ - Keep ESM namespaces free of inherited `__esModule` markers and preserve the own marker and live exports for `require(esm)`, fixing Vite/tsx namespace interop. Adapts [oven-sh/bun#33894](https://github.com/oven-sh/bun/pull/33894) and [oven-sh/WebKit#279](https://github.com/oven-sh/WebKit/pull/279). Thanks @robobun! - - Sync oven-sh/bun through `c7b06d94bac19817ba34b6677bb1099fb4f6d2be`, preserving fork fixes and incorporating TLS handshake shutdown, macOS split-DNS failover, file-body cloning, Buffer write validation, mimalloc 3.5.3 and idle-memory release. - Pin immutable [OpenClaw WebKit `42ab38d705`](https://github.com/openclaw/WebKit/releases/tag/autobuild-42ab38d705d4838748ccee77e7deb0e4e35515ee) by archive checksum, together with the required namespace facade integration from #106; retain fail-closed artifact selection. @@ -247,3 +246,5 @@ - Share exact owned source buffers on matching `node:vm` compilation-cache entries while preserving cold misses, cached-data validation, per-script origins and the existing byte budget. - Name package-target resolution options explicitly to satisfy the Rust Mordant lint without changing resolution behavior. + +- Destroy Duplex-backed TLS transports without calling `end()`, matching Node.js, while preserving graceful TLS shutdown. Adapts HTTP/2 transport coverage from [oven-sh/bun#38154](https://github.com/oven-sh/bun/pull/38154); thanks @robobun! diff --git a/docs/runtime/nodejs-compat.mdx b/docs/runtime/nodejs-compat.mdx index a3dc15f08a7b..a812f5e64060 100644 --- a/docs/runtime/nodejs-compat.mdx +++ b/docs/runtime/nodejs-compat.mdx @@ -266,6 +266,8 @@ Native `fetch()` uses `tls.setDefaultCACertificates()` overrides for new TLS con TLS starts reading paused Duplex and HTTP CONNECT transports after installing its listeners, including handshake bytes buffered before adoption. +Destroying a TLS socket destroys its underlying Duplex transport without calling `end()`, matching Node.js. Graceful TLS shutdown still ends the transport after sending its TLS shutdown data. + 🟡 Missing `pskCallback`, OCSP stapling (`requestOCSP`), the server `'newSession'`/`'resumeSession'` events and session ticket keys (`ticketKeys` is ignored). As a result, session resumption does not work across processes. Bun uses BoringSSL, so `tlsSocket.renegotiate()` always fails and `getEphemeralKeyInfo()`/`getSharedSigalgs()` return no information. ### [`node:util`](https://nodejs.org/api/util.html) diff --git a/src/js/node/_http2_upgrade.ts b/src/js/node/_http2_upgrade.ts index 98739ae42758..5f9c8c0e8a4f 100644 --- a/src/js/node/_http2_upgrade.ts +++ b/src/js/node/_http2_upgrade.ts @@ -5,7 +5,7 @@ const kSharedCreds = Symbol.for("::buntlssharedcreds::"); interface NativeHandle { resume(): void; close(): void; - end(): void; + shutdown(): void; $write(chunk: Buffer, encoding: string): boolean; alpnProtocol?: string; } @@ -109,6 +109,7 @@ function tlsSocketDestroy(this: TLSProxySocket, err: Error | null, callback: (er h.close(); this._ctx.nativeHandle = null; } + this._ctx.rawSocket.destroy(); // Must invoke pending write callback with error per Writable stream contract const writeCb = this._writeCallback; if (writeCb) { @@ -125,7 +126,7 @@ function tlsSocketFinal(this: TLSProxySocket, callback: () => void) { const h = this._ctx.nativeHandle; if (!h) return callback(); // Signal end-of-stream to the TLS layer - h.end(); + h.shutdown(); callback(); } diff --git a/src/js/node/net.ts b/src/js/node/net.ts index 385fc7f0f741..d864d281ea2e 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -2414,12 +2414,10 @@ Socket.prototype._destroy = function _destroy(err, callback) { // Node: after 'error', before 'close'. With no error it goes first, so the stream is errored before the EOF that the native close handler pushes can emit 'end'. if (canceledWrite && !err) process.nextTick(cancelWriteNT, canceledWrite); - // Tear down a wrapped generic duplex with this socket: the native handle's - // close only flushes close_notify and lets the wrapper drain; without an - // explicit destroy here a late RST on the underlying transport can surface - // as an unhandled error after this socket is gone. + // Stream-level TLS owns its transport; an adopted fd pair closes through closeOwedRaw. + // https://github.com/nodejs/node/blob/v24.21.0/lib/internal/js_stream_socket.js#L253 const upgraded = this[kupgraded]; - if (upgraded && !(upgraded instanceof Socket) && !upgraded.destroyed) { + if (upgraded && !upgraded.destroyed && (!(upgraded instanceof Socket) || !upgraded._handle?.[kAdoptedTLSRaw])) { upgraded.destroy?.(); } diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index 88f0656d7e09..cdf01416517c 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -220,8 +220,6 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); (self.handlers.on_close)(self.handlers.ctx); - // closes the underlying duplex - self.call_write_or_end(None, false); // Early teardown (struct itself is dropped later by parent). self.teardown(); diff --git a/src/runtime/socket/socket_body.rs b/src/runtime/socket/socket_body.rs index 84dfc2309afc..a8a40021795f 100644 --- a/src/runtime/socket/socket_body.rs +++ b/src/runtime/socket/socket_body.rs @@ -4483,11 +4483,7 @@ impl DuplexUpgradeContext { // is the ext-slot/owner pin). Null our pointer first so the // `deinit_in_next_tick` → `deinit` path doesn't deref it a second // time — that's the over-deref behind the cross-file - // `TLSSocket::finalize` use-after-poison. It also means a throw - // from `duplex.end()` (called right after this returns via - // `UpgradedDuplex::on_close` → `call_write_or_end`) hits the null-check - // in `on_error` instead of reading the Handlers that `TLSSocket::on_close` - // → `mark_inactive` just released. + // `TLSSocket::finalize` use-after-poison. let p = tls.into_this_ptr(); crate::dispatch::fold(TLSSocket::on_close(p, socket, 0, None)); } diff --git a/test/js/node/http2/node-http2-upgrade.test.mts b/test/js/node/http2/node-http2-upgrade.test.mts index 57a8b73d93bc..a400af9594af 100644 --- a/test/js/node/http2/node-http2-upgrade.test.mts +++ b/test/js/node/http2/node-http2-upgrade.test.mts @@ -16,6 +16,7 @@ import fs from "node:fs"; import http2 from "node:http2"; import net from "node:net"; import path from "node:path"; +import { Duplex } from "node:stream"; import { afterEach, describe, test } from "node:test"; import tls from "node:tls"; import { fileURLToPath } from "node:url"; @@ -461,6 +462,68 @@ describe("HTTP/2 upgrade — server TLS options", () => { }); }); +// Adapted from oven-sh/bun#38154: the peer deliberately keeps TCP half-open, +// so only destroying the accepted transport releases the server connection. +describe("HTTP/2 upgrade — session destruction releases the accepted transport", () => { + for (const withError of [false, true]) { + test(`destroyed ${withError ? "with" : "without"} an error`, async () => { + const h2Server = http2.createSecureServer(TLS); + const sessionClosed = new Promise(resolve => { + h2Server.once("session", (session: http2.ServerHttp2Session) => { + session.on("error", () => {}); + session.once("close", resolve); + session.destroy(withError ? new Error("test teardown") : undefined); + }); + }); + const accepted = Promise.withResolvers<{ raw: net.Socket; closed: Promise }>(); + const netServer = net.createServer(raw => { + const closed = new Promise(resolve => raw.once("close", resolve)); + accepted.resolve({ raw, closed }); + h2Server.emit("connection", raw); + }); + await new Promise(resolve => netServer.listen(0, "127.0.0.1", resolve)); + const tcp = net.connect({ + port: (netServer.address() as net.AddressInfo).port, + host: "127.0.0.1", + allowHalfOpen: true, + }); + tcp.on("error", () => {}); + const carrier = new Duplex({ + read() {}, + write(chunk, _encoding, callback) { + if (tcp.destroyed) return callback(); + tcp.write(chunk, () => callback()); + }, + final(callback) { + callback(); + }, + }); + tcp.on("data", chunk => carrier.push(chunk)); + tcp.on("end", () => carrier.push(null)); + const client = tls.connect({ socket: carrier, rejectUnauthorized: false, ALPNProtocols: ["h2"] }); + client.on("error", () => {}); + client.resume(); + try { + await sessionClosed; + const { raw, closed } = await accepted.promise; + assert.strictEqual(await closed, false); + assert.strictEqual(raw.destroyed, true); + const connections = await new Promise((resolve, reject) => { + netServer.getConnections((err, count) => (err ? reject(err) : resolve(count))); + }); + assert.strictEqual(connections, 0); + await new Promise(resolve => netServer.close(() => resolve())); + } finally { + client.destroy(); + carrier.destroy(); + tcp.destroy(); + (await accepted.promise).raw.destroy(); + if (netServer.listening) netServer.close(); + } + }); + } +}); + if (typeof Bun !== "undefined") { describe("Node.js compatibility", () => { test("tests should run on node.js", async () => { diff --git a/test/js/node/tls/node-tls-connect.test.ts b/test/js/node/tls/node-tls-connect.test.ts index f779be5ec939..27a5ec9540b5 100644 --- a/test/js/node/tls/node-tls-connect.test.ts +++ b/test/js/node/tls/node-tls-connect.test.ts @@ -2061,8 +2061,55 @@ describe("a TLS socket over a Duplex transport reports that transport's error", expect(events).toEqual(["error:transport failed", "close destroyed=true"]); }); - it("what a method or an accessor of the transport throws is reported", async () => { - // Out of process: a socket that has closed emits no 'error', so the throw is uncaught, as in node. + it("destroying TLS destroys its transport without ending it", async () => { + // Node 24 destroys the transport without accessing end, even if destroy is a no-op. + // Observe natural quiescence before cleanup so a delayed end getter cannot escape the assertion. + await using proc = Bun.spawn({ + cmd: [bunExe(), join(import.meta.dir, "tls-destroy-transport-fixture.cjs")], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + const rows = stdout + .trim() + .split("\n") + .map(line => JSON.parse(line)); + const scenarios = rows.filter(row => row.type === "scenario-result"); + expect(scenarios.map(row => row.scenario)).toEqual([ + "isolated-immediate", + "isolated-after-check", + "consecutive-immediate-after-check", + "consecutive-after-check-immediate", + ]); + expect(scenarios.flatMap(row => row.cases)).toHaveLength(6); + for (const scenario of scenarios) { + expect(scenario).toMatchObject({ + pass: true, + failures: [], + sequenceComplete: true, + cleanupComplete: true, + uncaughtCount: 0, + rejectionCount: 0, + }); + for (const result of scenario.cases) { + expect(result).toMatchObject({ + getterCount: 0, + destroyCalls: [{ phase: "observe", hasError: false, errorId: null }], + tlsCloses: [false], + continuation: true, + tlsErrors: [], + transportErrors: [], + transportCloses: [{ phase: "cleanup", cleanupRequested: true }], + }); + } + } + expect(rows.at(-1)).toMatchObject({ type: "parent-result", pass: true, timedOut: false }); + expect(stderr).toBe(""); + expect(exitCode).toBe(0); + }); + + it("what the transport write method or accessor throws is reported", async () => { const script = ` const tls = require("node:tls"); const { Duplex } = require("node:stream"); @@ -2070,13 +2117,11 @@ describe("a TLS socket over a Duplex transport reports that transport's error", const uncaught = []; process.on("uncaughtException", err => uncaught.push(err.message)); - async function run(method, kind, when) { + async function run(kind) { const transport = new Duplex({ read() {}, write(chunk, encoding, callback) { callback(); } }); - const name = method + " " + kind; + const name = "write " + kind; const fail = () => { throw new Error(name); }; - Object.defineProperty(transport, method, kind === "call" ? { value: fail } : { get: fail }); - // Only a transport that is still open is ended. - transport.destroy = function () { return this; }; + Object.defineProperty(transport, "write", kind === "call" ? { value: fail } : { get: fail }); const socket = tls.connect({ socket: transport, rejectUnauthorized: false }); const events = []; const closed = Promise.withResolvers(); @@ -2085,17 +2130,13 @@ describe("a TLS socket over a Duplex transport reports that transport's error", events.push("close:" + hadError); closed.resolve(); }); - if (when === "after the engine started") await new Promise(resolve => setImmediate(resolve)); // The ClientHello is the write that fails, and that destroys the socket. - if (method === "end") socket.destroy(); await closed.promise; - console.log(name + ", " + when + ": " + events.join("|") + " uncaught:" + uncaught.splice(0).join("|")); + console.log(name + ": " + events.join("|") + " uncaught:" + uncaught.splice(0).join("|")); } for (const kind of ["call", "getter"]) { - await run("end", kind, "before the engine starts"); - await run("end", kind, "after the engine started"); - await run("write", kind, "after the engine started"); + await run(kind); } `; await using proc = Bun.spawn({ @@ -2107,12 +2148,8 @@ describe("a TLS socket over a Duplex transport reports that transport's error", const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); expect({ stdout: stdout.trim().split("\n"), stderr, exitCode }).toEqual({ stdout: [ - "end call, before the engine starts: close:false uncaught:end call", - "end call, after the engine started: close:false uncaught:end call", - "write call, after the engine started: error:write call|close:true uncaught:", - "end getter, before the engine starts: close:false uncaught:end getter", - "end getter, after the engine started: close:false uncaught:end getter", - "write getter, after the engine started: error:write getter|close:true uncaught:", + "write call: error:write call|close:true uncaught:", + "write getter: error:write getter|close:true uncaught:", ], stderr: "", exitCode: 0, diff --git a/test/js/node/tls/node-tls-duplex-close-throw-uaf.test.ts b/test/js/node/tls/node-tls-duplex-close-throw-uaf.test.ts index e71a61c73b83..d909c945bd11 100644 --- a/test/js/node/tls/node-tls-duplex-close-throw-uaf.test.ts +++ b/test/js/node/tls/node-tls-duplex-close-throw-uaf.test.ts @@ -34,19 +34,15 @@ async function run(script: string, expected = "ok") { // build enables); release builds may read garbage without trapping. Each // test spawns an independent subprocess so they can run concurrently. describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) does not read freed Handlers", () => { - test("when duplex.end() throws after close", async () => { - // UpgradedDuplex.onClose → DuplexUpgradeContext.onClose → TLSSocket.onClose - // frees the Handlers; UpgradedDuplex.onClose then calls duplex.end(). If - // that throws, onError → tls.handleError → getHandlers() read the freed - // allocation. + test("closing a duplex does not call its throwing end method", async () => { + // Closing once called end() after freeing the Handlers. Like Node, destruction + // now avoids end() entirely, including when the transport emits close. await run( ` const tls = require("node:tls"); const { Duplex } = require("node:stream"); - // Minimal duplex: write/read are no-ops. The only thing that matters is - // that end() throws synchronously — UpgradedDuplex.callWriteOrEnd catches - // that and routes it through onError. + // A throwing end method exposes any accidental graceful shutdown during close. const duplex = new Duplex({ read() {}, write(chunk, enc, cb) { cb(); }, @@ -55,7 +51,7 @@ describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) d duplex.end = function () { throw new Error("end() throws during close"); }; - // Only a duplex that is still open gets that end(). + // Keep the transport open so destroyed-state guards cannot hide an end call. duplex.destroy = function () { return this; }; @@ -68,11 +64,7 @@ describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) d sock.on("error", () => {}); sock.on("close", () => {}); - // startTLS runs on the next tick; once onOpen has fired (is_open=true), - // emitting "close" on the duplex triggers the SSL wrapper's fast - // shutdown → UpgradedDuplex.onClose → DuplexUpgradeContext.onClose → - // TLSSocket.onClose (frees handlers) → callWriteOrEnd → duplex.end() - // throws → onError. + // Exercise transport close after the TLS engine has started. setImmediate(() => { setImmediate(() => { duplex.emit("close"); @@ -83,7 +75,7 @@ describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) d }); }); `, - "uncaught: end() throws during close\nok", + "ok", ); }); @@ -127,12 +119,9 @@ describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) d }); // Serial: four debug subprocesses at once reach the default test timeout on a loaded machine. - test.serial("when duplex.end() throws after a close that comes before StartTLS", async () => { - // No SSL wrapper exists yet, so the queued .StartTLS task carries the - // close out: TLSSocket.onClose frees the Handlers, then duplex.end() - // throws into onError. Only a duplex that is still open gets that - // end(), so destroy() leaves these open. One process runs both closes: - // each socket has its own Handlers. + test.serial("a close before StartTLS does not call duplex.end()", async () => { + // The queued StartTLS task carries out early closes. Both close paths + // must avoid end(), even when the transport's destroy() leaves it open. await run( ` const tls = require("node:tls"); @@ -176,13 +165,7 @@ describe.concurrent.skipIf(!isASAN && !isDebug)("tls.connect({socket: Duplex}) d }); }); `, - [ - "uncaught: end() throws when the duplex closes", - "uncaught: end() throws when the socket is destroyed", - "the duplex closes, end() calls: 1", - "the socket is destroyed, end() calls: 1", - "ok", - ].join("\n"), + ["the duplex closes, end() calls: 0", "the socket is destroyed, end() calls: 0", "ok"].join("\n"), ); }); diff --git a/test/js/node/tls/tls-destroy-transport-fixture.cjs b/test/js/node/tls/tls-destroy-transport-fixture.cjs new file mode 100644 index 000000000000..b331e922cca3 --- /dev/null +++ b/test/js/node/tls/tls-destroy-transport-fixture.cjs @@ -0,0 +1,320 @@ +"use strict"; + +const { spawn } = require("node:child_process"); +const { writeSync } = require("node:fs"); +const { Duplex } = require("node:stream"); +const tls = require("node:tls"); + +const originalTransportDestroy = Duplex.prototype.destroy; +const scenarios = { + "isolated-immediate": ["immediate"], + "isolated-after-check": ["after-check"], + "consecutive-immediate-after-check": ["immediate", "after-check"], + "consecutive-after-check-immediate": ["after-check", "immediate"], +}; + +function output(value) { + writeSync(1, JSON.stringify(value) + "\n"); +} + +function deferred() { + let resolve; + const promise = new Promise(done => { + resolve = done; + }); + return { promise, resolve }; +} + +async function parent() { + const results = []; + let active; + let timedOut = false; + output({ + type: "runtime", + version: process.version, + versions: process.versions, + platform: process.platform, + arch: process.arch, + }); + const watchdog = setTimeout(() => { + timedOut = true; + output({ type: "watchdog-failure", milliseconds: 5000 }); + active?.kill("SIGKILL"); + }, 5000); + watchdog.unref(); + for (const scenario of Object.keys(scenarios)) { + if (timedOut) break; + results.push( + await new Promise(resolve => { + let stderrBytes = 0; + let spawnError = null; + active = spawn(process.execPath, [__filename, "--scenario", scenario], { stdio: ["ignore", "pipe", "pipe"] }); + active.stdout.on("data", chunk => writeSync(1, chunk)); + active.stderr.on("data", chunk => { + stderrBytes += chunk.length; + writeSync(2, chunk); + }); + active.on("error", error => { + spawnError = error.message; + }); + active.on("close", (code, signal) => resolve({ scenario, code, signal, stderrBytes, spawnError })); + }), + ); + active = undefined; + } + clearTimeout(watchdog); + const pass = + !timedOut && + results.length === Object.keys(scenarios).length && + results.every(result => result.code === 0 && !result.signal && !result.stderrBytes && !result.spawnError); + output({ type: "parent-result", pass, timedOut, results }); + process.exitCode = pass ? 0 : 1; +} + +function child(scenario) { + const records = []; + const errorOwners = new Map(); + const trace = []; + const seenUncaught = []; + const failures = []; + let activeCase = null; + let phase = "observe"; + let sequenceComplete = false; + let cleanupComplete = false; + let pendingBoundaries = 0; + let uncaughtCount = 0; + let rejectionCount = 0; + + function event(owner, name, data = {}) { + const entry = { type: "event", scenario, seq: trace.length + 1, phase, owner, activeCase, name, ...data }; + trace.push(entry); + output(entry); + } + + function identity(error) { + return errorOwners.get(error)?.id ?? "unknown"; + } + + function boundaries(record, from) { + pendingBoundaries += 3; + process.nextTick(() => { + pendingBoundaries--; + event(record.id, "boundary.nextTick", { from }); + }); + queueMicrotask(() => { + pendingBoundaries--; + event(record.id, "boundary.microtask", { from }); + }); + setImmediate(() => { + pendingBoundaries--; + event(record.id, "boundary.check", { from }); + }); + } + + process.on("uncaughtException", (error, origin) => { + uncaughtCount++; + const owner = errorOwners.get(error); + seenUncaught.push(identity(error)); + event(owner?.id ?? null, "uncaughtException", { errorId: identity(error), message: error.message, origin }); + if (owner) boundaries(owner, "uncaughtException"); + }); + process.on("unhandledRejection", error => { + rejectionCount++; + event(null, "unhandledRejection", { errorId: identity(error), message: String(error) }); + }); + + async function runCase(timing, index) { + const id = `${scenario}/${index + 1}-${timing}`; + activeCase = id; + const record = { + id, + timing, + writes: 0, + writesAtDestroy: null, + getterCount: 0, + destroyCalls: [], + tlsCloses: [], + tlsErrors: [], + transportErrors: [], + transportCloses: [], + continuation: false, + cleanupRequested: false, + closed: deferred(), + transportClosed: deferred(), + }; + records.push(record); + const thrown = new Error(`unexpected end getter: ${id}`); + errorOwners.set(thrown, record); + event(id, "case.begin", { timing }); + const transport = new Duplex({ + read() {}, + write(chunk, encoding, callback) { + record.writes++; + event(id, "transport.write", { bytes: chunk.length }); + callback(); + }, + }); + record.transport = transport; + Object.defineProperty(transport, "end", { + get() { + record.getterCount++; + event(id, "transport.end.get", { errorId: id }); + throw thrown; + }, + }); + transport.destroy = function (error) { + const call = { phase, hasError: error != null, errorId: error == null ? null : identity(error) }; + record.destroyCalls.push(call); + event(id, "transport.destroy.noop", call); + return this; + }; + transport.on("error", error => { + record.transportErrors.push(identity(error)); + event(id, "transport.error", { errorId: identity(error), message: error.message }); + }); + transport.on("close", () => { + record.transportCloses.push({ phase, cleanupRequested: record.cleanupRequested }); + event(id, "transport.close", { cleanupRequested: record.cleanupRequested }); + record.transportClosed.resolve(); + }); + const socket = tls.connect({ socket: transport, rejectUnauthorized: false }); + socket.on("error", error => { + record.tlsErrors.push(identity(error)); + event(id, "tls.error", { errorId: identity(error), message: error.message }); + }); + socket.on("close", hadError => { + record.tlsCloses.push(hadError); + event(id, "tls.close", { hadError }); + record.closed.resolve(); + event(id, "close.promise.resolved"); + boundaries(record, "tls.close"); + }); + if (timing === "after-check") { + await new Promise(resolve => + setImmediate(() => { + event(id, "timing.check", { writes: record.writes }); + resolve(); + }), + ); + event(id, "timing.check.continuation", { writes: record.writes }); + } + record.writesAtDestroy = record.writes; + event(id, "tls.destroy.call", { writes: record.writes }); + socket.destroy(); + event(id, "tls.destroy.return"); + boundaries(record, "tls.destroy.return"); + await record.closed.promise; + record.continuation = true; + event(id, "close.promise.continuation", { seenUncaught: [...seenUncaught] }); + boundaries(record, "close.promise.continuation"); + } + + function checkCommon(at) { + if (!sequenceComplete) failures.push(`${at}: sequence incomplete`); + if (records.length !== scenarios[scenario].length) failures.push(`${at}: missing cases`); + if (pendingBoundaries) failures.push(`${at}: missing boundary callbacks ${pendingBoundaries}`); + if (uncaughtCount || rejectionCount) + failures.push(`${at}: uncaught/rejection counts ${uncaughtCount}/${rejectionCount}`); + for (const record of records) { + if (record.getterCount !== 0) + failures.push(`${at}: ${record.id}: end getter accessed ${record.getterCount} times`); + if ( + record.destroyCalls.length !== 1 || + record.destroyCalls[0].phase !== "observe" || + record.destroyCalls[0].hasError + ) + failures.push(`${at}: ${record.id}: expected one runtime destroy without error`); + if (record.tlsCloses.length !== 1 || record.tlsCloses[0] !== false || !record.continuation) + failures.push(`${at}: ${record.id}: TLS close/continuation contract`); + if (record.tlsErrors.length || record.transportErrors.length) + failures.push(`${at}: ${record.id}: unexpected error event`); + } + } + + async function cleanup() { + for (const record of records) { + record.cleanupRequested = true; + event(record.id, "cleanup.transport.destroy.call"); + originalTransportDestroy.call(record.transport); + } + await Promise.all(records.map(record => record.transportClosed.promise)); + cleanupComplete = true; + event(null, "cleanup.complete"); + } + + process.on("beforeExit", () => { + if (phase === "observe") { + event(null, "observation.natural-quiescence", { sequenceComplete, pendingBoundaries }); + checkCommon("before cleanup"); + for (const record of records) { + if (record.transportCloses.length) failures.push(`before cleanup: ${record.id}: premature transport close`); + } + phase = "cleanup"; + setImmediate(() => { + event(null, "cleanup.check"); + cleanup().catch(error => { + failures.push(`cleanup rejected: ${String(error)}`); + event(null, "cleanup.rejected", { errorId: identity(error) }); + }); + }); + return; + } + if (phase !== "cleanup") return; + event(null, "cleanup.natural-quiescence", { cleanupComplete, pendingBoundaries }); + checkCommon("after cleanup"); + if (!cleanupComplete) failures.push("cleanup incomplete at natural quiescence"); + for (const record of records) { + if ( + record.transportCloses.length !== 1 || + record.transportCloses[0].phase !== "cleanup" || + !record.transportCloses[0].cleanupRequested + ) + failures.push(`${record.id}: expected one reviewer-owned transport close`); + } + phase = "done"; + const cases = records.map(record => ({ + id: record.id, + timing: record.timing, + writesAtDestroy: record.writesAtDestroy, + writes: record.writes, + getterCount: record.getterCount, + destroyCalls: record.destroyCalls, + tlsCloses: record.tlsCloses, + continuation: record.continuation, + tlsErrors: record.tlsErrors, + transportErrors: record.transportErrors, + transportCloses: record.transportCloses, + })); + output({ + type: "scenario-result", + scenario, + pass: failures.length === 0, + failures, + sequenceComplete, + cleanupComplete, + uncaughtCount, + rejectionCount, + cases, + }); + process.exitCode = failures.length === 0 ? 0 : 1; + }); + + async function main() { + for (const [index, timing] of scenarios[scenario].entries()) await runCase(timing, index); + sequenceComplete = true; + event(null, "sequence.close-continuations.complete"); + } + main().catch(error => { + failures.push(`sequence rejected: ${String(error)}`); + event(null, "sequence.rejected", { errorId: identity(error) }); + }); +} + +if (process.argv[2] === "--scenario" && Object.hasOwn(scenarios, process.argv[3])) { + child(process.argv[3]); +} else { + parent().catch(error => { + writeSync(2, String(error) + "\n"); + process.exitCode = 1; + }); +} From 48ab860734885bee7c43d290f3acb3ca3a184958 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 8 Oct 2026 01:48:43 -0700 Subject: [PATCH 2/3] fix(tls): preserve wrapped transport ownership and half-close ordering Keep adopted-fd ownership on the TLS socket after its raw handle detaches. Separate peer EOF, graceful writable shutdown, and full TLS destruction; wait for transport completion without forwarding its error a second time. Use the inherited tls.Server connection path for injected HTTP/2 sockets instead of maintaining a second TLS transport adapter. Port HTTP/2 destroy-versus-close and final event ordering from oven-sh/bun#38195, and the destroyed-socket EOF guard from #43392. Retain the six-case destruction regression and add Node 24 controls for half-open replies, raw EOF, renegotiation shutdown, and transport ownership. Synchronize conformance assertions with the sessionError event itself. Co-authored-by: robobun <117481402+robobun@users.noreply.github.com> --- CHANGELOG.md | 3 +- docs/runtime/nodejs-compat.mdx | 4 + src/http/ProxyTunnel.rs | 1 + .../websocket_client/WebSocketProxyTunnel.rs | 1 + src/js/node/_http2_upgrade.ts | 398 ------------------ src/js/node/http2.ts | 97 ++--- src/js/node/net.ts | 42 +- src/runtime/socket/UpgradedDuplex.rs | 78 ++-- src/runtime/socket/WindowsNamedPipe.rs | 1 + src/uws/lib.rs | 51 ++- test/js/node/http2/h2-conformance.test.ts | 17 +- .../js/node/http2/node-http2-upgrade.test.mts | 29 +- test/js/node/tls/node-tls-connect.test.ts | 28 +- test/js/node/tls/node-tls-raw-end.test.ts | 100 +++++ test/js/node/tls/renegotiation.test.ts | 56 +++ .../tls/tls-half-close-transport-fixture.cjs | 207 +++++++++ .../tls/tls-renegotiate-shutdown-fixture.cjs | 248 +++++++++++ 17 files changed, 817 insertions(+), 544 deletions(-) delete mode 100644 src/js/node/_http2_upgrade.ts create mode 100644 test/js/node/tls/tls-half-close-transport-fixture.cjs create mode 100644 test/js/node/tls/tls-renegotiate-shutdown-fixture.cjs diff --git a/CHANGELOG.md b/CHANGELOG.md index 44dd7bf7249f..6cd26c61c953 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -192,6 +192,7 @@ - Keep ESM namespaces free of inherited `__esModule` markers and preserve the own marker and live exports for `require(esm)`, fixing Vite/tsx namespace interop. Adapts [oven-sh/bun#33894](https://github.com/oven-sh/bun/pull/33894) and [oven-sh/WebKit#279](https://github.com/oven-sh/WebKit/pull/279). Thanks @robobun! + - Sync oven-sh/bun through `c7b06d94bac19817ba34b6677bb1099fb4f6d2be`, preserving fork fixes and incorporating TLS handshake shutdown, macOS split-DNS failover, file-body cloning, Buffer write validation, mimalloc 3.5.3 and idle-memory release. - Pin immutable [OpenClaw WebKit `42ab38d705`](https://github.com/openclaw/WebKit/releases/tag/autobuild-42ab38d705d4838748ccee77e7deb0e4e35515ee) by archive checksum, together with the required namespace facade integration from #106; retain fail-closed artifact selection. @@ -247,4 +248,4 @@ - Name package-target resolution options explicitly to satisfy the Rust Mordant lint without changing resolution behavior. -- Destroy Duplex-backed TLS transports without calling `end()`, matching Node.js, while preserving graceful TLS shutdown. Adapts HTTP/2 transport coverage from [oven-sh/bun#38154](https://github.com/oven-sh/bun/pull/38154); thanks @robobun! +- Destroy Duplex-backed TLS transports without calling `end()`, preserve wrapped-socket close ordering, and honor half-open shutdown like Node.js. Ports HTTP/2 teardown from [oven-sh/bun#38195](https://github.com/oven-sh/bun/pull/38195) and adapts coverage from [oven-sh/bun#38154](https://github.com/oven-sh/bun/pull/38154); thanks @robobun! diff --git a/docs/runtime/nodejs-compat.mdx b/docs/runtime/nodejs-compat.mdx index a812f5e64060..53ebd26446b7 100644 --- a/docs/runtime/nodejs-compat.mdx +++ b/docs/runtime/nodejs-compat.mdx @@ -268,6 +268,10 @@ TLS starts reading paused Duplex and HTTP CONNECT transports after installing it Destroying a TLS socket destroys its underlying Duplex transport without calling `end()`, matching Node.js. Graceful TLS shutdown still ends the transport after sending its TLS shutdown data. +A TLS socket over a Duplex inherits the transport's `allowHalfOpen` setting. A peer TLS shutdown ends the readable side. With `allowHalfOpen: true`, the writable side remains available until the application calls `end()`. + +Destroying an HTTP/2 session releases its TLS transport without waiting for the peer to end its writable side. The session emits its final events after the transport closes. + 🟡 Missing `pskCallback`, OCSP stapling (`requestOCSP`), the server `'newSession'`/`'resumeSession'` events and session ticket keys (`ticketKeys` is ignored). As a result, session resumption does not work across processes. Bun uses BoringSSL, so `tlsSocket.renegotiate()` always fails and `getEphemeralKeyInfo()`/`getSharedSigalgs()` return no information. ### [`node:util`](https://nodejs.org/api/util.html) diff --git a/src/http/ProxyTunnel.rs b/src/http/ProxyTunnel.rs index 2f4e6f990116..04244a3afa21 100644 --- a/src/http/ProxyTunnel.rs +++ b/src/http/ProxyTunnel.rs @@ -605,6 +605,7 @@ impl ProxyTunnel { on_data, on_handshake, on_close, + on_end: None, write: write_encrypted, // fetch's proxy tunnel surfaces no 'session'/'keylog' events; // opting out keeps its SSL off the parked queues entirely. diff --git a/src/http_jsc/websocket_client/WebSocketProxyTunnel.rs b/src/http_jsc/websocket_client/WebSocketProxyTunnel.rs index 38f1e22846a2..f6b26ee19216 100644 --- a/src/http_jsc/websocket_client/WebSocketProxyTunnel.rs +++ b/src/http_jsc/websocket_client/WebSocketProxyTunnel.rs @@ -183,6 +183,7 @@ impl WebSocketProxyTunnel { on_data: Self::on_data, on_handshake: Self::on_handshake, on_close: Self::on_close, + on_end: None, write: Self::write_encrypted, // No JS TLSSocket fronts the tunnel; opting out keeps the // SSL off the parked session/keylog queues entirely. diff --git a/src/js/node/_http2_upgrade.ts b/src/js/node/_http2_upgrade.ts deleted file mode 100644 index 5f9c8c0e8a4f..000000000000 --- a/src/js/node/_http2_upgrade.ts +++ /dev/null @@ -1,398 +0,0 @@ -const { Duplex } = require("node:stream"); -const upgradeDuplexToTLS = $newRustFunction("runtime/socket/socket.rs", "jsUpgradeDuplexToTLS", 2); -const kSharedCreds = Symbol.for("::buntlssharedcreds::"); - -interface NativeHandle { - resume(): void; - close(): void; - shutdown(): void; - $write(chunk: Buffer, encoding: string): boolean; - alpnProtocol?: string; -} - -type UpgradeEventListener = (...args: any[]) => void; -type UpgradeEvents = [UpgradeEventListener, UpgradeEventListener, UpgradeEventListener, UpgradeEventListener]; - -interface UpgradeContextType { - connectionListener: (...args: any[]) => any; - server: Http2SecureServer; - rawSocket: import("node:net").Socket; - nativeHandle: NativeHandle | null; - events: UpgradeEvents | null; -} - -interface Http2SecureServer { - key?: Buffer; - cert?: Buffer; - ca?: Buffer; - passphrase?: string; - ALPNProtocols?: Buffer; - _requestCert?: boolean; - _rejectUnauthorized?: boolean; - emit(event: string, ...args: any[]): boolean; - [kSharedCreds](): { context: unknown }; -} - -type DuplexStream = import("node:stream").Duplex; - -interface TLSProxySocket extends DuplexStream { - _ctx: UpgradeContextType; - _writeCallback: ((err?: Error | null) => void) | null; - alpnProtocol: string | null; - authorized: boolean; - encrypted: boolean; - server: Http2SecureServer; - _requestCert: boolean; - _rejectUnauthorized: boolean | undefined; - _securePending: boolean; - secureConnecting: boolean; - _secureEstablished: boolean; - authorizationError?: string; -} - -/** - * Context object holding upgrade-time state for the TLS proxy socket. - * Attached as `tlsSocket._ctx` so named functions can reach it via `this._ctx` - * (Duplex methods) or via a bound `this` (socket callbacks). - */ -function UpgradeContext( - connectionListener: (...args: any[]) => any, - server: Http2SecureServer, - rawSocket: import("node:net").Socket, -) { - this.connectionListener = connectionListener; - this.server = server; - this.rawSocket = rawSocket; - this.nativeHandle = null; - this.events = null; -} - -// --------------------------------------------------------------------------- -// Duplex stream methods — called with `this` = tlsSocket (standard stream API) -// --------------------------------------------------------------------------- - -// _read: called by stream machinery when the H2 session wants data. -// Resume the native TLS handle so it feeds decrypted data via the data callback. -// Mirrors net.ts Socket.prototype._read which calls socket.resume(). -function tlsSocketRead(this: TLSProxySocket) { - const h = this._ctx.nativeHandle; - if (h) { - h.resume(); - } - this._ctx.rawSocket.resume(); -} - -// _write: called when the H2 session writes outbound frames. -// Forward to the native TLS handle for encryption, then back to rawSocket. -// Mirrors net.ts Socket.prototype._write which calls socket.$write(). -function tlsSocketWrite(this: TLSProxySocket, chunk: Buffer, encoding: string, callback: (err?: Error | null) => void) { - const h = this._ctx.nativeHandle; - if (!h) { - callback(new Error("Socket is closed")); - return; - } - // $write returns true if fully flushed, false if buffered - if (h.$write(chunk, encoding)) { - callback(); - } else { - // Store callback so drain event can invoke it (backpressure) - this._writeCallback = callback; - } -} - -// _destroy: called when the stream is destroyed (e.g. tlsSocket.destroy(err)). -// Cleans up the native TLS handle. -// Mirrors net.ts Socket.prototype._destroy. -function tlsSocketDestroy(this: TLSProxySocket, err: Error | null, callback: (err?: Error | null) => void) { - const h = this._ctx.nativeHandle; - if (h) { - h.close(); - this._ctx.nativeHandle = null; - } - this._ctx.rawSocket.destroy(); - // Must invoke pending write callback with error per Writable stream contract - const writeCb = this._writeCallback; - if (writeCb) { - this._writeCallback = null; - writeCb(err ?? new Error("Socket destroyed")); - } - callback(err); -} - -// _final: called when the writable side is ending (all data flushed). -// Shuts down the TLS write side gracefully. -// Mirrors net.ts Socket.prototype._final. -function tlsSocketFinal(this: TLSProxySocket, callback: () => void) { - const h = this._ctx.nativeHandle; - if (!h) return callback(); - // Signal end-of-stream to the TLS layer - h.shutdown(); - callback(); -} - -// --------------------------------------------------------------------------- -// Socket callbacks — called natively with `this` = native handle (not useful). -// All are bound to tlsSocket so `this` inside each = tlsSocket. -// --------------------------------------------------------------------------- - -// open: called when the TLS layer is initialized (before handshake). -// No action needed; we wait for the handshake callback. -function socketOpen() {} - -// data: called with decrypted plaintext after the TLS layer decrypts incoming data. -// Push into tlsSocket so the H2 session's _read() receives these frames. -function socketData(this: TLSProxySocket, _socket: NativeHandle, chunk: Buffer) { - if (!this.push(chunk)) { - this._ctx.rawSocket.pause(); - } -} - -// end: TLS peer signaled end-of-stream; signal EOF to the H2 session. -function socketEnd(this: TLSProxySocket) { - this.push(null); -} - -// drain: raw socket is writable again after being full; propagate backpressure signal. -// If _write stored a callback waiting for drain, invoke it now. -function socketDrain(this: TLSProxySocket) { - const cb = this._writeCallback; - if (cb) { - this._writeCallback = null; - cb(); - } -} - -// close: TLS connection closed; tear down the tlsSocket Duplex. -function socketClose(this: TLSProxySocket) { - if (!this.destroyed) { - this.destroy(); - } -} - -// error: TLS-level error (e.g. certificate verification failure). -// In server mode without _requestCert, the server doesn't request a client cert, -// so issuer verification errors on the server's own cert are non-fatal. -function socketError(this: TLSProxySocket, _socket: NativeHandle, err: NodeJS.ErrnoException) { - const ctx = this._ctx; - if (!ctx.server._requestCert && err?.code === "UNABLE_TO_GET_ISSUER_CERT") { - return; - } - this.destroy(err); -} - -// timeout: socket idle timeout; forward to the Duplex so H2 session can handle it. -function socketTimeout(this: TLSProxySocket) { - this.emit("timeout"); -} - -// handshake: TLS handshake completed. This is the critical callback that triggers -// H2 session creation. -// -// Mirrors the handshake logic in net.ts ServerHandlers.handshake: -// - Set secure-connection state flags on tlsSocket -// - Read alpnProtocol from the native handle (set by ALPN negotiation) -// - Handle _requestCert / _rejectUnauthorized for mutual TLS -// - Call connectionListener to create the ServerHttp2Session -function socketHandshake( - this: TLSProxySocket, - nativeHandle: NativeHandle, - success: boolean, - verifyError: NodeJS.ErrnoException | null, -) { - const tlsSocket = this; // bound - const ctx = tlsSocket._ctx; - - if (!success) { - const err = verifyError || new Error("TLS handshake failed"); - ctx.server.emit("tlsClientError", err, tlsSocket); - tlsSocket.destroy(err); - return; - } - - // Mark TLS handshake as complete on the proxy socket - tlsSocket._securePending = false; - tlsSocket.secureConnecting = false; - tlsSocket._secureEstablished = true; - - // Copy the negotiated ALPN protocol (e.g. "h2") from the native TLS handle. - // The H2 session checks this to confirm HTTP/2 was negotiated. - tlsSocket.alpnProtocol = nativeHandle?.alpnProtocol ?? null; - - // Handle mutual TLS: if the server requested a client cert, check for errors - const requestCert = tlsSocket._requestCert; - let rejectUnauthorized; - if (requestCert || (rejectUnauthorized = tlsSocket._rejectUnauthorized)) { - if (verifyError) { - tlsSocket.authorized = false; - tlsSocket.authorizationError = verifyError.code || verifyError.message; - ctx.server.emit("tlsClientError", verifyError, tlsSocket); - if (rejectUnauthorized ?? tlsSocket._rejectUnauthorized) { - tlsSocket.emit("secure", tlsSocket); - tlsSocket.destroy(verifyError); - return; - } - } else if (requestCert) { - tlsSocket.authorized = true; - } - } - - // Invoke the H2 connectionListener which creates a ServerHttp2Session. - // This is the same function passed to Http2SecureServer's constructor - // and is what normally fires on the 'secureConnection' event. - ctx.connectionListener.$call(ctx.server, tlsSocket); - - // Resume the Duplex so the H2 session can read frames from it. The accept - // path in net.ts ServerHandlers.handshake only read(0)s (node's manualStart - // server socket); the H2 session is the consumer here, so it flows. - tlsSocket.resume(); -} - -// --------------------------------------------------------------------------- -// Close-cleanup handler -// --------------------------------------------------------------------------- - -// onTlsClose: when the TLS socket closes (e.g. H2 session destroyed), clean up -// the raw socket listeners to prevent memory leaks and stale callback references. -// EventEmitter calls 'close' handlers with `this` = emitter (tlsSocket). -function onTlsClose(this: TLSProxySocket) { - const ctx = this._ctx; - const raw = ctx.rawSocket; - const ev = ctx.events; - if (!ev) return; - raw.removeListener("data", ev[0]); - raw.removeListener("end", ev[1]); - raw.removeListener("drain", ev[2]); - raw.removeListener("close", ev[3]); -} - -// --------------------------------------------------------------------------- -// Module-scope noop (replaces anonymous () => {} for the error suppression) -// --------------------------------------------------------------------------- - -// no-op handler used to suppress unhandled error events until -// the H2 session attaches its own error handler. -function noop() {} - -// --------------------------------------------------------------------------- -// Main upgrade function -// --------------------------------------------------------------------------- - -// Upgrades a raw TCP socket to TLS and initiates an H2 session on it. -// -// When a net.Server forwards an accepted TCP connection to an Http2SecureServer -// via `h2Server.emit('connection', socket)`, the socket has not been TLS-upgraded. -// Node.js Http2SecureServer expects to receive this and perform the upgrade itself. -// -// This mirrors the TLS server handshake pattern from net.ts ServerHandlers, but -// targets the H2 connectionListener instead of a generic secureConnection event. -// -// Data flow after upgrade: -// rawSocket (TCP) → upgradeDuplexToTLS (native TLS layer) → socket callbacks -// → tlsSocket.push() → H2 session reads -// H2 session writes → tlsSocket._write() → handle.$write() → native TLS layer → rawSocket -// -// CRITICAL: We do NOT set tlsSocket._handle to the native TLS handle. -// If we did, the H2FrameParser constructor would detect it as a JSTLSSocket -// and call attachNativeCallback(), which intercepts all decrypted data at the -// native level, completely bypassing our JS data callback and Duplex.push() path. -// Instead, we store the handle in _ctx.nativeHandle so _read/_write/_destroy -// can use it, while the H2 session sees _handle as null and uses the JS-level -// socket.on("data") → Duplex → parser.read() path for incoming frames. -function upgradeRawSocketToH2( - connectionListener: (...args: any[]) => any, - server: Http2SecureServer, - rawSocket: import("node:net").Socket, -): boolean { - // Create a Duplex stream that acts as the TLS "socket" from the H2 session's perspective. - const tlsSocket = new Duplex() as TLSProxySocket; - tlsSocket._ctx = new UpgradeContext(connectionListener, server, rawSocket); - - // Duplex stream methods — `this` is tlsSocket, no bind needed - tlsSocket._read = tlsSocketRead; - tlsSocket._write = tlsSocketWrite; - tlsSocket._destroy = tlsSocketDestroy; - tlsSocket._final = tlsSocketFinal; - - // Suppress unhandled error events until the H2 session attaches its own error handler - tlsSocket.on("error", noop); - - // Set TLS-like properties that connectionListener and the H2 session expect. - // These are set on the Duplex because we cannot use a real TLSSocket here — - // its internal state machine would conflict with upgradeDuplexToTLS. - tlsSocket.alpnProtocol = null; - tlsSocket.authorized = false; - tlsSocket.encrypted = true; - tlsSocket.server = server; - - // Only enforce client cert verification if the server explicitly requests it. - // tls.Server defaults _rejectUnauthorized to true, but without _requestCert - // the server doesn't actually ask for a client cert, so verification errors - // (e.g. UNABLE_TO_GET_ISSUER_CERT for the server's own self-signed cert) are - // spurious and must be ignored. - tlsSocket._requestCert = server._requestCert || false; - tlsSocket._rejectUnauthorized = server._requestCert ? server._rejectUnauthorized : false; - - // socket: callbacks — bind to tlsSocket since they are invoked with the native handle as `this` - let handle: NativeHandle, events: UpgradeEvents; - try { - // upgradeDuplexToTLS wraps rawSocket with a TLS layer in server mode (isServer: true). - // The native side will: - // 1. Read encrypted data from rawSocket via events[0..3] - // 2. Decrypt it through the TLS engine (with ALPN negotiation for "h2") - // 3. Call our socket callbacks below with the decrypted plaintext - // - // ALPNProtocols: server.ALPNProtocols is a Buffer in wire format (e.g. - // for ["h2"]). The native SSLConfig expects an ArrayBuffer, so we slice the underlying buffer. - [handle, events] = upgradeDuplexToTLS(rawSocket, { - isServer: true, - tls: { - secureContext: server[kSharedCreds]().context, - requestCert: server._requestCert, - rejectUnauthorized: server._rejectUnauthorized, - ALPNProtocols: server.ALPNProtocols - ? server.ALPNProtocols.buffer.slice( - server.ALPNProtocols.byteOffset, - server.ALPNProtocols.byteOffset + server.ALPNProtocols.byteLength, - ) - : null, - }, - socket: { - open: socketOpen, - data: socketData.bind(tlsSocket), - end: socketEnd.bind(tlsSocket), - drain: socketDrain.bind(tlsSocket), - close: socketClose.bind(tlsSocket), - error: socketError.bind(tlsSocket), - timeout: socketTimeout.bind(tlsSocket), - handshake: socketHandshake.bind(tlsSocket), - }, - data: {}, - }); - } catch (e) { - rawSocket.destroy(e as Error); - tlsSocket.destroy(e as Error); - return true; - } - - // Store handle in _ctx (NOT on tlsSocket._handle). - // This prevents H2FrameParser from attaching as native callback which would - // intercept data at the native level and bypass our Duplex push path. - tlsSocket._ctx.nativeHandle = handle; - tlsSocket._ctx.events = events; - - // Wire up the raw TCP socket to feed encrypted data into the TLS layer. - // events[0..3] are native event handlers returned by upgradeDuplexToTLS that - // the native TLS engine expects to receive data/end/drain/close through. - rawSocket.on("data", events[0]); - rawSocket.on("end", events[1]); - rawSocket.on("drain", events[2]); - rawSocket.on("close", events[3]); - - // When the TLS socket closes (e.g. H2 session destroyed), clean up the raw socket - // listeners to prevent memory leaks and stale callback references. - // EventEmitter calls 'close' handlers with `this` = emitter (tlsSocket). - tlsSocket.once("close", onTlsClose); - return true; -} - -export default { upgradeRawSocketToH2 }; diff --git a/src/js/node/http2.ts b/src/js/node/http2.ts index f41ce68605fc..1d95e7e38840 100644 --- a/src/js/node/http2.ts +++ b/src/js/node/http2.ts @@ -94,8 +94,6 @@ const StringPrototypeStartsWith = String.prototype.startsWith; const ObjectPrototypeHasOwnProperty = Object.prototype.hasOwnProperty; const H2FrameParser = $rust("h2_frame_parser.rs", "H2FrameParserConstructor"); -const { upgradeRawSocketToH2 } = require("node:_http2_upgrade"); -type UpgradableSecureServer = Parameters[1]; const kSettingIds: Record = { 0x1: "headerTableSize", @@ -444,11 +442,22 @@ function emitEventNT(self: any, event: string, ...args: any[]) { // The frame is passed in: destroy() clears it off the session before emitting // 'error', so a throwing listener on either event cannot leave a retained // session pinning the store. -function emitSessionCloseNT(self: Http2Session, frame) { +function emitSessionCloseNT(self: Http2Session, error: Error | null | undefined, frame) { + if (error) { + runInFrame(frame, self.emit, self, "error", error); + } if (self.listenerCount("close") > 0) { runInFrame(frame, self.emit, self, "close"); } } +// Node's emitClose: the session reports 'error' and 'close' once its socket has closed. +function emitSessionCloseAfterSocket(self: Http2Session, socket, error: Error | null | undefined, frame) { + if (socket && !socket.destroyed) { + socket.once("close", () => emitSessionCloseNT(self, error, frame)); + } else { + process.nextTick(emitSessionCloseNT, self, error, frame); + } +} function emitErrorNT(self: any, error: any, destroy: boolean) { if (destroy) { if (self.listenerCount("error") > 0) { @@ -4750,7 +4759,8 @@ class ServerHttp2Session extends Http2Session { return; } this.#destroying = true; - emitHttp2SessionPerf(this, this.#parser, this[bunHTTP2Socket]); + const socket = this[bunHTTP2Socket]; + emitHttp2SessionPerf(this, this.#parser, socket); try { const server = this[kServer]; if (server) { @@ -4774,7 +4784,6 @@ class ServerHttp2Session extends Http2Session { this[kSessionDestroyError] = error; } - const socket = this[bunHTTP2Socket]; if (!this.#connected) return; this.#closed = true; this.#connected = false; @@ -4785,24 +4794,7 @@ class ServerHttp2Session extends Http2Session { // a destroy(err) after close() must still put the error GOAWAY on the wire. this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); } - if (error) { - // node's finishSessionClose destroys the socket when the session dies - // with an error (a misbehaving peer must observe the connection going - // away) - but it still ends first and destroys a tick later, so the - // final GOAWAY flushes behind a FIN instead of an abortive close (see - // endThenDestroySessionSocket). - endThenDestroySessionSocket(socket, error); - } else { - // Node's finishSessionClose: "If we're gracefully closing the socket, - // call resume() so we can detect the peer closing in case - // binding.Http2Session is already gone." Without a reader, unread - // inbound bytes (a late GOAWAY from the peer) turn the close into an - // RST, which the peer surfaces as read ECONNRESET (routine on Windows - // loopback - the same reason Node delays the error-path destroy). - // https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 - socket.resume(); - socket.end(); - } + closeSessionSocket(socket, this.#closeCalled && !error, error); } const parser = this.#parser; if (parser) { @@ -4828,16 +4820,9 @@ class ServerHttp2Session extends Http2Session { } this[bunHTTP2Socket] = null; - // Read-and-clear the frame first: emitting 'error' with no listener throws, - // which would skip the clear and leave a retained session pinning the store. const asyncFrame = this[bunHTTP2AsyncContextFrame]; this[bunHTTP2AsyncContextFrame] = undefined; - if (error) { - runInFrame(asyncFrame, this.emit, this, "error", error); - } - // node emits the session 'close' event asynchronously (a listener attached right after - // close()/destroy() returns must still observe it). - process.nextTick(emitSessionCloseNT, this, asyncFrame); + emitSessionCloseAfterSocket(this, socket, error, asyncFrame); } } function emitTimeout(session: ClientHttp2Session) { @@ -4880,19 +4865,21 @@ function setSessionTimeout(this: Http2Session, msecs, callback) { return this; } -// Node's finishSessionClose error path: socket.end() flushes and sends the FIN -// first, and the hard destroy runs a tick later - "If session.destroy() was -// called, destroy the underlying socket. Delay it a bit to try to avoid -// ECONNRESET on Windows" - so the peer reads our final frames off a FIN'd -// socket instead of observing an abortive close. -// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 function destroySessionSocketDelayedNT(socket, error) { if (!socket.destroyed) { socket.destroy(error); } } -function endThenDestroySessionSocket(socket, error) { - socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error)); +// Node's finishSessionClose: end(), then destroy() unless close() asked for a graceful shutdown. +function closeSessionSocket(socket, graceful: boolean, error: Error | null | undefined) { + if (socket.destroyed) return; + if (graceful) { + // Unread inbound bytes would turn the close into an RST. + socket.resume(); + socket.end(); + } else { + socket.end(() => setImmediate(destroySessionSocketDelayedNT, socket, error)); + } } // node callTimeout (lib/internal/http2/core.js): when the timer expires while writes are still in // flight and bytes have reached the wire since the previous expiry, the session is not idle — @@ -5896,18 +5883,7 @@ class ClientHttp2Session extends Http2Session { // a destroy(err) after close() must still put the error GOAWAY on the wire. this.goaway(code || constants.NGHTTP2_NO_ERROR, 0, Buffer.alloc(0)); } - if (error) { - // See the client session: end first, destroy a tick later (node's - // finishSessionClose Windows-ECONNRESET avoidance). - endThenDestroySessionSocket(socket, error); - } else { - // See the client session's destroy: Node's finishSessionClose resumes - // the socket on a graceful close so unread inbound bytes cannot turn - // the FIN teardown into an RST. - // https://github.com/nodejs/node/blob/v26.3.0/lib/internal/http2/core.js#L1188 - socket.resume(); - socket.end(); - } + closeSessionSocket(socket, this.#closeCalled && !error, error); } const parser = this.#parser; if (parser) { @@ -5937,16 +5913,9 @@ class ClientHttp2Session extends Http2Session { this.#parser = null; this[bunHTTP2Socket] = null; - // Read-and-clear the frame first: emitting 'error' with no listener throws, - // which would skip the clear and leave a retained session pinning the store. const asyncFrame = this[bunHTTP2AsyncContextFrame]; this[bunHTTP2AsyncContextFrame] = undefined; - if (error) { - runInFrame(asyncFrame, this.emit, this, "error", error); - } - // node emits the session 'close' event asynchronously (a listener attached right after - // close()/destroy() returns must still observe it). - process.nextTick(emitSessionCloseNT, this, asyncFrame); + emitSessionCloseAfterSocket(this, socket, error, asyncFrame); } request(headers?: HeadersObject | any[] | null, options?: ClientRequestOptions) { @@ -6639,7 +6608,6 @@ function onErrorSecureServerSession(err, socket) { function emitFrameErrorEventNT(stream, frameType, errorCode) { stream.emit("frameError", frameType, errorCode); } -interface Http2SecureServer extends UpgradableSecureServer {} class Http2SecureServer extends (tls.Server as unknown as Http2SecureServerBase) { declare keepAliveTimeout: number | undefined; declare headersTimeout: number | undefined; @@ -6685,15 +6653,6 @@ class Http2SecureServer extends (tls.Server as unknown as Http2SecureServerBase) } this.on("tlsClientError", onErrorSecureServerSession); } - emit(event: string, ...args: any[]) { - if (event === "connection") { - const socket = args[0]; - if (socket && !(socket instanceof TLSSocket)) { - return upgradeRawSocketToH2(connectionListener, this, socket); - } - } - return super.emit(event, ...args); - } setTimeout(ms, callback) { if (callback !== undefined && typeof callback !== "function") { throw $ERR_INVALID_ARG_TYPE("callback", "function", callback); diff --git a/src/js/node/net.ts b/src/js/node/net.ts index d864d281ea2e..789ad827dd2f 100644 --- a/src/js/node/net.ts +++ b/src/js/node/net.ts @@ -299,6 +299,7 @@ const kOnUpgradedClose = Symbol("kOnUpgradedClose"); const kupgraded = Symbol("kupgraded"); // On the raw handle of an adopted fd: the TLS socket that adopted it. const kAdoptedTLSRaw = Symbol("kAdoptedTLSRaw"); +const kAdoptedTLSTransport = Symbol("kAdoptedTLSTransport"); // On that TLS socket: the fd closed, and the socket it wraps waits to be closed with it. const kOwesRawClose = Symbol("kOwesRawClose"); const ksocket = Symbol("ksocket"); @@ -390,12 +391,23 @@ function writeErrnoException(negErrno) { } return er; } -function endNT(socket, callback, err) { +function endNT(socket, callback, err, transport) { // Node's _final half-closes the writable side (sends FIN) and leaves the // readable side open; the Duplex's allowHalfOpen drives the eventual destroy. // https://github.com/nodejs/node/blob/614050b657e9757c1097aa85f92f2cb51149dc0d/lib/net.js#L500 - socket.shutdown(); - callback(err); + if (transport) { + const eos = require("internal/streams/end-of-stream"); + const cleanup = eos(transport, { readable: false }, () => { + cleanup(); + // Node's afterShutdown completes successfully; the transport forwards its error separately. + // https://github.com/nodejs/node/blob/v24.21.0/lib/net.js#L774-L779 + callback(err); + }); + socket.shutdown(); + } else { + socket.shutdown(); + callback(err); + } } function emitCloseNT(self, hasError) { self.emit("close", hasError); @@ -438,12 +450,16 @@ function destroyNT(self, err) { } // Node's wrap 'close' -> destroy(): https://github.com/nodejs/node/blob/v26.3.0/lib/internal/tls/wrap.js#L739-L741 function onUpgradedClose(self, connection) { - if (self[kupgraded] !== connection) return; + if (self[kupgraded] !== connection || self.destroyed) return; // The stream-level engine reads its transport with no backpressure, so the // transport can close after the peer's EOF with plaintext still unread. if ((self[kended] || self[kOnreadPendingEnd]) && !self.readableEnded) self.once("end", self[kOnUpgradedClose]); else self.destroy(); } +function closeDuplexTLSOwner(self) { + const connection = self[kupgraded]; + if (connection && !self[kAdoptedTLSTransport]) onUpgradedClose(self, connection); +} // Armed ahead of the stream-level engine's own 'close' thunk: that thunk aborts // a pending handshake, which a socket destroyed first does not report. function destroyWhenUpgradedCloses(self, connection) { @@ -687,6 +703,7 @@ const SocketHandlers = { //socket cannot be used after close detachSocket(self); if (!closeWithTLSSocket(self, socket)) SocketEmitEndNT(self, err); + closeDuplexTLSOwner(self); self.data = null; }, data(socket, buffer) { @@ -848,7 +865,7 @@ function unrefAfterDrain(self, handle) { } function finishSocketEnd(self) { - if (self[kended]) return; + if (self[kended] || self.destroyed) return; self[kended] = true; if (!self.allowHalfOpen) self.write = writeAfterFIN; self.push(null); @@ -1121,6 +1138,7 @@ const ServerHandlers = { //socket cannot be used after close detachSocket(data); if (!closeWithTLSSocket(data, socket)) SocketEmitEndNT(data, err); + closeDuplexTLSOwner(data); data.data = null; socket[owner_symbol] = null; } @@ -1626,6 +1644,7 @@ const SocketHandlers2 = { self[kwriteCallback] = null; pendingWrite($ERR_SOCKET_CLOSED()); } + closeDuplexTLSOwner(self); }, handshake(socket, success, verifyError) { $debug("Bun.Socket handshake"); @@ -1824,6 +1843,7 @@ function Socket(options?): void { this._parent = null; this._parentWrap = null; this[kupgraded] = null; + this[kAdoptedTLSTransport] = false; this[kStandaloneWrap] = false; this[kOnUpgradedClose] = undefined; this[kOwesRawClose] = false; @@ -2244,6 +2264,8 @@ Socket.prototype.connect = function connect(...args) { // https://github.com/nodejs/node/blob/c5cfdd48497fe9bd8dbd55fd1fca84b321f48ec1/lib/net.js#L311 // https://github.com/nodejs/node/blob/c5cfdd48497fe9bd8dbd55fd1fca84b321f48ec1/lib/net.js#L1126 this._undestroy(); + this[kupgraded] = connection; + this[kAdoptedTLSTransport] = false; const socket = connection._handle; if (!upgradeDuplex && socket) { // if is named pipe socket we can upgrade it using the same wrapper than we use for duplex @@ -2281,8 +2303,8 @@ Socket.prototype.connect = function connect(...args) { // replace socket connection._handle = raw; raw[kAdoptedTLSRaw] = this; + this[kAdoptedTLSTransport] = true; destroyWhenUpgradedCloses(this, connection); - this.once("end", this[kCloseRawConnection]); raw.connecting = false; this._handle = tls; } else { @@ -2332,8 +2354,8 @@ Socket.prototype.connect = function connect(...args) { // replace socket connection._handle = raw; raw[kAdoptedTLSRaw] = this; + this[kAdoptedTLSTransport] = true; destroyWhenUpgradedCloses(this, connection); - this.once("end", this[kCloseRawConnection]); raw.connecting = false; this._handle = tls; } else { @@ -2417,7 +2439,7 @@ Socket.prototype._destroy = function _destroy(err, callback) { // Stream-level TLS owns its transport; an adopted fd pair closes through closeOwedRaw. // https://github.com/nodejs/node/blob/v24.21.0/lib/internal/js_stream_socket.js#L253 const upgraded = this[kupgraded]; - if (upgraded && !upgraded.destroyed && (!(upgraded instanceof Socket) || !upgraded._handle?.[kAdoptedTLSRaw])) { + if (upgraded && !this[kAdoptedTLSTransport] && !upgraded.destroyed) { upgraded.destroy?.(); } @@ -2513,7 +2535,7 @@ Socket.prototype._final = function _final(callback) { if (!socket) return callback(); // emit FIN allowHalfOpen only allow the readable side to close first - process.nextTick(endNT, socket, callback); + process.nextTick(endNT, socket, callback, undefined, this[kAdoptedTLSTransport] ? undefined : this[kupgraded]); }; Object.defineProperty(Socket.prototype, "localAddress", { @@ -2692,8 +2714,8 @@ Socket.prototype[Symbol.for("::bunUpgradeServerTLS::")] = function (connection, const [raw, tlsHandle] = result; connection._handle = raw; raw[kAdoptedTLSRaw] = this; + this[kAdoptedTLSTransport] = true; destroyWhenUpgradedCloses(this, connection); - this.once("end", this[kCloseRawConnection]); raw.connecting = false; this._handle = tlsHandle; // Match Node's initRead(): an injected socket may be paused, but TLS still diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index cdf01416517c..493269c7a98b 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -62,9 +62,6 @@ pub(crate) struct UpgradedDuplex { pub pending_close: Cell, /// [`Self::shutdown`] arrived before the engine existed. [`Self::drain_pending`] replays it. pub pending_shutdown: Cell, - /// The transport delivered EOF (its 'end' event fired). Teardown payloads - /// (close_notify) are dropped after this; see [`Self::call_write_or_end`]. - pub transport_eof: Cell, } bun_event_loop::impl_timer_owner!(UpgradedDuplex; from_timer_ptr => event_loop_timer); @@ -212,6 +209,12 @@ impl UpgradedDuplex { unsafe { &*this }.finish_close(); } + fn on_end(this: *mut Self) { + // SAFETY: see handler note above. + let this = unsafe { &*this }; + (this.handlers.on_end)(this.handlers.ctx); + } + pub(super) fn finish_close(&self) { // Keep the wrapper (and so its visited `duplex*` slots) reachable // across `handlers.on_close`, which downgrades the socket's own strong @@ -226,7 +229,7 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); } - fn call_write_or_end(&self, data: Option<&[u8]>, msg_more: bool) { + fn call_write_or_end(&self, data: Option<&[u8]>) { // No JS duplex to talk to: the zeroed placeholder, or the owning // socket's finalizer abandoned it (`abandon_js_side`). let duplex = self.origin.get(); @@ -236,33 +239,7 @@ impl UpgradedDuplex { // global is set in `from()` whenever origin is set. let Some(global) = self.global else { return }; - // Teardown-phase bytes (close_notify / the trailing end()) aimed at a - // duplex whose write side already ended (TLS-inception teardown) only - // surface a spurious EPIPE - drop them. Ordinary data writes skip the - // probe so write-after-end still errors like node. - let teardown = data.is_none() || self.wrapper_ref().is_some_and(|w| w.is_shutdown()); - if teardown { - // A teardown payload (close_notify) after the transport's readable - // side ended has no reader behind it: node writes nothing there, - // and a transport that forwards into an auto-ended net.Socket - // throws writeAfterFIN (EPIPE). The trailing end() is not a write - // and still goes through the writableEnded probe below, so a - // half-open transport sees our FIN. - if data.is_some() && self.transport_eof.get() { - return; - } - // Node ends no destroyed stream. - for property in ["writableEnded", "destroyed"] { - match duplex.get(&global, property) { - Ok(Some(done)) if done.to_boolean() => return, - Ok(_) => {} - // Best-effort probe: consume the exception and fall through. - Err(err) => drop(global.take_exception(err)), - } - } - } - - let name = if msg_more { "write" } else { "end" }; + let name = if data.is_some() { "write" } else { "end" }; let write_or_end = match duplex.get(&global, name) { Ok(Some(f)) if f.is_callable() => f, Ok(_) => return, @@ -305,7 +282,7 @@ impl UpgradedDuplex { // Scenario 2: will not write if a exception is thrown (will be handled by onError) // Scenario 3: will be queued in memory and will be flushed later // Scenario 4: no write/end function exists (will be handled by onError) - self.call_write_or_end(Some(encoded_data), true); + self.call_write_or_end(Some(encoded_data)); } #[uws_callback(export = "UpgradedDuplex__flush")] @@ -415,7 +392,6 @@ impl UpgradedDuplex { pending_data: JsCell::new(Vec::new()), pending_close: Cell::new(false), pending_shutdown: Cell::new(false), - transport_eof: Cell::new(false), } } @@ -484,6 +460,7 @@ impl UpgradedDuplex { on_handshake: Self::on_handshake, on_data: Self::on_data, on_close: Self::on_close, + on_end: Some(Self::on_end), write: Self::internal_write, on_session: Some(Self::on_session), on_keylog: Some(Self::on_keylog), @@ -566,7 +543,20 @@ impl UpgradedDuplex { return; }; let _ = w.shutdown(false); - self.call_write_or_end(None, false); + if self.origin.get().is_empty() { + return; + } + let Some(global) = self.global else { return }; + let callback = bun_jsc::JSFunction::create( + &global, + "", + __jsc_host_end_transport, + 1, + Default::default(), + ); + if let Err(err) = JSValue::call_next_tick_1(callback, &global, self.js_wrapper) { + (self.handlers.on_error)(self.handlers.ctx, global.take_error(err)); + } } #[uws_callback(export = "UpgradedDuplex__shutdown_read")] @@ -579,7 +569,7 @@ impl UpgradedDuplex { /// `None` means `start_tls` has not run yet (teardown never clears the slot), not shut down. #[uws_callback(export = "UpgradedDuplex__is_shutdown", no_catch)] pub(crate) fn is_shutdown(&self) -> bool { - self.wrapper_ref().is_some_and(|w| w.is_shutdown()) + self.wrapper_ref().is_some_and(|w| w.is_write_shutdown()) } /// See [`Self::is_shutdown`] for the not-yet-started case. @@ -695,7 +685,6 @@ impl UpgradedDuplex { self.pending_data.set(Vec::new()); self.pending_close.set(false); self.pending_shutdown.set(false); - self.transport_eof.set(false); } } @@ -705,6 +694,22 @@ impl Drop for UpgradedDuplex { } } +// Node's JSStreamSocket.doShutdown ends the transport on the next tick, after its TLS writes. +// https://github.com/nodejs/node/blob/v24.21.0/lib/internal/js_stream_socket.js#L156-L161 +#[bun_jsc::host_fn] +fn end_transport(_global: &JSGlobalObject, frame: &CallFrame) -> JsResult { + let [owner] = frame.arguments_as_array::<1>(); + if let Some(socket) = js_TLSSocket::from_js(owner) { + // SAFETY: the queued argument roots the JS owner; close detaches its native handle. + let socket = unsafe { &*socket.as_ptr() }; + if let bun_uws::InternalSocket::UpgradedDuplex(duplex) = socket.socket.get().socket { + // SAFETY: this live handle was installed by DuplexUpgradeContext::duplex_socket. + unsafe { &*duplex.cast::() }.call_write_or_end(None); + } + } + Ok(JSValue::UNDEFINED) +} + // SAFETY (all four host fns): the function data is the `*mut UpgradedDuplex` // installed by `get_js_handlers`; `teardown` clears it before the storage is // freed, so a non-null data pointer is live for the call. @@ -754,7 +759,6 @@ fn on_end(_global: &JSGlobalObject, frame: &CallFrame) -> JsResult { // SAFETY: see host-fn note above. let this = unsafe { &*self_ptr.cast::() }; - this.transport_eof.set(true); // Like node's JSStreamSocket. Ahead of staged bytes too: no handshake can complete after it. (this.handlers.on_end)(this.handlers.ctx); } diff --git a/src/runtime/socket/WindowsNamedPipe.rs b/src/runtime/socket/WindowsNamedPipe.rs index 6c892b05f2a8..2b48113f4dfb 100644 --- a/src/runtime/socket/WindowsNamedPipe.rs +++ b/src/runtime/socket/WindowsNamedPipe.rs @@ -382,6 +382,7 @@ impl WindowsNamedPipe { on_handshake: Self::ssl_on_handshake, on_data: Self::ssl_on_data, on_close: Self::ssl_on_close, + on_end: None, write: Self::ssl_write, on_session: Some(Self::ssl_on_session), on_keylog: Some(Self::ssl_on_keylog), diff --git a/src/uws/lib.rs b/src/uws/lib.rs index 2e88b7eeff56..00921897cf83 100644 --- a/src/uws/lib.rs +++ b/src/uws/lib.rs @@ -406,6 +406,8 @@ pub mod ssl_wrapper { pub write: fn(T, &[u8]), pub on_data: fn(T, &[u8]), pub on_close: fn(T), + /// Stream owners handle peer EOF independently of their writable half. + pub on_end: Option, /// A new resumable TLS session arrived (serialized SSL_SESSION bytes) /// - node's `'session'` event. `None` opts the SSL out of session /// parking entirely (fetch / WebSocket tunnels have no consumer). @@ -838,6 +840,12 @@ pub mod ssl_wrapper { || self.flags.sent_ssl_shutdown() } + pub fn is_write_shutdown(&self) -> bool { + self.flags.closed_notified() + || self.flags.fatal_error() + || self.flags.sent_ssl_shutdown() + } + /// We sent and received the shutdown (fully closed) pub fn is_closed(&self) -> bool { self.flags.received_ssl_shutdown() && self.flags.sent_ssl_shutdown() @@ -983,6 +991,21 @@ pub mod ssl_wrapper { (handlers.on_close)(handlers.ctx); } + fn handle_peer_shutdown(&self) { + let handlers = self.handlers.get(); + if let Some(on_end) = handlers.on_end { + if self.flags.received_ssl_shutdown() { + return; + } + self.flags.set_received_ssl_shutdown(true); + on_end(handlers.ctx); + } else { + self.flags.set_received_ssl_shutdown(true); + let _ = self.shutdown(false); + self.trigger_close_callback(); + } + } + /// The SSL's X509 verdict. Shutdown state does not change it. fn verify_error(&self) -> us_bun_verify_error_t { let Some(ssl) = self.ssl.get() else { @@ -1013,11 +1036,7 @@ pub mod ssl_wrapper { & boring_sys::SSL_RECEIVED_SHUTDOWN) != 0 { - // we received a shutdown - self.flags.set_received_ssl_shutdown(true); - // 2-step shutdown - let _ = self.shutdown(false); - self.trigger_close_callback(); + self.handle_peer_shutdown(); return false; } @@ -1065,10 +1084,14 @@ pub mod ssl_wrapper { if err == boring_sys::SSL_ERROR_ZERO_RETURN { // Remotely-Initiated Shutdown // See: https://www.openssl.org/docs/manmaster/man3/SSL_shutdown.html - self.flags.set_received_ssl_shutdown(true); - // 2-step shutdown - let _ = self.shutdown(false); - self.handle_end_of_renegotiation(); + if self.handlers.get().on_end.is_some() { + self.handle_end_of_renegotiation(); + self.handle_peer_shutdown(); + } else { + self.flags.set_received_ssl_shutdown(true); + let _ = self.shutdown(false); + self.handle_end_of_renegotiation(); + } return false; } // as far as I know these are the only errors we want to handle @@ -1187,7 +1210,9 @@ pub mod ssl_wrapper { } else if err == boring_sys::SSL_ERROR_ZERO_RETURN { // Remotely-Initiated Shutdown // See: https://www.openssl.org/docs/manmaster/man3/SSL_shutdown.html - self.flags.set_received_ssl_shutdown(true); + if self.handlers.get().on_end.is_none() { + self.flags.set_received_ssl_shutdown(true); + } self.handle_end_of_renegotiation(); } if err == boring_sys::SSL_ERROR_SSL || err == boring_sys::SSL_ERROR_SYSCALL @@ -1213,10 +1238,10 @@ pub mod ssl_wrapper { return false; } if err == boring_sys::SSL_ERROR_ZERO_RETURN { - // 2-step shutdown, last: write_data fails once our close_notify is out. - let _ = self.shutdown(false); + self.handle_peer_shutdown(); + } else { + self.trigger_close_callback(); } - self.trigger_close_callback(); return false; } else { log!("wanna read/write just break"); diff --git a/test/js/node/http2/h2-conformance.test.ts b/test/js/node/http2/h2-conformance.test.ts index b121dbb06374..e3e2b168dfc3 100644 --- a/test/js/node/http2/h2-conformance.test.ts +++ b/test/js/node/http2/h2-conformance.test.ts @@ -1944,8 +1944,12 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou function respondingServer(options: Record = {}, rejectUploads = false) { const state = { handlers: 0, sessionErrorCode: undefined as string | undefined }; + const sessionError = Promise.withResolvers(); const server = http2.createServer(options); - server.on("sessionError", (e: any) => (state.sessionErrorCode = e.code)); + server.on("sessionError", (e: any) => { + state.sessionErrorCode = e.code; + sessionError.resolve(e.code); + }); server.on("session", s => s.on("error", () => {})); server.on("stream", (stream: any, headers: any) => { state.handlers++; @@ -1954,7 +1958,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou stream.respond({ ":status": 200 }); stream.end("x"); }); - return { server, state }; + return { server, state, sessionError: sessionError.promise }; } async function withClient(server: http2.Http2Server, body: (c: RawH2) => Promise): Promise { @@ -1972,10 +1976,11 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou } async function flood(opts: { options?: Record; count: number; kill?: (sid: number) => Buffer }) { - const { server, state } = respondingServer(opts.options); + const { server, state, sessionError } = respondingServer(opts.options); return withClient(server, async c => { c.send(pairs(opts.count, 1, opts.kill ?? rstStream)); const goaway = await c.waitForGoaway(10_000); + await sessionError; return { c, goaway, ...state }; }); } @@ -2143,20 +2148,21 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou // nghttp2 can exempt every reset after its GOAWAY because it ignores the streams a client // opens after that frame. This engine still opens them, so only the earlier streams are exempt. // It also accepts ids below the held stream's, so the GOAWAY's last stream id is no criterion. - const { server, state } = respondingServer(); + const { server, state, sessionError } = respondingServer(); await withClient(server, async c => { c.send(upload(4001)); await serverGoaway(c); c.send(pairs(1200, 1, rstStream)); const goaway = await c.waitFor(calm, 10_000); expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets"); + await sessionError; expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); }); test("server-sent resets stay charged after the server has sent its GOAWAY", async () => { // Every stream here is older than the GOAWAY. The exemption is for the peer's own resets only. - const { server, state } = respondingServer(); + const { server, state, sessionError } = respondingServer(); await withClient(server, async c => { const ids = Array.from({ length: 1300 }, (_, i) => 1 + 2 * i); c.send(Buffer.concat(ids.map(upload))); @@ -2164,6 +2170,7 @@ describe("stream-reset floods (CVE-2023-44487 rapid reset, CVE-2025-8671 MadeYou c.send(Buffer.concat(ids.map(madeYouReset["WINDOW_UPDATE with a 0 increment"]))); const goaway = await c.waitFor(calm, 10_000); expect(goaway.payload.subarray(8).toString()).toBe("too many stream resets"); + await sessionError; expect(state.sessionErrorCode).toBe("ERR_HTTP2_ERROR"); }); }); diff --git a/test/js/node/http2/node-http2-upgrade.test.mts b/test/js/node/http2/node-http2-upgrade.test.mts index a400af9594af..603f683c99c4 100644 --- a/test/js/node/http2/node-http2-upgrade.test.mts +++ b/test/js/node/http2/node-http2-upgrade.test.mts @@ -1,6 +1,6 @@ /** * Tests for the net.Server → Http2SecureServer upgrade path - * (upgradeRawSocketToH2 in _http2_upgrade.ts). + * through the inherited tls.Server connection listener. * * This pattern is used by http2-wrapper, crawlee, and other libraries that * accept raw TCP connections and upgrade them to HTTP/2 via @@ -468,16 +468,30 @@ describe("HTTP/2 upgrade — session destruction releases the accepted transport for (const withError of [false, true]) { test(`destroyed ${withError ? "with" : "without"} an error`, async () => { const h2Server = http2.createSecureServer(TLS); + const error = withError ? new Error("test teardown") : undefined; + const errors: Error[] = []; + const events: string[] = []; const sessionClosed = new Promise(resolve => { h2Server.once("session", (session: http2.ServerHttp2Session) => { - session.on("error", () => {}); - session.once("close", resolve); - session.destroy(withError ? new Error("test teardown") : undefined); + session.on("error", error => { + errors.push(error); + events.push("session error"); + }); + session.once("close", () => { + events.push("session close"); + resolve(); + }); + session.destroy(error); }); }); const accepted = Promise.withResolvers<{ raw: net.Socket; closed: Promise }>(); const netServer = net.createServer(raw => { - const closed = new Promise(resolve => raw.once("close", resolve)); + const closed = new Promise(resolve => + raw.once("close", hadError => { + events.push("raw close"); + resolve(hadError); + }), + ); accepted.resolve({ raw, closed }); h2Server.emit("connection", raw); }); @@ -508,6 +522,11 @@ describe("HTTP/2 upgrade — session destruction releases the accepted transport const { raw, closed } = await accepted.promise; assert.strictEqual(await closed, false); assert.strictEqual(raw.destroyed, true); + assert.deepStrictEqual( + events, + withError ? ["raw close", "session error", "session close"] : ["raw close", "session close"], + ); + assert.deepStrictEqual(errors, withError ? [error] : []); const connections = await new Promise((resolve, reject) => { netServer.getConnections((err, count) => (err ? reject(err) : resolve(count))); }); diff --git a/test/js/node/tls/node-tls-connect.test.ts b/test/js/node/tls/node-tls-connect.test.ts index 27a5ec9540b5..30a9b76f0ebc 100644 --- a/test/js/node/tls/node-tls-connect.test.ts +++ b/test/js/node/tls/node-tls-connect.test.ts @@ -857,6 +857,24 @@ it("a client and a server TLSSocket connected through a synchronous in-memory du }); }); +it.concurrent.each(["legacy-pair", "duplex-halfopen-false", "duplex-halfopen-true", "duplex-eof-halfopen-true"])( + "duplex TLS close_notify follows the transport's half-open policy: %s", + async scenario => { + await using proc = Bun.spawn({ + cmd: [bunExe(), join(import.meta.dirname, "tls-half-close-transport-fixture.cjs"), scenario], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ result: JSON.parse(stdout), stderr, exitCode }).toEqual({ + result: { pass: true, complete: scenario !== "legacy-pair", failure: null, failedChecks: [] }, + stderr: "", + exitCode: 0, + }); + }, +); + it("the last 'data' event fires before the close_notify reply is written to a duplex transport (tls.connect({ socket }))", async () => { // The peer's last application data and its close_notify reach the engine in // one chunk. The engine used to answer the close_notify before it emitted @@ -882,6 +900,8 @@ it("the last 'data' event fires before the close_notify reply is written to a du // close_notify, then pushed to the client as one chunk. let held: Buffer | null = null; const clientSide: Duplex = new Duplex({ + // Node inherits this flag from the transport; a half-open client sends no automatic reply. + allowHalfOpen: false, read() {}, write(chunk: Buffer, _encoding, callback) { if (recordTypes(chunk)?.includes(21)) log.push("write close_notify"); @@ -932,7 +952,7 @@ it("the last 'data' event fires before the close_notify reply is written to a du await once(client, "close"); // 23 is application data. One push carried it and the alert. - expect(log).toEqual(["push 23,21", "data last", "write close_notify", "transport end", "end"]); + expect(log).toEqual(["push 23,21", "data last", "end", "write close_notify", "transport end"]); }); describe("application data written over a Duplex transport before the handshake completes", () => { @@ -3480,11 +3500,7 @@ describe.each([ expect(await run("duplex-destroySoon", "after-first-flight")).toEqual(clientDestroyed); }); - // Known gap in bun: 'finish' does not wait for the transport's end(). In - // the turn that created the socket it fires before the engine exists, so - // destroySoon()'s destroy() runs first and destroys the transport: no - // ClientHello and no final(), the peer only sees the connection close. - (!exe ? it.skip : runtime === "bun" ? it.failing : it)( + it.skipIf(!exe)( "destroySoon() in the turn that created the socket ends the transport, then closes the socket", async () => { expect(await run("duplex-destroySoon", "same-turn")).toEqual(clientDestroyed); diff --git a/test/js/node/tls/node-tls-raw-end.test.ts b/test/js/node/tls/node-tls-raw-end.test.ts index 519695a2c520..1486fdcd9561 100644 --- a/test/js/node/tls/node-tls-raw-end.test.ts +++ b/test/js/node/tls/node-tls-raw-end.test.ts @@ -64,3 +64,103 @@ for (const when of ["nextTick", "setImmediate"]) { assert.ok(!events.includes("raw end"), events.join(", ")); }); } + +for (const route of ["client", "server", "injected-server"]) { + test(`a half-open ${route} wrap retains its raw socket until the application ends`, async () => { + const version = { minVersion: "TLSv1.2", maxVersion: "TLSv1.2" } as const; + const resources = []; + const closes = []; + const events = []; + const errors = []; + const targetReady = Promise.withResolvers(); + const peerReady = Promise.withResolvers(); + const reply = Promise.withResolvers(); + let raw; + let snapshot; + let endError; + let targetData = ""; + let peerData = ""; + function observe(socket, name) { + resources.push(socket); + closes.push(new Promise(resolve => socket.once("close", resolve))); + socket.on("close", () => events.push(`${name} close`)); + socket.on("error", error => errors.push(error.code)); + return socket; + } + function target(socket) { + observe(socket, "tls"); + socket.on("data", data => (targetData += data)); + socket.on("end", () => { + events.push("tls end"); + setImmediate(() => { + snapshot = { + allowHalfOpen: socket.allowHalfOpen, + writableEnded: socket.writableEnded, + destroyed: socket.destroyed, + rawDestroyed: raw.destroyed, + }; + events.push("application end"); + socket.end("tail", error => { + endError = error?.code ?? null; + reply.resolve(); + }); + }); + }); + targetReady.resolve(); + } + function peer(socket, secureEvent?) { + observe(socket, "peer"); + socket.on("data", data => (peerData += data)); + if (secureEvent) socket.once(secureEvent, () => socket.end("last")); + else socket.end("last"); + peerReady.resolve(); + } + const injected = route === "injected-server" ? tls.createServer({ key, cert, ...version }) : null; + injected?.on("secureConnection", target); + const server = + route === "client" + ? tls.createServer({ key, cert, ...version, allowHalfOpen: true }, socket => peer(socket)) + : net.createServer({ allowHalfOpen: true }, socket => { + raw = observe(socket, "raw"); + if (injected) injected.emit("connection", raw); + else target(new tls.TLSSocket(raw, { isServer: true, key, cert, ...version })); + }); + try { + await new Promise(resolve => server.listen(0, "127.0.0.1", resolve)); + const address = server.address() as net.AddressInfo; + if (route === "client") { + raw = observe(net.connect({ host: "127.0.0.1", port: address.port, allowHalfOpen: true }), "raw"); + target(tls.connect({ socket: raw, rejectUnauthorized: false, ...version })); + } else { + peer( + tls.connect({ + host: "127.0.0.1", + port: address.port, + allowHalfOpen: true, + rejectUnauthorized: false, + ...version, + }), + "secureConnect", + ); + } + await Promise.all([targetReady.promise, peerReady.promise]); + await Promise.all([reply.promise, ...closes]); + assert.deepStrictEqual( + { snapshot, targetData, peerData, endError, errors }, + { + snapshot: { allowHalfOpen: true, writableEnded: false, destroyed: false, rawDestroyed: false }, + targetData: "last", + peerData: "tail", + endError: null, + errors: [], + }, + ); + assert.ok(events.indexOf("application end") < events.indexOf("raw close"), events.join(", ")); + assert.ok(events.indexOf("application end") < events.indexOf("tls close"), events.join(", ")); + } finally { + for (const socket of resources) socket.destroy(); + server.close(); + injected?.close(); + } + }); +} diff --git a/test/js/node/tls/renegotiation.test.ts b/test/js/node/tls/renegotiation.test.ts index 0b7ae7ffae45..cd4ced27b507 100644 --- a/test/js/node/tls/renegotiation.test.ts +++ b/test/js/node/tls/renegotiation.test.ts @@ -29,6 +29,62 @@ afterAll(() => { process?.kill(); }); +it.concurrent.each(["renegotiate-native", "renegotiate-duplex"])( + "closes both transport halves after server-initiated TLS 1.2 renegotiation: %s", + async scenario => { + const fixture = join(import.meta.dir, "tls-renegotiate-shutdown-fixture.cjs"); + await using server = Bun.spawn({ + cmd: ["node", fixture, "--oracle-server", scenario], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const ready = Promise.withResolvers(); + const serverRows: any[] = []; + const serverOutput = (async () => { + let pending = ""; + for await (const chunk of server.stdout) { + pending += Buffer.from(chunk).toString(); + let end; + while ((end = pending.indexOf("\n")) !== -1) { + const row = JSON.parse(pending.slice(0, end)); + pending = pending.slice(end + 1); + serverRows.push(row); + if (row.type === "oracle-ready") ready.resolve(row.port); + } + } + ready.reject(new Error("TLS oracle exited before listening")); + })(); + const serverErrors = server.stderr.text(); + await using client = Bun.spawn({ + cmd: [bunExe(), fixture, "--client", scenario, String(await ready.promise)], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, clientExit, serverStderr, serverExit] = await Promise.all([ + client.stdout.text(), + client.stderr.text(), + client.exited, + serverErrors, + server.exited, + serverOutput, + ]); + const expected = { type: "role-result", pass: true, complete: true, failure: null, failedChecks: [] }; + expect({ client: JSON.parse(stdout), server: serverRows[1], stderr, serverStderr, clientExit, serverExit }).toEqual( + { + client: { ...expected, role: "client" }, + server: { ...expected, role: "oracle-server" }, + stderr: "", + serverStderr: "", + clientExit: 0, + serverExit: 0, + }, + ); + expect(serverRows).toHaveLength(2); + }, +); + it("allow renegotiation in fetch", async () => { const body = await fetch(url, { verbose: true, diff --git a/test/js/node/tls/tls-half-close-transport-fixture.cjs b/test/js/node/tls/tls-half-close-transport-fixture.cjs new file mode 100644 index 000000000000..3ad07d140993 --- /dev/null +++ b/test/js/node/tls/tls-half-close-transport-fixture.cjs @@ -0,0 +1,207 @@ +// Node 24: close_notify ends the readable half; stream policy owns the writable half. +"use strict"; +const fs = require("node:fs"); +const path = require("node:path"); +const tls = require("node:tls"); +const { Duplex } = require("node:stream"); +let mainFailure = null; +const out = value => fs.writeSync(1, JSON.stringify(value) + "\n"); +const fixture = () => ({ + key: fs.readFileSync(path.join(__dirname, "fixtures/agent1-key.pem")), + cert: fs.readFileSync(path.join(__dirname, "fixtures/agent1-cert.pem")), +}); +const version = { minVersion: "TLSv1.2", maxVersion: "TLSv1.2" }; +const waitEvent = (emitter, event) => new Promise(resolve => emitter.once(event, (...args) => resolve(args))); +async function child(name) { + const trace = []; + let complete = false; + let evaluate = () => ({ configured: false }); + let failure = null; + const event = (owner, eventName, data = {}) => { + const entry = { type: "event", case: name, seq: trace.length + 1, owner, event: eventName, ...data }; + trace.push(entry); + }; + const entries = (owner, eventName) => trace.filter(x => x.owner === owner && x.event === eventName); + const count = (owner, eventName) => entries(owner, eventName).length; + const one = (owner, eventName) => count(owner, eventName) === 1; + const pos = (owner, eventName) => entries(owner, eventName)[0]?.seq; + const before = (a, b, c, d) => one(a, b) && one(c, d) && pos(a, b) < pos(c, d); + const noUnexpectedErrors = () => + !trace.some( + x => + x.event === "error" || + x.event === "tlsClientError" || + x.event === "uncaughtException" || + x.event === "unhandledRejection", + ); + const observe = (stream, owner, { data = false } = {}) => { + for (const eventName of ["end", "finish", "close"]) + stream.on(eventName, value => event(owner, eventName, eventName === "close" ? { hadError: value ?? null } : {})); + stream.on("error", error => event(owner, "error", { code: error.code ?? null, message: error.message })); + if (data) stream.on("data", chunk => event(owner, "data", { text: chunk.toString() })); + return stream; + }; + process.on("uncaughtException", error => { + failure = String(error); + event("process", "uncaughtException", { message: String(error) }); + }); + process.on("unhandledRejection", error => { + failure = String(error); + event("process", "unhandledRejection", { message: String(error) }); + }); + process.once("beforeExit", () => { + const checks = evaluate(); + const pass = !failure && !mainFailure && Object.values(checks).every(Boolean); + event("process", "beforeExit", { complete }); + out({ + pass, + complete, + failure: failure ?? mainFailure, + failedChecks: Object.keys(checks).filter(key => !checks[key]), + }); + process.exitCode = pass ? 0 : 1; + }); + function recordTypes(buffer) { + const result = []; + let offset = 0; + while (offset + 5 <= buffer.length) { + const next = offset + 5 + buffer.readUInt16BE(offset + 3); + if (next > buffer.length) return null; + result.push(buffer[offset]); + offset = next; + } + return offset === buffer.length ? result : null; + } + if (name === "legacy-pair" || name.startsWith("duplex-")) { + const legacy = name === "legacy-pair", + rawEof = name === "duplex-eof-halfopen-true", + halfOpen = name === "duplex-halfopen-true" || rawEof; + let left, + right, + held = null; + const rawOptions = legacy ? {} : { allowHalfOpen: halfOpen }; + left = observe( + new Duplex({ + ...rawOptions, + read() {}, + write(chunk, enc, cb) { + const records = recordTypes(chunk); + event("client-transport", "write", { records }); + if (records?.includes(21)) event("client-transport", "alert"); + right.push(chunk); + cb(); + }, + final(cb) { + event("client-transport", "final"); + right.push(null); + cb(); + }, + }), + "client-transport", + ); + right = observe( + new Duplex({ + ...rawOptions, + read() {}, + write(chunk, enc, cb) { + if (held === null) left.push(chunk); + else { + held = Buffer.concat([held, chunk]); + if (recordTypes(held)?.at(-1) === 21) { + const burst = held; + held = null; + event("relay", "burst", { records: recordTypes(burst) }); + if (rawEof) { + const dataEnd = 5 + burst.readUInt16BE(3); + event("relay", "alert.dropped", { records: recordTypes(burst.subarray(dataEnd)) }); + left.push(burst.subarray(0, dataEnd)); + } else { + left.push(burst); + } + } + } + cb(); + }, + final(cb) { + event("server-transport", "final"); + left.push(null); + cb(); + }, + }), + "server-transport", + ); + const tlsOptions = legacy ? {} : { allowHalfOpen: halfOpen }; + const server = observe( + new tls.TLSSocket(right, { + ...tlsOptions, + isServer: true, + secureContext: tls.createSecureContext({ ...fixture(), ...version }), + }), + "server", + { data: true }, + ); + server.on("data", chunk => { + if (chunk.toString() === "go") server.end("last"); + }); + const client = observe( + tls.connect({ ...tlsOptions, socket: left, rejectUnauthorized: false, ...version }), + "client", + { data: true }, + ); + client.on("end", () => { + event("application", "client.end.observed", { + allowHalfOpen: client.allowHalfOpen, + transportAllowHalfOpen: left.allowHalfOpen, + writableEnded: client.writableEnded, + }); + if (halfOpen) { + event("application", "client.end.call", { text: "tail" }); + client.end("tail"); + } + }); + evaluate = () => ({ + handshake: one("client", "secureConnect"), + burst: one("relay", "burst") && JSON.stringify(entries("relay", "burst")[0].records) === "[23,21]", + lastData: one("client", "data") && entries("client", "data")[0].text === "last", + dataBeforeEnd: before("client", "data", "client", "end"), + legacyQuiescence: + !legacy || (count("client", "close") === 0 && !complete && count("client-transport", "alert") === 0), + completed: legacy || complete, + replyAfterData: legacy || before("client", "data", "client-transport", "alert"), + endBeforeClose: legacy || before("client", "end", "client", "close"), + bothTlsClosed: legacy || (one("client", "close") && one("server", "close")), + bothRawClosed: legacy || (one("client-transport", "close") && one("server-transport", "close")), + tailDelivered: !halfOpen || entries("server", "data").some(x => x.text === "tail"), + rawEofBeforeAlert: + !rawEof || + (one("relay", "alert.dropped") && + JSON.stringify(entries("relay", "alert.dropped")[0].records) === "[21]" && + before("client-transport", "end", "client-transport", "alert")), + writableOpenAtEnd: !halfOpen || entries("application", "client.end.observed")[0]?.writableEnded === false, + noErrors: noUnexpectedErrors(), + }); + await waitEvent(client, "secureConnect"); + event("client", "secureConnect", { + protocol: client.getProtocol(), + allowHalfOpen: client.allowHalfOpen, + transportAllowHalfOpen: left.allowHalfOpen, + }); + const closed = legacy + ? waitEvent(client, "close") + : Promise.all([ + waitEvent(client, "close"), + waitEvent(server, "close"), + waitEvent(left, "close"), + waitEvent(right, "close"), + ]); + held = Buffer.alloc(0); + client.write("go"); + await closed; + complete = true; + return; + } +} +child(process.argv[2]).catch(error => { + mainFailure = String(error); + process.exitCode = 1; +}); diff --git a/test/js/node/tls/tls-renegotiate-shutdown-fixture.cjs b/test/js/node/tls/tls-renegotiate-shutdown-fixture.cjs new file mode 100644 index 000000000000..d7228922e720 --- /dev/null +++ b/test/js/node/tls/tls-renegotiate-shutdown-fixture.cjs @@ -0,0 +1,248 @@ +"use strict"; +const fs = require("node:fs"); +const path = require("node:path"); +const net = require("node:net"); +const tls = require("node:tls"); +const { Duplex } = require("node:stream"); +const out = value => fs.writeSync(1, JSON.stringify(value) + "\n"); +const fixture = () => ({ + key: fs.readFileSync(path.join(__dirname, "fixtures/agent1-key.pem")), + cert: fs.readFileSync(path.join(__dirname, "fixtures/agent1-cert.pem")), +}); +const versions = { minVersion: "TLSv1.2", maxVersion: "TLSv1.2" }; +const waitEvent = (emitter, name) => new Promise(resolve => emitter.once(name, (...args) => resolve(args))); +function observer(name, role) { + const trace = []; + let failure = null, + complete = false; + let evaluate = () => ({ configured: false }); + const event = (owner, eventName, detail = {}) => { + const entry = { type: "event", case: name, role, seq: trace.length + 1, owner, event: eventName, ...detail }; + trace.push(entry); + }; + const entries = (owner, eventName) => trace.filter(x => x.owner === owner && x.event === eventName); + const count = (owner, eventName) => entries(owner, eventName).length; + const one = (owner, eventName) => count(owner, eventName) === 1; + const before = (a, b, c, d) => one(a, b) && one(c, d) && entries(a, b)[0].seq < entries(c, d)[0].seq; + const observe = (stream, owner, data = false) => { + for (const eventName of ["end", "finish", "close"]) + stream.on(eventName, value => event(owner, eventName, eventName === "close" ? { hadError: value ?? null } : {})); + stream.on("error", error => event(owner, "error", { code: error.code ?? null, message: error.message })); + if (data) stream.on("data", chunk => event(owner, "data", { text: chunk.toString() })); + return stream; + }; + const noErrors = () => + !trace.some(x => ["error", "tlsClientError", "uncaughtException", "unhandledRejection"].includes(x.event)); + process.on("uncaughtException", error => { + failure = String(error); + event("process", "uncaughtException", { message: failure }); + }); + process.on("unhandledRejection", error => { + failure = String(error); + event("process", "unhandledRejection", { message: failure }); + }); + process.once("beforeExit", () => { + const checks = { completed: complete, ...evaluate() }; + const pass = !failure && Object.values(checks).every(Boolean); + event("process", "beforeExit", { complete }); + out({ + type: "role-result", + role, + pass, + complete, + failure, + failedChecks: Object.keys(checks).filter(key => !checks[key]), + }); + process.exitCode = pass ? 0 : 1; + }); + return { + event, + entries, + count, + one, + before, + observe, + noErrors, + setChecks: fn => (evaluate = fn), + done: () => (complete = true), + fail: error => { + failure = String(error); + event("process", "main-failure", { message: failure, stack: error.stack }); + }, + }; +} +async function runOracle(name, o) { + const { event, entries, one, before, observe, noErrors } = o; + const server = tls.createServer({ ...fixture(), ...versions, allowHalfOpen: false }, sock => { + observe(sock, "tls", true); + event("tls", "secureConnection", { protocol: sock.getProtocol() }); + sock.on("secure", () => event("tls", "secure", { protocol: sock.getProtocol() })); + sock.on("data", chunk => { + if (chunk.toString() === "ready") { + event("application", "renegotiate.call"); + const accepted = sock.renegotiate({}, error => { + event("tls", "renegotiate.callback", { error: error?.message ?? null }); + if (error) { + o.fail(error); + return; + } + event("application", "marker.write"); + sock.write("renegotiated"); + }); + event("application", "renegotiate.return", { accepted }); + if (!accepted) o.fail(new Error("Oracle server renegotiation rejected")); + } else if (chunk.toString() === "go") { + event("application", "end.call", { text: "last" }); + sock.end("last"); + } + }); + sock.on("close", () => + server.close(() => { + event("listener", "closed"); + o.done(); + }), + ); + }); + server.on("tlsClientError", error => + event("listener", "tlsClientError", { code: error.code ?? null, message: error.message }), + ); + server.on("error", error => event("listener", "error", { code: error.code ?? null, message: error.message })); + o.setChecks(() => ({ + oneHandshake: one("tls", "secureConnection"), + tls12: entries("tls", "secureConnection")[0]?.protocol === "TLSv1.2", + request: + entries("tls", "data").length === 2 && + entries("tls", "data")[0].text === "ready" && + entries("tls", "data")[1].text === "go", + renegotiationAccepted: entries("application", "renegotiate.return")[0]?.accepted === true, + renegotiationSucceeded: + one("tls", "renegotiate.callback") && entries("tls", "renegotiate.callback")[0].error === null, + callbackBeforeMarker: before("tls", "renegotiate.callback", "application", "marker.write"), + markerBeforeEnd: before("application", "marker.write", "application", "end.call"), + endCallBeforeFinish: before("application", "end.call", "tls", "finish"), + finishBeforeClose: before("tls", "finish", "tls", "close"), + endBeforeClose: before("tls", "end", "tls", "close"), + closeFlag: entries("tls", "close")[0]?.hadError === false, + listenerClosed: one("listener", "closed"), + noErrors: noErrors(), + })); + server.listen(0, "127.0.0.1"); + await waitEvent(server, "listening"); + out({ type: "oracle-ready", port: server.address().port }); +} +function parseRecords(event, owner) { + let pending = Buffer.alloc(0); + return chunk => { + pending = Buffer.concat([pending, chunk]); + while (pending.length >= 5) { + const size = 5 + pending.readUInt16BE(3); + if (pending.length < size) break; + const type = pending[0]; + event(owner, "record", { recordType: type, bytes: size }); + if (type === 21) event(owner, "alert"); + pending = pending.subarray(size); + } + }; +} +async function runClient(name, port, o) { + const { event, entries, count, one, before, observe, noErrors } = o; + const generic = name === "renegotiate-duplex"; + const raw = observe(net.connect({ host: "127.0.0.1", port, allowHalfOpen: generic }), "raw"); + let transport = raw; + if (generic) { + const parse = parseRecords(event, "client-wire"); + transport = observe( + new Duplex({ + allowHalfOpen: false, + read() { + raw.resume(); + }, + write(chunk, encoding, cb) { + parse(chunk); + raw.write(chunk, encoding, cb); + }, + final(cb) { + event("transport", "final"); + raw.end(cb); + }, + destroy(error, cb) { + event("transport", "_destroy", { error: error?.message ?? null }); + if (raw.closed) { + cb(error); + return; + } + raw.once("close", () => cb(error)); + raw.destroy(error); + }, + }), + "transport", + ); + raw.on("data", chunk => { + if (!transport.push(chunk)) raw.pause(); + }); + raw.on("end", () => { + event("transport", "eof.forward"); + transport.push(null); + }); + raw.on("error", error => transport.destroy(error)); + } + const client = observe( + tls.connect({ socket: transport, rejectUnauthorized: false, allowHalfOpen: false, ...versions }), + "tls", + true, + ); + client.on("secureConnect", () => + event("tls", "secureConnect", { + protocol: client.getProtocol(), + allowHalfOpen: client.allowHalfOpen, + transportAllowHalfOpen: transport.allowHalfOpen, + rawAllowHalfOpen: raw.allowHalfOpen, + }), + ); + o.setChecks(() => ({ + handshakeCount: count("tls", "secureConnect") === 2, + tls12: entries("tls", "secureConnect").every(x => x.protocol === "TLSv1.2"), + markerBeforeRequest: before("application", "marker.received", "application", "request.write"), + dataSequence: + entries("tls", "data").length === 2 && + entries("tls", "data")[0].text === "renegotiated" && + entries("tls", "data")[1].text === "last", + dataBeforeEnd: before("application", "last.received", "tls", "end"), + endBeforeClose: before("tls", "end", "tls", "close"), + finishBeforeClose: before("tls", "finish", "tls", "close"), + closeFlag: entries("tls", "close")[0]?.hadError === false, + rawClosed: one("raw", "close"), + genericLifecycle: + !generic || + (one("transport", "final") && + one("transport", "finish") && + one("transport", "_destroy") && + one("transport", "close")), + genericDataBeforeReply: !generic || before("application", "last.received", "client-wire", "alert"), + genericReplyBeforeFinal: !generic || before("client-wire", "alert", "transport", "final"), + noErrors: noErrors(), + })); + client.on("data", chunk => { + if (chunk.toString() === "renegotiated") { + event("application", "marker.received"); + event("application", "request.write"); + client.write("go"); + } else if (chunk.toString() === "last") event("application", "last.received"); + }); + await waitEvent(client, "secureConnect"); + const closed = Promise.all([ + waitEvent(client, "close"), + waitEvent(raw, "close"), + ...(generic ? [waitEvent(transport, "close")] : []), + ]); + event("application", "ready.write"); + client.write("ready"); + await closed; + event("application", "all.closes"); + o.done(); +} + +const name = process.argv[3]; +const oracle = process.argv[2] === "--oracle-server"; +const observation = observer(name, oracle ? "oracle-server" : "client"); +(oracle ? runOracle(name, observation) : runClient(name, Number(process.argv[4]), observation)).catch(observation.fail); From 5da4013bc4c8fd0a5a802776fb667833059adbf8 Mon Sep 17 00:00:00 2001 From: Peter Steinberger Date: Thu, 8 Oct 2026 05:07:14 -0700 Subject: [PATCH 3/3] fix(tls): flush half-open writes after peer shutdown Continue draining encrypted output after close_notify while the transport remains open. Add a Node 24 parity case that writes outside the receive callback and waits for peer receipt before ending, so shutdown cannot mask a missing flush. --- src/uws/lib.rs | 6 ++- test/js/node/tls/node-tls-connect.test.ts | 37 ++++++++++--------- .../tls/tls-half-close-transport-fixture.cjs | 20 ++++++++-- 3 files changed, 42 insertions(+), 21 deletions(-) diff --git a/src/uws/lib.rs b/src/uws/lib.rs index 00921897cf83..410a90a882ba 100644 --- a/src/uws/lib.rs +++ b/src/uws/lib.rs @@ -1037,7 +1037,11 @@ pub mod ssl_wrapper { != 0 { self.handle_peer_shutdown(); - + // A half-open stream can still write after its readable side ends. + if !self.flags.closed_notified() { + let mut buffer = IoBuffer::uninit(); + self.handle_writing(&mut buffer); + } return false; } return true; diff --git a/test/js/node/tls/node-tls-connect.test.ts b/test/js/node/tls/node-tls-connect.test.ts index 30a9b76f0ebc..a38075384529 100644 --- a/test/js/node/tls/node-tls-connect.test.ts +++ b/test/js/node/tls/node-tls-connect.test.ts @@ -857,23 +857,26 @@ it("a client and a server TLSSocket connected through a synchronous in-memory du }); }); -it.concurrent.each(["legacy-pair", "duplex-halfopen-false", "duplex-halfopen-true", "duplex-eof-halfopen-true"])( - "duplex TLS close_notify follows the transport's half-open policy: %s", - async scenario => { - await using proc = Bun.spawn({ - cmd: [bunExe(), join(import.meta.dirname, "tls-half-close-transport-fixture.cjs"), scenario], - env: bunEnv, - stdout: "pipe", - stderr: "pipe", - }); - const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); - expect({ result: JSON.parse(stdout), stderr, exitCode }).toEqual({ - result: { pass: true, complete: scenario !== "legacy-pair", failure: null, failedChecks: [] }, - stderr: "", - exitCode: 0, - }); - }, -); +it.concurrent.each([ + "legacy-pair", + "duplex-halfopen-false", + "duplex-halfopen-true", + "duplex-eof-halfopen-true", + "duplex-write-after-end", +])("duplex TLS close_notify follows the transport's half-open policy: %s", async scenario => { + await using proc = Bun.spawn({ + cmd: [bunExe(), join(import.meta.dirname, "tls-half-close-transport-fixture.cjs"), scenario], + env: bunEnv, + stdout: "pipe", + stderr: "pipe", + }); + const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]); + expect({ result: JSON.parse(stdout), stderr, exitCode }).toEqual({ + result: { pass: true, complete: scenario !== "legacy-pair", failure: null, failedChecks: [] }, + stderr: "", + exitCode: 0, + }); +}); it("the last 'data' event fires before the close_notify reply is written to a duplex transport (tls.connect({ socket }))", async () => { // The peer's last application data and its close_notify reach the engine in diff --git a/test/js/node/tls/tls-half-close-transport-fixture.cjs b/test/js/node/tls/tls-half-close-transport-fixture.cjs index 3ad07d140993..fa230cb9d74a 100644 --- a/test/js/node/tls/tls-half-close-transport-fixture.cjs +++ b/test/js/node/tls/tls-half-close-transport-fixture.cjs @@ -75,7 +75,8 @@ async function child(name) { if (name === "legacy-pair" || name.startsWith("duplex-")) { const legacy = name === "legacy-pair", rawEof = name === "duplex-eof-halfopen-true", - halfOpen = name === "duplex-halfopen-true" || rawEof; + writeAfterEnd = name === "duplex-write-after-end", + halfOpen = name === "duplex-halfopen-true" || rawEof || writeAfterEnd; let left, right, held = null; @@ -142,6 +143,10 @@ async function child(name) { ); server.on("data", chunk => { if (chunk.toString() === "go") server.end("last"); + if (writeAfterEnd && chunk.toString() === "tail") { + event("application", "tail.received.before.end", { writableEnded: client.writableEnded }); + client.end(); + } }); const client = observe( tls.connect({ ...tlsOptions, socket: left, rejectUnauthorized: false, ...version }), @@ -155,8 +160,13 @@ async function child(name) { writableEnded: client.writableEnded, }); if (halfOpen) { - event("application", "client.end.call", { text: "tail" }); - client.end("tail"); + if (writeAfterEnd) { + // Leave the native receive callback before writing; end() must not flush the tail for us. + setImmediate(() => client.write("tail")); + } else { + event("application", "client.end.call", { text: "tail" }); + client.end("tail"); + } } }); evaluate = () => ({ @@ -172,6 +182,10 @@ async function child(name) { bothTlsClosed: legacy || (one("client", "close") && one("server", "close")), bothRawClosed: legacy || (one("client-transport", "close") && one("server-transport", "close")), tailDelivered: !halfOpen || entries("server", "data").some(x => x.text === "tail"), + tailBeforeEnd: + !writeAfterEnd || + (one("application", "tail.received.before.end") && + entries("application", "tail.received.before.end")[0].writableEnded === false), rawEofBeforeAlert: !rawEof || (one("relay", "alert.dropped") &&