跳到主要内容

AgentFlow V2 引擎:队列驱动的有状态解释器

30 秒导读: 这是 Flowise 的第二套执行引擎,只跑 type=AGENTFLOW 的图。它不像 经典引擎那样"把整张图拓扑排序、一次性构建成一条 LangChain 流水线再 run 一次", 而是像解释器一样,从起点节点开始、一个节点一个节点地执行:执行完一个节点,看它的输出决定下一步 把哪些节点推进队列。因为控制权始终握在引擎手里,所以循环、条件分支、人在环暂停/恢复、子流嵌套 这些"经典 DAG 表达不了"的东西,在这里都成了自然而然的事。


1. 为什么要新引擎(先讲清动机)

1.1 经典引擎的世界观:构建一次,终点跑一次

经典引擎的核心动作是拓扑构建:它把整张图按依赖排好序,从叶子往根依次 initNode,把每个节点实例化并用边把它们接成一条 LangChain 链,最后在终点节点上 run() 一次, 数据顺着这条预先接好的管道流到底。

这套世界观的隐含前提是:图是有向无环的(DAG),数据只朝一个方向流一次。 它非常适合 "检索 → 提示词拼装 → LLM → 输出"这类直线型 RAG 管道

1.2 DAG 表达不了的四件事

一旦 agent 编排的需求上来,DAG 就捉襟见肘。下面四类需求,经典引擎都很难自然表达:

需求为什么 DAG 难
循环(反复改写直到满意)边不能指回上游,否则不再是"无环",拓扑排序直接失败
条件分支(走 A 还是走 B)管道一旦接好就固定了,无法"运行到一半才决定不走某条边"
人在环(停下等人点"通过/拒绝")一次性 run 到底的模型里,没有"暂停、把状态存下来、以后再从这里续"的位置
长时运行 / 子流嵌套单次 run 的调用栈里,塞不下"跑另一个完整 agentflow 再回来"

1.3 V2 的答案:把"构建"换成"解释执行"

AgentFlow V2 换了世界观:引擎自己拿着一个待执行队列,一个节点一个节点地解释执行

  • 执行完一个节点,引擎当场看它的输出,动态决定下一步把哪些子节点推进队列(走哪条边)。
  • 想循环?把上游节点再推一次队列即可(用计数器防止无限循环)。
  • 想暂停等人?执行到人在环节点就把整个执行状态存进数据库并 return,人点了按钮再从存档续跑

一句话对照:

经典引擎: 图 ──拓扑构建──▶ 一条 LangChain 链 ──run 一次──▶ 结果
(管道预先焊死,数据流一遍)

V2 引擎: 图 ──▶ [起点入队] ──▶ 循环{ 出队一个节点 → 执行 → 看输出决定谁入队 } ──▶ 结果
(管道边跑边铺,控制权始终在引擎手里)

引擎主体全部在一个文件里:packages/server/src/utils/buildAgentflow.ts。下面逐层拆开。


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

整个引擎就是一个 while 循环 + 三个状态容器。先看这张图,再逐块讲。

executeAgentFlow() ← 入口(buildAgentflow.ts:1541)

┌───────────────────┴────────────────────┐
│ 初始化三个状态容器 (buildAgentflow.ts:1650-1652)│
│ nodeExecutionQueue : 待执行节点队列 │
│ waitingNodes(Map) : 还没凑齐输入的节点 │
│ loopCounts(Map) : 每个循环节点跑了几圈 │
└───────────────────┬────────────────────┘
│ 把起点(startAgentflow)入队

┌──────────────────────────────────────────┐
│ while(队列非空 && status=INPROGRESS) (2020)│
│ │
│ ① shift 出队一个节点 │
│ ② executeNode() 真正跑它 (1051) │
│ ③ 把这步结果 push 进 agentFlowExecutedData │
│ ④ processNodeOutputs() 决定下一步 (847) │
│ ├ 条件未命中的分支 → 剪掉 (799) │
│ ├ 子节点凑齐输入 → 入队 │
│ └ loop 节点 → 把目标节点再入队 │
│ ⑤ 人在环 / 等待工具确认 → shouldStop, break │
└──────────────────────────────────────────┘


updateExecution() 落库 + 组装响应 (2230)

怎么读这张图: 从上到下是一次调用的生命周期;中间那个方框是引擎的心脏——一个反复"出队→执行→ 决定下一步"的循环。左侧三个容器是它全程依赖的"记事本"。

