# LangGraph v3 → AG-UI 事件映射 本文描述本仓库如何将 LangGraph **Event Streaming v3** 的 `ProtocolEvent` 投影为 [AG-UI](https://docs.agui.com/) `BaseEvent`,供 CopilotKit 与 SSE 消费。 人在回路(`interrupt` / resume)全链路见 [HITL.md](./HITL.md)。写作编辑器对 CUSTOM / reasoning 的消费见 [文本编辑器.md](./文本编辑器.md)。 ## 目标与前提 - LangGraph `@langchain/langgraph`,`streamEvents(input, { version: 'v3' })` - AG-UI `@ag-ui/core` 事件类型(`EventType.*`) - 投影层只映射**过程事件**;`RUN_STARTED` / 成功·中断·错误收尾由 Server 编排 ## 架构分层 ```mermaid flowchart TB subgraph graph_pkg [packages/graph] Graph["graphs/*.ts 图定义"] AT[AguiTransformer] Mappers["mapTools / mapMessages / mapInterrupt"] Graph --> AT AT --> Mappers end subgraph server [apps/server/agent] SGE[streamGraphAguiEvents.ts] GA[graphAgents.ts] GTA[GraphTransformerAguiAgent] GA --> SGE GA --> GTA end subgraph client [apps/client] CK[CopilotKit /copilotkit] end AT -->|extensions.aguiEvents| SGE GTA --> CK ``` | 层级 | 路径 | 职责 | |------|------|------| | 图定义 | `packages/graph/src/graphs/*.ts` | Node、Edge、`interrupt()`;不含 HTTP/协议 | | 投影层 | `packages/graph/src/stream/` | `AguiTransformer` 将 v3 协议事件转为 AG-UI | | Server 编排 | `apps/server/src/agent/streamGraphAguiEvents.ts` | `RUN_STARTED`、drain 主协议流、读 `aguiEvents`、成功/中断/错误收尾 | | Agent 注册 | `apps/server/src/agent/graphAgents.ts` | `Graphs[name].compile` + `resolveStreamInput`(及可选 configurable / recursionLimit) | | CopilotKit 胶水 | `apps/server/src/agent/graphTransformerAguiAgent.ts` | `AsyncGenerator` → RxJS `Observable` | ## 注册 Transformer compile 时注入 factory;带 checkpointer 的图需传 `thread_id`: ```typescript import { aguiTransformerFactory } from '@agent/graph' const app = Graphs[name].compile({ checkpointer: getCheckpointer(), transformers: [aguiTransformerFactory], }) ``` `aguiTransformerFactory()` 每次返回新的 `AguiTransformer` 实例。Server 侧按图名缓存编译产物(`graphAgents.ts` 的 `aguiCache`)。 单次 run 也可追加 transformers: ```typescript await app.streamEvents(input, { version: 'v3', transformers: [aguiTransformerFactory] }) ``` ## v3 协议 → AG-UI 总览 `AguiTransformer.process(event: ProtocolEvent)` 按 **`event.method`** 分支(载荷在 `event.params.data`),默认 **`return true`**(原始 `tools` / `messages` / `custom` 仍出现在主协议流,便于调试;CopilotKit 应消费 `extensions.aguiEvents`): | `ProtocolEvent.method` | v3 载荷类型 | AG-UI 输出 | 实现 | |------------------------|-------------|------------|------| | `tools` | `ToolsEventData` | `TOOL_CALL_*` | `mapToolsEventDataToAgUi` | | `messages` | `MessagesEventData` | `TEXT_MESSAGE_*` + `REASONING_MESSAGE_*` | `mapMessagesEventDataToAgUi` | | `custom` | 节点 `config.writer({ name, payload })` | `CUSTOM { name, value }`;或 `name==='agui'` 时解包为原生 AG-UI 事件 | `pushCustomEvent` | `tasks` / `values` 不再由 Transformer 消费;中断收尾见 Server `streamGraphAguiEvents`。 `AguiMappedEvent` 联合类型:`AguiToolEvent | AguiTextMessageEvent | AguiReasoningEvent | CustomEvent | StateSnapshotEvent | RunFinishedEvent`。 ## 工具事件映射 来源:`packages/graph/src/stream/mapToolsToAgUi.ts` | `data.event` | 主要字段 | AG-UI 序列 | |--------------|----------|------------| | `tool-started` | `tool_name`, `tool_call_id`, `input` | `TOOL_CALL_START` + `TOOL_CALL_ARGS` | | `tool-finished` | `tool_call_id`, `output` | `TOOL_CALL_END` + `TOOL_CALL_RESULT` | | `tool-output-delta` | `tool_call_id`, `delta` | `TOOL_CALL_CHUNK` | | `tool-error` | `tool_call_id`, `message` | `TOOL_CALL_END` + `TOOL_CALL_RESULT`(带 `error`) | `AguiToolEvent` = `ToolCallStartEvent | ToolCallArgsEvent | ToolCallEndEvent | ToolCallResultEvent | ToolCallChunkEvent`。 **图侧约定(勿破)**:`tools` 节点必须直接挂 `ToolNode` 或 `LazyToolNode`(见 `packages/graph/src/nodes/lazyToolNode.ts`)。禁止在普通节点里再 `toolNode.invoke(...)`——会双发 `tool-started` / `TOOL_CALL_START`,AG-UI verify 失败。异步工具集用 `LazyToolNode`,不要外层 wrap。单测见 `lazyToolNode.test.ts`。 ## 文本与推理消息映射 来源:`packages/graph/src/stream/mapMessagesToAgUi.ts` 映射器同时处理两类内容块,维护一个 `TextMessageMapState`(`activeMessageId` + `activeReasoningMessageId`): ### 文本流(`TEXT_MESSAGE_*`) | `data.event` | AG-UI | |--------------|-------| | `message-start` | `TEXT_MESSAGE_START`(`role: 'assistant'`),记录 `activeMessageId` | | `content-block-delta`(`delta.type === 'text-delta'`) | `TEXT_MESSAGE_CONTENT`(`delta: text`);空 text 不发 | | `message-finish` | `TEXT_MESSAGE_END`(若 reasoning 块未 END,先补发 `REASONING_MESSAGE_END`) | ### 推理流(`REASONING_MESSAGE_*`) 把 DeepSeek 等模型的 `reasoning_content` 暴露为独立思考流,供 CopilotChat 原生渲染、供 writer 编辑器等自定义消费。所有 agent 共用此映射器,一处生效。 | `data.event` | AG-UI | |--------------|-------| | `content-block-start`(`content.type === 'reasoning'`) | `REASONING_MESSAGE_START`,`messageId = ${activeMessageId}:reasoning:${index}`,`role: 'reasoning'` | | `content-block-delta`(`delta.type === 'reasoning-delta'`) | `REASONING_MESSAGE_CONTENT`;缺 `content-block-start` 时按 `index` 隐式开一个思考消息 | | `content-block-finish`(`content.type === 'reasoning'`) | `REASONING_MESSAGE_END`,清空 `activeReasoningMessageId` | | `message-finish` | 若 `activeReasoningMessageId` 仍非空,兜底补发 `REASONING_MESSAGE_END` | `text` 类型的 `content-block-start` 不产生事件(text 块由 `message-start` + `text-delta` 隐式承载)。 > 通道分流:reasoning 事件**只进 `aguiEvents` 主通道**;`TEXT_MESSAGE_*` 额外进 `messageEvents`。文本优先走 `extensions.messageEvents` / `aguiEvents`,不必再单独消费 `stream.messages` 原始通道。 ### 旁路写入(非 LangChain messages 流) 部分节点不走 LC 消息通道,而是经 `config.writer` + `name: 'agui'` 直接推原生 AG-UI 事件: | 辅助 | 路径 | 行为 | |------|------|------| | `writeAguiAssistantText` | `stream/writeAguiAssistantText.ts` | 一次推送完整 `TEXT_MESSAGE_START/CONTENT/END`(如 document 路径的「编写中…」短提示) | | `runChatCompletion(..., mode: 'streamReasoning')` | `nodes/chatCompletion.ts` | 直连 Chat Completions;`reasoning_*` → `REASONING_MESSAGE_*`;正文只累积、**不**发 `TEXT_MESSAGE_*` | 二者均依赖 `AguiTransformer` 的 `name === 'agui'` 解包(见下节)。 ## 自定义事件 在图节点内调用 `config.writer`(v3 流式挂在 `config.writer`,非 `configurable.writer`): ```typescript import { ORDER_TOOL_PROGRESS_EVENT } from '@agent/graph' config.writer?.({ name: ORDER_TOOL_PROGRESS_EVENT, // 'order_tool_progress' payload: { orderId: '10086', step: 'dispatch' }, }) ``` `ToolNode` 内 tool 回调通常拿不到 graph `writer`;自定义进度事件应放在节点函数里(`ORDER_TOOL_PROGRESS_EVENT` 定义在 `packages/graph/src/tools/order.ts`)。 ### 投影规则(`pushCustomEvent`) 1. 从 `event.params.data` 解出 `name`(缺省 `'custom'`)与 `value`(`payload` 字段,缺省整个 record)。 2. **`name === 'agui'` 解包**(常量 `AGUI_WRITER_EVENT = 'agui'`,与 `@agent/claude` 一致):若 `value` 是对象且含 `type` 字段,视为原生 AG-UI 事件,直接推入 `aguiEvents`;并按 `type` 路由到 `toolEvents`(`TOOL_CALL_*`)或 `messageEvents`(`TEXT_MESSAGE_*`)。`claudeAgentGraph` 把 `config.writer` 透传给 `runQueryInGraphNode`,后者以 `{ name: 'agui', payload: }` 回传。 3. 其余 `name` 投影为 `{ type: EventType.CUSTOM, name, value, timestamp }`,同时进 `customEvents` 与 `aguiEvents`。 ### writer 润色修改摘要(`writer_change_summaries`) `writeEdit`(`nodes/writeEdit.ts`)在 document 润色完成后: 1. `computeHunks(baseline, polished)`(diff-match-patch 行级 diff)得到 `Hunk[]`; 2. LLM 结构化输出(`WriterHunkSummariesSchema`)按 hunk 生成 ≤20 字说明; 3. 经 `config.writer({ name: WRITER_CHANGE_SUMMARIES_EVENT, payload })` 发出。 `WRITER_CHANGE_SUMMARIES_EVENT = 'writer_change_summaries'`,schema 在 `packages/proto/src/writer.ts`。payload 形状: ```typescript { changes: WriterChangeSummary[] // hintFrom, originalText, newText, summary polished: string // 改稿全文,供客户端 Apply baseline: string // 生成时原文快照,Apply 时校验 } ``` server/client 共用 `computeHunks` / `hunkKey`,保证 hunk 划分一致。 ## 扩展通道 `AguiTransformer.init()` 创建四个 `StreamChannel`: | 通道 | 元素类型 | 用途 | |------|----------|------| | `extensions.toolEvents` | `AguiToolEvent` | 仅工具 AG-UI 事件;单测/调试 | | `extensions.messageEvents` | `AguiTextMessageEvent` | 仅 `TEXT_MESSAGE_*`(不含 reasoning) | | `extensions.customEvents` | `CustomEvent` | 仅 CUSTOM(不含 `agui` 解包出的原生事件) | | `extensions.aguiEvents` | `AguiMappedEvent` | **主通道**:过程事件合并(tools / messages / custom) | `aguiEvents` 内事件顺序与图执行中 `process` 调用顺序一致。 ## AguiTransformer 生命周期 1. **`init()`** — 创建四个 channel,重置 `#textMessageState`。 2. **`process(event)`** — 按 method 映射并 `push` 到对应 channel。 interrupt / `RUN_FINISHED` **不在** Transformer 内收尾,由 Server `streamGraphAguiEvents` 在 `aguiEvents` 结束后按 `stream.interrupted` 补发。 ## 中断 finalize 事件 来源:`packages/graph/src/stream/mapInterruptToAgUi.ts`(协议与 Client 行为见 [HITL.md](./HITL.md)) `buildInterruptFinalizeEvents(options)` 按顺序产出: 1. **`STATE_SNAPSHOT`**(仅当传入 `snapshot`) 2. **`RUN_FINISHED`** — `outcome: { type: 'interrupt', interrupts: Interrupt[] }`;CopilotKit 据此写入 `agent.pendingInterrupts` `mapInterruptPayloadToAgUi(lg)`: - `id` ← `lg.interruptId` - `reason` ← `payload.reason` → 否则 `type === 'approval'` 时 `INTERRUPT_REASON_CONFIRMATION = 'confirmation'` → 否则 `payload.type` → 否则 `'input_required'` - `message`、`toolCallId`、`responseSchema`(可选透传) - `metadata` ← `{ ...record, payload }` ## Resume 解析 来源:`packages/graph/src/stream/resolveResumeInput.ts` `resolveResumeFromRunAgentInput(input)` 从 AG-UI `RunAgentInput.resume[]` 解析 LangGraph `Command({ resume })` 的值(**不再**兼容 `forwardedProps.command.resume`): | `resume[]` | 返回值 | |------------|--------| | 单条 `resolved` | 该条 `payload` | | 多条 `resolved` | `Object.fromEntries([interruptId, payload])` | | 全部 `cancelled` | `{ approved: false, reason: '用户取消' }` | | 无 / 空 | `undefined` | `reactAgent`、`dev`、`tushare` 的 `resolveStreamInput` 在检测到 resume 时返回 `new Command({ resume })`。 ## Server 消费模式 `apps/server/src/agent/streamGraphAguiEvents.ts` 是 CopilotKit 路径的统一编排: ```typescript yield { type: EventType.RUN_STARTED, threadId, runId, input } const stream = await graph.streamEvents(streamInput, { version: 'v3', configurable: { thread_id: threadId, ...extraConfigurable }, ...(recursionLimit != null ? { recursionLimit } : {}), }) // 必须并行:drain 主协议流驱动 transformer,同时读 aguiEvents let streamError: unknown const protocolDone = (async () => { try { for await (const _ of stream) { /* drain */ } } catch (err) { streamError = err } })() for await (const event of stream.extensions.aguiEvents) { yield event as BaseEvent } await protocolDone if (streamError != null) throw streamError ``` `StreamGraphAguiOptions`: | 字段 | 用途 | |------|------| | `resolveStreamInput` | 必填;消息输入或 `Command({ resume })` | | `resolveConfigurable?` | 合并进 `configurable`(如 `kbId`、`userPrompt`、`editCase`) | | `resolveRecursionLimit?` | LangGraph `recursionLimit`(`reactAgent`:与 Lab `maxSteps` 1:1) | 收尾(`try` 内,drain 完成后): - 尚无 `RUN_FINISHED` 且 `stream.interrupted`:用 `stream.interrupts` + `await stream.output`(作 `snapshot`)调 `buildInterruptFinalizeEvents` - 尚无 `RUN_FINISHED` 且未中断:发 `RUN_FINISHED`,`outcome: { type: 'success' }` - 已有 `RUN_FINISHED`:不再重复发 异常(`catch`): - 打印 `[streamGraphAguiEvents]` 错误日志 - 若尚无 `RUN_FINISHED`,发 `RUN_ERROR`(**不再发 `RUN_FINISHED`**——`RUN_ERROR` 已是终态) - `serializeAgentError`(`serializeError.ts`)把 thrown 值分类为 `{ message, code, name, json }` 挂到事件上;ag-ui schema 为 passthrough,前端 `onError` 的 `context.event` 可拿到。常见 `code`:`GRAPH_STEP_LIMIT`、`KB_INFRA_DOWN`、`UPSTREAM_UNREACHABLE`、`LLM_AUTH_FAILED`、`LLM_RATE_LIMITED`、`UPSTREAM_TIMEOUT`、`AGENT_INTERNAL` `finally`:`ConversationService.touch(userId, threadId)` 刷新会话。 ### 当前注册的图 Agent `GRAPH_AGENT_DEFINITIONS` 与 `Graphs` 键一一对应: | 名称 | 说明 | |------|------| | `claudeAgent` | Claude Agent SDK + checkpoint + AG-UI | | `reactAgent` | 通用 ReAct(可配 prompt + ask_* + kb_search;支持 resume / maxSteps) | | `dev` | 开发演示:澄清分流 → 天气 / 订单 / KB / HITL approval | | `tushare` | A 股分析(Tushare MCP + ask_human;支持 resume) | | `writer` | 中文改写(inline / document) | | `editorChat` | 写作助手对话(Ask / Write) | | `kb` | 知识库 RAG | ## 测试 | 测试文件 | 断言内容 | |----------|----------| | `stream/mapMessagesToAgUi.test.ts` | 文本/reasoning 映射、空 delta、隐式开思考消息、message-finish 兜底 END | | `stream/aguiTransformer.test.ts` | `name=agui` 解包为原生 `TEXT_MESSAGE_*` / `TOOL_CALL_*` 并路由子通道 | | `stream/mapInterruptToAgUi.test.ts` | approval → `Interrupt`;`buildInterruptFinalizeEvents` 产出 snapshot + interrupt 终态 | | `stream/resolveResumeInput.test.ts` | 单/多 resolved、全 cancelled、忽略 legacy forwardedProps | | `stream/writeAguiAssistantText.test.ts` | writer 推送完整 `TEXT_MESSAGE_*` | | `nodes/lazyToolNode.test.ts` | 直接挂 ToolNode / LazyToolNode 只发一次 `TOOL_CALL_START`;外层 invoke 会双发 | | `graphs/dev.test.ts` 等 | 各图 + 对齐 server 的 interrupt finalize 补发 | 单测消费扩展通道示例: ```typescript const stream = await app.streamEvents(input, { version: 'v3', configurable: { thread_id } }) await Promise.all([ Array.fromAsync(stream.extensions.toolEvents), Array.fromAsync(stream.extensions.messageEvents), Array.fromAsync(stream.extensions.customEvents), ]) ``` ## 代码索引 | 模块 | 路径 | |------|------| | Transformer | `packages/graph/src/stream/aguiTransformer.ts` | | 工具映射 | `packages/graph/src/stream/mapToolsToAgUi.ts` | | 文本/推理映射 | `packages/graph/src/stream/mapMessagesToAgUi.ts` | | 旁路文本写入 | `packages/graph/src/stream/writeAguiAssistantText.ts` | | 中断映射 | `packages/graph/src/stream/mapInterruptToAgUi.ts` | | Resume 解析 | `packages/graph/src/stream/resolveResumeInput.ts` | | Chat Completions 旁路 | `packages/graph/src/nodes/chatCompletion.ts` | | 中性中断协议 + writer 域 | `packages/proto/src/index.ts`、`packages/proto/src/writer.ts` | | Server 编排 | `apps/server/src/agent/streamGraphAguiEvents.ts` | | 错误序列化 | `apps/server/src/agent/serializeError.ts` | | Agent 注册 | `apps/server/src/agent/graphAgents.ts` | | CopilotKit 胶水 | `apps/server/src/agent/graphTransformerAguiAgent.ts` | | 图注册表 | `packages/graph/src/index.ts`(`Graphs`) | | 示例图 | `graphs/claudeAgent.ts`、`reactAgent.ts`、`dev.ts`、`kb.ts`、`tushare.ts`、`writer.ts`、`editorChat.ts` |