跳到主要内容

数据截至 (上游 commit e741923f72c3)

暂停、恢复与触发:让一次执行跨越请求边界活下来

30 秒导读: 前四章讲的引擎(02)假设一次执行从头跑到尾。本章讲它怎么中途停住——停在任意一个块上,把整个执行态压成 JSON 存进数据库,几小时甚至 30 天后由人点一下按钮或由 cron 唤醒,接着往下跑。以及反过来:一次执行到底是被什么打起来的


1. 这章要解决的小问题

普通的工作流引擎是一个函数:run(workflow) → result。函数返回,执行结束。

但真实业务里有两类需求打破这个假设:

需求例子卡住的时长
等人发退款前要主管点"批准"分钟 ~ 天,不可预测
等时间发完邮件等 3 天再发跟进确定,但远超请求超时

HTTP 请求活不了 3 天,Serverless 函数更活不了。所以问题不是"怎么等",而是:

怎么把一个正在跑的执行,完整地从内存里拆下来存进数据库,再在另一个进程里原样装回去继续跑。

这就是本章的主线。Sim 的答案有三层:一个约定(块怎么表达"我要停")、一份快照(执行态怎么被完整序列化)、一套调度(谁负责把它捡回来,以及并发恢复怎么不打架)。


2. 顶层全景:一次执行的两段生命

先看整体。一次带暂停的执行会被切成两个(或更多)互不相干的进程生命周期,中间靠 Postgres 接力。

请求 A(毫秒级) 请求 B(几小时/几天后)
┌─────────────────────┐ ┌──────────────────────┐
│ ① 触发,跑到暂停块 │ │ ④ 恢复入口 │
│ 引擎停手 │ │ 人点按钮 / cron 到点│
└──────────┬──────────┘ └──────────┬───────────┘
│ ② 序列化整个执行态 │ ⑤ 反序列化 + 重建 DAG
▼ ▼
┌─────────────────────────────────────────────────────────────┐
│ ③ Postgres:paused_executions + resume_queue │
│ (快照 JSON、每个暂停点的状态、恢复排队) │
└─────────────────────────────────────────────────────────────┘


⑥ 删掉暂停块的出边 →
下游块变"就绪" → 继续跑

怎么读这张图: 左右是两次独立的进程执行,中间的数据库是唯一的连接物。⑥ 是全章最反直觉的一步——恢复不是"从暂停处继续",而是改图让下游重新就绪,§5 展开。

各部件一句话职责:

部件干什么在哪个文件
_pauseMetadata 约定块表达"我要停"的唯一方式apps/sim/executor/types.ts:228
ExecutionEngine.handleNodeCompletion看到该字段就停手、记暂停点apps/sim/executor/execution/engine.ts:536
serializePauseSnapshot把 Map/Set 全展平成可 JSON 的执行态apps/sim/executor/execution/snapshot-serializer.ts:186
PauseResumeManager落库、排队、单飞、恢复、取消apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:490
DAGExecutor + EdgeManager重建图、回灌失活边、重启队列apps/sim/executor/execution/executor.ts:87apps/sim/executor/execution/edge-manager.ts:142
触发面决定"谁在什么时候按下开始"apps/sim/triggers/registry.ts:482 + 四类 API 路由

3. 暂停侧:一个块怎么把执行按下暂停键

3.1 约定:输出里塞一个字段

Sim 没有为"暂停"设计特殊的返回类型、异常或信号量。它用了最轻的一种耦合:块正常返回一个对象,对象里多带一个 _pauseMetadata 字段

// 示意,非源码:任何块只要这么返回,引擎就会停
return {
status: 'waiting',
_pauseMetadata: { // 就这一个字段决定生死
contextId: 'block_7', // 这次暂停的唯一身份
blockId: 'block_7',
response: { /* 给调用方看的东西 */ },
timestamp: new Date().toISOString(),
pauseKind: 'time', // 'human' 还是 'time'
resumeAt: '2026-08-15T09:00:00.000Z',
},
}

这个字段声明在 NormalizedBlockOutput 上(apps/sim/executor/types.ts:203-229),并被登记进 EXECUTION_CONTROL_OUTPUT_FIELD_NAMES(apps/sim/executor/types.ts:231-236)——那是一张"控制用、不算业务数据"的字段名单,和 error / selectedOption / selectedRoute 并列。

为什么这么设计值得注意: 暂停能力对块作者是零仪式的。任何 handler,不管在循环里、并行分支里、还是普通链路上,只要在返回值里加这个字段就获得了暂停语义,不需要 implement 新接口。

3.2 三个类型:PauseMetadata / PausePoint / ResumeStatus

类型谁产生活在哪位置
PauseMetadata块 handler内存中的块输出里apps/sim/executor/types.ts:52-69
PausePoint引擎(由 PauseMetadata 转写)ExecutionResult 与数据库 pause_pointsapps/sim/executor/types.ts:73-92
ResumeStatus恢复调度器PausePoint.resumeStatusapps/sim/executor/types.ts:71

PausePointPauseMetadata 多两个字段,正是"从运行时进入持久世界"多出来的两件事:

  • resumeStatus——这个暂停点现在处于恢复流程的哪一步('paused' | 'resumed' | 'failed' | 'queued' | 'resuming')。
  • snapshotReady——快照是否已经写好了。没写好就不许恢复(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:716-722 会直接抛错)。

pauseKind 只有两种取值(apps/sim/executor/types.ts:50),它决定了后面由谁来唤醒:

pauseKind谁唤醒恢复入口
'human'人(或外部系统调 API)POST /api/resume/{workflowId}/{executionId}/{contextId}
'time'cron 轮询器GET /api/resume/poll

3.3 两个暂停源

其一:Human-in-the-loop 块。 HumanInTheLoopBlockHandler(apps/sim/executor/handlers/human-in-the-loop/human-in-the-loop-handler.ts:61)干三件事:

  1. 算出这次暂停的身份 contextId 与作用域(:90-91)。
  2. buildResumeApiUrl / buildResumeUiUrl 生成两条恢复链接(:102-103),一条给机器调、一条给人点。
  3. 可选地主动发通知——把审批链接塞进 Slack / 邮件工具的参数里发出去(executeNotificationTools,:161:446-518),_pauseContext 里携带 resumeApiUrl 和表单结构。

