流式Chat接口集成测试:压缩、中间件与SSE帧边界的坑
流式聊天端点的测试盲区:压缩中间件缓冲整个响应、nginx默认缓冲、TCP读取与SSE帧边界不对齐导致内容丢失;提供具体修复方案。
流式聊天端点的测试盲区:压缩中间件缓冲整个响应、nginx默认缓冲、TCP读取与SSE帧边界不对齐导致内容丢失;提供具体修复方案。
大多数流式测试调用 provider SDK 并检查 token 是否到达。这测的是 provider。而 bug 几乎总是出现在 provider socket 和 client 之间你自己的那五十行代码里,而这些代码只有在向你自己的端点发请求时才会运行。
流式聊天端点就是一个代理。它打开一个上游请求,读取 Server-Sent Events,可能对其做转换,再向下游写自己的事件。这条链上每一步都有一种失败模式,仅测 provider 调用是看不到的:
Compression middleware 会缓冲整个响应。默认的 compression() 或 gzip 层会收集输出直到有足够的内容才压缩。流在字节层面完全正确,但会一次性全部到达,延迟好几秒。
反向代理反而会缓冲它。nginx 默认就这么做;解决办法是 X-Accel-Buffering: no 响应头,得由你的端点来发送。
重新组帧会丢失帧。TCP 读取不会对齐 SSE 帧边界。代码对每次读取执行 chunk.toString().split("\n\n"),忘记处理跨两次读取的 remainder,每当帧跨越两次读取时就会静默丢失内容——这在长回答时会发生,但在短测试中几乎从不发生。
终端帧被转发了,但 socket 保持打开。客户端渲染完整个回答后会保持连接直到超时。
这些都不是 provider 的 bug,也都从你 token 处理函数的单元测试中看不到。它们需要从正门进去的请求才能触发。
上游必须被模拟,因为真实的 provider 调用会让测试变慢、不确定、还会计费。它还必须真正流式返回:一个一次性返回整个 body 的 mock 会把上面的缓冲 bug 藏起来,因为根本没有什么需要缓冲的。Mock Service Worker 2 可以做到这一点,因为它的 HttpResponse 接受与 Fetch Response 构造函数相同类型的 body,包括 ReadableStream。
// tests/streaming-endpoint.test.ts
import { afterAll, afterEach, beforeAll, expect, test } from "vitest";
import { http, HttpResponse } from "msw";
import { setupServer } from "msw/node";
function sseBody(frames: string[], gapMs = 0) {
const encoder = new TextEncoder();
return new ReadableStream({
async start(controller) {
for (const frame of frames) {
controller.enqueue(encoder.encode(frame));
if (gapMs) await new Promise((r) => setTimeout(r, gapMs));
}
controller.close();
},
});
}
const chunk = (delta: object, finish: string | null = null) =>
'data: ' +
JSON.stringify({
id: 'chatcmpl-test',
object: "chat.completion.chunk",
model: "gpt-4o-mini",
choices: [{ index: 0, delta, finish_reason: finish }],
}) +
'\n\n';
const upstream = setupServer(
http.post("https://api.openai.com/v1/chat/completions", () =>
new HttpResponse(
sseBody(
[
chunk({ role: "assistant", content: "" }),
chunk({ content: "Hel" }),
chunk({ content: "lo" }),
chunk({}, "stop"),
"data: [DONE]\n\n",
],
5,
),
{ headers: { "content-type": "text/event-stream" } },
),
),
);
beforeAll(() => upstream.listen({ onUnhandledRequest: "error" }));
afterEach(() => upstream.resetHandlers());
afterAll(() => upstream.close());
onUnhandledRequest: "error" 在这里起了真正的作用。没有它,触达第二个未 mock 的 provider 的代码路径会静默失败或打到网络;有了它,测试会失败并指出是哪个 URL。五毫秒的帧间隔使得缓冲变得可观测——会缓冲的代理会把四次独立到达变成一次。
不要在测试中使用 EventSource。它只会发 GET 请求,会自己重连,还会藏起你想断言的原始字节。直接读取响应 body 并自己做组帧——而且只写一次,写对,作为可复用的辅助函数,因为你在测试里写的 reader 同时也是你正在找 bug 的那个 reader。
async function readFrames(res: Response) {
const frames: string[] = [];
const arrivals: number[] = [];
const started = performance.now();
let buffer = "";
const reader = res.body!.pipeThrough(new TextDecoderStream()).getReader();
for (;;) {
const { value, done } = await reader.read();
if (done) break;
buffer += value;
let i: number;
while ((i = buffer.indexOf("\n\n")) !== -1) {
frames.push(buffer.slice(0, i));
arrivals.push(performance.now() - started);
buffer = buffer.slice(i + 2);
}
}
// Anything still here is a frame that was never terminated.
return { frames, arrivals, trailing: buffer };
}
TextDecoderStream 是重要的细节。它在多次读取之间携带解码器状态,所以跨两次网络读取的多字节字符能正确解码,而不会变成替换字符。每次对原始 chunk 执行 String(value) 是一个 bug,只有人在发送 emoji 或带重音的词时才会出现。这种失败模式有专门的文章:testing reassembly of streamed tokens。
文本来自生产中的模型和你的 fixture,所以断言它们只测了 fixture。断言传输的属性,而这些属性无论模型说什么都必须成立:
Headers。 content-type 以 text/event-stream 开头,cache-control 包含 no-cache,而 content-encoding 不存在——最后一条是压缩检查,对于在生产中极难诊断的 bug 这只是一行断言。
增量性。 到达了不止一帧,且第一帧和最后一帧的到达时间差非零。如果你的端点缓冲了,每条记录的到达时间都会在彼此的一毫秒之内。
组帧。 trailing 是空的。非空的 remainder 意味着最后一帧写的时候没有空行终止符,大多数客户端会永远卡住而不是交付它。
终止。 流以你的终端帧结束,且恰好一次,reader 达到 done。参见 testing that a stream closes cleanly。
与上游格式隔离。 如果你的端点发出自己的事件格式而不是透传 provider 的,断言没有上游字段泄露。对代理和透传都同样通过的测试没有在测你的翻译层。
在 beforeAll 中把你的应用启动在临时端口上,或者如果框架允许,导入路由 handler 并用 Request 调用它。临时端口多几行代码是值得的:它动了真实的 HTTP server,而 header 和缓冲行为实际就住在那里面。
安装上面的 MSW server 以拦截 provider 调用。让你的应用指向与 handler 匹配相同的 base URL。
用 fetch 和一个真实 body POST 到你的端点。先不要传 signal——取消要另外写测试。
在读取任何字节之前先断言四个 header。这里失败比六十帧之后失败读起来便宜。
用 readFrames 读取,然后断言 frames.length、trailing、终端帧和到达时间分布。
再加一个 case:fixture 把一帧拆到两次 enqueue 中发送——半行 data:,然后是剩下的。你的端点仍然必须产出相同的下游帧。这是唯一能捕获重新组帧 bug 的 case,只需要两行代码。
Header 期望取决于你的运行时和反向代理。无服务器平台可能在你的 handler 之外添加或剥离 content-encoding 和 transfer-encoding,所以在集成测试中断言你的 handler 能控制的,部署环境的边缘情况单独用针对真实环境的冒烟检查来覆盖。
如果你的端点代理了多个 provider,这个测试要乘以多份:上游 fixture 必须以每个 provider 的帧形状存在,你的翻译层也需要按形状写 case。Multigrid 在 provider 之间暴露了同一种流式格式,所以 fixture 和翻译层不再是按 provider 划分的 work——无论你用不用 gateway 都值得了解,因为替代方案是随着你添加的每个模型而增长的 fixture 矩阵。
Asserting on a Stream of Chunks Instead of a Final String
Testing That a Stream Closes Cleanly on the Happy Path
Contract Tests for Streaming Chunk Format