公开的 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 中隐藏所有故障更可靠。
评论