最后把 pauseMetadata 挂进输出(:179-192:220-224)。

其二:Wait 块。 WaitBlockHandler(apps/sim/executor/handlers/wait/wait-handler.ts:64)是个二选一:

模式怎么等上限代价
同步(默认)进程内 sleep,可被 abort 信号和 Redis 取消打断5 分钟(:14)占着进程和请求
异步(async: true)返回 pauseKind: 'time'_pauseMetadata,进程退出30 天(:17)需要 cron 捡回来

异步模式还禁掉了 seconds 单位(:139-141)——秒级的等待走异步(落库 + 轮询)比进程内 sleep 还贵,直接报错更诚实。

同步 sleep 的实现值得看一眼(:24-76):它不是裸 setTimeout,而是同时挂着三条退出路径——定时器到点、AbortSignal 触发、以及每 500ms 查一次 Redis 里的取消标记。返回 boolean 表示"是睡满了还是被打断了"。

3.4 contextId:暂停点在循环和并行里的身份问题

一个暂停块如果长在循环里,它会被执行 N 次,产生 N 个互不相干的暂停。所以暂停点的身份不能只是块 ID

generatePauseContextId(apps/sim/executor/human-in-the-loop/utils.ts:12-28)按需给块 ID 加后缀:

基础块 ID → block_7
在并行分支 2 里 → block_7_branch_2_ (并行分支前缀/后缀)
在循环第 3 轮里 → block_7_loop3
两者都有 → block_7_branch_2__loop3

配套的 mapNodeMetadataToPauseScopes(apps/sim/executor/human-in-the-loop/utils.ts:30-62)把当前循环轮次从 ctx.loopExecutions 里读出来,写进 loopScope.iteration

对应地,恢复侧有个反向函数 normalizePauseBlockId(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:3551-3559),用 id.replace(/_loop\d+/g, '') 把循环后缀剥掉,还原成原始块 ID——因为边是画在原始块之间的,查下游要用原始 ID,而记录状态要用带后缀的 ID。这两套 ID 在恢复代码里反复来回转换(stateBlockKey vs pauseBlockId vs contextId,见 :781-884),是这一带最容易看晕的地方。

节点 id 的完整编码规则(分支下标 ₍N₎、外层分支克隆 __obranch-N、轮次后缀 _loopN)属于子流程那一层,见 03 章。本章只用到"剥掉后缀还原原始块 ID"这一半。

3.5 引擎怎么判定与收尾

引擎的判定只有一个 if(apps/sim/executor/execution/engine.ts:536-545):

if (output._pauseMetadata) {
await this.nodeOrchestrator.handleNodeCompletion(this.context, nodeId, output)
const pauseMetadata = output._pauseMetadata
this.pausedBlocks.set(pauseMetadata.contextId, pauseMetadata)
this.context.metadata.status = 'paused'
this.context.metadata.pausePoints = Array.from(this.pausedBlocks.keys())
return
}

三个细节都很关键:

  1. 仍然调用 handleNodeCompletion——暂停块的输出会被正常写进 blockStates,所以它会进快照。
  2. return 得早,跳过了下面的 edgeManager.processOutgoingEdges(:488)。也就是说暂停块的出边一条都没被激活。这是整个恢复机制的地基。
  3. 暂停点存在一个 Map<contextId, PauseMetadata>(apps/sim/executor/execution/engine.ts:37),天然支持一次执行里有多个并发暂停(比如并行的三个分支各卡在一个审批上)。

引擎主循环跑干净后,如果 pausedBlocks 非空就走 buildPausedResult(apps/sim/executor/execution/engine.ts:203-205:493-524),产出:

  • status: 'paused'success: true(暂停不是失败)
  • pausePoints[]——每个都盖上 resumeStatus: 'paused'snapshotReady: true
  • snapshotSeed——序列化好的整个执行态
  • output——单个暂停就直接是它的 response,多个就包成 { pausedBlocks, pauseCount }(collectPauseResponses,:545-556)

4. 快照侧:把执行态压成一行 JSON

4.1 快照装了什么

ExecutionSnapshot(apps/sim/executor/execution/snapshot.ts:5)是个薄壳,六个只读字段 + toJSON() / fromJSON():元数据、工作流定义、输入、工作流变量、选中输出、以及 state

肉在 state 里,类型是 SerializableExecutionState(apps/sim/executor/execution/types.ts:86-124)。按用途分四组:

字段恢复时用来做什么
块结果blockStatesexecutedBlocksblockLogs让下游块能引用上游产出;避免重跑
分支决策decisions.routerdecisions.condition记住走过哪条岔路
子流程completedLoopsloopExecutionsparallelExecutionsparallelBlockMapping循环跑到第几轮、并行哪些分支已收
图与队列activeExecutionPathpendingQueueremainingEdgesdagIncomingEdgesdeactivatedEdgesnodesWithActivatedEdgecompletedPauseContexts重建"图现在长什么样、队列该从哪儿起"

serializePauseSnapshot(apps/sim/executor/execution/snapshot-serializer.ts:186-308)负责组装。它的主要工作是把运行时的 Map/Set 全部展平成 JSON 能表达的东西:

  • Object.fromEntries(context.blockStates)Array.from(context.executedBlocks)(:211-212)
  • 循环/并行作用域里还有嵌套的 Map,得单独处理(serializeLoopExecutions,:140;serializeParallelExecutions,:161)
  • 从 DAG 里抄出每个节点的入边集合 → dagIncomingEdges(:201-208)
  • 从 EdgeManager 里抄出失活边集合 → deactivatedEdges / nodesWithActivatedEdge(:225-226)

4.2 为什么要逐字段精算 JSON 字节数

序列化器里有一段看起来很怪的代码:一个手写的 getEscapedJsonStringByteLength(apps/sim/executor/execution/snapshot-serializer.ts:18-47),逐字符判断转义规则和 UTF-8 编码长度;外面套一个 getBoundedJsonByteLength(:69-126)递归走整棵对象树累加字节。

