Day 55 / 共 60 天 · 阶段 9 预制件与流式
ToolNode:把"执行工具"这件脏活干干净净
agent 节点负责"想",tools 节点负责"做"。这个 tools 节点就是 ToolNode。别看它在图里只是一个节点,源码有 2000 行——因为"执行工具"的真实世界充满脏活:并行跑多个工具、工具报错不能崩全图、参数校验失败要翻译成模型能懂的话、还要往工具里注入 state/store/运行时。今天把它的主干拆开。
📍 阶段 9 · 预制件与流式(6 天)你在这里
D53 总览→
D54 逐行①→
D55 ToolNode→
D56 校验→
D57 流式5模式→
D58 流式底层
💡 用一个类比先兜住今天
ToolNode 像一个后厨调度台。前台(模型)下了一叠订单(tool_calls),调度台先看懂订单(_parse_input)、把订单分给多个厨师同时做(线程池并行)、每个厨师做菜时万一翻车了不许喊崩整个厨房,而要写一张"这道菜出问题了:XXX"的纸条端回给前台(错误转 ToolMessage)、最后把所有成品和纸条装盘归位(合并输出)。今天看这套调度怎么写。
L01
ToolNode 是什么:一个会执行工具的 Runnable
🤔 痛点模型输出的
AIMessage.tool_calls 只是"我想调 weather(city=北京)"的意图,不是真调用。谁把意图变成真实函数调用、再把返回值包成模型能读的 ToolMessage?就是 ToolNode。它继承自 RunnableCallable(tool_node.py:622),所以能当图节点用。核心入口是 _func(同步)/_afunc(异步)。整体职责:输入含 tool_calls 的消息/状态 → 执行工具 → 输出 ToolMessage 列表。
# tool_node.py:622(类头,裁剪 docstring)
class ToolNode(RunnableCallable):
"""A node that runs the tools called in the last AIMessage."""
# 关键内部状态(在 __init__ 里建,见 L02):
# self._tools_by_name : dict[str, BaseTool] 工具名 → 工具
# self._handle_tool_errors : 错误处理策略
# self._injected_args : 每个工具需要注入哪些特殊参数
👶 一句话ToolNode = "把模型的调用意图,真正执行掉,并把结果/错误都变成 ToolMessage"。它是意图和现实之间的翻译官+执行者。
D53 里 create_react_agent 会自动帮你把 tools 列表包成 ToolNode。但你也能单独用它,配到任何自己的 StateGraph 里——它就是个普通节点。
L02
__init__:把工具建成"名字 → 工具"的索引
tool_node.py:743,构造时做两件准备:
# tool_node.py:743
def __init__(self, tools, *, name="tools", tags=None,
handle_tool_errors=_default_handle_tool_errors,
messages_key="messages", wrap_tool_call=None, awrap_tool_call=None):
super().__init__(self._func, self._afunc, name=name, tags=tags, trace=False)
self._tools_by_name = {}
self._injected_args = {}
self._handle_tool_errors = handle_tool_errors # 错误策略(L06)
self._messages_key = messages_key
self._wrap_tool_call = wrap_tool_call # 拦截器(重试/缓存)
for tool in tools:
if not isinstance(tool, BaseTool):
tool_ = create_tool(tool) # 普通函数 → 包成 BaseTool
else:
tool_ = tool
self._tools_by_name[tool_.name] = tool_ # ★ 建 名字→工具 索引
self._injected_args[tool_.name] = _get_all_injected_args(tool_) # 预扫注入参数
super().__init__(self._func, self._afunc)把同步 _func、异步 _afunc 注册给父类 RunnableCallable——这就是节点被调用时真正跑的东西。create_tool(tool)你传普通 Python 函数也行,这里自动包成 BaseTool(借 LangChain 的机制推断参数 schema)。self._tools_by_name[name] = tool_建 名字→工具 字典。执行时按模型给的 call["name"] 直接 O(1) 查表。_get_all_injected_args(tool_)提前一次性扫描每个工具"需要注入哪些特殊参数"(如 InjectedState/InjectedStore)。存起来避免每次执行都重扫。💡 设计取舍①:为什么在 __init__ 就预扫 injected_args,而不是执行时现扫?工具的注入参数(哪些形参要塞 state、哪些要塞 store)在整个 Agent 生命周期里不变。执行时现扫要用
inspect 反射签名——反射慢且每次工具调用都做一遍很浪费。放到 __init__ 扫一次、存进 self._injected_args,之后每次执行直接查字典。这是"不变量前移到初始化"的经典优化:一次成本换 N 次省。L03
_parse_input:从三种输入里挖出 tool_calls
节点被调用时,输入可能是消息列表、状态 dict、或 Send 传来的特殊结构。tool_node.py:1224 统一解析:
# tool_node.py:1224
def _parse_input(self, input) -> tuple[list[ToolCall], Literal["list","dict","tool_calls"]]:
if isinstance(input, list):
if isinstance(input[-1], dict) and input[-1].get("type") == "tool_call":
return cast(list[ToolCall], input), "tool_calls" # 直接给了 tool_call 列表
input_type = "list"; messages = input
elif isinstance(input, dict) and input.get("__type") == "tool_call_with_context":
input_with_ctx = cast(ToolCallWithContext, input) # ★ D53 v2 的 Send 载荷
return [input_with_ctx["tool_call"]], "tool_calls"
elif isinstance(input, dict) and (messages := input.get(self._messages_key, [])):
input_type = "dict"
elif messages := getattr(input, self._messages_key, []):
input_type = "dict" # dataclass 状态
else:
raise ValueError("No message found in input")
try:
latest_ai_message = next(m for m in reversed(messages) if isinstance(m, AIMessage))
except StopIteration:
raise ValueError("No AIMessage found in input")
return list(latest_ai_message.tool_calls), input_type
input[-1]... "tool_call"你直接喂一个 tool_call 列表(少见)→ 原样用,标记 "tool_calls"。__type == "tool_call_with_context"关键:D53 v2 用 Send 分发时,每个任务的输入就是这个结构,里面装着单个 tool_call + state 快照。这里取出那一个 tool_call。input.get(messages_key)最常见:输入是状态 dict,从 messages 里取。next(... reversed ... AIMessage)从后往前找最近的一条 AIMessage,它的 tool_calls 就是待执行清单。返回 input_type记住输入形态(list/dict/tool_calls),因为输出格式要跟输入对齐(L07)。💡 本质:一个节点要同时服务 v1 和 v2v1 整条状态进来(走 dict/list 分支,一次拿到全部 tool_calls);v2 每个 Send 单独进来(走 tool_call_with_context 分支,只拿一个)。
_parse_input 用输入结构区分两条路,让 _func 后续逻辑统一——不管几个 tool_call,都当"一个列表"处理。L04
_func:线程池并行执行每个工具
主流程 tool_node.py:793:
# tool_node.py:793
def _func(self, input, config, runtime):
tool_calls, input_type = self._parse_input(input) # L03
config_list = get_config_list(config, len(tool_calls))
tool_runtimes = []
for call, cfg in zip(tool_calls, config_list):
state = self._extract_state(input, cfg)
tool_runtime = ToolRuntime(state=state, tool_call_id=call["id"], config=cfg,
context=runtime.context, store=runtime.store, # 注入运行时资源
stream_writer=runtime.stream_writer, tools=list(self.tools_by_name.values()), ...)
tool_runtimes.append(tool_runtime)
input_types = [input_type] * len(tool_calls)
with get_executor_for_config(config) as executor:
outputs = list(executor.map(self._run_one, tool_calls, input_types, tool_runtimes)) # ★ 并行
return self._combine_tool_outputs(outputs, input_type) # L07 合并
for call in tool_calls: ToolRuntime(...)给每个工具调用造一个 ToolRuntime:把 state、store、stream_writer、config 等运行时资源打包,稍后按需注入进工具。get_executor_for_config(config)拿到一个线程池执行器(来自 config,可复用图的执行器配置)。executor.map(self._run_one, ...)今天的高潮:把每个 tool_call 丢进线程池并行跑 _run_one。三个工具就三个线程同时跑,不用串行等。异步版 _afunc(:828)对应异步:用 asyncio.gather(*coros) 并发 await 每个 _arun_one。同样是并行,只是换成协程。_combine_tool_outputs把 N 个工具的结果合并成节点的返回值(L07)。💡 设计取舍②:v2 已经用 Send 并行了,为什么 ToolNode 内部还要线程池并行?两者是不同层次的并行,且互补。v2(D53)在图层把每个 tool_call 拆成独立任务,好处是每个工具能独立中断/重试/记检查点。但当 ToolNode 被单独用在别人的图里(不走 create_react_agent 的 v2),或 v1 模式下,一个节点里就有多个 tool_call——此时靠节点内的线程池
executor.map 并行。ToolNode 自己保证"给我多少工具我都并行跑",不依赖上层怎么调度。这是"组件自洽"的设计:不假设调用方一定用 v2。⚠️ 边界:工具是 CPU 密集或阻塞 IO 时,线程池并行才有意义线程池对IO 型工具(调 API、查数据库)能真并行(GIL 在 IO 等待时释放)。但纯 CPU 计算的工具受 GIL 限制,线程并行提速有限。异步版
_afunc 要求工具本身是 async 的才能真正并发。选同步还是异步,取决于你的工具是什么类型——用错了会"看似并行实则排队"。L05
_execute_tool_sync:报错不崩,转成 ToolMessage
单个工具的执行 + 兜错在 tool_node.py:922:
# tool_node.py:922(裁剪)
def _execute_tool_sync(self, request, input_type, config):
call = request.tool_call; tool = request.tool
if tool is None: # 模型点了不存在的工具
if invalid_tool_message := self._validate_tool_call(call):
return invalid_tool_message # 返回"没这个工具"的 ToolMessage
injected_call = self._inject_tool_args(call, request.runtime, tool) # 注入 state/store
call_args = {**injected_call, "type": "tool_call"}
try:
try:
response = tool.invoke(call_args, config) # ★ 真正执行工具
except ValidationError as exc: # 参数校验失败
filtered_errors = _filter_validation_errors(exc, self._injected_args.get(call["name"]))
raise ToolInvocationError(call["name"], exc, call["args"], filtered_errors) from exc
return self._normalize_tool_response(response, request.tool_call, input_type)
except GraphBubbleUp: # interrupt() 等特殊信号
raise # ← 必须原样抛,不能吞
except Exception as e:
handled_types = ... # 根据配置定"哪些异常要接住"
if not self._handle_tool_errors or not isinstance(e, handled_types):
raise # 不在处理范围 → 抛出
content = _handle_tool_error(e, flag=self._handle_tool_errors) # 转成错误文本
return ToolMessage(content=content, name=call["name"],
tool_call_id=call["id"], status="error") # ★ 错误也变消息
tool is None → _validate_tool_call模型幻觉调了个不存在的工具,不崩,返回一条"这不是有效工具,可用的有[...]"的 ToolMessage(L06 会看模板)。_inject_tool_args(...)执行前把 state、store、tool_call_id 等注入进工具参数(对声明了 InjectedState 的工具)。模型看不到这些参数、也不用填。tool.invoke(call_args, config)真正执行工具函数。except GraphBubbleUp: raise关键边界:如果工具里调了 interrupt()(人在环,D41),会抛 GraphBubbleUp。这类信号绝不能被错误处理吞掉,必须原样往上抛,让图去暂停。except Exception → ToolMessage(status="error")今天的高潮:普通异常被接住,转成一条 status="error" 的 ToolMessage 喂回模型——模型看到报错文本,能自己纠正重试。工具翻车,Agent 不崩。💡 本质:错误也是给模型的"反馈"朴素思路:工具报错就抛异常、整个 Agent 挂掉。ToolNode 的思路:把错误翻译成模型能读的文本塞回对话。模型看到"Error: 城市名无效",下一轮可能就自己改成正确参数重试。把异常变成对话的一部分,是让 Agent 具备"自我纠错"能力的关键机制。
L06
handle_tool_errors:四种错误处理形态
到底怎么把异常变文本?看 _handle_tool_error(tool_node.py:394)和默认策略(tool_node.py:383):
# tool_node.py:383 默认策略
def _default_handle_tool_errors(e: Exception) -> str:
if isinstance(e, ToolInvocationError): # 只兜"参数校验类"错误
return e.message
raise e # 其它错误照样抛(不隐藏真 bug)
# tool_node.py:394 按配置类型生成错误文本
def _handle_tool_error(e, *, flag) -> str:
if isinstance(flag, (bool, tuple)) or (isinstance(flag, type) and issubclass(flag, Exception)):
content = TOOL_CALL_ERROR_TEMPLATE.format(error=repr(e)) # 用默认模板
elif isinstance(flag, str):
content = flag # 用你给的固定字符串
elif callable(flag):
content = flag(e) # 用你给的函数生成
else:
raise ValueError("Got unexpected type of `handle_tool_error`...")
return content
# TOOL_CALL_ERROR_TEMPLATE = "Error: {error}\n Please fix your mistakes." (:111)
| handle_tool_errors 传值 | 行为 |
|---|---|
True(或异常类型/元组) | 用模板 "Error: {repr(e)}\n Please fix your mistakes." |
"自定义文本"(str) | 不管什么错,都返回这句固定话 |
lambda e: ...(callable) | 调你的函数,按异常内容定制错误文本 |
False | 不处理,异常直接抛出(让图崩/上层处理) |
默认只兜 ToolInvocationError默认策略很克制:只把"参数校验失败"这类模型能自己修的错误转成消息;其它错误(如工具内部真 bug)照样抛,不掩盖真问题。flag 是异常类型/元组你可以只让某几类异常被接住(如 (ValueError, KeyError)),其余照抛——精细控制"哪些错该喂回模型、哪些该崩"。_infer_handled_types(:444)更妙:你传的自定义处理函数若标注了 def h(e: ValueError),源码用反射读你的类型注解,只接住 ValueError。函数签名即配置。⚠️ 边界:默认不是"接住一切",而是"只接住模型能修的"很多人以为 ToolNode 会吞掉所有异常。其实默认策略(
_default_handle_tool_errors)只对 ToolInvocationError(参数不对)返回消息,其它 raise e。这是刻意的:如果把"数据库连接失败"也悄悄转成 ToolMessage 喂给模型,模型只会一脸懵地重试,真正的基础设施故障却被隐藏了。可恢复的错误喂回模型,不可恢复的错误暴露出来——这条界线是默认策略的精髓。L07
_combine_tool_outputs:装盘归位 + 今日小结
并行跑完,N 个结果要合并成节点返回值。tool_node.py:862(裁剪):
# tool_node.py:862
def _combine_tool_outputs(self, outputs, input_type):
# 有工具返回了 list(多条消息)→ 先摊平
if any(isinstance(o, list) for o in outputs):
flat_outputs = []
for o in outputs:
flat_outputs.extend(o) if isinstance(o, list) else flat_outputs.append(o)
else:
flat_outputs = outputs
# 常见情况:没有 Command,直接按输入形态返回
if not any(isinstance(o, Command) for o in flat_outputs):
return flat_outputs if input_type == "list" else {self._messages_key: flat_outputs}
# 复杂情况:有工具返回 Command(要跳转/改状态),逐个处理并合并 parent goto
...
摊平 list某个工具可能返回多条 ToolMessage(一个 list),先展平成一维。没 Command → 按 input_type 返回输入是 list 就返回 list;输入是 dict 就返回 {"messages": [...]}。输出格式对齐输入(呼应 L03)。有 Command工具可以返回 Command(D15)来跳转/改状态。这里把多个 Command 的 goto 合并、和普通 ToolMessage 一起打包,交给 LangGraph 处理。图注:查表找工具 + 查注入清单,为每个调用构造 ToolRuntime。
图注:多个工具并行进 try 块,各自兜错,最后合并。
🧠 今天你应该能回答
- ToolNode 的职责?(把 tool_calls 意图变成真实执行,结果/错误都包成 ToolMessage)
- __init__ 为什么预扫 injected_args?(注入需求是不变量,前移到初始化省反射)
- _parse_input 如何区分 v1/v2?(v2 走 tool_call_with_context 分支,一次一个)
- 工具怎么并行?(_func 用线程池 executor.map;_afunc 用 asyncio.gather)
- 工具报错为什么不崩全图?(转成 status="error" 的 ToolMessage 喂回模型自纠)
- interrupt() 抛出的 GraphBubbleUp 为什么必须原样抛?(是暂停信号,不能被错误处理吞)
- 默认错误策略接住哪些错?(只 ToolInvocationError,其余照抛,不掩盖真 bug)
✋ 10 分钟动手
# 1. 三段核心逐行读
sed -n '793,826p' libs/prebuilt/langgraph/prebuilt/tool_node.py # _func 并行
sed -n '922,1013p' libs/prebuilt/langgraph/prebuilt/tool_node.py # 执行+兜错
sed -n '383,393p' libs/prebuilt/langgraph/prebuilt/tool_node.py # 默认错误策略
# 2. 亲手看"工具报错变消息"
python - <<'PY'
from langgraph.prebuilt import ToolNode
from langchain_core.tools import tool
from langchain_core.messages import AIMessage
@tool
def div(a:int,b:int)->float:
"除法"
return a/b
node = ToolNode([div], handle_tool_errors=True) # 接住一切错误
ai = AIMessage(content="", tool_calls=[{"name":"div","args":{"a":1,"b":0},"id":"1","type":"tool_call"}])
out = node.invoke({"messages":[ai]})
print(out["messages"][0].content, out["messages"][0].status) # 看到 Error: ... error
PY
明天预告 · Day 56:看
tools_condition(比 should_continue 更通用的路由函数)如何判断"该不该去 tools",再看 ValidationNode(tool_validator.py)——一个"只校验参数不真跑工具"的节点,用于结构化抽取和"让模型重填直到合法"的重试循环。