Skip to content
Closed
2 changes: 2 additions & 0 deletions packages/bun-usockets/src/context.c
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,7 @@ static void us_internal_init_listen_socket(struct us_listen_socket_t *ls,
s->flags.adopted = 0;
s->flags.allow_half_open = (options & LIBUS_SOCKET_ALLOW_HALF_OPEN);
s->unclassified_send_failures = 0;
s->read_eof = 0;
s->next = 0;
s->prev = 0;
s->connect_state = NULL;
Expand Down Expand Up @@ -508,6 +509,7 @@ static inline void us_internal_init_connect_socket(struct us_socket_t *s,
s->flags.adopted = 0;
s->flags.last_write_failed = 0;
s->unclassified_send_failures = 0;
s->read_eof = 0;
s->connect_state = NULL;
s->connect_next = NULL;
}
Expand Down
19 changes: 15 additions & 4 deletions packages/bun-usockets/src/eventing/epoll_kqueue.c
Original file line number Diff line number Diff line change
Expand Up @@ -280,8 +280,9 @@ static void us_internal_dispatch_ready_polls(struct us_loop_t *loop) {
uint8_t writable : 1;
uint8_t error : 1;
uint8_t eof : 1;
uint8_t send_eof : 1;
uint8_t skip : 1;
uint8_t _pad : 3;
uint8_t _pad : 2;
};

_Static_assert(sizeof(struct kevent_flags) == 1, "kevent_flags must be 1 byte");
Expand All @@ -305,7 +306,9 @@ static void us_internal_dispatch_ready_polls(struct us_loop_t *loop) {
#endif
.writable = (filter == EVFILT_WRITE),
.error = !!(flags & EV_ERROR),
.eof = !!(flags & EV_EOF),
/* EV_EOF on EVFILT_READ is the peer's FIN; on EVFILT_WRITE it is SS_CANTSENDMORE (peer gone, or our own shutdown) - not a read EOF (libuv kqueue.c ignores it there too). */
.eof = (flags & EV_EOF) && filter == EVFILT_READ,
.send_eof = (flags & EV_EOF) && filter == EVFILT_WRITE,
};

/* Look backward for a prior entry with the same poll to coalesce into.
Expand All @@ -317,6 +320,7 @@ static void us_internal_dispatch_ready_polls(struct us_loop_t *loop) {
coalesced[j].writable |= bits.writable;
coalesced[j].error |= bits.error;
coalesced[j].eof |= bits.eof;
coalesced[j].send_eof |= bits.send_eof;
coalesced[i] = (struct kevent_flags){ .skip = 1 };
merged = 1;
break;
Expand Down Expand Up @@ -344,9 +348,16 @@ static void us_internal_dispatch_ready_polls(struct us_loop_t *loop) {
int events = (bits.readable ? LIBUS_SOCKET_READABLE : 0)
| (bits.writable ? LIBUS_SOCKET_WRITABLE : 0);

int error = bits.error;
/* Write side dead without our own shutdown(): peer reset (or connect refused). Surface it as the poll error epoll would report as EPOLLERR. */
if (bits.send_eof) {
int type = us_internal_poll_type(poll);
if (type == POLL_TYPE_SOCKET || type == POLL_TYPE_SEMI_SOCKET) error = 1;
}

events &= us_poll_events(poll);
if (events || bits.error || bits.eof) {
us_internal_dispatch_ready_poll(poll, bits.error, bits.eof, events);
if (events || error || bits.eof) {
us_internal_dispatch_ready_poll(poll, error, bits.eof, events);
}
}
#endif
Expand Down
2 changes: 2 additions & 0 deletions packages/bun-usockets/src/internal/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -307,6 +307,8 @@ struct us_socket_t {
* the driver's epilogue via ssl_pending_detach. */
unsigned char ssl_in_use : 1;
unsigned char ssl_pending_detach : 1;
/* Peer FIN was dispatched as on_end on a half-open socket; readable interest is never re-added and on_end never re-fires. */
unsigned char read_eof : 1;
/* The close code passed to the deferred close (e.g. a reset requested from
* inside a handshake callback must still RST, not FIN, when it is finally
* performed). */
Expand Down
7 changes: 6 additions & 1 deletion packages/bun-usockets/src/loop.c
Original file line number Diff line number Diff line change
Expand Up @@ -540,6 +540,7 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in
s->flags.adopted = 0;
s->flags.last_write_failed = 0;
s->unclassified_send_failures = 0;
s->read_eof = 0;

/* We always use nodelay */
bsd_socket_nodelay(client_fd, 1);
Expand Down Expand Up @@ -851,7 +852,11 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in
s = us_internal_socket_close_raw(s, LIBUS_SOCKET_CLOSE_CODE_CLEAN_SHUTDOWN, NULL);
return;
}
if(s->flags.allow_half_open) {
if (s->flags.allow_half_open && s->read_eof) {
/* on_end already delivered (libuv UV_HANDLE_READ_EOF): just drop the readable interest that re-surfaced it. */
us_poll_change(&s->p, loop, us_poll_events(&s->p) & LIBUS_SOCKET_WRITABLE);
} else if(s->flags.allow_half_open) {
s->read_eof = 1;
/* EOF with half-open allowed: stop polling readable but KEEP
* polling writable. Masking with the current events dropped
* writable when the EOF landed before the poll had been
Expand Down
8 changes: 5 additions & 3 deletions packages/bun-usockets/src/socket.c
Original file line number Diff line number Diff line change
Expand Up @@ -415,7 +415,7 @@ struct us_socket_t *us_socket_pair(struct us_socket_group_t *group, unsigned cha
* deliver data the caller asked to defer. */
static void us_internal_rearm_writable(struct us_socket_t *s) {
us_poll_change(&s->p, s->group->loop,
LIBUS_SOCKET_WRITABLE | (s->flags.is_paused ? 0 : LIBUS_SOCKET_READABLE));
LIBUS_SOCKET_WRITABLE | ((s->flags.is_paused || s->read_eof) ? 0 : LIBUS_SOCKET_READABLE));
}

int us_socket_write2(struct us_socket_t *s, const char *header, int header_length, const char *payload, int payload_length) {
Expand Down Expand Up @@ -457,6 +457,7 @@ struct us_socket_t *us_socket_from_fd(struct us_socket_group_t *group, unsigned
s->flags.adopted = 0;
s->flags.last_write_failed = 0;
s->unclassified_send_failures = 0;
s->read_eof = 0;
s->connect_state = NULL;

/* We always use nodelay */
Expand Down Expand Up @@ -859,11 +860,12 @@ void us_socket_resume(struct us_socket_t *s) {
// closed cannot be resumed
if (us_socket_is_closed(s)) return;

int readable = s->read_eof ? 0 : LIBUS_SOCKET_READABLE;
if (us_socket_is_shut_down(s)) {
// we already sent FIN so we resume only readable side we are read-only
us_poll_change(&s->p, s->group->loop, LIBUS_SOCKET_READABLE);
us_poll_change(&s->p, s->group->loop, readable);
return;
}
// we are readable and writable so we resume everything
us_poll_change(&s->p, s->group->loop, LIBUS_SOCKET_READABLE | LIBUS_SOCKET_WRITABLE);
us_poll_change(&s->p, s->group->loop, readable | LIBUS_SOCKET_WRITABLE);
}
4 changes: 3 additions & 1 deletion src/bun_core/tty.rs
Original file line number Diff line number Diff line change
Expand Up @@ -89,7 +89,9 @@ impl RawModeGuard {
impl Drop for RawModeGuard {
#[inline]
fn drop(&mut self) {
let _ = self.state.set_mode(self.fd, Mode::Normal, SetAttrWhen::Drain);
let _ = self
.state
.set_mode(self.fd, Mode::Normal, SetAttrWhen::Drain);
Comment thread
cirospaciari marked this conversation as resolved.
}
}

Expand Down
6 changes: 5 additions & 1 deletion src/md/ansi_renderer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2535,7 +2535,11 @@ fn probe_kitty_graphics() -> bool {
Err(_) => return false,
};
let mut tty_state = bun_core::tty::State::new();
let _ = tty_state.set_mode(0, bun_core::tty::Mode::Raw, bun_core::tty::SetAttrWhen::Drain);
let _ = tty_state.set_mode(
0,
bun_core::tty::Mode::Raw,
bun_core::tty::SetAttrWhen::Drain,
);
let _restore = scopeguard::guard((saved_termios, tty_state), |(saved, mut state)| {
if bun_sys::posix::tcsetattr(0, bun_sys::posix::TCSA::Now, &saved).is_err() {
let _ = state.set_mode(
Expand Down
1 change: 1 addition & 0 deletions src/parsers/yaml.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3311,6 +3311,7 @@ impl CollectionData for E::Object {
}

impl<'i, Enc: Encoding> Parser<'i, Enc> {
#[allow(clippy::needless_pass_by_value)] // must_use token: binding consumes the anchor
fn bind_anchor(&mut self, anchor: PendingAnchor, node: Expr) -> Result<(), AllocError> {
self.anchors
.put(Enc::key_bytes(anchor.name.slice(self.input)), node)
Expand Down
2 changes: 1 addition & 1 deletion src/runtime/bake/bake_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,8 @@ use core::ptr::NonNull;
use bun_alloc::Arena; // = bumpalo::Bump
use bun_collections::ArrayHashMap;
use bun_core::Output;
use bun_jsc::{JSGlobalObject, JSValue, JsError, JsResult, ZigStringSlice};
use bun_core::{ZStr, strings};
use bun_jsc::{JSGlobalObject, JSValue, JsError, JsResult, ZigStringSlice};
use bun_options_types::schema as bun_schema;
use bun_paths::{self as paths, PathBuffer};

Expand Down
5 changes: 4 additions & 1 deletion src/runtime/server/server_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -492,7 +492,10 @@ pub mod BunInfo {
// `JSON.toAST(allocator, BunInfo, info)` — hand-expanded:
let platform_props = bun_alloc::AstAlloc::vec_from_iter([
prop(b"os", str_expr(os_tag_name(info.platform.os))),
prop(b"arch", str_expr(arch_tag_name(bun_core::Environment::ARCH))),
prop(
b"arch",
str_expr(arch_tag_name(bun_core::Environment::ARCH)),
),
prop(b"version", str_expr(info.platform.version)),
]);
let platform_expr = Expr::init(
Expand Down
4 changes: 4 additions & 0 deletions src/runtime/socket/uws_handlers.rs
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,10 @@
//! old `NewSocketHandler.configure`/`unsafeConfigure` machinery, which built
//! the same trampolines at runtime per `us_socket_context_t`.

// `swallow()` uniformly sinks `()` and `Result<(), E>` consumers; passing the
// unit-returning ones through it is the point, not an accident.
#![allow(clippy::unit_arg)]

use bun_ptr::ThisPtr;
use core::ffi::{c_int, c_void};
use core::ptr::NonNull;
Expand Down
195 changes: 195 additions & 0 deletions test/js/bun/net/half-open-peer-reset.mjs
Original file line number Diff line number Diff line change
@@ -0,0 +1,195 @@
// Body of test/js/bun/test/parallel/test-net-half-open-peer-reset-*: a backpressured socket whose peer stops reading, FINs, then RSTs must close (with `end` at most once), not spin.
import { tls as tlsCert } from "../../../harness";
import { once } from "node:events";
import http from "node:http";
import https from "node:https";

const big = Buffer.alloc(4 * 1024 * 1024, 0x78);
const { key, cert } = tlsCert;
const upgradeReq = "GET / HTTP/1.1\r\nHost: x\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n";
const getReq = "GET / HTTP/1.1\r\nHost: x\r\n\r\n";

export async function run(mode) {
const issued = Promise.withResolvers(); // victim has queued more than the kernel will buffer
const ended = Promise.withResolvers(); // victim saw the peer's FIN (half-open modes)
const closed = Promise.withResolvers(); // victim socket/request was torn down
let endCount = 0;
const onEnd = () => {
if (++endCount > 1) fail("end delivered " + endCount + " times");
ended.resolve();
};
const fail = why => {
console.error("FAIL", why);
process.exit(1);
};

// The peer is a raw Bun socket so FIN (shutdown) and RST (terminate) are exactly what hits the wire, in that order.
async function peer(port, { preface, useTls, waitForEnd }) {
const opened = Promise.withResolvers();
const s = await Bun.connect({
hostname: "127.0.0.1",
port,
allowHalfOpen: true,
tls: useTls ? { rejectUnauthorized: false } : false,
socket: {
open: s => (useTls ? undefined : opened.resolve(s)),
handshake: s => opened.resolve(s),
data() {},
end() {},
error() {},
close() {},
},
});
await opened.promise;
if (preface) s.write(preface);
s.flush();
s.pause();
await issued.promise;
s.shutdown();
if (waitForEnd) await Promise.race([ended.promise, closed.promise]);
// The FIN→RST gap is where the bug lived; an immediate RST can overtake the FIN and just read as ECONNRESET.
await Bun.sleep(50);
s.terminate();
const deadline = setTimeout(() => fail("still open 5s after peer reset"), 5000);
closed.promise.then(() => clearTimeout(deadline));
}

switch (mode) {
case "bun-listen":
case "bun-listen-tls": {
const useTls = mode === "bun-listen-tls";
using server = Bun.listen({
hostname: "127.0.0.1",
port: 0,
allowHalfOpen: true,
tls: useTls ? { key, cert } : undefined,
socket: {
open(s) {
s.write(big);
issued.resolve();
},
data() {},
drain(s) {
s.write(big);
},
end: onEnd,
error() {},
close: () => closed.resolve(""),
},
});
await peer(server.port, { useTls, waitForEnd: true });
await closed.promise;
break;
}
case "node-http-upgrade":
case "node-https-upgrade": {
const useTls = mode === "node-https-upgrade";
await using server = (useTls ? https : http).createServer(useTls ? { key, cert } : {}, (_q, res) => res.end());
server.on("upgrade", (_req, socket) => {
socket.on("error", () => {});
socket.on("close", () => closed.resolve(""));
socket.on("end", () => {
onEnd();
socket.end();
});
socket.write("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: websocket\r\n\r\n");
for (let i = 0; i < 8; i++) socket.write(big);
issued.resolve();
});
server.listen(0, "127.0.0.1");
await once(server, "listening");
await peer(server.address().port, { preface: upgradeReq, useTls, waitForEnd: true });
await closed.promise;
break;
}
case "node-http-response":
case "node-https-response": {
const useTls = mode === "node-https-response";
await using server = (useTls ? https : http).createServer(useTls ? { key, cert } : {}, (_req, res) => {
res.on("close", () => closed.resolve("res"));
res.writeHead(200);
for (let i = 0; i < 8; i++) res.write(big);
issued.resolve();
});
server.listen(0, "127.0.0.1");
await once(server, "listening");
await peer(server.address().port, { preface: getReq, useTls, waitForEnd: false });
await closed.promise;
break;
}
case "bun-serve":
case "bun-serve-tls": {
const useTls = mode === "bun-serve-tls";
using server = Bun.serve({
port: 0,
hostname: "127.0.0.1",
idleTimeout: 0,
tls: useTls ? { key, cert } : undefined,
fetch(req) {
req.signal.addEventListener("abort", () => closed.resolve("abort"));
let i = 0;
return new Response(
new ReadableStream({
pull(ctrl) {
if (i++ < 64) ctrl.enqueue(big);
else ctrl.close();
issued.resolve();
},
cancel: () => closed.resolve("cancel"),
}),
);
},
});
await peer(server.port, { preface: getReq, useTls, waitForEnd: false });
await closed.promise;
break;
}
case "fetch-upload": {
// Roles flipped: the server is the peer that stops reading / FINs / RSTs; fetch() is the victim with a pending body.
using server = Bun.listen({
hostname: "127.0.0.1",
port: 0,
allowHalfOpen: true,
socket: {
async open(s) {
s.pause();
await issued.promise;
s.shutdown();
await Bun.sleep(50);
s.terminate();
const deadline = setTimeout(() => fail("fetch still pending 5s after peer reset"), 5000);
closed.promise.then(() => clearTimeout(deadline));
},
data() {},
end() {},
error() {},
close() {},
},
});
let i = 0;
fetch(`http://127.0.0.1:${server.port}/`, {
method: "POST",
duplex: "half",
body: new ReadableStream({
pull(c) {
if (i++ < 64) c.enqueue(big);
else c.close();
issued.resolve();
},
}),
}).then(
r =>
r.text().then(
() => closed.resolve("resolved"),
e => closed.resolve("body rejected " + e?.code),
),
e => closed.resolve("rejected " + e?.code),
);
await closed.promise;
break;
}
default:
fail("unknown mode " + mode);
}
console.log("closed", await closed.promise);
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
// https://github.com/oven-sh/bun/pull/37077
import { run } from "../../net/half-open-peer-reset.mjs";
await run("bun-listen-tls");
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
// https://github.com/oven-sh/bun/pull/37077
import { run } from "../../net/half-open-peer-reset.mjs";
await run("bun-listen");
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
// https://github.com/oven-sh/bun/pull/37077
import { run } from "../../net/half-open-peer-reset.mjs";
await run("bun-serve-tls");
Loading
Loading