夜雨聆风学习资料网

ARTICLE · 1097074

从 Eino 源码学习一个基于图编排的 ReAct Agent

从 Eino 源码学习一个基于图编排的 ReAct Agent

1. 什么是 ReAct

ReAct 是 Reasoning(推理)与 Acting(行动)的组合。它让模型不必只依靠当前上下文一次性给出答案,而是可以在需要时调用搜索、计算、数据库查询等外部工具,再根据工具结果继续判断。

因此,ReAct 更像一种 Agent 的运行方式,而不是某个框架专有的功能:模型负责决定下一步,工具负责执行,执行结果再回到模型上下文中。

2. ReAct 的核心原理

ReAct 的名字来源于两个关键部分:

  • Reason(思考): LLM 会生成思维链(Chain-of-Thought),这是一种详细的推理过程,包括它对问题的理解、如何拆解问题、以及接下来打算做什么。这就像一个人在心里默念自己的思考过程。

  • Act(行动): LLM 会根据它的思考,调用外部工具或执行特定操作。这些“工具”可以是:

    • 搜索引擎:查找最新信息或特定事实。
    • 计算器:执行数学运算。
    • 代码解释器:运行代码来处理数据或验证逻辑。
    • API 接口:与外部系统交互,比如查询天气、发送邮件等。

它的的核心是:把一次生成答案,改成多轮决策。

每一轮中,模型只需要作出下面两种选择之一:

  • 信息足够,输出最终答案;
  • 信息不足,生成一个或多个工具调用。

工具执行后的结果会作为新的消息加入上下文。模型看到最新结果后再次决策,直到不再调用工具。为了避免异常情况下无限循环,工程实现通常还会设置最大执行步数。

3. ReAct 的工作流程

以“查询某地天气并给出出行建议”为例,模型第一次拿到问题时并不知道实时天气,于是生成天气工具调用;工具返回天气数据后,模型再结合用户问题组织最终答案。

对应到常见的 ReAct 术语,就是:接受任务、思考、行动、观察、循环、输出答案。这里的“思考”更准确地说是模型对下一步动作的判断;“观察”则是工具返回并重新进入上下文的结果。

4. Eino 源码解读

下面只讨论官方 react.go 中的实现(完整源码请去链接中阅读)。先看完整逻辑结构,再沿着一次请求的运行路径阅读 NewAgent()。

4.1 Eino 实现概述

Eino 把 ReAct Agent 编排成一张有环图。除去起止节点,图中有三个实际节点:

  • chat:调用模型,输入是消息列表,输出是一条 assistant 消息。
  • tools:读取 assistant 消息中的 ToolCalls,执行对应工具,输出一组 tool 消息。
  • direct_return:从工具结果中选出指定的一条,直接作为 Agent 结果返回。

这张图的主循环是 chat -> tools -> chat。模型输出中有工具调用时进入 tools;工具结果执行完后又回到 chat;模型不再输出工具调用时,图走向 END;如果在配置中指定了map[string]struct{} 类型的参数 ToolReturnDirectly,会在模型调用指定工具后直接返回结果,走向 END。

图的每次运行都会产生一个独立的运行时状态:

go
type state struct {    Messages                 []*schema.Message    ReturnDirectlyToolCallID string}

Messages 保存下一轮模型能够看到的消息历史;ReturnDirectlyToolCallID 记录哪一次工具调用的结果需要直接返回。这里存的是调用 ID,而不是工具名,因为同一轮中可能多次调用同一个工具,只有 ID 能定位到具体结果。

4.2 Agent 和 NewAgent()

先看 Agent 结构体,它本质上是一个已经编译好的图、编译前的图和相关配置项:

go
type Agent struct {    runnable         compose.Runnable[[]*schema.Message, *schema.Message]    graph            *compose.Graph[[]*schema.Message, *schema.Message]    graphAddNodeOpts []compose.GraphAddNodeOpt}
  • runnable 是编译完成、真正负责执行的图;
  • graph 保留原始图,供 ExportGraph() 把当前 Agent 继续嵌入其他图;
  • graphAddNodeOpts 保存把它作为子图加入其他图时需要沿用的编译配置。

我们重点关注 graph 即可,其他都是些框架特性

这里贴出 NewAgent() 的源码:

go
// NewAgent creates a ReAct agent that feeds tool response into next round of Chat Model generation.//// IMPORTANT!! For models that don't output tool calls in the first streaming chunk (e.g. Claude)// the default StreamToolCallChecker may not work properly since it only checks the first chunk for tool calls.// In such cases, you need to implement a custom StreamToolCallChecker that can properly detect tool calls.funcNewAgent(ctx context.Context, config *AgentConfig) (_ *Agent, err error) {var (    chatModel       model.BaseChatModel    toolsNode       *compose.ToolsNode    toolInfos       []*schema.ToolInfo    toolCallChecker = config.StreamToolCallChecker    messageModifier = config.MessageModifier    )    graphName := GraphNameif config.GraphName != "" {    graphName = config.GraphName    }    modelNodeName := ModelNodeNameif config.ModelNodeName != "" {    modelNodeName = config.ModelNodeName    }    toolsNodeName := ToolsNodeNameif config.ToolsNodeName != "" {    toolsNodeName = config.ToolsNodeName    }if toolCallChecker == nil {    toolCallChecker = firstChunkStreamToolCallChecker    }if toolInfos, err = genToolInfos(ctx, config.ToolsConfig); err != nil {returnnil, err    }if chatModel, err = agent.ChatModelWithTools(config.Model, config.ToolCallingModel, toolInfos); err != nil {returnnil, err    }if toolsNode, err = compose.NewToolNode(ctx, &config.ToolsConfig); err != nil {returnnil, err    }    graph := compose.NewGraph[[]*schema.Message, *schema.Message](compose.WithGenLocalState(func(ctx context.Context) *state {return &state{Messages: make([]*schema.Message, 0, config.MaxStep+1)}    }))    modelPreHandle := func(ctx context.Context, input []*schema.Message, state *state) ([]*schema.Message, error) {    state.Messages = append(state.Messages, input...)if config.MessageRewriter != nil {    state.Messages = config.MessageRewriter(ctx, state.Messages)    }if messageModifier == nil {return state.Messages, nil    }    modifiedInput := make([]*schema.Message, len(state.Messages))copy(modifiedInput, state.Messages)return messageModifier(ctx, modifiedInput), nil    }if err = graph.AddChatModelNode(nodeKeyModel, chatModel, compose.WithStatePreHandler(modelPreHandle), compose.WithNodeName(modelNodeName)); err != nil {returnnil, err    }if err = graph.AddEdge(compose.START, nodeKeyModel); err != nil {returnnil, err    }    toolsNodePreHandle := func(ctx context.Context, input *schema.Message, state *state) (*schema.Message, error) {if input == nil {return state.Messages[len(state.Messages)-1], nil// used for rerun interrupt resume    }    state.Messages = append(state.Messages, input)    state.ReturnDirectlyToolCallID = getReturnDirectlyToolCallID(input, config.ToolReturnDirectly)return input, nil    }if err = graph.AddToolsNode(nodeKeyTools, toolsNode, compose.WithStatePreHandler(toolsNodePreHandle), compose.WithNodeName(toolsNodeName)); err != nil {returnnil, err    }    modelPostBranchCondition := func(ctx context.Context, sr *schema.StreamReader[*schema.Message]) (endNode string, err error) {if isToolCall, err := toolCallChecker(ctx, sr); err != nil {return"", err    } elseif isToolCall {return nodeKeyTools, nil    }return compose.END, nil    }if err = graph.AddBranch(nodeKeyModel, compose.NewStreamGraphBranch(modelPostBranchCondition, map[string]bool{nodeKeyTools: true, compose.END: true})); err != nil {returnnil, err    }if err = buildReturnDirectly(graph); err != nil {returnnil, err    }    compileOpts := []compose.GraphCompileOption{compose.WithMaxRunSteps(config.MaxStep), compose.WithNodeTriggerMode(compose.AnyPredecessor), compose.WithGraphName(graphName)}    runnable, err := graph.Compile(ctx, compileOpts...)if err != nil {returnnil, err    }return &Agent{    runnable:         runnable,    graph:            graph,    graphAddNodeOpts: []compose.GraphAddNodeOpt{compose.WithGraphCompileOptions(compileOpts...)},    }, nil}

