跳到主要内容

数据截至 (上游 commit e741923f72c3)

调度引擎:就绪队列、边激活与分支剪枝

30 秒导读: 上一章把画布编译成了 DAG,但图不会自己跑。本章讲运行期的三件事——谁先跑(就绪队列)、谁能跑(边激活)、谁不该跑(分支剪枝)。最值得学的一点在第 6 节:一个分支没被选中,引擎不是简单地"跳过它",而是要把这条分支的整条下游链路逐级失活,下游的汇合节点才知道"别等了"。

引用约定: 本章所有源码路径写全路径(相对克隆根),行号锚定 frontmatter 里的 sourceCommit。同一段落内重复引用同一文件时用短名 + 行号(如 edge-manager.ts:306)。


1. 这章解决什么问题(零基础也能懂)

一张工作流图长这样:触发器 → 一个"判断"块 → 分两条路 → 两条路最后汇到一个"发 Slack"块。

图画好了,现在要它。跑图要回答三个问题:

问题白话本章对应机制
谁先跑有五个块都没有未满足的依赖,先跑哪个?能一起跑吗?就绪队列 readyQueue + 隐式并发
谁能跑一个块有三条入边,是三条都到齐才跑,还是到一条就跑?边激活 + isNodeReady
谁不该跑判断块选了「是」那条路,「否」那条路上的十个块怎么办?分支剪枝(级联失活)

第三个问题是本章的重头戏,也是最反直觉的。 直觉答案是"不选中就不入队,自然不会跑"——错。因为下游那个汇合块会一直等那条死掉的路:它数着"我还有 2 个上游没来",而其中一个永远不会来。

所以 Sim 的做法是:没被选中的边不是被忽略,而是被明确标记为"失活",并把这条失活沿着下游链路一路传播下去,直到碰到一个"还有别的活路可走"的节点为止。汇合节点因此能算出"我等的那条边已经死了,不用等了,可以跑"。

一句话直觉: 把边想成电线。选中的路通电,没选中的路要逐级拉闸;拉到某个还有别的电源的节点就停手。汇合节点判断自己能不能跑,看的不是"来了几个",而是"还有几条没断电的线还没到"。


2. 顶层全景:装配线与主循环

这节先看"大盘":谁把零件装起来、主循环长什么样。

2.1 五个零件由谁装配

入口是 DAGExecutor.execute(apps/sim/executor/execution/executor.ts:87),它干三件事:建图 → 造上下文 → 装配流水线,然后把控制权交给引擎。

DAGExecutor.execute(workflowId, triggerBlockId)

├─ ① 建图 DAGBuilder.build → DAG(节点 + 入边集合 + 出边表)
├─ ② 造上下文 createExecutionContext → ExecutionContext + ExecutionState
├─ ③ 装流水线 buildExecutionPipeline → ExecutionEngine
└─ ④ 跑 engine.run(triggerBlockId)

装配发生在 buildExecutionPipeline(executor.ts:346),一次性 new 出六个对象并串起来:

零件一句话职责文件 · 符号
VariableResolver把块参数里的 <blockName.field> 引用换成真值apps/sim/executor/variables/resolver.ts · VariableResolver
BlockExecutor找 handler、跑一个块、记日志、兜错apps/sim/executor/execution/block-executor.ts:78 · BlockExecutor
EdgeManager判定边激活/失活、算节点是否就绪apps/sim/executor/execution/edge-manager.ts:10 · EdgeManager
LoopOrchestrator / ParallelOrchestrator循环与并行的迭代语义(见 03)apps/sim/executor/orchestrators/loop.ts / parallel.ts
NodeExecutionOrchestrator分流:哨兵节点走子流程逻辑,普通块走 BlockExecutorapps/sim/executor/orchestrators/node.ts:51 · executeNode
ExecutionEngine主循环:队列、并发、取消、边传播apps/sim/executor/execution/engine.ts:31 · ExecutionEngine

装配顺序里有一处不能挪:edgeManager.restoreDeactivatedEdges(...)(executor.ts:372)必须在引擎创建之前执行——从快照恢复的那次执行,已经失活的边要先灌回 EdgeManager,否则恢复后的引擎会以为所有分支都还活着。快照恢复的完整故事见 05

createExecutionContext(executor.ts:386)则负责把快照里的一堆序列化结构还原成运行期的 Map/Set:blockStatesexecutedBlocksdecisions.routerdecisions.conditionloopExecutionsparallelExecutions。它同时构造 ExecutionState,这是唯一的块输出读写口(见 §7.4)。

2.2 主循环长什么样

ExecutionEngine.run(apps/sim/executor/execution/engine.ts:181)的骨架只有十几行,读一遍就懂:

run()
├─ initializeQueue(triggerBlockId) 放入起跑节点
├─ checkCancellationBackstop() 问 Redis:这次执行是不是已经被取消了

├─ while (hasWork()) 队列非空 或 还有在跑的
│ ├─ 若 已取消 / 已出错 / 提前停止 → break
│ └─ processQueue() 排空队列 → 全部起跑 → 等任意一个完成

├─ waitForAllExecutions() 未取消时,收尾等齐
└─ 出口三选一: 有暂停点 → paused
已取消 → cancelled
否则 → success

hasWork()(engine.ts:287)的定义是 readyQueue.length > 0 || executing.size > 0——只要还有在跑的节点就不算干完,因为在跑的节点完成后可能会往队列里塞新节点。

2.3 一个节点跑完之后发生什么(主线走一遍)

