数据截至 (上游 commit 7803d562546a)
流水线引擎:配置驱动的责任链 + 生成器分叉执行
30 秒导读: LangBot 把「处理一条消息」拆成十来个可插拔的阶段(Stage),串成一条责任链。链的顺序、开关、参数全由数据库里的配置决定,不在代码里写死。这条链最巧妙的地方是:某个阶段(比如真正生成回复的那个)可以返回一个异步生成器而不是单个结果——每
yield一个分片,引擎就递归地把后面所有阶段再跑一遍,从而让"流式回复"的每一小段都能完整走完"包装 → 切长文 → 发送"的下游流程。本章讲透这套机制。
本章是 LangBot 阅读地图 的第 2 站。上一章 一条消息的旅程 把消息从平台事件送到了流水线门口(controller.py 最终调用 pipeline.run(query));本章从这一脚踏进流水线开始,讲引擎本身。至于链条中段那个"真正处理消息"的阶段内部(本地 Agent 循环、工具调用),留给 03 本地 Agent 循环。
1. 这是什么(零基础也能懂)
一句话定义: 流水线引擎是 LangBot 处理每条消息的中枢——它把"检查权限 → 过滤内容 → 限流 → 生成回复 → 包装 → 发送"这一长串动作,做成一条可以按配置增删、按配置排序的阶段链。
它要解决的问题。 一个聊天机器人从"收到消息"到"发出回复",中间要做很多互不相关的杂事:
- 这个群/这个人,允不允许我回复?
- 消息里有没有敏感词,要不要拦掉?
- 这个用户是不是发得太频繁了,要限流吗?
- 真正去问大模型、生成回复。
- 回复太长了要不要切成图片 / 分条发?
- 最后把回复发回平台。
如果把这些全塞进一个大函数,改一处就要动全身。LangBot 的做法是:每件事做成一个独立的 Stage 类,引擎按一份配置把它们串起来跑。想给某个机器人关掉"内容过滤"、调换"限流"和"预处理"的顺序,改数据库配置即可,不碰代码。
一句话直觉/类比。 把它想成一条工厂流水线:query(一次请求的所有上下文)是传送带上的工件,每个 Stage 是一个工位,工件依次经过每个工位被加工;某个工位(发消息回来的那个)比较特殊——它不是一次性放行,而是一小段一小段地放行,每放一段,后面的工位就得为这一段重新加工一遍。这就是本章的主角:生成器分叉。
本节不出现代码。记住三个词就行:阶段(Stage)、责任链、生成器分叉。
2. 顶层全景(它大概怎么转)
2.1 两层对象:管理器造链,运行时跑链
引擎由两个类分工,都在 src/langbot/pkg/pipeline/pipelinemgr.py:
| 对象 | 职责 | 关键符号 |
|---|---|---|
PipelineManager | 启动时从数据库读所有流水线配置,把每份配置物化成一条运行时链,缓存在内存 | PipelineManager.load_pipelines_from_db / load_pipeline |
RuntimePipeline | 一条已经装配好的链。每来一条消息就 run 一次,负责触发插件事件、埋监控点、遍历阶段、统一处理输出 | RuntimePipeline.run / process_query / _execute_from_stage |
StageInstContainer | 一个薄包装:把「阶段名字」和「阶段实例」绑一起,链就是一个它的 list | StageInstContainer |
一句话分工:Manager 负责"装配"(一次性,启动时),Runtime 负责"执行"(每条消息一次)。
2.2 一张图看清全景
怎么读这张图:上半区是启动时发生一次的"装配";下半区是每条消息发生一次的"执行"。装配的产物(阶段链)被执行阶段反复使用。
┌──────────────────────── 启动时装配(一次)──────────────────────┐
数据库 │ │
legacy_pipelines ──┼─► load_pipelines_from_db │
(每行=一条流水线)│ │ 对每一行 │
│ ▼ │
│ load_pipeline │
│ ① coerce_pipeline_config 把 JSON 里的字符串"3"转成 int 3 │
│ ② 按 entity.stages 顺序,用 stage_dict 逐个 new 出阶段实例 │
│ ③ 每个实例 .initialize(config) │
│ ④ 打包成 RuntimePipeline,存进 self.pipelines │
└──────────────────────────────┬───────────────────────────────────┘
│ 产物:一条 StageInstContainer 链
┌───────────────────────────────▼──── 每条消息执行(多次)─────────┐
一条消息 │ RuntimePipeline.run(query) │
(见 01 章) ─────┼─► 写入 pipeline_config / 绑定的插件&MCP / 监控元数据 │
│ │ │
│ ▼ process_query │
│ 触发 MessageReceived 插件事件 ── 被拦截? ──► 直接 return │
│ │ 否 │
│ ▼ │
│ _execute_from_stage(0, query) ◄── 责任链遍历 + 生成器分叉 │
│ │ │
│ ▼ 每个阶段的输出都过一遍 │
│ _check_output 统一处理 user_notice / error_notice / 流式回复│
└──────────────────────────────────────────────────────────────────┘
2.3 链长什么样:默认的 12 个阶段
一条流水线跑哪些阶段、什么顺序,由 LegacyPipeline.stages(数据库里的一个 JSON 数组)决定。新建流水线时用的默认顺序写在 src/langbot/pkg/api/http/service/pipeline.py:14 的 default_stage_order:
| # | 阶段名(注册名) | 白话职责 | 阶段家族目录 |
|---|---|---|---|
| 1 | GroupRespondRuleCheckStage | 群里这条消息符合"该回复"的规则吗(@我、关键词…) | resprule/ |
| 2 | BanSessionCheckStage | 这个会话被封禁了吗 | bansess/ |
| 3 | PreContentFilterStage | 入站内容过滤(敏感词等) | cntfilter/ |
| 4 | PreProcessor | 预处理:组装上下文、附带图片等 | preproc/ |
| 5 | ConversationMessageTruncator | 按长度截断历史会话消息 | msgtrun/ |
| 6 | RequireRateLimitOccupancy | 申请一个限流令牌(占坑) | ratelimit/ |
| 7 | MessageProcessor | 真正生成回复(命令 or 聊天)——分叉源头 | process/ |
| 8 | ReleaseRateLimitOccupancy | 释放限流令牌(退坑) | ratelimit/ |
| 9 | PostContentFilterStage | 出站内容过滤 | cntfilter/ |
| 10 | ResponseWrapper | 把回复包装成消息链 | wrapper/ |
| 11 | LongTextProcessStage | 太长就转成图片/分条 | longtext/ |
| 12 | SendResponseBackStage | 发回平台 | respback/ |
注意第 7 个 MessageProcessor——它是唯一会返回生成器的阶段(流式回复时逐片 yield),也就是下面要解剖的"分叉源头"。它内部怎么调命令处理器/聊天处理器,是 03 章 的内容,本章不展开。
3. 核心机制一:配置驱动的装配(Manager 怎么造链)
3.1 要解决的小问题
链要长什么样,不该写死在代码里——不同机器人要能有不同的阶段组合和参数,而这些都存在数据库的一个 JSON 列里。装配就是把"死的配置数据"变成"活的、能跑的对象链"。
3.2 三步物化:coerce → 实例化 → initialize
load_pipeline 就干三件事(src/langbot/pkg/pipeline/pipelinemgr.py:618-665):
第一步:类型矫正。 数据库 JSON 列里,"3" 可能是字符串。coerce_pipeline_config 按 YAML 模板里声明的类型,把配置里的值就地转成 int / float / bool(src/langbot/pkg/pipeline/config_coercion.py:53 的 coerce_pipeline_config)。这样后续阶段拿到的就是干净的类型,不用到处 int(...)。
第二步:按名字实例化。 按 pipeline_entity.stages(那个顺序数组)逐个查表、new 出实例:
# 真实源码 src/langbot/pkg/pipeline/pipelinemgr.py:439-441
stage_containers: list[StageInstContainer] = []
for stage_name in pipeline_entity.stages:
stage_containers.append(StageInstContainer(inst_name=stage_name, inst=self.stage_dict[stage_name](self.ap)))
self.stage_dict 是"注册名 → 阶段类"的字典(下面 3.3 讲它从哪来)。链的顺序完全由配置里 stages 数组的顺序决定——这就是"配置驱动"。
第三步:逐个初始化。 每个实例 await inst.initialize(config)(pipelinemgr.py:655-656),让阶段有机会预建自己需要的处理器(例如 MessageProcessor.initialize 会建命令/聊天两个 handler,见 process/process.py:24)。
三步完成后打包成 RuntimePipeline 存进 self.pipelines(pipelinemgr.py:666-668)。启动时 load_pipelines_from_db 对数据库每一行都走一遍这个流程(pipelinemgr.py:549-560)。
3.3 阶段注册机制:装饰器 + 导入副作用
上面 stage_dict 里的映射,是靠一套很轻的"自注册"凑齐的。三个零件在 src/langbot/pkg/pipeline/stage.py:
零件一:一个全局字典。
# 真 实源码 src/langbot/pkg/pipeline/stage.py:11
preregistered_stages: dict[str, type[PipelineStage]] = {}
零件二:一个装饰器,把类塞进字典。
# 真实源码 src/langbot/pkg/pipeline/stage.py:14-19
def stage_class(name: str) -> typing.Callable[[type[PipelineStage]], type[PipelineStage]]:
def decorator(cls: type[PipelineStage]) -> type[PipelineStage]:
preregistered_stages[name] = cls
return cls
return decorator
于是每个阶段类只要挂一行装饰器就自动登记,例如 @stage.stage_class('MessageProcessor')(process/process.py:10)、@stage.stage_class('BanSessionCheckStage')(bansess/bansess.py:7)。
零件三:导入副作用(关键的一环)。 装饰器要执行,前提是那个模块被 import 过。LangBot 不逐个手写 import,而是在 pipelinemgr.py 顶部把所有阶段家族包一次性"扫进来":
# 真实源码 src/langbot/pkg/pipeline/pipelinemgr.py:34-47
importutil.import_modules_in_pkgs(
[resprule, bansess, cntfilter, process, longtext,
respback, wrapper, preproc, ratelimit, msgtrun]
)
import_modules_in_pkg 遍历包目录下每个 .py 并 importlib.import_module(src/langbot/pkg/utils/importutil.py,import_dir)。导入的副作用就是执行装饰器,preregistered_stages 随之填满。等 PipelineManager.initialize 跑到 self.stage_dict = {...preregistered_stages...}(pipelinemgr.py:545)时,字典已经齐了。
一句话:装饰器负责登记,批量导入负责"触发登记",二者合起来就是一张自动生成的"注册名 → 类"跳转表。 加一个新阶段,只要写好类、挂上装饰器、放进某个已被扫描的包,就自动可用。
阶段家族目录一览(都在 src/langbot/pkg/pipeline/ 下):
| 目录 | 家族职责 |
|---|---|
resprule/ | 群消息响应规则判定 |
bansess/ | 会话封禁检查 |
cntfilter/ | 内容过滤(前置 Pre / 后置 Post 两个阶段) |
preproc/ | 预处理(上下文/多媒体组装) |
ratelimit/ | 限流(占用 Require / 释放 Release 成对) |
msgtrun/ | 会话消息截断 |
longtext/ | 长文本转图片/分条 |
respback/ | 回复发送 |
wrapper/ | 回复包装成消息链 |
process/ | 真正的消息处理(命令/聊天分派,生成器源头) |
4. 核心机制二:责任链遍历与统一输出(Runtime 怎么跑链)
4.1 run → process_query:跑之前先做的三件事
RuntimePipeline.run(pipelinemgr.py:182)是每条消息的入口。它在真正遍历阶段前,先把上下文塞进 query:
- 把这条流水线的
config、绑定的插件/MCP 列表写进query(pipelinemgr.py:183-188)——供下游阶段和插件读取。 - 准备监控元数据(机器人名、流水线名),埋点用(
pipelinemgr.py:190-210)。
然后进 process_query(pipelinemgr.py:362),它在遍历前还有一道关卡:触发 MessageReceived 插件事件。
# 真实源码(节选)src/langbot/pkg/pipeline/pipelinemgr.py:326-336
event_ctx = await self.ap.plugin_connector.emit_event(event_obj, bound_plugins)
if event_ctx.is_prevented_default():
# 插件把这条消息"吃掉了",整条流水线直接不跑
return
await self._execute_from_stage(0, query)
含义:插件有机会在阶段链开跑之前拦截整条消息(is_prevented_default() 为真就 return)。放行后才 _execute_from_stage(0, query)——从第 0 个阶段开始遍历。链跑完再补记监控的成功/响应埋点(pipelinemgr.py:444-446),异常则记错误(pipelinemgr.py:470-472),finally 里把 query 从缓存池删除(pipelinemgr.py:483-484)。
4.2 _check_output:所有阶段共用的一个"出口"
每个阶段产出的结果都不自己发消息,而是交给 RuntimePipeline._check_output 统一处理(pipelinemgr.py:214)。它看结果对象上几个字段,分别落地:
| 结果字段 | _check_output 的动作 |
|---|---|
user_notice | 有就发给用户;支持流式(见下)、群里可自动 @ 发送者 |
debug_notice / console_notice | 打到 debug / info 日志 |
error_notice | 打 error 日志,并把 query 标记为出错、上报监控系统 |
其中"流式 vs 整发"的分支很关键(pipelinemgr.py:229-242):如果平台适配器支持流式输出、且 query 已经有回复消息,就调 reply_message_chunk 发一个分片(带 is_final 标记);否则调 reply_message 整条发。这正是下一节"生成器分叉"要配合的地方——分叉产生一串分片,_check_output 负责把每个分片流式地推给用户。
5. 核心机制三:生成器分叉执行(本章的精华)
这是整个流水线引擎"最巧妙"的地方。_execute_from_stage 一个方法同时实现了责任链遍历和生成器分叉(pipelinemgr.py:286)。