diff --git a/packages/bun-usockets/src/udp.c b/packages/bun-usockets/src/udp.c index fb12b7a52f5c..c1ebd7479abe 100644 --- a/packages/bun-usockets/src/udp.c +++ b/packages/bun-usockets/src/udp.c @@ -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; } + 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; } diff --git a/test/js/bun/udp/udp_socket.test.ts b/test/js/bun/udp/udp_socket.test.ts index 3217ccd41714..46a03b522f6e 100644 --- a/test/js/bun/udp/udp_socket.test.ts +++ b/test/js/bun/udp/udp_socket.test.ts @@ -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