这是全章最该记住的一条链路,共七步(handleNodeCompletion,engine.ts:519):

  1. 是否已被 Response 块终结? 是则整段跳过——不写状态、不传播边(engine.ts:530-534)。
  2. 输出里带 _pauseMetadata? 是则登记暂停点,metadata.status = 'paused',不传播边(engine.ts:536-545)。这是暂停的钩子,细节见 05
  3. 交给编排器做子流程记账(循环计数、并行分支收集)。
  4. 是 Response 块? 锁定 finalOutput、置 stoppedEarlyFlag,不传播边——工作流到此结束(engine.ts:549-557)。
  5. 若这个节点没有出边(isFinalOutput),它的输出就是工作流的最终输出。
  6. 若命中 stopAfterBlockId(调试用的"跑到这里为止"),且不是循环/并行的中途,提前停(engine.ts:563-573)。
  7. 传播边:edgeManager.processOutgoingEdges(node, output, false) 返回一批"现在可以跑了"的节点,批量入队(engine.ts:575-577)。

第 7 步是本章 §5、§6 的全部内容。前六步都是"提前 return",可以理解成:只有正常完成的普通块才会点亮下游


3. 就绪队列:一个没有 worker 池的并发模型

这节讲"谁先跑、能不能一起跑"。

3.1 三个字段就是全部状态

字段类型作用
readyQueuestring[]已就绪、还没起跑的节点 ID,FIFO
executingSet<Promise<void>>正在跑的节点的 promise 句柄
queueLockPromise<void>一条 promise 链,当互斥锁用

三个字段声明在 apps/sim/executor/execution/engine.ts:32-34没有 worker 数组、没有并发上限参数、没有信号量——并发是"隐式"的:队列里有几个就同时起几个。

3.2 一轮 processQueue 干什么

processQueue():
┌── while 队列非空 ──────────────────────┐
│ 出队一个 nodeId │
│ executeNodeAsync(nodeId) ← 不 await │
│ trackExecution(promise) → 塞进 executing
└────────────────────────────────────────┘
↓ 队列排空了
await Promise.race([...executing, abortPromise])
↓ 任意一个跑完(或被取消)就返回
回到 run() 的 while,再来一轮

看清楚两点:

  • 循环体里不 await,所以队列里的 N 个节点是同时起跑的。
  • 排空后只 race(等任意一个完成),不是 all。因为最先完成的那个可能已经往队列里塞了新节点,越早回到主循环、越早起跑新节点。

真实源码 processQueue(engine.ts:484-498)、trackExecution(:227)、waitForAnyExecution(:241)。trackExecution 还顺手做了两件事:给 promise 挂 .catch 记录第一个错误(errorFlag + executionError,后来的错误不覆盖),并在 .finally 里把自己从 executing 里摘掉。

下面这段示意代码把"隐式并发"的骨架剥出来看:

// 示意,非源码 —— 无 worker 池的调度骨架
async function processQueue() {
while (readyQueue.length > 0) { // 排空队列
const nodeId = readyQueue.shift()
track(runNode(nodeId)) // 起跑,不等它
}
if (executing.size > 0) {
await Promise.race([...executing, abortPromise]) // 等最快的那个
}
}
// 重点看:队列宽度 = 并发度,没有任何地方限制同时在跑的数量

3.3 并发度 = 图的宽度(以及那个常被误读的 20)

引擎不限制同时在跑的节点数。一个扇出到 50 个块的节点,会同时发起 50 次块执行(可能就是 50 个 LLM 请求)。

全仓唯一一个和"并行"沾边的数字是 DEFAULTS.MAX_PARALLEL_BRANCHES: 20(apps/sim/executor/constants.ts:179)。它很容易被读成"并行块最多展开 20 个分支"——不是。它夹的是批大小,不是分支数:

谁算出来有没有上限
并行分支总数 totalBranchesresolveBranchCount(apps/sim/executor/orchestrators/parallel.ts:202-218)没有count 型直接取 config.count ?? 1(:191);collection 型直接取 items.length(:200),全程无 clamp
一批同时展开几个分支 batchSizeresolveBatchSize(parallel.ts:254-261)Math.max(1, Math.min(DEFAULTS.MAX_PARALLEL_BRANCHES, parsed))(:260)
引擎同时在跑的节点数没人算没有

仓库自带测试正好反证了这一点:parallel.test.tscount 设成 MAX_PARALLEL_BRANCHES + 10(apps/sim/executor/orchestrators/parallel.test.ts:271),却只断言 batchSizecurrentBatchSize 被夹到 20——30 个分支原样留着,只是分两批跑完。画布侧的 clampParallelBatchSize(apps/sim/stores/workflows/workflow/utils.ts:11)夹的同样是 batchSize。

所以准确的说法是「并行批大小上限 20」。分批推进的机制见 03 §6。

代价与收益:省掉了 worker 池的全部复杂度(队列水位、空闲 worker 唤醒、公平性),换来的是并发上限由用户画的图决定。这是一个明确的取舍,不是疏忽——见 §10。

3.4 withQueueLock:为什么"完成处理"必须串行

并发起跑没问题,并发处理"完成"会出事。设想两个块 A、B 同时跑完,都指向同一个汇合块 M:

A ──┐
├──▶ M A、B 几乎同时完成
B ──┘

processOutgoingEdges 里对 M 的处理是"删掉一条入边 → 数还剩几条活的入边 → 为 0 则宣布就绪"。这是典型的读-改-判序列。若 A、B 的完成处理交错执行,可能两边都读到"还剩 1 条"从而都不入队(M 永远不跑),或都读到"剩 0 条"从而都入队(M 跑两次)。

解法是一把最朴素的 promise 链互斥锁(withQueueLock,engine.ts:339-351):

