From 4b771678ab49f58890be1ebe83cceeec70db02ad Mon Sep 17 00:00:00 2001
From: Alix-007
Date: Sat, 8 Aug 2026 11:50:18 +0800
Subject: [PATCH] fix(proxy-capture): time out stalled CONNECT upstreams
(#109269)
* fix(proxy-capture): time out stalled CONNECT upstreams
* fix(proxy-capture): preserve existing socket boundary
* fix(proxy-capture): clarify CONNECT timeout semantics
---------
Co-authored-by: Dallin Romney
---
src/proxy-capture/proxy-server.test.ts | 105 ++++++++++++++++++++++++-
src/proxy-capture/proxy-server.ts | 24 ++++++
2 files changed, 128 insertions(+), 1 deletion(-)
diff --git a/src/proxy-capture/proxy-server.test.ts b/src/proxy-capture/proxy-server.test.ts
index 22421300ad99..250603218b57 100644
--- a/src/proxy-capture/proxy-server.test.ts
+++ b/src/proxy-capture/proxy-server.test.ts
@@ -5,7 +5,7 @@ import {
createServer as createHttpServer,
type IncomingMessage,
} from "node:http";
-import net, { type AddressInfo } from "node:net";
+import net, { Socket, type AddressInfo } from "node:net";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { afterEach, describe, expect, it, vi } from "vitest";
@@ -413,7 +413,23 @@ async function rawSlowGetThroughProxy(params: {
});
}
+async function openConnectClient(proxyUrl: string, connectTarget: string): Promise {
+ const proxy = new URL(proxyUrl);
+ const socket = new Socket();
+ socket.on("error", () => {});
+ await new Promise((resolve, reject) => {
+ socket.once("error", reject);
+ socket.connect(Number(proxy.port), proxy.hostname, () => {
+ socket.off("error", reject);
+ resolve();
+ });
+ });
+ socket.write(`CONNECT ${connectTarget} HTTP/1.1\r\nHost: ${connectTarget}\r\n\r\n`);
+ return socket;
+}
+
afterEach(async () => {
+ vi.restoreAllMocks();
await cleanupTestRoot();
});
@@ -644,4 +660,91 @@ describe("startDebugProxyServer", () => {
await origin.stop();
}
});
+
+ it("returns 504 when a CONNECT upstream opening attempt times out", async () => {
+ const settings = await makeSettings();
+ const stalledUpstream = new Socket();
+ const setOpeningTimeout = vi.spyOn(stalledUpstream, "setTimeout");
+ let resolveConnectCalled!: (target: { hostname: string; port: number }) => void;
+ const connectCalled = new Promise<{ hostname: string; port: number }>((resolve) => {
+ resolveConnectCalled = resolve;
+ });
+ vi.spyOn(net, "connect").mockImplementation(((port: number, hostname: string) => {
+ resolveConnectCalled({ hostname, port });
+ return stalledUpstream;
+ }) as typeof net.connect);
+ const proxy = await startDebugProxyServer({ settings });
+ let client: Socket | undefined;
+
+ try {
+ const connectedClient = await openConnectClient(proxy.proxyUrl, "unreachable.example:443");
+ client = connectedClient;
+ let response = "";
+ connectedClient.setEncoding("utf8");
+ connectedClient.on("data", (chunk) => {
+ response += chunk.toString();
+ });
+ const clientClosed = new Promise((resolve) => {
+ connectedClient.once("close", resolve);
+ });
+
+ await expect(connectCalled).resolves.toMatchObject({
+ hostname: "unreachable.example",
+ port: 443,
+ });
+ expect(setOpeningTimeout).toHaveBeenCalledWith(30_000, expect.any(Function));
+ stalledUpstream.emit("timeout");
+ await clientClosed;
+
+ expect(response).toContain("504 Gateway Timeout");
+ expect(response).toContain("Gateway Timeout\n");
+ expect(stalledUpstream.destroyed).toBe(true);
+ expect(getDebugProxyCaptureStore().getSessionEvents(settings.sessionId, 10)).toContainEqual(
+ expect.objectContaining({
+ direction: "local",
+ errorText: "CONNECT upstream opening timed out after 30000ms of inactivity",
+ kind: "error",
+ protocol: "connect",
+ }),
+ );
+ } finally {
+ client?.destroy();
+ stalledUpstream.destroy();
+ await proxy.stop();
+ }
+ });
+
+ it("removes the CONNECT opening timeout after the upstream socket connects", async () => {
+ const settings = await makeSettings();
+ const upstream = new Socket();
+ const disableTimeout = vi.spyOn(upstream, "setTimeout");
+ let resolveConnectCalled!: () => void;
+ const connectCalled = new Promise((resolve) => {
+ resolveConnectCalled = resolve;
+ });
+ vi.spyOn(net, "connect").mockImplementation(((
+ _port: number,
+ _hostname: string,
+ onConnect: () => void,
+ ) => {
+ resolveConnectCalled();
+ upstream.once("connect", onConnect);
+ return upstream;
+ }) as typeof net.connect);
+ const proxy = await startDebugProxyServer({ settings });
+ let client: Socket | undefined;
+
+ try {
+ client = await openConnectClient(proxy.proxyUrl, "example.com:443");
+ await connectCalled;
+ upstream.emit("connect");
+
+ expect(disableTimeout).toHaveBeenCalledWith(0);
+ expect(upstream.listenerCount("timeout")).toBe(0);
+ } finally {
+ client?.destroy();
+ upstream.destroy();
+ await proxy.stop();
+ }
+ });
});
diff --git a/src/proxy-capture/proxy-server.ts b/src/proxy-capture/proxy-server.ts
index 2da5561b5cb2..1a21de4b1f72 100644
--- a/src/proxy-capture/proxy-server.ts
+++ b/src/proxy-capture/proxy-server.ts
@@ -17,6 +17,8 @@ const DEBUG_PROXY_DIRECT_CONNECT_OVERRIDE =
"OPENCLAW_DEBUG_PROXY_ALLOW_DIRECT_CONNECT_WITH_MANAGED_PROXY";
const CAPTURE_BODY_PREVIEW_BYTES = 8192;
const BAD_GATEWAY_BODY = "Bad Gateway\n";
+const DEBUG_PROXY_CONNECT_TIMEOUT_MS = 30_000;
+const GATEWAY_TIMEOUT_BODY = "Gateway Timeout\n";
type BodyPreviewCapture = {
chunks: Buffer[];
@@ -408,6 +410,10 @@ export async function startDebugProxyServer(params: {
return;
}
const upstreamSocket = net.connect(port, hostname, () => {
+ // This inactivity timeout only protects opening the upstream socket. CONNECT
+ // tunnels are intentionally long-lived and must not inherit it.
+ upstreamSocket.setTimeout(0);
+ upstreamSocket.off("timeout", onUpstreamConnectTimeout);
clientSocket.write("HTTP/1.1 200 Connection Established\r\n\r\n");
if (head.length > 0) {
upstreamSocket.write(head);
@@ -415,6 +421,24 @@ export async function startDebugProxyServer(params: {
clientSocket.pipe(upstreamSocket);
upstreamSocket.pipe(clientSocket);
});
+ function onUpstreamConnectTimeout() {
+ const message = `CONNECT upstream opening timed out after ${DEBUG_PROXY_CONNECT_TIMEOUT_MS}ms of inactivity`;
+ recordProxyEvent({
+ protocol: "connect",
+ direction: "local",
+ kind: "error",
+ flowId,
+ host: hostname,
+ path: req.url ?? "",
+ errorText: message,
+ });
+ upstreamSocket.destroy();
+ clientSocket.end(
+ `HTTP/1.1 504 Gateway Timeout\r\nConnection: close\r\nContent-Type: text/plain; charset=utf-8\r\nContent-Length: ${Buffer.byteLength(GATEWAY_TIMEOUT_BODY)}\r\n\r\n${GATEWAY_TIMEOUT_BODY}`,
+ () => clientSocket.destroy(),
+ );
+ }
+ upstreamSocket.setTimeout(DEBUG_PROXY_CONNECT_TIMEOUT_MS, onUpstreamConnectTimeout);
clientSocket.on("error", (error) => {
recordProxyEvent({
protocol: "connect",