跳到主要内容

捕获、算成本、投队列:边缘侧如何不阻塞地留存每次调用

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 流式碎片重新拼成完整响应 + 还原 tokenworker/src/lib/dbLogger/streamParsers/
DBLoggable核心日志对象:解析响应、抽 token、组装队列消息worker/src/lib/dbLogger/DBLoggable.ts
packages/cost + HeliconeProducer算这次调用的成本;把消息投 Kafka/SQS/HTTPpackages/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:56 onChunk)。
  • 流读完(done)或被取消(cancel)时,onDone 把攒好的 body 数组连同时间戳,通过 EventEmitter 发出一个 CompletedStream 事件(ReadableInterceptor.ts:38 onDone:4 interface CompletedStream)。

2.3 采集侧怎么"等"到完整响应

日志侧不关心流什么时候结束,它只要一个"给我最终完整 body"的 Promise。waitForStream() 就是:轮询 cachedChunk(那个 CompletedStream),没到就每秒重试,最长等到 chunkTimeoutMs(默认 30 分钟)(ReadableInterceptor.ts:111 waitForStream:25 默认超时)。

ProxyRequestHandler 里,这根拦截流被接进 DBLoggablegetResponseBody:日志要读响应时,本质是 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:199 isStream:209 model)。这几个值后面组装消息、选流式解析器都要用。
  • body 覆写带原型链防护: setBodyOverride 走递归合并,但显式跳过 __proto__/constructor/prototype 这几个危险 key(RequestBodyBuffer_InMemory.ts:53 DANGEROUS_KEYS:59 applyOverride)——这是防原型污染的小心思。
  • 落库统一成对存: 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()

4.2 三种流的还原手法

解析器provider核心动作位置
parseOpenAIStreamOpenAI 兼容逐行 JSON.parse 后交给 consolidateTextFields 累加 choices[].deltaopenAIStreamParser.ts:3
anthropicAIStreamAnthropic按事件类型(content_block_delta/message_delta…)递归合并成一个 messageanthropicStreamParser.ts:118
parseVercelStreamVercel AI SDK累加 type:"text"text,从 metadata 取 usage,再拼成 OpenAI 兼容形状vercelStreamParser.ts:3

三者共用的暗线是"同类型字段就累加":字符串拼接、数字相加、对象递归下钻。这套通用合并逻辑在 responseParserHelpers.ts:

  • consolidateTextFields 负责 OpenAI 形状:把每个 chunk 的 choices[i].delta.content 首尾相接,最后把 delta 摊平成 message(responseParserHelpers.ts:19)。
  • recursivelyConsolidate 是底层的"按类型合并"原语:number += numberstring += stringobject 递归(responseParserHelpers.ts:2)。
  • Anthropic 因为事件模型不同,单独有一套 recursivelyConsolidateAnthropicListForClaude,专门处理 content_block_start/delta/stoptool_usepartial_json 增量等 Claude 特有事件(anthropicStreamParser.ts:28)。

一个坑: OpenAI 流解析里,最后一行被当成结束标记直接跳过(return {}),不参与合并(openAIStreamParser.ts:9);而流式响应即便解析成功,usage 也先塞 -1 占位、标 helicone_calculated:false(openAIStreamParser.ts:22)——真正的 token 由后面 getDetailedUsage 从响应体里抽,或交给消费端兜底。


5. DBLoggable:把这次调用打包成一条队列消息

DBLoggable 是这半程的中心对象(DBLoggable.ts:282)。它持有 request/response/timing 三块信息,对外只暴露一个入口 log(),内部真正干活的是 useKafka()

5.1 log():先鉴权限流,再交给 useKafka

log() 做三件事:取组织鉴权参数、查一次限流状态、然后把活全交给 useKafka()(DBLoggable.ts:704)。注意限流这里只是打标——orgRateLimit 为真也不拦请求(请求早转发完了),只是后面给队列消息降优先级(DBLoggable.ts:736 注释明说"没做早退";实际降优先级在 useKafka:1029)。

5.2 useKafka():采集 + 抽 token + 决定要不要存 body + 组装消息

