# 人在回路(HITL) 本仓库如何实现 LangGraph `interrupt()` ↔ AG-UI / CopilotKit 人在回路全链路。 --- ## 心智模型 **LangGraph checkpoint 是中断的持久化真相;AG-UI / CopilotKit 是送到界面的唯一事件通道;客户端 HITL 只认 `agent.pendingInterrupts`。** ```mermaid flowchart TD A["① Graph interrupt() / hitl*()"] --> B[(Checkpoint
持久化真相)] B --> C["② AG-UI RUN_FINISHED
interrupts id + metadata
live run 或 connect 重放"] C --> D["③ pendingInterrupts → InterruptCard"] D --> E["④ runAgent resume
interruptId + payload"] E --> F["⑤ Graph 从 checkpoint 继续"] ``` - 会话 HTTP `threadState.pendingInterrupt` **仅**作服务端 connect 投影输入与 e2e 断言,**不进**客户端 HITL UI / `blockInput`。 - Connected 前不挂 HITL(`CopilotRuntimeReady`);空窗不用 REST 顶卡。 - 本期 UI **单中断**:取 `pendingInterrupts[0]`。 ### HTTP 无状态 + 两阶段 Run HTTP 无法在中途挂起连接等用户点按钮: 1. **Run 1**:执行到 `interrupt()` → 发中断 AG-UI 事件 → **SSE 正常结束**(状态在 checkpointer) 2. **Run 2**:新请求带 `resume[]` → `Command({ resume })` + 同 `thread_id` → 从 checkpoint 继续 --- ## 时序 ### 1. 运行中触发中断 ```mermaid sequenceDiagram autonumber actor U as 用户 participant Chat as CopilotChat / Shell participant Agent as CopilotKit Agent participant Runtime as Copilot Runtime participant Graph as LangGraph participant CP as Checkpoint participant Stream as AG-UI Stream participant Pending as agent.pendingInterrupts participant UI as AgentInterruptUi U->>Chat: 发送消息 Chat->>Agent: runAgent() Agent->>Runtime: POST run Runtime->>Graph: streamEvents() Graph->>Graph: hitl*() / interrupt() Graph->>CP: 持久化 interrupt CP-->>Graph: ok Graph-->>Runtime: interrupted Runtime->>Stream: STATE_SNAPSHOT Runtime->>Stream: RUN_FINISHED(outcome.interrupt) Stream->>Agent: 写入 pendingInterrupts[{ id, metadata }] Pending-->>UI: id + 展示载荷 UI->>Chat: setInterruptElement(card) UI-->>Chat: blockInput = pendingInterrupts.length > 0 Chat-->>U: 中断卡 ``` live 收尾:`streamGraphAguiEvents` 读 `stream.interrupted`,经 `buildInterruptFinalizeEvents` / `mapInterruptPayloadToAgUi` 发出标准 AG-UI `Interrupt`。UI **不**走 `useInterrupt` / `CUSTOM on_interrupt`。 ### 2. 刷新后恢复 ```mermaid sequenceDiagram autonumber actor U as 用户 participant Page as React 页面 participant Ready as CopilotRuntimeReady participant Runner as CheckpointConnectRunner participant CP as LangGraph Checkpoint participant Agent as CopilotKit Agent participant Pending as agent.pendingInterrupts participant UI as AgentInterruptUi U->>Page: 刷新 Page->>Ready: 等待 Connected Note over Ready: Connected 前不挂 HITL Ready->>Runner: connect(threadId) alt 内存中有 replay Runner-->>Agent: 重放进程内 AG-UI 事件 else 内存为空 Runner->>CP: hydrateThreadBundle CP-->>Runner: messages + pendingInterrupt Runner->>Agent: RUN_STARTED + MESSAGES_SNAPSHOT alt 有 pending Runner->>Agent: RUN_FINISHED(interrupt) Agent->>Pending: 重建 pendingInterrupts Pending-->>UI: 中断卡 + blockInput else 无 pending Runner->>Agent: RUN_FINISHED(success) end end ``` ```mermaid flowchart LR CP[(LangGraph Checkpoint)] CP --> Connect[CheckpointConnectRunner] Connect --> Events[AG-UI 事件] Events --> Agent[agent.pendingInterrupts] Agent --> Card[InterruptCard + blockInput] ``` 客户端视角:刷新后的中断就是 AG-UI 事件(connect 投影),不是 REST JSON。服务端 hydrate 仍用 `threadState.pendingInterrupt` 作投影**输入**。 ### 3. 用户 resume ```mermaid sequenceDiagram autonumber actor U as 用户 participant Card as InterruptCard participant UI as AgentInterruptUi participant Resume as useAgentInterruptResume participant Agent as CopilotKit Agent participant Runtime as Copilot Runtime participant Resolver as resolveResumeInput participant Graph as LangGraph participant CP as Checkpoint participant Stream as AG-UI Stream U->>Card: 批准/拒绝/填写 Card->>UI: onRespond(payload) UI->>Resume: resume(payload, interruptId) Resume->>Agent: runAgent({ resume: [{ interruptId, status: resolved, payload }] }) Agent->>Runtime: POST resume[] Runtime->>Resolver: resolveResumeFromRunAgentInput Resolver->>Graph: Command({ resume }) Graph->>CP: 载入挂起点 Graph->>Graph: interrupt() 返回 payload 并继续 alt 再次 interrupt Graph->>Stream: RUN_FINISHED(interrupt) Stream-->>UI: 下一张卡 else 完成 Graph->>Stream: RUN_FINISHED(success) Agent-->>UI: pendingInterrupts 空,blockInput=false end ``` - resume **必须**带 `interruptId`;不用 CopilotKit `resolve()` / `forwardedProps.command.resume`。 - resume 后 HITL **不** `reloadActiveThread`。 ### 4. 生命周期 ```mermaid stateDiagram-v2 [*] --> Connecting: 进入会话 Connecting --> Idle: Connected 且无 pendingInterrupts Connecting --> PendingUI: Connected 且 connect 重放出 pending Idle --> Running: 用户发送 / runAgent Running --> PendingUI: RUN_FINISHED(interrupt) Running --> Idle: RUN_FINISHED(success) Running --> Failed: RUN_ERROR PendingUI --> Responding: 用户提交 resume[] Responding --> Running: Graph 从 checkpoint 继续 Failed --> [*] Idle --> [*]: 离开会话 ``` --- ## 中性协议:`@agent/proto` `packages/proto` 定义 6 种交互形态(不绑定 LangGraph / CLI): | type | 请求要点 | resume payload | |------|----------|----------------| | `input` | `message`, `placeholder?` | `{ value: string }` | | `select` | `message`, `options` | `{ value: string }` | | `multiSelect` | `message`, `options` | `{ values: string[] }` | | `modal` | `title`, `body`, `actions` | `{ action: string }` | | `approval` | `message`, `details` | `{ approved: boolean, reason? }` | | `unlock` | `message`, `key` | `{}` | - `interruptId` 是恢复路由唯一依据。 - 图内 `interrupt(value)` 的 value **不含** id(由 LangGraph 生成);hydrate 时拼成 `PendingInterrupt`。 - `ThreadState = { pendingInterrupt: PendingInterrupt | null }`;真相在 Postgres checkpoint。 --- ## 图内核 ### hitl* helpers `packages/graph/src/tools/hitl/interrupt.ts`:`hitlInput` / `hitlSelect` / `hitlMultiSelect` / `hitlModal` / `hitlApproval` — 薄封装 `interrupt()`,payload 与协议同构(无 id)。 ### ask_* 工具 `packages/graph/src/tools/ask-tools.ts`:AI 可调用;内部走 `hitl*`,resume 后返回字符串 → ToolMessage。 | 工具 | type | 返回示例 | |------|------|----------| | `ask_input` | `input` | `用户回答:…` | | `ask_choice` | `select` | `用户选择:…` | | `ask_multi_choice` | `multiSelect` | `用户选择:a, b` | | `ask_confirm` | `modal` | `用户选择:…` | 导出 `ASK_TOOLS`;system 片段为 `@agent/proto` 的 `ASK_TOOLS_SYSTEM_PROMPT`。 ### 演示与业务图 | 图 | 路径 | HITL 用法 | |----|------|-----------| | `dev` | `graphs/dev.ts` | `collectHitlDemo`:input → select → multiSelect → approval | | `tushare` | `graphs/tushare.ts` | `resolve_stock` 缺名/多匹配 + `ASK_TOOLS` | | `reactAgent` | `graphs/reactAgent.ts` | 绑定 `ASK_TOOLS` + `kb_search` | 编译:**必须** `checkpointer: getCheckpointer()`(`apps/server/src/db/checkpointer.ts`)。 --- ## 投影与 Server ### 捕获与收尾 `apps/server/src/agent/streamGraphAguiEvents.ts`:流结束若无 `RUN_FINISHED` 且 `stream.interrupted`,调 `buildInterruptFinalizeEvents` 补发 `STATE_SNAPSHOT` + `RUN_FINISHED(interrupt)`。 `AguiTransformer` 只映射过程事件,不负责 interrupt 收尾。 ### mapInterruptToAgUi `packages/graph/src/stream/mapInterruptToAgUi.ts`: - `reason`:`payload.reason` → approval → `'confirmation'` → 其他 `type` → `'input_required'` - `metadata` 带完整 payload;`id` = LangGraph interrupt id ### Resume `resolveResumeFromRunAgentInput`(`packages/graph/src/stream/resolveResumeInput.ts`)**只认** `resume[]`: - 单条 resolved → 其 `payload` - 多条 → `{ [interruptId]: payload }` - 全部 cancelled → `{ approved: false, reason: '用户取消' }` `graphAgents.ts` 中带 HITL 的图:有 resume 则 `new Command({ resume })`,否则普通输入。 ### hydrate + connect - `extractPendingInterruptFromSnapshot`:`tasks[].interrupts[]` → `PendingInterrupt` - `hydrateThreadBundle` → `{ messages, threadState }` - `GET /conversations/messages` 返回 `threadState.pendingInterrupt`(e2e / 调试) - `CheckpointConnectRunner`:空内存时 `RUN_STARTED → MESSAGES_SNAPSHOT → RUN_FINISHED(interrupt|success)` --- ## Client | 模块 | 职责 | |------|------| | `AgentInterruptUi` | `pendingInterrupts[0]` → `narrowAgUiPendingInterrupt` → `InterruptCard` → `setInterruptElement` | | `InterruptCards` | 按 `type` 分发 6 种卡片 | | `useAgentInterruptResume` | `runAgent({ resume: [{ interruptId, status: 'resolved', payload }] })` | | `useAgentHasPendingInterrupt` | `pendingInterrupts.length > 0` → `blockInput` | | `interruptContracts` | AG-UI `{ id, metadata }` → `InterruptRequest` | | `CopilotRuntimeReady` | Connected 前不挂 Chat / HITL | 挂载:主 Chat(`routes/index.tsx`)、Agent Lab(`AgentLabPage.tsx`)。 --- ## 中断的 DB 结构 **没有单独 HITL 表。** 挂起中断在 `PostgresSaver` 表内(与 better-auth、`conversation_threads` 同库)。 | 表 | 与中断 | |--|--| | `checkpoints` | 图状态快照;`thread_id` = 会话 id | | `checkpoint_blobs` | channel 大值(如 messages) | | `checkpoint_writes` | **`channel=__interrupt__` / `__resume__`** | | `conversation_threads` | 仅会话元数据,不存 interrupt | `checkpoint_writes` 挂起 blob 示例: ```json { "id": "80fcfb93…", "value": { "type": "select", "message": "请问你离职的场景是?", "options": [{ "label": "…", "value": "corporate" }] } } ``` ```text checkpoint_writes (__interrupt__) → getState().tasks[].interrupts[] { id, value } → PendingInterrupt → connect:RUN_FINISHED(interrupt) / REST:threadState.pendingInterrupt ``` 排查: ```sql SELECT thread_id, checkpoint_id, task_id, idx, convert_from(blob, 'UTF8') AS interrupt_json FROM checkpoint_writes WHERE channel = '__interrupt__' ORDER BY checkpoint_id DESC LIMIT 20; ``` --- ## 设计要点 | 要点 | 实现 | |------|------| | 中断后释放连接 | `RUN_FINISHED` 后 SSE 结束;状态在 checkpointer | | 客户端单投影 | 只信 `agent.pendingInterrupts` | | resume 靠 id | `resume[].interruptId` 必填 | | 刷新不丢中断 | checkpoint → hydrate → connect 重放同套 AG-UI 事件 | | 扩展新工作流 | 新图 + `GRAPH_AGENT_DEFINITIONS`;复用 `hitl*` / `ASK_TOOLS` | ---