From 7d6df460862df5a36063d9ad353a38099b4bb84f Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 11 Sep 2026 14:56:48 +0000 Subject: [PATCH 1/4] tls: pause and resume a wrapped Duplex transport with the TLS socket's reads A TLS socket over a Duplex transport read it with no backpressure. The duplex arms of pause_stream and resume_stream returned false, so the handle's pause() and resume() did nothing and the transport stayed in flowing mode. A paused reader got the whole peer payload decrypted into its readable buffer. UpgradedDuplex now calls pause() and resume() on the wrapped stream, like node's JSStreamSocket readStop() and readStart(). A pause before the engine exists is left alone, because the handshake needs the reads and on_open forgets the paused flag. The close callback resumes a transport that it left paused, so a net.Socket transport still reads its peer's FIN and closes. A paused transport holds its 'end' event back, so the engine can now read the peer's close_notify long after the transport got its EOF. The check that drops the close_notify answer in that state read a flag set by the 'end' event. It now reads the stream's state, and the flag is gone. --- src/runtime/socket/UpgradedDuplex.rs | 101 ++++++++-- src/uws_sys/lib.rs | 10 + src/uws_sys/socket.rs | 4 +- test/js/node/tls/node-tls-connect.test.ts | 226 ++++++++++++++++++++++ 4 files changed, 327 insertions(+), 14 deletions(-) diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index 6782c6c1b5ee..847ff94cc3e6 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -64,9 +64,10 @@ pub(crate) struct UpgradedDuplex { /// Replayed by [`Self::drain_pending`] after the staged bytes, preserving /// the original data-then-EOF order. pub pending_end: 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, + /// [`Self::pause_stream`] paused `origin` and no [`Self::resume_stream`] + /// followed. [`Self::on_close`] resumes such a transport so it can still + /// drain to its own EOF once the engine is gone. + pub reads_paused: Cell, } bun_event_loop::impl_timer_owner!(UpgradedDuplex; from_timer_ptr => event_loop_timer); @@ -209,6 +210,11 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); (this.handlers.on_close)(this.handlers.ctx); + // A transport left paused would never read its peer's EOF and close. + // `teardown` neuters the thunks, so whatever it still delivers is dropped. + if this.reads_paused.get() { + this.resume_stream(); + } // closes the underlying duplex this.call_write_or_end(None, false); @@ -217,6 +223,57 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); } + /// node's `JSStreamSocket.readStop()` / `readStart()`: the transport only + /// emits 'data' while the TLS socket wants more. A chunk already handed to + /// the engine is still decrypted and delivered in full. + /// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/js_stream_socket.js#L117-L125 + /// + /// A pause before `start_tls` ran is left alone, like a socket that is + /// still connecting: `on_open` forgets the owner's paused flag, and the + /// handshake needs the reads. + #[uws_callback(export = "UpgradedDuplex__pause_stream")] + pub(crate) fn pause_stream(&self) -> bool { + if self.wrapper_ref().is_none() || !self.call_origin("pause") { + return false; + } + self.reads_paused.set(true); + true + } + + #[uws_callback(export = "UpgradedDuplex__resume_stream")] + pub(crate) fn resume_stream(&self) -> bool { + if !self.call_origin("resume") { + return false; + } + self.reads_paused.set(false); + true + } + + /// Calls `origin[name]()`. False when there is no JS duplex to talk to + /// (see [`Self::call_write_or_end`]) or the call threw (routed to `on_error`). + fn call_origin(&self, name: &str) -> bool { + let duplex = self.origin.get(); + if duplex.is_empty() { + return false; + } + let Some(global) = self.global else { + return false; + }; + let method = match duplex.get(&global, name) { + Ok(Some(f)) if f.is_callable() => f, + Ok(_) => return false, + Err(err) => { + (self.handlers.on_error)(self.handlers.ctx, global.take_error(err)); + return false; + } + }; + if let Err(err) = method.call(&global, duplex, &[]) { + (self.handlers.on_error)(self.handlers.ctx, global.take_error(err)); + return false; + } + true + } + fn call_write_or_end(&self, data: Option<&[u8]>, msg_more: bool) { // No JS duplex to talk to: the zeroed placeholder, or the owning // socket's finalizer abandoned it (`abandon_js_side`). @@ -234,12 +291,12 @@ impl UpgradedDuplex { 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() { + // side got its EOF 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::readable_got_eof(duplex, &global) { return; } match duplex.get(&global, "writableEnded") { @@ -276,6 +333,27 @@ impl UpgradedDuplex { } } + /// `duplex._readableState.ended`. The 'end' event comes too late to tell: + /// a paused transport holds it back until the unread bytes are consumed. + fn readable_got_eof(duplex: JSValue, global: &JSGlobalObject) -> bool { + let probe = || -> JsResult { + let Some(state) = duplex.get(global, "_readableState")? else { + return Ok(false); + }; + if !state.is_object() { + return Ok(false); + } + Ok(state + .get(global, "ended")? + .is_some_and(|ended| ended.to_boolean())) + }; + // Best-effort probe: consume the exception and report no EOF. + probe().unwrap_or_else(|err| { + let _ = global.take_exception(err); + false + }) + } + fn internal_write(this: *mut Self, encoded_data: &[u8]) { // SAFETY: see handler note above. unsafe { &*this }.write_encrypted(encoded_data); @@ -408,7 +486,7 @@ impl UpgradedDuplex { current_timeout: Cell::new(0), pending_data: JsCell::new(Vec::new()), pending_end: Cell::new(false), - transport_eof: Cell::new(false), + reads_paused: Cell::new(false), } } @@ -679,7 +757,7 @@ impl UpgradedDuplex { self.ssl_error.set(CertError::default()); self.pending_data.set(Vec::new()); self.pending_end.set(false); - self.transport_eof.set(false); + self.reads_paused.set(false); } } @@ -738,7 +816,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); if this.wrapper_ref().is_some() { (this.handlers.on_end)(this.handlers.ctx); } else { diff --git a/src/uws_sys/lib.rs b/src/uws_sys/lib.rs index ad9170bca813..f0ff0b1338ab 100644 --- a/src/uws_sys/lib.rs +++ b/src/uws_sys/lib.rs @@ -221,6 +221,8 @@ unsafe extern "C" { safe fn UpgradedDuplex__shutdown_read(this: &mut UpgradedDuplex); safe fn UpgradedDuplex__close(this: &mut UpgradedDuplex); safe fn UpgradedDuplex__abandon_js_side(this: &mut UpgradedDuplex); + safe fn UpgradedDuplex__pause_stream(this: &mut UpgradedDuplex) -> bool; + safe fn UpgradedDuplex__resume_stream(this: &mut UpgradedDuplex) -> bool; } impl UpgradedDuplex { #[inline] @@ -282,6 +284,14 @@ impl UpgradedDuplex { pub(crate) fn abandon_js_side(&mut self) { UpgradedDuplex__abandon_js_side(self) } + #[inline] + pub(crate) fn pause_stream(&mut self) -> bool { + UpgradedDuplex__pause_stream(self) + } + #[inline] + pub(crate) fn resume_stream(&mut self) -> bool { + UpgradedDuplex__resume_stream(self) + } } // ── WindowsNamedPipe (cycle-break shim) ───────────────────────────────────── diff --git a/src/uws_sys/socket.rs b/src/uws_sys/socket.rs index 016e77932502..7e85301b21c9 100644 --- a/src/uws_sys/socket.rs +++ b/src/uws_sys/socket.rs @@ -523,7 +523,7 @@ impl NewSocketHandler { connected s => if s.is_established() { s.pause(); true } else { false }, connecting _c => false, detached => true, - duplex _d => false, // TODO: pause/resume upgraded duplex + duplex d => d.pause_stream(), pipe p => p.pause_stream(), ) } @@ -533,7 +533,7 @@ impl NewSocketHandler { connected s => if s.is_established() { s.resume(); true } else { false }, connecting _c => false, detached => true, - duplex _d => false, // TODO: pause/resume upgraded duplex + duplex d => d.resume_stream(), pipe p => p.resume_stream(), ) } diff --git a/test/js/node/tls/node-tls-connect.test.ts b/test/js/node/tls/node-tls-connect.test.ts index cbb1c842a4ba..c95f6c89f51c 100644 --- a/test/js/node/tls/node-tls-connect.test.ts +++ b/test/js/node/tls/node-tls-connect.test.ts @@ -843,6 +843,232 @@ it("a client and a server TLSSocket connected through a synchronous in-memory du }); }); +describe("a TLS socket over a Duplex transport reads it with backpressure", () => { + // Node reads such a transport through a JSStreamSocket, whose readStop() and + // readStart() pause and resume it, so it only flows while the TLS socket + // takes more data: + // https://github.com/nodejs/node/blob/v26.3.0/lib/internal/js_stream_socket.js#L117-L125 + // Each side's _write pushes straight into the other side and never waits, + // so only the reading TLS socket can slow the transport down. + function inMemoryPair() { + const makeSide = (peer: () => Duplex) => + new Duplex({ + read() {}, + write(chunk, _encoding, callback) { + peer().push(chunk); + callback(); + }, + final(callback) { + peer().push(null); + callback(); + }, + }); + const clientSide: Duplex = makeSide(() => serverSide); + const serverSide: Duplex = makeSide(() => clientSide); + return { clientSide, serverSide }; + } + const serverContext = () => ({ isServer: true, secureContext: tls.createSecureContext(COMMON_CERT_) }); + const payload = Buffer.alloc(1024 * 1024, "x"); + // Rejects with the first 'error' any of `sockets` emits. Meant to be raced. + function firstErrorOf(...sockets: TLSSocket[]) { + const { promise, reject } = Promise.withResolvers(); + promise.catch(() => {}); + for (const socket of sockets) socket.on("error", reject); + return promise; + } + // Destroys every socket in `sockets` when the test leaves its scope. + function destroyedOnExit(sockets: { destroy(): unknown }[]) { + return { + [Symbol.dispose]() { + for (const socket of sockets) socket.destroy(); + }, + }; + } + + // Both wraps run the same engine; the reader is the side under test. + function securePair(reader: "client" | "server") { + const { clientSide, serverSide } = inMemoryPair(); + const server = new TLSSocket(serverSide, serverContext()); + const client = tls.connect({ socket: clientSide, rejectUnauthorized: false }); + const failed = firstErrorOf(server, client); + const secured = Promise.all([once(server, "secure"), once(client, "secureConnect")]); + const [readerSocket, transport, writerSocket] = + reader === "client" ? [client, clientSide, server] : [server, serverSide, client]; + return { readerSocket, transport, writerSocket, failed, secured, cleanup: destroyedOnExit([client, server]) }; + } + + describe.each(["client", "server"] as const)("%s reader", reader => { + // The writer does not end: what happens to unread data once the peer has + // closed is a separate matter from how much the socket takes in. + it("a paused socket pauses the transport once its buffer is full", async () => { + const { readerSocket, transport, writerSocket, failed, secured, cleanup } = securePair(reader); + using _ = cleanup; + await Promise.race([secured, failed]); + readerSocket.pause(); + // The pair is synchronous, so the callback runs with every byte handed + // to the reader's transport. + const written = Promise.withResolvers(); + writerSocket.write(payload, err => (err ? written.reject(err) : written.resolve())); + await Promise.race([written.promise, failed]); + // Without backpressure the transport keeps flowing and the paused socket + // holds the whole payload. + expect({ + transportFlowing: transport.readableFlowing, + socketIsFull: readerSocket.readableLength >= readerSocket.readableHighWaterMark, + mostOfItWaitsInTheTransport: transport.readableLength > payload.length / 2, + }).toEqual({ transportFlowing: false, socketIsFull: true, mostOfItWaitsInTheTransport: true }); + + const chunks: Buffer[] = []; + let total = 0; + const received = Promise.withResolvers(); + readerSocket.on("data", (chunk: Buffer) => { + chunks.push(chunk); + if ((total += chunk.length) >= payload.length) received.resolve(); + }); + readerSocket.resume(); + await Promise.race([received.promise, failed]); + expect(Buffer.concat(chunks).equals(payload)).toBe(true); + }); + + it("a slow reader gets the whole payload while the transport is paused and resumed under it", async () => { + const { readerSocket, transport, writerSocket, failed, secured, cleanup } = securePair(reader); + using _ = cleanup; + await Promise.race([secured, failed]); + let pauses = 0; + transport.on("pause", () => pauses++); + writerSocket.write(payload); + + const read = (async () => { + let total = 0; + let mostHeld = 0; + for await (const chunk of readerSocket) { + total += chunk.length; + mostHeld = Math.max(mostHeld, readerSocket.readableLength); + if (total >= payload.length) break; + // Yield a macrotask per chunk so the socket's buffer fills up. + await new Promise(resolve => setImmediate(resolve)); + } + return { total, heldLessThanHalfAtAnyTime: mostHeld < payload.length / 2 }; + })(); + expect({ ...(await Promise.race([read, failed])), pausedMoreThanOnce: pauses > 1 }).toEqual({ + total: payload.length, + heldLessThanHalfAtAnyTime: true, + pausedMoreThanOnce: true, + }); + }); + }); + + it("does not answer a close_notify that it reads after the transport's EOF", async () => { + // A net.Socket that has read its peer's FIN fails a later write with EPIPE + // (writeAfterFIN in net.ts), and so does this transport. A paused reader + // makes the engine read the peer's close_notify long after that EOF, and + // the transport's 'end' event, held back by the pause, comes later still. + const transportErrors: string[] = []; + const gotEOF = new WeakSet(); + const makeSide = (peer: () => Duplex) => + new Duplex({ + read() {}, + write(chunk, _encoding, callback) { + if (gotEOF.has(this)) + return callback(Object.assign(new Error("write after the peer's FIN"), { code: "EPIPE" })); + peer().push(chunk); + callback(); + }, + final(callback) { + gotEOF.add(peer()); + peer().push(null); + callback(); + }, + }); + const clientSide: Duplex = makeSide(() => serverSide); + const serverSide: Duplex = makeSide(() => clientSide); + clientSide.on("error", (err: NodeJS.ErrnoException) => transportErrors.push(`${err.code}`)); + const server = new TLSSocket(serverSide, serverContext()); + const client = tls.connect({ socket: clientSide, rejectUnauthorized: false }); + using _ = destroyedOnExit([client, server]); + const failed = firstErrorOf(server, client); + await Promise.race([Promise.all([once(server, "secure"), once(client, "secureConnect")]), failed]); + + client.pause(); + // Like a TLS server over TCP: the payload, the close_notify, then the FIN. + server.end(payload); + await Promise.race([once(server, "finish"), failed]); + serverSide.end(); + + let total = 0; + client.on("data", (chunk: Buffer) => (total += chunk.length)); + const closed = once(client, "close"); + client.resume(); + await Promise.race([closed, failed]); + expect({ total, transportErrors }).toEqual({ total: payload.length, transportErrors: [] }); + }); + + it("a pause() issued before the engine exists does not stall the handshake", async () => { + // The engine is created on a later event-loop turn. A transport paused + // ahead of it would never deliver the server's flight. + const { clientSide, serverSide } = inMemoryPair(); + const server = new TLSSocket(serverSide, serverContext()); + const client = tls.connect({ + socket: clientSide, + rejectUnauthorized: false, + // Only an onread socket stops its handle from pause(). + onread: { buffer: Buffer.alloc(64), callback: () => {} }, + }); + using _ = destroyedOnExit([client, server]); + client.pause(); + const secured = Promise.all([once(server, "secure"), once(client, "secureConnect")]); + await Promise.race([secured, firstErrorOf(server, client)]); + expect({ server: server.getProtocol(), client: client.getProtocol() }).toEqual({ + server: "TLSv1.3", + client: "TLSv1.3", + }); + }); + + it("destroying a paused socket lets a net.Socket transport read its peer's close", async () => { + // A server wrap over a net.Socket with unflushed plain writes cannot take + // the fd over, so it reads the net.Socket like any other Duplex. Nothing + // destroys that net.Socket along with the TLS socket, so one left paused + // would never read the client's FIN and would stay open for good. + const sockets: { destroy(): unknown }[] = []; + const accepted = Promise.withResolvers<{ transport: net.Socket; wrapped: TLSSocket }>(); + await using listener = net.createServer(transport => { + transport.cork(); + transport.write("!"); + const wrapped = new TLSSocket(transport, serverContext()); + transport.uncork(); + sockets.push(wrapped, transport); + accepted.resolve({ transport, wrapped }); + }); + // Declared after the listener, so it runs first: close() waits for them. + using _ = destroyedOnExit(sockets); + await once(listener.listen(0, "127.0.0.1"), "listening"); + const clientTransport = net.connect((listener.address() as AddressInfo).port, "127.0.0.1"); + sockets.push(clientTransport); + const [{ transport, wrapped }, [greeting]] = await Promise.all([accepted.promise, once(clientTransport, "data")]); + expect(String(greeting)).toBe("!"); + const client = tls.connect({ socket: clientTransport, rejectUnauthorized: false }); + sockets.push(client); + const failed = firstErrorOf(wrapped, client); + await Promise.race([Promise.all([once(wrapped, "secure"), once(client, "secureConnect")]), failed]); + + wrapped.pause(); + // Without backpressure the engine takes the whole payload and the client's + // close_notify, answers it, and the client closes. + const outcome = Promise.race([ + once(transport, "pause").then(() => "transport paused"), + once(client, "close").then(() => "client closed"), + failed, + ]); + client.end(payload); + expect(await outcome).toBe("transport paused"); + + const closed = once(transport, "close"); + wrapped.destroy(); + await closed; + expect(transport.destroyed).toBe(true); + }); +}); + describe("application data written over a Duplex transport before the handshake completes", () => { // Node parks such a write (TLSWrap's pending cleartext) and sends it right // after the handshake: the write is still pending when 'secureConnect' / From 1d2e7dd019b28dfbbe9295f5e01bb7be8f0b5489 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 11 Sep 2026 15:15:15 +0000 Subject: [PATCH 2/4] tls: route a failed transport state probe to the error handler readable_got_eof took the exception of a throwing _readableState getter and dropped it. It now returns JsResult, and call_write_or_end hands the error to on_error like the write and end calls next to it. --- src/runtime/socket/UpgradedDuplex.rs | 36 ++++++++++++++-------------- 1 file changed, 18 insertions(+), 18 deletions(-) diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index 847ff94cc3e6..13ea56ff5701 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -296,8 +296,15 @@ impl UpgradedDuplex { // 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::readable_got_eof(duplex, &global) { - return; + if data.is_some() { + match Self::readable_got_eof(duplex, &global) { + Ok(false) => {} + Ok(true) => return, + Err(err) => { + (self.handlers.on_error)(self.handlers.ctx, global.take_error(err)); + return; + } + } } match duplex.get(&global, "writableEnded") { Ok(Some(ended)) if ended.to_boolean() => return, @@ -335,23 +342,16 @@ impl UpgradedDuplex { /// `duplex._readableState.ended`. The 'end' event comes too late to tell: /// a paused transport holds it back until the unread bytes are consumed. - fn readable_got_eof(duplex: JSValue, global: &JSGlobalObject) -> bool { - let probe = || -> JsResult { - let Some(state) = duplex.get(global, "_readableState")? else { - return Ok(false); - }; - if !state.is_object() { - return Ok(false); - } - Ok(state - .get(global, "ended")? - .is_some_and(|ended| ended.to_boolean())) + fn readable_got_eof(duplex: JSValue, global: &JSGlobalObject) -> JsResult { + let Some(state) = duplex.get(global, "_readableState")? else { + return Ok(false); }; - // Best-effort probe: consume the exception and report no EOF. - probe().unwrap_or_else(|err| { - let _ = global.take_exception(err); - false - }) + if !state.is_object() { + return Ok(false); + } + Ok(state + .get(global, "ended")? + .is_some_and(|ended| ended.to_boolean())) } fn internal_write(this: *mut Self, encoded_data: &[u8]) { From fc69389f7e9b84cfa65b84e3bbc9f5b870ee332d Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 11 Sep 2026 15:20:47 +0000 Subject: [PATCH 3/4] tls: cut the new UpgradedDuplex comments to one line each --- src/runtime/socket/UpgradedDuplex.rs | 33 ++++++++++------------------ 1 file changed, 11 insertions(+), 22 deletions(-) diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index 13ea56ff5701..1ea2d9461aba 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -64,9 +64,7 @@ pub(crate) struct UpgradedDuplex { /// Replayed by [`Self::drain_pending`] after the staged bytes, preserving /// the original data-then-EOF order. pub pending_end: Cell, - /// [`Self::pause_stream`] paused `origin` and no [`Self::resume_stream`] - /// followed. [`Self::on_close`] resumes such a transport so it can still - /// drain to its own EOF once the engine is gone. + /// [`Self::pause_stream`] paused `origin`; [`Self::on_close`] undoes it. pub reads_paused: Cell, } @@ -210,8 +208,7 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); (this.handlers.on_close)(this.handlers.ctx); - // A transport left paused would never read its peer's EOF and close. - // `teardown` neuters the thunks, so whatever it still delivers is dropped. + // Left paused, a net.Socket transport never reads its peer's FIN and stays open. if this.reads_paused.get() { this.resume_stream(); } @@ -223,16 +220,10 @@ impl UpgradedDuplex { js_wrapper.ensure_still_alive(); } - /// node's `JSStreamSocket.readStop()` / `readStart()`: the transport only - /// emits 'data' while the TLS socket wants more. A chunk already handed to - /// the engine is still decrypted and delivered in full. - /// https://github.com/nodejs/node/blob/v26.3.0/lib/internal/js_stream_socket.js#L117-L125 - /// - /// A pause before `start_tls` ran is left alone, like a socket that is - /// still connecting: `on_open` forgets the owner's paused flag, and the - /// handshake needs the reads. + /// node's `JSStreamSocket.readStop()`: https://github.com/nodejs/node/blob/v26.3.0/lib/internal/js_stream_socket.js#L117-L125 #[uws_callback(export = "UpgradedDuplex__pause_stream")] pub(crate) fn pause_stream(&self) -> bool { + // Before `start_tls` the handshake still needs the reads, and `on_open` clears the owner's paused flag. if self.wrapper_ref().is_none() || !self.call_origin("pause") { return false; } @@ -249,8 +240,7 @@ impl UpgradedDuplex { true } - /// Calls `origin[name]()`. False when there is no JS duplex to talk to - /// (see [`Self::call_write_or_end`]) or the call threw (routed to `on_error`). + /// Calls `origin[name]()`. A throw goes to `on_error`; false when the call did not complete. fn call_origin(&self, name: &str) -> bool { let duplex = self.origin.get(); if duplex.is_empty() { @@ -291,11 +281,11 @@ impl UpgradedDuplex { 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 got its EOF 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. + // 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() { match Self::readable_got_eof(duplex, &global) { Ok(false) => {} @@ -340,8 +330,7 @@ impl UpgradedDuplex { } } - /// `duplex._readableState.ended`. The 'end' event comes too late to tell: - /// a paused transport holds it back until the unread bytes are consumed. + /// `_readableState.ended`, not the 'end' event: a paused transport holds 'end' back. fn readable_got_eof(duplex: JSValue, global: &JSGlobalObject) -> JsResult { let Some(state) = duplex.get(global, "_readableState")? else { return Ok(false); From fc6876bb2da28195501fdbe78165e2b7b1e2c958 Mon Sep 17 00:00:00 2001 From: robobun <117481402+robobun@users.noreply.github.com> Date: Fri, 11 Sep 2026 15:35:57 +0000 Subject: [PATCH 4/4] tls: mark the transport paused before pause() runs pause() on the transport is user code. A 'pause' listener that destroys the TLS socket closed the engine while reads_paused was still false, so the close callback did not resume the transport and a net.Socket transport stayed open. The flag is now set before the call and stays set when the call fails. The two net.Socket transport tests share one setup, and the wait for the transport's close also fails on a socket error. --- src/runtime/socket/UpgradedDuplex.rs | 7 +- test/js/node/tls/node-tls-connect.test.ts | 117 ++++++++++++++-------- 2 files changed, 81 insertions(+), 43 deletions(-) diff --git a/src/runtime/socket/UpgradedDuplex.rs b/src/runtime/socket/UpgradedDuplex.rs index 1ea2d9461aba..c64cb4acb3b2 100644 --- a/src/runtime/socket/UpgradedDuplex.rs +++ b/src/runtime/socket/UpgradedDuplex.rs @@ -64,7 +64,7 @@ pub(crate) struct UpgradedDuplex { /// Replayed by [`Self::drain_pending`] after the staged bytes, preserving /// the original data-then-EOF order. pub pending_end: Cell, - /// [`Self::pause_stream`] paused `origin`; [`Self::on_close`] undoes it. + /// [`Self::pause_stream`] called `origin.pause()`; [`Self::on_close`] undoes it. pub reads_paused: Cell, } @@ -224,11 +224,12 @@ impl UpgradedDuplex { #[uws_callback(export = "UpgradedDuplex__pause_stream")] pub(crate) fn pause_stream(&self) -> bool { // Before `start_tls` the handshake still needs the reads, and `on_open` clears the owner's paused flag. - if self.wrapper_ref().is_none() || !self.call_origin("pause") { + if self.wrapper_ref().is_none() { return false; } + // Set first and kept on failure: `pause()` is user code that can close this socket, or throw after it paused. self.reads_paused.set(true); - true + self.call_origin("pause") } #[uws_callback(export = "UpgradedDuplex__resume_stream")] diff --git a/test/js/node/tls/node-tls-connect.test.ts b/test/js/node/tls/node-tls-connect.test.ts index c95f6c89f51c..0310b3779dcc 100644 --- a/test/js/node/tls/node-tls-connect.test.ts +++ b/test/js/node/tls/node-tls-connect.test.ts @@ -1024,48 +1024,85 @@ describe("a TLS socket over a Duplex transport reads it with backpressure", () = }); }); - it("destroying a paused socket lets a net.Socket transport read its peer's close", async () => { - // A server wrap over a net.Socket with unflushed plain writes cannot take - // the fd over, so it reads the net.Socket like any other Duplex. Nothing - // destroys that net.Socket along with the TLS socket, so one left paused - // would never read the client's FIN and would stay open for good. - const sockets: { destroy(): unknown }[] = []; - const accepted = Promise.withResolvers<{ transport: net.Socket; wrapped: TLSSocket }>(); - await using listener = net.createServer(transport => { - transport.cork(); - transport.write("!"); - const wrapped = new TLSSocket(transport, serverContext()); - transport.uncork(); - sockets.push(wrapped, transport); - accepted.resolve({ transport, wrapped }); + describe("a server wrap over a net.Socket with unflushed plain writes", () => { + // Such a wrap cannot take the fd over, so it reads the net.Socket like any + // other Duplex. Nothing destroys that net.Socket along with the TLS + // socket, so one left paused would never read the client's FIN and would + // stay open for good. + async function connectedWrap() { + const sockets: { destroy(): unknown }[] = []; + const accepted = Promise.withResolvers<{ transport: net.Socket; wrapped: TLSSocket }>(); + const listener = net.createServer(transport => { + transport.cork(); + transport.write("!"); + const wrapped = new TLSSocket(transport, serverContext()); + transport.uncork(); + sockets.push(wrapped, transport); + accepted.resolve({ transport, wrapped }); + }); + // close() waits for the connections, so they go first. + const dispose = async () => { + for (const socket of sockets) socket.destroy(); + await listener[Symbol.asyncDispose](); + }; + try { + await once(listener.listen(0, "127.0.0.1"), "listening"); + const clientTransport = net.connect((listener.address() as AddressInfo).port, "127.0.0.1"); + sockets.push(clientTransport); + const [{ transport, wrapped }, [greeting]] = await Promise.all([ + accepted.promise, + once(clientTransport, "data"), + ]); + expect(String(greeting)).toBe("!"); + const client = tls.connect({ socket: clientTransport, rejectUnauthorized: false }); + sockets.push(client); + const failed = firstErrorOf(wrapped, client); + await Promise.race([Promise.all([once(wrapped, "secure"), once(client, "secureConnect")]), failed]); + return { transport, wrapped, client, failed, [Symbol.asyncDispose]: dispose }; + } catch (err) { + await dispose(); + throw err; + } + } + + it("destroying the paused socket lets the transport read its peer's close", async () => { + await using wrap = await connectedWrap(); + const { transport, wrapped, client, failed } = wrap; + wrapped.pause(); + // Without backpressure the engine takes the whole payload and the + // client's close_notify, answers it, and the client closes. + const outcome = Promise.race([ + once(transport, "pause").then(() => "transport paused"), + once(client, "close").then(() => "client closed"), + failed, + ]); + client.end(payload); + expect(await outcome).toBe("transport paused"); + + const closed = once(transport, "close"); + wrapped.destroy(); + await Promise.race([closed, failed]); + expect(transport.destroyed).toBe(true); }); - // Declared after the listener, so it runs first: close() waits for them. - using _ = destroyedOnExit(sockets); - await once(listener.listen(0, "127.0.0.1"), "listening"); - const clientTransport = net.connect((listener.address() as AddressInfo).port, "127.0.0.1"); - sockets.push(clientTransport); - const [{ transport, wrapped }, [greeting]] = await Promise.all([accepted.promise, once(clientTransport, "data")]); - expect(String(greeting)).toBe("!"); - const client = tls.connect({ socket: clientTransport, rejectUnauthorized: false }); - sockets.push(client); - const failed = firstErrorOf(wrapped, client); - await Promise.race([Promise.all([once(wrapped, "secure"), once(client, "secureConnect")]), failed]); - - wrapped.pause(); - // Without backpressure the engine takes the whole payload and the client's - // close_notify, answers it, and the client closes. - const outcome = Promise.race([ - once(transport, "pause").then(() => "transport paused"), - once(client, "close").then(() => "client closed"), - failed, - ]); - client.end(payload); - expect(await outcome).toBe("transport paused"); - const closed = once(transport, "close"); - wrapped.destroy(); - await closed; - expect(transport.destroyed).toBe(true); + it("a destroy from the transport's 'pause' listener does the same", async () => { + await using wrap = await connectedWrap(); + const { transport, wrapped, client, failed } = wrap; + wrapped.pause(); + // The engine is still inside the transport's pause() when the socket closes. + let destroyedFromPause = false; + transport.once("pause", () => { + destroyedFromPause = true; + wrapped.destroy(); + }); + const closed = once(transport, "close"); + client.end(payload); + await Promise.race([closed, failed]); + expect({ destroyedFromPause, transportDestroyed: transport.destroyed }).toEqual({ + destroyedFromPause: true, + transportDestroyed: true, + }); + }); }); });