OncallAgent 的普通聊天是一条由 LangChain Agent 驱动的流式链路。后端不会先无条件执行 RAG,而是把知识检索、当前时间、渐进式 Skill 和当前用户已发现的 MCP 工具交给模型,由模型决定是否调用。最终答案、模型提供的推理元数据、工具生命周期和引用通过统一 SSE 合同发送给 Vue 前端。

📷 [图片 token=Ih9AbC9jtoZLN2x6QZtcbNy9nrg(未能下载,见飞书原文)]

SSE 的意义不只是“打字机效果”。它把一次长模型执行拆成可判别事件:客户端能知道当前收到的是答案字符、推理片段、工具启动、工具完成、引用、终态还是结构化错误。服务端同时把用户消息、完整助手消息、引用、推理和工具调用审计写入 SQLite,使流式临时状态最终回归持久事实。

这条链路也有清晰边界:聊天流不是 durable job event log,断线后不能用 Last-Event-ID 从中间续传;前端在收到 complete 后会重新读取会话与审计,完成最终对账。流中失败可能已经把部分字符显示给用户,但不会把部分助手消息保存到历史。

📷 [图片 token=VADtbfHfZoSTPBxBK3ycwJZxnTe(未能下载,见飞书原文)]

学习目标

  • 理解 create_agentastream_events 到共享 SSE 的适配过程。

  • 掌握 SSE 的 type 判别字段和每类 payload。

  • 追踪 tool.call 的 stable ID、started、completed、failed 与 SQLite 审计。

  • 理解逐字符 content、推理、引用、complete 和 error 的顺序与持久化语义。

  • 识别浏览器解析、前端草稿、终态对账和断线恢复的实际边界。

功能入口与完整调用链

前端 createChatClient.streamMessagePOST /chat/sessions/{session_id}/messages:stream 发送 JSON,并通过 createSseClient 加上 bearer token 和 Accept: text/event-stream。后端路由 stream_chat_message 先按当前 user.id 查询会话;会话不属于当前用户时,直接返回 HTTP 403,Agent 不会运行。

ChatStreamingService.stream_message 校验内容、读取历史、装配系统提示词和 Skills,并在上下文上限检查通过后先持久化用户消息。它构造 ChatAgentRequest,把 owner、session、模型上下文、授权知识库和配置交给 LangChainChatAgentRunner

Runner 创建当前请求专属工具集合,调用 langchain.agents.create_agent,再遍历 agent.astream_events(version="v2")_agent_event_from_langchain_event 把 LangChain 原始事件转成内部 dataclass;服务进一步生成共享 SSE payload,路由用 encode_sse 编码为 event/data 帧。前端 store 按 event.type 更新草稿、推理、实时工具和引用,收到 complete 后重新读取会话与审计。

📷 [图片 token=ESBfbjSeroBRQ3xkO7gcBX54nJh(未能下载,见飞书原文)]

ChatView → useChatStore.send
  → ChatClient.streamMessage
  → POST messages:stream
  → owner-scoped session check
  → ChatStreamingService 持久化 user message
  → LangChainChatAgentRunner
  → create_agent + astream_events v2
  → 内部 Agent event
  → shared SSE payload
  → encode_sse
  → sseClient.parseSseFrames
  → Pinia 草稿状态
  → complete 后重新加载 SQLite 会话与工具审计

核心源码地图

