Merge pull request #1372 from pikasTech/fix/issue1368-1369-edge-proxy-bad-port
fix(v03): 修复 edge-proxy 6667 upstream 转发
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 () => {
|
test("edge proxy tunnels websocket upgrade to upstream", async () => {
|
||||||
const upstream = await startWebSocketUpstream();
|
const upstream = await startWebSocketUpstream();
|
||||||
const proxy = await startEdgeProxy(upstream.port);
|
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;
|
let captured = null;
|
||||||
const server = createServer(async (request, response) => {
|
const server = createServer(async (request, response) => {
|
||||||
const chunks = [];
|
const chunks = [];
|
||||||
@@ -67,7 +80,7 @@ async function startUpstream() {
|
|||||||
status: "live"
|
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 {
|
return {
|
||||||
port: server.address().port,
|
port: server.address().port,
|
||||||
get captured() {
|
get captured() {
|
||||||
|
|||||||
@@ -1,4 +1,8 @@
|
|||||||
#!/usr/bin/env bun
|
#!/usr/bin/env bun
|
||||||
|
import { request as httpRequest } from "node:http";
|
||||||
|
import { request as httpsRequest } from "node:https";
|
||||||
|
import { Readable } from "node:stream";
|
||||||
|
|
||||||
import {
|
import {
|
||||||
healthPayload,
|
healthPayload,
|
||||||
resolveHostPort,
|
resolveHostPort,
|
||||||
@@ -128,27 +132,82 @@ for (const signal of ["SIGINT", "SIGTERM"]) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async function proxyFetch(request: Request, upstreamUrl: string, timeoutMs: number) {
|
async function proxyFetch(request: Request, upstreamUrl: string, timeoutMs: number) {
|
||||||
const target = upstreamHttpUrl(request, upstreamUrl);
|
const target = new URL(upstreamHttpUrl(request, upstreamUrl));
|
||||||
const controller = new AbortController();
|
return await proxyNodeRequest(request, target, upstreamUrl, timeoutMs);
|
||||||
const timeout = setTimeout(() => controller.abort(), timeoutMs);
|
}
|
||||||
try {
|
|
||||||
return await fetch(target, {
|
function proxyNodeRequest(request: Request, target: URL, upstreamUrl: string, timeoutMs: number): Promise<Response> {
|
||||||
method: request.method,
|
const method = request.method || "GET";
|
||||||
headers: request.headers,
|
let timedOut = false;
|
||||||
body: request.method === "GET" || request.method === "HEAD" ? undefined : request.body,
|
return new Promise((resolve) => {
|
||||||
redirect: "manual",
|
let settled = false;
|
||||||
signal: controller.signal
|
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) {
|
timeout = setTimeout(() => {
|
||||||
const timedOut = controller.signal.aborted;
|
timedOut = true;
|
||||||
const code = timedOut ? "proxy_timeout" : "upstream_unavailable";
|
upstreamRequest.destroy(new Error(`upstream timed out after ${timeoutMs}ms`));
|
||||||
const userMessage = timedOut
|
}, timeoutMs);
|
||||||
? `Code Agent 代理等待上游超过 ${timeoutMs}ms;输入已保留,可稍后重试。`
|
upstreamRequest.on("error", fail);
|
||||||
: "Code Agent 代理暂时无法连接上游;输入已保留,可稍后重试。";
|
if (method === "GET" || method === "HEAD" || !request.body) {
|
||||||
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);
|
upstreamRequest.end();
|
||||||
} finally {
|
return;
|
||||||
clearTimeout(timeout);
|
}
|
||||||
|
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) {
|
function upstreamHttpUrl(request: Request, upstreamUrl: string) {
|
||||||
|
|||||||
Reference in New Issue
Block a user