Skip to content
Open
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
4 changes: 4 additions & 0 deletions docs/runtime/networking/udp.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,10 @@ const client = await Bun.udpSocket({});
client.send("Hello!", server.port, "127.0.0.1");
```

Bun reads at most 32 datagrams from a socket in one iteration of the event loop, as Node.js does. The operating system holds the rest in the receive buffer of the socket, and Bun reads them in the next iterations. Timers and other sockets keep running when a peer sends without pause.

When the receive buffer is full, the operating system drops new datagrams and does not tell Bun. To keep up with a fast sender, keep the `data` callback short and move long work to a [`Worker`](/runtime/workers).

### Connections

UDP has no concept of a connection, but many UDP exchanges (especially as a client) involve only one peer.
Expand Down
4 changes: 2 additions & 2 deletions packages/bun-usockets/src/bsd.c
Original file line number Diff line number Diff line change
Expand Up @@ -148,7 +148,7 @@ int bsd_sendmmsg(LIBUS_SOCKET_DESCRIPTOR fd, struct udp_sendbuf* sendbuf, int fl
int bsd_recvmmsg(LIBUS_SOCKET_DESCRIPTOR fd, struct udp_recvbuf *recvbuf, int flags, int max_packets) {
if (max_packets > LIBUS_UDP_RECV_COUNT) max_packets = LIBUS_UDP_RECV_COUNT;
#if defined(_WIN32)
for (int i = 0; i < LIBUS_UDP_RECV_COUNT; i++) {
for (int i = 0; i < max_packets; i++) {
while (1) {
socklen_t addr_len = sizeof(struct sockaddr_storage);
ssize_t ret = recvfrom(fd, recvbuf->buf + (size_t) i * LIBUS_UDP_MAX_SIZE,
Expand All @@ -169,7 +169,7 @@ int bsd_recvmmsg(LIBUS_SOCKET_DESCRIPTOR fd, struct udp_recvbuf *recvbuf, int fl
break;
}
}
return LIBUS_UDP_RECV_COUNT;
return max_packets;
#elif defined(__APPLE__)
if (Bun__doesMacOSVersionSupportSendRecvMsgX()) {
while (1) {
Expand Down
4 changes: 4 additions & 0 deletions packages/bun-usockets/src/internal/networking/bsd.h
Original file line number Diff line number Diff line change
Expand Up @@ -57,6 +57,10 @@ struct bsd_addr_t {
};

#define LIBUS_UDP_RECV_COUNT (LIBUS_RECV_BUFFER_LENGTH / LIBUS_UDP_MAX_SIZE)
/* The most datagrams that one readable event of a UDP socket hands over. It is
* libuv's count, and like libuv's it counts datagrams, not receive calls:
* https://github.com/libuv/libuv/blob/895cd04bf7cadd2f8993c0d6fc7e55eaaec60bc9/src/unix/udp.c#L309-L312 */
#define LIBUS_UDP_MAX_RECV_PER_EVENT 32

#ifdef __APPLE__
/*
Expand Down
20 changes: 15 additions & 5 deletions packages/bun-usockets/src/loop.c
Original file line number Diff line number Diff line change
Expand Up @@ -1024,13 +1024,22 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in
#else
const int run_recv = events & LIBUS_SOCKET_READABLE;
#endif
/* Reading until EAGAIN lets the peer decide when this loop ends: a
* sender at or above the rate on_data drains at, or one that
* answers what on_data sends, keeps the queue non-empty, and no
* timer, immediate or other poll runs until it stops. The poll is
* level-triggered on every backend, so what is left raises the
* next event. */
int recv_budget = LIBUS_UDP_MAX_RECV_PER_EVENT;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🔴 Operators of Bun.serve({ http3: true }) and node:quic servers get kernel-side packet loss under concurrent load that the base branch does not have. The cap at packages/bun-usockets/src/loop.c:1033 is per socket, and a QUIC listener is one socket for every connection, so 32 packets per loop iteration is the whole server's inbound ceiling. Fix: keep the bound where on_data runs user JS, but let QUIC sockets drain what the base drained, e.g. a per-socket budget that quic.c raises, a time-bounded continuation past 32, or a larger SO_RCVBUF on QUIC sockets. The PR's Downsides put the loss at 6.5 % to 16.2 % for 64 uploads; the author calls it accepted, but us_quic_udp_on_data runs no user code.

Why this was flagged

An HTTP/3 server started with Bun.serve({ http3: true }) listens through one QUIC UDP socket for all its connections (packages/bun-uws/src/Http3Context.h:100 calls us_quic_socket_context_listen, packages/bun-usockets/src/quic.c:1017 creates the single us_udp_socket_t); node:quic listen() does the…

Verification: normal — acknowledged in diff: the PR description's "Downsides" section states "An HTTP/3 listener is one socket. 64 concurrent uploads: kernel loss 6.5 % to 16.2 %, CPU per MiB +10 %" and the Notes leave "the form of the bound" as an open question; the mechanism the note describes is accurate and nothing in the code bounds it further. Triggering condition: an HTTP/3 (`Bun.serve({ http3:…

Copy link
Copy Markdown
Collaborator Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This cost is real and the Downsides and the Notes of the description state it with the numbers. The bound applies to the QUIC sockets on purpose, and the choice of its form is the open question for a maintainer at the end of the description.

Why the same bound: us_quic_udp_on_data runs no JS, but each packet goes through lsquic_engine_packet_in (decrypt, frame parse, stream buffering) before the event ends, and loop_post then runs process_conns over all of it. The starvation was reproduced on node:quic, not only on node:dgram: a 4 MiB body next to a loop that blocks 20 ms between turns read 181 to 555 packets in one turn on main (comment above, Aug 11). A budget that quic.c raises brings that back. Node's QUIC endpoint receives through libuv's uv_udp_t and so has the same count per event.

What reduces the loss on a listener without moving the bound: a larger receive buffer for the QUIC sockets (quic.c sets none, so the listener has the default 208 KiB for all its connections). The Notes list it as the follow-up that comes after the bound, because on main a larger buffer makes one event longer. The other form in the Notes, 32 as a floor and then a short time bound, also removes the ceiling for a cheap handler. Both are small changes on top of this one. I leave this thread open for the maintainer's answer on the form.

if (run_recv && !u->closed) {

do {
struct udp_recvbuf recvbuf;
bsd_udp_setup_recvbuf(&recvbuf, u->loop->data.recv_buf, LIBUS_RECV_BUFFER_LENGTH);
int npackets = bsd_recvmmsg(us_poll_fd(p), &recvbuf, MSG_DONTWAIT, u->shared_fd ? 1 : LIBUS_UDP_RECV_COUNT);
int max_packets = u->shared_fd ? 1 : LIBUS_UDP_RECV_COUNT;
if (max_packets > recv_budget) max_packets = recv_budget;
int npackets = bsd_recvmmsg(us_poll_fd(p), &recvbuf, MSG_DONTWAIT, max_packets);
if (npackets > 0) {
recv_budget -= npackets;
u->on_data(u, &recvbuf, npackets);
} else {
if (npackets == LIBUS_SOCKET_ERROR) {
Expand Down Expand Up @@ -1077,7 +1086,7 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in

break;
}
} while (!u->closed);
} while (!u->closed && recv_budget > 0);
Comment thread
robobun marked this conversation as resolved.
}

if (events & LIBUS_SOCKET_WRITABLE && !u->closed) {
Expand All @@ -1100,8 +1109,9 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in
* EAGAIN (which means the error queue is already drained,
* leaving a residual EPOLLERR). Otherwise the socket stays
* open so the user can keep sending/receiving after a
* transient ICMP error. */
if (error && !recv_error_surfaced && !recv_would_block_only && !u->closed) {
* transient ICMP error. A read that stopped on its budget never
* got as far as either answer: the next event decides. */
if (error && !recv_error_surfaced && !recv_would_block_only && recv_budget > 0 && !u->closed) {
us_udp_socket_close(u);
}
#else
Expand Down
68 changes: 68 additions & 0 deletions test/_util/loop-iterations.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,68 @@
// Helpers for fixtures that measure what one iteration of the event loop does:
// how many datagrams, reads or connections a socket hands over before the loop
// goes on to its timers, immediates and other sockets. For processes that run
// with bunEnv, which makes bun:internal-for-testing available.
import { getEventLoopStats } from "bun:internal-for-testing";

/**
* Counts events by the iteration of the event loop they happen in. Call
* `count()` from the callback to measure. `perIteration` has one entry for each
* iteration that counted something, in order.
*/
export function iterationCounter() {
const counts = new Map<number, number>();
let total = 0;
return {
count(events = 1) {
const { iteration } = getEventLoopStats();
counts.set(iteration, (counts.get(iteration) ?? 0) + events);
total += events;
},
get total() {
return total;
},
summary() {
const perIteration = [...counts.values()];
return { total, max: Math.max(0, ...perIteration), perIteration };
},
};
}

/**
* Lets the event loop iterate until `done()` holds. Resolves to false when it
* still does not hold after `seconds`, so that the caller can report what it
* has and not hang.
*/
export async function iterateUntil(done: () => boolean, seconds = 10) {
const deadline = performance.now() + seconds * 1000;
while (!done()) {
if (performance.now() > deadline) return false;
await new Promise<void>(resolve => setImmediate(resolve));
}
return true;
}

/**
* Runs `scenario` until a run is `usable`, at most `attempts` times and within
* `seconds` in all, and resolves to that run. A scenario that depends on what
* the kernel had queued before the loop polled cannot promise it on every
* platform and under every load, so a run that did not get there is set up
* again. When no run is usable the last one is the result, for the test to
* fail on.
*
* The scenario gets the seconds that are left of the budget. It passes them to
* `iterateUntil`, so that a run whose datagrams never arrive ends with the
* budget and not 20 deadlines later.
*/
export async function firstUsable<T extends object>(
scenario: (secondsLeft: number) => Promise<T>,
usable: (run: T) => boolean,
{ attempts = 20, seconds = 20 } = {},
) {
const deadline = performance.now() + seconds * 1000;
for (let attempt = 1; ; attempt++) {
Comment thread
robobun marked this conversation as resolved.
const secondsLeft = Math.max(0, (deadline - performance.now()) / 1000);
const run = await scenario(secondsLeft);
if (usable(run) || attempt === attempts || performance.now() > deadline) return { ...run, attempt };
}
}
58 changes: 58 additions & 0 deletions test/js/bun/udp/dgram-cluster-recv-budget-fixture.ts

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

62 changes: 61 additions & 1 deletion test/js/bun/udp/dgram.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ import { describe, expect, jest, test } from "bun:test";
import { createSocket } from "dgram";
import { Worker } from "node:worker_threads";

import { bunEnv, bunExe, bunRun, disableAggressiveGCScope, isWindows } from "harness";
import { bunEnv, bunExe, bunRun, disableAggressiveGCScope, isLinux, isWindows } from "harness";
import path from "path";
import { nodeDataCases } from "./testdata";

Expand Down Expand Up @@ -261,6 +261,16 @@ describe.skipIf(isWindows)("cluster", () => {
expect(stdout).toStartWith("ok: all 4 workers adopted and released the shared descriptor ");
expect(exitCode).toBe(0);
}, 40_000);

// A worker reads a shared descriptor one datagram per receive call. The
// bound of a readable event counts datagrams, not calls, so a backlog
// arrives 32 to an iteration of the event loop, as it does in Node.
test("a shared socket hands over at most 32 datagrams per readable event", async () => {
const { stdout, stderr, exitCode } = await runClusterFixture("dgram-cluster-recv-budget-fixture.ts", 20_000);
expect(stderr).toBe("");
expect(JSON.parse(stdout)).toMatchObject({ finished: true, total: 100, max: 32 });
expect(exitCode).toBe(0);
});
});

describe("after close()", () => {
Expand Down Expand Up @@ -915,6 +925,56 @@ for (const [kind, bind] of Object.entries(icmpBindModes)) {
});
}

// One readable event hands over at most 32 datagrams, as in Node, and the rest
// waits in the kernel for the next iteration of the event loop. Each scenario
// of the fixture queues its backlog before the loop polls and reports the
// datagrams of every iteration.
describe.concurrent("a readable event hands over at most 32 datagrams", () => {
const fixture = path.join(import.meta.dir, "udp-recv-budget-fixture.ts");

test("of a backlog of 100", async () => {
const result = await bunRun([fixture, "dgram-backlog"]);
expect(result).toSpawn();
expect(JSON.parse(result.stdout)).toMatchObject({ finished: true, total: 100, max: 32 });
});

test("of each socket", async () => {
const result = await bunRun([fixture, "dgram-two-sockets"]);
expect(result).toSpawn();
expect(JSON.parse(result.stdout)).toMatchObject({ finished: true, total: 200, max: 64, maxOfEach: [32, 32] });
});

// An event with EPOLLERR keeps its socket open when the receive failed with
// the error or ran dry. An event that stops at its bound saw neither. The
// flag can be stale then: a send from an earlier callback of the same
// iteration takes the pending error. The socket has to stay open.
describe.skipIf(!isLinux)("and keeps the socket open when a send took the pending error", () => {
// A descriptor bun did not create has no error queue to find a report on.
test("of an adopted descriptor", async () => {
const result = await bunRun([fixture, "residual-adopted"]);
expect(result).toSpawn();
expect(JSON.parse(result.stdout)).toMatchObject({
residual: true,
finished: true,
closed: false,
total: 40,
max: 32,
});
});

// The kernel queues the report only when the receive buffer has room.
test("of a socket with a full receive buffer", async () => {
const result = await bunRun([fixture, "residual-full-buffer"]);
expect(result).toSpawn();
const run = JSON.parse(result.stdout);
expect(run).toMatchObject({ residual: true, finished: true, closed: false, max: 32 });
// More than one event's worth arrived, and the buffer did overflow.
expect(run.total).toBeGreaterThan(32);
expect(run.total).toBeLessThan(run.sent);
});
});
});

// A worker that calls process.exit() from the FIRST 'message' of a batch
// leaves a TerminationException pending for the rest of that poll dispatch:
// the remaining on_data iterations, the drain, and any recv error must all
Expand Down
Loading
Loading