各部件一句话职责:

部件干什么位置(buildAgentflow.ts
executeAgentFlow引擎入口:初始化、跑主循环、落库、组装响应:1541
nodeExecutionQueue待执行节点的 FIFO 队列,主循环从这里取活:1650
waitingNodes多入边节点的"等待室":还没凑齐上游输入就先待在这:1651
loopCounts每个 loop 节点已经绕了几圈,用来防死循环:1652
executeNode执行单个节点:解析变量 → 调节点 run() → 处理人在环/迭代:1051
processNodeOutputs看节点输出,把该走的子节点入队、该剪的分支剪掉:847
IAgentFlowRuntime跨节点流转的运行态(state / chatHistory / form / webhook):97
addExecution / updateExecution把每一步执行数据写进 Execution 表,支持回看与恢复:184 / :212

3. 三个状态容器(引擎的"记事本")

主循环之所以能表达复杂控制流,全靠这三个数据结构。它们在 executeAgentFlow 开头一起初始化 (buildAgentflow.ts:1650-1652)。

3.1 nodeExecutionQueue —— 待办队列

一个数组,元素是 INodeQueuebuildAgentflow.ts:75):

字段含义
nodeId该执行哪个节点
data合并后的上游输入(喂给节点)
inputs各上游来源分开的原始输入({sourceNodeId: result}

主循环每轮 shift() 取队首,跑完再把新就绪的子节点 push 进来。队列空了,执行就结束。

3.2 waitingNodes —— 等待室

不是每个节点被"点到"就能立刻跑。有多条入边的节点,得等所有该到的上游都到齐才行。 waitingNodes 是个 Map<nodeId, IWaitingNode>IWaitingNodebuildAgentflow.ts:67)记着:

  • expectedInputs:这个节点必须等到的上游集合。
  • receivedInputs:已经到了的上游 → 结果。
  • isConditional / conditionalGroups:如果某些入边来自条件分支,它们归成一组——同组 只要到一个就算数(因为条件分支本来就只会走其中一条)。

判断"能不能跑了"的逻辑在 hasReceivedRequiredInputsbuildAgentflow.ts:770):必需输入 全到齐、且每个条件组至少来一个,才返回 true

3.3 loopCounts —— 圈数计数器

Map<nodeId, number>,记每个 loopAgentflow 节点已经循环了多少圈。每次循环回跳前 +1, 到达上限就停——防止无限循环的兜底就靠它(详见 §5.3)。上限常量:

// buildAgentflow.ts:174 —— 允许用环境变量覆盖,默认 10 圈
const MAX_LOOP_COUNT = process.env.MAX_LOOP_COUNT ? parseInt(process.env.MAX_LOOP_COUNT) : 10

4. 执行主循环(executeAgentFlow

这是引擎的骨架。入口 buildAgentflow.ts:1541,主循环 buildAgentflow.ts:2020

4.1 起点从哪来

图先被 constructGraphs 拆成正向 graphnodeDependencies(每个节点的入边数), 再各求一份反向图(buildAgentflow.ts:1595-1596)。源码注释直接给了这两个结构的形态 (buildAgentflow.ts:1624-1642):

graph { nodeDependencies {
startAgentflow_0: [ 'conditionAgentflow_0' ], startAgentflow_0: 0, ← 入度0 = 起点
conditionAgentflow_0: [ 'llmAgentflow_0', conditionAgentflow_0: 1,
'llmAgentflow_1' ], llmAgentflow_0: 1,
llmAgentflow_0: [ 'llmAgentflow_2' ], llmAgentflow_1: 1,
llmAgentflow_1: [ 'llmAgentflow_2' ], llmAgentflow_2: 2 ← 两条入边
llmAgentflow_2: []
} }
  • graph[nodeId] 告诉引擎"跑完这个节点,候选下一步是谁"。
  • nodeDependencies 里入度为 0 的就是起点;getStartingNodeserver/src/utils/index.ts:213) 从中挑出起点,checkForMultipleStartNodes 保证非递归调用只有一个起点 (buildAgentflow.ts:1509)。起点入队在 buildAgentflow.ts:1904
  • llmAgentflow_2nodeDependencies=2 正是它要在 waitingNodes 里等两个上游的原因。

4.2 循环体一轮做什么