// 示意,非源码 —— promise 链互斥锁
async function withQueueLock(fn) {
const prev = queueLock // 抢到上一位的"排队号"
let release
queueLock = new Promise(r => (release = r)) // 立刻把自己挂成新队尾
await prev // 再等前面的人做完
try { return await fn() } finally { release() }
}
// 重点看:先换 queueLock 再 await prev —— 这个顺序保证了严格 FIFO,不会插队

executeNodeAsync(engine.ts:500-517)把整个 handleNodeCompletion 包在这把锁里。所以:块执行是并发的,边记账是串行的。

executeNodeAsync 还有一处细节值得记:它先取 wasAlreadyExecuted = executedBlocks.has(nodeId),只有这次才第一次执行的节点才做完成处理。这挡住了"缓存命中的节点重复点亮下游"。

3.5 入队的两道过滤

addToQueue(engine.ts:291-300)只有 10 行,但有两道过滤:

过滤条件为什么
恢复触发器node.metadata.isResumeTrigger && !allowResumeTriggers恢复用的虚拟触发器只在"从快照恢复"这一模式下允许起跑
去重!readyQueue.includes(nodeId)同一节点被多条边同时点亮时,只入队一次

去重用的是 Array.includes 线性扫描。队列通常很短,这是合理取舍;若图极宽会变成 O(n²)。

3.6 起跑点有四种

initializeQueue(engine.ts:353-453)是一串 early-return,按优先级从高到低:

模式触发条件入队什么
从任意块起跑context.runFromBlockContext 存在只入队 startBlockId(见 §8)
恢复:释放挂起边metadata.remainingEdges 非空逐条处理暂停时挂着的边,把因此就绪的汇合节点入队
恢复:待办队列metadata.pendingBlocks 非空全部入队,然后清空 pendingBlocks
正常起跑以上都没有指定的 triggerBlockId,否则扫全图找 START_TRIGGER/STARTER

第二种模式里藏着一个细腻的判断(resolveRemainingEdgeHandle,engine.ts:461-482):持久化的 remainingEdges 可能没记 handle,要回 DAG 里查;若同一对 (source, target) 之间既有正常边又有 error,正常边优先。原因很实在:一次成功的恢复绝不能把自己的正常出路当成 error 边给剪掉。


4. 取消:三个入口,一个开关

取消的难点不是"设个标志位",而是正在等一个 30 秒的 LLM 请求时怎么立刻醒过来

┌─────────────────────┐
入口① AbortSignal(同进程) ─┐
入口② Redis pub/sub(跨进程)├─▶ signalCancelled() ─┬─▶ cancelledFlag = true
入口③ Redis key 启动兜底 ─┘ └─▶ abortResolve()

Promise.race([...executing, abortPromise])

race 立刻返回,主循环 break

三个入口各堵一个洞:

入口代码堵的洞
AbortSignalinitializeAbortHandler(engine.ts:85-100)同进程内的请求中断;构造时若已 aborted 立即取消
Redis pub/subsubscribeToCancellationChannel(engine.ts:75-83)执行跑在 A 进程、取消按钮点在 B 进程
Redis key 兜底checkCancellationBackstop(engine.ts:124-179)取消事件发布在本引擎订阅之前(典型场景:从快照恢复)

发布侧在 apps/sim/lib/execution/cancellation.ts:markExecutionCancelled(:43)先写 Redis key 再 publish,顺序是刻意的——这样即便订阅者还没上线,它启动时也能靠 isExecutionCancelled(:64)读到那个 key。

取消的三条重要性质

  • 取消不杀正在跑的节点。 signalCancelled(engine.ts:86-90)只置标志 + resolve abortPromise,让引擎不再等。已经发出的 HTTP/LLM 请求继续跑完然后被丢弃。真正能中断等待/请求的,是把 ctx.abortSignal 透传下去的那几个 handler(如 Wait 块的 sleepUntilAborted(waitMs, ctx.abortSignal),apps/sim/executor/handlers/wait/wait-handler.ts:121;Agent 块透传给 provider 请求,apps/sim/executor/handlers/agent/agent-handler.ts:2596)。
  • 取消优先于错误。 run() 的 catch 分支先判 cancelledFlag,是则返回 status: 'cancelled' 而不是抛错(engine.ts:237-247)。
  • 半截日志要补齐。 finalizeIncompleteLogs(engine.ts:677-687)给所有没有 endedAt 的块日志补上结束时间和真实耗时,否则 UI 上会留下一堆"永远在跑"的块。

收尾在 run()finally 里调 cleanup()(engine.ts:280-285)退订 pub/sub——引擎是一次性对象,不退订就是订阅泄漏。


5. 边激活:handle 命名协议就是分支协议

从这节开始进入本章重头戏。先讲"谁能跑"。

5.1 handle 是什么

DAG 里每条边带一个可选的 sourceHandle: string(apps/sim/executor/dag/types.ts:1-5)。它是画布上"从哪个小圆点拉出来的线"的 ID,而在运行期,它是一个协议字符串:引擎靠字符串比对决定这条边要不要激活。

全部协议常量集中在 EDGE(apps/sim/executor/constants.ts:69-82):

handle谁产生激活条件
undefined普通块之间的连线无条件激活
condition-<条件ID>图编译期由 Condition 块的出边按序生成(apps/sim/executor/dag/construction/edges.ts:166)output.selectedOption === 条件ID
router-<路由ID>Router V2 由 UI 指定;旧版 Router 用目标块 ID 当路由 ID(edges.ts:192)output.selectedRoute === 路由ID
error块的错误端口output.error 为真
source块的正常端口output.error 为假
loop_continue / loop-continue-source循环回边output.selectedRoute === 'loop_continue'
loop_exit循环出口output.selectedRoute === 'loop_exit'
parallel_continue / parallel_exit并行的下一批 / 出口同上,对应 PARALLEL_*

