Skip to content
kscarrot edited this page Jul 31, 2026 · 1 revision

人在回路(HITL)

本仓库如何实现 LangGraph interrupt() ↔ AG-UI / CopilotKit 人在回路全链路。


心智模型

LangGraph checkpoint 是中断的持久化真相;AG-UI / CopilotKit 是送到界面的唯一事件通道;客户端 HITL 只认 agent.pendingInterrupts。

flowchart TD
  A["① Graph interrupt() / hitl*()"] --> B[(Checkpoint<br/>持久化真相)]
  B --> C["② AG-UI RUN_FINISHED<br/>interrupts id + metadata<br/>live run 或 connect 重放"]
  C --> D["③ pendingInterrupts → InterruptCard"]
  D --> E["④ runAgent resume<br/>interruptId + payload"]
  E --> F["⑤ Graph 从 checkpoint 继续"]
Loading
  • 会话 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. 运行中触发中断

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: 中断卡
Loading

live 收尾:streamGraphAguiEvents 读 stream.interrupted,经 buildInterruptFinalizeEvents / mapInterruptPayloadToAgUi 发出标准 AG-UI Interrupt。UI 不走 useInterrupt / CUSTOM on_interrupt。

2. 刷新后恢复

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
Loading
flowchart LR
    CP[(LangGraph Checkpoint)]
    CP --> Connect[CheckpointConnectRunner]
    Connect --> Events[AG-UI 事件]
    Events --> Agent[agent.pendingInterrupts]
    Agent --> Card[InterruptCard + blockInput]
Loading

客户端视角:刷新后的中断就是 AG-UI 事件(connect 投影),不是 REST JSON。服务端 hydrate 仍用 threadState.pendingInterrupt 作投影输入。

3. 用户 resume

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
Loading
  • resume 必须带 interruptId;不用 CopilotKit resolve() / forwardedProps.command.resume。
  • resume 后 HITL 不 reloadActiveThread。

4. 生命周期

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 --> [*]: 离开会话
Loading

中性协议:@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 示例:

{
  "id": "80fcfb93…",
  "value": {
    "type": "select",
    "message": "请问你离职的场景是?",
    "options": [{ "label": "…", "value": "corporate" }]
  }
}
checkpoint_writes (__interrupt__)
  → getState().tasks[].interrupts[] { id, value }
  → PendingInterrupt
  → connect:RUN_FINISHED(interrupt)  /  REST:threadState.pendingInterrupt

排查:

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

Clone this wiki locally