SSE

什么是SSE?

SSE(Server-Sent Events,服务器推送事件) 是一种基于 HTTP 的长连接技术。 核心思想极其朴素:

普通的 HTTP 响应是"服务器一次性把整个 body 发完就断开"。 SSE 则是服务器保持连接不关闭,把数据分成一小块一小块持续发

OcEasy 用的是一种简化变体(也叫 NDJSON over SSE / 逐行 JSON):

  • 响应头标注 Content-Type: text/event-stream
  • 响应体是一行一个 JSON 对象,行之间用 \n 换行符分隔
  • 每次 yield 一行,数据就立刻到达浏览器,不等后面的行

与 WebSocket 的区别

特性SSEWebSocket
方向服务器 → 客户端(单向)双向
协议纯 HTTP,无需升级需要 101 协议升级
自动重连浏览器原生支持需自己实现
断点续传容易(带游标)较难
适用场景流式文本、进度、通知游戏、实时协作、聊天

聊天 AI 的逐字输出是 SSE 的典型场景:只需要服务端往客户端推文本,不需要客户端往服务器推数据(发消息走普通 POST)。

架构

img

即后端 LLM 每次生成一个 token → 包成 Packet → 序列化成"一行 JSON + 换行" → 立刻 flush 给浏览器; 前端用 response.body 的流式 reader 逐块读取 → 按换行切分 → JSON.parse → 累加到界面。

后端

这是流式聊天的 HTTP 端点。注意返回值类型是 StreamingResponse | ChatFullResponse,根据请求参数 stream 决定走流式还是非流式。

python
# chat_backend.py:681
def stream_generator() -> Generator[str, None, None]:
    state_container = ChatStateContainer()
    try:
        for obj in handle_stream_message_objects(...):
            yield get_json_line(obj.model_dump())   # 每个对象  一行 JSON
    except Exception as e:
        yield json.dumps({"error": str(e)})
    finally:
        logger.debug("Stream generator finished")

return StreamingResponse(stream_generator(), media_type="text/event-stream")

FastAPI 的 StreamingResponse 是流式输出的"发动机"

  • 它接受一个生成器(generator)
  • 每次生成器 yield 一个字符串,FastAPI(底层是 Starlette)就把它立即作为一块 HTTP chunk 发出去,不会攒起来。
  • 连接保持打开,直到生成器结束(StopIteration)或客户端断开。

这就是"后端不一次给完、前端能收到增量"的根源。

序列化

python
def get_json_line(
    json_dict: dict[str, Any], encoder: type[json.JSONEncoder] = OnyxJSONEncoder
) -> str:
    return json.dumps(json_dict, cls=encoder) + "\n"
  • json.dumps 把 Pydantic 对象 model_dump() 后的字典序列化。
  • 末尾加 \n——这个换行符是前端的"消息分隔符"(协议关键)。
  • 自定义 OnyxJSONEncoder 处理 datetimeUUID 等非原生 JSON 类型(server/utils.py:20)。

每个包在线上就是一行文本,例如: {"placement":{"turn_index":0},"obj":{"type":"agent_response_delta","content":"你"}}

数据模型

python
class Packet(BaseModel):
    placement: Placement                      # 定位信息(哪个 turn / 哪个模型)
    obj: Annotated[PacketObj, Field(discriminator="type")]   # 具体内容,按 type 区分

PacketObj 是一个type 字段判别的联合类型,几十种包包括:

