HITL 源码:Eino 的中断与恢复机制(第55篇-E41)
·
系列「企业级 AI Agent 实现拆解」E41 篇,Part 9 源码深度篇第六章。E39 提到工具可以抛出 interrupt,E40 提到 ReAct 的
SetReturnDirectly。这篇完整拆 Eino 的 HITL(Human-In-The-Loop)机制——中断如何产生、如何保存状态、如何恢复。
读完这篇你会知道
- 三种中断方式:配置级节点拦截 vs 运行时组件中断 vs 复合中断
- 中断信号怎么从工具冒泡到 Graph 层
GetInterruptState+GetResumeContext:组件怎么感知自己被恢复Resume/ResumeWithData:外部怎么恢复一个中断的 Agent- 地址系统:如何唯一定位嵌套结构中的断点
- 一个完整的人工审批流实现
三种中断方式
方式一:配置级节点拦截(最粗粒度)
runnable, err := graph.Compile(ctx,
compose.WithInterruptBeforeNodes([]string{"PaymentNode"}), // 执行前暂停
compose.WithInterruptAfterNodes([]string{"KYCNode"}), // 执行后暂停
)
编译时告诉执行引擎:凡是要执行 PaymentNode,先暂停等审批。不需要修改节点代码,适合全局策略。
方式二:运行时组件中断(细粒度,可带状态)
// 无状态中断(只是暂停,不需要恢复时传数据)
return compose.Interrupt(ctx, map[string]any{
"reason": "需要人工确认",
"amount": input.Amount,
})
// 有状态中断(暂停,恢复时组件能读回自己保存的 state)
return compose.StatefulInterrupt(ctx,
map[string]any{"reason": "需要人工确认"}, // info:给外部看的信息
&TransferState{Amount: input.Amount, From: input.From}, // state:恢复时用
)
方式三:复合中断(ToolsNode 使用)
当一批并行工具里只有部分 interrupt:
// ToolsNode 内部
return compose.CompositeInterrupt(ctx,
rerunExtra, // 给外部的信息(包含哪些 tool 成功了,哪些需要重跑)
rerunState, // 内部状态(已执行成功的 tool 结果)
errs..., // 每个 interrupt 的 InterruptSignal
)
// 已成功的 tool 结果存进 rerunState,恢复时不重跑
中断信号的冒泡路径
组件内: return StatefulInterrupt(ctx, info, state)
↓ 产生 *core.InterruptSignal{ID, Address, Info, State}
↓
Graph 执行引擎(graph_run.go):
completedTask.err = InterruptSignal
resolveInterruptCompletedTasks() 收集所有中断任务
↓
handleInterrupt() 把 state + 断点地址 保存进 checkpoint
返回 *interruptError{Info: &InterruptInfo{InterruptContexts: [...]}}
↓
外部调用方:
_, err = agent.Generate(ctx, messages)
info, ok := compose.ExtractInterruptInfo(err)
// info.InterruptContexts[0].InterruptID ← 恢复用的唯一 ID
// info.InterruptContexts[0].Info ← 组件提供的人工阅读信息
地址系统:精确定位断点
每个断点都有一个层级地址,格式是 []AddressSegment{Type, ID}:
AgentRunnable/
graph_node:ChatModel/
tool_call:transfer:call_abc123 ← 第 3 层,精确到具体工具调用
这个地址就是 InterruptID,ResumeWithData 用它来路由恢复信号:
// 中断时 info 里有 InterruptID
info.InterruptContexts[0].InterruptID // "AgentRunnable/ChatModel/transfer:call_abc123"
// 恢复时用这个 ID 定向恢复
ctx = compose.ResumeWithData(ctx, interruptID, approvalResult)
_, err = agent.Generate(ctx, messages) // 重新执行,只有指定断点会被激活
Resume API
// 恢复一个断点(不带数据)
ctx = compose.Resume(ctx, interruptID)
// 恢复一个断点并传数据(人工审批结果)
ctx = compose.ResumeWithData(ctx, interruptID, &ApprovalResult{
Approved: true,
Approver: "alice",
Comment: "金额合规,批准",
})
// 批量恢复多个断点
ctx = compose.BatchResumeWithData(ctx, map[string]any{
interruptID1: &ApprovalResult{Approved: true},
interruptID2: &ApprovalResult{Approved: false},
})
组件内感知恢复
组件被恢复后执行的代码:
func (t *TransferTool) InvokableRun(ctx context.Context, args string, opts ...tool.Option) (string, error) {
// Step 1:检查自己是否曾经中断过
wasInterrupted, hasState, savedState := compose.GetInterruptState[*TransferState](ctx)
if !wasInterrupted {
// 正常执行路径
input := parseArgs(args)
if input.Amount > 10000 {
// 需要审批 → 中断,保存状态
return "", compose.StatefulInterrupt(ctx,
map[string]any{"reason": "超额转账", "amount": input.Amount},
&TransferState{Input: input},
)
}
return doTransfer(input), nil
}
// Step 2:被恢复了。确认自己是被定向恢复,还是"顺带恢复"
isTarget, hasData, approval := compose.GetResumeContext[*ApprovalResult](ctx)
if !isTarget {
// 不是这次 Resume 的目标 → 重新中断,保持自己的状态
return "", compose.StatefulInterrupt(ctx,
map[string]any{"reason": "超额转账,等待其他审批完成"},
savedState,
)
}
if !hasData || !approval.Approved {
return "", fmt.Errorf("转账被拒绝:%s", approval.Comment)
}
// 审批通过 → 执行真正的转账
return doTransfer(savedState.Input), nil
}
两个关键函数:
| 函数 | 作用 |
|---|---|
GetInterruptState[T] |
检查自己是否在上一次中断了;返回 (wasInterrupted, hasState, state) |
GetResumeContext[T] |
检查自己是否是这次 Resume 的目标;返回 (isResumeFlow, hasData, data) |
两种 Resume 策略
策略一:隐式 “Resume All”(适合只有一个审批点)
// 不管哪个断点被中断,Resume 就认为所有断点都该继续
ctx = compose.Resume(ctx, interruptIDs...)
// 组件里只检查 wasInterrupted,不检查 isResumeFlow
策略二:显式 “Targeted Resume”(适合多个独立审批点)
// 只恢复指定的断点,其他断点继续等待
ctx = compose.ResumeWithData(ctx, interruptID1, data1)
// 组件里:isResumeFlow=false 的断点必须重新中断自己,保持等待
if !isTarget {
return "", compose.StatefulInterrupt(ctx, info, savedState)
}
DeepFlux 使用策略二:多个工具可能同时等待不同审批人,每个审批独立进行。
完整流程示意
用户请求 → Agent.Generate(ctx, messages)
↓ 执行到 TransferTool
↓ return StatefulInterrupt(...)
↓ 执行引擎保存 checkpoint
↓ 返回 interruptError
应用层:
info := compose.ExtractInterruptInfo(err)
interruptID := info.InterruptContexts[0].InterruptID
userInfo := info.InterruptContexts[0].Info // {"reason": "超额转账"}
// 异步等待人工审批(可以存 DB、发消息通知)
approval := waitForApproval(interruptID, userInfo)
// 恢复执行
ctx = compose.ResumeWithData(ctx, interruptID, approval)
result, err = agent.Generate(ctx, messages) // 从 checkpoint 恢复,跳过已执行节点
↓
↓ TransferTool 重新执行,GetInterruptState 拿回 savedState
↓ GetResumeContext 确认自己是目标,拿到 approval
↓ 执行真正的转账
↓ 返回结果
小结
Eino 的 HITL 机制三个核心设计:
- 地址系统:每个断点都有层级地址(图→节点→工具调用),支持嵌套结构中精确定位
- StatefulInterrupt:中断时可以把状态存进 checkpoint,恢复时通过
GetInterruptState拿回,组件不需要外部存储 - 定向恢复:
ResumeWithData(interruptID, data)只激活指定断点,其他断点用GetResumeContext检查isResumeFlow,不是目标的重新中断自己
三者配合实现了部分恢复:一个 batch 里 5 个工具,2 个需要审批,另外 3 个已经跑完——Resume 只重跑那 2 个,3 个已完成的结果直接从 checkpoint 读取。
更多推荐


所有评论(0)