源码位置关键符号职责
apps/backend/src/super_ai/chat/streaming.pyChatStreamingServiceLangChainChatAgentRunnerencode_sseAgent 运行、事件适配、消息持久化和 SSE 编码。
apps/backend/src/super_ai/api/app.pystream_chat_message先做 owner 会话鉴权,再返回 StreamingResponse。
packages/api-contracts/src/sse.tsSseEventSSE_EVENT_TYPESToolCallSseEvent定义事件判别联合、channel 和各 payload。
packages/api-contracts/src/chat.tsStreamChatMessageRequestChatStreamCompleteResultToolCallAudit定义请求、最终结果和持久审计 DTO。
apps/frontend/src/api/sseClient.tscreateSseClientparseSseFrames认证 fetch、流式解码、frame 分割和基础事件校验。
apps/frontend/src/chat/chatClient.tscreateChatClient组合普通 HTTP 与聊天 SSE API。
apps/frontend/src/stores/chat.tsuseChatStore.sendupdateLiveToolCalluniqueReferences逐字渲染、实时工具、引用聚合和 complete 对账。
apps/frontend/src/views/ChatView.vueactiveTitleopenCitationDocumentsaveChatConfiguration呈现会话、消息、工具、引用与输入状态。
apps/backend/tests/test_stream_rag_chat_api.pytest_streaming_chat_emits_sse_events_and_persists_messages端到端验证顺序、逐字符、持久化、审计与安全日志。
apps/frontend/tests/chatStore.test.tsreconciles streamed content, references, and tool calls with persisted history验证临时流状态最终以服务端历史为准。
openspec/specs/stream-rag-chat/spec.mdChat stream event sequence规定 Agent 决策工具、事件顺序、持久化与失败语义。

代码调用流程图

聊天链路包含两个事实面:SSE 负责即时过程,SQLite 保存完成后的消息和审计。complete 到达后,前端重新加载持久数据完成对账。

画板

关键实现拆解

从 LangChain 原始事件到稳定领域事件

**看什么:**看 Runner 如何把 LangChain v2 原始 event 收敛成项目内部 dataclass;未知框架事件不会直接穿透到共享 SSE。

        # 1. 每次请求独立消费 LangChain v2 事件协议。
        async for raw_event in agent.astream_events(
            cast(Any, {"messages": messages}),
            version="v2",
        ):
            # 2. 适配层把外部事件收敛为稳定的项目领域事件。
            parsed = _agent_event_from_langchain_event(cast(Mapping[str, object], raw_event))
            if parsed is None:
                continue
            if isinstance(parsed, list):
                for item in parsed:
                    yield item
            else:
                yield parsed

适配器只识别工具开始、完成、失败以及模型内容/推理流;未知事件被忽略,不会迫使前端跟随 LangChain 内部命名变化。run_id 通常成为稳定 tool call ID,缺失时才生成新值;适配层没有生成 tool delta 状态。

📷 [图片 token=CytCbRt34oYmhQxdpamcszU3nlk(未能下载,见飞书原文)]

LangChainChatAgentRunner 每次请求创建知识检索工具和当前时间工具;选中 Skills 时加入 load_skill;配置了用户 MCP connection service 时,只加载该用户已启用且发现成功的工具。Runner 把 SQLite 消息转成 user/assistant 消息列表,并把最终 system prompt 传给 create_agent。普通聊天没有自定义 LangGraph 状态图。

📷 [图片 token=KCAab5Frioh06WxNelUcjX7wnBb(未能下载,见飞书原文)]

_agent_event_from_langchain_event 只处理明确事件:on_tool_start 生成 started,on_tool_end 生成 completed 并从工具输出抽取 citations,on_tool_error 生成 failed,模型或 model chain 的 stream 事件抽取 content 和 reasoning。工具调用 ID 来自 LangChain run_id;缺失时才生成新的 ID,因此开始和终态通常能用同一 ID 对账。

📷 [图片 token=AMsSbpkD9oycZUxOxSJcqfk3nyd(未能下载,见飞书原文)]

共享契约允许 tool status 为 started、delta、completed、failed,但当前 LangChain 适配器实际生成 started、completed 和 failed,没有生成 delta 的分支。契约预留能力不能被描述成已经存在的运行时行为。类似地,推理只从模型 chunk 的 reasoning_contentreasoning 读取;模型不提供时,服务不会合成“思考过程”。

SSE 判别字段与事件粒度

**看什么:**先看模型内容 delta 在领域服务中如何被拆成单字符事件;sequence 同时被内容和 reasoning 共享,但工具与引用不带 sequence。

            async for event in self._agent_runner.stream(request):
                if isinstance(event, ChatAgentContentDelta):
                    if event.delta == "":
                        continue
                    answer_parts.append(event.delta)
                    # 1. 模型 delta 可能多字符,服务主动逐字符拆分。
                    for character in event.delta:
                        sequence += 1
                        yield _sse_event(
                            "content.delta",
                            {
                                "delta": character,
                                "sequence": sequence,
                            },
                        )