主循环体(buildAgentflow.ts:2020-2211)每一轮:

  1. 出队一个节点,找到它的 reactFlowNode;跳过便签节点(:2032-2036)。
  2. 执行 executeNode(...):2048)——真正调用节点实现,见 §5。
  3. 若返回 shouldStop(人在环 / 等待工具确认),status=STOPPEDbreak:2102)。
  4. 否则把这步结果 push 进 agentFlowExecutedData,标 FINISHED:2110)。
  5. 把运行态回写:节点返回的 state / chatHistory / form 就地更新到 agentflowRuntime:2127-2137,见 §6)。
  6. 分流 processNodeOutputs(...):2144)——决定谁入队、谁被剪、要不要循环回跳,见 §5。

有两道防线:MAX_ITERATIONS(默认 1000,:1912/2028)防止主循环本身跑飞;每轮开头查 abortController.signal.aborted:2041)以支持用户中止。


5. 单节点执行与分流

主循环的两个关键子程序:executeNode(跑一个节点)和 processNodeOutputs(决定下一步)。

5.1 executeNode —— 跑一个节点

buildAgentflow.ts:1051。核心步骤(都在 try 块内,出错就抛给主循环的 catch 统一处理):

  1. 动态加载节点实现:按 componentNodes[name].filePath import 节点类并 new:1113-1115)——所以每个节点跑的是 packages/components/nodes/agentflow/* 里各自的 run()
  2. 解析变量:调 resolveVariables(...):1169)把 {{ ... }} 引用替换成真值(见 §6.2)。
  3. 组装 runParamsnewNodeInstance.run(reactFlowNodeData, finalInput, runParams):1259)。
  4. 迭代节点特判:若是 iterationAgentflow 且输出带数组,就对每个元素递归调用 executeAgentFlow 跑子块(:1262-1406,见 §7.4)。
  5. 人在环特判:若是 humanInputAgentflow 且本轮没带 humanInput,就构造带 "通过/拒绝"按钮的 humanInputAction,把这步标 STOPPED 落进执行数据,返回 shouldStop:true:1409-1452)。agentAgentflow 在"用工具前要人批准"时同理(:1456-1500)。

5.2 processNodeOutputs —— 决定下一步

buildAgentflow.ts:847。这是"边跑边铺管道"的地方:

processNodeOutputs(当前节点, 结果):
childNodeIds = graph[nodeId] # 候选下一步
ignoreNodeIds = determineNodesToIgnore(...) # 条件没命中要剪的分支 (799)
for childId in childNodeIds:
if childId 在 ignore 里: continue # 剪枝
waitingNode = waitingNodes[childId] 或新建 # 进等待室 (885-891)
waitingNode.receivedInputs[nodeId] = 结果 # 记下"我这条边到了"
if hasReceivedRequiredInputs(waitingNode): # 该到的都到齐了?(897)
从等待室删除, 合并输入, 入队 # 就绪 → push 队列 (899-904)
else:
继续待在等待室 # 还差上游, 先等着

多入边节点的输入用 combineNodeInputsbuildAgentflow.ts:971)按来源合并成一个对象再喂下去。

5.3 determineNodesToIgnore —— 条件分支剪枝

buildAgentflow.ts:799条件节点如何"不走某条边"就靠它。

  • 只对决策节点conditionAgentflow / conditionAgentAgentflow / humanInputAgentflow)生效。
  • 节点输出里带一组 conditions,每个有 isFulfilled 布尔。没命中的条件对应的那条出边 (sourceHandle = ${nodeId}-output-${index})的目标节点,就被加进 ignoreNodeIds 剪掉 (:824-838)。
  • 兜底:如果一个条件都没命中,就把最后一条当 else 分支强制置为命中,保证至少走一条 (:818-822)——避免"全剪光、流程凭空断掉"。

5.4 循环回跳

还是在 processNodeOutputs 里(buildAgentflow.ts:918):如果刚跑完的是 loopAgentflow 且输出带 nodeID(要跳回的目标),就:

loopCount = loopCounts[nodeId] + 1
if loopCount < maxLoop: # 没到上限
loopCounts[nodeId] = loopCount
把 result.output.nodeID 再 push 进队列 # ← 关键:目标节点重新入队 = 循环
清掉 humanInput 防止被重复消费
else: # 到上限
流式发一条 fallbackMessage, 停止循环

这就是循环的全部秘密:不修改图、不加环,只是"把一个已经跑过的节点再推一次队列"。计数器守住上限。


6. 运行态:数据怎么在节点间流转

经典引擎里数据顺着焊死的管道流;V2 里节点之间不直接连数据——它们共享一本"运行态记事本"。

6.1 IAgentFlowRuntime

buildAgentflow.ts:97。四个字段,全程被主循环携带、被节点读写:

字段装什么谁写
state跨节点共享的键值状态(工作流的"全局变量")节点返回 state 时回写(:2127
chatHistory本次执行内累积的对话轮次节点返回 chatHistory 时追加(:2131
form表单输入(formInput 起点持久化在 session 内):2135:1734
webhook触发本次执行的 webhook 原始载荷:1770

它在 executeAgentFlow 开头初始化成空壳(:1655-1660),可能从上一次执行恢复(持久化 state、 表单值、webhook 重放,见 :1683-1757)。

6.2 resolveVariables 如何读运行态

buildAgentflow.ts:241。节点输入里写的 {{ ... }} 模板,跑之前由它替换成真值。它认得一整套前缀, 其中几个直接读运行态与执行历史:

模板解析成源码
{{ $flow.state.* }}flowConfig.state(即 agentflowRuntime.state:390
{{ $form.* }}当前/运行态表单值:303
{{ $webhook.body.* }}webhook 载荷:315
{{ $iteration.* }}当前迭代项(迭代子流内):356
{{ $loopCount }}最近 loop 节点当前圈数(读 loopCounts:338
{{ nodeId.output.path }}agentFlowExecutedData 里翻该节点最近一次输出:402

注意最后一条:引用上游节点的输出,是去执行历史里"倒着找最近一次匹配":411:438), 而不是顺着边拿——这正是解释执行模型下"节点间不直连数据"的取数方式。

6.3 webhook 引用的预解析

webhook 触发时,起点节点的 webhookDefaultInput 模板里的 {{ $webhook.* }}任何节点跑之前 就先被 resolveWebhookRefsbuildAgentflow.ts:110)解析一遍(:1763),保证起点 run() 和下游 看到的是同一个值。它还防原型链穿越:118 挡掉 __proto__/constructor/prototype)。


7. 特色能力(各节点如何借引擎实现控制流)

引擎本身只提供"队列 + 等待室 + 计数器 + 剪枝"这几样原语。真正的控制流语义,是各节点在 run()返回特定形状的输出、再由引擎解读出来的。下面这些节点都在 packages/components/nodes/agentflow/*

7.1 Condition / ConditionAgent —— 条件分支

  • ConditionCondition/Condition.ts:274):按用户配的 if/else 规则逐条判定,命中的置 isFulfilled:true;一条都没中就补一个 else 分支置真(:340-359)。输出 output.conditions 正是 §5.3 剪枝的依据。
  • ConditionAgentConditionAgent/ConditionAgent.ts:249):把"走哪条"交给 LLM 判断, 输出同样是一组带 isFulfilled 的 conditions,下游剪枝逻辑复用同一套。

7.2 Loop / Iteration —— 两种循环

节点循环方式关键输出
LoopLoop/Loop.ts:104回跳式:输出 nodeID(跳回哪个节点)+ maxLoopCount,引擎把该节点重新入队output.nodeID:143
IterationIteration/Iteration.ts:39遍历式:输出一个数组,引擎对每个元素递归跑子流output.iterationInput

Iteration 的递归执行在 executeNode 里(buildAgentflow.ts:1262-1406):为迭代块内的子节点 单独造一份 flowData,逐元素带 iterationContextexecuteAgentFlow(..., isRecursive:true), 再把子流结果和运行态合并回父流。

7.3 HumanInput —— 人在环暂停/恢复

HumanInput/HumanInput.ts:136。一个节点,两种角色,看有没有带 humanInput

第一次跑到(没带 humanInput): 人点了按钮再调一次(带 humanInput):
├ 流式发出问题描述 ├ run() 走 humanInput 分支 (HumanInput.ts:155)
├ engine 造 "通过/拒绝" action ├ 按 type 把 proceed/reject 之一置 isFulfilled
├ 标 STOPPED, 存进 Execution 表 ├ 引擎从上次 STOPPED 存档恢复 (1775-1873)
└ return shouldStop:true → 主循环 break └ 从该节点续跑, 剪掉没选的分支

暂停侧在引擎的 executeNodebuildAgentflow.ts:1409-1452);恢复侧在 executeAgentFlow 开头一大段 "找上次 STOPPED 节点、校验、恢复 state、从它续跑"(:1775-1873)。这正是"长时运行"能成立的地方—— 状态全在数据库里,跨请求。

7.4 ExecuteFlow —— 子流调用

ExecuteFlow/ExecuteFlow.ts:162。把另一个 agentflow当一步来调:向 /api/v1/prediction/{selectedFlowId} 发预测请求(:197),可选把子流结果并进本流 state:238)。它拦截"调用自己"防递归自噬(:189)。

7.5 DirectReply —— 直接回话

DirectReply/DirectReply.ts:40。最简单的节点:把配置好的文本直接流式发给前端 (streamTokenEvent:50)并作为 output.content 返回——用于"不经 LLM、直接回一句固定话"。


8. 执行持久化(回看与恢复的地基)

V2 把每一步都写进数据库的 Execution 表。这既支撑前端"执行回放",也是人在环恢复的存档。

  • 创建:非递归执行一开始就 addExecutionbuildAgentflow.ts:184),插一条 state=INPROGRESSexecutionData 为 JSON 化的 agentFlowExecutedData 的记录(:1899)。
  • 更新:出错、迭代进展、人在环暂停、整体收尾,都调 updateExecutionbuildAgentflow.ts:212) 覆盖 executionDatastate;置 STOPPED 时还盖上 stoppedDate:230-232)。收尾处按 执行数据里有没有 TERMINATED/ERROR/STOPPED 决定最终 state:2213-2233)。
  • agentFlowExecutedData 是什么:一个数组,每个节点跑完 push 一项,含 nodeIdnodeLabel、 节点结果 datapreviousNodeIdsstatus:2110)。它既是落库内容,也是 §6.2 里 {{ nodeId.output }} 取数的来源——执行历史本身就是数据总线。

递归子流(迭代块)不新建执行记录,而是复用 parentExecutionId 把数据并回父记录(:1360-1369:1874-1892)。


9. 边界与坑(诚实)

  • 单起点:非递归执行只允许一个起点节点,多起点直接抛错(buildAgentflow.ts:1519)。
  • 两道循环护栏MAX_LOOP_COUNT(默认 10,管 loop 回跳,:174)与 MAX_ITERATIONS (默认 1000,管主循环总轮数,:1912)。设计上宁可提前停也不放任跑飞。
  • 恢复态严格:只有上次是 STOPPED(或 STOPPED 后紧跟一个 ERROR)的执行才能续跑, 其它状态直接拒绝(:1783-1818)——防止从一个不一致的中间态胡乱恢复。
  • 恢复要求图没改:续跑时校验 startNodeId 仍存在于上次执行数据,否则报"可能是无效恢复或图被改过" (:1836-1843)。
  • 条件全不命中会强走最后一条:这是刻意的兜底(§5.3),但也意味着"最后一条分支"隐含承担了 else 语义,配图时要留意。

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

引擎主体(packages/server/src/utils/buildAgentflow.ts):

主题符号名
引擎入口 / 主循环executeAgentFlow1541 / 循环 2020
待办队列元素类型INodeQueue75
等待室元素类型IWaitingNode67
三容器初始化nodeExecutionQueue/waitingNodes/loopCounts1650-1652
循环上限常量MAX_LOOP_COUNT174
单节点执行executeNode1051
分流 / 下一步决策processNodeOutputs847
条件剪枝determineNodesToIgnore799
就绪判定hasReceivedRequiredInputs770
依赖分析setupNodeDependencies672
多入边输入合并combineNodeInputs971
运行态类型IAgentFlowRuntime97
变量解析resolveVariables241
webhook 引用预解析resolveWebhookRefs110
执行落库addExecution / updateExecution184 / 212

特色节点(packages/components/nodes/agentflow/):

节点符号 / 文件run
条件分支Condition/Condition.ts (conditionAgentflow)274
LLM 条件分支ConditionAgent/ConditionAgent.ts (conditionAgentAgentflow)249
回跳循环Loop/Loop.ts (loopAgentflow)104
遍历循环Iteration/Iteration.ts (iterationAgentflow)39
人在环HumanInput/HumanInput.ts (humanInputAgentflow)136
子流调用ExecuteFlow/ExecuteFlow.ts (executeFlowAgentflow)162
直接回话DirectReply/DirectReply.ts (directReplyAgentflow)40

同组其它章: Flowise 全景 · 数据模型与节点系统 · 经典执行引擎 · 请求生命周期与流式/队列 · 前端画布。节点接口通用形态见 01,SSE 流式与队列扩展见 04,本章不重复。