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
461 changes: 248 additions & 213 deletions packages/bun-usockets/src/crypto/openssl.c

Large diffs are not rendered by default.

19 changes: 14 additions & 5 deletions packages/bun-usockets/src/internal/internal.h
Original file line number Diff line number Diff line change
Expand Up @@ -253,6 +253,8 @@ void us_internal_ssl_socket_left_group(us_socket_r s);
struct us_socket_t *us_internal_ssl_on_open(us_socket_r s, int is_client, char *ip, int ip_length);
struct us_socket_t *us_internal_ssl_on_data(us_socket_r s, char *data, int length);
struct us_socket_t *us_internal_ssl_on_writable(us_socket_r s);
/* The socket's timeout fired: ends a close that still waits for queued ciphertext, else dispatches on_timeout. */
struct us_socket_t *us_internal_ssl_on_timeout(us_socket_r s);
struct us_socket_t *us_internal_ssl_on_close(us_socket_r s, int code, void *reason);
struct us_socket_t *us_internal_ssl_on_end(us_socket_r s);
int us_internal_ssl_is_low_prio(us_socket_r s);
Expand Down Expand Up @@ -313,20 +315,22 @@ struct us_socket_t {
* not finished. ssl_write_wants_read cannot tell: every pending handshake
* sets it. */
unsigned char ssl_write_parked : 1;
unsigned char ssl_read_wants_write : 1;
/* Sealed ciphertext of this socket waits for the kernel (openssl.c us_ssl_out_queue_t). */
unsigned char ssl_out_queued : 1;
unsigned char ssl_fatal_error : 1;
unsigned char ssl_is_server : 1;
/* If set, us_internal_ssl_on_data() first dispatches the still-encrypted
* bytes via us_dispatch_ssl_raw_tap() before feeding them to SSL_read.
* Used by Bun's `socket.upgradeTLS()` so the returned [raw, tls] pair's
* `raw` half can observe ciphertext (node:net Duplex.ondata semantics). */
unsigned char ssl_raw_tap : 1;
/* A graceful TLS shutdown arrived while batched ciphertext was still
* spilled (see ssl_flush_write_batch); the shutdown re-runs once the
* spill drains so those records are not cut off by our FIN/close_notify. */
/* A graceful TLS shutdown arrived while ciphertext was still queued; the
* shutdown (or only its FIN, when the queue holds the close_notify) runs
* from the writable event once the queue drains. */
unsigned char ssl_shutdown_after_spill : 1;
/* Same as ssl_shutdown_after_spill but for us_internal_ssl_close: the
* close re-runs from the writable event once the spill drains. */
* close re-runs from the writable event once the queue drains, or from the
* socket's timeout when the peer never takes it. */
unsigned char ssl_close_after_spill : 1;
/* The plaintext EOF (peer close_notify or the raw TCP FIN behind it) was
* already dispatched to the user layer; both EOF paths can fire for one
Expand Down Expand Up @@ -361,6 +365,11 @@ struct us_socket_t {
* inside a handshake callback must still RST, not FIN, when it is finally
* performed). */
unsigned char ssl_pending_close_code : 2;
/* The deferred close armed the socket's timeout itself (its holder had none). */
unsigned char ssl_close_timeout_armed : 1;
/* us_internal_ssl_close sent the close_notify and waits for the peer's: the
* socket's timeout ends that wait too. */
unsigned char ssl_close_awaits_peer : 1;
/* Consecutive send() failures with an errno that is neither
* would-block/transient nor a known peer-gone error (see
* us_socket_write_check_error). Reset by any send that makes progress.
Expand Down
4 changes: 2 additions & 2 deletions packages/bun-usockets/src/libusockets.h
Original file line number Diff line number Diff line change
Expand Up @@ -698,8 +698,8 @@ int us_socket_is_ssl_handshake_finished(us_socket_r s) nonnull_fn_decl;
int us_socket_ssl_handshake_callback_has_fired(us_socket_r s) nonnull_fn_decl;
/* TLS ciphertext bytes already sealed for this socket and reported as
* written by us_socket_write(), still waiting on a writable event to reach
* the kernel (the loop-wide spill slot owned by this socket). 0 for
* plain-TCP sockets and for TLS sockets with nothing spilled. */
* the kernel (the socket's own queue). 0 for plain-TCP sockets and for TLS
* sockets with nothing queued. */
unsigned int us_socket_ssl_spill_pending(us_socket_r s) nonnull_fn_decl;

struct us_socket_t *us_socket_close(us_socket_r s, int code, void *reason) __attribute__((nonnull(1)));
Expand Down
18 changes: 14 additions & 4 deletions packages/bun-usockets/src/loop.c
Original file line number Diff line number Diff line change
Expand Up @@ -238,6 +238,14 @@ int us_loop_close_all_groups(struct us_loop_t *loop) {
}

/* This functions should never run recursively */
static void us_internal_dispatch_timeout(struct us_socket_t *s) {
if (s->ssl) {
us_internal_ssl_on_timeout(s);
} else {
us_dispatch_timeout(s);
}
}

void us_internal_timer_sweep(struct us_loop_t *loop) {
struct us_internal_loop_data_t *loop_data = &loop->data;
/* For all socket groups in this loop */
Expand Down Expand Up @@ -272,7 +280,7 @@ void us_internal_timer_sweep(struct us_loop_t *loop) {

if (short_ticks == s->timeout) {
s->timeout = 255;
us_dispatch_timeout(s);
us_internal_dispatch_timeout(s);
}
/* An owner must not deinit the embedding group from a timeout handler
* (see us_socket_group_deinit). Survive one that closed every socket
Expand Down Expand Up @@ -317,7 +325,7 @@ void us_internal_timer_sweep(struct us_loop_t *loop) {
unsigned char long_stamp = s->group->long_timestamp;
if (stamp == s->timeout) {
s->timeout = 255;
us_dispatch_timeout(s);
us_internal_dispatch_timeout(s);
if (loop_data->low_prio_iterator != s) continue;
}
if (long_stamp == s->long_timeout) {
Expand Down Expand Up @@ -622,8 +630,10 @@ void us_internal_dispatch_ready_poll(struct us_poll_t *p, int error, int eof, in
return;
}

/* If we have no failed write or if we shut down, then stop polling for more writable */
if (!s->flags.last_write_failed || us_socket_is_shut_down(s)) {
/* If we have no failed write or if we shut down, then stop polling for more writable.
* A TLS socket that sent its close_notify counts as shut down while the alert can
* still wait in its queue: that one keeps polling. */
if (!s->flags.last_write_failed || (us_socket_is_shut_down(s) && !us_socket_ssl_spill_pending(s))) {
us_poll_change(&s->p, loop, us_poll_events(&s->p) & LIBUS_SOCKET_READABLE);
} else {
#ifdef LIBUS_USE_KQUEUE
Expand Down
2 changes: 1 addition & 1 deletion packages/bun-usockets/src/socket.c
Original file line number Diff line number Diff line change
Expand Up @@ -847,7 +847,7 @@ void us_socket_resume(struct us_socket_t *s) {
if (us_socket_is_closed(s)) return;

int events = s->read_eof ? 0 : LIBUS_SOCKET_READABLE;
if (!us_socket_is_shut_down(s)) {
if (!us_socket_is_shut_down(s) || us_socket_ssl_spill_pending(s)) {
// still writable: a FIN of ours would have left the socket read-only
events |= LIBUS_SOCKET_WRITABLE;
}
Expand Down
6 changes: 3 additions & 3 deletions packages/bun-uws/src/AsyncSocket.h
Original file line number Diff line number Diff line change
Expand Up @@ -199,9 +199,9 @@ struct AsyncSocket {
}

/* Whether every byte handed to us_socket_write() has reached the kernel.
* For TLS, us_socket_write() can report a batch as written while its
* ciphertext still sits in the loop's spill slot (openssl.c
* ssl_flush_write_batch); the close-after-drain gates in HttpResponse /
* For TLS, us_socket_write() can report records as written while their
* ciphertext still sits in the socket's queue (openssl.c ssl_emit);
* the close-after-drain gates in HttpResponse /
* HttpContext must wait for that too. Kept separate from
* getBufferedAmount() so WebSocket's maxBackpressure policy and the
* JS-exposed bufferedAmount stay a plaintext count. */
Expand Down
2 changes: 1 addition & 1 deletion packages/bun-uws/src/HttpContext.h
Original file line number Diff line number Diff line change
Expand Up @@ -485,7 +485,7 @@ struct HttpContext {
/* node:http also queues behind responses that were dispatched but
* are not the connection's current response yet, and behind a
* response that has ended but not finished: its bytes are still
* in the outgoing buffer (or the TLS spill slot), and it owns the
* in the outgoing buffer (or the TLS queue), and it owns the
* connection (Node's socket._httpMessage) until they have been
* written out, with later responses queued behind it (Node's
* state.outgoing). A write or uncork can empty the buffer before
Expand Down
8 changes: 4 additions & 4 deletions test/js/bun/net/socket-syscall-fault.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,10 +67,10 @@ test.concurrent.skipIf(skip)(
// 2 rounds x (1 primer + 2 waves x 32).
opened: 130,
closed: 130,
// Every wave socket is closed mid-handshake by the unread-ciphertext
// guard, which reports the handshake as failed; the primers are
// closed by stop(true) and report nothing.
handshakeFailed: 128,
// Every socket is closed mid-handshake, which reports the handshake
// as failed: the wave sockets by the fatal alert the child sends, the
// primers when the child closes them.
handshakeFailed: 130,
handshakeOk: 0,
data: 0,
errors: 0,
Expand Down
213 changes: 213 additions & 0 deletions test/js/bun/net/socket.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1508,6 +1508,219 @@ it("TLS mid-read boundary dispatch: writing to another TLS socket from data() do
}
}, 60_000);

describe("TLS write() that ends short while another TLS socket is stalled", () => {
// write() reports how many bytes it took, and the caller may send anything
// next. So the peer must receive exactly the reported bytes, then the next
// write: no byte that write() did not report, and none missing.
//
// `stalled` stays under backpressure for the whole test: its peer stops
// reading, and it fills again on every drain. A TLS socket in that state
// changes how the other TLS sockets on the event loop send their records.
const STEP = 1024 * 1024;
const FIRST = 0x61;
const OTHER = 0x62;

// Writes `step` until a write is short. Returns the sum write() reported.
function fill(socket: Socket, step: Buffer) {
let total = 0;
for (let wrote = step.length; wrote === step.length && total < 64 * STEP; ) {
wrote = socket.write(step);
total += Math.max(wrote, 0);
}
return total;
}

// `next` is what the caller does after the short write: "end" ends the
// socket in the same tick; bytes are written once every reported byte has
// arrived (after setMaxSendFragment(fragment), when given), then the socket
// ends. Returns what write() reported for the first buffer, and what the
// peer received: its length, where the OTHER bytes start, and where the
// last FIRST byte is.
async function shortWriteThen(next: "end" | { bytes: Buffer; fragment?: number }) {
const step = Buffer.alloc(STEP, FIRST);
const stalledStep = Buffer.alloc(STEP, 0x7a);
const stalledFilled = Promise.withResolvers<void>();
const reportedArrived = Promise.withResolvers<void>();
const peerClosed = Promise.withResolvers<void>();
const received = { length: 0, otherStartsAt: -1, firstEndsAt: -1 };
let reported = -1;
let connections = 0;

const server = Bun.listen<{ stalled: boolean }>({
hostname: "127.0.0.1",
port: 0,
tls,
socket: {
open(peer) {
peer.data = { stalled: connections++ === 0 };
},
data(peer, chunk) {
if (peer.data.stalled) {
peer.pause();
return;
}
const otherStart = chunk.indexOf(OTHER);
if (otherStart !== -1 && received.otherStartsAt === -1) received.otherStartsAt = received.length + otherStart;
const firstEnd = chunk.lastIndexOf(FIRST);
if (firstEnd !== -1) received.firstEndsAt = received.length + firstEnd;
received.length += chunk.length;
if (received.length >= reported) reportedArrived.resolve();
},
close(peer) {
if (!peer.data.stalled) peerClosed.resolve();
},
error() {},
},
});

let stalled: Socket | undefined;
let socket: Socket | undefined;
let sent = -1;
const sendRest = (s: Socket, bytes: Buffer) => {
while (sent < bytes.length) {
const wrote = s.write(bytes.subarray(sent));
if (wrote <= 0) return;
sent += wrote;
}
s.end();
};
try {
stalled = await Bun.connect({
hostname: "127.0.0.1",
port: server.port,
tls: { ...tls, rejectUnauthorized: false },
socket: {
handshake(s) {
fill(s, stalledStep);
stalledFilled.resolve();
},
drain(s) {
fill(s, stalledStep);
},
data() {},
close() {},
error() {},
},
});
await stalledFilled.promise;

socket = await Bun.connect({
hostname: "127.0.0.1",
port: server.port,
tls: { ...tls, rejectUnauthorized: false },
socket: {
handshake(s) {
reported = fill(s, step);
if (next === "end") s.end();
},
drain(s) {
if (next !== "end" && sent !== -1) sendRest(s, next.bytes);
},
data() {},
close() {},
error() {},
},
});

if (next !== "end") {
// Every reported byte has arrived, so the socket is idle again and
// takes a write.
await reportedArrived.promise;
if (next.fragment) socket.setMaxSendFragment(next.fragment);
const wrote = socket.write(next.bytes);
expect(wrote).toBeGreaterThan(0);
sent = wrote;
sendRest(socket, next.bytes);
}
await peerClosed.promise;

return { reported, received };
} finally {
stalled?.terminate();
socket?.terminate();
server.stop(true);
}
}

it.each([
["a later write of other data arrives byte for byte", 256 * 1024],
["a later write that is smaller than one TLS record is sent", 100],
])("%s", async (_, length) => {
const { reported, received } = await shortWriteThen({ bytes: Buffer.alloc(length, OTHER) });
expect(received).toEqual({ length: reported + length, otherStartsAt: reported, firstEndsAt: reported - 1 });
});

it("the same stream continues after setMaxSendFragment()", async () => {
const length = 64 * 1024;
const { reported, received } = await shortWriteThen({ bytes: Buffer.alloc(length, FIRST), fragment: 512 });
expect(received).toEqual({ length: reported + length, otherStartsAt: -1, firstEndsAt: reported + length - 1 });
});

it("end() in the same tick sends every reported byte first", async () => {
const { reported, received } = await shortWriteThen("end");
expect(received).toEqual({ length: reported, otherStartsAt: -1, firstEndsAt: reported - 1 });
});
});

describe.concurrent("a TLS end() whose peer stopped reading closes at the socket's timeout", () => {
// end() closes the socket once the ciphertext that write() already counted
// has reached the kernel and the peer has answered the close_notify. The
// peer stops reading here, so that never happens, and the socket's timeout
// is what ends the wait.
it.each([
["with ciphertext the kernel did not take", true],
["with everything sent", false],
])(
"%s",
async (_, fillKernel) => {
const step = Buffer.alloc(1024 * 1024, "a");
const peerPaused = Promise.withResolvers<void>();
const closed = Promise.withResolvers<void>();
const server = Bun.listen({
hostname: "127.0.0.1",
port: 0,
tls,
socket: {
data(peer) {
peer.pause();
peerPaused.resolve();
},
close() {},
error() {},
},
});
let client: Socket | undefined;
try {
client = await Bun.connect({
hostname: "127.0.0.1",
port: server.port,
tls: { ...tls, rejectUnauthorized: false },
socket: {
handshake(s) {
if (!fillKernel) return void s.write("a");
for (let wrote = step.length; wrote === step.length; ) wrote = s.write(step);
},
close() {
closed.resolve();
},
data() {},
error() {},
},
});
await peerPaused.promise;
client.timeout(1);
client.end();
await closed.promise;
} finally {
client?.terminate();
server.stop(true);
}
},
// The timeout fires on a 4 second tick.
15_000,
);
});

describe.concurrent("TLS server: write() to the accepted socket from inside its own selection callback", () => {
// alpnCallback / serverName are the listener hooks node:tls's ALPNCallback /
// SNICallback go through. Both run from inside the read that is processing
Expand Down
Loading
Loading