useKafka() 是全章最密的一段(DBLoggable.ts:763),按顺序:

① 拿完整响应体。 await this.response.getResponseBody()——这一步内部就是 §2.3 那个 waitForStream(),在此才真正等流结束(DBLoggable.ts:804)。

② 从响应体里抽 token 和 model。 getDetailedUsage() 兼容三种形状:OpenAI 的 usage、Anthropic 的 cache token、Gemini 的 usageMetadata,抽出 prompt/completion/cache read/cache write/audio/reasoning 六类 token(DBLoggable.ts:472 getDetailedUsage)。抽这些是为了即使 body 不落 S3,队列消息里也带着用量

③ 决定要不要把 body 存 S3。 一组布尔逻辑决定跳过与否(DBLoggable.ts:838-850):

情形是否存 body 到 S3
免费额度超限 且 非 PTB跳过(只留 metadata)
免费额度超限 + PTB 且已拿到用量跳过(计费信息已够)
免费额度超限 + PTB 但没拿到用量(Jawn 得靠 body 抽用量计费)
两个 omit header 都设了跳过

存的时候,若是 AI Gateway 请求还会先把 provider 原生响应归一化成用户请求的格式再存(DBLoggable.ts:867 normalizeAIGatewayResponse)。

④ 组装 MessageData 一个大对象,分 heliconeMeta(模型覆写、omit 开关、prompt 版本、PTB 标志、posthog/lytix 集成 key 等)和 log(request 元数据 + response 元数据 + 上面抽出的 token)(DBLoggable.ts:931 kafkaMessage、消息类型见 producers/types.ts:39 MessageData)。

⑤ 投递。 超限则先降优先级,最后 await db.producer.sendMessage(kafkaMessage)(DBLoggable.ts:1032)。

注意成本字段的来龙去脉: 消息里的 response.cost 取自 this.response.cost(DBLoggable.ts:1012),而它只有异步日志(SDK 直传 provider 用量)那条路才会被填(DBLoggable.ts:257)。走代理/网关时这里通常是空的——边缘另算的成本用在别处(见 §6),最终结构化成本由消费端 Jawn 兜底。


6. 算成本:边缘为什么也要算一遍

6.1 谁在边缘算、算给谁用

投队列的 log() 是在 ProxyForwarder 里被 ctx.waitUntil() 拉起的(ProxyForwarder.ts:437)。那个内部 log 函数把两件事并行跑,最后 Promise.all 等齐(ProxyForwarder.ts:734):

ctx.waitUntil( log(...) )

├─ logPromise ───────────▶ loggable.log() → 组装消息 → 投队列 (§5)

└─ responseProcessingPromise ─▶ 读原始响应 → 算 cost → 结算 escrow / 记限流用量

关键区分:边缘算的这份 cost 主要不是给队列消息用的,而是给两件必须在边缘当场完成的事:

  1. PTB 钱包扣费(escrow 结算): Pass-Through Billing 请求先冻结一笔押金,请求完成后按真实成本 finalizeEscrowAndSyncSpend 扣款(ProxyForwarder.ts:665-682);成本从 USD 转成分(* 100)。
  2. 成本型限流计数: 把这次的花费喂给 bucket 限流器 recordBucketUsage(ProxyForwarder.ts:700-720)。

6.2 成本怎么算出来

responseProcessingPromise 分两条路(ProxyForwarder.ts:555):

  • AI Gateway 请求(BYOK/PTB):getUsageProcessor(provider) 拿到 provider 专属的用量解析器,parse 出标准化 ModelUsage,再交给 modelCostBreakdownFromRegistry 按新版模型注册表算出分项成本 totalCost(ProxyForwarder.ts:563-581)。
  • 普通代理请求:loggable.parseRawResponse 拿 model,尝试走同一套注册表;算不出再回退到 legacy 的 costOfPrompt(ProxyForwarder.ts:634-650)。

packages/cost 内部两代并存:

入口说明
新版注册表modelCostBreakdownFromRegistrycalculateModelCostBreakdown按 modality 分项(输入/输出/缓存/思考/图像/音频…)算,产出 CostBreakdown
legacycostOfPrompt / costOfproviders/mappings 里的老价目表,按 token × 单价累加
  • 新版聚合入口:costCalc.ts:51 modelCostBreakdownFromRegistry;分项结构 models/calculate-cost.ts:13 CostBreakdown
  • legacy 逐项累加:index.ts:61 costOfPrompt,含 Anthropic 缓存写 token 去重(5m/1h 分开计价、避免重复计)(index.ts:101-114)。
  • provider→解析器分派:usage/getUsageProcessor.ts:13,一个 switch 把十几个 OpenAI 兼容 provider 都指向 OpenAIUsageProcessor
  • ClickHouse 里成本是被乘过 COST_PRECISION_MULTIPLIER(10 亿)的整数,读出来要除回真实美元(costCalc.ts:8)。

路由/PTB/BYOK 的更多规则见 packages/cost/FLOWS.md(BYOK 全部先试、再按成本排序试 PTB)。


7. 投队列:三条通道,一个门面

HeliconeProducer 是投递门面(HeliconeProducer.ts:28)。它在构造时由工厂按环境变量选一条底层通道(HeliconeProducer.ts:6 MessageProducerFactory):

QUEUE_PROVIDER == "sqs" ─────────────▶ SQSProducerImpl (AWS SQS)
QUEUE_PROVIDER == "dual" ─────────────▶ DualWriteProducer (Kafka + SQS 双写)
否则 且 Upstash Kafka 环境齐全 ───────▶ KafkaProducerImpl (Upstash Kafka)
否则 ────────────────────────────────▶ null → 退化成 HTTP 直发

sendMessage() 还有一条旁路:如果配了 heliconeManualAccessKey(自托管场景)或根本没有 producer,就绕过队列,直接 HTTP POST 到 Jawn 的 /v1/log/request(HeliconeProducer.ts:46 sendMessage:58 sendMessageHttp)。这让 Helicone 在没有 Kafka/SQS 时也能跑。

各通道细节:

通道topic/队列消息形状重试降优先级
Kafkarequest-response-logs-prod,key = requestId{value: JSON.stringify(msg)} 再套一层 JSON3 次,间隔 1sno-op(已弃用给 SQS)
SQS正常 / 低优先级两条队列 URLMessageBody = JSON.stringify(msg)3 次,间隔 1s切到低优先级队列 URL
HTTPJawn /v1/log/requestPOST body 带 log/authorization/heliconeMeta无(catch 打日志)
  • Kafka 用 requestId 作 partition key,保证同一请求有序(KafkaProducerImpl.ts:42);重试循环 KafkaProducerImpl.ts:36
  • SQS 的降优先级就是把 queueUrl 换成 lowerPriorityQueueUrl(SQSProducer.ts:35 setLowerPriority)——§5.1 那个"超限打标"最终就落到这里,让被限流的组织的日志走慢车道,不挤占正常流量。

8. 不阻塞是怎么被保证的(把线索接起来)

三个机制叠起来,才让"留存每次调用"完全不出现在用户的关键路径上:

  1. 流分叉(tee): 转发用 enqueue,采集用旁路 push,同一次 read() 里转发先走(§2)。
  2. ctx.waitUntil: Worker return 响应后,Cloudflare 仍让 log(...) 这个 Promise 继续跑到完(ProxyForwarder.ts:437)——用户早拿到结果了,解析/定价/入队都在"之后"。
  3. 投递内建重试 + 超时上限: 队列 producer 各自 3 次重试(KafkaProducerImpl.ts:36SQSProducer.ts:44);waitForStream 有 30 分钟硬上限,卡死也不会永久挂着(ReadableInterceptor.ts:116)。

