跳到主要内容

图即工具:PipecatEngine 状态机(核心)

30 秒导读: 一张对话流程图(节点 = 对话阶段,边 = 阶段之间的跳转条件)怎么变成一通"会说话、会自己决定往哪走"的电话?答案是一句话:把每条出边包装成一个 LLM 能调用的函数。LLM 每次生成时,看着"当前节点的 system prompt + 当前节点的出边函数列表",它调用哪个函数,就等于选择走哪条边。函数被调用时,引擎抽取变量、播报过渡语、切到新节点、换上新节点的 prompt 和工具,再触发下一轮生成。本章讲这套"图 → 工具 → 跳转"的运行时机制,以及它最精妙的时序细节。

本章聚焦 PipecatEngine 内部的运行时逻辑。相关分工:


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 全部装进 LLMpipecat_engine.py:506
_create_transition_func为一条边生产出真正被 LLM 调用的 transition_funcpipecat_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/tagspipecat_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_reasonEndFrame/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_contextif 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 或预录音频)节点配了 greeting 且 previous_node_id != node_id
"llm"排了一次初始 LLM 生成没问候语但 generate_if_no_greeting=True
"none"什么都没排以上都不满足

问候语的取法在 get_node_greeting(:614):优先预录音频(greeting_type == "audio" 且有 greeting_recording_id),否则用文本问候(node.greeting,经 _format_prompt 渲染变量)。get_start_greeting(:637)是它对起始节点的便捷封装。

文本问候走 TTS 时有个易错点,源码注释点破了:

# pipecat_engine.py:692
await self.task.queue_frame(
TTSSpeakFrame(greeting_value, append_to_context=True)
)
# append_to_context=True 让助手聚合器在 TTS 结束后把问候语写进 LLM context,
# 否则 LLM 首次生成时会"以为还没打过招呼",重复问候一次。

7.2 什么时候闭掉用户的麦

should_mute_user(pipecat_engine.py:771)是给 CallbackUserMuteStrategy 用的判定函数,每帧被问一次"现在要不要静音用户"。它先从帧里更新机器人说话状态,再按优先级判断:

收到 BotStartedSpeaking → _bot_is_speaking=True;若队列语音在 waiting → 转 playing
收到 BotStoppedSpeaking → _bot_is_speaking=False;队列语音回 idle

▼ 判定(命中即返回 True)
1. _mute_pipeline 为真(通话正在收尾) → 静音
2. _queued_speech_mute_state != "idle" → 静音(过渡语/工具语音正在排或播)
3. 机器人在说话 且 当前节点 allow_interrupt=False → 静音(该节点不许打断)
否则 → 不静音

第 2 条就是 §4.4 埋下的伏笔:过渡语一旦排队,用户就被闭麦,直到这句播完回到 idle。第 3 条把"能不能打断机器人"下放到每个节点allow_interrupt 配置——有的阶段(如念免责声明)不许打断,有的可以。

7.3 各种回调工厂

引擎把一批边缘行为的回调抽到了 pipecat_engine_callbacks.py,PipecatEngine 只留薄封装(pipecat_engine.py:807-830):

引擎方法造出的回调作用源码
create_user_idle_handlerUserIdleHandler用户静默时逐级升级:第 1 次追问"还在吗",第 2 次道别并结束通话pipecat_engine_callbacks.py:31
create_max_duration_callbackhandle_max_duration通话超过硬上限时 abort_immediately=True 立刻掐断:75
create_generation_started_callbackhandle_generation_started每轮生成开始时清空上轮的参考文本:93
create_aggregation_correction_callbackcorrect_aggregation用 LLM 原始文本纠正 TTS 供应商回传的、被打乱的对齐文本:104

其中聚合纠错(correct_corrupted_aggregation,:107)是个纯函数小算法:当 ElevenLabs 之类返回的转写与 LLM 实际生成文本"字母相同但空格/标点错位"时,用一个双指针扫描,以 LLM 参考文本为准补回结构字符——并且带了多重安全兜底(短于 10 字符、字母数不匹配就原样返回),避免"越纠越错"。