两个派生集合把这些 handle 分了组:

  • SUBFLOW_CONTROL_EDGE_HANDLES(constants.ts:78-84)= 五个循环/并行控制 handle。不在子流程语境下,它们一律不激活
  • CONTROL_BACK_EDGE_HANDLES(constants.ts:86-90)= 三个"回边"。级联失活时不许沿回边往回传播,否则会把循环体自己剪掉。

5.2 判定顺序是一个决策梯

shouldActivateEdge(apps/sim/executor/execution/edge-manager.ts:253-304)是一串 early-return,顺序即优先级:

output.selectedRoute 是 loop_exit / loop_continue / parallel_* ?
│ 是 → 只有 handle 完全匹配的那条边活;其余全死。子流程控制权最高
↓ 否
handle 属于子流程控制集合? → 死(不在子流程语境里,控制边不该点亮)
↓ 否
handle 为空? → 活(普通连线)
↓ 否
handle 以 'condition-' 开头? → 比对 output.selectedOption
↓ 否
handle 以 'router-' 开头? → 比对 output.selectedRoute
↓ 否
handle 是 'error' ? → 有 error 才活
handle 是 'source'? → 没 error 才活
其它 → 活(未知 handle 默认放行)

最后那条"未知 handle 默认放行"是保守设计:新增一种 handle 而忘了在这里加分支,结果是多跑而不是卡死

5.3 handler 如何反过来驱动边

引擎不认识"条件"和"路由",它只认识 output 里的两个字段。块 handler 用返回值反向驱动边激活——这就是 handler 层与调度层的全部接口:

handler文件 · 符号写进 output 的字段引擎据此点亮
Conditionapps/sim/executor/handlers/condition/condition-handler.ts:240 · ConditionBlockHandlerselectedOption: 命中条件的 id(:161)condition-<id>
Router(旧版)apps/sim/executor/handlers/router/router-handler.ts:42 · RouterBlockHandlerselectedRoute: String(chosenBlock.id)(:173)router-<目标块ID>
Router V2同上 · executeV2(:185)selectedRoute: chosenRoute.id(:354)router-<路由ID>
Function / APIhandlers/function/function-handler.ts:40 / handlers/api/api-handler.ts:14不写这两个字段全部无 handle 出边 + source
循环/并行哨兵apps/sim/executor/orchestrators/node.ts:51 · executeNodeselectedRouteloop_continue / loop_exit / parallel_* 之一对应控制边

旧版 Router 那一格值得停一下: 它没有独立的"路由 ID"概念,而是让 LLM 直接选一个目标块,于是图编译期把 handle 写成 router-<目标块ID>(edges.ts:190-193),运行期再拿 selectedRoute(就是那个块 ID)去比对。同一套 ROUTER_PREFIX 协议,同时兼容了"选端口"和"选块"两种心智,不用给引擎加分支。

Condition handler 还额外把决策写进 ctx.decisions.condition(condition-handler.ts:293:151),Router 同理写 decisions.router——这份决策记录会进快照,恢复时用来复现分支选择(见 05)。

5.4 error 端口:失败不一定是失败

一个块抛错时,是"整条工作流失败",还是"走 error 分支继续"?判据只有一条:这个块画了 error 出边没有

块执行抛错

├─ 写 errorOutput = { error: 消息 } 到状态、日志、SSE 回调

└─ hasErrorPortEdge(node) ?
│ 有 → blockLog.errorHandled = true,把 errorOutput 当正常返回值返回
│ 引擎照常传播边 → error 边激活、source 边失活并级联
│ 无 → 抛 buildBlockExecutionError → trackExecution 捕获 → 整段执行失败

实现在 handleBlockError(apps/sim/executor/execution/block-executor.ts:534-748),判据函数是 hasErrorPortEdge(:430-437)——它就是扫一遍出边找 sourceHandle === EDGE.ERROR

妙在:"错误可恢复"这件事不是配置项,是图的形状本身。用户画了错误分支 = 声明"我要自己处理这个错误",不用再填任何开关。


6. 分支剪枝:级联失活(本章的核心)

前一节讲了"哪条边活"。这节讲**"死掉的边要怎么处理"**——这是 Sim 调度里最不显然的一段设计。

6.1 为什么"不跑"是不够的

看一个最小的菱形:

┌──[condition-true]──▶ A ──┐
Cond ───┤ ├──▶ M
└──[condition-false]─▶ B ──┘

条件选了 true。A 跑,B 不跑。M 有两条入边(来自 A、来自 B)。

如果引擎只是"不把 B 入队",那么 A 跑完后 M 的账本是:还欠 B 一条入边。而 B 永远不会跑,于是 M 永远不就绪——工作流静默卡死。而且 hasWork() 返回 false(队列空、没有在跑的),run() 会正常返回 success,输出里少了 M 那一截。这是最难查的一类 bug。

所以正确做法是:边死了要说出来,而且要往下游传染。

Cond 选了 true:
[true] 边 → 激活 → A 的入边里删掉 Cond → A 就绪
[false] 边 → 失活 → 记进 deactivatedEdges
→ 级联:B 还有别的活路吗?没有
→ 把 B→M 这条边也记为失活
→ M:入边还剩 {B},但 B→M 已失活 → 活跃入边数 = 0 → 就绪 ✓

6.2 就绪的两个口径

isNodeReady(apps/sim/executor/execution/edge-manager.ts:111-113)只有一行,但它是两个口径的或:

// 真实源码 apps/sim/executor/execution/edge-manager.ts:110-112
isNodeReady(node: DAGNode): boolean {
return node.incomingEdges.size === 0 || this.countActiveIncomingEdges(node) === 0
}

