数据截至 (上游 commit 689ca048bb0a)
一次通话的端到端编排与引擎旁路能力
30 秒导读: 前三章分别讲 了"图长什么样"(01)、"帧怎么在处理器间流"(02)、"图怎么变成状态机工具"(03)。这一章把它们串成一根线:一个 WebSocket / WebRTC / 电话连接进来,系统怎么一步步把它变成一通真实通话,跑完,再干净地收尾。核心是
run_pipeline.py里的_run_pipeline_impl—— 通话的"总装配车间"。
本章覆盖三件事:
- 入口:连接从哪几个口进来,进来后各自做什么(建 run、校配额、装 transport)。
- 编排主体:
_run_pipeline_impl如何用运行时快照装配一整套 services→引擎→管线→Worker,注册观测器,跑完,收尾。 - 旁路/带外能力:变量抽取、上下文压缩、知识库检索、pre-call fetch、语音信箱检测、录音路由 —— 这些"不在主音频管线里"的能力,是在这条编排线的哪一步接上去的。
不重复的部分:管线内部帧结构看 02-voice-pipeline;图→工具的状态机机制看 03-pipecat-engine;供应商注册表看 05-extensibility-registries。
1. 这是什么(零基础也能懂)
一句话定义: 这是 Dograh 的"通话主控程序" —— 从一个连接建立,到一通语音通话跑完并把录音、转写、用量落库,全程由它编排。
想象一家呼叫中心。接 线员(入口路由)接起电话,核对工号和额度;调度台(_run_pipeline_impl)把这通电话需要的所有设备(麦克风、喇叭、翻译、大脑)一次性配齐、接线,然后按下"开始";质检员(观测器)全程记录;通话结束,清洁工(收尾逻辑)关掉设备、封存录音、结账。
这一章讲的就是调度台 + 接线员 + 清洁工,不是"大脑怎么想"(那是引擎,03 章)。
它要解决的真问题: 一通语音通话涉及十几个组件(STT、TTS、LLM、VAD、转写聚合、录音、观测、集成会话、MCP……),而且入口五花八门(浏览器 WebRTC、七八家电话运营商、通用 WebSocket)。如果每个入口各写一套装配,代码会爆炸。Dograh 把入口收敛成薄薄的适配层,把装配逻辑全部收敛到一个 _run_pipeline_impl。
一句话直觉: 入口负责"把连接变成一个 transport 对象 + 一个 workflow_run_id",剩下的所有脏活累活都交给同一个编排函数。换入口不改编排,加供应商不改编排。
2. 顶层全景(它大概怎么转)
一通通话的生命周期,是"连接 → 编排 → 观测 → 收尾"四段:
┌──────────── 入口层(薄适配) ────────────┐
│ 浏览器 Web Call 通用 WebSocket 电话 │
│ webrtc_signaling agent_stream telephony/providers/* │
└───────┬──────────────┬─────────────────┬──┘
│ 建 workflow_run + 校配额 + 装 transport
└──────────────┼─────────────────┘
▼
┌────────────────────────────────────────┐
│ _run_pipeline_impl (编排主体) │
│ │
│ 1 取运行时快照定义(pinned definition) │
│ 2 解析 run_configs + effective model cfg │
│ 3 判定 is_realtime,造 services │
│ 4 造 PipecatEngine(注入图/上下文/回调) │
│ 5 装管线 → PipelineWorker(task) │
│ 6 engine.initialize()(系统提示+工具+MCP) │
│ 7 注册观测器(反馈/延迟/turn log/buffer) │
│ 8 run_pipeline_worker(task) ── 跑完一通 │
│ 9 finally: close_mcp_sessions + cleanup │
└────────────────────────────────────────┘
▲ ▲ ▲
带外能力在装配期接线(不进主音频帧流):
变量抽取 · 上下文压缩 · 知识库 · pre-call fetch
· 语音信箱检测 · 录音路由
部件一句话职责:
| 部件 | 干什么 | 在哪 |
|---|---|---|
| 入口路由 | 接连接、建 run、校配额、造 transport | routes/agent_stream.py、routes/webrtc_signaling.py、services/telephony/providers/* |
_run_pipeline_impl | 编排主体:装配→运行→收尾 | services/pipecat/run_pipeline.py:549 |
register_active_call | 通话计数,供发布时优雅排空(drain) | services/pipecat/active_calls.py |
run_pipeline_worker | 走 pipecat v1.3 的 WorkerRunner 生命周期 | services/pipecat/worker_runner.py:7 |
| 观测器 | 反馈事件、延迟、turn log 落 buffer/WS | services/pipecat/event_handlers.py、realtime_feedback_observer.py |
| 收尾 | 落库、封存录音/转写、enqueue 后处理 | event_handlers.py:245 on_pipeline_finished |
主线走一遍(高层): 连接进来 → 入口建 workflow_run 并校配额 → 入口造出 transport → 调 _run_pipeline_impl → 它按快照装好一整套 → engine.initialize() 备好系统提示和工具 → run_pipeline_worker 让音频真正开跑 → 用户挂断 / 到时 / 语音信箱触发结束 → 收尾落库 → finally 关 MCP。
3. 三个入口:连接怎么变成一通通话
所有入口最终都汇聚到 _run_pipeline_impl,但它们建 run、校配额、造 transport 的方式不同。
| 入口 | 路径 | 谁用 | 凭证从哪来 |
|---|---|---|---|
| 通用 agent-stream | routes/agent_stream.py | 任意外部方,凭证内联在查询串 | query string(含 provider 凭证) |
| WebRTC 信令 | routes/webrtc_signaling.py | 浏览器 Web Call(SmallWebRTC) | 登录用户 / embed session token |
| 电话 | services/telephony/providers/* | Twilio/Plivo/Telnyx/Vonage/Vobiz/ARI/Cloudonix | 运营商 config 行 |
3.1 通用入口 agent_stream.py:内联凭证
/agent-stream/{provider_name}/{workflow_uuid} 是一个"万能口":provider 名是 URL 路径段,provider 特有的通话元数据(主被叫号、凭证等)则从该 provider 自己的流协议里读,不需要在组织里预存 TelephonyConfigurationModel 行(agent_stream.py:1-11 的模块 docstring 点明了这层区别)。
它进来后按顺序做四件事(agent_stream.py:72-127):
- 建 workflow_run ——
db_client.create_workflow_run(...),把 provider、direction="inbound"存进initial_context(主被叫号等元数据由 provider 协议解析,agent_stream.py:72-92)。 - 设运行上下文 ——
set_current_run_id/set_current_org_id,让后续日志和 trace 带上 run/org(agent_stream.py:93-94)。 - 校配额 ——
authorize_workflow_run_start(...);has_quota为假就用错误信息关闭 WebSocket(agent_stream.py:96-113)。 - 派发 ——
provider_instance.handle_external_websocket(...),把这个 socket 交给注册表里对应 provider 的实现(agent_stream.py:116-127),后者最终会调run_pipeline_telephony。
注意边界: provider 名是必填路径段,未注册的 provider 直接以 1008 关闭(
agent_stream.py:48-53);OLD 版"无?provider=的裸音频分支"随路由改版一并消失。
3.2 浏览器 Web Call webrtc_signaling.py:SmallWebRTC 信令
浏览器打电话走 WebRTC。这个文件用 WebSocket 信令 + ICE trickling(边收集边发候选)取代 HTTP PATCH,因为多 worker 部署下本地 _pcs_map 无法共享(webrtc_signaling.py:1-15)。
关键动作在 _handle_offer(webrtc_signaling.py:515):
- 收到
offer后设 run/org 上下文,并先校配额authorize_workflow_run_start(webrtc_signaling.py:552-569)。 - 新建
SmallWebRTCConnection,用get_ice_servers(user_id=...)塞入按用户生成的时限 TURN 凭证(webrtc_signaling.py:669)。 - 注册 WS 反馈通道
register_ws_sender(workflow_run_id, ws_sender)—— 这就是后面观测器把实时反馈推给浏览器的那根管子(webrtc_signaling.py:683)。 - 后台起管线
asyncio.create_task(run_pipeline_smallwebrtc(...)),不阻塞信令(webrtc_signaling.py:708);随后把 answer 发回浏览器,ICE 候选另行 trickle。
还有一个 public/signaling/{session_token} 公开口(embed 嵌入用),多一层 token 校验 + 来源域校验 validate_origin,防止泄露的 token 从任意站点接入(webrtc_signaling.py:874-951)。
3.3 电话入口:经 telephony providers
七家运营商的 provider 各自处理完自己的信令握手后,统一调 run_pipeline_telephony(如 services/telephony/providers/twilio/provider.py:328、plivo/provider.py:350 等,共七家)。配额校验发生在更上游的运营商回调里;这里只负责把 socket 变成 transport。供应商如何插进注册表,见 05-extensibility-registries。
3.4 三个入口的会合点:register_active_call
三条路各有一个薄包装函数,职责相同:在任何异步 setup 之前先登记这通活跃通话,finally 里注销。
# services/pipecat/run_pipeline.py:165 run_pipeline_telephony(节选,真实源码)
register_active_call(workflow_run_id) # 先登记,再干活
try:
await _run_pipeline_telephony_impl(...) # 解析 run、算 is_realtime、造 transport
finally:
unregister_active_call(workflow_run_id) # 无论如何注销
为什么先登记? 注释点明:发布(deploy)时要优雅排空在跑的通话;必须让排空逻辑也看得见那些"还在解析 DB/config/transport 状态"的通话,所以登记要早于一切 async setup(
run_pipeline.py:268-270)。
三个包装(run_pipeline_telephony:165、run_pipeline_smallwebrtc:292、_run_pipeline:386)各自解析出 transport 后,都汇入下面的 _run_pipeline_impl。
4. 编排主体:_run_pipeline_impl 主流程
这是本章的心脏(run_pipeline.py:549)。它拿到 transport + workflow_run_id,把一通通话需要的一切从零装好。按代码顺序,它是这样一步步走的: