跳到主要内容

对话即图:工作流数据模型与校验

30 秒导读: Dograh 让你在一块画布上拖节点、拉连线,搭出一通电话该怎么聊。这一章只讲最浅的一层:你在画布上搭出来的东西,到了后端到底长什么样、要过哪几道校验才算「合法的工作流」。前端发来的一坨 ReactFlow JSON,先被 Pydantic 逐个节点 model_validate 成结构化的 ReactFlowDTO,再被翻译成一张真正的图 WorkflowGraph(邻接表)。图合不合法——起点有几个、Agent 节点必须有人连进来、End 不能再往外连——全部由每种节点自己的「规格(spec)」说了算,而不是散落在代码各处的 if。至于「怎么把这张图跑起来」,不在本章,留给 03-pipecat-engine.md


1. 这是什么(零基础也能懂)

一句话定义: 一个 Dograh「工作流」就是一张有向图——每个节点(Node)是对话里的一个步骤(打招呼、问预算、挂电话……),每条边(Edge)是「在什么条件下从这步走到下一步」。本章讲的,就是这张图在后端的静态数据模型和它必须满足的校验规则

解决什么问题 / 给谁用: 画布是给人用的,但机器不能直接信任画布传来的 JSON。用户可能连出一张没有起点的图、给 End 节点接了根出边、或者把两个「开始通话」拖进同一张画布。后端需要一个可信任的中间层:把杂乱的前端 JSON 收进来 → 严格校验 → 变成一个后续运行时敢直接用的对象。这一层就是本章的主角。

它管什么、不管什么:

管(本章)不管(别的章)
节点/边的数据结构长啥样图怎么被一帧帧执行 → 03
一张图合不合法(起点数、连接度)一通电话怎么端到端编排 → 04
每种节点类型的契约(spec)供应商/工具注册表 → 05
模板变量 {{first_name}} 怎么被抽出来变量在运行时怎么被填 → 03

一句话直觉: 把这一层想成编译器的前端——它做的是「词法/语法/类型检查」,产出一棵干净的语法树(这里是 WorkflowGraph);至于「解释执行这棵树」是后端(运行时)的事。

本节不出现底层代码。目标:知道「这一层是把画布 JSON 变成一张校验过的图」。


2. 顶层全景(它大概怎么转)

从「前端 JSON」到「可信任的图」,数据要过两道关,产出三种形态

怎么读这张图: 从上往下是数据流;左边是数据的三种形态,右边是每一步做的事。

前端 ReactFlow JSON │ 形态①:一坨 dict
{ nodes:[{id,type,data,position},…], │ (未校验,不可信)
edges:[{id,source,target,data},…] } │
│ │
▼ ReactFlowDTO.model_validate │ ── 第一道关:结构 / 类型 / 引用 ──
┌───────────────────────────────────────┐ │
│ RFNodeDTO 逐个节点:type 认识吗? │ │ · 未知 type → 报错
│ data 按该类型的模型校验 │ │ · 该类型缺 prompt → 报错
│ RFEdgeDTO label/condition 非空? │ │ · 边指向不存在的节点 → 报错
└───────────────────────────────────────┘ │
│ │
▼ │ 形态②:ReactFlowDTO
┌───────────────────────────────────────┐ │ (结构合法的强类型对象)
│ WorkflowGraph(dto) │ │
│ · 建邻接表:Node.out / out_edges │ │ ── 第二道关:业务不变量 ──
│ · _validate_graph(): │ │ · 实例基数(起点恰好 1 个…)
│ 实例基数 + 入/出度 + 节点配置 │ │ · 每类节点入/出度上下限
└───────────────────────────────────────┘ │ · (环:目前有意放行)
│ │
▼ │ 形态③:WorkflowGraph
可信任的图对象,交给运行时 │ (邻接表 + 起点/全局节点引用)

部件一句话职责:

部件干什么在哪
NodeType七种核心节点类型的机器名枚举services/workflow/dto.py:25
RFNodeDTO / RFEdgeDTO单个节点/边的入站校验(第一道关)dto.py:985 / dto.py:1024
ReactFlowDTO整张图的入站 DTO + 引用完整性dto.py:1031 ReactFlowDTO._referential_integrity
WorkflowGraph把 DTO 变成邻接表 + 业务校验(第二道关)services/workflow/workflow_graph.py:173
Node / Edge图里的节点/边对象workflow_graph.py:65 / :42
NodeSpec每种节点类型的契约(暴露给前端/MCP/SDK)services/workflow/node_specs/_base.py:252
spec 注册表register / get_spec / all_specsnode_specs/__init__.py:28

主线走一遍(高层): 前端 POST 一份工作流 JSON → ReactFlowDTO.model_validate(payload) 逐节点逐边过第一道关 → 拿这个 DTO 构造 WorkflowGraph(dto),构造函数里建邻接表并跑第二道关 → 任一关失败就抛错(附带一串 WorkflowError),全过则得到一个后续能直接用的图。


3. 核心机制(逐个看,由浅入深)

3.1 两道关:结构校验 vs 业务不变量

要解决的小问题: 「合法」其实有两种含义。一种是结构合法——JSON 字段齐不齐、类型对不对、边有没有指向鬼节点。另一种是语义合法——这张图作为一个工作流讲不讲得通(有没有起点、Agent 有没有入口)。Dograh 把这两种关注点拆到两层

第一道关在 DTO 层。 入口是 ReactFlowDTO,它的字段就是 List[RFNodeDTO]List[RFEdgeDTO](dto.py:1031)。Pydantic 在 model_validate 时会递归校验每个节点、每条边。整张图层面它只加一条规则:边的两端必须指向存在的节点

# services/workflow/dto.py:1035 ReactFlowDTO._referential_integrity
node_ids = {n.id for n in self.nodes}
for idx, edge in enumerate(self.edges):
for endpoint in (edge.source, edge.target):
if endpoint not in node_ids:
line_errors.append(dict(loc=("edges", idx), type="missing_node", ...))

这段在做:把「悬空的边」当成 Pydantic 校验错误抛出。注意它只查引用完整性,不查「起点有几个」——那是第二道关的事。

第二道关在图层。 WorkflowGraph 的类文档一句话点破了这个分工:

# services/workflow/workflow_graph.py:173
class WorkflowGraph:
"""*All* business invariants (acyclic, cardinality, etc.) are verified here.
The constructor accepts a validated ReactFlowDTO."""

也就是说:结构错误在 DTO 层就被挡下;一切「业务不变量」统一在 WorkflowGraph 里查。 构造函数先建邻接表,再调 self._validate_graph(...)(workflow_graph.py:211),校验不过直接 raise ValueError(errors)

关键细节:两层的报错形态不同。 第一道关抛的是 Pydantic 的 ValidationError;第二道关抛的是 ValueError,里面装一串 WorkflowError(services/workflow/errors.py:12,一个带 kind/id/field/messageTypedDict),这样前端能把错误精确标到具体节点/边上。


3.2 节点类型 = 逐类判别联合;data 按类型校验

要解决的小问题: 一个「开始通话」节点和一个「Webhook」节点,data 里该有的字段完全不同。怎么让每种类型只接受自己该有的字段、且缺了必填就报错?

思路: 每种节点类型配一个专属的 Pydantic 数据模型,type 字段决定用哪个模型去校验 data。核心的七种类型列在 NodeType 枚举(dto.py:25):

NodeType 成员机器名(字符串值)语义
startNode"startCall"起点:打招呼、开场,可选调外部 API 取上下文
agentNode"agentNode"中间对话步:LLM 跑一轮,可挂工具/文档
endNode"endCall"终点:收尾挂断,可在挂断前抽变量
globalNode"globalNode"全局人设:拼进每个节点的 prompt
trigger"trigger"对外 HTTP 触发入口
webhook"webhook"通话结束后回调外部系统
qa"qa"通话后跑 LLM 质检