对应两种"上游满足":

口径含义由谁造成
incomingEdges.size === 0所有上游都已经真的跑完并激活了边激活时 targetNode.incomingEdges.delete(node.id)(:63)
countActiveIncomingEdges === 0剩下没来的上游,它们的边全被失活了失活时不删入边,而是把边记进 deactivatedEdges

关键是这个不对称:激活 = 删入边;失活 = 记边。为什么不都删?因为失活是可撤销的——循环进入下一轮时,clearDeactivatedEdgesForNodes(:178)要把循环体内的失活记录抹掉,让同一批边下一轮重新可用。如果当初直接删了入边,就无从恢复了。

countActiveIncomingEdges(:369-388)遍历"还没到的上游",对每个上游只要找到一条没失活的边就计数 1 并 break(同一对节点之间可能有多条边,比如同时有 condition-truecondition-false 都指向 M)。

6.3 递归怎么停

deactivateEdgeAndDescendants(edge-manager.ts:306-346)是级联的本体。它的正确性全在几个停止条件上:

停止条件代码为什么必须停
这条边已经失活过:309-311幂等 + 防止图里有环时无限递归
目标节点还有别的活跃入边hasActiveIncomingEdges(:346-367)它还有活路,不能剪
目标节点已经收到过激活边nodesWithActivatedEdge.has(targetId)(:325)它已经被点亮了,更不能剪
出边是回边isBackwardsEdge(:236-238)沿 loop_continue 往回剪会剪掉循环体自己
// 示意,非源码 —— 级联失活的骨架
function deactivate(src, tgt, handle) {
const key = edgeKey(src, tgt, handle)
if (deactivatedEdges.has(key)) return // ① 已剪过
deactivatedEdges.add(key)

const node = nodes.get(tgt)
if (hasOtherActiveIncoming(node, key)) return // ② 还有别的活路
if (nodesWithActivatedEdge.has(tgt)) return // ③ 已经被点亮

for (const e of node.outgoingEdges.values()) {
if (!isBackwardsEdge(e.sourceHandle)) deactivate(tgt, e.target, e.sourceHandle)
}
}
// 重点看:三个 return 决定了"剪到哪停手";剪过头 = 该跑的不跑,剪不够 = 汇合点卡死

边的身份用 JSON 数组编码(createEdgeKey,:390-392):JSON.stringify([sourceId, targetId, handle ?? 'default'])。历史上用的是 ${src}-${tgt}-${handle} 拼接,那种格式在节点 ID 含连字符或互为前缀时会串味——clearDeactivatedEdgesForNodes 会把别人的失活记录一起清掉。normalizeSerializedEdgeKey(:415-430)专门负责把快照里的旧格式 key 迁移成新格式,老执行恢复后不至于全盘失效。

6.4 三个补丁,各堵一个坑

processOutgoingEdges(edge-manager.ts:16-109)主体只有四十行,但里面有三处一眼看不出用意的补丁。它们各自对应一类真实的卡死,edge-manager.test.ts 里都有对应用例。

补丁代码位置堵的坑
nodesWithActivatedEdge:40-42:94-105多个 error 端口汇到同一个块时,后成功的源头无法唤醒它
cascadeTargets + isTerminalControlNode:44-55:75-92分支在循环体内死掉,循环的 sentinel-end 永远等不到,循环不再前进
isRoutedDeadEnd + isEnclosingSentinel:72-90:209-230上一条补丁过度触发:把下游别的子流程的哨兵也点着了

补丁一:nodesWithActivatedEdge(已被点亮过的节点集合)

场景:两个块 P、Q 都画了 error 边指向同一个告警块 M。P 出错(error 边激活,M 被点亮);随后 Q 成功(Q 的 error 边失活)。此时 M 的活跃入边数刚刚归零,它就绪了——但这一轮 processOutgoingEdges 处理的是 Q,而 M 并不在 Q 的 activatedTargets 里,常规路径不会把它入队。

补丁的逻辑(:94-105):遍历本轮失活的目标,如果它 ① 曾被别人激活过(在 nodesWithActivatedEdge 里)、② 现在已就绪,就补一次入队。

反面用例同样重要:如果 P、Q 都成功,M 从未被任何边点亮,nodesWithActivatedEdge 里没有它,补丁不触发,M 正确地不跑(测试 should NOT mark target ready when all sources succeed,edge-manager.test.ts:1255)。

这个集合同时还是级联的刹车(§6.3 条件三),并且要进快照——所以 EdgeManager 暴露了 getNodesWithActivatedEdge / markNodeWithActivatedEdge / restoreDeactivatedEdges 三个方法(:133:148:141)。

补丁二:cascadeTargets + 终端控制节点

isTerminalControlNode(:240-250)的定义:这个节点有出边,且出边全部是子流程控制 handle——翻译成人话就是"循环/并行的 sentinel-end"。

场景:循环体里有个 condition,某轮它选中的那条路是死胡同(什么也不接)。这轮循环体走到头了,但 sentinel-end 是靠"有人激活我的入边"才跑的,现在没人激活它,循环卡在这一轮

补丁:级联失活过程中每碰到一个终端控制节点,就把它记进 cascadeTargets(:318-320);全部剪完后,如果这一轮一个目标都没激活(isDeadEnd)且该哨兵已就绪,就把它入队(:75-92),让循环继续或退出。另有一个直连特例(:49-55):失活目标本身就是终端控制节点时,不用等级联也直接收集。

补丁三:isRoutedDeadEnd——补丁二的刹车

补丁二写完后出现了新问题:一个 condition 有意选了死胡同分支时,它会把下游所有能碰到的哨兵都点着,包括后面另一个循环/并行的哨兵——那些子流程根本还没开始,就被推着"退出"了。

