图即工具:PipecatEngine 状态机(核心)
30 秒导读: 一张对话流程图(节点 = 对话阶段,边 = 阶段之间的跳转条件)怎么变成一通"会说话、会自己决定往哪走"的电 话?答案是一句话:把每条出边包装成一个 LLM 能调用的函数。LLM 每次生成时,看着"当前节点的 system prompt + 当前节点的出边函数列表",它调用哪个函数,就等于选择走哪条边。函数被调用时,引擎抽取变量、播报过渡语、切到新节点、换上新节点的 prompt 和工具,再触发下一轮生成。本章讲这套"图 → 工具 → 跳转"的运行时机制,以及它最精妙的时序细节。
本章聚焦 PipecatEngine 内部的运行时逻辑。相关分工:
- 图的静态结构(节点/边的数据模型与校验)在 01-workflow-model。
- 帧在处理器之间怎么流动(管线本身)在 02-voice-pipeline。
- 引擎的装配与生命周期(谁创建它、怎么接上管线)、以及引擎旁路能力在 04-call-orchestration。
- transition 函数之外的可扩展工具(自定义工具、MCP、知识库)在 05-extensibility-registries。
1. 这是什么(先建直觉)
一句话定义
PipecatEngine 是一台由 LLM 函数调用驱动的对话状态 机:它把工作流图跑起来,让语言模型在每个节点上"边说话边决定下一步去哪个节点"。
它要解决的问题
假设你要搭一个外呼电话机器人:先问候 → 确认身份 → 介绍产品 → 如果感兴趣就预约、不感兴趣就礼貌挂断。
传统做法是写一堆 if/else 状态机,手工判断用户说了什么、该跳哪。但用户的话是自然语言,判断条件千变万化,硬编码根本写不完。
Dograh 的思路是把"判断走哪条边"这件事外包给 LLM:
- 你在图里给每条边写一个条件描述(例:"用户表示有兴趣了解更多")。
- 引擎把这条边变成一个 LLM 可调用的函数,函数的描述就是那句条件。
- LLM 在对话中觉得"用户确实感兴趣了",就去调用那个函数——于是对话跳到下一个节点。
一句话类比
把它想成带对讲机的导游:导游(引擎)站在某个展厅(节点)里,手上有几张写着"去往 X 厅"的门卡(transition 函数)。游客(用户)说的话让导游判断该带去哪个厅,导游刷对应门卡,一行人就走到新展厅,墙上的讲解词(system prompt)和可用设施(工具)也随之全换成新厅的。
用起来什么样(一次跳转的直观样子)
[节点: 确认身份] system prompt: "你在确认对方是不是本人……"
可调用函数:
- identity_confirmed (描述: 用户确认了自己是本人)
- wrong_person (描述: 接电话的不是目标本人)
用户: "对,我就是张先生本人。"
LLM 决策: 调用 identity_confirmed()
│
▼
引擎: 抽取变量 → 播过渡语"好的,谢谢确认" → set_node(产品介绍节点)
→ 换上"产品介绍"的 prompt 和函数 → 触发下一轮 LLM 生成
[节点: 产品介绍] ……LLM 开始用新 prompt 说话
本节不谈代码细节。记住一个核心等式即可:LLM 调用哪个 transition 函数 = 对话走哪条边。
2. 顶层全景(引擎大概怎么转)
部件职责
| 部件 | 干什么 | 在哪 |
|---|---|---|
set_node | 进入一个节点的总入口:更新当前节点、发跳转事件、按类型分派 | pipecat_engine.py:549 |
_setup_llm_context | 把当前节点的出边、工具、prompt 全部装进 LLM | pipecat_engine.py:506 |
_create_transition_func | 为一条边生产出真正被 LLM 调用的 transition_func | pipecat_engine.py:224 |
compose_system_prompt_for_node | 拼系统提示(全局 prompt + 节点 prompt + 录音模式指令) | pipecat_engine_context_composer.py:49 |
compose_functions_for_node | 拼函数列表(知识库 + 自定义工具 + 出边 transition schema) | pipecat_engine_context_composer.py:86 |
end_call_with_reason | 结束节点收尾:排 EndFrame、落库 disposition/tags | pipecat_engine.py:720 |
queue_node_opening / get_node_greeting | 节点开场:播文本 TTS、预录音频,或触发首轮生成 | pipecat_engine.py:641 / :614 |
should_mute_user | 判定此刻要不要静音用户输入 | pipecat_engine.py:771 |
pipecat_engine_callbacks.py | 各种回调工厂(用户静默、超时、聚合纠错) | pipecat_engine_callbacks.py |
主线走一遍(高层)
一次"从当前节点跳到下一个节点"的完整流程,只有一条主干:
┌─────────────────────────────────────────┐
进入某节点 ──────► │ set_node(node_id) │
│ • _current_node = 新节点 │
│ • 发 node_transition 事件 │
│ • 按类型分派 ↓ │
└───────┬─────────────┬────────────┬───────┘
is_start │ is_end │ agent │
▼ ▼ ▼
_handle_start_node _handle_end_node _handle_agent_node
└─────────────┴────────────┘
│ 都调用
▼
┌─────────────────────────────────────────┐
│ _setup_llm_context(node) │
│ ① 为每条出边注册 transition 函数 │
│ ② 注册自定义工具 / 知识库函数 │
│ ③ compose_system_prompt_for_node │
│ ④ compose_functions_for_node │
│ ⑤ _update_llm_context(prompt, funcs) │
└───────────────────┬─────────────────────┘
▼
LLM 带着"新 prompt + 新函数列表"生成下一句
│
┌───────────────────┴─────────────────────┐
│ LLM 调用了某个 transition 函数 │
│ → transition_func 执行 → 又一次 set_node │ ← 回到顶部,循环
└─────────────────────────────────────────┘
一句话:set_node 装填节点 → LLM 调 transition 函数 → transition 函数再 set_node,如此循环,直到走进结束节点。
3. 核心机制一:图 → 工具(每条出边变成一个函数)
3.1 它要解决的小问题
图里的"边"是静态数据,LLM 看不见。要让 LLM 能"选边",必须把每条出边翻译成 LLM 世界里唯一能"主动触发外部动作"的东西——函数调用(tool call)。
3.2 函数名从哪来、描述从哪来
一条边 翻译成函数时,两个关键字段:
| LLM 看到的 | 来自边的什么 | 源码 |
|---|---|---|
函数名(name) | 边的 label 做小写化 + 非字母数字替成 _ | workflow_graph.py:53 Edge.get_function_name |
函数描述(description) | 边的 condition(你写的那句跳转条件) | workflow_graph.py:48 Edge.condition |
get_function_name 的真实实现很短:
# workflow_graph.py:53 Edge.get_function_name
def get_function_name(self):
return re.sub(r"[^a-z0-9]", "_", self.label.lower())
它把边的标签(如 Identity Confirmed)规范成合法函数名(identity_confirmed)。LLM 靠 condition 这句描述来判断"什么时候该调这个函数"——所以图里边的条件写得越清楚,LLM 选边越准。
3.3 装填时机:_setup_llm_context
每次进入节点,_setup_llm_context 遍历该节点的 out_edges,把每条出边注册成一个函数:
# pipecat_engine.py:514 _setup_llm_context 内
if not node.is_end:
for outgoing_edge in node.out_edges:
await self._register_transition_function_with_llm(
outgoing_edge.get_function_name(), # 函数名 ← 边 label
outgoing_edge.target, # 跳去哪个节点
outgoing_edge.transition_speech, # 过渡语(可选)
outgoing_edge.data.transition_speech_type,
outgoing_edge.data.transition_speech_recording_id,
)
注意 if not node.is_end:结束节点没有出边,自然不注册任何 transition 函数——这是对话能"停下来"的前提(见 §5)。
注册分两步:_register_transition_function_with_llm(pipecat_engine.py:325)先用 _create_transition_func 造出闭包函数,再 self.llm.register_function(name, transition_func)(:347)把它挂到 LLM 上。
同时,函数的schema(给 LLM 看的名字+描述)由 composer 单独拼(见 §6):
# pipecat_engine_context_composer.py:126 compose_functions_for_node 内
for outgoing_edge in node.out_edges:
function_schema = get_function_schema(
outgoing_edge.get_function_name(), outgoing_edge.condition
)
functions.append(function_schema)
这里体现了一个分工:register_function 挂的是"被调用时干什么"(行为),get_function_schema 生成的是"LLM 看到的样子"(声明)。名字必须两边一致(都来自 get_function_name),LLM 才能把"它看到的声明"对上"引擎挂的行为"。
4. 核心机制二:transition 函数的精妙时序(全章最巧的地方)
这是整章最值得细读的一段。_create_transition_func 返回的内层 transition_func(pipecat_engine.py:232)就是 LLM 真正调用的那个函数。它 被调用时,做了一串有严格先后顺序的事。
4.1 四步时序
LLM 调用 transition_func(args)
│
① 变量抽取 await _perform_variable_extraction_if_needed(_current_node)
│ 从"即将离开的旧节点"抽取变量(默认后台任务) pipecat_engine.py:242
│
② 排过渡语音 若边配了 transition_speech / 预录音频:
│ 标记 _queued_speech_mute_state="waiting" 并排队播放 :245-281
│
③ set_node await set_node(transition_to_node) :286
│ → _current_node 换成新节点
│ → _setup_llm_context 换上新 prompt + 新函数
│
④ 回调 + 交还 properties = FunctionCallResultProperties(
on_context_updated=on_context_updated) :307
await result_callback({"status":"done"}, properties=...) :314
关键点:②③ 都在"把函数调用结果交还给框架(④)"之前完成。 也就是说,当框架把这次函数调用的结果写进 context、准备触发下一轮 LLM 生成时,新节点的 system prompt 和函数列表早已换好——下一句话就会用新节点的身份来生成。
4.2 为什么需要 on_context_updated 回调
问题在于:LLM 的"下一轮生成"不是 transition_func 里同步触发的,而是由管线里的聚合器收到函数调用结果帧之后才触发的。引擎需要一个"结果已经写进 context 了"的精确信号,于是用了 FunctionCallResultProperties.on_context_updated:
# pipecat_engine.py:288 transition_func 内
async def on_context_updated() -> None:
# pipecat 框架会在"函数调用结果已写入 context 后"运行这个回调。
# 这样带着新 system prompt 去做 LLM 补全时,context 里已经有了函数调用结果。
# FIXME: 存在潜在竞态——见下。
if self._current_node.is_end:
await self.end_call_with_reason(EndTaskReason.USER_QUALIFIED.value)
它的实际用途很聚焦:如果这次跳转刚好跳进了结束节点,就在"结果落定"这一刻排 EndFrame 结束通话(见 §5)。把 EndFrame 的排队推迟到 on_context_updated,是为了不让结束帧抢在 context 更新之前。
4.3 源码里明写的竞态 FIXME
作者在 on_context_updated 上留了一条 诚实的 FIXME(pipecat_engine.py:294-297):
当我们从
UserContextAggregator带着FunctionCallResultFrame触发 LLM 补全,同时又end_call_with_reason排EndFrame/CancelFrame时,如果EndFrame先于ContextFrame到达 LLM 处理器,那次生成可能永远不会跑(而这有时正是想要的)。
换句话说:结束节点到底"跑不跑最后一次生成",取决于两个帧谁先到——这是一个已知的、未彻底消除的时序竞态。把它写进注释而不是假装不存在,正是这段代码值得学习的地方。
4.4 静音配合:过渡语播放期间闭麦
②里排过渡语时把 _queued_speech_mute_state 置为 "waiting",是为了在"机器人正在播过渡语"时别让用户插嘴打断。这个状态由 should_mute_user 消费(见 §7),形成一个小小的状态机:idle → waiting → playing → idle。
5. 核心机制三:结束节点如何把通话停下来
结束节点(is_end)是对话的终点。它的特殊之处有两层。
5.1 结束节点不注册出边函数
§3.3 已提到:_setup_llm_context 里 if not node.is_end 守卫住了 transition 函数注册。结束节点没有可调用的"下一步",LLM 无处可跳——这是从机制上保证对话能收束,而不是靠 prompt 说服 LLM 停下。
5.2 排 EndFrame,顺带落库
真正结束由 end_call_with_reason(pipecat_engine.py:720)完成。它做三件事:
① 幂等 + 闭麦。 用 _call_disposed 防重入,并 _mute_pipeline = True 彻底静音后续用户输入(:728-735)。
② 结束前补抽变量。 只要不是错误/语音信箱类结束,先 await 掉后台在飞的抽取任务,再同步抽一次当前节点变量,确保变量不因通话结束而丢失:
# pipecat_engine.py:741-747
await self._await_pending_extractions()
await self._perform_variable_extraction_if_needed(
self._current_node, run_in_background=False # 同步,等它抽完
)
③ 决定帧类型 + 落库 disposition/tags + 排帧。 正常结束用 EndFrame(排空管线后优雅结束),需要立刻掐断则用 CancelFrame;call disposition 优先取对话中抽出来的,否则退回到结束原因,并把它并入 call_tags:
# pipecat_engine.py:749-769
frame_to_push = (
CancelFrame(reason=reason) if abort_immediately else EndFrame(reason=reason)
)
call_disposition = self._gathered_context.get("call_disposition", "") or reason
self._gathered_context["call_disposition"] = call_disposition
self._gathered_context["mapped_call_disposition"] = call_disposition
# … 把 disposition 并入 call_tags …
await self.task.queue_frame(frame_to_push)
EndTaskReason(pipecat.utils.enums)是结束原因的枚举:走进结束节点是 USER_QUALIFIED(:302),超时是 CALL_DURATION_EXCEEDED,用户长时间静默是 USER_IDLE_MAX_DURATION_EXCEEDED(见 §8)。
6. 核心机制四:prompt 与工具的组装(composer)
每次进节点,_setup_llm_context 都会重新拼一遍"这个节点该用什么 system prompt、能调用哪些函数",委托给 pipecat_engine_context_composer.py。
6.1 system prompt 的三段拼接
compose_system_prompt_for_node(pipecat_engine_context_composer.py:49)把 system prompt 拼成最多三段:
[ 全局节点 prompt ] ← 若 workflow.global_node_id 存在且 node.add_global_prompt
+
[ 本节点 prompt ] ← format_prompt(node.prompt),已渲染模板变量
+
[ 录音响应模式指令 ] ← 仅当 has_recordings 且本节点 prompt 里含 "RECORDING_ID:"
用 "\n\n".join(parts) 拼接(:83)。前两段就是"全局背景 + 本阶段任务"。第三段是个巧妙设计:当工作流用到了预录音频,composer 会追加一段 RECORDING_RESPONSE_MODE_INSTRUCTIONS(:27),用两个 marker 字符约束 LLM 的输出格式:
| Marker | 含义 | 常量 |
|---|---|---|
▸ | 动态 TTS:后面跟要合成语音的文本 | TTS_MARKER (:21) |
● | 预录音频:后面跟 recording_id + 转写 | RECORDING_MARKER (:20) |
指令要求 LLM 每条回复必须以 ▸ 或 ● 开头、绝不混用——引擎据此决定"这句是现场 TTS 还是放录音"。这是一种把"选择输出通道"塞进 prompt、让 LLM 自己标注的轻量协议。
6.2 函数列表的三类来源
compose_functions_for_node(:86)把一个节点能调用的函数按顺序拼成一个 list:
functions = [
① 知识库检索函数 ← 若 node.document_uuids 非空 :107
② 自定义 / 内置工具 ← 若 node.tool_uuids(计算器、MCP…) :118
③ 各出边 transition ← 遍历 node.out_edges,每条一个 schema :126
]
①②的细节属于可扩展工具,展开在 05-extensibility-registries;③正是 §3 讲的"图 → 工具"。三类函数最终被 ToolsSchema 包起来交给 LLM:
# pipecat_engine.py:205 _update_llm_context
if functions:
tools_schema = ToolsSchema(standard_tools=functions)
self.context.set_tools(tools_schema)
# …
await self.llm._update_settings(LLMSettings(system_instruction=system_prompt))
注意 schema(声明)与 handler(行为)是两条平行的注册路径:§3.3 里 _setup_llm_context 负责 llm.register_function 挂 handler,这里的 composer 负责生成 schema。二者靠同一个函数名对齐。
7. 核心机制五:开场行为与静音判定
7.1 一个节点怎么"开口"
进入节点后,谁先说话?由 queue_node_opening(pipecat_engine.py:641)统一决定,返回三种结果之一:
| 返回值 | 含义 | 触发条件 |
|---|---|---|
"greeting" | 播了节点问候语(文本 TTS 或预录音频) |