第一步:准备模型和工具节点

NewAgent() 先读取所有工具的描述:

go
toolInfos, err = genToolInfos(ctx, config.ToolsConfig)

genToolInfos() 会依次调用每个工具的 Info()。得到的 ToolInfo 包含工具名、说明和参数结构,这些信息随后会绑定给模型:

go
chatModel, err = agent.ChatModelWithTools(    config.Model,    config.ToolCallingModel,    toolInfos,)

与此同时,真正的工具实现被交给 ToolsNode:

go
toolsNode, err = compose.NewToolNode(ctx, &config.ToolsConfig)

这里有一个很重要的分工:模型拿到的是工具说明,ToolsNode 拿到的是可执行的工具。 模型只负责生成“调用哪个工具、参数是什么”,并不会亲自执行函数。

如果配置了推荐使用的 ToolCallingModel,ChatModelWithTools() 会调用 WithTools() 得到绑定工具后的模型实例;旧的 Model 字段则通过 BindTools() 修改模型。前者不会改动原模型实例,更适合并发场景,因此 Eino 已经把后者标记为废弃。如果两个字段同时提供,会优先使用 ToolCallingModel。

第二步:创建 Graph 和本地状态

图的输入是消息切片,输出是一条消息:

go
graph := compose.NewGraph[[]*schema.Message, *schema.Message](    compose.WithGenLocalState(func(ctx context.Context) *state {return &state{            Messages: make([]*schema.Message, 0, config.MaxStep+1),        }    }),)

在 NewGraph() 时通过 WithGenLocalState() 传入 State 的创建方法。这个请求维度的全局状态在一次请求的各环节可读写使用。

你可以把 State 看作 Graph 的一组全局变量,它自带互斥锁,随 Agent Runtime 产生和消亡。

第三步:添加模型节点

模型节点的前置处理器 modelPreHandle,负责处理输入和更新 state:

go
modelPreHandle := func(    ctx context.Context,    input []*schema.Message,    state *state,) ([]*schema.Message, error) {    state.Messages = append(state.Messages, input...)if config.MessageRewriter != nil {        state.Messages = config.MessageRewriter(ctx, state.Messages)    }if messageModifier == nil {return state.Messages, nil    }    modifiedInput := make([]*schema.Message, len(state.Messages))copy(modifiedInput, state.Messages)return messageModifier(ctx, modifiedInput), nil}

第一次进入模型节点时,input 是调用者传入的用户消息;工具执行完成后再次进入时,input 则是本轮产生的 tool 消息。它们都会先被追加到 state.Messages,再一起交给模型。

配置项中,两个 Message 处理相关的参数类型相同(都是 MessageModifier)但职责是不同的:

  • MessageRewriter 的返回值会写回 state.Messages,影响后面的所有轮次,适合裁剪或压缩历史消息。
  • MessageModifier 接收复制后的切片,只影响当前这次模型调用,适合临时补充 system message 等内容。

指定messageModifier时复制的只是切片,切片中的 *schema.Message 仍指向原来的消息对象。因此可以把它理解成“隔离对切片结构的修改”,而不是深拷贝整个消息对象。即最终模型节点中嵌套的 LLM 看到的是messageModifier修改过的文本。

随后,把模型节点加入图,并连接入口:

go
graph.AddChatModelNode(    nodeKeyModel,    chatModel,    compose.WithStatePreHandler(modelPreHandle),    compose.WithNodeName(modelNodeName),)graph.AddEdge(compose.START, nodeKeyModel)

第四步:添加工具节点

工具节点也有一个前置处理器:

go
toolsNodePreHandle := func(    ctx context.Context,    input *schema.Message,    state *state,) (*schema.Message, error) {if input == nil {return state.Messages[len(state.Messages)-1], nil    }    state.Messages = append(state.Messages, input)    state.ReturnDirectlyToolCallID = getReturnDirectlyToolCallID(        input,        config.ToolReturnDirectly,    )return input, nil}

这里追加的是带有 ToolCalls 的 assistant 消息。这样,当工具结果回到模型节点时,历史会形成完整的三段:

user 消息 | assistant 的工具调用消息 | tool 的执行结果消息

这是 Tool Calling 能正确进行多轮交互的基础。模型既能看到自己刚才发起了什么调用,也能看到对应的执行结果。