持久化答案使用未拆分的 answer_parts 拼接,因此字符动画不会改变最终正文。sequence 只在当前请求内递增,前端目前也不按 sequence 重排或补洞;它不是跨连接可重放的全局游标。

**看什么:**再看每个领域事件如何补齐四个基础判别字段,并编码成标准 event/data/空行帧。

def encode_sse(event: Mapping[str, object]) -> str:
    """Encode a shared SSE event payload as one text/event-stream frame."""
    event_type = str(event["type"])
    # 1. event 行与 data 内 type 来自同一个值。
    return f"event: {event_type}\ndata: {json.dumps(event, separators=(',', ':'))}\n\n"

def _sse_event(event_type: str, payload: Mapping[str, object]) -> dict[str, object]:
    # 2. 所有聊天事件共享身份、类型、channel 与时间戳。
    return {
        "id": f"evt_{uuid4().hex}",
        "type": event_type,
        "channel": "chat",
        "timestamp": _now_iso(),
        **payload,
    }

前端 parser 实际忽略 event 行并信任 data JSON 的 type,所以后端必须保持二者一致。基础字段构成运行时最低守卫,逐类 payload 主要依靠共享 TypeScript 判别联合和后端契约测试,当前没有完整 runtime schema 验证。

每个事件都有 idtypechanneltimestamp。聊天 channel 固定为 chat;type 是判别联合的核心。content.deltareasoning.delta 带 delta 与 sequence;tool.call 带 toolCall;reference.source 带 reference;complete 带可选 result;error 复用统一 ApiErrorMessage

📷 [图片 token=OQiSbX7yAobokHx2btlcS7G5nLP(未能下载,见飞书原文)]

Agent runner 可能一次给出多字符内容,但 ChatStreamingService 会遍历每个字符,逐一发出 content.delta。sequence 是服务内的递增计数,reasoning 事件也会推进它;工具和引用事件不带 sequence。答案字符被加入 answer_parts,推理加入 reasoning_parts,引用转成详细 payload,工具 ID 去重累积。

📷 [图片 token=SBhNbrd3UorsDvxFJNYciBnAnXb(未能下载,见飞书原文)]

encode_sse 同时写 event: {type} 和一行 JSON data。前端 parser 以空行切帧,拼接多条 data 行并 JSON.parse。运行时守卫目前只核对 id、type、channel、timestamp 是字符串,并没有逐类验证整个 payload;静态 TypeScript 联合提供编译期约束,后端测试负责契约形状。若将来面对不可信 SSE 源,应增加完整 runtime schema 校验。

工具审计、引用和完成对账

**看什么:**看 tool event 在发 SSE 前先尝试写审计,以及审计失败为何被刻意隔离,不阻塞主答案。

        try:
            if event.status == "started":
                # 1. started 先创建 owner-scoped 审计,再允许上层发 SSE。
                await repository.create_for_chat_session(
                    owner_user_id=owner_user_id,
                    audit_id=event.id,
                    chat_session_id=session_id,
                    tool_name=event.name,
                    arguments=_json_dict_or_empty(event.input),
                )
                return
            if event.status == "completed":
                audit = await repository.finalize(
                    owner_user_id=owner_user_id,
                    audit_id=event.id,
                    status="completed",
                    result_summary=_audit_summary(event.output),
                )
                if audit is None:
                    # 2. 缺少 started 记录时补建并完成同一 ID。
                    await self._create_and_finalize_missing_audit(
                        owner_user_id=owner_user_id,
                        session_id=session_id,
                        event=event,
                        result_summary=_audit_summary(event.output),
                    )
                return
            # … 省略 failed 终态的脱敏摘要分支
        except Exception:
            # 3. 审计是旁路可观察性,失败不抑制聊天输出。
            return

审计和 SSE 不是同一事务,因此极端情况下 UI 看到了 tool.call 而历史审计缺项。相反,assistant message 是 complete 的主事实:只有答案、引用和工具 ID 成功写入 SQLite 后,服务才发 complete;引用从同一工具 output.citations 结构转换,不重新执行检索。

