Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
22 changes: 14 additions & 8 deletions packages/bun-usockets/src/bsd.c
Original file line number Diff line number Diff line change
Expand Up @@ -1444,8 +1444,8 @@ LIBUS_SOCKET_DESCRIPTOR bsd_create_listen_socket_unix(const char *path, size_t l

/* Receive-path options every UDP socket needs, whether freshly created or
* adopted from an existing fd: destination-address and TOS reporting for
* recvmmsg, Windows ICMP-reset suppression, and Linux IP_RECVERR. */
static void bsd_apply_udp_recv_options(LIBUS_SOCKET_DESCRIPTOR fd, int family) {
* recvmmsg, Windows ICMP-reset suppression, and Linux IP_RECVERR (opt-in). */
static void bsd_apply_udp_recv_options(LIBUS_SOCKET_DESCRIPTOR fd, int family, int options) {
/* We need destination address for udp packets in both ipv6 and ipv4 */

/* On FreeBSD this option seems to be called like so */
Expand Down Expand Up @@ -1488,15 +1488,21 @@ static void bsd_apply_udp_recv_options(LIBUS_SOCKET_DESCRIPTOR fd, int family) {
#if defined(__linux__)
/* IP_RECVERR/IPV6_RECVERR queues ICMP errors on the socket's error queue
* for on_recv_error to drain. libuv gates this on UV_UDP_LINUX_RECVERR
* (Node's dgram never passes it); Bun opts in for sockets it creates. */
* (Node's dgram never passes it). Opt-in only: on a shared unconnected
* socket (the HTTP/3 fetch client) it also makes a queued ICMP fail the
* next send to a different, live peer. */
if (options & LIBUS_UDP_LINUX_RECVERR) {
#ifdef IP_RECVERR
setsockopt(fd, IPPROTO_IP, IP_RECVERR, &enabled, sizeof(enabled));
setsockopt(fd, IPPROTO_IP, IP_RECVERR, &enabled, sizeof(enabled));
#endif
#ifdef IPV6_RECVERR
if (family == AF_INET6) {
setsockopt(fd, IPPROTO_IPV6, IPV6_RECVERR, &enabled, sizeof(enabled));
}
if (family == AF_INET6) {
setsockopt(fd, IPPROTO_IPV6, IPV6_RECVERR, &enabled, sizeof(enabled));
}
#endif
}
#else
(void) options;
#endif
}

Expand Down Expand Up @@ -1631,7 +1637,7 @@ LIBUS_SOCKET_DESCRIPTOR bsd_create_udp_socket(const char *host, int port, int op
}
#endif

bsd_apply_udp_recv_options(listenFd, listenAddr->ai_family);
bsd_apply_udp_recv_options(listenFd, listenAddr->ai_family, options);

/* We bind here as well */
if (bind(listenFd, listenAddr->ai_addr, (socklen_t) listenAddr->ai_addrlen)) {
Expand Down
5 changes: 5 additions & 0 deletions packages/bun-usockets/src/libusockets.h
Original file line number Diff line number Diff line change
Expand Up @@ -155,6 +155,11 @@ enum {
* Safe for HTTP/TLS where the client always sends first; do not use for protocols where
* the server sends the first bytes. */
LIBUS_LISTEN_DEFER_ACCEPT = 64,
/* Enable IP_RECVERR on a UDP socket so ICMP errors land on the error
* queue for on_recv_error to drain. Off by default: on a shared
* unconnected socket it also makes the next send fail for a datagram
* bound to a different, live peer. */
LIBUS_UDP_LINUX_RECVERR = 128,
};

/* Library types publicly available */
Expand Down
36 changes: 29 additions & 7 deletions packages/bun-usockets/src/quic.c
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
#include "quic.h"

#include "internal/internal.h"
#include "internal/fault_inject.h"
#if defined(_WIN32) && !defined(WIN32)
/* lsquic.h gates on WIN32 (not _WIN32) to pick <vc_compat.h> over <sys/uio.h>. */
#define WIN32 1
Expand Down Expand Up @@ -236,8 +237,11 @@ static int us_quic_send_one(LIBUS_SOCKET_DESCRIPTOR fd, const struct lsquic_out_
msg.msg_namelen = sa_len(spec->dest_sa);
msg.msg_iov = spec->iov;
msg.msg_iovlen = spec->iovlen;
ssize_t r;
do { r = sendmsg(fd, &msg, 0); } while (r < 0 && errno == EINTR);
ssize_t r = 0; int unused = 0;
if (!US_FAULT_CHECK(US_FAULT_SENDMSG, fd, r, unused)) {
do { r = sendmsg(fd, &msg, 0); } while (r < 0 && errno == EINTR);
}
(void) unused;
return r < 0 ? -1 : 1;
#endif
}
Expand Down Expand Up @@ -269,23 +273,41 @@ static int us_quic_packets_out(void *out_ctx, const struct lsquic_out_spec *spec
k++;
}
int r;
do { r = sendmmsg(fd, mm, k, 0); } while (r < 0 && errno == EINTR);
if (r < 0) break;
sent += (unsigned) r;
{
ssize_t injected = 0; int unused = 0;
if (US_FAULT_CHECK(US_FAULT_SENDMSG, fd, injected, unused)) {
r = (int) injected;
} else {
do { r = sendmmsg(fd, mm, k, 0); } while (r < 0 && errno == EINTR);
}
(void) injected; (void) unused;
}
/* sendmmsg(2) BUGS: on a short return the error code is lost and the
* caller is expected to retry starting at the first failed message.
* udp(7): an unconnected socket surfaces async ICMP from an earlier
* datagram on the next send — on the shared client socket that means
* a packet to a live peer can fail mid-batch with an error that
* belongs to a prior dead peer. So loop instead of breaking; r >= 1
* here so `sent` advances and the retry's first message either
* consumes the stale error (returns -1, handled below) or succeeds. */
* consumes the stale error (returns -1, handled below) or succeeds.
* The r < 0 path gets one retry for the same reason: the failing
* read cleared sk_err, so the retry sends cleanly unless this is
* real backpressure. EAGAIN/ENOBUFS (send buffer full) stays a
* break — that's the backpressure lsquic's pause is for. */
if (r < 0 && !(errno == EAGAIN || errno == EWOULDBLOCK || errno == ENOBUFS)) {
do { r = sendmmsg(fd, mm, k, 0); } while (r < 0 && errno == EINTR);
}
if (r < 0) break;
sent += (unsigned) r;
}
#else
for (; sent < n; sent++) {
us_quic_listen_socket_t *ls = (us_quic_listen_socket_t *) specs[sent].peer_ctx;
if (!ls->udp) { errno = EBADF; break; }
if (us_quic_send_one(us_poll_fd((struct us_poll_t *) ls->udp), &specs[sent]) < 0) break;
if (us_quic_send_one(us_poll_fd((struct us_poll_t *) ls->udp), &specs[sent]) < 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK || errno == ENOBUFS) break;
if (us_quic_send_one(us_poll_fd((struct us_poll_t *) ls->udp), &specs[sent]) < 0) break;
}
}
#endif

Expand Down
4 changes: 4 additions & 0 deletions packages/bun-usockets/src/udp.c
Original file line number Diff line number Diff line change
Expand Up @@ -184,6 +184,10 @@ struct us_udp_socket_t *us_create_udp_socket(
void *user
) {

/* IP_RECVERR is only useful when there is an on_recv_error handler to
* drain the error queue; without one it only poisons subsequent sends. */
if (recv_error_cb) flags |= LIBUS_UDP_LINUX_RECVERR;
Comment thread
coderabbitai[bot] marked this conversation as resolved.
else flags &= ~LIBUS_UDP_LINUX_RECVERR;
LIBUS_SOCKET_DESCRIPTOR fd = bsd_create_udp_socket(host, port, flags, err);
if (fd == LIBUS_SOCKET_ERROR) {
return 0;
Expand Down
34 changes: 34 additions & 0 deletions patches/lsquic/requeue-unsent-coalesced.patch
Original file line number Diff line number Diff line change
@@ -0,0 +1,34 @@
send_batch: requeue every packet in an unsent coalesced datagram

`off` is `unsigned`, so when the first unsent spec in a batch is at
index 0 and coalesces multiple packets, `&batch->packets[off - 1]`
indexes with UINT_MAX and `end` lands far past the array. The
`--packet_out > end` condition is then false after the first iteration
and only the last packet of the coalesced group is returned to the
connection; the earlier ones (typically the INIT ACK and the HSK
CRYPTO carrying the client Finished) are silently dropped and never
retransmitted, so the peer never completes the handshake.

Rewriting the bounds as [off, off+count) avoids the underflow while
preserving the reverse iteration order that send_ctl_sched_prepend
relies on.

--- a/src/liblsquic/lsquic_engine.c
+++ b/src/liblsquic/lsquic_engine.c
@@ -2739,12 +2739,12 @@
off = batch->pack_off[i];
count = batch->outs[i].iovlen;
assert(count > 0);
- packet_out = &batch->packets[off + count - 1];
- end = &batch->packets[off - 1];
+ packet_out = &batch->packets[off + count];
+ end = &batch->packets[off];
do
batch->conns[i]->cn_if->ci_packet_not_sent(batch->conns[i],
- *packet_out);
- while (--packet_out > end);
+ *--packet_out);
+ while (packet_out > end);
if (!(batch->conns[i]->cn_flags & (LSCONN_COI_ACTIVE|LSCONN_EVANESCENT)))
coi_reactivate(sb_ctx->conns_iter, batch->conns[i]);
}
1 change: 1 addition & 0 deletions scripts/build/deps/lsquic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -106,6 +106,7 @@ export const lsquic: Dependency = {
"patches/lsquic/allow-no-sni.patch",
"patches/lsquic/skip-priority-walk.patch",
"patches/lsquic/disable-gquic.patch",
"patches/lsquic/requeue-unsent-coalesced.patch",
],

fetchDeps: ["zlib", "lshpack", "lsqpack", "boringssl"],
Expand Down
15 changes: 12 additions & 3 deletions test/js/bun/http/serve-protocols.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,7 +75,12 @@ function fixtureFor(serve: object) {
},
});
console.error("PORT=" + server.port);
process.stdin.on("data", () => {});
// Graceful stop on stdin close so the client's pooled QUIC session sees
// CONNECTION_CLOSE. An abrupt kill leaves the client retransmitting to
// an unbound port, and that ICMP noise contends with every other
// connection on the shared HTTP/3 client engine.
process.stdin.on("end", () => { server.stop(true); setTimeout(() => process.exit(0), 50); });
process.stdin.resume();
`;
}

Expand Down Expand Up @@ -114,8 +119,12 @@ async function withServer(serve: object, fn: (origin: string) => Promise<void>)
await fn(`127.0.0.1:${port}`);
} finally {
proc.stdin?.end();
proc.kill();
await proc.exited;
const killTimer = setTimeout(() => proc.kill(), 500);
try {
await proc.exited;
} finally {
clearTimeout(killTimer);
}
}
}

Expand Down
131 changes: 131 additions & 0 deletions test/js/web/fetch/fetch-http3-syscall-fault.test.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,131 @@
/**
* HTTP/3 fetch under injected UDP send faults. Exercises lsquic's
* packets_out short-return path, which is otherwise only reachable when the
* UDP send buffer is genuinely full or an ICMP from a dead peer is queued on
* the shared socket.
*
* The first case pins the lsquic send_batch requeue-underflow patch: lsquic
* coalesces the client's INIT-ACK, HSK CRYPTO (TLS Finished) and a SHORT
* packet into one datagram with pack_off[0]==0 and iovlen>1. If packets_out
* returns 0 for that spec, the unpatched requeue loop computed
* &batch->packets[off - 1] with unsigned off and only returned the last
* packet of the group to the connection; the Finished was silently dropped
* and the server could never complete the handshake.
*/
import { socketFaultInjection as fault } from "bun:internal-for-testing";
import { afterEach, expect, test } from "bun:test";
import { bunEnv, bunExe, isWindows, tempDir, tls } from "harness";

const skip = !fault.available() || isWindows;

afterEach(() => fault.clear());

async function withServer(fn: (port: number) => Promise<void>) {
using dir = tempDir("h3-fault", {
"server.mjs": `
const server = Bun.serve({
port: 0, hostname: "127.0.0.1",
...${JSON.stringify({ tls, http3: true, http1: false })},
fetch: () => new Response("ok"),
});
console.error("PORT=" + server.port);
process.stdin.on("end", () => { server.stop(true); setTimeout(() => process.exit(0), 50); });
process.stdin.resume();
`,
});
const proc = Bun.spawn({
cmd: [bunExe(), "server.mjs"],
cwd: String(dir),
env: bunEnv,
stdout: "ignore",
stderr: "pipe",
stdin: "pipe",
});
let port = 0;
let buf = "";
for await (const chunk of proc.stderr) {
buf += new TextDecoder().decode(chunk);
const m = buf.match(/PORT=(\d+)/);
if (m) {
port = Number(m[1]);
break;
}
if (buf.length > 4096) break;
}
if (!port) {
proc.kill();
await proc.exited;
throw new Error("server did not report a port:\n" + buf);
}
try {
await fn(port);
} finally {
proc.stdin?.end();
const killTimer = setTimeout(() => proc.kill(), 500);
try {
await proc.exited;
} finally {
clearTimeout(killTimer);
}
}
}

const h3 = (port: number, init: RequestInit = {}) =>
fetch(`https://127.0.0.1:${port}/`, {
...init,
protocol: "http3",
tls: { rejectUnauthorized: false },
signal: AbortSignal.timeout(8000),
} as RequestInit);

