数据截至 (上游 commit 460c729002dc)
第 4 章 · 运行时底座
本章讲什么: 前三章反复出现的
RunContext.get()、.middleware(...)、ctx.emitter.emit(...)到底是什么。这一层是 BeeAI 所有可观测性、审批、流式输出的共同地基,值得单独拆开。
4.1 三个概念,先分清
| 概念 | 是什么 | 一句话 |
|---|---|---|
RunContext | 一次执行的上下文节点 | 带 run_id、abort 信号、自己的 emitter;父子成树 |
Run | run() 的返回值 | 一个惰性 awaitable:先挂钩子,再 await;也能当异步迭代器 |
Emitter | 事件总线节点 | 有命名空间,子节点向父节点冒泡 |
三者的关系:
Agent.run() ──▶ RunContext.enter(...) ──▶ 返回 Run
│
├─ 新建 RunContext(挂到父 context 下)
└─ RunContext.emitter = 实例 emitter 的子节点,并 pipe 到父 context 的 emitter
await run ──▶ ① 依次执行注册的钩子(middleware / on / context)
──▶ ② 真正跑 handler
4.2 Run:为什么 run() 不直接返回协程
看这个真实用法:
response = await agent.run("What to do in Boston?").middleware(GlobalTrajectoryMiddleware())
如果 run() 返回协程,.middleware() 就无处可挂。所以它返回一个实现了 Awaitable 的对象:
# python/beeai_framework/context.py:47-57 结构
class Run(Generic[R], Awaitable[R]):
def __init__(self, handler, context):
self.handler = ensure_async(handler)
self._tasks: list[tuple[Callable, list]] = [] # 待执行的钩子
self._run_context = context
self._events = Queue()
四个链式方法都只是往 _tasks 里塞条目(context.py:88-108):observe(fn)、on(matcher, cb)、context(dict)、middleware(*fns)。真正 await 时才一次性跑完再进 handler:
# python/beeai_framework/context.py:110-120
async def _run_tasks(self) -> R:
tasks = self._tasks[:]
self._tasks.clear()
for fn, params in tasks:
await ensure_async(fn)(*params)
try:
return await self.handler()
finally:
await self._events.put(None) # 关闭事件流
它还是个异步迭代器
构造时就订阅了自己 context 的所有事件、塞进队列(context.py:59-64),于是可以:
# 示意,非源码
async for data, meta in model.run(messages, tools=tools):
if meta.name == "new_token":
print(data.value.get_text_content(), end="") # 边跑边消费
__aiter__(context.py:69-86)把 handler 丢进 asyncio.create_task,自己从队列里取事件 yield,取到 None 收尾。LiteAgent 的主循环就是这么写的(agents/lite/agent.py:112-117)。
同一个对象,await 拿结果,async for 拿过程。
4.3 RunContext:树、追踪、取消
进入一次执行
RunContext.enter(context.py:187-273)做了四件事:
① 读 ContextVar 拿到父 context(有就挂树上)
② 新建 RunContext:run_id / parent_id / group_id / 继承的 context dict
③ 建子 emitter,带上 EventTrace(id=group_id, run_id, parent_run_id),并 pipe 到父 emitter
④ 返回 Run,其 handler 里:
发 start → 跑任务 → 发 success/error → finally 发 finish 并销毁 context
group_id 从父节点继承,run_id 每次新生成 —— 这样一整棵执行树共享一个 group_id,方便把一次用户请求的所有事件归到一起。
取消
每个 context 有自己的 AbortController,并把父的信号和用户传的信号一起注册(context.py:156-162)。跑的时候起两个任务赛跑:
# python/beeai_framework/context.py:235-258 摘要
runner_task = asyncio.create_task(_context_storage_run(), name="run-task")
abort_task = asyncio.create_task(_context_signal_aborted(), name="abort-task")
done, pending = await asyncio.wait([runner_task, abort_task], return_when=asyncio.FIRST_COMPLETED)
谁先完成谁说了算:业务先完成就取消 abort 任务;abort 先触发就取消业务任务并抛 AbortError。取消沿树往下传播,因为子 context 注册了父的信号。
上下文变量
storage: ContextVar["RunContext"] 在 _context_storage_run 里设置(context.py:216),所以任意深处都能 RunContext.get() 拿到当前上下文 —— 第 1 章 RequirementAgent.run 里那句 run_context=RunContext.get() 就是这么来的。
@runnable_entry:样板代码收进装饰器
# python/beeai_framework/runnable.py:120-129 摘要
return (RunContext.enter(self, inner,
signal=runnable_kwargs.get("signal"),
run_params={"input": args[1], **exclude_keys(kwargs, {"signal", "input"})})
.middleware(*self.middlewares)
.context(runnable_kwargs.get("context") or {}))
任何 Runnable 的 run 加上这个装饰器,就自动拥有:上下文树、abort 信号、实例级中间件、上下文字典。ChatModel、RequirementAgent、LiteAgent 都用它。
4.4 Emitter:分层事件总线
结构
Emitter.root() ← functools.cache 的全局单例
▲ pipe
┌────────┴────────┬──────────────┐
agent.requirement backend.ollama.chat tool.think
▲ pipe (各组件 _create_emitter 建的)
每次 run 的 context emitter
▲ pipe
run 内部事件 emitter(namespace 前缀 "run")
child()(emitter.py:88-109)建子节点时:命名空间前置拼接(namespace + self.namespace)、context 合并、然后立刻 pipe 回自己。pipe 的实现就是注册一个转发监听(emitter.py:111-121)。事件永远往上冒泡。
匹配器:五种写法
_create_matcher(emitter.py:197-234)支持:
| 写法 | 含义 | match_nested 默认 |
|---|---|---|
"*" | 本节点自己的事件 | False |
"*.*" | 一切事件 | True |
"success"(无点) | 本节点的具名事件 | False |
"agent.requirement.start"(有点) | 按完整路径 | True |
| 正则 / 任意函数 | 自定义 | 正则 True,函数 False |
match_nested=False 时会额外插一个守卫:
# python/beeai_framework/emitter/emitter.py:227-232
def match_same_run(event: EventMeta) -> bool:
return self.trace is None or (
self.trace.run_id == event.trace.run_id if event.trace is not None else False
)
matchers.insert(0, match_same_run)
这是「只听我这次运行的事件,不听子运行的」的实现方式 —— 靠 EventTrace.run_id 比对,而不是靠对象引用。
派发
# python/beeai_framework/emitter/emitter.py:267-277
async with asyncio.TaskGroup() as tg:
for listener in reversed(list(self._listeners)):
if not listener.match(event):
continue
if listener.options and listener.options.once:
self._listeners.remove(listener)
task = tg.create_task(run(listener))
if listener.options and listener.options.is_blocking:
_ = await task
两点:监听器按 priority 有序插入(bisect.insort_left,emitter.py:189-193);is_blocking=True 的监听器同步等待——审批需求要靠这个才能在工具执行前把人问完(第 2.6 节)。
4.5 中间件:能观察,更能干预
中间件的接口小到不能再小(context.py:36-44):要么是个 (ctx) -> None 的函数,要么是个有 bind(ctx) 方法的对象。
它的能力全部来自一个约定:RunContext.enter 会在真正执行前发一个内部 start 事件,并读回事件对象上的两个字段。
# python/beeai_framework/context.py:209-220 摘要
start_event = RunContextStartEvent(input=context.run_params, output=output)
await emitter.emit("start", start_event)
context.run_params = start_event.input # ← 监听者改了 input,会被采纳
...
if start_event.output is not None:
return start_event.output # ← 监听者填了 output,直接短路
else:
return await fn(context)
配合 runnable_entry 里那句 modified_input = ctx.run_params.get("input", ...)(runnable.py:107),得到两项能力:
| 能力 | 怎么做 | 现实用例 |
|---|---|---|
| 改写输入 | 在 start 监听里改 data.input | 注入上下文、脱敏、重写提示 |
| 短路执行 | 在 start 监听里设 data.output | 审批拒绝、缓存命中、mock 测试 |
这两件事对 agent、工具、模型调用一视同仁,因为它们都走同一个 RunContext.enter。
内部事件的识别方式
框架自己发的 run 事件带 context={"internal": True}(context.py:199-204),create_internal_event_matcher(emitter/utils.py:27-62)据此过滤,还能进一步按 parent_run_id、按实例、按事件名筛。第 2 章审批需求用的就是它。
4.6 案例一:GlobalTrajectoryMiddleware
把整棵执行树打印成缩进日志(middleware/trajectory.py:41)。
怎么做到跨层级? 它在自己绑定的 emitter 上注册两类监听(trajectory.py:95-132):
- 本层的
start/success/error/finish—— 打印。 - 嵌套的
start—— 发现是新的子 context,就递归地把自己绑到那个子 emitter 上:
# python/beeai_framework/middleware/trajectory.py:126-130 节选
async def handle_nested_event(data: Any, meta: EventMeta) -> None:
if meta.creator.emitter is not emitter:
await handle_top_level_event(data, meta)
self._bind_emitter(meta.creator.emitter) # 顺着树往下长
缩进怎么算? 维护一张 run_id -> TraceLevel 表(trajectory.py:144-159),每个 TraceLevel 有两个深度:
| 字段 | 含义 |
|---|---|
absolute | 距根的真实深度 |
relative | 只数被过滤器放行的层的深度 |
所以 GlobalTrajectoryMiddleware(included=[Tool]) 只打印工具、而且缩进是连续的,不会因为跳过了中间层出现一堆空档。
默认前缀表也贴心:{BaseAgent: "🤖 ", ChatModel: "💬 ", Tool: "🛠️ ", Requirement: "🔎 "}(trajectory.py:79)。输出目标可以是 stdout、任意有 write 的对象、框架 Logger,或者 False 丢弃(trajectory.py:327-340)。
还有个 emitter_priority 参数默认 -1(trajectory.py:84),注释写明用意:故意排在后面执行,以便看到其它中间件修改后的最终值。
4.7 案例二:StreamToolCallMiddleware
要解决的小问题
最终答案是通过 final_answer 工具调用返回的。工具调用参数是一整段 JSON, 按常规要等它完整才能解析——用户就得干等到最后一个 token。
思路
边流边用 json_repair 的流式稳定模式解析残缺 JSON,把 response 字段已经成型的部分持续吐出来。
# python/beeai_framework/middleware/stream_tool_call.py:97-108 摘要
parsed_args = parse_broken_json(args, fallback={}, stream_stable=True)
output_structured = self._target.input_schema.model_validate(parsed_args) # 校验不过就静默跳过
if output_structured and hasattr(output_structured, self._key):
output = getattr(output_structured, self._key) or ""
self._delta = output[len(self._buffer):] # 只发增量
self._buffer = output
stream_stable=True 是 json_repair 的关键开关:它保证「补全出来的结果不会随着后续 token 到来而回退」,于是增量计算 output[len(buffer):] 才安全。
Runner 把它接到自己的事件上(_runner.py:92-108),转发成 final_answer 事件带 delta。它还处理了非流式的退路:_handle_success 里如果一个 token 事件都没收到,就把完整响应当一个 chunk 补喂进去(stream_tool_call.py:123-125)。
4.8 错误模型
FrameworkError(errors.py:34)是所有错误的基类,带两个决策位:致命与可重试。
三套名字,别记混
同一个决策位在三个层面上写法不同,照抄错了会静默失效:
| 层面 | 写法 | 出处 |
|---|---|---|
| 构造关键字 | FrameworkError(msg, is_fatal=True, is_retryable=False) | errors.py:40-46 |
| 实例字段 | err.fatal / err.retryable | errors.py:52-53 |
| 判定入口 | FrameworkError.is_fatal(err) / FrameworkError.is_retryable(err) | 静态方法,errors.py:63-76 |
is_fatal / is_retryable 挂在类上、且是 staticmethod,不是实例属性。写 err.is_fatal 取到的是那个函数对象本身(恒为真),不是你要的布尔值——要判就写 FrameworkError.is_fatal(err),要读字段就写 err.fatal。
静态方法比读字段多干一件事:对非 FrameworkError 的普通异常也给默认答案 —— is_fatal 一律 False,is_retryable 只把 CancelledError 判成不可重试(errors.py:63-76)。
两个决策位分别被谁读
| 决策位 | 谁读它 | 判为 True 时 |
|---|---|---|
| 致命 | 工具层的 on_error(tools/tool.py:149-150) | 立刻上抛,不进重试 |
| 可重试 | Retryable(retryable.py:127、:138) | 才允许再试一次 |
explain() 与 ensure()
最实用的是 explain()(errors.py:107-118):沿着 cause 链逐层格式化,每层多缩进两格,还会把 context 字典序列化出来。第 1 章工具出错时喂给模型的就是这段文本 —— 同一套错误描述,人和模型共用。
FrameworkError.ensure(errors.py:120-134)负责把任意异常规范化,并特判 CancelledError → AbortError。
4.9 关键细节与坑
Emitter.root()是进程级单例(functools.cache,emitter.py:83-86)。在 root 上挂"*.*"监听会收到进程内所有组件的事件 —— 调试神器,生产慎用(多租户下会串台)。- 事件回调抛异常会打断执行。
_invoke把回调异常包成EmitterError上抛(emitter.py:255-265),阻塞式监听尤其危险。中间件里要自己吞异常。 context.destroy()在 finally 里必定执行(context.py:270),它会 abort 信号并清空监听。跨 run 复用 context 对象不可行。- 中间件绑定会先清掉上一次的注册。
GlobalTrajectoryMiddleware.bind开头就while self._cleanups: self._cleanups.pop(0)()(trajectory.py:135-136)——同一个中间件实例重复用于多次 run 是安全的,但不能并发用于两次 run。
4.10 代码地图
| 主题 | 文件路径 | 符号名 |
|---|---|---|
| 惰性 awaitable | python/beeai_framework/context.py | Run、Run.__aiter__、Run._run_tasks |
| 执行上下文 | python/beeai_framework/context.py | RunContext、RunContext.enter、RunContext.get |
| 短路 / 改写输入的载体 | python/beeai_framework/context.py | RunContextStartEvent |
| 统一入口装饰器 | python/beeai_framework/runnable.py | runnable_entry、Runnable |
| 事件总线 | python/beeai_framework/emitter/emitter.py | Emitter、Emitter.child、Emitter.pipe |
| 匹配器与运行隔离 | python/beeai_framework/emitter/emitter.py | _create_matcher、match_same_run |
| 事件元数据 | python/beeai_framework/emitter/emitter.py | EventMeta |
| 监听选项 | python/beeai_framework/emitter/types.py | EmitterOptions、EventTrace |
| 内部事件过滤 | python/beeai_framework/emitter/utils.py | create_internal_event_matcher |
| 轨迹日志中间件 | python/beeai_framework/middleware/trajectory.py | GlobalTrajectoryMiddleware、TraceLevel |
| 流式最终答案 | python/beeai_framework/middleware/stream_tool_call.py | StreamToolCallMiddleware |
| 错误基类与决策位 | python/beeai_framework/errors.py | FrameworkError、is_fatal、is_retryable、explain、ensure |
| 取消信号 | python/beeai_framework/utils/cancellation.py | AbortController、register_signals |