服务在发出每个 ChatAgentToolCall 前先调用 _persist_tool_call_audit。started 创建 owner-scoped 审计,保存 session、tool name 和 arguments;completed 或 failed 以相同 ID finalize,保存有限结果摘要或脱敏错误摘要。若只收到终态而没有已有记录,服务会补建再完成。审计持久化异常被吞掉,以免抑制聊天输出;因此 SSE 生命周期和审计属于尽力对齐,而不是同一数据库事务。

📷 [图片 token=VHjHbkyuMo2ZaBxH2q2cWz4pnkg(未能下载,见飞书原文)]

知识检索工具完成时,适配器先发 completed tool event,再从 output.citations 转换为一个或多个 ChatAgentReference。reference 可携带 chunk、document、knowledge base、来源、metadata、excerpt、knowledgeType,以及 vector、BM25、RRF、rerank 分数与排名。前端 uniqueReferences 按 ID 去重,优先 rerankScore、其次 score 排序,最多保留 5 条。

📷 [图片 token=MUX1bqxXLoegCgx9FTsc6Xc5n3l(未能下载,见飞书原文)]

只有 Agent 迭代正常结束且助手消息成功写入 SQLite 后,服务才发 complete。complete.result 包含刷新后的 session 和记忆状态,以及真正持久化的 assistant message。前端把流中内容视为草稿,收到 complete 只标记 finished;循环结束后并行 loadSessionreloadSessions,用后端历史和工具审计替换临时状态。

📷 [图片 token=JqoAbebGToyVQ0xTgxGcxRcknfb(未能下载,见飞书原文)]

按时间线观察一次工具型回答

**看什么:**这张序列图展示一次可能的工具型回答,而不是规定所有请求都必须出现每一种事件。

画板

工具开始和终态通常用同一 run_id 对账,引用紧随包含 citations 的工具完成事件。纯聊天可跳过工具和引用,模型不提供 reasoning 就没有 reasoning.delta;前端必须按 type 分支,不能依赖固定位置。

一次典型知识问答可能先产生 reasoning.delta,表示模型实际提供了准备检索的推理元数据;随后 on_tool_start 产生 tool.call started,审计记录在对应 SSE 发出前创建。知识工具完成时先出现 tool.call completed,接着是若干 reference.source。最后模型依据工具结果生成答案,服务把一个模型 chunk 拆成逐字符 content.delta,并在助手落库后发 complete。

事件顺序由 Agent 实际运行决定,并非所有回答都必须包含上述类型。纯聊天可以只有 content.delta 和 complete;工具失败可能出现 started、failed,随后 Agent 选择继续回答,也可能让整个流进入 error;模型不返回 reasoning 时不会出现 reasoning.delta。前端因此必须按 type 独立处理,不能依赖固定数组位置或假设第一帧一定是内容。

📷 [图片 token=FUVabNUoVoc4Izxo2dcc5IK2nag(未能下载,见飞书原文)]

sequence 主要帮助保持内容与推理增量的顺序,但 SSE 基于单一 HTTP 流,本身已经按字节顺序到达。当前前端并没有检查 sequence 连续性或重排事件;测试确保服务生成递增值。若未来引入并行模型通道或断点续传,sequence 的范围和去重语义需要进一步定义,不能直接把当前请求内计数当作全局事件编号。

前端草稿为何还要终态刷新

**看什么:**看 Pinia store 如何把流式事件写入临时草稿,却把 complete 仅当作完成标志;流结束后仍以服务端 session 和消息列表为准。

        let finished = false;
        for await (const event of client.streamMessage(targetSessionId, { content })) {
          if (event.type === "content.delta") {
            // 1. 后端虽已逐字符,前端仍对 delta 做兼容遍历。
            for (const character of event.delta) {
              updateAssistantDraft(targetSessionId, draftId, character, messages);
              await waitForTypewriterTick();
            }
          }
          //  省略 reasoning.delta 分支
          if (event.type === "reference.source") {
            references.value = uniqueReferences([...references.value, event.reference]);
          }
          if (event.type === "tool.call") {
            updateLiveToolCall(event.toolCall, liveToolCalls);
          }
          if (event.type === "error") {
            throw new ApiClientError(event.error);
          }
          if (event.type === "complete") {
            // 2. complete 只标记流完整,随后仍重新读取后端事实。
            finished = true;
          }
        }
        if (!finished) {
          throw new Error("回答流在完成前意外中断。");
        }
        await Promise.all([loadSession(targetSessionId), reloadSessions()]);

