跳到主要内容

数据截至 (上游 commit 460c729002dc)

第 4 章 · 运行时底座

本章讲什么: 前三章反复出现的 RunContext.get().middleware(...)ctx.emitter.emit(...) 到底是什么。这一层是 BeeAI 所有可观测性、审批、流式输出的共同地基,值得单独拆开。


4.1 三个概念,先分清

概念是什么一句话
RunContext一次执行的上下文节点带 run_id、abort 信号、自己的 emitter;父子成树
Runrun()返回值一个惰性 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 {}))

任何 Runnablerun 加上这个装饰器,就自动拥有:上下文树、abort 信号、实例级中间件、上下文字典。ChatModelRequirementAgentLiteAgent 都用它。


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):

  1. 本层的 start/success/error/finish —— 打印。
  2. 嵌套的 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=Truejson_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.retryableerrors.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 代码地图

主题文件路径符号名
惰性 awaitablepython/beeai_framework/context.pyRunRun.__aiter__Run._run_tasks
执行上下文python/beeai_framework/context.pyRunContextRunContext.enterRunContext.get
短路 / 改写输入的载体python/beeai_framework/context.pyRunContextStartEvent
统一入口装饰器python/beeai_framework/runnable.pyrunnable_entryRunnable
事件总线python/beeai_framework/emitter/emitter.pyEmitterEmitter.childEmitter.pipe
匹配器与运行隔离python/beeai_framework/emitter/emitter.py_create_matchermatch_same_run
事件元数据python/beeai_framework/emitter/emitter.pyEventMeta
监听选项python/beeai_framework/emitter/types.pyEmitterOptionsEventTrace
内部事件过滤python/beeai_framework/emitter/utils.pycreate_internal_event_matcher
轨迹日志中间件python/beeai_framework/middleware/trajectory.pyGlobalTrajectoryMiddlewareTraceLevel
流式最终答案python/beeai_framework/middleware/stream_tool_call.pyStreamToolCallMiddleware
错误基类与决策位python/beeai_framework/errors.pyFrameworkErroris_fatalis_retryableexplainensure
取消信号python/beeai_framework/utils/cancellation.pyAbortControllerregister_signals