runner := adk.NewTypedRunner[M](adk.TypedRunnerConfig[M]{
Agent: agent,
EnableStreaming: true,
CheckPointStore: adkstore.NewInMemoryStore(), // 内存存储
})
// 在 Tool 中触发中断
func myTool(ctx context.Context, args string) (string, error) {
// wasInterrupted: 是否是 Resume 后的第二次调用(第一次为 false,Resume 后为 true)
// storedArgs: 第一次调用时通过 StatefulInterrupt 保存的参数,Resume 后可取回
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
if !wasInterrupted {
// 第一次调用:触发中断,同时保存 args 供 Resume 后使用
return "", tool.StatefulInterrupt(ctx, &ApprovalInfo{
ToolName: "my_tool",
ArgumentsInJSON: args,
}, args) // 第三个参数是要保存的状态(Resume 后通过 storedArgs 取回)
}
// Resume 后的第二次调用:读取用户审批结果
// isTarget: 本次 Resume 是否针对当前 Tool(一次 Resume 只针对一个 Tool)
// hasData: Resume 时是否携带了审批结果数据
// data: 用户传入的审批结果
isTarget, hasData, data := tool.GetResumeContext[*ApprovalResult](ctx)
if isTarget && hasData {
if data.Approved {
return doSomething(storedArgs) // 使用保存的参数执行实际操作
}
return "Operation rejected by user", nil
}
// 其他情况(isTarget=false 意味着本次 Resume 目标不是当前 Tool):重新中断
return "", tool.StatefulInterrupt(ctx, &ApprovalInfo{
ToolName: "my_tool",
ArgumentsInJSON: storedArgs,
}, storedArgs)
}
其中,CheckPointStore: adkstore.NewInMemoryStore()
CheckPointStore 是检查点快照存储接口,专门用于存放 Agent 的 Checkpoint 断点快照;adkstore.NewInMemoryStore () 用于创建内存版快照存储器,底层依托程序内存中的 map 保存全部 checkpoint 数据。当触发 StatefulInterrupt 中断时,框架会调用该存储的 Put 方法,将快照写入内存;执行 Resume 恢复会话时,框架调用 Get 方法从内存读取快照,恢复并继续运行状态机。
// 生产用redis存储,进程重启快照还在
CheckPointStore: adkstore.NewRedisStore(redisClient)
该实现仅将快照保存在当前进程内存中,一旦程序重启或进程退出,所有快照数据都会丢失,仅适合本地调试与 Demo 演示,不可直接用于生产环境;它只是 CheckPointStore 接口的其中一种实现,该接口还有其他实现例如 RedisStore,切换到生产环境只需替换一行代码即可。
type CheckPointStore interface {
// 保存检查点
Put(ctx context.Context, key string, checkpoint *Checkpoint) error
// 获取检查点
Get(ctx context.Context, key string) (*Checkpoint, error)
}
tool.StatefulInterrupt用于抛出中断信号,CheckPointStore.Put由框架内部自动调用,负责将快照持久化存储,二者处于不同层级,是先后配套的协作关系。
时序流程如下:工具代码执行时,如果未处于恢复状态,就调用tool.StatefulInterrupt(ctx, &ApprovalInfo{...}, args)主动抛出中断信号;该方法会在上下文写入中断标记、审批信息以及需要暂存的args,终止当前工具与 Runner 执行流程,向上抛出中断事件给 Eino Runner。业务代码中不需要手动调用 Put。当 Eino Runner 捕获到这个中断事件后,会在框架内部隐式执行checkpointStore.Put(ctx, checkPointID, checkpointObj),把组装完成的*Checkpoint快照,以 checkPointID 作为 Key 存入内存或 Redis 等存储介质。Put 是 CheckPointStore 接口定义的存储层方法,只负责快照数据的写入。
对 Eino Eino 智能体启动的程序来说,它其实一直都没有中断的
对于 Eino 智能体程序,要区分两层执行逻辑,避免混淆:Go 主程序进程本身始终保持运行,不会退出;而 Eino 的 Runner/TurnLoop 本轮 Agent 执行流程会被暂停冻结。当触发 StatefulInterrupt 时,当前runner.Run()执行链路直接终止,TurnLoop 不再继续推进后续节点,框架保存 checkpoint 快照,将本轮 Agent 任务冻结,Runner 把执行权交还给外层业务代码。
在终端示例里,main 协程会卡在控制台输入等待逻辑,进程保持存活,只是原地等待用户输入审批结果;用户确认后调用 ResumeWithParams,重新启动 Runner 加载快照,从暂停的节点继续执行 Agent。
而 Web HTTP 版本则不一样,触发中断后本次 HTTP 请求直接返回、该请求对应的 Go 协程退出,但主进程依旧存活,Agent 任务冻结在 checkpoint 中,待前端发起新请求时,再调用 ResumeWithParams 恢复执行。
其实终端对话中断就是保存和恢复 checkpoint 结构体信息
在 ch09/ch05 终端示例中,触发中断时 Checkpoint 结构体存放于 InMemoryStore 内存存储,并非本地 json 文件;jsonl 文件保存的是 session 会话聊天记录,InMemoryStore 则是程序内部的 map 结构,用来存放 Checkpoint 断点快照,一旦 Ctrl+C 退出 Go 程序,内存内所有 Checkpoint 数据会直接消失。
完整流程为:执行go run ./cmd/ch09启动程序,初始化checkpointStore := adkstore.NewInMemoryStore();Agent 执行并调用工具、触发tool.StatefulInterrupt后,Eino Runner 捕获中断,组装 * Checkpoint 快照并调用 checkpointStore.Put,以 sessionID 为 key 存入内存 map,当前 Runner 流程终止,回到终端代码进入 handleInterrupt 等待输入审批;
当输入 y 调用 ResumeWithParams,内部会通过 checkpointStore.Get 从内存 map 取出 Checkpoint,注入审批结果,重新驱动 Agent 执行并再次进入工具与中间件逻辑,全部执行完毕若无新中断,该 Checkpoint 快照依旧保留在内存,直到程序退出或被新中断覆盖。
任务正常跑完、没有再次中断,这个旧 Checkpoint不会自动删除,一直留在内存 map;只有同 session 再次触发 StatefulInterrupt,新快照 Put 时,会覆盖这个 sessionID 对应的旧 Checkpoint。
Checkpoint 断点和 session 会话文件是一回事吗?
session 会话:存于磁盘data目录 jsonl 文件,持久化保存聊天记录,仅 Agent 完整跑完无中断时才追加消息,中断状态不会写入助手回答。
Checkpoint 断点:默认使用InMemoryStore,底层是程序内存 map,触发tool.StatefulInterrupt中断时立刻保存断点快照,用于 Resume 恢复任务;程序 Ctrl+C 退出后内存断点全部丢失。 如需服务重启后断点不丢失,需将InMemoryStore替换为 Redis 版 CheckPointStore。
Agent 流程中断,不是截断大模型生成
StatefulInterrupt 属于 Eino Agent 流程层面的中断,并非大模型 API 层面的流式截断,大模型本身完全感知不到中断事件。该中断无法在大模型流式输出的中途触发,只会发生在大模型本轮推理完成、返回结果之后,框架准备执行工具的阶段。
一轮 TurnLoop 迭代流程:Eino 将消息历史发给大模型,大模型返回文本或者工具调用指令;如果返回工具调用指令,Eino 准备执行工具前触发 StatefulInterrupt,保存 checkpoint 快照,暂停 Agent 执行流程,等待人工审批。审批通过并 Resume 恢复后,框架执行工具拿到结果,把工具结果追加到消息历史,再次调用大模型进行下一轮推理。
整个过程中,大模型每一次都是接收完整消息批量推理,中间的暂停、人工审批逻辑,对大模型是透明的。
Web 场景下,Eino Agent 触发工具审批中断,推荐的无阻塞方案是怎样的?两次 HTTP 请求之间,服务端是否有协程阻塞等待用户审批操作?
Web 场景下 Eino Agent 工具审批推荐采用无阻塞方案,全程不会有协程长时间阻塞等待用户审批。前端发起聊天 HTTP 请求,服务端执行runner.Run()运行 Agent,当 Agent 触发StatefulInterrupt中断后,框架保存 checkpoint 快照,runner 返回中断信息;当前这条 HTTP 协程直接向前端返回{"code":"need_approve","session_id":"xxx"}的 JSON 响应,请求结束、协程销毁,释放服务器资源。
前端收到提示后弹出审批按钮,用户点击同意或拒绝,会发起全新的POST /api/approve接口请求。/api/approve的新协程拿到 sessionID 与审批结果,调用ResumeWithParams加载 Redis 中存储的 checkpoint,恢复 Agent 继续执行。整个链路是两次相互独立的 HTTP 请求,对应两个完全不同的协程,服务端没有协程卡在原地等待审批;依靠 sessionID 关联两次请求,checkpoint 和会话数据持久化保存在 Redis 中,不需要 channel 传递审批结果。
Web 场景下采用两次 HTTP 请求方案实现工具审批中断:第一次聊天接口触发 Agent 中断并保存 Checkpoint,协程销毁;用户前端审批后发起第二次审批接口请求。后台如何知道 Agent 需要从第二个节点继续执行,而不是从头重新执行节点 1?
依靠存储在 Redis 中的 Checkpoint 快照记录 Agent 运行状态。快照内部保存了 TurnLoop 状态机信息,标记 Agent 上次执行停留在节点 2(准备执行工具、等待审批),同时附带消息历史、迭代信息、中断 ID 以及工具入参 args。
第一次聊天请求 POST /chat:启动 runner.Run 执行 TurnLoop,节点 1 调用大模型并拿到工具调用指令,进入节点 2 准备执行工具时触发StatefulInterrupt。Eino 框架自动组装完整 Checkpoint 快照并存入 Redis,key 为 sessionID。随后 runner 返回中断事件,当前 HTTP 协程销毁,向前端返回等待审批提示与 sessionId。
type CheckPointStore interface {
// 保存检查点
Put(ctx context.Context, key string, checkpoint *Checkpoint) error
// 获取检查点
Get(ctx context.Context, key string) (*Checkpoint, error)
}
用户在前端完成审批,发起第二次请求 POST /approve。新的 HTTP 协程携带 sessionID 和审批结果,调用ResumeWithParams。该方法内部通过 CheckPointStore.Get,根据 sessionID 从 Redis 读取 Checkpoint 快照。框架解析快照,识别上一轮停在节点 2,节点 1 已执行完毕,无需重复调用大模型;接着重建 TurnLoop 运行状态,将审批结果注入中断上下文,直接从节点 2 继续向后执行,读取审批结果、运行工具,再依次执行后续节点。
中断无需阻塞,终端阻塞是因为要原地等用户输入结果
Eino TurnLoop 内层的中断逻辑在终端、Web HTTP、gRPC 场景下完全一致:触发工具审批中断后,框架自动保存 Checkpoint 快照,记录当前停留在当前节点、准备执行工具的运行状态,随后本轮runner.Run执行结束,这部分状态保存能力由 Eino 框架负责,业务代码无需改动。
而等待用户审批属于外层业务逻辑,和 Eino 框架本身解耦,有两种实现方式。终端场景中主线程阻塞在键盘输入等待用户确认,拿到审批结果后直接调用 Resume 恢复执行;Web 场景采用非阻塞方案,触发中断后直接向前端返回响应并释放当前协程,用户点击审批按钮会发起全新 HTTP 请求,新协程凭借 sessionID 加载 Redis 中持久化的 Checkpoint,调用 Resume 继续 Agent 流程。无论外层采用阻塞还是非阻塞模式,Eino 内部状态机、Checkpoint 断点恢复逻辑都不会变化,只要快照保存在 Redis,就算原协程销毁,隔一段时间也能正常恢复执行。
恢复中断时,当前节点重新执行
func toolFunc(ctx, args) {
// 第1行,你自己写的判断
wasInterrupted, _, storedArgs := tool.GetInterruptState[string](ctx)
if !wasInterrupted {
// 首次执行:在这里触发中断,保存args,return中断信号
return tool.StatefulInterrupt(...)
}
// Resume之后,直接走这里,执行工具业务(对应50步之后的逻辑)
runBusiness(storedArgs)
}
外层 TurnLoop 对应的大节点状态由 Eino 框架自动保存,无需我们手动编码处理;而工具函数内部区分首次调用与恢复调用的分支跳转逻辑,则需要自行编写 if 判断实现。一般是将中断逻辑放在工具业务执行之前,在函数入口处判断 wasInterrupted,尽量不要在函数中间代码行里触发中断,否则需要额外维护大量自定义子状态,增加复杂度。