为什么不直接 JSON.stringify(x).length? 因为目的不是"知道多大",而是"判断是否超过 8MB 阈值"(LARGE_VALUE_THRESHOLD_BYTES,apps/sim/lib/execution/payloads/large-value-ref.ts:3)。用 stringify 测量,意味着为了确认一个可能有 500MB 的对象太大,先真的把它变成 500MB 的字符串——这一步本身就可能把进程打爆。

所以这个测量函数带预算,一旦超支立刻掉头:

// 示意,非源码:带预算的字节计数,超了就停,不再往下走
function measure(value, budget) {
let bytes = 2 // {} 或 []
for (const [key, child] of entries(value)) {
bytes += keySize(key) + 1 // "key":
bytes += measure(child, budget - bytes) ?? 4
if (bytes > budget) return bytes // 重点:超了立刻返回,不遍历剩下的
}
return bytes
}

真源码里这两处提前退出在 :100:121

守卫函数 assertSnapshotValueIsCompact(:128-133)一旦发现超标就直接抛 Cannot serialize pause snapshot with oversized ...; compact it first.

"必须外置"是什么意思? 大块数据(长文本、大数组、文件)不进快照,而是被写进对象存储,快照里只留一个 __simLargeValueRef 指针。块输出的这一步在执行时就做完了——每个块跑完立刻走 compactExecutionPayload(apps/sim/executor/execution/block-executor.ts:330-337)。所以到序列化快照时,blockStates 里已经是瘦身过的。

但守卫只覆盖两处:workflowVariablesstate.loopExecutions(:229-230)。这两条路径不经过块级压缩,所以要在这里兜底。这是个诚实的、有针对性的窄守卫,不是全量校验。

4.3 大值的引用计数

外置的大值不能随便删——一份大值可能同时被"执行日志"和"暂停快照"引用,而暂停快照可能存活 30 天。所以有一组表做引用计数:

干什么位置
execution_large_values大值本体的元数据(key、大小、拥有者、软删时间)packages/db/schema.ts:535
execution_large_value_references(key, executionId, source) 三元主键;source 枚举是 execution_log / paused_snapshotpackages/db/schema.ts:567
execution_large_value_dependencies大值之间的父子引用(嵌套指针)packages/db/schema.ts:589

每次写快照,persistPauseResult 都会重算这份快照引用了哪些 key(collectLargeValueReferenceKeys),然后用 replaceLargeValueReferenceKeysWithClient 在同一个事务里整体替换(executionId, 'paused_snapshot') 的引用集合(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:592-603:361-372)。恢复后快照变了,再替换一次(:1699-1710)。


5. 恢复 = 重建 DAG + 回灌失活边 + 重启队列

这是全章最需要建直觉的一节。

5.1 直觉:恢复不是"跳到某一行继续"

传统的 durable execution(比如 Temporal)靠重放:重新执行代码,遇到已完成的步骤就用记录的结果替代。Sim 不重放,它重建图状态

回忆 02 章的调度规则:一个节点的 incomingEdges 集合空了(或者说该激活的都激活了),它就进就绪队列。所以:

暂停块之所以卡住下游,是因为它的出边一条都没被删。恢复要做的,就是把这些边删掉,下游自然就绪。

这就是为什么 §3.5 里那个提前 return 是地基。

5.2 三步走

快照 JSON


┌───────────────────────────────────────────────────────┐
│ ① 重建 DAG │
│ DAGBuilder.build(workflow, { savedIncomingEdges }) │
│ 从工作流定义重新编译一遍图,再用快照里的 │
│ dagIncomingEdges 覆盖每个节点的入边集合 │
└────────────────────────┬──────────────────────────────┘

┌───────────────────────────────────────────────────────┐
│ ② 回灌失活边 │
│ edgeManager.restoreDeactivatedEdges( │
│ deactivatedEdges, nodesWithActivatedEdge) │
│ 把"哪些分支已经被剪掉"重新装回剪枝器 │
└────────────────────────┬──────────────────────────────┘

┌───────────────────────────────────────────────────────┐
│ ③ 重启队列 │
│ engine.initializeQueue(): │
│ 逐条删掉 remainingEdges(暂停块的出边), │
│ 删完就检查目标节点是否 ready,ready 就入队 │
└───────────────────────────────────────────────────────┘

第①步 有两处代码做同一件事:DAGBuilder.build 内部(apps/sim/executor/dag/builder.ts:83-95)和 DAGExecutor.restoreSavedIncomingEdges(apps/sim/executor/execution/executor.ts:269-278),后者在 build 之后又覆盖一遍(:96)。中间夹着 restoreSnapshotParallelBatches(:278-309)——并行分支的克隆节点是运行时展开的,不在原始工作流定义里,必须按快照里的 currentBatchStart / currentBatchSize 重新展开一遍,否则那些 xxx_branch_2_ 节点在新 DAG 里根本不存在。

第②步buildExecutionPipeline 里(apps/sim/executor/execution/executor.ts:374-377)。restoreDeactivatedEdges(apps/sim/executor/execution/edge-manager.ts:142-147)把字符串数组变回 Set,顺手做了一次格式归一化

这里有个真实的兼容性补丁值得拎出来:边的 key 现在是 JSON.stringify([sourceId, targetId, handle])(createEdgeKey,apps/sim/executor/execution/edge-manager.ts:395-397),但历史上是 `${source}-${target}-${handle}` 这种拼接串。老格式在节点 ID 本身含 - 时会有歧义。normalizeSerializedEdgeKey(:415-430)的做法是:先试着 JSON.parse,能解析就是新格式直接用;不能解析就遍历整张 DAG 的边,按老格式重新拼一遍字符串去比对,命中了就换成新 key。这是为"数据库里躺着 30 天前写的老快照"付的代价。

第③步engine.initializeQueue()(apps/sim/executor/execution/engine.ts:353-453)。这个函数总共有四种起跑模式(完整清单见 02 章 §3.6),与恢复有关的是中间三种:

条件行为行号
remainingEdges 非空逐条删边,删完检查目标是否 ready,ready 就入队:282-326
pendingBlocks 非空直接把这些块塞进队列:328-346
resumeFromSnapshot === true 但两者都空什么都不做(终端暂停块,下游本来就没有):348-351

