@router:把 if/else 分支和循环塞进事件模型
前三天 @start/@listen 只能表达"A 完了做 B"这种无条件的边。真实流程要"通过就发布、不通过就打回重写"——这需要根据运行时数据选路。@router 就是干这个的:它的返回值本身变成一个事件名,谁 @listen 那个名字就被触发。今天读 @router 装饰器、emit 怎么从返回注解自动推断、引擎 _execute_listeners 里"先跑路由、把路由结果当新触发、直到没有路由"的循环,以及分支与循环在事件模型里到底怎么落地。
@router 就像铁路的道岔(扳道员)。火车(流程)开到这,扳道员看一眼情况(读 state),喊一声"走左线!"或"走右线!"——喊出的那句话就是一个信号(事件名)。哪条线的信号灯(@listen("左线"))亮了,火车就往那开。如果扳道员喊的是"回上一站",火车就循环回去重跑。@router 不自己干活,它只负责"喊一个方向",方向即事件。痛点:if/else 怎么塞进"谁监听谁"的模型?
@listen("A") 写死了"等 A"。可分支是运行时才知道走哪条:研究质量高就发布、低就重写。你没法预先写死"等哪个方法"——因为走哪条取决于结果。怎么让"运行时的选择"接回"静态的监听"?@router 让方法的返回值动态地"变成一个事件名"。return "approved" 就等于"广播了 approved 这个事件",于是 @listen("approved") 的方法被触发。分支 = router 根据 state 返回不同字符串;循环 = router 返回一个能触发 @start("retry") 的字符串,把流程勾回起点。静态的边 + 动态的事件名 = 完整的控制流。@router 装饰器:多挂两个字段
dsl/_router.py 的核心(flow/dsl/_router.py:142):
# flow/dsl/_router.py:142
def decorator(func: Callable[P, R]) -> RouterMethod[P, R]:
wrapper = RouterMethod(func)
if emit is not None:
router_events = _normalize_router_emit(emit) # 显式给了 emit
else:
router_events = _get_router_return_events(func) or [] # 否则从返回注解推
method_definition_kwargs: dict[str, Any] = {
"do": _method_action(func),
"router": True, # ★标记为路由
"emit": router_events or None, # 可能吐出的事件名列表
}
if condition is not None:
method_definition_kwargs["listen"] = _to_definition_condition(condition)
_merge_flow_method_definition(wrapper, FlowMethodDefinition(**method_definition_kwargs))
return wrapper
RouterMethod第三种包装类型。引擎用它区分"这是路由方法",行为和 listener 不同。router=True★关键标记。FlowMethodDefinition.router 置真,引擎据此把它的返回值当事件名,而非当普通输出。emit声明"我可能吐出哪些事件名"。纯为静态可视化 + 校验服务——画图时能画出分支箭头,校验时能查"有没有人接这些事件"。不影响运行。condition → listenrouter 也可以有 listen 条件(router 也是"等谁触发的")。@router("research") 表示"research 完成后我来扳道岔"。@router 的 condition 可以是 None——这允许把 @router 叠在 @start() 或 @listen() 上(一个方法既是起点/监听者,又用返回值路由)。_merge_flow_method_definition(D42)保证叠放时字段不互相覆盖。router=True 被引擎特殊对待——返回值当事件名;emit 声明可能的分支,供画图/校验。listener 只是订阅某个事件名。emit:从返回类型注解自动推断分支
不显式写 emit 时,从函数的返回注解里抠(flow/dsl/_router.py:86):
# flow/dsl/_router.py:86
def _get_router_return_events(function: Any) -> list[str] | None:
values = _string_values_from_annotation(_return_annotation(function))
return list(dict.fromkeys(values)) if values else None # 去重、保序
# :45(节选)从 Literal / Enum / Union 注解里提取字符串常量
def _string_values_from_annotation(annotation: Any) -> list[str]:
if isinstance(annotation, type) and issubclass(annotation, Enum):
return [m.value for m in annotation if isinstance(m.value, str)] # 枚举 → 取值
origin = get_origin(annotation)
...
if origin is Literal or getattr(origin, "__name__", "") == "Literal":
return [arg for arg in args if isinstance(arg, str)] # Literal → 取字面量
...
_return_annotation用 get_type_hints 读函数的 -> ... 返回注解。Literal["approved","revise"]★如果你写 def route(self) -> Literal["approved","revise"]:,框架自动推出 emit=["approved","revise"]。类型注解一举两得:既给 IDE 补全,又给框架静态图。Enum 支持返回 Enum 类型也行,取每个成员的 value。运行时 router 返回 Enum 会被 .value 拆出字符串(见 L05)。dict.fromkeys 去重保序去掉重复事件名但保持声明顺序——可视化时分支顺序稳定。@router(emit=["approved","revise"]) 手写一遍。但用户在函数签名里本来就该写 -> Literal["approved","revise"](这是好的类型习惯)。让框架从注解自动推,就消除了"注解和 emit 两处重复、还可能写不一致"的隐患——单一事实来源。emit= 参数保留给"返回类型是动态字符串、注解表达不了"的兜底情况。能从既有信息推导的,就不要求用户重复提供。引擎:路由先跑,且循环处理到没有路由为止
一个方法完成后,_execute_listeners 先把所有 router 处理完(runtime/__init__.py:2751):
# runtime/__init__.py:2751
router_results = []
current_trigger = trigger_method
current_result = result
while True: # ★循环处理 router
routers_triggered = self._find_triggered_methods(
current_trigger, router_only=True) # 只找被当前触发点触发的 router
if not routers_triggered:
break # 没有 router 了 → 退出
for router_name in routers_triggered:
router_result, current_triggering_event_id = await self._execute_single_listener(
router_name, router_input, current_triggering_event_id) # 跑这个 router
if router_result is None:
current_trigger = FlowMethodName(""); continue
router_result = router_result.value if isinstance(router_result, enum.Enum) else router_result
router_result_event = FlowMethodName(str(router_result)) # ★返回值 → 事件名
router_results.append(router_result_event)
current_trigger = router_result_event # 路由结果又可能触发下一个 router
router_only=True★先只找 router。路由决定"接下来往哪走",必须在普通 listener 之前算清楚,否则不知道该触发谁。while True 循环router 可以串联:一个 router 的结果触发另一个 router。循环处理直到"当前触发点再也触发不了任何 router"。Enum → .value返回枚举成员就取它的 .value 字符串。和 L03 的 emit 推断对应上。router_result_event★把返回的字符串包成 FlowMethodName——它现在是一个事件名,将和原触发点一起去找普通 listener。路由结果 + 原触发点,一起去找 listener
router 处理完,把"原触发点 + 所有路由结果"合成触发列表,逐个找 listener(runtime/__init__.py:2803):
# runtime/__init__.py:2803
all_triggers = [trigger_method, *router_results] # 原触发点 + 路由吐出的事件
for idx, current_trigger in enumerate(all_triggers):
if current_trigger:
listeners_triggered = self._find_triggered_methods(
current_trigger, router_only=False) # ★这次只找普通 listener
if listeners_triggered:
listener_result = router_result_payloads.get(str(current_trigger), result)
...
tasks = [self._execute_single_listener(name, listener_result, ...)
for name in listeners_triggered]
await asyncio.gather(*tasks) # ★普通 listener 并行执行
all_triggers不光路由结果能触发 listener,原方法名本身也能——所以一个 router 方法完成,既广播"我(方法名)完成",也广播"我路由到的事件"。router_only=False这一轮找普通 listener(router 上一步已处理完)。两轮分开,职责清晰。asyncio.gather 并行★被同一事件触发的多个 listener 并行跑——这是 Flow"扇出"的性能来源。router 之所以顺序处理,正因为它决定路径不能乱;listener 无此约束故可并行。router_result_payloads传给 listener 的"结果"要对:监听路由事件的 listener 拿到的是路由前那个方法的输出(或人在环结果),而非事件名字符串。条件判定:and_/or_ 的事件到底怎么算满足
D42 的条件树,在这里被真正"求值"(runtime/__init__.py:172):
# runtime/__init__.py:172
def _condition_satisfied(condition: FlowDefinitionCondition, events: set[str]) -> bool:
if isinstance(condition, str):
return condition in events # 叶子:这个事件出现过吗?
operator, branches = _condition_branches(condition) # "and" / "or"
combine = all if operator == "and" else any
return combine(_condition_satisfied(branch, events) for branch in branches) # 递归
而"是否满足"要靠累积已见事件来判断(runtime/__init__.py:2860):
# runtime/__init__.py:2860
def _condition_met(self, condition, trigger_method, subscription_key) -> bool:
seen = self._pending_events.setdefault(subscription_key, set())
seen.add(str(trigger_method)) # ★记下"这个 listener 又收到了一个事件"
if not _condition_satisfied(condition, seen):
return False # 还没集齐 → 先不触发
del self._pending_events[subscription_key] # 集齐了 → 清空、准备下次
return True
_condition_satisfied 递归and → 所有子条件都满足(all);or → 任一满足(any)。叶子就是"事件名在已见集合里吗"。_pending_events 累积★and_("a","b") 的精髓:a 先完成时,seen={a},不满足,先记着;b 完成时 seen={a,b},满足,才触发。这就是"等多路汇合"。subscription_key每个 listener 一个键,各自独立累积——不会串味。满足后 del集齐即触发并清空,让下一轮循环能重新累积(配合循环流程)。@listen(and_(translate, illustrate)) def layout(self)。translate 先完成 → seen={translate},all 判定不满足,layout 不跑;illustrate 后完成 → seen={translate,illustrate},满足,layout 触发一次。"等两个都好了再排版"就这么实现,你没写一句同步代码。循环:路由回起点,让方法重跑
路由结果如果能触发一个 @start(条件),引擎会重新执行那个起点——这就是循环(runtime/__init__.py:2846):
# runtime/__init__.py:2846
if current_trigger in router_results:
for method_name in self._start_method_names():
if self._start_condition_triggered_by(method_name, current_trigger):
if method_name in self._completed_methods:
# ★循环重执行:临时清掉"恢复中"标志,让方法真的再跑一遍
was_resuming = self._is_execution_resuming
self._is_execution_resuming = False
await self._execute_start_method(method_name)
self._is_execution_resuming = was_resuming
else:
await self._execute_start_method(method_name)
current_trigger in router_results只有"路由吐出的事件"才有资格重新勾起起点——普通方法完成不会。这限制了循环只能被显式路由发起。_start_condition_triggered_by查有没有 @start("revise") 这样的条件起点被当前路由事件命中。命中就重跑它。已在 completed 里 → 特殊处理★方法第一次跑完会进 _completed_methods。循环要它再跑,就得临时关掉"恢复模式",否则引擎以为是断点恢复会跳过执行。_execute_start_method重新执行该起点,事件又扩散一轮——循环就转起来了。revise 变 approved)来跳出。万一你的 router 逻辑写错、永远返回 revise,流程就无限重写。D43 的 max_method_calls=100(flow/flow_definition.py:220)是最后一道刹车:方法执行次数超限,引擎强制停。写带循环的 Flow,务必确保 router 的退出分支真的可达,别把退出条件写成永远不成立。取舍 + 今日小结
👶 小白:@listen(my_router) 传 router 方法本身,能触发吗?
👨🏫 老师:一般不要这么写。你该监听 router 吐出的事件名(@listen("approved")),而不是 router 方法名。虽然 router 完成也会广播它的方法名(L05 的 all_triggers 含原触发点),但那表示"router 跑过了",不含"走了哪条分支"的信息。监听 emit 值,才拿得到分支语义。
🧠 今天你应该能回答
- 运行时才知道的分支,怎么接回静态的 @listen?router 的返回值扮演什么角色?
@router比@listen多挂了哪两个字段?emit 影响运行吗?- emit 为什么能从返回类型注解自动推?好处是什么?
- 引擎为什么"router 顺序处理、listener 并行执行"?
and_("a","b")的"等两个都到"是靠什么数据结构实现的?- 循环流程怎么实现?为什么循环重执行时要临时关掉"恢复模式"?靠什么防死循环?
- "返回事件名"和"返回节点名"两种路由风格,耦合上有何不同?
✋ 10 分钟动手
P=lib/crewai/src/crewai/flow
sed -n '97,165p' $P/dsl/_router.py # @router + emit 推断
sed -n '2751,2859p' $P/runtime/__init__.py # _execute_listeners: router 循环 + listener
sed -n '164,178p' $P/runtime/__init__.py # _condition_satisfied
# 跑一个带分支 + 循环的 Flow
python -c "
from typing import Literal
from crewai.flow.flow import Flow, start, listen, router
class Review(Flow):
tries = 0
@start('revise')
@start()
def research(self): self.tries+=1; return self.tries
@router(research)
def gate(self) -> Literal['approved','revise']:
return 'approved' if self.tries>=3 else 'revise'
@listen('approved')
def publish(self): return f'发布(第{self.tries}稿)'
print(Review().kickoff()) # 研究3次后 approved
"
self.state 决策。明天专讲 state 管理与持久化:状态怎么初始化(dict / Pydantic)、@persist 怎么在每个方法完成后落盘、SQLiteFlowPersistence 的存取、以及 kickoff 传 id 怎么"接着上次跑"。