于是加了第二层判据(:73:79-87):

  • isRoutedDeadEnd = 死胡同 output 里有 selectedOptionselectedRoute(说明这是一次有意的路由,不是普通块自然走到头)。
  • 这种情况下只放行 isEnclosingSentinel(:209-230)为真的哨兵——即和当前节点同属一个子流程(subflowTypesubflowId 都相等)的那个哨兵。

一句话总结这三个补丁:补丁一管"汇合点被别人点过",补丁二管"死胡同也要让循环走下去",补丁三管"别让死胡同踢到别人家的循环"。

6.5 复习:一次条件分支的完整账本

起点: Cond 完成,output.selectedOption = 'c-true'

① 遍历出边 [condition-c-true → A] 活 → activatedTargets = [A]
[condition-c-false → B] 死 → edgesToDeactivate = [B]
② 记已点亮 nodesWithActivatedEdge += {A}
③ 级联失活 Cond→B 记入 deactivatedEdges
B 无其它活跃入边、未被点亮 → 继续剪 B→M
④ 删已激活入边 A.incomingEdges.delete(Cond)
⑤ 就绪判定 A: 入边空 → 就绪 ✓
⑥ 补丁扫描 cascadeTargets 空;M 不在本轮失活目标里 → 不处理
返回 [A] ——— M 会在 A 完成、A→M 激活那一轮才就绪

7. 块执行器:调度层与业务层的接缝

BlockExecutor 不做调度决策,但它决定了调度层看到什么样的 output,所以必须一起看。

7.1 一次 execute 的时间线

BlockExecutor.execute(apps/sim/executor/execution/block-executor.ts:96-420):

findHandler(block) 没找到 → 直接抛

记时间锚点 + 建 blockLog + 发 onBlockStart(哨兵节点跳过这一段)

解析输入 VariableResolver 失败 → handleBlockError(phase='input_resolution')

handler.execute(...) 失败 → handleBlockError(phase='execution')

是流式输出? → handleStreamingExecution 边转发给客户端边攒全文

normalizeOutput → 大对象外置(compactExecutionPayload) → 写 state → 发 onBlockComplete

几处值得记的细节:

  • startedAtstartTime 在同一个同步瞬间取(:102-103),让日志时间戳和 performance.now() 算出来的耗时共享同一个原点。
  • 哨兵节点不写日志、不发回调(isSentinel 分支,:107-112),它们是基础设施不是用户块。
  • 输入解析失败与执行失败走同一个错误出口,只用 phase 区分文案(:159-173:271-284)。

7.2 findHandler:顺序即优先级

// 真实源码 apps/sim/executor/execution/block-executor.ts:325-327
private findHandler(block: SerializedBlock): BlockHandler | undefined {
return this.blockHandlers.find((h) => h.canHandle(block))
}

线性扫描 + 首个命中。handler 数组由 createBlockHandlers()(apps/sim/executor/handlers/registry.ts:33-53)按固定顺序给出,GenericBlockHandler 必须排最后——它的 canHandle 直接 return true(apps/sim/executor/handlers/generic/generic-handler.ts:148-150),是所有"块 = 一次工具调用"的兜底路径(见 04)。注册表刻意不含哨兵:哨兵由 NodeExecutionOrchestrator 处理,不是用户块。

7.3 normalizeOutput:为什么调度层敢直接读 output.selectedRoute

// 真实源码 apps/sim/executor/execution/block-executor.ts:497-507
private normalizeOutput(output: unknown): NormalizedBlockOutput {
if (output === null || output === undefined) return {}
if (typeof output === 'object' && !Array.isArray(output)) return output as NormalizedBlockOutput
return { result: output }
}

三行规约保证了 handler 无论返回 null、返回字符串还是返回数组,到了 EdgeManager 手上都是一个可以安全取属性的对象。shouldActivateEdge 才敢无脑写 output.selectedRoute

7.4 状态只有一个入口

ExecutionState(apps/sim/executor/execution/state.ts:56)实现 BlockStateController(读写合并接口,apps/sim/executor/execution/types.ts:389)。写只有一处:setBlockOutput(state.ts:150-153)——写输出的同时把块加入 executedBlocks,这两件事绑死,避免出现"有输出但没标记执行过"的中间态。

BlockExecutor.setNodeOutput(executor/execution/block-executor.ts:437-470,文件已挪进 executor/execution/)在此之上做了一件调度相关的事:并行分支节点的输出会额外写两个别名 ID(全局分支 ID、外层分支作用域 ID),让别的分支用普通块名就能引用到。别名规则属于 03 的变量解析范畴。


8. 从任意块起跑:脏集计算

"改了中间某个块,只重跑它和它下游"——UI 上的一个按钮,调度上是 executeFromBlock(apps/sim/executor/execution/executor.ts:129-267)。

8.1 三个集合

computeExecutionSets(apps/sim/executor/utils/run-from-block.ts:76-139)做三次 BFS:

集合怎么算干什么用
dirtySet从起点沿出边BFS需要重跑的块
upstreamSet从起点沿入边BFS起点的祖先(诊断用)
reachableUpstreamSet每个脏块沿入边 BFS,跳过脏块本身要从快照里保留输出的块

第三个集合是关键。设想 A→CB→C,C 的表达式写着 A.result || B.result。从 A 重跑时,B 不脏也不是 A 的祖先——但 C 要引用它。只保留 upstreamSet 会让 C 读到 undefined;reachableUpstreamSet 把这类兄弟分支兜住了(executor.ts:145-148 的注释点明了这个用例)。

若起点是循环/并行容器,BFS 起点换成它的 sentinel-start(resolveContainerToSentinelStart,run-from-block.ts:57-65),容器 ID 本身也算脏。

