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
31 changes: 31 additions & 0 deletions src/boringssl_sys/boringssl.rs
Original file line number Diff line number Diff line change
Expand Up @@ -813,6 +813,11 @@ pub struct BIO {
pub num_write: usize,
}

/// `#define BIO_CTRL_FLUSH 11` — the `cmd` of `BIO_flush`.
pub const BIO_CTRL_FLUSH: c_int = 11;
/// `#define BIO_TYPE_SOURCE_SINK 0x0400` — OR-ed into the type of a BIO that ends a chain.
pub const BIO_TYPE_SOURCE_SINK: c_int = 0x0400;

// ═══════════════════════════════════════════════════════════════════════════
// Additional opaque handles
// ═══════════════════════════════════════════════════════════════════════════
Expand Down Expand Up @@ -971,6 +976,32 @@ unsafe extern "C" {
pub fn BIO_new_mem_buf(buf: *const c_void, len: ossl_ssize_t) -> *mut BIO;
pub fn BIO_set_mem_eof_return(bio: *mut BIO, eof_value: c_int) -> c_int;

// ── Custom BIOs ──────────────────────────────────────────────────────
pub safe fn BIO_get_new_index() -> c_int;
/// Opaque: set its hooks with `BIO_meth_set_*`, never through the `BIO_METHOD` fields above.
pub fn BIO_meth_new(r#type: c_int, name: *const c_char) -> *mut BIO_METHOD;
pub fn BIO_meth_set_create(
method: *mut BIO_METHOD,
create_func: Option<unsafe extern "C" fn(*mut BIO) -> c_int>,
) -> c_int;
pub fn BIO_meth_set_write(
method: *mut BIO_METHOD,
write_func: Option<unsafe extern "C" fn(*mut BIO, *const c_char, c_int) -> c_int>,
) -> c_int;
pub fn BIO_meth_set_read(
method: *mut BIO_METHOD,
read_func: Option<unsafe extern "C" fn(*mut BIO, *mut c_char, c_int) -> c_int>,
) -> c_int;
pub fn BIO_meth_set_ctrl(
method: *mut BIO_METHOD,
ctrl_func: Option<unsafe extern "C" fn(*mut BIO, c_int, c_long, *mut c_void) -> c_long>,
) -> c_int;
pub fn BIO_set_data(bio: *mut BIO, ptr: *mut c_void);
pub fn BIO_get_data(bio: *mut BIO) -> *mut c_void;
pub fn BIO_set_init(bio: *mut BIO, init: c_int);
pub fn BIO_set_retry_read(bio: *mut BIO);
pub fn BIO_clear_retry_flags(bio: *mut BIO);

// ── RAND ─────────────────────────────────────────────────────────────
/// Fills `buf[0..len]` from BoringSSL's thread-local CTR-DRBG and returns 1.
/// In the event that sufficient random data can not be obtained, `abort`
Expand Down
5 changes: 1 addition & 4 deletions src/runtime/socket/UpgradedDuplex.rs
Original file line number Diff line number Diff line change
Expand Up @@ -348,10 +348,7 @@ impl UpgradedDuplex {
let staged = self.pending_data.replace(Vec::new());
self.reset_timeout();
// Feed in bounded slices rather than one concatenated buffer. Each JS
// chunk was originally delivered on its own; `receive_data` casts the
// length to `c_int` with a panicking `expect`, so handing it the sum of
// every chunk staged in the window would turn a large pre-start burst
// into a process abort. Re-check the engine each round: BoringSSL can
// chunk was originally delivered on its own. Re-check the engine each round: BoringSSL can
// re-enter and tear it down partway through, and `teardown()` neuters
// in place (frees the SSL, keeps the Option `Some`), so the live
// signal is the SSL handle, not the Option.
Expand Down
3 changes: 1 addition & 2 deletions src/runtime/socket/socket_body.rs
Original file line number Diff line number Diff line change
Expand Up @@ -155,8 +155,7 @@ extern "C" fn select_alpn_callback(
// handler, the selection's `toString`, the scope's checkpoint), and
// restore it on every path back to BoringSSL — after the scope guard
// below has exited, since this guard is declared first. Connected
// usockets only: UpgradedDuplex/Pipe own mem BIOs whose BIO_get_data
// is a BUF_MEM*, not loop_ssl_data.
// usockets only: the BIO of an UpgradedDuplex/Pipe SSL does not hold loop_ssl_data.
const LOOP_STATE_SLOTS: usize = 6; // US_SSL_LOOP_STATE_SLOTS
debug_assert_eq!(
LOOP_STATE_SLOTS as core::ffi::c_int,
Expand Down
374 changes: 237 additions & 137 deletions src/uws/lib.rs

Large diffs are not rendered by default.

58 changes: 58 additions & 0 deletions test/js/bun/http/proxy.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -1117,6 +1117,64 @@ test("HTTPS over HTTP proxy preserves TLS record order with large bodies", async
}
});

// The tunnel's TLS engine queued the ciphertext of a request body in a
// BoringSSL memory BIO. That buffer never shrinks, so a pooled tunnel kept an
// allocation the size of the largest body it ever sent. ASAN keeps freed
// memory resident, so RSS shows nothing there.
test.skipIf(isASAN)("a pooled HTTPS proxy tunnel does not keep the memory of a large request body", async () => {
using origin = Bun.serve({
port: 0,
tls: tlsCert,
async fetch(req) {
let received = 0;
for await (const chunk of req.body!) received += chunk.byteLength;
return new Response(String(received));
},
});
const MiB = 1024 * 1024;
const bodySize = 64 * MiB;
const bound = bodySize / 2;
const fixture = `
const url = ${JSON.stringify(origin.url.href)};
const proxy = ${JSON.stringify(httpProxyServer.url)};
async function post(size) {
const res = await fetch(url, { method: "POST", body: Buffer.alloc(size, "a"), proxy, tls: { rejectUnauthorized: false } });
const received = await res.text();
if (received !== String(size)) throw new Error("the origin received " + received + " of " + size + " bytes");
}
// The first request opens the tunnel. Every later one reuses it.
await post(1024);
Bun.gc(true);
const before = process.memoryUsage.rss();
await post(${bodySize});
// Small requests on the same tunnel, until the upload's memory is released.
let retained = Infinity;
for (let attempt = 0; attempt < 100 && retained > ${bound}; attempt++) {
await post(1024);
Bun.gc(true);
retained = Math.min(retained, process.memoryUsage.rss() - before);
}
console.log(JSON.stringify({ retainedMiB: retained / ${MiB} }));
`;
httpProxyServer.log.length = 0;
await using proc = Bun.spawn({
cmd: [bunExe(), "-e", fixture],
env: { ...bunEnv, ...proxyFreeEnv },
stdout: "pipe",
stderr: "pipe",
});
const [stdout, stderr, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
if (exitCode !== 0) console.error("stderr:", stderr);
const connects = httpProxyServer.log.filter(line => line === `CONNECT localhost:${origin.port}`).length;
Comment thread
claude[bot] marked this conversation as resolved.

// One CONNECT: every request went through the same tunnel, which is still
// pooled when the child reads its RSS. The memory BIO kept more than the body.
expect(stdout).toStartWith("{");
expect(connects).toBe(1);
expect(JSON.parse(stdout).retainedMiB).toBeLessThan(bound / MiB);
expect(exitCode).toBe(0);
});

test("HTTPS origin close-delimited body via HTTP proxy does not ECONNRESET", async () => {
// Inline raw HTTPS origin: 200 + no Content-Length then close
const originServer = tls.createServer(
Expand Down
227 changes: 227 additions & 0 deletions test/js/node/tls/node-tls-connect.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@ import {
bunRun,
tls as COMMON_CERT_,
isASAN,
isWindows,
nodeExe,
rejectUnauthorizedScope,
tempDir,
Expand Down Expand Up @@ -855,6 +856,232 @@ it("a client and a server TLSSocket connected through a synchronous in-memory du
});
});

describe("large writes and reads of a TLSSocket over a Duplex transport", () => {
// The TLS engine behind a Duplex kept its ciphertext in two BoringSSL memory
// BIOs. A memory BIO moves every byte that is still unread to its front after
// each read, so both directions cost quadratic time in what was queued: a
// write was copied out 64 KiB at a time, and a chunk from the transport was
// taken out one TLS record at a time.
const pattern = Buffer.from(Array.from({ length: 251 }, (_, i) => i));

// A client and a server TLSSocket joined in memory. `wire` sees every chunk
// of ciphertext the named side writes and returns what to deliver to the peer.
function connectedPair(wire: (from: "client" | "server", chunk: Buffer) => Buffer | null = (_, chunk) => chunk) {
const makeSide = (name: "client" | "server", peer: () => Duplex) =>
new Duplex({
read() {},
write(chunk: Buffer, _encoding, callback) {
const forwarded = wire(name, chunk);
if (forwarded !== null) peer().push(forwarded);
callback();
},
final(callback) {
peer().push(null);
callback();
},
});
const clientSide: Duplex = makeSide("client", () => serverSide);
const serverSide: Duplex = makeSide("server", () => clientSide);
const server = new TLSSocket(serverSide, { isServer: true, secureContext: tls.createSecureContext(COMMON_CERT_) });
server.on("error", () => {});
server.on("end", () => server.end());
const client = tls.connect({ socket: clientSide, rejectUnauthorized: false });
return { client, server, clientSide };
}

// Resolves with the next `length` bytes `socket` emits.
function receive(socket: TLSSocket, length: number) {
const { promise, resolve, reject } = Promise.withResolvers<Buffer>();
const chunks: Buffer[] = [];
let received = 0;
const onData = (chunk: Buffer) => {
chunks.push(chunk);
received += chunk.length;
if (received < length) return;
socket.off("data", onData);
socket.off("error", reject);
resolve(Buffer.concat(chunks));
};
socket.on("data", onData);
socket.once("error", reject);
return promise;
}

it("one write() reaches the transport in at most two chunks, like node", async () => {
const payload = Buffer.alloc(1024 * 1024, pattern);
const transportWrites: number[] = [];
let recording = false;
const { client, server } = connectedPair((from, chunk) => {
if (recording && from === "client") transportWrites.push(chunk.length);
return chunk;
});
await once(client, "secureConnect");

const received = receive(server, payload.length);
recording = true;
client.write(payload);
const plaintext = await received;
recording = false;
client.end();
await once(client, "close");

expect(plaintext.equals(payload)).toBe(true);
// 64 KiB pieces made this 17 writes. Node gives the transport 2 chunks.
expect(transportWrites.length).toBeLessThanOrEqual(2);
});

// process.cpuUsage() advances in scheduler ticks on Windows, about 15.6 ms
// each. Both measurements here are shorter than one tick on a release
// build, so both read 0 and the ratio is NaN.
it.skipIf(isWindows)("one large chunk from the transport costs the same CPU per byte as small chunks", async () => {
const payload = Buffer.alloc(8 * 1024 * 1024, pattern);
let held: Buffer[] | null = null;
const { client, server, clientSide } = connectedPair((from, chunk) => {
if (held === null || from === "client") return chunk;
held.push(chunk);
return null;
});
await once(client, "secureConnect");

// The CPU time this process takes to decrypt `payload`, with its ciphertext
// handed to the client's transport in pieces of `pieceSize` bytes.
async function cpuTimeToReceive(pieceSize: number) {
held = [];
await new Promise<void>((resolve, reject) => server.write(payload, error => (error ? reject(error) : resolve())));
const ciphertext = Buffer.concat(held);
held = null;

const received = receive(client, payload.length);
const before = process.cpuUsage();
for (let offset = 0; offset < ciphertext.length; offset += pieceSize) {
clientSide.push(ciphertext.subarray(offset, offset + pieceSize));
}
const plaintext = await received;
const { user, system } = process.cpuUsage(before);
expect(plaintext.equals(payload)).toBe(true);
return user + system;
}

// The best of two runs each, interleaved, so that a busy machine does not
// decide the ratio.
let small = Infinity;
let large = Infinity;
for (let run = 0; run < 2; run++) {
small = Math.min(small, await cpuTimeToReceive(64 * 1024));
large = Math.min(large, await cpuTimeToReceive(Infinity));
}
client.end();
await once(client, "close");

// About 1 when the cost is linear. With the memory BIO it was 36 or more.
expect(large / small).toBeLessThan(4);
});

it("writes made from 'data' and from the transport's write() keep both streams in order", async () => {
// Both writes are made while the engine is handing ciphertext to the
// transport or plaintext to the socket, and each is larger than one pass of
// either. A TLS record that leaves out of order fails the connection.
const fromServer = Buffer.alloc(320 * 1024, pattern);
const fromData = Buffer.alloc(200 * 1024, pattern.subarray(7));
const fromTransportWrite = Buffer.alloc(150 * 1024, pattern.subarray(13));

let writeFromTransport = false;
const { client, server } = connectedPair((from, chunk) => {
if (from === "client" && writeFromTransport) {
writeFromTransport = false;
client.write(fromTransportWrite);
}
return chunk;
});
await once(client, "secureConnect");

server.once("data", () => server.write(fromServer));
const atServer = receive(server, 2 + fromData.length + fromTransportWrite.length);

const chunks: Buffer[] = [];
let depth = 0;
let nestedDataEvents = 0;
let received = 0;
let repliedFromData = false;
const allAtClient = Promise.withResolvers<void>();
client.on("data", (chunk: Buffer) => {
depth++;
if (depth > 1) nestedDataEvents++;
chunks.push(chunk);
received += chunk.length;
if (!repliedFromData) {
repliedFromData = true;
writeFromTransport = true;
client.write(fromData);
}
if (received >= fromServer.length) allAtClient.resolve();
depth--;
});
client.on("error", allAtClient.reject);
client.write("go");

const [serverGot] = await Promise.all([atServer, allAtClient.promise]);
client.end();
await once(client, "close");

expect({
nestedDataEvents,
clientGotAll: Buffer.concat(chunks).equals(fromServer),
serverGotAll: serverGot.equals(Buffer.concat([Buffer.from("go"), fromData, fromTransportWrite])),
}).toEqual({ nestedDataEvents: 0, clientGotAll: true, serverGotAll: true });
});

it("ciphertext that arrives while an earlier chunk is partly decrypted is read in order", async () => {
// The first chunk ends inside a TLS record. The rest arrives in small
// pieces from inside each 'data' event, so it queues behind bytes the
// engine has not read yet, after it has read part of the first chunk. The
// engine's queue is a ring: the record that was cut comes back in two reads.
const payload = Buffer.alloc(1024 * 1024, pattern);
let held: Buffer[] | null = null;
const { client, server, clientSide } = connectedPair((from, chunk) => {
if (held === null || from === "client") return chunk;
held.push(chunk);
return null;
});
await once(client, "secureConnect");

held = [];
await new Promise<void>((resolve, reject) => server.write(payload, error => (error ? reject(error) : resolve())));
const wire = Buffer.concat(held);
held = null;

// The middle of the last record that starts in the first two thirds.
let cut = 0;
for (let offset = 0; offset < (wire.length * 2) / 3; ) {
const recordLength = 5 + wire.readUInt16BE(offset + 3);
cut = offset + (recordLength >> 1);
offset += recordLength;
}
const pieces: Buffer[] = [];
for (let offset = cut; offset < wire.length; offset += 7001) pieces.push(wire.subarray(offset, offset + 7001));

const received = receive(client, payload.length);
let next = 0;
let pushedFromData = 0;
client.on("data", () => {
for (let pushed = 0; pushed < 4 && next < pieces.length; pushed++, pushedFromData++) {
clientSide.push(pieces[next++]);
}
});
clientSide.push(wire.subarray(0, cut));
while (next < pieces.length) clientSide.push(pieces[next++]);
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const plaintext = await received;
client.end();
await once(client, "close");

expect({
somePiecesCameFromData: pushedFromData > 0,
length: plaintext.length,
intact: plaintext.equals(payload),
}).toEqual({ somePiecesCameFromData: true, length: payload.length, intact: true });
});
});

it("the last 'data' event fires before the close_notify reply is written to a duplex transport (tls.connect({ socket }))", async () => {
// The peer's last application data and its close_notify reach the engine in
// one chunk. The engine used to answer the close_notify before it emitted
Expand Down
Loading
Loading