test.skipIf(skip)("EAGAIN on the coalesced handshake datagram is requeued and the fetch completes", async () => {
await withServer(async port => {
// Arm an EAGAIN on the second UDP send the client engine makes. The
// first is the padded Initial (CRYPTO ClientHello). The second is the
// response to the server's flight: an INIT ACK coalesced with the HSK
// CRYPTO (Finished) and a SHORT NEW_CONNECTION_ID, i.e. the
// pack_off[0]==0, iovlen>1 spec whose requeue the patch fixes. The
// retry-once in us_quic_packets_out is gated on non-EAGAIN, so EAGAIN
// reaches lsquic as a genuine 0-of-N return.
fault.set({ syscall: "sendmsg", action: "errno", errno: "EAGAIN", after: 1, repeat: 1 });

const res = await h3(port);
expect(await res.text()).toBe("ok");
expect(res.status).toBe(200);
});
});

test.skipIf(skip)(
"a non-backpressure send error on the first datagram recovers without stalling the engine",
async () => {
await withServer(async port => {
// ECONNREFUSED on the very first send is what a stale ICMP on the shared
// client socket looks like. The errno is remapped to EAGAIN for lsquic,
// the UDP poll is re-armed writable, and on_drain → send_unsent_packets
// resends once the single-shot fault is consumed.
fault.set({ syscall: "sendmsg", action: "errno", errno: "ECONNREFUSED", after: 0, repeat: 1 });

const res = await h3(port);
expect(await res.text()).toBe("ok");
expect(res.status).toBe(200);
});
},
);

test.skipIf(skip)(
"repeated EAGAIN over several loop iterations recovers via on_drain without stalling the fetch",
async () => {
await withServer(async port => {
// Fail the first handful of sends with EAGAIN. Each failure re-arms the
// UDP poll's writable interest; on_drain → send_unsent_packets runs on
// the next iteration, so progress resumes as soon as the rule disarms.
// The 8s abort is well above lsquic's one-second resume_sending_at
// failsafe, so the only way to time out is an engine-level stall.
fault.set({ syscall: "sendmsg", action: "errno", errno: "EAGAIN", after: 0, repeat: 5 });

const res = await h3(port);
expect(await res.text()).toBe("ok");
expect(res.status).toBe(200);
});
},
);
Loading