包类型作用
agent_response_delta回答内容增量(真正打字的包)
agent_response_start回答开始(携带最终文档、预处理耗时)
reasoning_delta思考/推理内容增量
search_tool_*_delta搜索工具参数/结果增量
tool_call_argument_delta工具调用参数流式生成
citation_info引用文献信息
chat_heartbeat心跳(见第 6 节)
stream_stop流结束(含 stop_reason
error错误(StreamingError

Pydantic 的 discriminator="type" 意味着反序列化时前端必须带 type 字段, 前端 PacketType(lib.tsx:92)是同样的判别联合。

生成器链:包怎么从 LLM 流到 HTTP

调用链是层层嵌套的生成器,每一层把上层 yield 的东西透传给下层:

text
stream_generator (chat_backend.py:682)
   └─ handle_stream_message_objects (process_message.py:1752)
        └─ _stream_chat_turn (process_message.py:1528)
             └─ _run_models (process_message.py:1062)
                  └─ _read_stream (process_message.py:1496)  ← 实际对外 yield

_stream_chat_turn(process_message.py:1582-1661)的关键逻辑:

  1. 在一个短命 DB session 里执行 build_chat_turn,把准备阶段产生的包(如 session id、message id)先 yield 出去;
  2. 构造 StreamBufferWriter(流缓冲);
  3. 调用 _run_models 得到 run 流并 yield from

为什么先发准备阶段包? 前端需要立刻知道"这条消息分配到了哪个 session / 哪个 message id",才能在本地消息树上建出节点来承接后续的内容包。

逐Token

这是"字一个个出来"的最底层源头

python
for packet in llm.stream(
    prompt=llm_msg_history,
    tools=tool_definitions,
    tool_choice=tool_choice,
    max_tokens=max_tokens,
    reasoning_effort=reasoning_effort,
    ...
):
    delta = packet.choice.delta

    if not answer_start:
        yield Packet(obj=AgentResponseStart(final_documents=..., ...))
        answer_start = True

    accumulated_answer += content_chunk                    # 后端也累加一份
    state_container.set_answer_tokens(accumulated_answer)  # 用于最终持久化
    yield Packet(
        placement=_current_placement(),
        obj=AgentResponseDelta(content=content_chunk),     # 每块一个增量包
    )
  • llm.stream() 的实现封装在 backend/onyx/llm/multi_llm.py:981,底层是 litellmCustomStreamWrappermulti_llm.py:992),再往下是各家 LLM 的 HTTP 流。
  • 每次迭代只有一个 token(或一小段)。这就是"一个一个"的最小粒度。
  • 回答开始前还有一个 AgentResponseStart 包,携带 final_documents(检索到的文档)和 pre_answer_processing_seconds(从发消息到开始作答的耗时)。
  • 同时 accumulated_answer += content_chunk:后端维护一个累积文本,用于流结束后一次性持久化到数据库(前端是增量的,数据库要的是完整文本)。

回答带引用

真实项目里回答文本中常带 [1][2] 这样的引文标记。 代码通过 citation_processor.process_token(content_chunk)(llm_step.py:1258)逐 token 扫描,把属于引文标记的部分拆出来单独发 CitationInfo 包,正文里的标记换成占位符。前端再把引文号和文档 ID 对应起来,渲染成可点击的 1

并发

多个模型同时生产数据,一个writer统一收集和保存,一个reader负责把数据发给浏览器

_run_models(process_message.py:1062)是后端流式的"调度中心"。核心设计:

  • 每个模型一个 worker 线程_run_model,process_message.py:1209),各自的 Emitter 把 Packet 推进一个共享的无界队列 merged_queue
  • 一个 writer 线程_drain_to_completion)从队列 drain,负责两件事:
    1. 把每个包写入 StreamBufferWriter(供断线回放);
    2. 转发给 reader 的 tee 队列;
  • reader 生成器_read_stream,process_message.py:1496)从 tee 里取包 yield 给 HTTP。
text
worker0 ──Emitter──▶ merged_queue ──┐
worker1 ──Emitter──▶ merged_queue ──┼─▶ writer线程 ──▶ tee队列 ──▶ reader → yield → HTTP
worker2 ──Emitter──▶ merged_queue ──┘      │
                                          └─▶ StreamBufferWriter(持久化缓冲)