8.2 三处状态手术

算完集合后要做三件事,少一件都跑不起来:

  1. 删掉脏块指向非脏源的入边(executor.ts:222-237)。汇合型脏块的上游是缓存值、这辈子不会再"激活"一次,不删就永远不就绪——本质上是 §6.1 那个卡死问题的另一种形态。
  2. 把脏块从 executedBlocks 里摘掉(:402-409)。否则 NodeExecutionOrchestrator.executeNode(apps/sim/executor/orchestrators/node.ts:68-75)会走"已执行,直接返回缓存"的快捷路径,脏块根本不会重跑。
  3. 过滤快照(:159-211)。只保留 reachableUpstreamSet / 可达容器 / 外层分支别名对应的块状态、循环与并行记账,其余丢弃。

跑起来之后,非脏节点仍会被入队并"执行",但 executeNodenode.ts:57-65 直接返回缓存输出——它们参与边传播,但不真的干活。这是个聪明的取巧:边的账本逻辑一行都不用改。

8.3 允许从哪些块起跑

validateRunFromBlock(apps/sim/executor/utils/run-from-block.ts:156-239)的规则:

规则为什么
块必须在 DAG 里,或者是循环/并行容器容器不是真节点,要换算成哨兵
不能在循环内部循环内部的块没有独立的迭代上下文可以复原
不能在并行内部同上,分支索引无从确定
不能是哨兵节点哨兵是基础设施
直接上游必须都已执行过否则没有缓存输出可用;触发器节点和哨兵免检

最后一条的免检逻辑很实在:触发器的判据是"没有入边"(:227),哨兵靠 metadata.isSentinel(:223)——它们本来就不进 executedBlocks


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

  1. 失活集合而不是删边。 激活删入边、失活记边(edge-manager.ts:64 vs :313)。这个不对称让"循环下一轮重置"变成一次集合清理(clearDeactivatedEdgesForNodes,:178),而不是重建图。
  2. 边身份用 JSON 数组当 key。 JSON.stringify([src, tgt, handle])(:390)天然规避了字符串拼接的分隔符歧义,并且带一个显式的旧格式迁移器(:415)。
  3. "错误可恢复"由图的形状表达。 画了 error 出边就等于开启了错误处理,没有第二个开关(hasErrorPortEdge,block-executor.ts:750)。
  4. 块执行并发、边记账串行。 一把 20 行的 promise 链互斥锁(withQueueLock,engine.ts:339)就解决了整个并发正确性,不需要任何并发原语库。
  5. 取消先写 Redis key 再 publish。 顺序保证了"发布早于订阅"的竞态也能被启动兜底查到(apps/sim/lib/execution/cancellation.ts:47-76 + engine.ts:124)。
  6. 未知 handle 默认放行。 shouldActivateEdge 的最后一个 default: return true(edge-manager.ts:302)——新增 handle 忘了改这里,后果是多跑一个块,不是整条流卡死。
  7. run-from-block 复用同一套边传播。 非脏节点照样入队、照样"完成",只是返回缓存(node.ts:57-65),调度代码零改动。

10. 边界与局限(诚实说)

  • 没有并发上限。 引擎会把就绪队列一次性全部起跑(processQueue,engine.ts:485-493),同时在跑的节点数完全由图的宽度决定。想限流只能靠图的形状;MAX_PARALLEL_BRANCHES: 20 帮不上忙——它夹的是并行块的批大小(parallel.ts:254-261),既不是分支数也不是全局并发。
  • 并行分支数本身没有上限。 一个 collection 型并行块拿到 500 条数据就展开 500 个分支(resolveBranchCount,parallel.ts:185-201,无 clamp),只是按批大小分 25 批推进;克隆出来的节点全部进同一张 DAG。
  • 取消不中断已发出的请求。 引擎只是不再等(signalCancelled,engine.ts:86);是否真的中断取决于各 handler 有没有把 ctx.abortSignal 透传下去,目前只有 wait / agent / pi / workflow / mothership 这几个 handler 传了。
  • 就绪队列去重是线性扫描。 readyQueue.includes(nodeId)(engine.ts:297),极宽的图上是 O(n²)。
  • 调试模式没实现。 DAGExecutor.continueExecution(executor.ts:105-122)直接返回失败,注释写明"重构后的执行器尚不支持"。
  • remainingEdges 走的是 any 逃逸。 恢复用的挂起边通过 (context.metadata as any).remainingEdges 传递(executor.ts:512engine.ts:365),没进类型系统。
  • 两处 API 没有生产调用方。 EdgeManager.restoreIncomingEdge(edge-manager.ts:115)只在测试里被调用(循环重置的入边恢复实际写在 apps/sim/executor/orchestrators/loop.ts:611restoreLoopEdges 里);processOutgoingEdgesskipBackwardsEdge = true 分支同样只有测试覆盖——引擎唯一调用点固定传 false(engine.ts:575)。
  • 卡死是静默的。 若某个补丁没覆盖到的形状导致汇合节点永不就绪,hasWork() 返回 false,run()正常返回 success,只是输出里少了一截。引擎没有"图没跑完"的自检。

11. 横向对比:同一张图,四种"这条路不走"

同货架的可视化工作流平台都要回答本章那三个问题,但"分支和环怎么表达"的答案分得很开:

项目图上怎么表达"这条路不走"调度怎么推进对应 doc
Sim(本章)边失活 + 沿下游级联剪枝;就绪判据 = 活跃入边为 0进程内就绪队列,队列宽度即并发度本章
Rivet一个 control-flow-excluded 排除值像毒药一样沿边传染,没有独立控制流图拉取式(pull)数据流引擎rivet/03
Langflow强连通分量识别环 + 环上顶点专用放行判据 + 两套并存的剪枝状态机分层拓扑排序出名单,RunnableVerticesManager 做前驱对账,一层 asyncio.gatherlangflow/03langflow/04
Kestra声明式 YAML 的状态推进,不靠画布连线剪枝没有常驻执行循环:每条队列消息在执行级锁里把 Execution 推进一步kestra/02
Dify调度器整个抽成外部包 graphon,本仓库只守节点工厂/队列/Layer 钩子队列流水线 + 可暂停恢复dify/index

三个可以带走的判断:

  • Sim 与 Rivet 是同一个想法的两种写法。 都用"沿边传播一个否定信号"代替显式分支状态机;区别在于 Rivet 把否定信号做成数据(排除值跟着数据流走),Sim 把它做成调度器旁边的一张集合(deactivatedEdges)。后者的直接好处是可以整份序列化进暂停快照(见 05)。
  • Sim 与 Langflow 都要处理"图有环",但落点不同。 Langflow 在拓扑排序层识别 SCC、给环上顶点特判;Sim 在建图层就让回边只登记出边、不登记入边(见 01 §6.8),于是运行期的就绪判定完全不必知道环的存在。
  • Sim 的调度是进程内的,Kestra 是队列驱动、无常驻循环的。 这条差别直接决定了:Sim 只能靠快照跨请求存活,Kestra 天然跨进程。

再往上一层:本货架把可视化 workflow 平台归在"长成什么产品"这一支,分支章见 guide/agent-products.md,货架总览见 ../index.md


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

主题文件关键符号
装配流水线apps/sim/executor/execution/executor.tsDAGExecutor.executebuildExecutionPipelinecreateExecutionContext
从任意块起跑apps/sim/executor/execution/executor.tsexecuteFromBlockrestoreSavedIncomingEdges
主循环与队列apps/sim/executor/execution/engine.tsExecutionEngine.runprocessQueuehasWorkaddToQueue
并发与互斥apps/sim/executor/execution/engine.tsexecuteNodeAsynctrackExecutionwithQueueLockwaitForAnyExecution
完成处理与边传播apps/sim/executor/execution/engine.tshandleNodeCompletion
入队起点apps/sim/executor/execution/engine.tsinitializeQueueresolveRemainingEdgeHandle
取消apps/sim/executor/execution/engine.tsinitializeAbortHandlersubscribeToCancellationChannelsignalCancelledcheckCancellationBackstopcleanup
取消的跨进程通道apps/sim/lib/execution/cancellation.tsmarkExecutionCancelledisExecutionCancelledgetCancellationChannel
边激活判定apps/sim/executor/execution/edge-manager.tsshouldActivateEdgeisSubflowControlEdgeisBackwardsEdge
边传播主体apps/sim/executor/execution/edge-manager.tsprocessOutgoingEdgesisNodeReadyisTargetReady
级联失活apps/sim/executor/execution/edge-manager.tsdeactivateEdgeAndDescendantshasActiveIncomingEdgescountActiveIncomingEdges
三个补丁apps/sim/executor/execution/edge-manager.tsnodesWithActivatedEdgeisTerminalControlNodeisEnclosingSentinel
边身份与迁移apps/sim/executor/execution/edge-manager.tscreateEdgeKeyparseEdgeKeynormalizeSerializedEdgeKey
循环重置钩子apps/sim/executor/execution/edge-manager.tsclearDeactivatedEdgesForNodesdeactivateResumedEdge
handle 协议常量apps/sim/executor/constants.tsEDGESUBFLOW_CONTROL_EDGE_HANDLESCONTROL_BACK_EDGE_HANDLESBlockTypeDEFAULTS.MAX_PARALLEL_BRANCHES
handle 生成apps/sim/executor/dag/construction/edges.tsEdgeConstructor.generateSourceHandle
并行的分支数与批大小apps/sim/executor/orchestrators/parallel.tsresolveBranchCountresolveBatchSizeinitializeParallelScope
块执行apps/sim/executor/execution/block-executor.tsBlockExecutor.executefindHandlernormalizeOutputsetNodeOutput
错误端口apps/sim/executor/execution/block-executor.tshandleBlockErrorhasErrorPortEdge
日志与流式apps/sim/executor/execution/block-executor.tscreateBlockLoghandleStreamingExecutionsanitizeInputsForLog
块状态apps/sim/executor/execution/state.tsExecutionStatesetBlockOutputgetBlockOutput
handler 注册表apps/sim/executor/handlers/registry.tscreateBlockHandlers
分支决策 handlerapps/sim/executor/handlers/router/router-handler.tscondition/condition-handler.tsRouterBlockHandler.executeV2ConditionBlockHandler.execute
兜底 handlerapps/sim/executor/handlers/generic/generic-handler.tsGenericBlockHandler.canHandle
节点分流apps/sim/executor/orchestrators/node.tsNodeExecutionOrchestrator.executeNodeisFinalSentinelOutput
脏集计算apps/sim/executor/utils/run-from-block.tscomputeExecutionSetsvalidateRunFromBlockresolveContainerToSentinelStart
行为规格(读测试最快)apps/sim/executor/execution/edge-manager.test.tsengine.test.tsapps/sim/executor/orchestrators/parallel.test.tsCascade deactivationMultiple error ports to same targetCancellation via Redisnormalizes %s(批大小 clamp)

接着读: 循环/并行的迭代语义与变量解析在 03 子流程编排与变量解析;块 handler 内部与模型 Provider 在 04;暂停点、快照恢复与 remainingEdges 的完整故事在 05;本章之前的图编译在 01;本组导览见 index.md