真实实现: RFNodeDTOtype 是普通字符串,data 先声明成 Any,再由 model_validator 动态挑模型校验:

# services/workflow/dto.py:996 RFNodeDTO._validate
data_model = get_node_data_model(self.type) # 按 type 找到对应数据类
if data_model is None:
raise ValueError(f"Unknown node type: {self.type!r}")
self.data = data_model.model_validate(self.data) # 用它把 data 校验成强类型
prompt_label = _PROMPT_REQUIRED_NODE_TYPES.get(self.type)
if prompt_label:
_require_prompt(self.data, prompt_label) # start/agent/end/global 必须有 prompt

这里有三个巧妙点:

  • 类型是「动态判别」而非硬编码联合。 get_node_data_model(dto.py:1073)先查核心字典 _CORE_NODE_DATA_CLASSES,查不到再回落到集成包 get_integration_node_data_model。所以第三方集成节点无需改这里就能被认(详见 05)。
  • 未知类型早失败。 还有个 field_validator(dto.py:989)在校 type 时就先拒掉不认识的类型,报 Unknown node type
  • prompt 必填是「按类型」的。 只有 _PROMPT_REQUIRED_NODE_TYPES(dto.py:977,含 start/agent/end/global)四类强制非空 prompt;trigger/webhook/qa 不需要。

边的数据也有硬约束: EdgeDataDTO(dto.py:1016)的 labelcondition 都是 min_length=1——一条边至少要有名字和跳转条件,transition_speech(过场话术)才是可选的。


3.3 start/end 由「节点类型」定,而非持久化的 legacy 标志

要解决的小问题: data 里其实躺着 is_start / is_end 两个布尔字段(node_data.py:18,BaseNodeData 上,且 spec_exclude=True)。它们是 UI/运行时的历史遗留状态,可能过期。那到底谁说了算「这是不是起点」?

答案:节点类型说了算,持久化标志不作数。Node 时,起点/终点语义是node_type 现算的:

# services/workflow/workflow_graph.py:65 Node.__init__
# Start/end semantics are defined by node type. The persisted
# data flags are legacy UI/runtime state and may be stale.
self.is_start = node_type == NodeType.startNode.value # == "startCall"
self.is_end = node_type == NodeType.endNode.value # == "endCall"

注意这里比的是 NodeType.startNode.value(字符串 "startCall"),因为 Node.node_type 存的就是原始字符串(来自 dto.nodes[*].type)。

类型专属字段用 getattr 兜底读。 因为 Node 要能装下七种形态各异的 data,它对所有类型专属字段一律 getattr(data, "prompt", None) 这样读(workflow_graph.py:79-97)——某类型没这字段就拿默认值,不会炸。这让一个 Node 类就能通吃整个判别联合。

构造完顺手取两个关键引用:

# workflow_graph.py:214
self.start_node_id = [n.id for n in dto.nodes if n.type == NodeType.startNode.value][0]
try:
self.global_node_id = [n.id for n in dto.nodes if n.type == NodeType.globalNode.value][0]
except IndexError:
self.global_node_id = None # 全局节点可有可无

start_node_id 直接取 [0]——它敢这么写,是因为第二道关已经保证了起点恰好一个(见 3.5)。global_node_id 允许缺席,所以用 try/except 兜 None


3.4 spec 驱动的节点类型系统:暴露给前端/MCP/SDK 的契约

要解决的小问题: 「一个节点有哪些字段、哪个必填、下拉框有哪些选项、这类节点在图里最多几个」——这些信息前端要拿去渲染表单,MCP 工具要拿去让 LLM 拼工作流,SDK 要拿去生成代码。如果散落在各处,三边就会各写一份、迟早对不上。