原逻辑表明草稿 ID、owner 和字符节奏都不是持久事实。若流中断或 error,draft 会被删除;重新读取后,引用和工具审计也从持久消息与审计 API 对账。

📷 [图片 token=YPJcbbuCiot1qNxbGtuceyJfnRf(未能下载,见飞书原文)]

Pinia store 在发送前创建 optimistic user message,ownerUserId 临时写为 current,ID 带 optimistic 前缀。助手草稿也使用本地 draft ID。它们只用于即时界面,真正的 message ID、owner、createdAt、完整 metadata 都由服务端 SQLite 记录决定。complete 到达后重新加载,能消除本地 ID、字符节奏和服务端最终持久化之间的差异。

references 在流中独立维护,assistant draft 的 metadata 初始 citations 为空。最终 reload 后,setReferencesFromMessages 从最近一条持久化 assistant message 重新构造引用。工具调用同样分为 liveToolCallstoolAudits:前者服务进行中状态,流结束后清空;后者来自审计 API,可在重新打开会话时恢复。

📷 [图片 token=LxUQb4w2LoOHakxUnt8cXwGqnab(未能下载,见飞书原文)]

逐字符效果有两层。后端已经保证每个 content.delta 只有一个字符;前端仍对 event.delta 再遍历字符,并在每个字符后等待 28 毫秒。这使测试注入多字符事件时仍能平滑显示,也为契约演进提供容错。代价是长回答在网络已结束后仍可能等待 UI 动画;store 只有消费完全部字符后才处理后续事件,所以视觉 complete 可能晚于服务端完成时间。

错误发生在不同阶段时会留下什么

**看什么:**这张状态图把错误发生点与 SQLite 中可留下的消息分开,重点观察 user message 持久化前后的边界。

画板

鉴权失败是普通 HTTP 错误;准备阶段失败是 SSE error 且不写 user message。进入 Agent 后失败会保留 user message但不保留部分 assistant,前端见过的草稿会被删除;只有 assistant 写入成功的路径能够发 complete。

鉴权失败发生在 StreamingResponse 建立前,客户端收到普通 HTTP 403 和统一错误 envelope,不会进入 SSE parser。空消息或 95% 上下文硬上限则由 stream service 生成 error SSE;此时没有用户消息落库。用户消息成功保存后,Agent 或工具再失败,错误作为 SSE 发出,历史保留用户问题。

如果已经输出部分答案后发生异常,浏览器先看到若干 content.delta,再看到 error。服务不会创建 assistant message,前端 catch 会删除 draft。重新加载会话时只剩 user message。这种行为避免把不完整输出长期当成正式答案,但用户短暂看到过的内容无法从屏幕体验中“撤销已阅读”,所以 UI 必须明确标出失败,而不能只悄悄消失。

更晚的失败是 assistant persistence 自身出错。即使模型已经完整返回,append_message 失败仍会进入 error,绝不发 complete。相反,如果工具审计写入失败,服务刻意继续 SSE 和答案持久化,因为审计是旁路可观察性;最终审计列表可能缺项。这两个选择体现了主事实优先级:正式 assistant message 是 complete 的前提,审计不是。

📷 [图片 token=JE3KbzwSnolv7ixMHB0c1xp9nrc(未能下载,见飞书原文)]

传输协议的细节与限制