删边那段还有个 error 边的特判(:293-306):如果暂停块和目标之间既有正常边又有 error 边,恢复成功意味着 error 路径应该被剪掉而不是被放行,所以走 deactivateResumedEdge。持久化的 remainingEdges 可能没记 sourceHandle,resolveRemainingEdgeHandle(:376-397)会回查当前 DAG 补上,并且明确规定正常边优先于 error 边

5.3 恢复输入怎么灌回暂停块

删边只解决了"放行",还得解决"审批人填的表单去哪儿"。这在 runResumeExecution 里(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:1034-1387),核心是把暂停块的 blockState 就地改写成"已完成 + 带提交数据":

  1. resumeInput 规整成对象(字符串就试着 JSON.parse,失败就包成 { value })(:790-803)。
  2. 与暂停时留下的 response 合并,补上 submission / submittedAt(:826-838)。
  3. 产出 mergedOutput,盖上 _resumed: true_resumedFrom_pauseDurationMs(:840-849),并把提交字段平铺到顶层(:868-870),这样下游可以直接写 <审批块.reason>
  4. 写回 blockStates[stateBlockKey],并把旧的 pauseBlockId / contextId 键删掉(:876-884)——避免同一个块留下两份状态。
  5. 同步块日志(:894-915)和 executedBlocks(:917-925)。

第 6 步是最容易漏的:聚合缓冲区。 如果暂停块长在循环或并行里,它的输出还被另外抄了一份loopExecutions[].currentIterationOutputsparallelExecutions[].branchOutputs——那才是循环结束时聚合结果的来源。只改 blockStates 的话,聚合出来的还是"暂停中"的那份旧输出。

updateResumeOutputInAggregationBuffers(:61-105)专门补这一刀。它靠 isPausedOutputForContext(:55-59)——检查缓冲区里那条记录的 _pauseMetadata.contextId 是否匹配——来定位要替换的那一项,而不是靠位置或键名猜。

最后组装恢复用的快照(:1022-1029),metadata.resumeFromSnapshot = true,交给 executeWorkflowCore(:1308-1316)。

5.4 一个特例:终端暂停块

如果暂停块没有下游(工作流以审批结尾),edgesToRemove 为空,resumeTerminalNoop 置为 true(:997)。这时恢复执行会跑一个空队列、什么块都不执行。为了让调用方仍然看到有意义的输出,结果的 output 被替换成暂停块的 mergedOutput(:1318-1323)。

executeWorkflowCore 那边也有对应的判断:恢复执行必须跳过触发块解析,哪怕队列是空的(apps/sim/lib/workflows/executor/execution-core.ts:609-618)——"空队列"在恢复语境下是有意义的,不能被当成"没找到起点"而回退去找 Start 块。


6. 持久化与恢复调度:PauseResumeManager

PauseResumeManager(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:490,2316 行,全是 static 方法)是整个暂停/恢复的控制平面。它管两张表:

一行代表关键列位置
paused_executions一次被暂停的执行executionSnapshot(JSONB 快照)、pausePoints(JSONB,contextId → 暂停点)、totalPauseCount / resumedCountstatusnextResumeAtpackages/db/schema.ts:612-641
resume_queue一次恢复请求parentExecutionIdcontextIdresumeInputstatus(pending/claimed/completed/failed)packages/db/schema.ts:643-668

paused_executions 上有个部分索引值得注意:nextResumeAtIdx 只覆盖 status = 'paused' AND next_resume_at IS NOT NULL 的行(packages/db/schema.ts:637-639)。cron 轮询器每次扫的就是这个索引,所以纯人工审批的行(nextResumeAt 为 NULL)完全不进这个索引——一张可能有百万行的表,轮询只碰到"确实在等时间"的那几行。

下面各小节的裸行号(:241-376 这类)都指 apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts,不再重复写路径。

6.1 落库:persistPauseResult

:241-376。在一个事务里 SELECT ... FOR UPDATE 锁住行,然后分叉:

  • 没有既有行 → 插入,顺手写大值引用(:286-313)。
  • 已有行(说明这是"恢复后又暂停了",比如 wait1 → agent → wait2)→ 合并 pausePoints:先把所有 resumeStatus === 'resuming' 的旧点改判为 'resumed'(:317-324),再把新的暂停点盖上去(:326-328),重算计数与状态。

metadata合并而非替换的(:352-355),注释解释得很直白:表格单元格任务把 cellContext 外键存在同一列上,重新暂停时不能把它冲掉。

方法末尾无条件调一次 processQueuedResumes(:375)——刚落库就看看有没有排队的恢复请求可以启动。

6.2 并发恢复:enqueueOrStartResume 的单飞

问题: 一次执行有三个并行审批,三个人同时点了"批准"。三个恢复执行同时从同一份快照起跑,写同一批状态——必然打架。

Sim 的答案:同一个 parentExecutionId 下,永远只有一个 claimed 的恢复在跑。

enqueueOrStartResume(:378-511)整个跑在一个事务里,FOR UPDATE 锁住 paused_executions 行,然后:

请求进来


┌───────────────────────┐
│ 校验:执行可恢复? │ status ∈ {paused, partially_resumed}
│ 暂停点存在? │ 且 resumeStatus === 'paused'
│ 快照就绪? │ 且 snapshotReady
│ pauseKind 允许? │ allowedPauseKinds 白名单
└───────────┬───────────┘

已有 claimed 的恢复在跑?
┌─────┴─────┐
是 否
│ │
▼ ▼
插入 pending 插入 claimed
点状态→queued 点状态→resuming
返回队列位次 返回 'starting'

两处状态更新用的是 Postgres 的 jsonb_set 原地改一个 JSON 路径(:455:495),而不是读出整个 pausePoints 再整体写回——避免了并发写互相覆盖。

allowedPauseKinds 白名单(:415-420)是个安全边界:人工恢复端点只允许 ['human'](apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.ts:145-151),cron 轮询器只允许 ['time'](apps/sim/app/api/resume/poll/route.ts:496-503)。人不能提前戳醒一个定时等待,cron 也不能替人批准。

6.3 排队的消化:processQueuedResumes

:2059-2179。这是把队列往前推的泵,在恢复完成、恢复失败、以及新暂停落库后都会被调一次。

