diff --git a/cmd/hwlab-edge-proxy/main.test.ts b/cmd/hwlab-edge-proxy/main.test.ts index eb75192b..b7aa9d85 100644 --- a/cmd/hwlab-edge-proxy/main.test.ts +++ b/cmd/hwlab-edge-proxy/main.test.ts @@ -37,6 +37,19 @@ test("edge proxy reports local health and proxies live health to upstream", asyn } }); +test("edge proxy proxies upstream port 6667 without using fetch bad-port handling", async () => { + const upstream = await startUpstream(6667); + const proxy = await startEdgeProxy(upstream.port); + try { + const live = await fetchJson(`http://127.0.0.1:${proxy.port}/health/live`); + assert.equal(live.serviceId, "hwlab-cloud-api"); + assert.equal(upstream.captured.url, "/health/live"); + } finally { + await proxy.stop(); + await upstream.stop(); + } +}); + test("edge proxy tunnels websocket upgrade to upstream", async () => { const upstream = await startWebSocketUpstream(); const proxy = await startEdgeProxy(upstream.port); @@ -50,7 +63,7 @@ test("edge proxy tunnels websocket upgrade to upstream", async () => { } }); -async function startUpstream() { +async function startUpstream(port = 0) { let captured = null; const server = createServer(async (request, response) => { const chunks = []; @@ -67,7 +80,7 @@ async function startUpstream() { status: "live" })); }); - await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + await new Promise((resolve) => server.listen(port, "127.0.0.1", resolve)); return { port: server.address().port, get captured() { diff --git a/cmd/hwlab-edge-proxy/main.ts b/cmd/hwlab-edge-proxy/main.ts index 5c0e044a..95c0ee19 100644 --- a/cmd/hwlab-edge-proxy/main.ts +++ b/cmd/hwlab-edge-proxy/main.ts @@ -1,4 +1,8 @@ #!/usr/bin/env bun +import { request as httpRequest } from "node:http"; +import { request as httpsRequest } from "node:https"; +import { Readable } from "node:stream"; + import { healthPayload, resolveHostPort, @@ -128,27 +132,82 @@ for (const signal of ["SIGINT", "SIGTERM"]) { } async function proxyFetch(request: Request, upstreamUrl: string, timeoutMs: number) { - const target = upstreamHttpUrl(request, upstreamUrl); - const controller = new AbortController(); - const timeout = setTimeout(() => controller.abort(), timeoutMs); - try { - return await fetch(target, { - method: request.method, - headers: request.headers, - body: request.method === "GET" || request.method === "HEAD" ? undefined : request.body, - redirect: "manual", - signal: controller.signal + const target = new URL(upstreamHttpUrl(request, upstreamUrl)); + return await proxyNodeRequest(request, target, upstreamUrl, timeoutMs); +} + +function proxyNodeRequest(request: Request, target: URL, upstreamUrl: string, timeoutMs: number): Promise { + const method = request.method || "GET"; + let timedOut = false; + return new Promise((resolve) => { + let settled = false; + let timeout: ReturnType | null = null; + const fail = () => { + if (settled) return; + settled = true; + if (timeout) clearTimeout(timeout); + resolve(proxyFailureResponse({ timedOut, timeoutMs, upstreamUrl })); + }; + const transport = target.protocol === "https:" ? httpsRequest : httpRequest; + const upstreamRequest = transport(target, { method, headers: nodeRequestHeaders(request.headers, target) }, (upstreamResponse) => { + if (settled) { + upstreamResponse.resume(); + return; + } + settled = true; + if (timeout) clearTimeout(timeout); + const status = upstreamResponse.statusCode ?? 502; + const body = method === "HEAD" || status === 204 || status === 304 ? null : Readable.toWeb(upstreamResponse) as unknown as ReadableStream; + resolve(new Response(body, { status, headers: responseHeadersFromNode(upstreamResponse.headers) })); }); - } catch (error) { - const timedOut = controller.signal.aborted; - const code = timedOut ? "proxy_timeout" : "upstream_unavailable"; - const userMessage = timedOut - ? `Code Agent 代理等待上游超过 ${timeoutMs}ms;输入已保留,可稍后重试。` - : "Code Agent 代理暂时无法连接上游;输入已保留,可稍后重试。"; - return json({ status: "failed", error: { code, layer: "proxy", category: timedOut ? "timeout" : "proxy", retryable: true, userMessage, message: userMessage, blocker: { code, layer: "proxy", retryable: true, summary: userMessage } }, upstream: upstreamUrl, message: userMessage }, 502); - } finally { - clearTimeout(timeout); + timeout = setTimeout(() => { + timedOut = true; + upstreamRequest.destroy(new Error(`upstream timed out after ${timeoutMs}ms`)); + }, timeoutMs); + upstreamRequest.on("error", fail); + if (method === "GET" || method === "HEAD" || !request.body) { + upstreamRequest.end(); + return; + } + const bodyStream = Readable.fromWeb(request.body as unknown as ReadableStream); + bodyStream.on("error", (error) => upstreamRequest.destroy(error)); + bodyStream.pipe(upstreamRequest); + }); +} + +function nodeRequestHeaders(headers: Headers, target: URL): Record { + const next: Record = {}; + for (const [key, value] of headers.entries()) { + if (isHopByHopHeader(key)) continue; + next[key] = value; } + next.host = target.host; + return next; +} + +function responseHeadersFromNode(headers: Record): Headers { + const next = new Headers(); + for (const [key, value] of Object.entries(headers)) { + if (value === undefined || isHopByHopHeader(key)) continue; + if (Array.isArray(value)) { + for (const item of value) next.append(key, item); + continue; + } + next.set(key, String(value)); + } + return next; +} + +function isHopByHopHeader(key: string): boolean { + return ["connection", "keep-alive", "proxy-authenticate", "proxy-authorization", "te", "trailer", "transfer-encoding", "upgrade"].includes(key.toLowerCase()); +} + +function proxyFailureResponse({ timedOut, timeoutMs, upstreamUrl }: { timedOut: boolean; timeoutMs: number; upstreamUrl: string }) { + const code = timedOut ? "proxy_timeout" : "upstream_unavailable"; + const userMessage = timedOut + ? `Code Agent 代理等待上游超过 ${timeoutMs}ms;输入已保留,可稍后重试。` + : "Code Agent 代理暂时无法连接上游;输入已保留,可稍后重试。"; + return json({ status: "failed", error: { code, layer: "proxy", category: timedOut ? "timeout" : "proxy", retryable: true, userMessage, message: userMessage, blocker: { code, layer: "proxy", retryable: true, summary: userMessage } }, upstream: upstreamUrl, message: userMessage }, 502); } function upstreamHttpUrl(request: Request, upstreamUrl: string) {