数据截至 (上游 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:87、apps/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_points 列 | apps/sim/executor/types.ts:73-92 |
ResumeStatus | 恢复调度器 | PausePoint.resumeStatus | apps/sim/executor/types.ts:71 |
PausePoint 比 PauseMetadata 多两个字段,正是"从运行时进入持久世界"多出来的两件事:
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)干三件事:
- 算出这次暂停的身份
contextId与作用域(:90-91)。 - 用
buildResumeApiUrl/buildResumeUiUrl生成两条恢复链接(:102-103),一条给机器调、一条给人点。 - 可选地主动发通知——把审批链接塞进 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
}
三个细节都很关键:
- 仍然调用
handleNodeCompletion——暂停块的输出会被正常写进blockStates,所以它会进快照。 return得早,跳过了下面的edgeManager.processOutgoingEdges(:488)。也就是说暂停块的出边一条都没被激活。这是整个恢复机制的地基。- 暂停点存在一个
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: truesnapshotSeed——序列化好的整个执行态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)。按用途分四组:
| 组 | 字段 | 恢复时用来做什么 |
|---|---|---|
| 块结果 | blockStates、executedBlocks、blockLogs | 让下游块能引用上游产出;避免重跑 |
| 分支决策 | decisions.router、decisions.condition | 记住走过哪条岔路 |
| 子流程 | completedLoops、loopExecutions、parallelExecutions、parallelBlockMapping | 循环跑到第几轮、并行哪些分支已收 |
| 图与队列 | activeExecutionPath、pendingQueue、remainingEdges、dagIncomingEdges、deactivatedEdges、nodesWithActivatedEdge、completedPauseContexts | 重建"图现在长什么样、队列该从哪儿起" |
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 里已经是瘦身过的。
但守卫只覆盖两处:workflowVariables 和 state.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_snapshot | packages/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 就地改写成"已完成 + 带提交数据":
- 把
resumeInput规整成对象(字符串就试着JSON.parse,失败就包成{ value })(:790-803)。 - 与暂停时留下的
response合并,补上submission/submittedAt(:826-838)。 - 产出
mergedOutput, 盖上_resumed: true、_resumedFrom、_pauseDurationMs(:840-849),并把提交字段平铺到顶层(:868-870),这样下游可以直接写<审批块.reason>。 - 写回
blockStates[stateBlockKey],并把旧的pauseBlockId/contextId键删掉(:876-884)——避免同一个块留下两份状态。 - 同步块日志(
:894-915)和executedBlocks(:917-925)。
第 6 步是最容易漏的:聚合缓冲区。 如果暂停块长在循环或并行里,它的输出还被另外抄了一份进 loopExecutions[].currentIterationOutputs 或 parallelExecutions[].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 / resumedCount、status、nextResumeAt | packages/db/schema.ts:612-641 |
resume_queue | 一次恢复请求 | parentExecutionId、contextId、resumeInput、status(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 循环,每轮开一个事务:
- 锁住
paused_executions,状态不可恢复就直接返回。 - 有
claimed的就返回(还在跑,别插队)。 - 按
queuedAt升序取最老的一条pending(FOR UPDATE)。 - 对账:如果这个暂停点的
resumeStatus已经不是'queued'(被别的路径改掉了),把这条队列项判failed,continue下一轮。 - 否则改成
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)的完整逻辑:
- 全局分布式锁(Redis,
LOCK_KEY/ TTL 180s,:22-24、:43-49),多副本部署时只有一个在跑。 - 查
nextResumeAt <= now且状态 ∈{paused, partially_resumed}的行,最多 200 条(:54-74)。 - 每行挑出真正到点的
pauseKind: 'time'暂停点,逐个enqueueOrStartResume(:116-139)。 - 拿到
'starting'就调executeResumeJob(apps/sim/background/resume-execution.ts:46,签名多了可选signal)而不是直接startResumeExecution——注释解释了原因(:142-146):前者额外处理表格单元格上下文恢复和级联循环,而且在 trigger.dev 开关开与关时行为一致。 - 重算
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),把快照里的恢复信息一路透传下去:
| 传给执行器的字段 | 来源 | 谁消费 |
|---|---|---|
resumeFromSnapshot | metadata.resumeFromSnapshot | 引擎(决定 allowResumeTriggers、initializeQueue 分支) |
resumePendingQueue | snapshot.state.pendingQueue | 引擎队列初始化 |
remainingEdges | snapshot.state.remainingEdges | 引擎删边逻辑 |
dagIncomingEdges | snapshot.state.dagIncomingEdges | DAGBuilder / restoreSavedIncomingEdges |
snapshotState | snapshot.state | EdgeManager.restoreDeactivatedEdges、并行批次重建 |
跑之前还有一步 warmLargeValueRefs(:672-683):把快照里那些 __simLargeValueRef 指针预热回真实值,否则下游块读到的是指针不是数据。
finalizeExecutionOutcome(apps/sim/lib/workflows/executor/execution-core.ts:265-331)按 result.status 三选一收尾日志:cancelled → safeCompleteWithCancellation,paused → safeCompleteWithPause(: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 四类触发入口
| 入口 | 路由 / 任务 | 特点 | 位置 |
|---|---|---|---|
| Webhook | POST /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)→ executeScheduleJob | 用 croner 算 nextRunAt,靠 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 模式自己发 SSE | apps/sim/app/api/workflows/[id]/execute/route.ts:355 |
| Chat | POST /api/chat/[identifier] | 面向对话前端,走 executeWorkflow | apps/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:755、apps/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——人工动作不限流 |
用量额度 skipUsageLimits | false | true——恢复是已授权执行的延续 |
计费主体解析 actorUserId | 总是 | 总是(见下) |
| 账号封禁 | 总是 | 总是 |
恢复路径的这组设置在 apps/sim/lib/workflows/executor/human-in-the-loop-manager.ts:1436-1451,每一行都带注释解释为什么关。唯独不关的是计费主体解析和封禁检查——恢复执行仍然要归属到一个真实用户头上,并且这个用户被封了就不许继续。解析出来的 actorUserId 会覆盖快照里的 metadata.userId(:1085-1087)。