diff --git a/src/runtime/node/node_net_binding.rs b/src/runtime/node/node_net_binding.rs index 23c77f09f614..fab67451d739 100644 --- a/src/runtime/node/node_net_binding.rs +++ b/src/runtime/node/node_net_binding.rs @@ -152,6 +152,7 @@ pub(crate) fn new_detached_socket(global: &JSGlobalObject, frame: &CallFrame) -> server_name: JsCell::new(None), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), diff --git a/src/runtime/socket/Listener.rs b/src/runtime/socket/Listener.rs index a49b90ca412f..c8ac7e0d3f79 100644 --- a/src/runtime/socket/Listener.rs +++ b/src/runtime/socket/Listener.rs @@ -608,6 +608,7 @@ impl Listener { server_name: JsCell::new(None), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -654,6 +655,7 @@ impl Listener { server_name: JsCell::new(None), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -1202,6 +1204,7 @@ impl Listener { ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -1288,6 +1291,7 @@ impl Listener { ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -1530,6 +1534,7 @@ fn connect_finish( ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: Default::default(), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(crate::socket::NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), diff --git a/src/runtime/socket/socket_body.rs b/src/runtime/socket/socket_body.rs index 11853f20771f..f57e10be5c01 100644 --- a/src/runtime/socket/socket_body.rs +++ b/src/runtime/socket/socket_body.rs @@ -276,6 +276,9 @@ pub struct NewSocket { pub(crate) server_name: JsCell>>, pub(crate) buffered_data_for_node_net: JsCell>, pub(crate) bytes_written: Cell, + /// First fatal `send()` errno a JS write observed (peer RST while the loop + /// was blocked). Consumed by `on_close` as the close error; gates `on_end`. + pub(crate) fatal_write_errno: Cell, pub(crate) native_callback: JsCell, /// `upgradeTLS` produces two `TLSSocket` wrappers over one @@ -912,6 +915,7 @@ impl NewSocket { &sys::Error::from_code_int(fatal_send_errno, sys::Tag::write), &global, ); + this.fatal_write_errno.set(0); // delivered below; on_close must not re-report it let _ = handlers.call_error_handler(this_value, &[this_value, err_value]); // The error handler can destroy the socket itself; only close a // still-attached socket. Close without detaching so on_close runs @@ -1278,6 +1282,7 @@ impl NewSocket { self.socket.set(SocketHandler::::DETACHED); self.buffered_data_for_node_net .with_mut(|b| b.clear_and_free()); + self.fatal_write_errno.set(0); self.detach_native_callback(); old.close(uws::CloseCode::Failure); self.poll_ref.with_mut(|p| p.unref(js_loop_ctx())); @@ -1585,6 +1590,13 @@ impl NewSocket { if this.socket.get().is_detached() { return; } + if this.fatal_write_errno.get() != 0 { + // Not a peer FIN. Close so `on_close` reports the errno; under + // allow_half_open nothing else would (sk_err already consumed). + let _keepalive = this.ref_guard(); + this.socket.get().close(uws::CloseCode::Normal); + return; + } let handlers = this.get_handlers(); log!( "onEnd {}", @@ -1966,6 +1978,8 @@ impl NewSocket { reason: Option<*mut c_void>, ) -> JsResult<()> { jsc::mark_binding!(); + // Per-transport; consumed here so a reconnected wrapper starts clean. + let fatal_write_errno = this.fatal_write_errno.replace(0); // A late close on a socket that already released its Handlers through // a path that did not route back through this dispatch - e.g. a // JS-side destroy on a TLS socket driven by an upgraded duplex. There @@ -2054,6 +2068,12 @@ impl NewSocket { &sys::Error::from_code_int(err, sys::Tag::read), &global, ); + } else if fatal_write_errno != 0 { + // Loop saw a clean HUP but a JS write already observed the RST. + js_error = ::to_js( + &sys::Error::from_code_int(fatal_write_errno, sys::Tag::write), + &global, + ); } if let Err(e) = callback.call(&global, this_value, &[this_value, js_error]) { @@ -2263,7 +2283,8 @@ impl NewSocket { Ok( match this.write_or_end::(global, args.mut_(), false) { WriteResult::Fail => JSValue::ZERO, - WriteResult::Success { wrote, .. } => JSValue::js_number_from_int32(wrote), + // `wrote < -1` is a fatal errno (recorded); native API returns -1. + WriteResult::Success { wrote, .. } => JSValue::js_number_from_int32(wrote.max(-1)), }, ) } @@ -2439,6 +2460,10 @@ impl NewSocket { // Kernel rejected the send (peer gone): return the negative errno so // JS fails the write; never close from under the caller's stack, and // leave the undeliverable buffer to the caller (aliasing). + #[cfg(not(windows))] // same quarantine as on_writable (a5e7ba5905) + if self.fatal_write_errno.get() == 0 { + self.fatal_write_errno.set(fatal_errno); + } return -fatal_errno; } let uwrote: usize = usize::try_from(res.max(0)).expect("int cast"); @@ -2946,6 +2971,10 @@ impl NewSocket { // buffer, stop re-arming the writable retry, and report the // errno so the event-loop caller surfaces it (the data was // already acknowledged to JS, so only an 'error' can). + #[cfg(not(windows))] // same quarantine as on_writable (a5e7ba5905) + if self.fatal_write_errno.get() == 0 { + self.fatal_write_errno.set(fatal_errno); + } self.buffered_data_for_node_net .with_mut(|b| b.clear_and_free()); return fatal_errno; @@ -3119,7 +3148,7 @@ impl NewSocket { if wrote >= 0 && usize::try_from(wrote).expect("int cast") == total { let _ = this.internal_flush(); } - JSValue::js_number(wrote as f64) + JSValue::js_number(f64::from(wrote.max(-1))) } }; Ok(result) @@ -3505,6 +3534,7 @@ impl NewSocket { ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -3612,6 +3642,7 @@ impl NewSocket { ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), @@ -4610,6 +4641,7 @@ pub fn js_upgrade_duplex_to_tls( ref_pollref_on_connect: Cell::new(true), buffered_data_for_node_net: JsCell::new(Vec::new()), bytes_written: Cell::new(0), + fatal_write_errno: Cell::new(0), native_callback: JsCell::new(NativeCallbacks::None), twin: JsCell::new(None), verify_error: JsCell::new(None), diff --git a/test/js/bun/net/socket.test.ts b/test/js/bun/net/socket.test.ts index 8b0bfe27de99..c507e292f59f 100644 --- a/test/js/bun/net/socket.test.ts +++ b/test/js/bun/net/socket.test.ts @@ -3330,3 +3330,105 @@ describe("TLS handshake callback throw", () => { } }); }); + +// Windows: the fatal-send detection in usockets is gated out there (see +// on_writable in socket_body.rs), so the write-side RST path is POSIX-only. +it.concurrent.skipIf(isWindows)( + "native write() on a peer-RST'd socket returns -1 and the close reports ECONNRESET", + async () => { + // The server lives in its own process so the RST arrives while this process + // is inside a synchronous write burst (the loop can't deliver it first). + using dir = tempDir("socket-rst-write", { + "server.mjs": ` + const server = Bun.listen({ + hostname: "127.0.0.1", + port: 0, + socket: { + open(s) { setTimeout(() => { try { s.terminate(); } catch {} }, 100); }, + data() {}, close() {}, error() {}, drain() {}, + }, + }); + console.log("PORT", server.port); + setTimeout(() => process.exit(0), 60000); + `, + }); + + await using child = Bun.spawn({ + cmd: [bunExe(), "server.mjs"], + env: bunEnv, + cwd: String(dir), + stdout: "pipe", + stderr: "inherit", + }); + + let port = 0; + { + const rd = child.stdout.getReader(); + let acc = ""; + while (!port) { + const { value, done } = await rd.read(); + if (done) throw new Error("server exited before reporting its port"); + acc += new TextDecoder().decode(value); + const m = acc.match(/PORT (\d+)/); + if (m) port = +m[1]; + } + rd.releaseLock(); + } + + const closed = Promise.withResolvers<{ err: unknown; events: string[] }>(); + const events: string[] = []; + const negatives: number[] = []; + + const sock = await Bun.connect({ + hostname: "127.0.0.1", + port, + socket: { + open() {}, + data() {}, + drain() {}, + error(_s, e) { + events.push("error:" + (e as any)?.code); + }, + end() { + events.push("end"); + }, + close(_s, err) { + events.push("close"); + closed.resolve({ err, events }); + }, + }, + }); + + const chunk = Buffer.alloc(65536, 1); + // Synchronous burst: keep writing until write() reports the dead peer. + // The server RSTs ~100 ms after accept; give the burst a generous deadline. + const deadline = Date.now() + 5000; + while (negatives.length < 3 && Date.now() < deadline) { + const r = sock.write(chunk); + if (r < 0) negatives.push(r); + Bun.sleepSync(5); + } + // Yield so the loop can poll the HUP and close the socket. + const { err: closeErr } = await closed.promise; + child.kill(); + + // Every negative return is the documented -1 sentinel; the raw errno must + // not leak to JS. + expect(negatives).toEqual([-1, -1, -1]); + // A peer RST is not a clean FIN: `end` must not fire. + expect(events).not.toContain("end"); + // The close carries the errno so the reset is observable. Which syscall + // observed it is platform-dependent: on Linux send() consumes sk_err + // (ECONNRESET then EPIPE) and the close reports the latched write errno; + // on macOS send() returns EPIPE without clearing so_error, so kqueue's + // recv() reports ECONNRESET first. + expect(closeErr).toBeInstanceOf(Error); + if (isLinux) { + expect((closeErr as any).code).toBe("ECONNRESET"); + expect((closeErr as any).syscall).toBe("write"); + } else { + expect(["ECONNRESET", "EPIPE"]).toContain((closeErr as any).code); + expect(["read", "write"]).toContain((closeErr as any).syscall); + } + }, +);