它是个 while 循环,每轮开一个事务:

  1. 锁住 paused_executions,状态不可恢复就直接返回。
  2. claimed 的就返回(还在跑,别插队)。
  3. queuedAt 升序取最老的一条 pending(FOR UPDATE)。
  4. 对账:如果这个暂停点的 resumeStatus 已经不是 'queued'(被别的路径改掉了),把这条队列项判 failed,continue 下一轮。
  5. 否则改成 claimed,暂停点改 'resuming',跳出循环。

出了事务才真正启动执行,而且是 fire-and-forget(:2165-2178,.catch 只记日志)——因为恢复可能跑几分钟,不能占着调用方的请求。

6.4 完成与失败的三种收尾

方法什么时候用对暂停点做什么位置
markResumeCompleted恢复正常跑完resumeStatus → 'resumed',resumedCount + 1:1490-1549
markResumeFailed恢复真的炸了resumeStatus → 'failed',日志判 failed:1551-1578
markResumeAttemptFailed基础设施暂时不可用(Redis 缓冲区拿不到、终端事件发不出去)resumeStatus → 'paused',退回可重试状态:1580-1618

第三个是设计上的分水岭:区分"工作流失败"和"这次尝试没能开始"。后者把暂停点放回 'paused',人可以再点一次;如果混成 failed,一次 Redis 抖动就永久毁掉一个待审批。

markResumeCompleted 末尾还有个状态回填(:1519-1547):看还剩几个未恢复的暂停点,全恢复完就置 fully_resumed,否则把执行日志的状态从 running 改回 pending——UI 上这次运行又变回"等待中"。

6.5 时间暂停:cron 轮询

computeEarliestResumeAt(:224-238)从一堆暂停点里挑出最早的那个 resumeAt(可带 after 过滤),写进 paused_executions.nextResumeAt。一行可能有多个时间暂停,只存最早那个就够了——到点了会重算。

GET /api/resume/poll(apps/sim/app/api/resume/poll/route.ts:59)的完整逻辑:

  1. 全局分布式锁(Redis,LOCK_KEY / TTL 180s,:22-24:43-49),多副本部署时只有一个在跑。
  2. nextResumeAt <= now 且状态 ∈ {paused, partially_resumed} 的行,最多 200 条(:54-74)。
  3. 每行挑出真正到点的 pauseKind: 'time' 暂停点,逐个 enqueueOrStartResume(:116-139)。
  4. 拿到 'starting' 就调 executeResumeJob(apps/sim/background/resume-execution.ts:46,签名多了可选 signal)而不是直接 startResumeExecution——注释解释了原因(:142-146):前者额外处理表格单元格上下文恢复和级联循环,而且在 trigger.dev 开关开与关时行为一致。
  5. 重算 nextResumeAt,用 { after: now } 排掉刚发出去的(:181-184)。

第 5 步前面那段注释是本章最值得抄走的运维观点(:177-180):

失败的派发永远不自动重试——工作流块不是幂等的,重跑可能重复发邮件、重复扣款。宁可把行留在那里让人工排查。

6.6 取消一个暂停中的执行

beginPausedCancellation(:1719-1765)返回一个布尔,含义是"能不能立刻标记为已取消":

  • claimed 的恢复正在跑 → 置 cancelling,返回 false。取消是延后的,等那次恢复自己跑完/被中断,由 startResumeExecution 里的收尾逻辑(:589-598)调 completePausedCancellation 补刀。
  • 没有在跑的 → 置 cancelling,返回 true,调用方可以直接推进到 cancelled

blockQueuedResumesForCancellation(:1808-1851)把所有 pending 的队列项一次性判 failed,防止取消过程中又有排队的被捡起来。

cancelling 这个中间态在 SQL 里被反复保护:好几处 status 更新写成 CASE WHEN status = 'cancelling' THEN 'cancelling' ELSE ... END(:1514:1528:1599),确保取消意图不会被并发的恢复完成流程覆盖掉。

6.7 状态机全貌

paused_executions.status:

persistPauseResult


┌─────────┐ 部分暂停点恢复 ┌───────────────────┐
│ paused │─────────────────>│ partially_resumed │
└────┬────┘ └─────────┬─────────┘
│ │ 全部恢复
│ 取消请求 ▼
│ ┌───────────────┐
├─────────────────────────>│ fully_resumed │
▼ └───────────────┘
┌────────────┐ 收尾完成 ┌───────────┐
│ cancelling │──────────────>│ cancelled │
└────────────┘ └───────────┘

PausePoint.resumeStatus(每个暂停点独立):

paused ──(有恢复在跑)──> queued ──(轮到它)──> resuming ──> resumed
▲ │ │
│ │ └──> failed
└── markResumeAttemptFailed(基础设施抖动) ───┘

7. 执行入口与触发面

前面讲的都是"停下来之后",这一节讲"最开始是被谁打起来的",以及暂停这套机制如何嵌进正常的执行流水线。

7.1 两层入口函数

各类触发源


executeWorkflow(workflow, requestId, input, actorUserId, opts)
│ 造 ExecutionMetadata + ExecutionSnapshot
│ apps/sim/lib/workflows/executor/execute-workflow.ts:61

executeWorkflowCore({ snapshot, callbacks, loggingSession, ... })
│ 拉工作流状态、解密环境变量、解析触发块、序列化、
│ 组 contextExtensions、new Executor().execute()
│ apps/sim/lib/workflows/executor/execution-core.ts:307

finalizeExecutionOutcome ──> handlePostExecutionPauseState
execution-core.ts:223 pause-persistence.ts:27

executeWorkflow(apps/sim/lib/workflows/executor/execute-workflow.ts:92-260)是给"我有一个工作流,帮我跑一遍"的调用方用的便利层:生成 executionId、建 LoggingSession、造快照、跑核心、打点、最后调 handlePostExecutionPauseState(:163)。注意它在 status === 'paused'不打 workflow_executed 事件(:143)——暂停不算跑完一次。

executeWorkflowCore(apps/sim/lib/workflows/executor/execution-core.ts:376)是真正的核心。与本章相关的是它组装 contextExtensions 时那一段(:638-670),把快照里的恢复信息一路透传下去:

传给执行器的字段来源谁消费
resumeFromSnapshotmetadata.resumeFromSnapshot引擎(决定 allowResumeTriggersinitializeQueue 分支)
resumePendingQueuesnapshot.state.pendingQueue引擎队列初始化
remainingEdgessnapshot.state.remainingEdges引擎删边逻辑
dagIncomingEdgessnapshot.state.dagIncomingEdgesDAGBuilder / restoreSavedIncomingEdges
snapshotStatesnapshot.stateEdgeManager.restoreDeactivatedEdges、并行批次重建

跑之前还有一步 warmLargeValueRefs(:672-683):把快照里那些 __simLargeValueRef 指针预热回真实值,否则下游块读到的是指针不是数据。

finalizeExecutionOutcome(apps/sim/lib/workflows/executor/execution-core.ts:265-331)按 result.status 三选一收尾日志:cancelledsafeCompleteWithCancellation,pausedsafeCompleteWithPause(:245-253),其余 → safeComplete。暂停的收尾不写 finalOutput——执行还没结束,没有最终输出可写。

handlePostExecutionPauseState(apps/sim/lib/workflows/executor/pause-persistence.ts:28-73)是所有调用方的强制收尾,TSDoc 写明"每个 executeWorkflowCore 的调用方都必须调这个"。它只有两个分支:暂停了就 persistPauseResult(没有 snapshotSeed 就把执行判失败,:34-36——宁可失败也不留一个恢复不了的暂停),没暂停就 processQueuedResumes

7.2 SSE 事件

apps/sim/lib/workflows/executor/execution-events.ts:10-40 定义了 11 种事件类型。与暂停相关的是 'execution:paused'(:63-74),它和 'execution:completed' 结构几乎一样(都带 output、duration、finalBlockLogs),但语义上是"这一段跑完了,还没结束"。

值得注意的是在事件缓冲区层面,暂停被当成终态写入:writeBufferedEvent(..., 'complete')(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:1882-1899)。也就是说前端这条 SSE 流会正常关闭,恢复时会开一条新的流(resetExecutionStreamBuffer,:1098-1101)。

7.3 四类触发入口

入口路由 / 任务特点位置
WebhookPOST /api/webhooks/trigger/[path]executeWebhookJob有准入闸(tryAdmit),先做 provider 挑战验证,再排后台任务apps/sim/app/api/webhooks/trigger/[path]/route.ts:53;apps/sim/background/webhook-execution.ts:290
定时GET /api/schedules/execute(cron)→ executeScheduleJobcronernextRunAt,靠 lastQueuedAt 防重复排队apps/sim/app/api/schedules/execute/route.ts:1800;apps/sim/background/schedule-execution.ts:796
API 调用POST /api/workflows/[id]/execute支持 sync / stream / async 三种模式,stream 模式自己发 SSEapps/sim/app/api/workflows/[id]/execute/route.ts:355
ChatPOST /api/chat/[identifier]面向对话前端,走 executeWorkflowapps/sim/app/api/chat/[identifier]/route.ts:48(调用点 :283)

四类入口最后都会落到 handlePostExecutionPauseState,只是路径不同:

入口怎么调到
Webhook后台任务里显式调(apps/sim/background/webhook-execution.ts:463)
定时后台任务里显式调(apps/sim/background/schedule-execution.ts:677)
API 调用路由里显式调,同步与异步两条分支各一处(apps/sim/app/api/workflows/[id]/execute/route.ts:960:1562)
Chat不自己调,靠 executeWorkflow 内部那一次(apps/sim/lib/workflows/executor/execute-workflow.ts:221)

这就是"任何入口进来的执行都能暂停"的保证。

webhook 和定时这两类还有一个共同参数:triggerBlockId(apps/sim/background/webhook-execution.ts:755apps/sim/background/schedule-execution.ts:600)。它们不从 Start 块进入,而是从触发块本身进入——因为一个工作流可以挂多个触发器,进来的是哪一个决定了从哪儿起跑。

7.4 触发器注册表

TriggerConfig(apps/sim/triggers/types.ts:11-41)是每个第三方触发器的声明:

字段作用
id / provider / name / version身份
subBlocks配置界面长什么样(复用块的 SubBlock 体系,见 01 章)
outputs: Record<string, TriggerOutput>这个触发器往工作流里灌什么数据,决定下游 <trigger.xxx> 能引用什么
webhook?HTTP 方法与固定 header
polling?是拉(cron 去问)还是推(对方来调)

TRIGGER_REGISTRY(apps/sim/triggers/registry.ts:482-884)是一张扁平的 id → config 大表,as-of 本 commit 共 336 个触发器条目,覆盖约 55 个服务商(Slack、Airtable、Attio、Azure DevOps、Zendesk、Zoom……)。注册表本身没有逻辑,就是一张查找表——这正是它能长到 336 条还不腐烂的原因。

7.5 起跑前的闸门:preprocessExecution

apps/sim/lib/execution/preprocessing.ts:163。所有入口在真正执行前都过这一关。它做的事和它的可关开关同样重要:

检查默认恢复执行怎么设
工作流存在且启用总是总是
部署状态 checkDeployment非 manual 触发才查false——恢复的是既有执行,不重查部署
限流 checkRateLimit非 manual/chat 才查false——人工动作不限流
用量额度 skipUsageLimitsfalsetrue——恢复是已授权执行的延续
计费主体解析 actorUserId总是总是(见下)
账号封禁总是总是

恢复路径的这组设置在 apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:1436-1451,每一行都带注释解释为什么关。唯独不关的是计费主体解析和封禁检查——恢复执行仍然要归属到一个真实用户头上,并且这个用户被封了就不许继续。解析出来的 actorUserId 会覆盖快照里的 metadata.userId(:1085-1087)。

7.6 触发相关的表

一行代表关键约束位置
webhook一个部署版本上的一个 webhook 端点(path, deploymentVersionId) 唯一(未归档)packages/db/schema.ts:1003
workflow_schedule一条定时规则(workflowId, blockId, deploymentVersionId) 唯一;两个"到期"部分索引分别服务 workflow 与 jobpackages/db/schema.ts:889
workflow_deployment_version一个不可变的已部署版本(workflowId, version) 唯一;isActive 标当前版本packages/db/schema.ts:3274
idempotency_key一次已处理过的外部事件key 主键,result 存返回值packages/db/schema.ts:3379
workflow_execution_snapshots一份工作流状态的去重存档(workflowId, stateHash) 唯一packages/db/schema.ts:398

