捕获、算成本、投队列:边缘侧如何不阻塞地留存每次调用
30 秒导读: 上一章(01)把一次 LLM 请求截获并转发给了 provider。这一章讲转发之后到消息进队列这半程:边缘 Worker 怎么在不拖慢把响应流回用户的前提下,无损攒下整个响应体,把流式碎片拼回完整内容与 token,算出这次调用的钱,最后把一条结构化消息塞进 Kafka/SQS。真正把消息加工成数据库行是消费端 Jawn 的活,见 03。
1. 这半程在全景里的位置(先看大盘)
一次调用在边缘侧其实分两条时间线:用户看得见的那条必须最快,观测那条躲在后面慢慢做。
用户请求 ─▶ 代理转发 provider ─▶ 响应流 ─┬─▶ 原样转发回用户 (关键路径:越快越好)
│
└─▶ 旁路副本(tee)──▶ [本章] 攒 chunk → 解析 → 算成本 → 投队列
▲
ctx.waitUntil:响应已返回,这里还在跑
关键就一句:采集不能站在用户和响应之间。Helicone 的做法是把 provider 的响应流分叉(tee)——一份直接 pipe 回用户,一份进旁路缓冲;等 Worker 已经 return 了响应,才在 Cloudflare 的 ctx.waitUntil() 里做解析、定价、入队这些重活。
本章覆盖的五个部件:
| 部件 | 干什么 | 在哪个文件 |
|---|---|---|
ReadableInterceptor | 把响应流一分为二,旁路那份逐 chunk 攒进内存 | worker/src/lib/util/ReadableInterceptor.ts |
RequestBodyBuffer | 缓冲请求体,按大小选内存/远端容器,供转发与落库复用 | worker/src/RequestBodyBuffer/ |
| streamParsers | 把 SSE 流式碎片重新拼成完整响应 + 还原 token | worker/src/lib/dbLogger/streamParsers/ |
DBLoggable | 核心日志对象:解析响应、抽 token、组装队列消息 | worker/src/lib/dbLogger/DBLoggable.ts |
packages/cost + HeliconeProducer | 算这次调用的成本;把消息投 Kafka/SQS/HTTP | packages/cost/、worker/src/lib/clients/producers/HeliconeProducer.ts |
2. 不阻塞的地基:一根 tee 出来的旁路流
2.1 它要解决的小问题
Worker 要同时满足两件互相打架的事:
- 把响应尽快、原样流回用户——流式场景下用户在等第一个 token,任何缓冲都会被感知成卡顿。
- 完整留下响应体用于观测——但响应流是一次性的,读了转发就没了给日志,读了日志就没了给用户。
2.2 思路:边转发边偷偷复制每个 chunk
答案是包一层 ReadableStream:每次 provider 那侧 read() 出一个 chunk,先 enqueue 给下游(用户),再顺手 push 一份进内部数组。转发一步都不等旁路。
provider 响应流
│ read()
▼
┌─────────────────────────┐
│ pull(controller){ │
│ enqueue(value) ──────────▶ 下游 = 用户 (先走,不阻塞)
│ onChunk(value) ──────────▶ responseBody.push(decode(chunk)) (顺手攒)
│ } │
└─────────────────────────┘
│ done
▼
emit("done", { body, endTimeUnix, firstChunkTimeUnix })
真实实现里,pull() 就是先 enqueue 再记录的顺序,采集绝不插在 转发前面:
- 新流的
pull():controller.enqueue(value)之后才onChunk(value)(ReadableInterceptor.ts:75-87,方法interceptStream)。 onChunk把 chunk 解码后推进responseBody数组,并在第一个 chunk 落地时记下firstChunkTimeUnix——这就是首 token 延迟的原始数据(ReadableInterceptor.ts:56onChunk)。- 流读完(
done)或被取消(cancel)时,onDone把攒好的body数组连同时间戳,通过EventEmitter发出一个CompletedStream事件(ReadableInterceptor.ts:38onDone、:4interface CompletedStream)。
2.3 采集侧怎么"等"到完整响应
日志侧不关心流什么时候结束,它只要一个"给我最终完整 body"的 Promise。waitForStream() 就是:轮询 cachedChunk(那个 CompletedStream),没到就每秒重试,最长等到 chunkTimeoutMs(默认 30 分钟)(ReadableInterceptor.ts:111 waitForStream、:25 默认超时)。
在 ProxyRequestHandler 里,这根拦截流被接进 DBLoggable 的 getResponseBody:日志要读响应时,本质是 await interceptor.waitForStream(),拿到 body 数组和 endTime(ProxyRequestHandler.ts:187 getResponseBody)。首 token 延迟同理,用 firstChunkTimeUnix - startTime 现算(ProxyRequestHandler.ts:207 timeToFirstToken)。
一句话精华: 转发和采集共用同一次
read(),但顺序上转发永远先走;采集只是"路过时抄一份",且抄的结果通过事件异步交付,谁都不用等谁。
3. 请求体的缓冲:小的进内存,大的进容器
响应体靠 tee 攒,请求体则在更早、请求刚进来时就被 RequestBodyBuffer 接管。它要同时服务三个下游:转发给 provider、落库到 S3、以及事后读出来算 body 大小/抽 model。所以它必须能被安全地反复读,又不能把一个几十 MB 的 body 一股脑吃进 Worker 内存。
3.1 两种实现,一个接口
IRequestBodyBuffer 定义了统一契约(IRequestBodyBuffer.ts:5),两种实现:
| 实现 | 适用 | body 存哪 |
|---|---|---|
RequestBodyBuffer_InMemory | 小 body / GET·HEAD / 无容器绑定 | Worker 内存里的 cachedText 字符串 |
RequestBodyBuffer_Remote | 大 body(需容器绑定) | 远端 Durable Object / 容器 |
选哪个由 RequestBodyBufferBuilder 一个不消费 body的启发式决定(RequestBodyBufferBuilder.ts:80):
method 是 GET/HEAD 或无 body ───────────────▶ InMemory
Content-Length 已知 且 ≤ 20 MiB ───────────▶ InMemory
Content-Length 已知 且 > 20 MiB 且有容器 ──▶ Remote
Content-Length 未知 / 无容器绑定 ──────────▶ InMemory (兜底)
那个 20 MiB 阈值写死在 MAX_INMEMORY_BYTES,注释里还留着调参人的名字和日期(RequestBodyBufferBuilder.ts:94)。注意:size 未知且有容器的分支,代码目前直接回退 InMemory,原本用 tee() 探测大小的逻辑被注释掉了(RequestBodyBufferBuilder.ts:119-136)——这是"代码里看得出来"的当前行为,别被上面的注释文档误导。
3.2 InMemory 的几个要点
- 懒读一次、缓存复用:
unsafeGetRawText()第一次 才request.text(),之后都吃cachedText(RequestBodyBuffer_InMemory.ts:87)。方法名带unsafe是提醒:这会把整个 body 拉进内存。 - 从 body 里抽字段:
isStream()、model()、userId()都是JSON.parse后取一个字段(RequestBodyBuffer_InMemory.ts:199isStream、:209model)。这几个值后面组装消息、选流式解析器都要用。 - body 覆写带原型链防护:
setBodyOverride走递归合并,但显式跳过__proto__/constructor/prototype这几个危险 key(RequestBodyBuffer_InMemory.ts:53DANGEROUS_KEYS、:59applyOverride)——这是防原型污染的小心思。 - 落库统一成对存:
uploadS3Body把请求体和响应体打包成一个 JSON 存 S3,新版(version 2)同时存 OpenAI 归一化格式和 provider 原生格式(RequestBodyBuffer_InMemory.ts:214)。
4. 流式响应:把一地碎片拼回一句话
4.1 它要解决的小问题
非流式响应是一个完整 JSON,直接 parse 就有 usage。但流式(SSE)响应是几十上百行 data: {...} 碎片,每行只带一小截 delta——文本要拼起来,token 用量散落在最后几个 chunk。要落库,就得先把这堆碎片还原成一个等价的完整响应对象。
DBLoggable.parseResponse() 是总调度:先问 RequestBodyBuffer 这次是不是流、是什么 model,再按 (provider × 是否流式) 分派到不同解析器(DBLoggable.ts:308 parseResponse)。
parseResponse
├─ 4xx/3xx 错误 ──────────────▶ 原样 JSON.parse
├─ 非流 + ANTHROPIC/GOOGLE ──▶ 直接读 usage,必要时自算 total
├─ 流 + ANTHROPIC ───────────▶ anthropicAIStream()
├─ 流 (其它) ────────────────▶ parseOpenAIStream()
└─ VERCEL 且 body 像 SSE ────▶ parseVercelStream()