Skip to content
Closed
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
29 changes: 17 additions & 12 deletions packages/bun-usockets/src/udp.c
Original file line number Diff line number Diff line change
Expand Up @@ -52,21 +52,26 @@ int us_udp_socket_send(struct us_udp_socket_t *s, void** payloads, size_t* lengt

int total_sent = 0;
while (total_sent < num) {
int count = bsd_udp_setup_sendbuf(buf, LIBUS_SEND_BUFFER_LENGTH, payloads, lengths, addresses, num);
payloads += count;
lengths += count;
addresses += count;
num -= count;
// TODO nohang flag?
/* The send buffer holds a fixed number of mmsghdr slots, so a batch bigger
* than that takes several passes, each starting at the first unsent one. */
bsd_udp_setup_sendbuf(buf, LIBUS_SEND_BUFFER_LENGTH, payloads + total_sent, lengths + total_sent, addresses + total_sent, num - total_sent);
int sent = bsd_sendmmsg(fd, buf, MSG_DONTWAIT);
if (sent < 0) {
return sent;
}
total_sent += sent;
if (0 <= sent && sent < num) {
// if we couldn't send all packets, register a writable event so we can call the drain callback
if (sent <= 0) {
/* sendmmsg reports "not one datagram could be sent" as -1 with errno,
* the sendmsg/sendto fallbacks as 0. Only a full send buffer is
* recoverable; every other errno belongs to the caller. */
if (sent < 0 && !bsd_would_block()) {
return sent;
}
/* Out of send buffer space: ask for a writable event so the drain
* callback fires, and report the datagrams that did go out. */
us_poll_change((struct us_poll_t *) s, s->loop, LIBUS_SOCKET_READABLE | LIBUS_SOCKET_WRITABLE);
break;
Comment thread
robobun marked this conversation as resolved.
}
total_sent += sent;
/* A short count means the datagram at index `sent` failed and sendmmsg(2)
* dropped its errno. The next pass retries from that datagram, which
* surfaces the errno as -1 instead of hiding it in the count. */
}
return total_sent;
}
Expand Down
52 changes: 52 additions & 0 deletions test/js/bun/udp/udp_socket.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -364,6 +364,58 @@ describe("udpSocket()", () => {
}
}

// The send buffer only holds ~204 sendmmsg slots, so a bigger batch needs
// several passes. us_udp_socket_send used to subtract the batch size from
// `num` before testing it, which ended the loop after the first pass: an
// idle socket reported a short count (204 of 300) and never armed the
// writable poll, so the documented "if fewer were sent, wait for drain and
// resend the rest" protocol had nothing to wait for.
describe("sendMany sends every packet of a batch larger than one send buffer", () => {
test.each([205, 300, 500])("connected, %i packets", async count => {
using server = await udpSocket({ socket: { data() {} } });
using client = await udpSocket({
connect: { port: server.port, hostname: "127.0.0.1" },
});

const packets = Array.from({ length: count }, (_, i) => `packet ${i}`);
expect(client.sendMany(packets)).toBe(count);
});

test("unconnected, 300 packets", async () => {
using server = await udpSocket({ socket: { data() {} } });
using client = await udpSocket({});

const packets = Array.from({ length: 300 }, (_, i) => [`packet ${i}`, server.port, "127.0.0.1"]).flat();
expect(client.sendMany(packets)).toBe(300);
});
});

// sendmmsg(2) reports a datagram that fails mid-batch as a short count and
// drops its errno, so an oversized payload used to throw EMSGSIZE at index 0
// but come back as backpressure ("1 of 3 sent") anywhere after it.
describe("sendMany surfaces a per-packet error whatever its position", () => {
for (const connected of [true, false]) {
test.each([0, 1, 2])(`${connected ? "connected" : "unconnected"}, oversized packet at index %i`, async index => {
using server = await udpSocket({ socket: { data() {} } });
using client = await udpSocket(connected ? { connect: { port: server.port, hostname: "127.0.0.1" } } : {});

// 70000 > the 65507-byte maximum UDP payload, so the kernel always
// rejects this one datagram with EMSGSIZE.
const payloads: any[] = ["a", "b", "c"];
payloads[index] = Buffer.alloc(70000, 1);
const packets = connected ? payloads : payloads.flatMap(p => [p, server.port, "127.0.0.1"]);

let error: any;
try {
client.sendMany(packets);
} catch (e) {
error = e;
}
expect({ code: error?.code, syscall: error?.syscall }).toEqual({ code: "EMSGSIZE", syscall: "send" });
});
}
});

// send()/sendMany() capture a pointer into the payload's backing store and
// then run user JS (port `valueOf()`, address `toString()`, and for
// sendMany also array index getters on later iterations). That JS can
Expand Down
Loading