8. 巧妙之处(可带走的精华)

  • 把"选边"变成"选函数"。 整套引擎最核心的一招:图的边不进 prompt 里靠自然语言描述让 LLM"说"往哪走,而是变成结构化的函数调用,LLM"调"往哪走。选择因此可被框架可靠地捕获、并驱动确定性的节点切换。workflow_graph.py:53pipecat_engine_context_composer.py:126

  • 声明与行为分两路、靠函数名对齐。 schema(get_function_schema)和 handler(llm.register_function)分别在 composer 和 _setup_llm_context 里注册,名字都来自 Edge.get_function_name。改一处名字规则,两路自动一致。pipecat_engine.py:514:126

  • on_context_updated 卡住结束时机。 结束帧不在函数体里立刻排,而是推迟到"函数结果已写入 context"的回调里,尽量让最后一轮生成建立在完整 context 上。pipecat_engine.py:288-303

  • 诚实标注竞态,而非假装无懈可击。 FIXME(:294-297)明写了 EndFrameContextFrame 的到达顺序竞态。留痕比藏拙更利于后人维护。

  • 从机制上让对话能停:结束节点无出边函数。 不靠 prompt 求 LLM 别再跳,而是根本不给它可调的下一步。pipecat_engine.py:514

  • 过渡语期间的三态静音。 idle/waiting/playing 让"机器人正在播过渡语"与"用户能否插话"精确联动,避免过渡语被用户打断截断。pipecat_engine.py:114:785-796

  • 录音/TTS 模式用一个前导 marker 让 LLM 自标输出通道。 / 把"这句现场合成还是放录音"的决定权交给 LLM,引擎只按首字符分流。pipecat_engine_context_composer.py:20-46


9. 边界与局限

  • on_context_updated 的竞态未彻底消除。 见 §4.3 的 FIXME:结束节点是否跑完最后一次生成,取决于 EndFrameContextFrame 的到达先后,代码承认这一点未被完全控制。

  • 选边质量 = LLM + 条件描述质量。 走哪条边完全由 LLM 依据边的 condition 判断。条件写得含糊、或多条边语义重叠,LLM 就可能选错——引擎不做二次仲裁。

  • 变量抽取默认后台异步。 _perform_variable_extraction_if_needed 默认 run_in_background=True(pipecat_engine.py:403),抽取在后台任务里跑。只有结束时才强制同步兜底(§5.2)。中途若抽取慢于跳转,新节点 prompt 里引用的变量可能尚未就绪。

  • 实时(speech-to-speech)LLM 不能做带外推理。 构造器注释指出:realtime 模式下管线 LLM 不实现 run_inference,变量抽取/上下文摘要必须靠单独传入的 inference_llm(pipecat_engine.py:84-89)。

  • MCP 会话的任务亲和性约束。 打开 MCP 会话的任务和关闭它的任务必须是同一个,否则 anyio 取消域会报"跨任务退出";因此清理逻辑刻意不放进 cleanup()(pipecat_engine.py:953-977)。这属于引擎生命周期,详见 04-call-orchestration


10. 代码地图(导航索引)

主题文件符号
进入节点总入口 / 类型分派api/services/workflow/pipecat_engine.py:549PipecatEngine.set_node
起始 / 结束 / 普通节点处理api/services/workflow/pipecat_engine.py:600 :710 :715_handle_start_node / _handle_end_node / _handle_agent_node
装填节点的 prompt+工具+出边函数api/services/workflow/pipecat_engine.py:506_setup_llm_context
造出被 LLM 调用的 transition 函数api/services/workflow/pipecat_engine.py:224_create_transition_func / transition_func
把 transition 函数挂到 LLMapi/services/workflow/pipecat_engine.py:325_register_transition_function_with_llm
结果落定回调 / 结束节点排 EndFrameapi/services/workflow/pipecat_engine.py:288on_context_updated
更新 LLM 的 system prompt + 工具api/services/workflow/pipecat_engine.py:205_update_llm_context
结束通话 / disposition 落库api/services/workflow/pipecat_engine.py:720end_call_with_reason
节点开场行为api/services/workflow/pipecat_engine.py:641 :614queue_node_opening / get_node_greeting
静音判定api/services/workflow/pipecat_engine.py:771should_mute_user
变量抽取调度api/services/workflow/pipecat_engine.py:402_perform_variable_extraction_if_needed
system prompt 拼接api/services/workflow/pipecat_engine_context_composer.py:49compose_system_prompt_for_node
函数列表拼接api/services/workflow/pipecat_engine_context_composer.py:86compose_functions_for_node
录音模式 marker / 指令api/services/workflow/pipecat_engine_context_composer.py:20RECORDING_MARKER / TTS_MARKER / RECORDING_RESPONSE_MODE_INSTRUCTIONS
用户静默逐级升级api/services/workflow/pipecat_engine_callbacks.py:31UserIdleHandler
最大时长 / 生成开始 / 聚合纠错回调api/services/workflow/pipecat_engine_callbacks.py:75 :93 :104create_max_duration_callback / create_generation_started_callback / create_aggregation_correction_callback
边 → 函数名 / 条件api/services/workflow/workflow_graph.py:53 :48Edge.get_function_name / Edge.condition