Skip to content

流式解码 ​

stream 模块回答一个问题:把 OpenAI 兼容 / Anthropic 的 SSE 字节流,解成一套统一的事件。

0.1.4 起可用,在默认的 chat feature 里,没有新增依赖,移动端也能编。

为什么放进库里 ​

每个做对话的应用都要写一套 SSE 解析,而各家网关的怪癖随服务商变化:

  • 换行用 \r\n,或者中途插 : keep-alive 心跳行
  • 工具调用缺 index,或者 id、名字晚到
  • 只以 [DONE] 收尾、不给结束原因
  • 流断在半路,却没有任何报错
  • 流里夹着 {"error":…}
  • 返回 2xx,内容却是网页

同一份知识复制 N 份,必然漏改。所以把「字节 → 事件」这一层收进来,让所有应用共用一份修过的实现。

🔴 只吃字节、吐事件 ​

本模块不发 HTTP、不持有连接、不绑任何异步运行时。这是它没有变成又一个 LLM SDK 的原因。

在库里留在你的应用里
SSE 帧解析、各家网关怪癖请求体构造
两种协议统一成同一套事件HTTP 客户端、代理、超时
断流 / 取消 / 流内错误的收尾判定取消信号怎么接线
工具调用按块号拼装成完整 content推给前端、执行工具、存历史
is_stream_options_rejected 判定函数「某地址不带 stream_options」的记忆状态

接入 ​

rust
use ai_profile::stream::{StreamDecoder, StreamEvent};
use ai_profile::Protocol;

let mut dec = StreamDecoder::new(Protocol::OpenAiCompatible).with_stream_id("chat-42");

// 用你自己的 HTTP 客户端读响应,读到一段就喂一段,切在哪里都行
while let Some(chunk) = body_stream.next().await {
    for ev in dec.push(&chunk?) {
        match ev {
            StreamEvent::TextDelta { text, .. } => ui.append(&text),
            StreamEvent::ToolUseStart { name, .. } => ui.show_tool(&name),
            _ => {} // 事件枚举会新增成员,别漏了这个兜底
        }
    }
}

// 读到 EOF:补发最后的事件,拿到收尾结果
let (tail, outcome) = dec.finish();

用户点了「停止」时,不调 finish,改调 abort:

rust
let outcome = dec.abort(); // 返回已收到的文字,工具调用一律丢

事件 ​

事件含义
TextDelta正文增量,带块号
ReasoningDelta思考增量,不占块号、不进最终内容
ToolUseStart一个工具调用开始(id、工具名)
ToolUseDelta工具参数的 JSON 片段,拼起来才是完整参数
Usage用量变化(累计值,不是增量)
Finish流正常结束时的最后一个事件,带归一化的结束原因
Error流内错误,终态,之后的输入被忽略

序列化为 JSON 时,事件用 kind 区分,字段名是 camelCase,可以直接推给前端。

块号 ​

协议文字第 k 个工具调用
OpenAI 兼容固定是块 0块 k+1
Anthropic协议自带的 index协议自带的 index

两种协议因此对齐成同一套编号,前端不用区分。

工具调用何时开始 ​

ToolUseStart 等到工具名已知(并且 id 已到或参数已经开始)才发,所以前端一收到就能显示「正在调用 xxx」。 始终没等到的,在 finish() 时补发。

网关没给 id 的,库会补一个 call_<流id>_<块号>,补出来之后不再改变。 想让补出来的 id 在应用重启后也不重复,传应用自己的流 id(with_stream_id)。

收尾 ​

finish() 返回的 StreamOutcome 里,end 说明流是怎么结束的:

情形end内容
有 finish_reason / stop_reason / message_stopComplete全部块
只有 [DONE],没有结束原因Complete有工具调用记为 tool_use,否则 end_turn
既无结束原因也无 [DONE]Truncated🔴 只留文字,丢全部工具调用
调用了 abortCancelled同上
流内 {"error":…} / Anthropic error 事件Failed只留文字,之后的输入忽略
2xx,但整段没有一行 SSENotEventStream带原始响应体(最多 1 MiB)

断流时不要执行工具

Truncated / Cancelled / Failed 时,此前已经发出的 ToolUseStart 对应的工具调用是不完整的。 按 end 判断,收回即可,不要执行。

结果里的其它字段:

字段内容
stop_reason归一化的结束原因,只有 Complete 才有。OpenAI 的 stop / tool_calls / length 对应 end_turn / tool_use / max_tokens
contentAnthropic 风格的 content block 数组,与历史裁剪用的消息是同一种形状,可以直接存进历史
usage取流里最后一次给出的值,覆盖不累加
reasoning累积的思考文字
model流里上报的模型名
skipped_frames解析失败被跳过的帧数,排查网关问题用

判断要不要跑工具,以 content 里有没有 tool_use 块为准。 个别网关带着工具调用却给 finish_reason: "stop",stop_reason 照实透传,不替它改写。

已经吸收的网关怪癖 ​

类别处理
换行与格式\r\n、单独的 \r、BOM、data: 后没有空格、多行 data、事件之间漏空行
工具调用缺 index(按 id 或位置分块)、index 全是 0 但 id 不同、id 或名字晚到、arguments 是对象而不是字符串
结束finish_reason 为空串、末尾只带 usage 的帧、只以 [DONE] 收尾
用量后面的值覆盖前面的,不累加
其它多个 choices 只读第一个、reasoning_content / thinking_delta 思考内容、Anthropic 缺 event: 行时按 data.type 分派
流内错误统一成终态事件

多字节字符不会丢字 ​

解码按完整一行进行,所以把一个汉字或一个 \r\n 切在两个网络包中间也不会乱码, 无论怎么分包,结果完全一致(有专门的守卫测试)。

不做的事 ​

不做原因
Ollama 原生 NDJSON不是 SSE,留给应用
剥 <think> 标签、「正文为空就把思考提升为正文」产品取舍
Anthropic 的 thinking 块及其 signature开启扩展思考并配合工具调用时,要把 thinking 块原样回传的应用需要自己处理
缓存用量字段暂不需要

辅助函数 ​

函数用途
is_stream_options_rejected(status, body)服务端回 400 / 422 且提到 stream_options 或 include_usage 时为真。去掉该参数重试一次,并记住这个地址不支持
looks_like_html(body)响应体像不像一个网页,用来区分「地址填成了网站首页」
StopReason::from_openai(s)OpenAI 兼容的 finish_reason 归一化

其他语言 ​

规范里有对应的 stream.json,含分包、断流、取消、流内错误等用例,见其他语言实现。

相关 ​

MIT 协议开源 · 文档同样欢迎 PR · B 站视频介绍 · QQ 交流群 1087715758