公开的 LLM API 入口不能只是转发请求,还要同时考虑滥用、上游故障与流式响应。本文用一个可直接部署的 Cloudflare Worker,完成按 IP 限流、主备模型降级和 SSE 透明转发,并说明自动切换不能跨越的边界。

网关结构与准备工作

这里假设主、备供应商都提供 OpenAI 兼容的 POST /v1/chat/completions 接口。Worker 对外只暴露 /v1/chat/completions,调用方不需要知道实际模型和供应商。

整体链路如下:

环节处理内容失败行为
请求入口校验方法、路径和 JSON返回 400 或 404
Durable Object按来源 IP 固定窗口计数超限返回 429
主模型设置首包响应超时超时、429、408、5xx 时降级
备用模型接管尚未开始的响应失败后返回 502 或原始错误
SSE 转发直接传递字节流不拆分、不重新拼接事件

创建一个 Worker 项目后,使用下面的 wrangler.toml。Durable Object 用来保证同一 IP 的计数更新集中到同一个对象中,避免多个 Worker 实例各自计数。

name = "llm-gateway"
main = "src/index.js"
compatibility_date = "2025-01-01"

[vars]
MAIN_URL = "https://main.example.com/v1/chat/completions"
MAIN_MODEL = "main-model"
BACKUP_URL = "https://backup.example.com/v1/chat/completions"
BACKUP_MODEL = "backup-model"
UPSTREAM_TIMEOUT_MS = "8000"

[[durable_objects.bindings]]
name = "RATE_LIMITER"
class_name = "RateLimiter"

[[migrations]]
tag = "v1"
new_sqlite_classes = ["RateLimiter"]

域名和模型名需要替换成供应商真实值。密钥不要写入配置文件,通过 Wrangler Secret 保存:

npx wrangler secret put MAIN_API_KEY
npx wrangler secret put BACKUP_API_KEY
npx wrangler deploy

用 Durable Object 按 IP 限流

仅在 Worker 全局变量里保存计数并不可靠:不同请求可能落到不同实例,实例也会被回收。下面使用固定窗口算法,每个 IP 每分钟允许 20 次请求。它实现简单,但窗口交界处可能出现突发流量;如果需要更平滑的限制,可以进一步改成令牌桶。

完整的 src/index.js 如下:

const WINDOW_MS = 60_000;
const REQUEST_LIMIT = 20;

export class RateLimiter {
  constructor(ctx, env) {
    this.ctx = ctx;
    this.env = env;
  }

  async fetch() {
    const now = Date.now();

    const result = await this.ctx.storage.transaction(async (tx) => {
      let state = await tx.get("window");

      if (!state || now >= state.resetAt) {
        state = { count: 0, resetAt: now + WINDOW_MS };
      }

      if (state.count >= REQUEST_LIMIT) {
        return {
          allowed: false,
          remaining: 0,
          resetAt: state.resetAt
        };
      }

      state.count += 1;
      await tx.put("window", state);

      return {
        allowed: true,
        remaining: REQUEST_LIMIT - state.count,
        resetAt: state.resetAt
      };
    });

    return Response.json(result);
  }
}

export default {
  async fetch(request, env) {
    const url = new URL(request.url);

    if (request.method !== "POST" || url.pathname !== "/v1/chat/completions") {
      return jsonError(404, "not_found", "接口不存在");
    }

    const contentLength = Number(request.headers.get("content-length") || 0);
    if (contentLength > 1024 * 1024) {
      return jsonError(413, "payload_too_large", "请求体不能超过 1 MiB");
    }

    const ip = request.headers.get("CF-Connecting-IP") || "local";
    const limiterId = env.RATE_LIMITER.idFromName(ip);
    const limiter = env.RATE_LIMITER.get(limiterId);
    const limitResult = await limiter
      .fetch("https://rate-limit.internal/check")
      .then((response) => response.json());

    if (!limitResult.allowed) {
      const retryAfter = Math.max(
        1,
        Math.ceil((limitResult.resetAt - Date.now()) / 1000)
      );

      const response = jsonError(429, "rate_limited", "请求过于频繁");
      response.headers.set("Retry-After", String(retryAfter));
      response.headers.set("X-RateLimit-Remaining", "0");
      return response;
    }

    let input;
    try {
      input = await request.json();
    } catch {
      return jsonError(400, "invalid_json", "请求体必须是合法 JSON");
    }

    if (!Array.isArray(input.messages) || input.messages.length === 0) {
      return jsonError(400, "invalid_request", "messages 必须是非空数组");
    }

    const providers = [
      {
        name: "main",
        url: env.MAIN_URL,
        model: env.MAIN_MODEL,
        apiKey: env.MAIN_API_KEY
      },
      {
        name: "backup",
        url: env.BACKUP_URL,
        model: env.BACKUP_MODEL,
        apiKey: env.BACKUP_API_KEY
      }
    ];

    const wantsStream = input.stream === true;
    const timeoutMs = Number(env.UPSTREAM_TIMEOUT_MS || 8000);
    let lastError;

    for (let index = 0; index < providers.length; index += 1) {
      const provider = providers[index];
      const canFallback = index < providers.length - 1;

      try {
        const upstream = await fetchWithHeaderTimeout(
          provider,
          { ...input, model: provider.model },
          timeoutMs
        );

        const contentType = upstream.headers.get("content-type") || "";
        const invalidStream =
          wantsStream &&
          upstream.ok &&
          !contentType.toLowerCase().includes("text/event-stream");

        const retryableStatus =
          upstream.status === 408 ||
          upstream.status === 429 ||
          upstream.status >= 500;

        if (canFallback && (retryableStatus || invalidStream)) {
          await upstream.body?.cancel();
          continue;
        }

        if (invalidStream) {
          await upstream.body?.cancel();
          return jsonError(502, "invalid_upstream", "上游未返回 SSE 流");
        }

        return forwardResponse(upstream, wantsStream, provider.name);
      } catch (error) {
        lastError = error;
        if (!canFallback) break;
      }
    }

    console.error("All upstream providers failed", lastError);
    return jsonError(502, "upstream_unavailable", "上游模型暂时不可用");
  }
};