几个关键点:

  • 多生产者单消费者:多个模型并发产包,但最终合成一个 SSE 流(按到达顺序交错)。
  • reader 可以"死"而 writer 必须活:客户端断开(GeneratorExit)时,reader 直接 return,但 writer 继续把剩余包写进缓冲(process_message.py:1515-1523),这样断线的客户端重连后能回放到完整内容
  • reader_gone 事件标记 reader 已离开,tee 队列停止累积无人消费的包,避免内存泄漏。

前端

发消息

发消息:sendMessage(app/services/lib.tsx:145)

ts
const response = await fetch(`/api/chat/send-chat-message`, {
  method: "POST",
  headers: { "Content-Type": "application/json" },
  body,
  signal,                       //  取消信号
});

if (!response.ok) { /* 抛错 */ }

yield* withoutHeartbeats(handleSSEStream<PacketType>(response, signal));

最反直觉的一点fetch 在这里返回的是一个 Response,但 await fetch(...) 只等响应头,不等 body。因为 sendMessageasync function*(异步生成器),它把 Response 交给 handleSSEStream,调用方再 for await 逐包消费。

withoutHeartbeats(lib.tsx:207)是一个过滤层:

ts
async function* withoutHeartbeats(stream) {
  for await (const packet of stream) {
    if ("obj" in packet && packet.obj.type === "chat_heartbeat") continue;
    yield packet;
  }
}

心跳包不承载运行状态,消费端一律看不见它(心跳本身干什么,见下文)。

核心解析器

步骤原因
getReader()response.body 是一个 ReadableStream,用 reader 才能异步增量读取,而不是 await response.json() 一次性拿完
TextDecoder网络传的是字节(Uint8Array),需要解码成字符串
buffer网络数据可能切在任意位置(一个 JSON 的后半段),需要暂存"不完整的尾巴"
reader.cancel()AbortController 触发时,必须主动 cancel 流,否则底层连接挂起(见第 5 节)
{ stream: true }UTF-8 多字节字符可能被切到字节边界stream:true 让解码器缓存不完整字节、下次补全。不传会乱码!
⑦⑧ 按 \n 切行协议规定一行一个 JSON。split 后最后一段没有换行符结尾,说明是半包,留回 buffer
finally cancel消费方 break/return 提前退出生成器时,只有 cancel 才会真正断开 SSE

💡 步骤 ⑦⑧ 是所有流式协议通用的"粘包/半包处理"思路,在 WebSocket 消息、NDJSON 日志流、日志采集里都能见到。

缓冲层

ts
export class CurrentMessageFIFO {
  private stack: PacketType[] = [];
  isComplete: boolean = false;
  error: string | null = null;

  push(packetBunch) { this.stack.push(packetBunch); }
  nextPacket() { return this.stack.shift(); }   // FIFO:先入先出
  isEmpty() { return this.stack.length === 0; }
}

export async function updateCurrentMessageFIFO(stack, params) {
  try {
    for await (const packet of sendMessage(params)) {
      if (params.signal?.aborted) throw new Error("AbortError");
      stack.push(packet);
    }
  } catch (error) { /* 记入 stack.error */ }
  finally { stack.isComplete = true; }
}

为什么需要 FIFO?因为 for await 的生产速度远快于 React 渲染速度

如果直接在每个包上 setState,React 会被几万个重渲染压垮。

FIFO 把"生产"(后台异步迭代,持续 push)和"消费"(主循环按帧取包渲染)解耦。

消费循环

主循环 while (!stack.isComplete || !stack.isEmpty()) 是前端流式消费的心脏:

Object.hasOwn(packet, "xxx") 判断包类型——不同包有不同的顶层字段,这是项目里判别"这是哪种包"的惯用手法(对应后端 Pydantic 的判别联合)。

渲染性能三件套

flushViaRAF 批量刷新

requestAnimationFrame 把一帧内的所有包合并成一次 React 更新:

ts
await delay(50);
while (...) {
  if (stack.isEmpty()) {
    if (pendingFlush) await flushViaRAF();  // 攒到下一帧再刷
  }
  ...
}

packetCount 代替数组做 memo(AgentMessage.tsx:73)

React.memo 比较 props。packets 数组被原地 push(引用不变),所以 memo 直接比较数组会认为"没变"。项目用 packetCount 这个原始数字做比较:

tsx
prev.packetCount === next.packetCount &&
prev.chatState.agent === next.chatState.agent &&
... // 其余字段
ts
// useChatSessionController.ts:319 的注释点明了设计意图:
// "AgentMessage's memo compares packetCount, not the packets array."
node.packetCount = accumulated.length;

③ trailing flush:burst 结束不饿死最后几个包useChatSessionController.ts:346:如果距上次 flush 不足 100ms 就设一个 120ms 的定时器,保证"突发数据的最后几个包"即使等不到下一个包,也会被定时刷出来。

渲染层

AgentMessage 拿到原始 packets 后交给这个 hook 处理成 UI 数据:

ts
// 处理状态放在 ref 里:增量、同步、不触发双渲染
const stateRef = useRef<ProcessorState>(createInitialState(nodeId));
// 只有真正的 UI 状态才用 useState
const [renderComplete, setRenderComplete] = useState(false);

// 增量处理:只处理新到的包(避免重复处理已消费的)
if (rawPackets.length > stateRef.current.nextPacketIndex) {
  stateRef.current = processPackets(stateRef.current, rawPackets);
}

processPackets(packetProcessor.ts)遍历新包:

  • agent_response_delta → 累加进 potentialDisplayGroups 的文本块;
  • 工具类包 → 归类成 toolGroups(步骤/轮次分组,供时间线 UI 展示"思考了几步");
  • message_start → 标记 finalAnswerComing
  • stream_stop → 标记 stopPacketSeenstopReason
  • citation_info → 进 citationMap

最终回答文本由 displayGroups 渲染成 Markdown,每次新 delta 到达 → rawPackets.length 增加 → 重新 processPackets → 重渲染。 打字机效果 = 增量数据 + 高频重渲染

心跳机制

LLM 生成可能长时间"没动静"(比如调工具、检索文档),但 HTTP 长连接如果长时间没有数据,中间代理/负载均衡可能把连接掐掉。 心跳(chat_heartbeat)就是用来保活的。

产生方(后端)

主 stream 的 reader 在**空闲超过 CHAT_HEARTBEAT_INTERVAL_S(默认 15 秒)**时发一个心跳包(process_message.py:1500-1508):

python
while True:
    try:
        item = tee.get(timeout=_CANCEL_POLL_INTERVAL_S)
    except queue.Empty:   # 队列空 = 没有真实包到来
        now = time.monotonic()
        if now - last_packet_yield >= CHAT_HEARTBEAT_INTERVAL_S:
            yield heartbeat_packet()          # 发个心跳保活
            last_packet_yield = now
        continue
    ...

heartbeat_packet()(streaming_models.py:496):

python
def heartbeat_packet() -> Packet:
    """Keepalive for silent stretches; carries no run state."""
    return Packet(placement=Placement(turn_index=0), obj=ChatHeartbeat())

CHAT_HEARTBEAT_INTERVAL_S = int(os.environ.get("CHAT_HEARTBEAT_INTERVAL_S") or "15")(chat_configs.py:35)——可用环境变量调整。

消费方(前端)

  • withoutHeartbeats(lib.tsx:207)在发消息链路里直接丢弃心跳;
  • resume 链路里心跳被用作liveness tick(useChatSessionController.ts:341): 收到心跳说明连接还活着、会话还是当前焦点,用于决定是否继续刷新;
  • 渲染层完全看不见心跳包。

断线重连与流回放

SSE 是单向长连接,任何网络抖动都可能断。OcEasy 用"持久化缓冲 + 游标回放"实现断线续传。