**看什么:**看浏览器端如何用 streaming TextDecoder 保留半个 UTF-8 字符或半帧,并只在空行边界解析完整 SSE frame。

      const reader = response.body.getReader();
      const decoder = new TextDecoder();
      let buffer = "";
      for (;;) {
        const { done, value } = await reader.read();
        if (done) {
          break;
        }
        // 1. streaming decode 避免网络 chunk 截断 UTF-8 字符。
        buffer += decoder.decode(value, { stream: true });
        const parsed = parseSseFrames(buffer);
        buffer = parsed.remainder;
        yield* parsed.events;
      }
      buffer += decoder.decode();
      const parsed = parseSseFrames(buffer);
      yield* parsed.events;
    }
  };
}

function parseSseFrames(buffer: string): {
  readonly events: readonly SseEvent[];
  readonly remainder: string;
} {
  const frames = buffer.split(/\r?\n\r?\n/);
  // 2. 最后一段不完整 frame 留给下一次读取。
  const remainder = frames.pop() ?? "";

客户端用 fetch 而非 EventSource,以便发送 POST JSON 和 Authorization header;当前没有自动重连、Last-Event-ID、心跳或主动 AbortController,可靠恢复依赖 complete 后重新读取持久会话,而不是 SSE 帧重放。

后端设置 media_type="text/event-stream"Cache-Control: no-cache。每帧包含 event 行、data 行和空行。data JSON 使用紧凑分隔符,非 ASCII 默认会被 JSON 转义,但前端 JSON.parse 后恢复原字符。事件 id 是随机 UUID 前缀,不与 message ID、tool ID 或 sequence 相同。

前端 decoder 使用 streaming 模式处理 UTF-8 字节被网络 chunk 截断的情况,buffer 保留不完整 frame。HTTP 非成功状态先走普通错误解析;成功但 response.body 为空映射为系统不可用。frame JSON 无效或缺少四个基础字段时,统一转成系统内部错误。parser 忽略 event 行,仅信任 data 内的 type,因此后端必须确保二者一致。

📷 [图片 token=W51GbZCAMoCxPAxNWMQck6Mqn4d(未能下载,见飞书原文)]

浏览器端使用 fetch 读取 ReadableStream,而不是原生 EventSource,因为请求需要 POST JSON 和 Authorization header。当前客户端没有主动取消控制器、心跳帧、自动重连或空闲超时逻辑;这些能力依赖浏览器和上游连接行为。把 SSE 作为协议选择不自动获得可靠重放,持久化 complete 后的会话读取才是现有恢复路径。

Agent 工具集合与事件可见性

**看什么:**看请求级 Runner 如何装配固定工具、选中 Skill 和 owner-scoped MCP 工具;前端不能通过请求体直接指定任意工具名。

        langchain_tool = create_langchain_knowledge_retrieval_tool(
            self._retrieval_tool,
            owner_user_id=request.owner_user_id,
            accessible_knowledge_base_ids=request.accessible_knowledge_base_ids,
        )
        # 1. 知识检索与当前时间始终进入普通聊天工具集合。
        tools = [langchain_tool, create_current_time_tool()]
        if request.skills:
            tools.append(create_load_skill_tool(request.skills))
        mcp_client = self._mcp_client
        if self._mcp_client_provider is not None:
            # 2. MCP client 由当前 owner 的连接服务提供。
            mcp_client = await self._mcp_client_provider.client_for_user(
                owner_user_id=request.owner_user_id
            )
        if mcp_client is not None:
            await mcp_client.discover_tools()
            tools.extend(await mcp_client.get_langchain_tools())
        agent = _create_langchain_agent(
            model=cast(Any, self._llm_provider.create_chat_model()),
            tools=tools,
            system_prompt=(request.system_prompt),
        )

MCP 工具必须真实发现成功才会注册,Skill 工具也只在本次请求有选中 Skills 时出现。所有 LangChain 工具共享同一 started/completed/failed 适配路径;只有输出中实际含 citations 的工具才产生引用,provider 不提供 reasoning 时也不会合成推理事件。

知识检索和当前时间工具始终加入普通聊天 Agent;load_skill 只有选择了 Skill 时加入。MCP 工具来自当前 owner 的 connection service:先创建用户 client,再真实 discover,之后转成 LangChain tools。发现失败不会伪造工具定义。模型看到的工具集合因此由请求身份和服务器配置共同决定,而不是前端提交任意工具名。

📷 [图片 token=XmMZbjOwComlwmxrfefct4mrnBg(未能下载,见飞书原文)]

所有 LangChain 工具共享 on_tool_start、on_tool_end、on_tool_error 适配路径,所以知识、时间、Skill 与 MCP 都使用同一种 tool.call SSE 和审计生命周期。适配器并不按工具名称限制引用:任意工具输出只要包含可解析的 citations 数组,就会额外产生 reference.source;当前典型来源是知识检索。普通时间工具的现有输出不含 citations,因此不会凭空生成引用。load_skill 的完整正文也被压缩成首行 summary 再放入 UI 可见 output,减少敏感指令正文直接展示。

工具 input 和 output 可以通过实时 SSE 到达前端,审计也保存 arguments 与有限摘要。应用日志规则更严格,不记录完整工具参数值或输出。前端组件应以折叠摘要展示过程,不能把任意原始对象直接当 HTML 渲染;聊天 Markdown 组件另有安全渲染测试,引用与工具结果走结构化组件。

📷 [图片 token=BLawb3B4nomdGDxDzMfcjkEOn6d(未能下载,见飞书原文)]

reasoning.delta 的处理同样强调真实来源。适配器支持 OpenAI-compatible chunk 的 additional_kwargs,并从明确字段取值;服务只累积实际返回的文本。推理随 assistant metadata 持久化,可在重开会话后显示,但它与最终答案内容分离。若 provider 不提供推理,界面应没有该区块,而不是使用模板文案假装模型进行过某些思考。

数据、契约与状态

助手消息 metadata 保存 citationsreasoningtoolCallIds。最终正文只由 content delta 拼接,reasoning 不混入回答。工具输出在 SSE 中可以是结构化 object,但审计只保存最多 2000 字符的 JSON 摘要;前端历史读取的工具详情来自单独 GET /chat/sessions/{session_id}/tool-call-audits

前端发送时先乐观追加 user message,再为 assistant 创建 draft ID。每个内容字符之间等待 28 毫秒形成可感知的打字机节奏;reasoning 追加到 draft metadata;tool.call 按 ID 合并状态;reference 单独显示。没有 complete 即使网络正常结束也被当作异常,草稿会被移除并显示统一错误。

权限、安全与失败边界

路由在构造服务前按 owner_user_id 读取 session,跨用户访问返回 403 且 runner.requests 保持为空。知识检索工具与 MCP client 也绑定同一 owner。用户消息在 Agent 执行前写入;如果模型或工具失败,历史会保留这条 user message,帮助用户知道哪次请求失败,但不保存部分 assistant message。

流中异常由 _error_event 转换为统一错误目录中的安全结构。模型已经产生的字符可能先到达浏览器,随后才出现 error;测试明确覆盖这种情形。错误响应不会把 provider 的 sk- 密钥带入 SSE。工具审计的失败摘要额外对常见 sk- 和 AKID 形态做替换,但这不是任意秘密格式的完整检测。

聊天 SSE 没有持久事件序列、重放 API 或 Last-Event-ID 处理。断线后的恢复方式是重新读取已完成的 SQLite 会话;尚未完成且未持久化的 assistant 草稿无法续传。不要把 AIOps durable job 的事件恢复能力外推到普通聊天。

📷 [图片 token=Wq54bXZzNoUSg0xRdgjcn96CnX8(未能下载,见飞书原文)]

阅读顺序与小结

  1. 先读 sse.ts,把 type 当作整个协议的主键。

  2. 再读 _agent_event_from_langchain_event,理解外部框架事件如何被收敛。

  3. 随后跟进 ChatStreamingService.stream_message 的持久化与发流顺序。

  4. 最后阅读 sseClient 与 chat store,确认浏览器中的流式草稿如何回归服务端持久事实。

可靠的流式 Agent 不是简单地把 token 推到浏览器。OncallAgent 用判别式 SSE 建立过程协议,用 owner-scoped SQLite 保存最终事实,用工具审计和引用连接过程与证据,并在 complete 后让前端重新对账。它既提供即时体验,也明确承认普通聊天流尚不具备事件重放能力。