async function fetchWithHeaderTimeout(provider, payload, timeoutMs) {
  const controller = new AbortController();
  const timer = setTimeout(() => controller.abort(), timeoutMs);

  try {
    return await fetch(provider.url, {
      method: "POST",
      headers: {
        "Authorization": `Bearer ${provider.apiKey}`,
        "Content-Type": "application/json",
        "Accept": payload.stream ? "text/event-stream" : "application/json"
      },
      body: JSON.stringify(payload),
      signal: controller.signal
    });
  } finally {
    clearTimeout(timer);
  }
}

function forwardResponse(upstream, wantsStream, providerName) {
  const headers = new Headers();
  headers.set(
    "Content-Type",
    upstream.headers.get("content-type") ||
      (wantsStream ? "text/event-stream; charset=utf-8" : "application/json")
  );
  headers.set("Cache-Control", "no-store");
  headers.set("X-Content-Type-Options", "nosniff");
  headers.set("X-LLM-Provider", providerName);

  return new Response(upstream.body, {
    status: upstream.status,
    headers
  });
}

function jsonError(status, code, message) {
  return Response.json(
    { error: { code, message } },
    {
      status,
      headers: { "Cache-Control": "no-store" }
    }
  );
}

CF-Connecting-IP 是 Cloudflare 在已部署环境中提供的客户端地址。代码中的 local 只用于本地调试;不要接受客户端自定义的其他 IP 请求头,否则攻击者可以自行轮换地址。

主备模型降级的判定

并非所有错误都应该切换供应商。401、403 通常意味着网关密钥或权限配置错误,切换备用模型会掩盖问题;400 则多半是调用参数不合法。因此示例只对以下情况降级:

  • 建立响应超时;
  • HTTP 408 或 429;
  • HTTP 5xx;
  • 请求流式响应,但上游成功响应不是 text/event-stream

代码会为每个供应商重新执行 JSON.stringify(payload),而不是复用已经消费过的请求流。模型名由网关覆盖,避免客户端绕过路由策略直接指定昂贵模型。

这里的超时准确说是“等待响应头超时”。fetch() 返回后,Worker 已经取得状态码和响应头,此时才会把响应体交给客户端。这个节点非常重要:在任何字节发给客户端之前,可以安全切换供应商;一旦 SSE 已开始输出,再切换就可能把两个模型的事件混入同一条连接。

SSE 转发为什么不要按数据块解析

SSE 事件用空行分隔,但网络数据块不保证与事件边界一致。一个 data: 事件可能被拆成多个块,多个事件也可能合并在一个块里。因此不能假设每次 reader.read() 都得到完整事件。

本文没有修改事件内容,最稳妥的方式就是直接把 upstream.body 交给新的 Response。这样既不会错误处理 UTF-8 跨块字符,也不会因为等待完整事件而增加额外缓冲。

转发时也没有复制上游全部响应头。尤其不应手工保留旧的 Content-Length,因为流式响应长度未知;Connection 等逐跳头同样不需要设置。这里只保留内容类型,并明确禁止缓存。

直接转发也带来一个明确限制:如果主模型已经返回 SSE 响应头,随后在生成中途停住,网关不能无缝改用备用模型继续回答。要处理这种情况,只能选择中止连接,让客户端重试;或者先在网关缓冲完整答案,但那会失去流式输出的意义。生产系统应把“响应头超时”和“流中断”作为两类故障分别监控。

部署验证与运维边界

部署后可以先验证非流式请求,再验证 SSE:

curl https://your-worker.example.com/v1/chat/completions \
  -H 'Content-Type: application/json' \
  -d '{"messages":[{"role":"user","content":"只回复 ok"}],"stream":false}'

curl -N https://your-worker.example.com/v1/chat/completions \
  -H 'Content-Type: application/json' \
  -d '{"messages":[{"role":"user","content":"写两句话"}],"stream":true}'

curl -N 会关闭客户端输出缓冲,便于观察事件是否逐步到达。还应分别测试主模型超时、主模型返回 429、备用模型也失败和连续请求触发限流等场景。

按 IP 限流只能作为第一层防刷:公司出口、校园网络可能共享 IP,攻击者也可能使用代理池。公开服务通常还需要用户身份、API Token、总预算上限和并发限制。本文的 1 MiB 检查依赖 Content-Length,它适合拦截正常客户端的明显大请求,但不能替代更严格的流式请求体大小控制。

总结

一个可用的 LLM 网关至少要守住三条边界:

  • 用 Durable Object 集中维护同一 IP 的计数,避免实例内存限流失效;
  • 只在响应尚未开始时,根据超时、429、408 和 5xx 切换备用模型;
  • SSE 透明转发时直接传递字节流,不把网络数据块误当成完整事件。

主备切换不是无限兜底:流已经发送后无法安全换模型,按 IP 限流也无法识别真实用户。把这些限制明确写进设计,再叠加身份、预算和监控,通常比试图在一个 Worker 中隐藏所有故障更可靠。