OpenAI 兼容接口能返回 200,不代表它真的兼容:响应头迟迟不到、SSE 被缓冲、工具参数残缺都很常见。本文用 Cloudflare Worker 构建透明代理,并在客户端加入分层超时、对冲请求和非流式降级。
先区分中转层的几类“成功失败”
接入中转服务时,我通常不会只记录状态码,而是把一次请求拆成四个阶段:建立连接、收到响应头、收到首个有效事件、完整结束。很多“接口偶尔卡死”,实际是其中某个阶段没有截止时间。
| 现象 | 常见原因 | 不能只靠什么解决 |
|---|---|---|
| 长时间没有响应头 | 上游连接阻塞、DNS 或中转队列拥堵 | 单纯增加总超时 |
声明 stream=true,但迟迟没有内容 | 中转读取完整响应后才转发,或压缩、缓存链路缓冲 | 看到 text/event-stream 就判定成功 |
| 返回若干 SSE 事件后停住 | 上游断流,中转没有正确关闭连接 | 只设置首包超时 |
| 普通文本正常,工具调用残缺 | 中转没有完整拼接 tool_calls 增量,或提前结束 | 对残缺 JSON 做容错执行 |
| 200 响应中携带错误对象 | 兼容实现没有映射上游状态码 | 只检查 HTTP 200 |
因此,健康判定至少应检查 HTTP 状态、内容类型和首个可解析的 SSE data 事件。若涉及工具调用,还要等参数完整拼接并通过 JSON 解析后,才能进入执行阶段。
用 Cloudflare Worker 做尽量透明的代理
下面的 Worker 只允许转发 /v1/ 路径,给“等待上游响应头”设置 15 秒期限,并直接把上游 ReadableStream 放入响应。关键点是不要调用 response.text() 或 response.json(),否则流式响应会被完整读取后再返回,看起来就像假死。
// worker.js
export default {
async fetch(request, env) {
const incoming = new URL(request.url);
if (!incoming.pathname.startsWith("/v1/")) {
return new Response("Not Found", { status: 404 });
}
const base = env.UPSTREAM_BASE.endsWith("/")
? env.UPSTREAM_BASE
: env.UPSTREAM_BASE + "/";
const target = new URL(
incoming.pathname.slice(1) + incoming.search,
base
);
const headers = new Headers(request.headers);
headers.delete("host");
headers.delete("content-length");
headers.delete("connection");
if (env.UPSTREAM_API_KEY) {
headers.set("authorization", `Bearer ${env.UPSTREAM_API_KEY}`);
}
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), 15_000);
try {
const upstream = await fetch(target, {
method: request.method,
headers,
body: ["GET", "HEAD"].includes(request.method)
? null
: request.body,
redirect: "manual",
signal: controller.signal
});
clearTimeout(timer);
const responseHeaders = new Headers(upstream.headers);
responseHeaders.delete("connection");
responseHeaders.delete("keep-alive");
responseHeaders.delete("transfer-encoding");
responseHeaders.set("x-proxy-layer", "cloudflare-worker");
return new Response(upstream.body, {
status: upstream.status,
statusText: upstream.statusText,
headers: responseHeaders
});
} catch (error) {
clearTimeout(timer);
const timeout = error instanceof Error && error.name === "AbortError";
return Response.json(
{ error: { message: timeout ? "upstream header timeout" : "upstream unavailable" } },
{ status: timeout ? 504 : 502 }
);
}
}
};
最小配置如下,可用 npx wrangler dev 本地运行,再用 npx wrangler deploy 部署:
# wrangler.toml
name = "openai-compatible-proxy"
main = "worker.js"
compatibility_date = "2025-01-01"
[vars]
UPSTREAM_BASE = "https://api.openai.com"
密钥不要写入配置文件,可执行 npx wrangler secret put UPSTREAM_API_KEY。如果上游本身要求调用方的 Authorization,则移除 Worker 中覆盖该请求头的代码。
这段实现解决的是响应头超时和透明转发,并没有解决流开始后的永久停顿。流的空闲超时最好放在客户端,因为客户端最清楚模型、请求类型以及可以接受的等待时间。
超时要按阶段设置,而不是共用一个数字
我会为中转调用设置四种预算:响应头超时限制连接和排队;首事件超时要求出现可解析的 SSE 数据;空闲超时限制相邻事件之间的间隔;总超时作为最终保险。不同模型的首字延迟差异较大,这些值应通过自己的日志调整,不宜照搬固定数字。
观测上至少记录目标中转、模型、是否流式、收到响应头的耗时、首个有效事件耗时、最后事件时间、HTTP 状态和结束原因。不要记录完整提示词、密钥及工具参数;确需排障时,可记录请求体大小、工具数量和脱敏后的错误类别。
SSE 的结束不能仅依赖 TCP 断开。客户端应识别 data: [DONE],同时处理响应中出现的 error 对象。连接直接关闭但没有 [DONE] 时,可以标记为“不完整流”,而不是当作正常完成。
客户端对冲:谁先证明可用,就采用谁
对冲不是同时向所有中转广播。更克制的做法是先请求主线路;在短暂等待后仍未拿到有效 SSE,再启动备用线路。胜者必须返回成功状态、正确内容类型和至少一个可解析事件,随后立即取消败者,避免继续计费。
以下脚本可在 Node.js 20 运行。设置两个以逗号分隔的服务根地址后执行 node hedge.mjs;两条流式线路都失败时,它会对主线路发起一次非流式降级请求。
// hedge.mjs
const bases = (process.env.OPENAI_BASES || "")
.split(",")
.map(x => x.trim())
.filter(Boolean);
const apiKey = process.env.OPENAI_API_KEY;
const model = process.env.MODEL;
if (bases.length === 0 || !apiKey || !model) {
throw new Error("请设置 OPENAI_BASES、OPENAI_API_KEY 和 MODEL");
}
const body = stream => JSON.stringify({
model,
stream,
messages: [{ role: "user", content: "用一句话解释什么是幂等性。" }]
});
const sleep = ms => new Promise(resolve => setTimeout(resolve, ms));
function launch(base) {
const controller = new AbortController();
const promise = candidate(base, controller);
return { controller, promise };
}
async function candidate(base, controller) {
const timer = setTimeout(() => controller.abort(), 15_000);
try {
const response = await fetch(
`${base.replace(/\/$/, "")}/v1/chat/completions`,
{
method: "POST",
headers: {
authorization: `Bearer ${apiKey}`,
"content-type": "application/json",
accept: "text/event-stream"
},
body: body(true),
signal: controller.signal
}
);
if (!response.ok) throw new Error(`HTTP ${response.status}`);
if (!(response.headers.get("content-type") || "").includes("text/event-stream")) {
throw new Error("响应不是 SSE");
}
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = "";
let scanned = 0;
while (true) {
const { done, value } = await reader.read();
if (done) throw new Error("首个有效事件前流已结束");
buffer += decoder.decode(value, { stream: true });
while (true) {
const match = buffer.slice(scanned).match(/\r?\n\r?\n/);
if (!match) break;
const end = scanned + match.index;
const event = buffer.slice(scanned, end);
scanned = end + match[0].length;
const data = event
.split(/\r?\n/)
.filter(line => line.startsWith("data:"))
.map(line => line.slice(5).trimStart())
.join("\n");
if (!data) continue;
if (data === "[DONE]") throw new Error("流未返回有效数据");
const parsed = JSON.parse(data);
if (parsed.error) throw new Error(parsed.error.message || "上游错误");
clearTimeout(timer);
return { reader, decoder, prefix: buffer, controller };
}
}
} catch (error) {
clearTimeout(timer);
throw error;
}
}
async function readWithIdleTimeout(reader, controller, ms) {
let timer;
try {
return await Promise.race([
reader.read(),
new Promise((_, reject) => {
timer = setTimeout(() => {
controller.abort();
reject(new Error("流空闲超时"));
}, ms);
})
]);
} finally {
clearTimeout(timer);
}
}
async function fallback() {
const response = await fetch(
`${bases[0].replace(/\/$/, "")}/v1/chat/completions`,
{
method: "POST",
headers: {
authorization: `Bearer ${apiKey}`,
"content-type": "application/json"
},
body: body(false),
signal: AbortSignal.timeout(30_000)
}
);
if (!response.ok) throw new Error(`降级请求失败:HTTP ${response.status}`);
console.log(JSON.stringify(await response.json(), null, 2));
}
const attempts = [launch(bases[0])];
let winner;
try {
const early = await Promise.race([
attempts[0].promise.then(value => ({ value }), error => ({ error })),
sleep(800).then(() => ({ timeout: true }))
]);
if (early.value) {
winner = early.value;
} else if (bases.length > 1) {
attempts.push(launch(bases[1]));
winner = await Promise.any(attempts.map(x => x.promise));
} else {
throw early.error || new Error("主线路首事件超时");
}
for (const attempt of attempts) {
if (attempt.controller !== winner.controller) attempt.controller.abort();
}
process.stdout.write(winner.prefix);
while (true) {
const { done, value } = await readWithIdleTimeout(
winner.reader,
winner.controller,
20_000
);
if (done) break;
process.stdout.write(winner.decoder.decode(value, { stream: true }));
}
process.stdout.write(winner.decoder.decode());
} catch (error) {
for (const attempt of attempts) attempt.controller.abort();
console.error(`流式线路失败,执行非流式降级:${error.message}`);
await fallback();
}
对冲会增加请求量,应只对超时或明确的网关错误启用,并设置并发上限。已经向用户输出内容后再切换线路,可能产生重复或语义跳变,因此示例只在首个有效事件前竞争,胜者确定后不再切换。
工具调用要延迟执行
流式工具调用的 function.arguments 通常分散在多个增量中。客户端应按 choice 和工具索引拼接字符串,等结束原因为 tool_calls 后再执行 JSON.parse。解析失败、工具名缺失或调用 ID 不一致时,应丢弃本次候选,并改用非流式请求重试。
不要在工具已经执行后进行普通重试或切换中转,否则付款、发信、写数据库等副作用可能重复发生。可靠做法是在业务工具层传递幂等键并持久化执行结果;不能假设兼容中转会自动对模型请求去重。
总结
- 200 状态和 SSE 响应头都不足以证明中转可用,首个可解析事件才是更可靠的流式就绪信号。
- Cloudflare Worker 应直接转发
ReadableStream,避免读取完整响应,同时限制路径并设置响应头超时。 - 客户端需要分别管理响应头、首事件、流空闲和总超时,并把不完整结束纳入监控。
- 对冲应延迟启动备用线路,验证胜者后取消败者;两条流式线路都失败时,可降级为非流式请求。
- 工具调用必须完整拼接、校验后再执行,副作用操作还需要业务层幂等,不能依赖中转兼容性兜底。
评论