9. 边界与坑(诚实说)

  • size 未知 + 有容器分支当前退化成 InMemory。 探测大小的 tee() 逻辑被注释掉了(RequestBodyBufferBuilder.ts:119-136),所以"未知大小的大 body 走远端容器"这条路目前不生效,会吃进内存。
  • 流式 usage 在解析器里是占位的 -1 parseOpenAIStream 给的 usage 全是 -1(openAIStreamParser.ts:22),真 token 靠 getDetailedUsage 从 body 抽或消费端兜底;若响应体格式意外,failedToGetUsage 会为真,Jawn 得再从 S3 body 抽(DBLoggable.ts:824)。
  • 边缘成本 ≠ 队列消息成本。 边缘算的 cost 主要服务 escrow 扣费与限流(§6),代理路径下队列消息里的 cost 常为空,最终成本以 Jawn 计算为准(见 03)。
  • unsafeGetRawText 名副其实。 它把整个 body 读进内存,只该用在已知小 body 的场景(RequestBodyBuffer_InMemory.ts:87)。
  • 投递失败大多只打 console。 HTTP 直发失败仅 console.error(HeliconeProducer.ts:73),队列三次重试耗尽才返回 err——观测数据在极端情况下可能丢,这是"观测不许拖垮代理"的自觉取舍。

10. 代码地图(导航索引)

主题文件路径符号名
响应流分叉/旁路采集worker/src/lib/util/ReadableInterceptor.tsReadableInterceptorinterceptStreamonChunkwaitForStream
拦截流接入日志对象worker/src/lib/HeliconeProxyRequest/ProxyRequestHandler.tshandleProxyRequestgetResponseBodytimeToFirstToken
请求体缓冲接口worker/src/RequestBodyBuffer/IRequestBodyBuffer.tsIRequestBodyBufferValidRequestBody
缓冲策略选择worker/src/RequestBodyBuffer/RequestBodyBufferBuilder.tsRequestBodyBufferBuilderMAX_INMEMORY_BYTEStryInitRemote
内存缓冲实现worker/src/RequestBodyBuffer/RequestBodyBuffer_InMemory.tsunsafeGetRawTextisStreammodeluploadS3BodyDANGEROUS_KEYS
OpenAI 流还原worker/src/lib/dbLogger/streamParsers/openAIStreamParser.tsparseOpenAIStream
Anthropic 流还原worker/src/lib/dbLogger/streamParsers/anthropicStreamParser.tsanthropicAIStreamrecursivelyConsolidateAnthropicListForClaude
Vercel 流还原worker/src/lib/dbLogger/streamParsers/vercelStreamParser.tsparseVercelStream
通用字段合并worker/src/lib/dbLogger/streamParsers/responseParserHelpers.tsconsolidateTextFieldsrecursivelyConsolidate
核心日志对象worker/src/lib/dbLogger/DBLoggable.tsDBLoggableparseResponsegetDetailedUsageloguseKafka
队列消息类型worker/src/lib/clients/producers/types.tsMessageDataHeliconeMetaLog
投递门面 + 工厂worker/src/lib/clients/producers/HeliconeProducer.tsHeliconeProducerMessageProducerFactorysendMessageHttp
Kafka 通道worker/src/lib/clients/producers/KafkaProducerImpl.tsKafkaProducerImplsendMessage
SQS 通道worker/src/lib/clients/producers/SQSProducer.tsSQSProducerImplsetLowerPriority
非阻塞调度 + 边缘成本worker/src/lib/HeliconeProxyRequest/ProxyForwarder.tsloglogPromiseresponseProcessingPromise
成本聚合入口packages/cost/costCalc.tsmodelCostBreakdownFromRegistryCOST_PRECISION_MULTIPLIER
分项成本结构packages/cost/models/calculate-cost.tsCostBreakdowncalculateModelCostBreakdown
legacy 定价packages/cost/index.tscostOfPromptcostOf
用量解析器分派packages/cost/usage/getUsageProcessor.tsgetUsageProcessor
路由/PTB/BYOK 规则packages/cost/FLOWS.md(文档)

相邻章节: 上游怎么截获转发看 01 边缘代理与 AI 网关;这条消息进队列之后怎么被 Jawn 消费成数据库行看 03 消费与责任链;最终落到哪些库看 04 三处落库。全景导读见 index