思路:一个类型一份 spec,做单一数据源**。** 这份契约就是 NodeSpec(node_specs/_base.py:252),它把一种节点类型的全部元信息收在一起:

NodeSpec 字段含义
name / display_name机器名(= 判别值)/ 人看的名字
properties: list[PropertySpec]每个字段的类型、描述、必填、选项、校验边界
category: NodeCategoryUI 里归到哪一组(call_node/global_node/trigger/integration)
examples给 LLM 照抄的样例
graph_constraints图级约束(下一节的主角)

配套的几个受控词表也都在 _base.py:PropertyType(:25,渲染器 switch 的字段类型词表,如 string/boolean/options/mention_textarea)、NodeCategory(:53)、DisplayOptions(:62,条件显隐规则)、GraphConstraints(:239)。

spec 不是手写的,是从数据模型「长」出来的。 每个数据类头上挂一个 @node_spec(...) 装饰器,把元信息塞进 __node_spec_metadata__;build_spec 再遍历模型字段,把每个 spec_field(...) 的元数据编译成 PropertySpec:

# services/workflow/node_specs/model_spec.py:120 build_spec
for name, field in model_cls.model_fields.items():
prop = _build_property_spec(model_cls, name, field) # 逐字段 → PropertySpec
if prop is not None: # spec_exclude 的字段返回 None
properties.append(prop)
properties = _sort_properties(metadata.name, properties, metadata.property_order)

举个具体的:startCall 的数据类 StartCallNodeData(dto.py:319)头上的 @node_spec 就带了 graph_constraints=GraphConstraints(min_incoming=0, max_incoming=0, min_instances=1, max_instances=1)(dto.py:199)——「起点没有入边、全图恰好一个」这条规则,就写在这里。

注册表把所有 spec 收成一张表。 node_specs/__init__.py 提供 register(:28,重复名字会抛错)、get_spec(:40)、all_specs(:50)。核心 spec 是懒加载的:第一次访问时 _ensure_core_registered(:78)遍历 _CORE_NODE_DATA_CLASSES,对每个类跑 build_spec 注册进 REGISTRY。集成节点的 spec 则从集成注册表合并进来。

关键细节:MCP 投影是精简版。 NodeSpec.to_mcp_dict(_base.py:275)会把 display_name/icon/version 这类纯 UI 元信息、以及 PropertySpec 里的 placeholder/display_options/editor 全部丢掉,只留下 LLM 拼工作流真正需要的部分——省 token,也避免噪声。完整 spec 仍然原样发给前端渲染器和 SDK。


3.5 图级校验:完全由 spec 的 graph_constraints 驱动

要解决的小问题: 「起点恰好一个」「Agent 至少有一根入边」「End 不能再往外连」——这些规则怎么查,而且不散在代码各处?

思路:所有规则都从各节点 spec 的 graph_constraints 读,校验代码本身不认识任何具体类型。 _validate_graph 依次跑三个校验器(workflow_graph.py:280-287)。

校验器一:实例基数。 validate_node_instance_constraints(workflow_graph.py:120)把全图节点类型 Counter 一下,再遍历 all_specs(),按每个 spec 的 min_instances/max_instances 比:

# workflow_graph.py:138
count = counts.get(spec.name, 0)
if gc.max_instances is not None and count > gc.max_instances:
errors.append(WorkflowError(kind=ItemKind.workflow, ..., message="…at most one…"))
if enforce_min_instances and gc.min_instances is not None and count < gc.min_instances:
errors.append(WorkflowError(..., message="…must have at least one…"))

这就是「起点恰好 1 个」(startCall: min=max=1)、「全局节点至多 1 个」(globalNode: max=1)、「触发器至多 1 个」(trigger: max=1)的执行处。

校验器二:连接度数。 _assert_connection_counts(workflow_graph.py:306)先数每个节点的入度/出度,再按该节点 spec 的 min/max_incomingmin/max_outgoing 逐条查。没有 graph_constraints 的类型 = 不受约束(比如 agentNode 的出边侧就没设,可以随便连几条出去)。