input == nil 的分支用于工具中断后重新运行、恢复执行的场景。正常路径中,处理器还会检查当前 assistant 消息里是否调用了需要直接返回的工具,并保存对应的 ToolCall ID。

第五步:添加条件分支

模型节点之后的分支决定是调用工具还是结束:

go
modelPostBranchCondition := func(    ctx context.Context,    sr *schema.StreamReader[*schema.Message],) (string, error) {    isToolCall, err := toolCallChecker(ctx, sr)if err != nil {return"", err    }if isToolCall {return nodeKeyTools, nil    }return compose.END, nil}

有 ToolCalls 就走向 tools,没有就走向 END。ReAct 的一次“是否继续思考”,对应的就是该分支选择。

没有配置StreamToolCallChecker时会使用默认的 firstChunkStreamToolCallChecker 跳过开头的空 chunk,并在遇到工具调用或第一段非空文本时立即作出判断。这对一开始就输出 ToolCall 的模型没有问题;但有些模型会先输出文本、后输出 ToolCall(Claude 就是这样的),此时默认检查器可能提前判断为“没有工具调用”。所以实际业务中,应该根据使用的 LLM 来判断是否需要自定义StreamToolCallChecker,并要求自定义检查器在返回前关闭收到的流。

工具节点之后还有第二处分支:普通结果回到 chat,需要直返的结果进入 direct_return。这部分由 buildReturnDirectly() 完成,后面小节再细嗦。

第六步:编译图

节点和边都添加完后,NewAgent() 编译这张图:

go
compileOpts := []compose.GraphCompileOption{    compose.WithMaxRunSteps(config.MaxStep),    compose.WithNodeTriggerMode(compose.AnyPredecessor),    compose.WithGraphName(graphName),}runnable, err := graph.Compile(ctx, compileOpts...)

chat 既可以由 START 触发,也可以由 tools 触发。AnyPredecessor 表示任一前驱在上一个执行步完成,当前节点就可以运行,正适合这种带回边的图。

MaxStep 用来给循环兜底,防止模型不断调用工具而无法结束。要注意,它限制的是图运行步数,不是“最多思考几轮”或“最多调用几次工具”。一次 chat -> tools -> chat 会经过多个图节点,二者不能直接画等号。当 MaxStep 为零时,具体默认值由 compose 层按图结构计算,不必在业务代码中硬编码一个轮数。

最后,编译后的 runnable、原始 graph 和编译选项一起被装进 Agent 返回。

4.3 AgentConfig

AgentConfig 的字段很多,但按职责来讨论并就不复杂了。先看源码:

go
// AgentConfig is the config for ReAct agent.type AgentConfig struct {// ToolCallingModel is the chat model to be used for handling user messages with tool calling capability.// This is the recommended model field to use.    ToolCallingModel model.ToolCallingChatModel// Deprecated: Use ToolCallingModel instead.    Model model.ChatModel// ToolsConfig is the config for tools node.    ToolsConfig compose.ToolsNodeConfig// MessageModifier.// modify the input messages before the model is called, it's useful when you want to add some system prompt or other messages.    MessageModifier MessageModifier// MessageRewriter modifies message in the state, before the ChatModel is called.// It takes the messages stored accumulated in state, modify them, and put the modified version back into state.// Useful for compressing message history to fit the model context window,// or if you want to make changes to messages that take effect across multiple model calls.// NOTE: if both MessageModifier and MessageRewriter are set, MessageRewriter will be called before MessageModifier.    MessageRewriter MessageModifier// MaxStep.// default 12 of steps in pregel (node num + 10).    MaxStep int`json:"max_step"`// Tools that will make agent return directly when the tool is called.// When multiple tools are called and more than one tool is in the return directly list, only the first one will be returned.    ToolReturnDirectly map[string]struct{}// StreamToolCallChecker is a function to determine whether the model's streaming output contains tool calls.// Different models have different ways of outputting tool calls in streaming mode:// - Some models (like OpenAI) output tool calls directly// - Others (like Claude) output text first, then tool calls// This handler allows custom logic to check for tool calls in the stream.// It should return:// - true if the output contains tool calls and agent should continue processing// - false if no tool calls and agent should stop// Note: This field only needs to be configured when using streaming mode// Note: The handler MUST close the modelOutput stream before returning// Optional. By default, it checks if the first chunk contains tool calls.// Note: The default implementation does not work well with Claude, which typically outputs tool calls after text content.// Note: If your ChatModel doesn't output tool calls first, you can try adding prompts to constrain the model from generating extra text during the tool call.    StreamToolCallChecker func(ctx context.Context, modelOutput *schema.StreamReader[*schema.Message]) (bool, error)// GraphName is the graph name of the ReAct Agent.// Optional. Default `ReActAgent`.    GraphName string// ModelNodeName is the node name of the model node in the ReAct Agent graph.// Optional. Default `ChatModel`.    ModelNodeName string// ToolsNodeName is the node name of the tools node in the ReAct Agent graph.// Optional. Default `Tools`.    ToolsNodeName string}
类别字段作用
模型ToolCallingModel推荐使用的工具调用模型,通过 WithTools() 绑定工具
模型Model已废弃,通过 BindTools() 绑定工具,会修改原模型实例
工具ToolsConfig配置工具列表以及 ToolsNode 的执行行为
消息MessageRewriter改写并保存累计消息,影响后续模型调用
消息MessageModifier临时修改当前一次模型输入,不替换状态中的消息切片
循环MaxStep限制图的最大运行步数,避免无限循环
分支ToolReturnDirectly指定哪些工具执行后直接结束,不再回到模型
流式StreamToolCallChecker自定义如何从模型流中判断是否存在工具调用,未指定时默认使用 firstChunkStreamToolCallChecker
命名GraphName自定义图在日志、回调等场景中显示的名称
命名ModelNodeName自定义模型节点的显示名称
命名ToolsNodeName自定义工具节点的显示名称

ToolsConfig 本身又包含几项常用配置:

  • Tools:可被调用的工具实现;
  • ExecuteSequentially:多个 ToolCall 是串行还是并行执行,默认并行;
  • UnknownToolsHandler:模型调用不存在的工具时,是否用自定义逻辑处理;
  • ToolArgumentsHandler:工具执行前改写或校验模型生成的参数;
  • ToolCallMiddlewares:给工具调用增加中间件。

一个常见配置大致如下:

go
ra, err := react.NewAgent(ctx, &react.AgentConfig{    ToolCallingModel: chatModel,    ToolsConfig: compose.ToolsNodeConfig{        Tools:               tools,        ExecuteSequentially: false,    },    MaxStep: 20,    ToolReturnDirectly: map[string]struct{}{"download_report": {}, // 当调用了下载结果报告的工具就直接返回    },})

4.4 buildReturnDirectly()

正常的 ReAct 路径是:工具给出观察结果,模型再读取结果并生成自然语言答案。但有些工具已经能产出最终结果,例如下载地址、生成好的文件信息或无需改写的查询结果。这时再调用一次模型不仅增加延迟和消耗,还有可能让模型改动本应原样返回的内容。

buildReturnDirectly() 就是这条快捷路径的定义。

go
funcbuildReturnDirectly(graph *compose.Graph[[]*schema.Message, *schema.Message]) (err error) {    directReturn := func(ctx context.Context, msgs *schema.StreamReader[[]*schema.Message]) (*schema.StreamReader[*schema.Message], error) {return schema.StreamReaderWithConvert(msgs, func(msgs []*schema.Message) (*schema.Message, error) {var msg *schema.Message    err = compose.ProcessState[*state](ctx, func(_ context.Context, state *state)error {for i := range msgs {if msgs[i] != nil && msgs[i].ToolCallID == state.ReturnDirectlyToolCallID {    msg = msgs[i]returnnil    }    }returnnil    })if err != nil {returnnil, err    }if msg == nil {returnnil, schema.ErrNoValue    }return msg, nil    }), nil    }    nodeKeyDirectReturn := "direct_return"if err = graph.AddLambdaNode(nodeKeyDirectReturn, compose.TransformableLambda(directReturn)); err != nil {return err    }// this branch checks if the tool called should return directly. It either leads to END or back to ChatModel    err = graph.AddBranch(nodeKeyTools, compose.NewStreamGraphBranch(func(ctx context.Context, msgsStream *schema.StreamReader[[]*schema.Message]) (endNode string, err error) {    msgsStream.Close()    err = compose.ProcessState[*state](ctx, func(_ context.Context, state *state)error {iflen(state.ReturnDirectlyToolCallID) > 0 {    endNode = nodeKeyDirectReturn    } else {    endNode = nodeKeyModel    }returnnil    })if err != nil {return"", err    }return endNode, nil    }, map[string]bool{nodeKeyModel: true, nodeKeyDirectReturn: true}))if err != nil {return err    }return graph.AddEdge(nodeKeyDirectReturn, compose.END)}

1. 确定哪一个工具调用应该直接返回

静态方式是在 AgentConfig.ToolReturnDirectly 中配置工具名。工具节点运行前,getReturnDirectlyToolCallID() 会按 ToolCall 的顺序查找第一个命中的工具,并记录它的调用 ID:

go
funcgetReturnDirectlyToolCallID(input *schema.Message, toolReturnDirectly map[string]struct{})string {iflen(toolReturnDirectly) == 0 {return""    }for _, toolCall := range input.ToolCalls {if _, ok := toolReturnDirectly[toolCall.Function.Name]; ok {return toolCall.ID    }    }return""}

动态方式是在工具执行期间调用:

go
funcSetReturnDirectly(ctx context.Context)error {return compose.ProcessState(ctx, func(        ctx context.Context,        s *state,    )error {        s.ReturnDirectlyToolCallID = compose.GetToolCallID(ctx)returnnil    })}

它会取得当前工具调用的 ID 并写入状态。按照源码约定,动态设置的优先级高于 ToolReturnDirectly;同一步里有多个工具调用它时,最后一次设置生效。由于 ToolsNode 默认并行执行多个工具,这里的“最后”指实际写入状态的先后,而不是 ToolCalls 在消息中的排列顺序。

2. 决定回到模型还是直接结束

buildReturnDirectly() 在 tools 节点后添加分支。分支读取状态中的 ReturnDirectlyToolCallID:

  • ID 为空,说明工具结果仍要交给模型,下一站是 chat;
  • ID 非空,说明结果可以直接返回,下一站是 direct_return。

所以,所谓“工具直返”并不是把 tools 节点简单连接到 END,中间还需要一个结果选择节点。它也不是在执行工具前提前结束:如果同一轮包含多个 ToolCall,ToolsNode 仍会执行它们,只是最终仅向调用者输出被选中的那一条结果。

最后选出对应的工具结果

ToolsNode 可能一次执行多个工具,因此会输出 []*schema.Message;Agent 对外的输出类型却是单个 *schema.Message。direct_return 是一个 Lambda 节点,它同时完成结果筛选和类型转换:

go
for i := range msgs {if msgs[i] != nil &&        msgs[i].ToolCallID == state.ReturnDirectlyToolCallID {        msg = msgs[i]returnnil    }}

找到 ID 对应的 tool 消息后,这条消息进入 END。当前流片段中没有找到时,转换函数返回 schema.ErrNoValue;在 StreamReaderWithConvert 中,这个值表示跳过当前片段,而不是把它当成普通执行错误抛出。

这里再次说明了为什么状态中保存的是 ToolCall ID:模型可能在一条消息里调用多个工具,也可能多次调用同名工具。只凭工具名无法判断究竟应该返回哪一个结果。

还有一个容易忽略的行为:direct_return 原样返回 ToolsNode 产生的消息,所以结果的*schema.Message角色仍然是 tool,并不会被包装成一条新的 assistant 消息。如果上层业务只读取 Content,通常没有影响;如果还依赖消息角色,就要考虑这点。

5. 总结

从 Eino 的实现看,要实现一个 ReAct Agent 逻辑上并不复杂:模型节点负责决定下一步,工具节点负责执行,消息状态保存后续决策所需的上下文,两处分支决定继续循环还是结束。

研究 Eino 的 ReAct 的相关源码最重要的其实是学习思想。虽然官方 react.go 看起来功能封装得很完美,但是在实际开发 AI 应用中,该封装可能还没有灵活到能适应每一种场景。不过 ReAct 的基本原理是不变的,我们可以按需自己用基础的组件节点来编排属于自己的 ReAct Agent。

相关学习资料