-
Notifications
You must be signed in to change notification settings - Fork 0
AGUI
本文描述本仓库如何将 LangGraph Event Streaming v3 的 ProtocolEvent 投影为 AG-UI BaseEvent,供 CopilotKit 与 SSE 消费。
人在回路(interrupt / resume)全链路见 HITL.md。写作编辑器对 CUSTOM / reasoning 的消费见 文本编辑器.md。
- LangGraph
@langchain/langgraph,streamEvents(input, { version: 'v3' }) - AG-UI
@ag-ui/core事件类型(EventType.*) - 投影层只映射过程事件;
RUN_STARTED/ 成功·中断·错误收尾由 Server 编排
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<BaseEvent> → RxJS Observable
|
compile 时注入 factory;带 checkpointer 的图需传 thread_id:
import { aguiTransformerFactory } from '@agent/graph'
const app = Graphs[name].compile({
checkpointer: getCheckpointer(),
transformers: [aguiTransformerFactory],
})aguiTransformerFactory() 每次返回新的 AguiTransformer 实例。Server 侧按图名缓存编译产物(graphAgents.ts 的 aguiCache)。
单次 run 也可追加 transformers:
await app.streamEvents(input, { version: 'v3', transformers: [aguiTransformerFactory] })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):
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) |
把 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原始通道。
部分节点不走 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):
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)。
- 从
event.params.data解出name(缺省'custom')与value(payload字段,缺省整个 record)。 -
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: <AG-UI event> }回传。 - 其余
name投影为{ type: EventType.CUSTOM, name, value, timestamp },同时进customEvents与aguiEvents。
writeEdit(nodes/writeEdit.ts)在 document 润色完成后:
-
computeHunks(baseline, polished)(diff-match-patch 行级 diff)得到Hunk[]; - LLM 结构化输出(
WriterHunkSummariesSchema)按 hunk 生成 ≤20 字说明; - 经
config.writer({ name: WRITER_CHANGE_SUMMARIES_EVENT, payload })发出。
WRITER_CHANGE_SUMMARIES_EVENT = 'writer_change_summaries',schema 在 packages/proto/src/writer.ts。payload 形状:
{
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 调用顺序一致。
-
init()— 创建四个 channel,重置#textMessageState。 -
process(event)— 按 method 映射并push到对应 channel。
interrupt / RUN_FINISHED 不在 Transformer 内收尾,由 Server streamGraphAguiEvents 在 aguiEvents 结束后按 stream.interrupted 补发。
来源:packages/graph/src/stream/mapInterruptToAgUi.ts(协议与 Client 行为见 HITL.md)
buildInterruptFinalizeEvents(options) 按顺序产出:
-
STATE_SNAPSHOT(仅当传入snapshot) -
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 }
来源: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 })。
apps/server/src/agent/streamGraphAguiEvents.ts 是 CopilotKit 路径的统一编排:
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 streamErrorStreamGraphAguiOptions:
| 字段 | 用途 |
|---|---|
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) 刷新会话。
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 补发 |
单测消费扩展通道示例:
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
|