两个设计点值得单独说:

webhook 与 schedule 都挂在 deploymentVersionId 上,不是挂在 workflow 上。 这意味着"部署新版本"天然带来一套新的端点与定时规则,老版本的还能并存到归档为止——触发面和版本是绑死的。

workflow_execution_snapshots 与本章的暂停快照不是一回事。 它按 stateHash 去重存"工作流当时长什么样",供执行日志引用(workflow_execution_logs.stateSnapshotId,packages/db/schema.ts:427-429);暂停快照存的是"执行跑到哪儿了",在 paused_executions.executionSnapshot。同名不同物,读代码时容易混。


8. 巧妙之处

  • 暂停是一个字段,不是一套接口。 _pauseMetadata 让暂停能力对块作者零成本,也让"暂停"和"错误"、"分支选择"一样只是输出里的一个控制字段(EXECUTION_CONTROL_OUTPUT_FIELD_NAMES,apps/sim/executor/types.ts:231)。想加第三种暂停源?写个 handler 返回这个字段就行。

  • 恢复靠删边,不靠重放。 因为暂停块的出边从未被激活(apps/sim/executor/execution/engine.ts:536-545 提前 return),恢复只要把这些边删掉,02 章那套就绪判定就会自动把下游推进队列。不需要"重放已完成步骤"这类复杂度。

  • 测大小不许先变成字符串。 getBoundedJsonByteLength 带预算递归、超支立刻掉头(apps/sim/executor/execution/snapshot-serializer.ts:101:121),用来判断"是否超过 8MB"而不会为了判断先分配 8MB+。配套的 getEscapedJsonStringByteLength 逐字符按 JSON 转义规则和 UTF-8 编码算,包括代理对(:27-36)。

  • 单飞 + 排队,而不是加锁重试。 同一执行下永远只有一个 claimed 恢复,其余进 resume_queue;推进靠 processQueuedResumes 这个泵。所有状态跃迁都在 SELECT ... FOR UPDATE 事务里,JSON 子字段用 jsonb_set 原地改(集中在 updatePausePointResumeStateSql 一族辅助函数,apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:214-241,另一处在 :2102),避免读-改-写覆盖。

  • 区分"失败"和"没能开始"。 markResumeAttemptFailed(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:2210)把基础设施抖动导致的失败退回 'paused',不烧掉这个待审批;真正的执行错误才走 markResumeFailed

  • 明确拒绝自动重试。 时间暂停的派发失败不重试,理由写在 apps/sim/app/api/resume/poll/route.ts:542-547:工作流块不幂等,重跑会重复副作用,宁可让人来看。

  • 聚合缓冲区靠 contextId 认领,不靠位置。 updateResumeOutputInAggregationBuffers(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:143-187)用 _pauseMetadata.contextId 在循环/并行的输出缓冲里找那一条,替换成恢复后的输出——这样即使数组顺序或键名对不上也不会改错。

  • 老快照的边 key 兼容路径。 normalizeSerializedEdgeKey(apps/sim/executor/execution/edge-manager.ts:420-435)在 JSON key 解析失败时,遍历 DAG 用老拼接格式反查——一段专为"数据库里躺着 30 天前的快照"存在的代码。


9. 边界与局限

  • isResumeTrigger 是死的。 NodeMetadata.isResumeTrigger 有两处读(apps/sim/executor/execution/engine.ts:293 的入队拦截、apps/sim/executor/dag/builder.ts:106 的日志),但全仓没有任何一处写它。也就是说 allowResumeTriggers 那道闸目前恒不生效。

  • pauseTriggerMapping 永远是空 Map。 NodeConstructor.execute 建了它就直接返回(apps/sim/executor/dag/construction/nodes.ts:18:36),从没 .set() 过;消费方 apps/sim/executor/dag/construction/edges.ts:304pauseTriggerMapping.get(...) ?? source 因此恒走 ?? source。与上一条一样,像是一次重构后留下的骨架。

  • 进程内等待硬顶 5 分钟。 非 async 的 Wait 块超过 MAX_INPROCESS_WAIT_MS 直接抛错并提示开异步(apps/sim/executor/handlers/wait/wait-handler.ts:114-118)。异步顶到 30 天(:17)。

  • 快照体积守卫只覆盖两处。 assertSnapshotValueIsCompact 只查 workflowVariablesloopExecutions(apps/sim/executor/execution/snapshot-serializer.ts:251-252)。blockStates / blockLogs 靠块执行时的 compactExecutionPayload(apps/sim/executor/execution/block-executor.ts:330)提前瘦身,不在这里兜底——如果那条路径漏了,超大快照会直接写进 Postgres 而不报错。

  • 恢复现在是新 executionId。 enqueueOrStartResumeconst resumeExecutionId = generateId()(apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:745),排进 resume_queue 时同时记 parentExecutionIdnewExecutionId(:753-757)。与旧版"复用父执行 ID"相反,现在每次恢复尝试有独立 ID,父子关系靠队列表还原;排查时顺着 resume_queue 就能还原尝试序列。

  • 恢复失败没有自动补偿。 派发失败不重试(见 §6.5),markResumeFailed 也不会把工作流回滚——已经执行过的块的副作用留在那里。

  • 本章不覆盖的相邻主题: 暂停块的输出如何被下游引用见 03 章;块本身怎么被调度、边怎么被激活/剪枝见 02 章;Agent 块内部的多轮工具调用(它不产生暂停,是块内的循环)见 04 章


10. 横向对比:同类项目怎么让一次执行跨越请求边界

"跑到一半停下来等人/等时间"是 workflow-builder 这一派的公共难题(总览见总库 长成什么产品 §2 流派一,以及 货架地图 分支 E)。同 shelf 的兄弟项目给出了四种不同答案:

