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
125 changes: 37 additions & 88 deletions src/http/ProxyTunnel.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,7 +41,6 @@ pub use bun_uws::MaybeAnySocket as Socket;
#[derive(bun_ptr::CellRefCounted)]
pub struct ProxyTunnel {
pub(crate) wrapper: Option<ProxyTunnelWrapper>,
pub(crate) shutdown_err: Cell<Error>,
/// active socket is the socket that is currently being used
pub(crate) socket: Socket,
pub(crate) write_buffer: bun_io::StreamBuffer,
Expand All @@ -64,7 +63,6 @@ impl Default for ProxyTunnel {
fn default() -> Self {
Self {
wrapper: None,
shutdown_err: Cell::new(crate::Error::ConnectionClosed),
socket: Socket::None,
write_buffer: bun_io::StreamBuffer::default(),
did_have_handshaking_error: false,
Expand Down Expand Up @@ -113,27 +111,6 @@ impl ProxyTunnel {
unsafe { &mut *addr_of_mut!((*this.as_ptr()).write_buffer) }
}

/// Shared access to `shutdown_err` (a `Cell<Error>`; disjoint from
/// `wrapper`). Callers use `.get()`/`.set()` — no `&mut` needed.
#[inline]
fn shutdown_err_of<'a>(this: NonNull<Self>) -> &'a Cell<Error> {
// SAFETY: see [`Self::socket_of`].
unsafe { &*addr_of!((*this.as_ptr()).shutdown_err) }
}

/// Callback-safe close: sets `shutdown_err` then drives `wrapper.shutdown()`.
/// Takes `NonNull<Self>` so the SSLWrapper close callback (which reenters
/// `on_close` and reborrows tunnel fields via the disjoint accessors above)
/// does not alias a held `&mut ProxyTunnel`.
///
/// Module-level INVARIANT: `this` is a live intrusive-refcounted tunnel and
/// the caller's `&mut HTTPClient`/`&mut ProxyTunnel` borrows are NLL-dead
/// before this call (every callsite in this module follows that shape).
#[inline]
fn close_from_callback(this: NonNull<Self>, err: Error) {
Self::close_raw(this, err);
}

#[inline]
fn wrapper_ssl(this: NonNull<Self>) -> Option<NonNull<bun_boringssl_sys::SSL>> {
Self::wrapper_ref(this.as_ptr()).and_then(|w| w.ssl.get())
Expand Down Expand Up @@ -174,7 +151,7 @@ impl ProxyTunnel {
/// `*mut HTTPClient` registered in [`ProxyTunnel::start`]/`adopt`. The client
/// is embedded in its `AsyncHTTP` and outlives the tunnel. Each callback's
/// outer `&mut HTTPClient` must be NLL-dead before any reentrant call that
/// re-derives it via raw ptr (`close_from_callback`, `progress_update_*`,
/// re-derives it via raw ptr (`fail_request`, `progress_update_*`,
/// `on_writable`); call sites are shaped accordingly.
#[inline]
fn client_from_ctx<'a, 'c>(ctx: *mut HTTPClient<'c>) -> &'a mut HTTPClient<'c> {
Expand Down Expand Up @@ -249,19 +226,19 @@ fn on_data(ctx: *mut HTTPClient, decoded_data: &[u8]) {
// SAFETY: see on_open. `&mut HTTPClient` is disjoint from the caller's
// `&SSLWrapper` (HTTPClient holds the tunnel only by pointer). NLL
// ends this borrow before any reentrant call below that re-derives
// `&mut *ctx` (close → on_close, progress_update).
// `&mut *ctx` (fail_request, progress_update).
let this = client_from_ctx(ctx);
let Some(proxy_nn) = this.proxy_tunnel_ptr() else {
return;
};
let _guard = ProxyTunnel::ref_guard(proxy_nn);
let guard = ProxyTunnel::ref_guard(proxy_nn);
// While parked waiting for the JS `checkServerIdentity` verdict no request
// has been written through the tunnel, so any decrypted application data
// arriving here is unexpected.
if this.state.flags.is_waiting_for_cert_check {
scoped_log!(http_proxy_tunnel, "ProxyTunnel onData while parked");
// SAFETY: `this` dead (NLL); reenter via raw ptr.
ProxyTunnel::close_from_callback(proxy_nn, crate::Error::UnexpectedData);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, crate::Error::UnexpectedData);
return;
}
match this.state.response_stage {
Expand All @@ -273,10 +250,8 @@ fn on_data(ctx: *mut HTTPClient, decoded_data: &[u8]) {
let report_progress = match this.handle_response_body(decoded_data, false) {
Ok(v) => v,
Err(err) => {
// `this` is dead (NLL); reenter via raw ptr so on_close's
// fresh `&mut *ctx` / `&mut *proxy_ptr` do not alias us.
// SAFETY: tunnel pinned by ref_raw above.
ProxyTunnel::close_from_callback(proxy_nn, err);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, err);
return;
}
};
Expand All @@ -295,8 +270,8 @@ fn on_data(ctx: *mut HTTPClient, decoded_data: &[u8]) {
let report_progress = match this.handle_response_body_chunked_encoding(decoded_data) {
Ok(v) => v,
Err(err) => {
// SAFETY: see Body arm.
ProxyTunnel::close_from_callback(proxy_nn, err);
// `this` dead (NLL); see Body arm.
fail_request(ctx, guard, err);
return;
}
};
Expand Down Expand Up @@ -327,8 +302,8 @@ fn on_data(ctx: *mut HTTPClient, decoded_data: &[u8]) {
}
_ => {
scoped_log!(http_proxy_tunnel, "ProxyTunnel onData unexpected data");
// SAFETY: `this` dead (NLL); reenter via raw ptr.
ProxyTunnel::close_from_callback(proxy_nn, crate::Error::UnexpectedData);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, crate::Error::UnexpectedData);
}
}
}
Expand All @@ -345,7 +320,7 @@ fn on_handshake(
};
scoped_log!(http_proxy_tunnel, "ProxyTunnel onHandshake");
// Do NOT form `&mut ProxyTunnel` (see ALIASING NOTE).
let _guard = ProxyTunnel::ref_guard(proxy_nn);
let guard = ProxyTunnel::ref_guard(proxy_nn);
this.state.response_stage = HTTPStage::ProxyHeaders;
this.state.request_stage = HTTPStage::ProxyHeaders;
this.state.request_sent_len = 0;
Expand All @@ -357,9 +332,8 @@ fn on_handshake(
// only reject the connection if reject_unauthorized == true
if this.flags.reject_unauthorized && this.flags.did_have_handshaking_error {
let err = crate::get_cert_error_from_no(handshake_error.error_no);
// SAFETY: `this` dead (NLL); reenter via raw ptr so on_close's
// fresh `&mut *ctx` does not alias us.
ProxyTunnel::close_from_callback(proxy_nn, err);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, err);
return;
}
if this.wants_server_identity_check() {
Expand Down Expand Up @@ -428,12 +402,12 @@ fn on_handshake(
});
if this.flags.reject_unauthorized && peer_sent_certificate && handshake_error.error_no > 0 {
let err = crate::get_cert_error_from_no(handshake_error.error_no);
// SAFETY: `this` dead (NLL); reenter via raw ptr.
ProxyTunnel::close_from_callback(proxy_nn, err);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, err);
return;
}
// SAFETY: `this` dead (NLL); reenter via raw ptr.
ProxyTunnel::close_from_callback(proxy_nn, crate::Error::TLSHandshakeFailed);
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(ctx, guard, crate::Error::TLSHandshakeFailed);
return;
}
}
Expand Down Expand Up @@ -478,10 +452,6 @@ pub(crate) fn write_encrypted(ctx: *mut HTTPClient, encoded_data: &[u8]) {
}

fn on_close(ctx: *mut HTTPClient) {
// on_close is fired from inside SSLWrapper::shutdown (via close_raw) whose
// caller may itself be a callback that already held `&mut *ctx`; that
// outer borrow is required to be NLL-dead before close_raw is invoked
// (see on_data/on_handshake), so this fresh `&mut` is sole.
let this = client_from_ctx(ctx);
scoped_log!(
http_proxy_tunnel,
Expand All @@ -492,6 +462,7 @@ fn on_close(ctx: *mut HTTPClient) {
"tunnel exists"
}
);
// The client detaches the tunnel before it closes it, so its own closes return here.
let Some(proxy_nn) = this.proxy_tunnel_ptr() else {
return;
};
Expand All @@ -517,24 +488,25 @@ fn on_close(ctx: *mut HTTPClient) {
}
}

// Otherwise, treat as failure. `close_and_fail` de-tags the outer socket
// before `fail()` frees the AsyncHTTP that embeds `self` (the uSockets ext
// still points here until then).
let err = fail_err.unwrap_or_else(|| ProxyTunnel::shutdown_err_of(proxy_nn).get());
match ProxyTunnel::socket_of(proxy_nn) {
&Socket::Ssl(socket) => {
this.close_and_fail::<true>(err, socket);
}
&Socket::Tcp(socket) => {
this.close_and_fail::<false>(err, socket);
}
Socket::None => {
if fail_err.is_some() {
this.fail(err);
}
}
// `this` dead (NLL); reborrow via `client_from_ctx` inside.
fail_request(
ctx,
keepalive,
fail_err.unwrap_or(crate::Error::ConnectionClosed),
);
}

/// Fails the request. The client may be freed when this returns.
fn fail_request(ctx: *mut HTTPClient, keepalive: RefPtr<ProxyTunnel>, err: Error) {
let proxy = keepalive.as_non_null();
let this = client_from_ctx(ctx);
// `close_and_fail` de-tags the outer socket before `fail()` frees the client.
match ProxyTunnel::socket_of(proxy) {
&Socket::Ssl(socket) => this.close_and_fail::<true>(err, socket),
&Socket::Tcp(socket) => this.close_and_fail::<false>(err, socket),
Socket::None => this.fail(err),
}
ProxyTunnel::set_socket(proxy_nn, Socket::None);
ProxyTunnel::set_socket(proxy, Socket::None);
crate::http_thread().schedule_proxy_deref(keepalive);
}

Expand Down Expand Up @@ -653,29 +625,6 @@ impl ProxyTunnel {
}
}

/// Raw-pointer close: sets `shutdown_err` then drives `wrapper.shutdown()`.
/// Takes `NonNull<Self>` so the SSLWrapper close callback (which reenters
/// on_close and reborrows tunnel fields via raw projection) does not alias
/// a held `&mut ProxyTunnel`.
///
/// All field access goes through the disjoint-field accessors
/// ([`Self::shutdown_err_of`], [`Self::wrapper_ref`]), which already
/// encode the module INVARIANT that `this` is a live intrusive-refcounted
/// tunnel and no whole-struct `&mut ProxyTunnel` is held across the call.
/// Callers satisfy that by construction (see [`Self::close_from_callback`]).
pub(crate) fn close_raw(this: NonNull<Self>, err: Error) {
// `shutdown_err` is a `Cell<Error>` disjoint from `wrapper`; safe set.
Self::shutdown_err_of(this).set(err);
// shutdown() fires on_close synchronously, which accesses only
// disjoint tunnel fields via `addr_of!` (see on_close), so the
// `&SSLWrapper` from `wrapper_ref` is never aliased by a `&mut`
// across the reentrant call.
if let Some(wrapper) = ProxyTunnel::wrapper_ref(this.as_ptr()) {
// fast shutdown the connection
let _ = wrapper.shutdown(true);
}
}

pub(crate) fn shutdown(this: NonNull<Self>) {
// SAFETY: module INVARIANT — `this` is a live intrusive-refcounted
// tunnel and the caller's `&mut` borrows are NLL-dead before this call.
Expand Down
96 changes: 96 additions & 0 deletions test/cli/install/bun-install.test.ts
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
import { file, listen, Socket, spawn, write } from "bun";
import { afterAll, beforeAll, describe, expect, it, jest, setDefaultTimeout, test } from "bun:test";
import { randomBytes } from "crypto";
import { once } from "events";
import { readFileSync, readlinkSync, realpathSync, statSync } from "fs";
import { access, cp, exists, mkdir, readlink, rm, stat, writeFile } from "fs/promises";
import {
Expand All @@ -13,11 +15,15 @@ import {
runBunInstall,
tempDir,
textLockfile,
tls as tlsCert,
toBeValidBin,
toBeWorkspaceLink,
toHaveBins,
} from "harness";
import { type AddressInfo, createServer as createTcpServer, connect as tcpConnect } from "net";
import { basename, join, resolve, sep } from "path";
import { createServer as createTlsServer } from "tls";
import { gzipSync } from "zlib";
import {
createTestContext,
destroyTestContext,
Expand Down Expand Up @@ -513,6 +519,96 @@ describe.concurrent("bun-install", () => {
});
});

it("does not use a manifest whose gzip checksum is wrong, through a CONNECT tunnel", async () => {
await withContext(defaultOpts, async ctx => {
const requests: string[] = [];
// Big and incompressible, so the last chunk reaches the client in a later read than the response head.
const readme = randomBytes(300_000).toString("base64");
const registry = createTlsServer(tlsCert, socket => {
socket.on("error", () => {});
socket.once("data", data => {
const requestLine = data.toString("latin1").split("\r\n", 1)[0];
requests.push(requestLine);
if (requestLine !== "GET /bar HTTP/1.1") {
socket.end("HTTP/1.1 404 Not Found\r\nContent-Length: 0\r\nConnection: close\r\n\r\n");
return;
}
const manifest = gzipSync(
JSON.stringify({
name: "bar",
"dist-tags": { latest: "0.0.2" },
versions: {
"0.0.2": { name: "bar", version: "0.0.2", dist: { tarball: `${registryUrl}bar-0.0.2.tgz` } },
},
readme,
}),
);
// One bit of the CRC-32 in the gzip trailer. The HTTP framing around the stream is whole.
manifest[manifest.length - 5] ^= 1;
socket.write(
"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Encoding: gzip\r\nTransfer-Encoding: chunked\r\n\r\n",
);
socket.write(`${manifest.length.toString(16)}\r\n`);
socket.write(manifest);
socket.end("\r\n0\r\n\r\n");
});
});
registry.listen(0, "127.0.0.1");
await once(registry, "listening");
const registryPort = (registry.address() as AddressInfo).port;
const registryUrl = `https://localhost:${registryPort}/`;

// A forward proxy that only speaks CONNECT.
const proxy = createTcpServer(client => {
client.on("error", () => {});
client.once("data", () => {
const upstream = tcpConnect(registryPort, "127.0.0.1", () => {
client.write("HTTP/1.1 200 Connection Established\r\n\r\n");
client.pipe(upstream);
upstream.pipe(client);
});
upstream.on("error", () => client.destroy());
client.on("close", () => upstream.destroy());
});
});
proxy.listen(0, "127.0.0.1");
await once(proxy, "listening");

try {
await writeFile(
join(ctx.package_dir, "bunfig.toml"),
Bun.TOML.stringify({ install: { cache: false, registry: registryUrl } }),
);
await writeFile(
join(ctx.package_dir, "package.json"),
JSON.stringify({ name: "foo", version: "0.0.1", dependencies: { bar: "0.0.2" } }),
);
await using proc = spawn({
cmd: [bunExe(), "install", "--ca", tlsCert.cert],
cwd: ctx.package_dir,
stdout: "pipe",
stderr: "pipe",
env: {
...env,
HTTPS_PROXY: `http://127.0.0.1:${(proxy.address() as AddressInfo).port}`,
https_proxy: undefined,
NO_PROXY: undefined,
no_proxy: undefined,
BUN_CONFIG_HTTP_RETRY_COUNT: "0",
},
});
const [, err, exitCode] = await Promise.all([proc.stdout.text(), proc.stderr.text(), proc.exited]);
expect(err).toContain("error: ZlibError downloading package manifest bar");
// The tarball that the manifest names is never requested.
expect(requests).toEqual(["GET /bar HTTP/1.1"]);
expect(exitCode).toBe(1);
} finally {
registry.close();
proxy.close();
}
});
});

it("should support --registry CLI flag", async () => {
await withContext(defaultOpts, async ctx => {
const connected = jest.fn();
Expand Down
Loading
Loading