各核心类型的图级约束(全部读自它们 @node_spec 里的 GraphConstraints):

类型入边出边实例数
startCall恰好 0不限恰好 1
agentNode至少 1不限不限
endCall至少 1恰好 0不限
globalNode恰好 0恰好 0至多 1
trigger恰好 0恰好 0至多 1
webhook恰好 0恰好 0不限
qa恰好 0恰好 0不限

(读法:「恰好 0」= min 与 max 都为 0;「不限」= 该方向没设约束。globalNode/trigger/webhook/qa 都是孤立节点——不参与主对话链的连线。)

校验器三:节点配置。 _assert_node_configs(workflow_graph.py:371)是留给「节点内部字段级」业务校验的钩子,当前基本是空壳(start 节点分支里只有 pass)。

关键细节:环被有意放行。 无环校验 _assert_acyclic(workflow_graph.py:291,标准白/灰/黑 DFS 找回边)代码写好了,但在 _validate_graph 里被注释掉:

# services/workflow/workflow_graph.py:270
# TODO: Figure out what kind of cyclic contraints can be applied, since there can be a cycle in the graph
# try:
# self._assert_acyclic()

也就是说 Dograh 允许工作流里有环——对话本就可能「没答上来就绕回去再问一遍」。这是一个刻意的设计取舍,不是漏写。


3.6 Edge:把 label 规约成合法函数名

要解决的小问题: 运行时会把每条出边暴露成一个「LLM 可调用的转移工具」(细节在 03)。工具得有个合法的函数名,但边的 label 是人随手写的中文/带空格/带标点的字符串。

真实实现: Edge.get_function_name(workflow_graph.py:53)一行搞定——小写化,再把所有非 a-z0-9 的字符替成下划线:

# services/workflow/workflow_graph.py:53
def get_function_name(self):
return re.sub(r"[^a-z0-9]", "_", self.label.lower())

举例:label "User says Yes!""user_says_yes_"

另一个细节:边的身份只看两端。 Edge.__eq__ / __hash__(workflow_graph.py:56-62)只用 (source, target)——同一对节点之间在集合语义里被当作同一条边。


3.7 模板变量抽取:{{first_name}} 从哪儿来

要解决的小问题: 节点 prompt、开场白、过场话术里会写 {{first_name}} 这种占位符。系统需要知道「这个工作流到底期望上游数据提供哪些变量」,好在跑之前校验数据齐不齐。注意:本章只讲抽取这一步(静态扫描),变量在运行时怎么被填是 03 的事。

思路:正则扫文本,但要滤掉三类「不算外部变量」的。 核心是 extract_template_variables(workflow_graph.py:19),用 TEMPLATE_VAR_PATTERN 匹配 {{ var | filter:value }},然后跳过:

  • 嵌套路径(名字里带 .,如 gathered_context.city)——那是运行时才解析的;
  • 带 fallback 过滤器的(如 {{name | Guest}})——它自带默认值,不算「必需」;
  • 系统注入变量 _SYSTEM_VARIABLES = {"campaign_id", "provider", "source_uuid"}(workflow_graph.py:16)——系统运行时自己会塞。
# services/workflow/workflow_graph.py:27
if "." in var_name: # 嵌套路径:运行时解析,跳过
continue
if filter_name is not None: # 带 fallback:有默认值,跳过
continue
if var_name in _SYSTEM_VARIABLES: # 系统注入,跳过
continue
variables.add(var_name)

图级汇总: WorkflowGraph.get_required_template_variables(workflow_graph.py:229)遍历全图,扫这几处文本并求并集:start/agent/end/global 四类节点的 prompt、start 节点的 greeting、以及所有边的 transition_speech。返回的就是「这份工作流真正要求上游提供的顶层变量名集合」。

小注:这里节点类型的判断写成 node.node_type in (NodeType.startNode, …),而 node.node_type 是字符串——之所以成立,是因为 NodeType 继承自 str,枚举成员与其字符串值相等。