项目停下来时存什么恢复怎么接上取舍
Sim(本章)整个执行态 + 边状态(失活边、已点亮节点、剩余入边)进 Postgres paused_executions重建 DAG → 回灌失活边 → 删掉暂停块出边让下游就绪快照要连"分支剪枝的历史"一起存,换来恢复后分支语义不走样
Dify05-human-in-the-loopGraphRuntimeState + generate entity 序列化进对象存储Celery 任务读回快照、重建引擎、从断点继续快照放对象存储不放主库,代价是多一跳存储依赖
PySpur04-human-in-loopPauseError 异常把整条流冻住,下游节点标 PENDING 而非 CANCELED/FAILEDresumed_node_ids + precomputed_outputs 从断点续跑靠"状态往 PENDING 掰"的贯穿约定守可恢复性,不需要序列化整个引擎
Kestra02-executor-state-machine什么都不用另存:Executor 本身无状态,执行态一直在数据库里队列消息回灌自己,换个实例接着推进天生跨进程、可水平扩展,代价是每一步都要读写 DB + 加执行级锁

两点值得单独提:

  • Sim 的特殊之处在"存边不存指令指针"。 其余三家存的都是"跑到哪个节点了"(节点态或异常态),Sim 存的是"调度器当时那张图长什么样"。因为它的分支语义完全由边的激活/失活表达(02 章),不把剪枝历史一起存,恢复后条件分支就会重新变成"两条路都活着"。
  • 触发面的形态也分两派。 Sim 把 336 条触发器做成一张静态注册表 + 四类 HTTP/后台入口;Kestra 则有一个独立的 Scheduler 服务每秒扫触发器状态、自己生成 Execution(见 04-scheduler-and-triggers)。前者贴 Serverless 部署形态,后者贴常驻集群。

11. 代码地图

主题文件符号
暂停三类型 + 输出约定apps/sim/executor/types.tsPauseKindPauseMetadataPausePointResumeStatusNormalizedBlockOutput._pauseMetadataEXECUTION_CONTROL_OUTPUT_FIELD_NAMESExecutionResult
人工审批块apps/sim/executor/handlers/human-in-the-loop/human-in-the-loop-handler.tsHumanInTheLoopBlockHandler(executeWithNodeexecuteNotificationTools)
等待块(同步 / 异步)apps/sim/executor/handlers/wait/wait-handler.tsWaitBlockHandlerMAX_INPROCESS_WAIT_MSMAX_ASYNC_WAIT_MSsleep
暂停身份与作用域apps/sim/executor/human-in-the-loop/utils.tsgeneratePauseContextIdmapNodeMetadataToPauseScopesbuildTriggerBlockId
恢复链接构造apps/sim/executor/constants.tsPAUSE_RESUMEbuildResumeApiUrlbuildResumeUiUrl
引擎的暂停判定与收尾apps/sim/executor/execution/engine.tshandleNodeCompletionbuildPausedResultcollectPauseResponsesinitializeQueueresolveRemainingEdgeHandle
快照壳apps/sim/executor/execution/snapshot.tsExecutionSnapshot(toJSONfromJSON)
可序列化执行态apps/sim/executor/execution/types.tsSerializableExecutionStateExecutionMetadataContextExtensions
快照序列化与体积精算apps/sim/executor/execution/snapshot-serializer.tsserializePauseSnapshotgetEscapedJsonStringByteLengthgetBoundedJsonByteLengthassertSnapshotValueIsCompact
大值阈值apps/sim/lib/execution/payloads/large-value-ref.tsLARGE_VALUE_THRESHOLD_BYTESLARGE_VALUE_REF_MARKER
失活边的存取与归一化apps/sim/executor/execution/edge-manager.tsgetDeactivatedEdgesrestoreDeactivatedEdgesnormalizeSerializedEdgeKeycreateEdgeKeydeactivateResumedEdge
DAG 入边恢复apps/sim/executor/dag/builder.tsDAGBuildOptions.savedIncomingEdgesDAGBuilder.build
恢复时的图与管线组装apps/sim/executor/execution/executor.tsDAGExecutor.executerestoreSavedIncomingEdgesrestoreSnapshotParallelBatchesbuildExecutionPipeline
暂停/恢复控制平面apps/sim/lib/workflows/executor/human-in-the-loop-manager.tsPauseResumeManager(persistPauseResultenqueueOrStartResumestartResumeExecutionrunResumeExecutionprocessQueuedResumesmarkResumeCompletedmarkResumeAttemptFailedupdateSnapshotAfterResumebeginPausedCancellationsetNextResumeAtnormalizePauseBlockId)、computeEarliestResumeAtupdateResumeOutputInAggregationBuffers
执行后的强制收尾apps/sim/lib/workflows/executor/pause-persistence.tshandlePostExecutionPauseState
执行入口(便利层)apps/sim/lib/workflows/executor/execute-workflow.tsexecuteWorkflowExecuteWorkflowOptions
执行核心apps/sim/lib/workflows/executor/execution-core.tsexecuteWorkflowCorefinalizeExecutionOutcomewasExecutionFinalizedByCore
SSE 事件类型apps/sim/lib/workflows/executor/execution-events.tsExecutionEventTypeExecutionPausedEvent
时间暂停轮询apps/sim/app/api/resume/poll/route.tsGETdispatchRowLOCK_KEYPOLL_BATCH_LIMIT
人工恢复端点apps/sim/app/api/resume/[workflowId]/[executionId]/[contextId]/route.tsPOSTGET
恢复后台任务apps/sim/background/resume-execution.tsexecuteResumeJob
触发器声明与注册表apps/sim/triggers/types.tsapps/sim/triggers/registry.tsTriggerConfigTriggerOutputTriggerRegistryTRIGGER_REGISTRY
起跑闸门apps/sim/lib/execution/preprocessing.tspreprocessExecutionPreprocessExecutionOptionsPreprocessExecutionResult
Webhook / 定时后台任务apps/sim/background/webhook-execution.tsapps/sim/background/schedule-execution.tsexecuteWebhookJobexecuteScheduleJob
暂停与触发相关表packages/db/schema.tspausedExecutionsresumeQueueworkflowExecutionSnapshotsexecutionLargeValuesexecutionLargeValueReferencesexecutionLargeValueDependencieswebhookworkflowScheduleworkflowDeploymentVersionidempotencyKey