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,避免读取完整响应,同时限制路径并设置响应头超时。
  • 客户端需要分别管理响应头、首事件、流空闲和总超时,并把不完整结束纳入监控。
  • 对冲应延迟启动备用线路,验证胜者后取消败者;两条流式线路都失败时,可降级为非流式请求。
  • 工具调用必须完整拼接、校验后再执行,副作用操作还需要业务层幂等,不能依赖中转兼容性兜底。