4. 巧妙之处(可带走的技术)

  • 两段式校验分离关注点。 结构/类型/引用完整性交给 Pydantic(DTO 层),业务不变量集中到一个类(WorkflowGraph)。想加一条「图必须怎样」的规则,只有一个地方要动。
  • 规则即数据,不是代码。 「起点几个、谁能连谁」不写成 if 链,而是各节点 spec 里的 GraphConstraints;校验器对具体类型一无所知,只按约束比大小。加一种新节点类型,连校验逻辑都不用碰。
  • spec 是唯一数据源,还能分投影。 同一份 NodeSpec 既发完整版给前端渲染器,又用 to_mcp_dict 发精简版给 LLM——省 token 且不重复维护。
  • spec 从模型「长」出来。 spec_field + @node_spec + build_spec 让「字段声明」和「字段契约」是同一处真相,杜绝模型与 spec 漂移。
  • getattr 兜底 + str 枚举,让一个 Node 类通吃七种判别联合的数据形态,类型比较又能字符串/枚举混用不出错。
  • 有意允许环。 别把「无环」当成图校验的默认信仰——对话流天然会绕回。这里代码写好却注释掉,是记录了一个明确的产品取舍。

5. 边界与局限(诚实)

  • _assert_node_configs 目前近乎空壳(workflow_graph.py:371,start 分支只有 pass)——字段间的跨字段业务校验还没在图层落地。
  • 无环校验未启用:含环的图能通过校验,是否/如何终止由运行时保证,不在本层。
  • 本章不覆盖持久化与草稿保存路径。 sanitize_workflow_definition(dto.py:1085)是另一条更宽松的路径——它只剥字段、不跑 model_validator,让「半成品草稿」也能存下来;和本章讲的「严格校验成图」是两套不同入口。
  • 本章止于静态模型 + 校验。 图怎么被逐帧执行、变量怎么在运行时被填、边工具怎么触发跳转,全部在 03-pipecat-engine.md

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

主题文件符号
入站 DTO(整张图)services/workflow/dto.pyReactFlowDTOReactFlowDTO._referential_integrity
单节点校验 + 按类型校 dataservices/workflow/dto.pyRFNodeDTORFNodeDTO._validate_PROMPT_REQUIRED_NODE_TYPES
单边校验services/workflow/dto.pyRFEdgeDTOEdgeDataDTO
类型枚举 / 类型→数据类services/workflow/dto.pyNodeTypeget_node_data_model_CORE_NODE_DATA_CLASSES
各类型数据模型 + 图约束声明services/workflow/dto.pyStartCallNodeDataAgentNodeDataEndCallNodeDataGlobalNodeData
图对象构造 + 邻接表services/workflow/workflow_graph.pyWorkflowGraph.__init__
节点 / 边对象services/workflow/workflow_graph.pyNodeEdgeEdge.get_function_name
图级校验编排services/workflow/workflow_graph.pyWorkflowGraph._validate_graph
实例基数校验services/workflow/workflow_graph.pyvalidate_node_instance_constraints
连接度数校验services/workflow/workflow_graph.pyWorkflowGraph._assert_connection_counts
无环校验(未启用)services/workflow/workflow_graph.pyWorkflowGraph._assert_acyclic
模板变量抽取services/workflow/workflow_graph.pyextract_template_variablesget_required_template_variablesTEMPLATE_VAR_PATTERN
节点契约 / 受控词表services/workflow/node_specs/_base.pyNodeSpecPropertyTypeNodeCategoryDisplayOptionsGraphConstraints
spec 注册表services/workflow/node_specs/__init__.pyregisterget_specall_specs_ensure_core_registered
spec 从模型编译services/workflow/node_specs/model_spec.pynode_specspec_fieldbuild_spec
校验错误结构services/workflow/errors.pyWorkflowErrorItemKind