Day 44 / 共 60 天 · 阶段7 Flow 事件驱动

@router:把 if/else 分支和循环塞进事件模型

前三天 @start/@listen 只能表达"A 完了做 B"这种无条件的边。真实流程要"通过就发布、不通过就打回重写"——这需要根据运行时数据选路。@router 就是干这个的:它的返回值本身变成一个事件名,谁 @listen 那个名字就被触发。今天读 @router 装饰器、emit 怎么从返回注解自动推断、引擎 _execute_listeners 里"先跑路由、把路由结果当新触发、直到没有路由"的循环,以及分支与循环在事件模型里到底怎么落地。

📍 你在 60 天里的位置(阶段7 Flow 事件驱动 · 共 8 天)
阶段6 记忆 D41 Flow 总览 D42 装饰器 D43 定义契约 D44 路由跳转 D45 状态持久化 D46 表达式 D47 对话式 D48 选型 阶段8 LLM
💡 先用一个类比兜住今天 @router 就像铁路的道岔(扳道员)。火车(流程)开到这,扳道员看一眼情况(读 state),喊一声"走左线!"或"走右线!"——喊出的那句话就是一个信号(事件名)。哪条线的信号灯(@listen("左线"))亮了,火车就往那开。如果扳道员喊的是"回上一站",火车就循环回去重跑。@router 不自己干活,它只负责"喊一个方向",方向即事件。
L01

痛点:if/else 怎么塞进"谁监听谁"的模型?

🤔 痛点事件模型里边是静态的:@listen("A") 写死了"等 A"。可分支是运行时才知道走哪条:研究质量高就发布、低就重写。你没法预先写死"等哪个方法"——因为走哪条取决于结果。怎么让"运行时的选择"接回"静态的监听"?
💡 一句话本质 @router 让方法的返回值动态地"变成一个事件名"return "approved" 就等于"广播了 approved 这个事件",于是 @listen("approved") 的方法被触发。分支 = router 根据 state 返回不同字符串;循环 = router 返回一个能触发 @start("retry") 的字符串,把流程勾回起点。静态的边 + 动态的事件名 = 完整的控制流。
大白话普通方法完成,广播的是"我(方法名)完成了"。router 完成,广播的是"我说的那句话"。所以监听 router 的人不该监听 router 的方法名,而该监听 router 可能吐出的那些"话"(emit 值)。
L02

@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 完成后我来扳道岔"。
@routercondition 可以是 None——这允许把 @router 叠在 @start()@listen() 上(一个方法既是起点/监听者,又用返回值路由)。_merge_flow_method_definition(D42)保证叠放时字段不互相覆盖。
数据结构:router 方法的定义 vs 普通 listener router 方法 MethodDef router = True ★ emit = ["approved", "revise"] ★ listen = "research" 返回值 → 事件名 普通 listener MethodDef router = False emit = None listen = "approved" 监听 router 吐的事件名
图注:router 靠 router=True 被引擎特殊对待——返回值当事件名;emit 声明可能的分支,供画图/校验。listener 只是订阅某个事件名。
L03

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_annotationget_type_hints 读函数的 -> ... 返回注解。
Literal["approved","revise"]★如果你写 def route(self) -> Literal["approved","revise"]:,框架自动推出 emit=["approved","revise"]。类型注解一举两得:既给 IDE 补全,又给框架静态图。
Enum 支持返回 Enum 类型也行,取每个成员的 value。运行时 router 返回 Enum 会被 .value 拆出字符串(见 L05)。
dict.fromkeys 去重保序去掉重复事件名但保持声明顺序——可视化时分支顺序稳定。
💡 设计取舍①:为什么优先从"返回类型注解"推 emit,而不是强制手写? 朴素做法是逼用户 @router(emit=["approved","revise"]) 手写一遍。但用户在函数签名里本来就该写 -> Literal["approved","revise"](这是好的类型习惯)。让框架从注解自动推,就消除了"注解和 emit 两处重复、还可能写不一致"的隐患——单一事实来源。emit= 参数保留给"返回类型是动态字符串、注解表达不了"的兜底情况。能从既有信息推导的,就不要求用户重复提供。
L04

引擎:路由先跑,且循环处理到没有路由为止

一个方法完成后,_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。
控制流:router 把返回值变成事件、再触发 listener research 完成 @router 扳道 读 state 决定 事件 "approved" 事件 "revise" publish rewrite rewrite 完成 → 再次触发 research/router → 循环重来(直到 approved)
图注:router 读 state 后返回 "approved"/"revise",各自触发不同 listener;revise 分支跑完可再触发 router,形成"打回重写"的循环。
L05

路由结果 + 原触发点,一起去找 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 拿到的是路由前那个方法的输出(或人在环结果),而非事件名字符串。
💡 为什么 router 顺序、listener 并行?router 决定"走哪条路",多个 router 若并行、结果互相影响,路径就不确定了——必须顺序、可预测。而 listener 只是"响应某事件干活",彼此独立,并行更快。这是"控制流串行、数据流并行"的经典切分。
L06

条件判定: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集齐即触发并清空,让下一轮循环能重新累积(配合循环流程)。
📝 例子:and_("translate","illustrate") 的汇合 @listen(and_(translate, illustrate)) def layout(self)。translate 先完成 → seen={translate},all 判定不满足,layout 不跑;illustrate 后完成 → seen={translate,illustrate},满足,layout 触发一次。"等两个都好了再排版"就这么实现,你没写一句同步代码。
L07

循环:路由回起点,让方法重跑

路由结果如果能触发一个 @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重新执行该起点,事件又扩散一轮——循环就转起来了。
⚠️ 边界:循环靠 max_method_calls 兜底,不然真会转到死 事件驱动的循环没有天然的"退出条件"——全靠 router 某次返回不同的事件名(如 reviseapproved)来跳出。万一你的 router 逻辑写错、永远返回 revise,流程就无限重写。D43 的 max_method_calls=100flow/flow_definition.py:220)是最后一道刹车:方法执行次数超限,引擎强制停。写带循环的 Flow,务必确保 router 的退出分支真的可达,别把退出条件写成永远不成立。
L08

取舍 + 今日小结

💡 设计取舍②:为什么用"返回值当事件名"而不是像别的框架那样返回"下一个节点对象"? LangGraph 的条件边是返回"下一个节点名",直接指定跳去哪。CrewAI 的 router 返回一个事件名(字符串),由"谁监听这个名字"来决定跳去哪。区别在于耦合方向:返回节点名 = router 知道下游是谁(紧耦合);返回事件名 = router 只管"喊个信号",下游自己订阅(松耦合,发布/订阅)。松耦合的好处是同一个路由事件可以被多个 listener 订阅、可以后加订阅者而不改 router;代价是"读代码时要跳着找谁监听了这个事件",不如直接指定直观。CrewAI 押注可扩展性,选了发布/订阅。

👶 小白:@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
"
明日预告 · Day 45:今天 router 反复读 self.state 决策。明天专讲 state 管理与持久化:状态怎么初始化(dict / Pydantic)、@persist 怎么在每个方法完成后落盘、SQLiteFlowPersistence 的存取、以及 kickoff 传 id 怎么"接着上次跑"。
← Day 43 定义契约 Day 45 · state 管理与持久化 →