Day 42 / 共 60 天 · 阶段7 Flow 事件驱动
@start / @listen:把普通方法"贴标签"变成流程图的节点
Day 41 我们看到 @start/@listen 把方法连成了图。今天钻进 flow/dsl/,逐行读这几个装饰器:它们不改变函数的行为,只是往函数上贴一张"元数据标签"(FlowMethodDefinition),记下"我是不是起点、我监听谁"。我们会看到装饰器怎么把函数包成 StartMethod/ListenMethod、描述符协议 __get__ 怎么让包装后的方法还能正常 self.xxx() 调用、以及 or_/and_ 怎么把"任一/全部"写成一棵条件树。
📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆→
D41 Flow 总览→
D42 装饰器→
D43 定义契约→
D44 路由跳转→
D45 状态持久化→
D46 表达式→
D47 对话式→
D48 选型→
阶段8 LLM
💡 先用一个类比兜住今天
装饰器就像给员工别一枚工牌。
@start() 是"我是第一个上班的",@listen("A") 是"A 交班给我我才开工"。别工牌不改变员工怎么干活(函数体不变),只是让调度台(引擎)扫一眼工牌就知道"谁是起点、谁等谁"。今天读的就是"怎么别工牌"(装饰器)和"工牌里记了啥"(FlowMethodDefinition)。L01
痛点:装饰器没改函数,那它到底干了啥?
🤔 痛点你写
@start() / @listen(fetch),函数体一个字没动,可 Flow 就"知道"了执行顺序。这个"知道"存在哪、长什么样?引擎是怎么读到的?如果不理解这层,Flow 对你就是黑魔法。💡 一句话本质
装饰器做两件事:① 把你的函数包进一个
FlowMethod 子类对象(StartMethod/ListenMethod/RouterMethod),让引擎能通过类型认出它;② 往这个对象上挂一个 FlowMethodDefinition,用 start=True / listen=条件 记下角色和触发条件。装饰器 = "包装 + 贴元数据",纯声明,不改行为。引擎日后扫描类,看类型 + 读这份元数据,就还原出整张图。大白话函数还是那个函数,装饰器只是在它头上贴了张便利贴:"我是起点"或"我等 fetch"。真正干活时,引擎照着便利贴排班。
L02
@start:逐行拆一个最简单的装饰器
dsl/_start.py 全文很短(flow/dsl/_start.py:18):
# flow/dsl/_start.py:18
def start(condition: FlowTrigger | None = None) -> FlowMethodDecorator:
"""Marks a method as a flow's starting point."""
def decorator(func: Callable[P, R]) -> StartMethod[P, R]:
wrapper = StartMethod(func) # ① 包成 StartMethod 对象
_merge_flow_method_definition( # ② 挂上元数据
wrapper,
FlowMethodDefinition(
do=_method_action(func), # 记住"这方法怎么调"
start=( # 记住"起点 + 触发条件"
_to_definition_condition(condition)
if condition is not None
else True # 无条件 → start=True
),
),
)
return wrapper
return cast(FlowMethodDecorator, decorator)
start(condition=None)这是"装饰器工厂"——先调一次 @start() 传参,返回真正的 decorator。所以要写 @start() 带括号,不是 @start。StartMethod(func)把你的函数塞进 StartMethod 包装对象。引擎日后用 isinstance(x, StartMethod) 就能认出它是起点。_method_action(func)记下"这方法怎么被调用"——具体是 FlowCodeActionDefinition(ref="模块:限定名")(见 L07)。start = 条件 or True★没给条件 → start=True(无条件起点,kickoff 就跑);给了条件 → 转成 Definition 条件("某事件发生才作为起点启动")。📝 例子:两种 @start
@start() def a(self): ... → 元数据 start=True,一 kickoff 立即执行。@start("retry") def a(self): ... → 元数据 start="retry",表示"当 retry 事件出现时,把 a 当起点重跑"——这是实现循环重试的关键写法(配合 D44 router)。L03
@listen:和 @start 几乎一样,只差一个字段
dsl/_listen.py(flow/dsl/_listen.py:18),对比着看:
# flow/dsl/_listen.py:18
def listen(condition: FlowTrigger) -> FlowMethodDecorator: # 注意:condition 必填
def decorator(func: Callable[P, R]) -> ListenMethod[P, R]:
wrapper = ListenMethod(func) # 包成 ListenMethod
_merge_flow_method_definition(
wrapper,
FlowMethodDefinition(
do=_method_action(func),
listen=_to_definition_condition(condition), # ← 填的是 listen 字段
),
)
return wrapper
return cast(FlowMethodDecorator, decorator)
condition 必填★和 @start 不同:@listen 必须说明"等谁",一个没有触发条件的 listener 毫无意义。类型强制了这一点。ListenMethod换了个包装类型。引擎区分 Start / Listen / Router 全靠这三个类型 + 元数据里的字段。listen=... 而非 start=...唯一实质差别:填的是 listen 字段。@start 填 start,@listen 填 listen——引擎据此判断这方法是"起点"还是"监听者"。💡 设计取舍①:为什么 @start 和 @listen 长得几乎一样,还要拆成两个装饰器?
两者代码高度相似,本可以合并成一个
@node(start=..., listen=...)。但 CrewAI 选择拆开,因为装饰器是给人读的 DSL,语义清晰比代码 DRY 更重要。@start() 一眼就是"入口",@listen(x) 一眼就是"响应"。合并成一个通用装饰器会逼使用者每次都想"这个参数填哪个"。在"面向使用者的 API 层",可读性 > 去重;真正的去重发生在下面共享的 _merge_flow_method_definition / FlowMethod。L04
FlowMethod:包装类怎么"假装自己还是原函数"
三个装饰器都用 FlowMethod 作基类(flow/flow_wrappers.py:48):
# flow/flow_wrappers.py:56
def __init__(self, meth: Callable[P, R], instance: Any = None) -> None:
self._meth = meth
self._instance = instance
functools.update_wrapper(self, meth, updated=[]) # ① 复制 __name__/__doc__ 等
self.__name__ = FlowMethodName(self.__name__)
self.__signature__ = inspect.signature(meth) # ② 保留原签名
if instance is not None:
self.__self__ = instance
if inspect.iscoroutinefunction(meth): # ③ 异步方法也要能被认出
inspect.markcoroutinefunction(self)
# 把 @human_feedback / @persist 等其它装饰器留下的标记也一并搬过来
for attr in ["__human_feedback_config__", "__conversational_only__",
"__flow_persistence_config__", "__flow_method_definition__"]:
if hasattr(meth, attr):
setattr(self, attr, getattr(meth, attr))
# flow_wrappers.py:90
def __call__(self, *args, **kwargs) -> R:
if self._instance is not None:
return self._meth(self._instance, *args, **kwargs) # 绑定了实例:补上 self
return self._meth(*args, **kwargs)
update_wrapper把原函数的 __name__、__doc__、__module__ 拷到包装对象上——让它"看起来还是原函数",调试/日志里名字不乱。__signature__ 保留保住原始签名,IDE 补全和引擎的参数检查(listener 要不要收 result)才准。markcoroutinefunction★如果你的 Flow 方法是 async def,得让 iscoroutinefunction(wrapper) 仍返回 True,引擎才会 await 它。搬运其它装饰器标记★关键:@start 常和 @persist/@human_feedback 叠用。这段把它们留下的元数据一起搬过来,装饰器叠放顺序才不会丢信息。__call__ 补 self包装对象要能像方法一样被调用;绑定了实例就自动补上 self。大白话
FlowMethod 是一件"贴满标签的外套",套在你的函数上。它努力让外套穿起来和没穿一样(名字、签名、能调用都不变),只是多了一堆标签给引擎读。L05
描述符 __get__:让 self.my_method() 依然可用
🤔 痛点你的方法现在是一个
StartMethod 对象,不是函数了。可 Python 的"方法自动绑定 self"只对函数生效。那 self.fetch() 里的 self 从哪来?靠实现描述符协议 __get__(flow/flow_wrappers.py:112):
# flow/flow_wrappers.py:112
def __get__(self, instance: Any, owner: type | None = None) -> Self:
if instance is None:
return self # 从类上访问 → 返回未绑定的自己
bound = type(self)(self._meth, instance) # ★从实例访问 → 造一个"绑定了 instance"的新包装
skip = {"_meth", "_instance", "__name__", "__doc__", "__signature__", ...}
for attr, value in self.__dict__.items():
if attr not in skip:
setattr(bound, attr, value) # 把元数据(含 __flow_method_definition__)复制给绑定版
return bound
instance is None当你写 MyFlow.fetch(从类访问)时 instance 是 None,返回未绑定的包装本身——引擎扫描类结构时走这条。type(self)(self._meth, instance)★当你写 flow.fetch(从实例访问),造一个新的同类型包装,并把 instance 记进去。之后 __call__ 就能自动补 self。复制元数据到 bound绑定版也要带着那份 FlowMethodDefinition,否则从实例拿到的方法就"丢了工牌"。skip 集合避免覆盖已单独处理的核心属性。💡 描述符是什么一个对象只要实现了
__get__,放到类里当属性,Python 访问它时就会调 __get__ 而不是直接返回对象。普通函数正是靠这个机制实现"自动绑定 self"的。FlowMethod 手动实现 __get__,就是模仿函数的绑定行为,让被包装后的方法用起来和普通方法一模一样。图注:从类访问返回未绑定包装(供引擎读结构),从实例访问返回带 self 的绑定副本(供正常调用)——模仿函数的绑定行为。
L06
or_ / and_:把"任一 / 全部"写成一棵条件树
dsl/_conditions.py(flow/dsl/_conditions.py:22):
# flow/dsl/_conditions.py:22
def or_(*triggers: FlowTrigger) -> FlowCondition:
"""任一 trigger 触发就满足。"""
return _condition_tree(OR_CONDITION, triggers)
def and_(*triggers: FlowTrigger) -> FlowCondition:
"""所有 trigger 都触发才满足。"""
return _condition_tree(AND_CONDITION, triggers)
# :65
def _condition_tree(condition_type, triggers) -> FlowCondition:
return {
"type": condition_type, # "OR" 或 "AND"
"conditions": [_coerce_trigger(t) for t in triggers], # 每个 trigger 归一成"名字"或子树
}
方法引用会被 _coerce_trigger 归一成方法名字符串(flow/dsl/_conditions.py:32):
# flow/dsl/_conditions.py:32
def _trigger_name(value: Any) -> str | None:
if isinstance(value, str):
return value # 直接给字符串(方法名/路由标签)
name = getattr(value, "__name__", None)
if callable(value) and isinstance(name, str):
return name # 给的是方法引用 → 取它的 __name__
return None
*triggers 可变参and_("a", "b", "c") 接收任意多个触发源。返回一个 dict 树条件不是代码,是数据:{"type":"AND","conditions":["a","b"]}。数据才能序列化、才能存进 Definition。_coerce_trigger 归一不管你传字符串还是方法引用,都归一成"方法名";传嵌套的 or_/and_ 结果则原样保留 → 支持任意嵌套。_trigger_name 取 __name__★所以 @listen(fetch)(传方法本身)和 @listen("fetch")(传名字)等价——前者靠 __name__ 取到 "fetch"。图注:or_/and_ 可任意嵌套成树。D44/引擎用
_condition_satisfied 递归遍历这棵树,判断当前收到的事件集合是否满足条件。L07
合并定义 + do 动作:多个装饰器叠放怎么不打架
_merge_flow_method_definition 让多个装饰器往同一个方法上"累加"元数据(flow/dsl/_utils.py:109):
# flow/dsl/_utils.py:109
def _merge_flow_method_definition(wrapper, definition) -> None:
existing = _get_flow_method_definition(wrapper)
if existing is None:
_set_flow_method_definition(wrapper, definition) # 第一个装饰器:直接设
return
updates = { # 后续装饰器:只覆盖"显式设过的字段"
field_name: getattr(definition, field_name)
for field_name in definition.model_fields_set
}
_set_flow_method_definition(
wrapper, existing.model_copy(deep=True, update=updates))
# :89 do 动作只是记一个"import 路径"
def _method_action(method: Any) -> FlowActionDefinition:
return FlowCodeActionDefinition(ref=f"{method.__module__}:{method.__qualname__}")
existing is None方法上还没元数据 → 直接放。第一个执行的装饰器(离函数最近的)走这。model_fields_set★只取"这次装饰器显式设置过的字段"来更新。没设的字段保留旧值——@router 叠在 @listen 上时,router 只加 router=True,不会抹掉已有的 listen 条件。model_copy(update=...)用 Pydantic 的不可变式拷贝合并,干净、无副作用。_method_action = ref 字符串★"这方法怎么调"不存函数对象,而存 "模块:限定名" 字符串。这样 Definition 才能序列化成 YAML/JSON——用时再按 ref 反查函数。💡 设计取舍②:do 动作为什么存"字符串 ref"而不是函数本身?
朴素做法是把函数对象直接塞进定义。但那样 Definition 就不可序列化了——函数对象没法变成 JSON。CrewAI 存
module:qualname 字符串(如 my_pkg.flows:ResearchFlow.fetch),付出的代价是"用时要 import 反查",换来的是整个 Flow 契约能落盘、能可视化、能声明式重建(D43)。这正是 D41 说的"三层解耦"能成立的技术基础:写法层产出的东西必须是纯数据。L08
边界 + 今日小结
⚠️ 边界:@start 忘了括号 / listen 引用自己
两个常见坑:①
@start 写成不带括号——因为 start 是"装饰器工厂",必须 @start();不带括号会把你的函数当成 condition 参数传进去,行为错乱。② listener 监听自己:@listen("a") def a(self)。这会被 Definition 的校验器直接拦下——_validate_trigger_namespace 里有 if _condition_references(method.listen, method_name): raise ValueError("...listen must not reference itself")(flow/flow_definition.py:774)。为什么拦?方法监听自己会形成"我完成→触发我→我完成……"的死循环,框架宁可在构建契约时报错,也不让你在运行时无限空转。🧠 今天你应该能回答
- 装饰器改变函数行为了吗?它到底做了哪两件事?
@start和@listen的代码差别只在哪一个字段?FlowMethod用哪些手段"假装自己还是原函数"?- 方法被包成对象后,
self.fetch()的self靠什么补上? or_/and_返回的是代码还是数据?为什么?@listen(fetch)和@listen("fetch")为何等价?- do 动作为什么存字符串 ref 而不是函数对象?
✋ 10 分钟动手
P=lib/crewai/src/crewai/flow
sed -n '18,70p' $P/dsl/_start.py # @start
sed -n '18,57p' $P/dsl/_listen.py # @listen
sed -n '56,148p' $P/flow_wrappers.py # FlowMethod: __init__/__call__/__get__
sed -n '22,86p' $P/dsl/_conditions.py # or_/and_/条件归一
# 看装饰后方法的元数据
python -c "
from crewai.flow.flow import Flow, start, listen, or_
class D(Flow):
@start()
def a(self): ...
@listen(or_('a','b'))
def c(self): ...
print(type(D.__dict__['a']).__name__) # StartMethod
print(D.a.__flow_method_definition__) # start=True ...
print(D.c.__flow_method_definition__.listen) # {'or': ['a','b']}
"
明日预告 · Day 43:今天装饰器往方法上挂的那份
FlowMethodDefinition,以及它们汇总成的 FlowDefinition,就是"契约层"。明天读 flow_definition.py:这份契约的完整字段、校验规则,以及 build_flow_definition 怎么扫描整个 Flow 类把图还原出来——还有它怎么支持"不写 Python、纯 YAML 声明一个 Flow"。