fix(v03): proxy cloud api through bad upstream ports
This commit is contained in:
@@ -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() {
|
||||
|
||||
@@ -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<Response> {
|
||||
const method = request.method || "GET";
|
||||
let timedOut = false;
|
||||
return new Promise((resolve) => {
|
||||
let settled = false;
|
||||
let timeout: ReturnType<typeof setTimeout> | 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<Uint8Array>;
|
||||
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<Uint8Array>);
|
||||
bodyStream.on("error", (error) => upstreamRequest.destroy(error));
|
||||
bodyStream.pipe(upstreamRequest);
|
||||
});
|
||||
}
|
||||
|
||||
function nodeRequestHeaders(headers: Headers, target: URL): Record<string, string> {
|
||||
const next: Record<string, string> = {};
|
||||
for (const [key, value] of headers.entries()) {
|
||||
if (isHopByHopHeader(key)) continue;
|
||||
next[key] = value;
|
||||
}
|
||||
next.host = target.host;
|
||||
return next;
|
||||
}
|
||||
|
||||
function responseHeadersFromNode(headers: Record<string, string | string[] | number | undefined>): 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) {
|
||||
|
||||
Reference in New Issue
Block a user