系列「企业级 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 层,精确到具体工具调用

这个地址就是 InterruptIDResumeWithData 用它来路由恢复信号:

// 中断时 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 机制三个核心设计:

  1. 地址系统:每个断点都有层级地址(图→节点→工具调用),支持嵌套结构中精确定位
  2. StatefulInterrupt:中断时可以把状态存进 checkpoint,恢复时通过 GetInterruptState 拿回,组件不需要外部存储
  3. 定向恢复ResumeWithData(interruptID, data) 只激活指定断点,其他断点用 GetResumeContext 检查 isResumeFlow,不是目标的重新中断自己

三者配合实现了部分恢复:一个 batch 里 5 个工具,2 个需要审批,另外 3 个已经跑完——Resume 只重跑那 2 个,3 个已完成的结果直接从 checkpoint 读取。


代码来源:eino/compose/interrupt.go · eino/compose/resume.go

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