CrewStructuredTool:真正被执行循环调用的那层
Day 27 的 BaseTool 是"配置态"——描述一个工具长什么样。但 Agent 执行循环(Day 08)里那句 tool.invoke(input=..., config=...) 调的其实是 to_structured_tool() 转出来的 CrewStructuredTool。今天读 tools/structured_tool.py:它怎么把 LLM 吐出来的字符串/字典参数解析成能喂给函数的 kwargs、怎么同步和异步各准备一套、以及 result_schema 怎么把杂乱的返回值格式化成给 Agent 看的规整文本。这是"配置"变"执行"的那一层。
BaseTool 是"菜谱卡片"(写着菜名、原料、步骤),那 CrewStructuredTool 就是后厨那口真正开火的锅。顾客(LLM)点单时递来的往往是一张手写的、格式随意的纸条(可能是 JSON 字符串,可能是字典)。锅(invoke)先让服务员(_parse_args)把纸条翻译成标准订单(校验过的 kwargs),再决定用大火(同步)还是慢炖(异步)开做,最后把菜摆盘(result_schema 格式化)端出去。痛点:为什么不直接用 BaseTool 执行?
BaseTool.run() 明明也能执行工具,为啥还要多一个 CrewStructuredTool?它们不是重复了吗?而且 LLM 给的参数经常是一坨字符串(比如 '{"city": "上海"}'),甚至是有点破损的 JSON——谁来把它变成 Python 函数能收的关键字参数?CrewStructuredTool 是"执行态"的统一外壳:不管背后是 @tool 函数、继承 BaseTool 的类、还是 MCP 远程工具,最后都被 to_structured_tool() 归一成 CrewStructuredTool,让执行循环只需面对一种调用接口 invoke()。它比 BaseTool 多干的活:把 LLM 的字符串参数解析成 dict、根据函数是不是协程自动切同步/异步、用 result_schema规整输出。BaseTool 管"定义",CrewStructuredTool 管"跑"。类的骨架(structured_tool.py:115):
# structured_tool.py:115
class CrewStructuredTool(BaseModel):
name: str = Field(default="")
description: str = Field(default="")
args_schema: ... = Field(default=None) # 参数 schema(带序列化器)
result_schema: ... = Field(default=None) # 输出 schema
func: Any = Field(default=None, exclude=True) # ★被包裹的真正函数
result_as_answer: bool = Field(default=False)
max_usage_count: int | None = Field(default=None)
current_usage_count: int = Field(default=0)
cache_function: Any = Field(default=None, exclude=True)
_original_tool: Any = PrivateAttr(default=None) # 指回原 BaseTool
def invoke(self, input, config=None, **kwargs) -> Any: ... # :340 主入口
async def ainvoke(self, input, config=None, **kwargs): ... # :296 异步入口
def _parse_args(self, raw_args) -> dict: ... # :272 解析参数
BaseTool 很像,但多了一个 func——就是"到底调哪个函数"。还有个 _original_tool 指针,指回它是从哪个 BaseTool 转来的(后面格式化输出时要用回原工具的方法)。exclude=True 表示 func 不参与序列化(函数没法存进 JSON)。从 BaseTool 到 StructuredTool 的那座桥
回看 Day 27 的 to_structured_tool(base_tool.py:393)——就是这座桥:
# base_tool.py:393
def to_structured_tool(self) -> CrewStructuredTool:
self._set_args_schema() # 确保 args_schema 就绪
structured_tool = CrewStructuredTool(
name=self.name, description=self.description,
args_schema=self.args_schema, result_schema=self.result_schema,
func=self._run, # ★把 _run 作为要调的函数
result_as_answer=self.result_as_answer,
max_usage_count=self.max_usage_count,
current_usage_count=self.current_usage_count,
cache_function=self.cache_function)
structured_tool._original_tool = self # ★留一根指针指回自己
return structured_tool
func=self._run把配置态工具的"真正干活方法" _run 交给执行态工具当 func。这样 invoke 最终调的还是你写的 _run。把 result_as_answer / usage 都搬过来所有行为开关一起搬,保证执行态工具行为和你配置的一致。_original_tool = self★关键指针。执行完输出格式化时(L07)会用 _original_tool.format_output_for_agent 回到原工具的逻辑,也用它把 usage 计数同步回原工具。BaseTool 面向用户(好写、好读、能序列化存档),CrewStructuredTool 面向执行引擎(只管跑得对、跑得稳)。执行前统一 to_structured_tool() 归一。代价是多一次转换、多一根 _original_tool 指针;收益是执行循环只认一种接口,MCP 工具、langchain 工具都能被归一进来。from_function:直接从裸函数造执行态工具
除了从 BaseTool 转,也能直接从一个函数造(structured_tool.py:151):
# structured_tool.py:151
@classmethod
def from_function(cls, func, name=None, description=None,
return_direct=False, args_schema=None,
result_schema=None, infer_schema=True, **kwargs):
name = name or func.__name__ # 没名字 → 用函数名
description = description or inspect.getdoc(func) # 没描述 → 用 docstring
if description is None:
raise ValueError(f"Function {name} must have a docstring ...")
description = textwrap.dedent(description).strip()
if args_schema is not None:
schema = args_schema
elif infer_schema:
schema = cls._create_schema_from_function(name, func) # ★从签名推断
else:
raise ValueError("Either args_schema must be provided or infer_schema must be True.")
return cls(name=name, description=description, args_schema=schema,
result_schema=result_schema or _infer_result_schema_from_callable(func),
func=func, result_as_answer=return_direct, **kwargs)
推断 schema 的 _create_schema_from_function(structured_tool.py:217)和 Day 27 思路一致,但用 get_type_hints:
# structured_tool.py:217
@staticmethod
def _create_schema_from_function(name, func):
sig = inspect.signature(func)
type_hints = get_type_hints(func) # 解析"字符串形式的类型标注"
fields = {}
for param_name, param in sig.parameters.items():
if param_name in ("self", "cls"):
continue
annotation = type_hints.get(param_name, Any)
default = ... if param.default == param.empty else param.default
fields[param_name] = (annotation, Field(default=default))
schema_name = f"{name.title()}Schema"
return create_model(schema_name, **fields)
name / description 兜底不传就用函数名和 docstring。和 @tool 一样强制要有描述——没描述直接报错。get_type_hints(func)比 Day 27 的 param.annotation 更进一步:能正确解析 from __future__ import annotations 下"字符串形式"的类型标注。_infer_result_schema_from_callable看函数返回标注是不是一个 BaseModel 子类,是就当 result_schema——为 L07 的结构化输出铺路。return_direct → result_as_answerlangchain 世界叫 return_direct,这里映射成 result_as_answer。命名兼容历史生态。_parse_args:把"一坨字符串"翻译成参数
'{"city":"上海"}',也可能已经是 dict。函数收的却是关键字参数。中间必须有个"翻译 + 校验"的环节。_parse_args(structured_tool.py:272)就干这个:
# structured_tool.py:272
def _parse_args(self, raw_args: str | dict) -> dict:
if isinstance(raw_args, str):
try:
raw_args = json.loads(raw_args) # ① 字符串 → dict
except json.JSONDecodeError as e:
raise ValueError(f"Failed to parse arguments as JSON: {e}") from e
if not self.args_schema:
return raw_args if isinstance(raw_args, dict) else {}
try:
validated_args = self.args_schema.model_validate(raw_args) # ② schema 校验
return dict(validated_args.model_dump()) # ③ 返回清洗后的 dict
except Exception as e:
hint = build_schema_hint(self.args_schema) # ④ 校验失败 → 附上期望参数
raise ValueError(f"Arguments validation failed: {e}{hint}") from e
isinstance(raw_args, str)先判断是不是字符串。是 → json.loads 转成 dict。不是(已是 dict)→ 跳过。JSONDecodeError → ValueErrorJSON 坏了明确报错。注意这里不尝试修复破损 JSON——修复逻辑在 D29 的 _validate_tool_input 层,这里只做干净解析。model_validate + model_dump用 schema 校验并做类型转换,返回规范化的 dict。多余字段会被丢弃、类型会被强制转换。build_schema_hint和 Day 27 同款:失败时附上"期望哪些参数、哪些必填",这段提示帮 LLM 下一圈改对。invoke:主执行流程
同步主入口 invoke(structured_tool.py:340):
# structured_tool.py:340
def invoke(self, input, config=None, **kwargs) -> Any:
parsed_args = self._parse_args(input) # ① 解析参数
if self.has_reached_max_usage_count(): # ② 查用量上限
raise ToolUsageLimitExceededError(
f"Tool '...' has reached its maximum usage limit of {self.max_usage_count}. ...")
self._increment_usage_count() # ③ 计数 +1
if inspect.iscoroutinefunction(self.func): # ④ func 本身是 async?
return asyncio.run(self.func(**parsed_args, **kwargs))
result = self.func(**parsed_args, **kwargs) # ⑤ 普通函数:直接调
if asyncio.iscoroutine(result): # ⑥ 返回了协程?跑完它
return asyncio.run(result)
return result
用量检查 has_reached_max_usage_count 与计数同步(structured_tool.py:366):
# structured_tool.py:366
def has_reached_max_usage_count(self) -> bool:
return (self.max_usage_count is not None
and self.current_usage_count >= self.max_usage_count)
def _increment_usage_count(self) -> None:
self.current_usage_count += 1
if self._original_tool is not None:
self._original_tool.current_usage_count = self.current_usage_count # ★同步回原工具
解析 → 查上限 → 计数顺序和 Day 27 run 一致,但这里超限是 raise 异常(ToolUsageLimitExceededError),由上层 tool_usage.py 捕获转成给 LLM 的话(D29/D30 会看到)。iscoroutinefunction 分流★两处判断:函数本身是 async → 用 asyncio.run;函数是普通的但返回了协程 → 也 asyncio.run 跑完。双保险,保证 invoke 永远返回实际结果而非协程。_increment 同步回 _original_tool★执行态工具计数 +1 时,把数字同步回原 BaseTool。这样无论你查哪个对象的 current_usage_count,看到的都是一致的。asyncio.run 会新建一个事件循环,如果当前线程已经在一个事件循环里(比如在 async 上下文调 invoke),会抛 RuntimeError: asyncio.run() cannot be called from a running event loop。所以异步场景要走 ainvoke(L06)。MCP 工具遇到这问题时的解法是"开个新线程跑 asyncio.run"(D31 会看到 MCPNativeTool._run 正是这么干的)。坑:在异步代码里误用同步 invoke。ainvoke:为什么异步要单独写一份
异步入口 ainvoke(structured_tool.py:296):
# structured_tool.py:296
async def ainvoke(self, input, config=None, **kwargs) -> Any:
parsed_args = self._parse_args(input)
if self.has_reached_max_usage_count():
raise ToolUsageLimitExceededError(...)
self._increment_usage_count()
try:
if inspect.iscoroutinefunction(self.func):
return await self.func(**parsed_args, **kwargs) # ① async 函数 → 直接 await
import asyncio
return await asyncio.get_event_loop().run_in_executor( # ② 同步函数 → 丢线程池
None, lambda: self.func(**parsed_args, **kwargs))
except Exception:
raise
iscoroutinefunction → await如果 func 本身是异步的,直接 await——不阻塞事件循环,这正是异步的意义。run_in_executor★如果 func 是同步的(可能很慢,比如读文件),直接调会卡住整个事件循环。所以丢进线程池执行,用 await 等结果——同步函数也能"异步友好"。参数解析/限流复用前三步(解析、查上限、计数)和 invoke 完全一样。只有"怎么跑函数"这步分叉。def get_weather(...)),逼他们全改 async 门槛太高;而且很多场景(简单脚本、非并发)根本不需要异步。源码做法:两套接口并存,且互相兜底——invoke 遇到协程会 asyncio.run 跑完,ainvoke 遇到同步函数会丢线程池。用户写哪种都能跑。代价是代码重复(两份几乎一样的前置逻辑)和"哪个场景用哪个"的认知负担。收益是对用户"写什么都行"的最大包容度——对一个面向广大开发者的框架,这个取舍很值。result_schema:把返回值"摆盘"成规整文本
工具返回的可能是对象、字典、任意东西,但喂给 LLM 的必须是字符串。_format_tool_output_for_agent(structured_tool.py:54)负责这一步:
# structured_tool.py:54
def _format_tool_output_for_agent(tool, raw_result):
original_tool = getattr(tool, "_original_tool", None)
if original_tool is not None:
return original_tool.format_output_for_agent(raw_result) # ① 有原工具 → 用它的逻辑
result_schema = getattr(tool, "result_schema", None)
if not (isinstance(result_schema, type) and issubclass(result_schema, BaseModel)):
return str(raw_result) # ② 没 result_schema → 直接 str()
try:
validation_input = raw_result
if isinstance(raw_result, BaseModel) and not isinstance(raw_result, result_schema):
validation_input = raw_result.model_dump()
validated = result_schema.model_validate(validation_input)
return validated.model_dump_json() # ③ 有 schema → 校验后转 JSON
except Exception as exc:
warnings.warn(f"Failed to validate ... Falling back to str(raw_result).", ...)
return str(raw_result) # ④ 校验失败 → 优雅降级为 str()
先看 _original_tool如果是从 BaseTool 转来的,回到原工具的 format_output_for_agent(它最终还是调这个函数,但确保用原工具的 schema)。没 result_schema → str()大多数简单工具不定义输出 schema,直接把返回值 str() 化。够用。有 schema → 校验 + JSON定义了 result_schema 时:先用它校验返回值(保证结构对),再序列化成 JSON 字符串喂给 LLM。让下游能稳定解析。校验失败 → warn + str()★关键的优雅降级:schema 不匹配时不崩,只发个警告,退回 str(raw_result)。宁可给个不规整的结果,也不让工具因为格式化失败而报错。WeatherReport(temp=18, desc="晴") 对象,且 result_schema=WeatherReport:→ 输出被格式化为
{"temp":18,"desc":"晴"}(JSON 字符串),LLM/下游可稳定解析。若没定义 result_schema:→ 输出是
str(WeatherReport(...)),可能是 temp=18 desc='晴' 这种不好解析的形式。边界 + 今日小结
CrewStructuredTool 也保留了一个 _run(structured_tool.py:332),标注是 "Legacy method for compatibility"。它把位置参数 args 按 args_schema 的字段顺序zip 成关键字,再转调 invoke。这是为了兼容"有些地方还按位置参数调工具"的老代码。坑:如果你的 args_schema 字段顺序和函数签名顺序不一致,按位置调用会串位。所以现代用法一律走 invoke(input=dict) 显式传参,别依赖位置。👶 小白:func 为什么设了 exclude=True?
👨🏫 老师:因为函数(可调用对象)没法序列化成 JSON。工具在存档/传输时(比如 Crew 状态持久化)需要能 dump 成 JSON,而 func 存不了。所以标 exclude=True 让它不参与序列化——存档时只存"这是什么工具"(name/description/schema),重新加载时再从原始定义恢复 func。cache_function 同理。
🧠 今天你应该能回答
- 为什么要有"执行态"的
CrewStructuredTool,它比BaseTool多干什么? to_structured_tool这座桥搬了什么?_original_tool有什么用?_parse_args怎么把字符串/dict 变成校验过的 kwargs?invoke里两处asyncio.run分别处理什么情况?- 为什么同步/异步两套接口并存?
run_in_executor解决什么? result_schema格式化失败时为什么"降级"而不是报错?
✋ 10 分钟动手
P=lib/crewai/src/crewai/tools
sed -n '115,143p' $P/structured_tool.py # 类字段
sed -n '272,295p' $P/structured_tool.py # _parse_args
sed -n '340,378p' $P/structured_tool.py # invoke + 用量
sed -n '54,84p' $P/structured_tool.py # 输出格式化
python -c "
from crewai.tools import tool
@tool
def add(a:int,b:int)->int:
'''相加'''
return a+b
st = add.to_structured_tool()
print(type(st).__name__)
print(st.invoke(input={'a':2,'b':3})) # 传 dict
print(st.invoke(input='{\"a\":10,\"b\":20}')) # 传 JSON 字符串也行
"
invoke 调用?明天读 tool_calling.py(ToolCalling 数据结构)和 tool_usage.py 的 _tool_calling / _select_tool / _validate_tool_input:模型给的工具名对不上怎么模糊匹配、参数是破损 JSON 怎么修、解析失败怎么重试。