后端

写路径:_run_models 的 writer 线程把每个包 get_json_line同时写入缓冲:

python
# process_message.py:1650
for pre_run_packet in pre_run_packets:
    stream_buffer.append_line(get_json_line(pre_run_packet.model_dump()))

append_line(stream_buffer.py:80)攒到一定字节再批量 flushflushzlib 压缩后按 chunk_count 分块写进共享缓存(Redis),带 TTL:

python
def flush(self) -> None:
    payload = zlib.compress("".join(self._pending).encode("utf-8"))
    self._cache.set(
        _chunk_key(self._chat_session_id, self._run_id, self._meta.chunk_count),
        payload,
        ex=CHAT_STREAM_BUFFER_TTL_S,
    )

关键属性:

  • truncated:缓冲超上限(CHAT_STREAM_BUFFER_MAX_BYTES)就标记截断,降级为"不可恢复";
  • done:流正常结束后 mark_done()(stream_buffer.py:131),延长数据保留期, 让刷新页面后的客户端能回放完整内容。
python
def resume_chat_stream(
    session_id: UUID,
    cursor: int = Query(0, ge=0),     # 客户端上次读到哪
    user: User = Depends(...),
) -> StreamingResponse:
    ...
    run_id = get_processing_run_id(session_id, cache)
    if run_id is None or not has_stream_buffer(cache, session_id, run_id):
        raise OnyxError(NOT_FOUND, "No resumable run for this chat session")

    def stream_buffered_run() -> Generator[str, None, None]:
        chunk_cursor = cursor
        while True:
            read = read_stream_chunks(cache, session_id, run_id, chunk_cursor, ...)
            if read is None or read.gap: return      # 缓冲丢了  终止,客户端重新拉全量
            if read.blocks:
                yield "".join(read.blocks)           # 回放缓冲内容
                chunk_cursor = read.next_cursor
                continue
            if read.done: return                     # 正常结束
            if not is_chat_session_processing(session_id, cache):
                # writer 死了但没标 done  再补一次 drain
                ...
            # 缓冲读完了  实时 tail 新的包
            time.sleep(CHAT_RESUME_POLL_INTERVAL_S)
            if now - last_emit >= CHAT_HEARTBEAT_INTERVAL_S:
                yield get_json_line(heartbeat_packet().model_dump())  # 轮询中也发心跳
    ...

设计要点:

  • 缓冲在共享缓存里,任何 pod 都能服务(注释:"Serves any pod");
  • 404(无可恢复 run)是正常业务信号,前端回退到重新拉取落库 session;
  • 回放结束进入实时 tail:轮询新 chunk + 心跳保活。

前端

ts
export async function* resumeStream(chatSessionId, cursor, signal?) {
  const response = await fetch(
    `/api/chat/chat-session/${chatSessionId}/resume-stream?cursor=${cursor}`,
    { signal }
  );
  ...
  yield* handleSSEStream<PacketType>(response, signal);
}

页面刷新后,如果 session 有 current_run 且是 assistant 节点,就"粘"上去续读:

ts
const abortController = new AbortController();
try {
  for await (const rawPacket of resumeStream(sessionId, 0, abortController.signal)) {
    if (!stillCurrent()) return;              // 用户已切走会话      if (packet.obj.type === "chat_heartbeat") continue;
    accumulated.push(packet);
    //  100ms  trailing 120ms 刷一次到消息树
  }
} finally {
  abortController.abort();
  // 流结束后拉一次落库的完整 session 做"结算"
  const settled = await (await fetch(`/api/chat/get-chat-session/${sessionId}`)).json();
  updateSessionAndMessageTree(sessionId, processRawChatHistory(settled.messages, ...));
}

"流式 + 落库"双保险:在线期间用增量流;断线/刷新后用回放流;回放结束用落库数据最终对齐。

新故事即将发生
python_essay

评论区

评论加载中...