一次通话的端到端编排与引擎旁路能力
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:415 |
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:219 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/{workflow_uuid} 是一个"万能口":调用方把 provider 名、主被叫号、甚至 provider 凭证全塞在查询串里,不需要在组织里预存 TelephonyConfigurationModel 行(agent_stream.py:1-13 的模块 docstring 点明了这层区别)。
它进来后按顺序做四件事(agent_stream.py:80-125):
- 建 workflow_run ——
db_client.create_workflow_run(...),把 provider、主被叫号、direction="inbound"存进initial_context(agent_stream.py:73-93)。 - 设运行上下文 ——
set_current_run_id/set_current_org_id,让后续日志和 trace 带上 run/org(agent_stream.py:95-96)。 - 校配额 ——
authorize_workflow_run_start(...);has_quota为假就用错误信息关闭 WebSocket(agent_stream.py:98-110)。 - 派发 ——
provider_instance.handle_external_websocket(...),把这个 socket 交给注册表里对应 provider 的实现(agent_stream.py:116-125),后者最终会调run_pipeline_telephony。
注意边界: 没有
?provider=的"裸音频"分支目前未实现,直接以 1011 关闭(agent_stream.py:50-56)—— 代码里明写"reserved for a future protocol decision"。
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:396-537):
- 收到
offer后设 run/org 上下文,并先校配额authorize_workflow_run_start(webrtc_signaling.py:429-445)。 - 新建
SmallWebRTCConnection,用get_ice_servers(user_id=...)塞入按用户生成的时限 TURN 凭证(webrtc_signaling.py:477-485)。 - 注册 WS 反馈通道
register_ws_sender(workflow_run_id, ws_sender)—— 这就是后面观测器把实时反馈推给浏览器的那根管子(webrtc_signaling.py:491-495)。 - 后台起管线
asyncio.create_task(run_pipeline_smallwebrtc(...)),不阻塞信令(webrtc_signaling.py:513-522);随后把 answer 发回浏览器,ICE 候选另行 trickle。
还有一个 public/signaling/{session_token} 公开口(embed 嵌入用),多一层 token 校验 + 来源域校验 validate_origin,防止泄露的 token 从任意站点接入(webrtc_signaling.py:643-697)。
3.3 电话入口:经 telephony providers
七家运营商的 provider 各自处理完自己的信令握手后,统一调 run_pipeline_telephony(如 services/telephony/providers/twilio/provider.py:323、plivo/provider.py:346 等,共七家)。配额校验发生在更上游的运营商回调里;这里只负责把 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:176-178)。
三个包装(run_pipeline_telephony:165、run_pipeline_smallwebrtc:292、_run_pipeline:386)各自解析出 transport 后,都汇入下面的 _run_pipeline_impl。
4. 编排主体:_run_pipeline_impl 主流程
这是本章的心脏(run_pipeline.py:415)。它拿到 transport + workflow_run_id,把一通通话需要的一切从零装好。按代码顺序,它是这样一步步走的:
4.1 幂等闸门 + 运行时快照(最关键的设计)
先挡掉重复运行:workflow_run 已完成就抛 400(run_pipeline.py:442-443)。
然后取 pinned definition —— 运行时快照定义。这是整个编排最值得记住的一点:
# run_pipeline.py:456 真实源码
# Use the run's pinned definition for graph + configs (not the workflow's current)
run_definition = workflow_run.definition
run_workflow_json = run_definition.workflow_json
run_configs = run_definition.workflow_configurations or {}
为什么要快照? 用户随时在编辑器里改工作流图。如果一通正在跑的通话读"当前的图",用户一保存,通话中途图就变了 —— 灾难。所以每个 workflow_run 在创建时就钉死了当时那一版定义(workflow_json),运行期只认这份快照,不认最新图。这份 JSON 随后喂给 WorkflowGraph(ReactFlowDTO.model_validate(run_workflow_json)) 构图(run_pipeline.py:576-579),图模型细节见 01-workflow-model。
4.2 解析配置 + 有效模型配置
从 run_configs(即快照里的 workflow_configurations)里抠出运行参数(run_pipeline.py:462-486):
| 参数 | 默认 | 作用 |
|---|---|---|
max_call_duration | 300s | 最长通话时长 |
max_user_idle_timeout | 10.0s | 用户静默多久算 idle |
smart_turn_stop_secs | 2.0s | 智能轮次结束超时 |
turn_stop_strategy | transcription | 轮次结束检测策略 |
dictionary | 无 | STT 关键词增强(逗号分隔→keyterms) |
再算 effective AI model config:get_effective_ai_model_configuration_for_workflow(...) 把这一版工作流的 model_overrides 叠加到组织级用户配置上(run_pipeline.py:490-501)。入口层若已算过就直接复用(resolved_user_config),避免二次拉取。
4.3 判定 is_realtime,装配 services
is_realtime = user_config.is_realtime and user_config.realtime is not None(run_pipeline.py:517)—— 决定走"语音到语音"(OpenAI Realtime、Gemini Live)还是传统 STT→LLM→TTS 三段式。这个布尔值会贯穿后面几乎每个分支。
is_realtime = True is_realtime = False
llm = 实时语音服务 stt / tts / llm 三件套
stt = tts = None inference_llm = None
inference_llm = 独立文本 LLM ◄── 关键:实时服务不实现
(变量抽取/语音信箱用) run_inference,带外推理要另开一个
见 run_pipeline.py:520-544。这个 inference_llm(旁路文本 LLM)是后面变量抽取、上下文压缩、语音信箱检测的共同"侧信道大脑"—— 记住它,4.7 节全靠它。
装完 services 后,把本次真正用的 provider/model 盖章进 initial_context.runtime_configuration 并落库,供事后分析用(run_pipeline.py:546-574)。
4.4 造引擎:PipecatEngine
把图、上下文变量、各种回调注入,造出状态机引擎(run_pipeline.py:677-692):
# run_pipeline.py:677 真实源码(节选)
engine = PipecatEngine(
llm=llm,
inference_llm=inference_llm, # 旁路文本 LLM
workflow=workflow_graph, # 4.1 的快照图
call_context_vars=merged_call_context_vars,
node_transition_callback=node_transition_callback, # 节点切换→观测
has_recordings=has_recordings,
context_compaction_enabled=context_compaction_enabled,
...
)
引擎内部如何把图变成 LLM 工具、如何驱动节点转移,是 03-pipecat-engine 的主题。这里只需知道:编排负责"造好并喂料",引擎负责"想"。
4.5 装管线 + 造 Worker
按 is_realtime 走 build_realtime_pipeline 或 build_pipeline(run_pipeline.py:882-906),再 create_pipeline_task 造出 PipelineWorker(run_pipeline.py:909)。轮次策略(realtime / 外部信号 / 智能轮次 / 转写)在此按 STT 和配置分流(run_pipeline.py:725-780)。管线内部结构见 02-voice-pipeline。
集成运行时会话在此 attach(task)(run_pipeline.py:911-917),然后把 task 和 transport 输出回填给引擎(run_pipeline.py:920-921)。
4.6 初始化引擎:系统提示 + 工具 + MCP
# run_pipeline.py:923 真实源码
# Initialize the engine to set the initial context with System Prompt and Tools
await engine.initialize()
engine.initialize() 在这里做两件对收尾至关重要的事:装好起始节点的系统提示与工具、打开本次通话的持久 MCP 会话(pipecat_engine.py:193 _open_mcp_sessions)。记住"MCP 会话在 initialize() 里打开",4.9 节的清理谜题就靠它。