Workflow:Worker 线程里的脚本编排
追踪一段 workflow 脚本从 start() 校验到 Worker 线程内 vm.Context 沙箱执行,再到子 agent RPC、并发槽位、取消/grace/terminate 的完整流程。
引言:Workflow 不是 Tool Call
你在前面章节中已经看到 Agent 通过 tool call 驱动一个一个动作。Workflow 走的是截然不同的路线——它把一段预写好的 JavaScript 脚本放进隔离的 Worker 线程,用 vm.Context 沙箱执行。脚本不是模型实时输出的 tool call,而是编排层面的确定性逻辑:声明哪些子 agent 并行跑、哪些流水线串联、什么时候取消。
下面按一次 workflow 的真实生命周期走:从 start() 的同步校验开始,到 Worker 线程里的脚本执行、子 agent RPC、并发槽位、取消、grace timer 和 terminate,逐帧看它怎么收束。
第一帧:start() 同步校验
当调用方执行 workflowEngine.start(request) 时,引擎做三件事:
- Meta 校验 —
validateMeta()对request.meta做形状检查(name、description 必须非空字符串,phases 数组可选),违规抛META_INVALID。 - Body 预编译 —
assertBodyParses()用与 Worker 侧相同的包装器(async () => {\n${body}\n})()做一次new vm.Script(),确保语法正确。如果脚本以export const meta开头,直接报错——meta 通过 request 数据传递,不由脚本导出。 - Provider 解析 — 确认
ctx.subagents上注册了请求的 provider,否则抛AGENT_START。
校验通过后,引擎组装 WorkerInit 载荷:
interface WorkerInit {
meta: WorkflowMeta
body: string // 脚本正文
args?: unknown // 调用方传入的参数(plain JSON)
limits: WorkerLimits // 并发上限、总 agent 上限、同步超时
}
limits.maxConcurrentAgents 默认自动解析为 min(16, max(1, availableParallelism() - 2))——让 workflow 不会把全部 CPU 吃光,但也不至于退化成串行。
第二帧:Worker 线程诞生
WorkerRun 构造函数拿着 WorkerInit 作为 workerData 创建一个新的 Worker:
const { entry, options } = resolveWorkerSpawn(init)
this.worker = new Worker(entry, options)
关键安全设计:Worker 的环境被清洗。
execArgv: []— 不继承宿主进程的 Node 参数(--inspect、--require等全部清零)。env— 只有 Windows 的 TMP/TEMP 路径以及未编译形态下的TSX_TSCONFIG_PATH,没有 API Key、没有凭证。
这不是安全边界(代码文档明确说明),而是隔离保障——防止模型写出的脚本意外读到宿主环境变量或影响 Node 调试端口。
Worker 入口极简——worker.ts 只有三行:导入 parentPort 和 workerData,调用 runWorkerSession()。
第三帧:Session 握手与 Go 信号
Worker 侧启动后,runWorkerSession() 做以下事情:
- 构建
ChildRpcBridge(子 agent 的 RPC 管道)和ExecutionObserver(通过port.postMessage把 phase/log/agentStart/agentEnd 发回 host)。 - 创建
WorkflowExecution——编译脚本、构建vm.Context。 - 发送
Ready消息,然后等待 gate。
Host 侧收到 Ready 后回复 Go,gate 放行,execution.drive() 开始执行脚本。
这个握手存在的意义是:如果外部 signal 在 Worker 启动的瞬间已经 abort,host 会发送 Cancel 而非 Go,脚本体一行都不执行。
第四帧:vm.Context 沙箱内的全局 API
WorkflowExecution 构造函数创建一个空的 vm.Context 并注入全局变量:
| 全局名 | 类型 | 作用 |
|---|---|---|
agent(prompt, opts?) | 异步函数 | 启动一个子 agent 并等待结果 |
parallel(thunks) | 异步函数 | 并行执行一组 thunk,每个失败返回 null |
pipeline(items, ...stages) | 异步函数 | 每个 item 串过多个 stage,无跨 stage barrier |
phase(title) | 同步函数 | 标记当前阶段(影响后续 agent 的 label) |
log(message) | 同步函数 | 向 observer 输出叙述文本 |
args | 数据 | 调用方传入的参数(已经过 structured clone 隔离) |
函数全局用 Object.freeze() 冻结,脚本覆盖自己的 hook 只会伤害自己。args 是数据属性——workerData 的 structured clone 已完成了跨线程的深拷贝,脚本对其做任何 mutation 不影响调用方。
第五帧:drive() 的永不 reject 契约
drive() 是脚本的执行入口。它有一个核心设计原则:永远 resolve,永不 reject。所有故障都映射为 WorkflowResult 的变体:
interface WorkflowResult {
value: unknown // completed 时的返回值
stopReason: 'completed' | 'error' | 'cancelled'
error?: string // 非 completed 时的人类可读描述
agentsStarted: number // 本次执行启动了多少 agent
}
执行流程:
- 如果在 body 执行前已 cancelled → 直接返回 cancelled result。
this.compiled.runInContext(this.context, { timeout: syncTimeoutMs })— 启动脚本的同步前缀(超时保护)。await脚本返回的 Promise。- 如果 await 结束后发现已 cancelled → 返回 cancelled(防止 cancel 信号到达时脚本恰好自然结束)。
- 对返回值做
materializeFromRealm()序列化为 plain JSON → 返回 completed。 - 任何 catch 分支:如果已 cancelled 返回 cancelled,否则返回 error。
第六帧:agent() 调用的完整路径
当脚本执行 const result = await agent("summarize this", { schema }) 时,会经过以下阶段:
6.1 入口校验
throwIfCancelled()— 如果已取消,立即抛出CANCELLED。- 验证 prompt 为非空字符串。
readAgentOptions()把 opts 从 vm realm 物化为 plain JSON,校验支持的 key(label/phase/schema/provider/model)。- 检查
started < limits.maxTotalAgents(总量上限,防死循环)。
6.2 并发槽位获取
private acquireSlot(): Promise<void> {
if (this.activeSlots < this.limits.maxConcurrentAgents) {
this.activeSlots += 1
return Promise.resolve()
}
return new Promise((resolve, reject) => {
this.slotWaiters.push({ resolve: ..., reject })
})
}
FIFO 队列。当 cancel() 到来时,所有排队的 waiter 被 reject — 等待中的 agent 不会白白启动。
6.3 子 agent RPC(跨线程)
获得槽位后,执行路径跨越线程边界:
- Worker 侧
ChildRpcBridge.startAgent()分配一个callId,发送ChildStart消息。 - Host 侧
onChildStart()接收消息,调用this.subagents.start(provider, {...})启动真正的子 agent。 - 启动成功 → host 注册
ChildRecord,发回ChildStarted { callId, childId }。 - Worker 侧 bridge 拿到
childId,构造RpcChildHandle,返回给agent()函数。
然后 agent() 等待 run.result:
- 子 agent 完成 → host 将 result 做
snapshotJsonValue序列化,发送ChildSettled。 - 子 agent 失败 → host 发送
ChildFailed,worker 侧 reject。
6.4 结果处理
completed+ 有 schema → 返回result.structured(结构化数据)。completed+ 无 schema → 返回outputText(result.output)(拼接文本块)。- 非 completed → 返回
null(脚本可以.filter(Boolean)过滤)。
最后 releaseSlot() 唤醒下一个排队者。
第七帧:parallel() 和 pipeline() 组合子
parallel(thunks)
接收一个函数数组,Promise.all 并行执行。每个 thunk 内部抛出非 fatal 错误 → 该项返回 null;抛出 fatal WorkflowError(如 CANCELLED、AGENT_CAP)→ 整个 parallel() 向上传播,脚本终止。
// 脚本示例
const [a, b, c] = await parallel([
() => agent("task A"),
() => agent("task B"),
() => agent("task C"),
])
// 如果 B 的子 agent 失败,b === null,A 和 C 正常返回
pipeline(items, …stages)
每个 item 独立地串过所有 stage——没有跨 stage barrier。这意味着 item[0] 可以在 item[1] 还没进入 stage[0] 的时候就已经完成了 stage[2]。
const results = await pipeline(
documents,
(doc) => agent(`extract key points from: ${doc}`),
(points, doc) => agent(`write summary based on: ${points}`)
)
每个 item 链中的 stage 失败 → 该 item 返回 null,不影响其它 item。Fatal 错误终止所有。
Fatal vs Non-Fatal 的判定
关键设计:isFatalWorkflowError(error) 检测 instanceof WorkflowError。由于脚本运行在不同的 vm realm,脚本内部无法构造真正的 WorkflowError 实例——它拿不到 host realm 的类引用。这意味着 fatality 判断不可被脚本伪造。
第八帧:Realm 边界与物化
脚本 vm.Context 是一个独立的 JavaScript realm。从这个 realm 流出的值(agent options、脚本返回值)不能直接跨线程 postMessage,必须先物化为 plain JSON。
materializeFromRealm() 做递归遍历:
- 基础类型(boolean/string/finite number/null)→ 直接返回。
- 数组 → 逐元素递归,拒绝稀疏数组和非索引属性。
- 对象 → 检查原型链长度 ≤ 2(plain object),拒绝 Date/Map/class 实例。
- 拒绝 bigint、function、symbol、undefined(嵌套)、循环引用。
如果脚本返回了不可序列化的值,drive() 会抛出 RESULT_UNSERIALIZABLE,最终映射为 stopReason: 'error'。
第九帧:取消的三层防线
取消是 workflow 引擎最复杂的部分。三个层面协作:
层一:Worker 侧 — hook 边界拦截
cancel(reason) 设置 cancelReason。此后每一个 hook(agent/parallel/pipeline/phase/log)在入口调用 throwIfCancelled(),脚本在下一个 await 处死亡。
排队中的 slot waiter 被立即 reject,避免 cancel 后还有 agent 启动。
层二:Host 侧 — 子 agent abort
cancel() 同时 abort 共享的 AbortController(所有子 agent start 携带的 signal)。已启动的子 agent 收到 abort 信号,按各自 provider 的逻辑终止。
层三:Grace Timer + Terminate
如果 cancel 后脚本在 disposeGraceMs(默认 5s)内仍不结束(比如卡在一个没有 hook 的死循环),host 强制:
- 合成缺失的 agent-end 事件(
endStrandedAgents())。 settleResult(cancelledResult)。worker.terminate()— 杀死整个线程。
cancel() → [脚本在 hook 边界死亡 OR 脚本不响应] → grace 到期 → terminate
第十帧:Host 侧的消息分发
Host 通过 worker.on('message', ...) 接收 Worker 的消息,按 type 分发:
| Worker → Host | 处理 |
|---|---|
ready | 回复 go |
phase / log | 转发给 observer(cancel 后抑制) |
agent-start | 记入 liveAgents ledger + 通知 observer |
agent-end | 从 ledger 删除 + 通知 observer |
child-start | 调用 subagents.start() 启动子 agent |
child-dispose | 调用子 agent 的 dispose() |
result | 竞争 terminal 结果(见下文) |
反方向 Host → Worker 的消息:go、cancel、child-started、child-start-error、child-settled、child-failed、child-disposed。
第十一帧:Terminal 竞争与 Exactly-Once 配对
一次 workflow 运行有三种终止来源:
- Worker Result 消息 — 脚本正常完成或出错。
- Worker Death — 线程崩溃(
error/exit事件)。 - Grace Timer 到期 — cancel 后脚本不响应。
它们通过 terminalClaimed 标志做先到先得竞争。第一个 claim 成功的来源决定最终 result。
Agent 配对保证
Host 维护 liveAgents: Map<seq, WorkflowAgentInfo>。每个 workflow/agent-start 事件必须恰好配一个 workflow/agent-end。
- 正常路径:Worker 发
agent-end,hostendAgent()从 ledger 删除并通知 observer。 - 异常路径(Worker 死亡/grace 到期):
endStrandedAgents()遍历 ledger 中所有未配对的 start,合成outcome: 'cancelled'的 end 事件。
这保证无论 Worker 以何种方式终止,observer 看到的 start/end 总是 exactly-once paired。
第十二帧:子 agent 的生命周期管理
Host 为每个子 agent 维护 ChildRecord:
interface ChildRecord {
readonly run: SubagentRun
disposal?: Promise<void> // 幂等:第一次 dispose 启动,后续 join
}
三个触发 dispose 的路径:
- Worker 发
child-disposeRPC(正常路径)。 reapChildren()— cancel/death 时批量 abort + dispose。dispose()— 外部调用 WorkflowRun.dispose()。
所有路径汇聚到 disposeChild(),它的 disposal Promise 是幂等的——多次调用同一 callId 会 join 同一个 Promise。
Quiescence
childQuiescence() 在 pendingStarts.size === 0 && children.size === 0 时 resolve。dispose() 方法用 Promise.race([result + quiescence, sleep(grace)]) 等待——如果子 agent 在 grace 内没清理完,直接放弃(有限泄漏)。
第十三帧:contain() — 防止 unhandled rejection 杀死 Worker
脚本可能 await 一个 hook 返回的 Promise,也可能丢弃它(fire-and-forget)。如果丢弃的 Promise 后来因为 cancel 被 reject,Node 的 unhandled rejection 会杀死整个 Worker 线程。
解决方案是 contain():
private contain<T>(promise: Promise<T>): Promise<T> {
promise.catch(() => { /* consumed */ })
return promise
}
挂一个空的 .catch() 消费 rejection,但返回原始 Promise——如果脚本确实 await 了它,脚本仍能观察到 rejection。
第十四帧:事件发射与结果获取
Workflow 引擎发射以下 Cordis 事件:
workflow/start— 运行开始。workflow/phase— 进入新阶段。workflow/log— 脚本日志。workflow/agent-start/workflow/agent-end— 子 agent 生命周期。workflow/end— 运行结束(只携带 stopReason + agentsStarted,不携带 value)。
结果不通过事件传递。你必须 await workflowRun.result 获取返回值。这个设计防止事件监听者意外持有或修改结果数据的引用。
第十五帧:Worker 死亡处理
Worker 可能因多种原因死亡:
- 脚本内
process.exit()(不该发生但模型可能写出)。 - 同步超时(
syncTimeoutMs触发 vm 级 timeout exception,但如果同步代码在 native 层卡住则可能直接崩溃)。 - Node 内部错误。
Host 对 error/messageerror/exit 事件的处理:
- 第一个死亡信号关闭消息准入(
workerDeathObserved = true)——之后收到的 queued message 被静默丢弃。 reapChildren()— abort + dispose 所有注册子 agent。endStrandedAgents()— 合成缺失的 agent-end。- 如果 terminal 尚未 claim → 以 error result 结算(或如果之前已发起 cancel,则以 cancelled 结算)。
exit 事件额外做一轮 disposal sweep——确保所有 ChildRecord 都启动了 dispose。
第十六帧:配置全景
| 配置项 | 默认值 | 作用 |
|---|---|---|
provider | 'spawn' | 子 agent 使用的 subagent provider |
maxConcurrentAgents | 0(自动) | 并发 agent 上限 |
maxTotalAgents | 1000 | 单次 run 总 agent 上限(防死循环) |
maxItemsPerCall | 4096 | parallel/pipeline 单次调用的 items 上限 |
syncTimeoutMs | 5000 | vm 同步切片超时 |
disposeGraceMs | 5000 | cancel 后的宽限期,到期 terminate |
maxTotalAgents 是双层保护:引擎配置是天花板,request.maxTotalAgents 可以设更低但不能更高。
第十七帧:完整流程图
调用方 Host 主线程 Worker 线程
│ │ │
│── start(request) ──────────▶│ │
│ │── validate meta/body ────────│
│ │── new Worker(workerData) ───▶│
│ │ │── compile vm.Script
│ │ │── vm.createContext
│ │ │── inject globals
│ │◀── Ready ────────────────────│
│ │── Go ───────────────────────▶│
│ │ │── runInContext()
│ │ │ └── await agent(...)
│ │◀── ChildStart {callId} ──────│
│ │── subagents.start() ─────────│
│ │── ChildStarted ─────────────▶│
│ │ │ └── await run.result
│ │◀── (子agent完成) ────────────│
│ │── ChildSettled ─────────────▶│
│ │ │ └── agent() 返回
│ │ │── 脚本 return value
│ │◀── Result ───────────────────│
│◀── workflowRun.result ──────│ │
│ │── workflow/end event ─────────│
第十八帧:安全模型收口
Worker 线程不是安全边界,但提供三层隔离:
- 线程隔离 — 同步阻塞不影响 host 事件循环;可 terminate 强杀。
- 环境隔离 — scrubbed env 不暴露凭证;空 execArgv 不暴露调试端口。
- Realm 隔离 — vm.Context 内的对象不能跨 realm 冒充 host 类型(
instanceof检查自然成立)。
信任前提是:脚本由模型撰写并经过审查(“model-written”),不是任意用户输入。getter 和 proxy trap 在 materialize 期间会执行——这被视为可接受的,因为脚本作者是受信任的模型。
要点回顾
- Workflow 是确定性编排(脚本),不是模型实时决策(tool call)。
- 脚本跑在独立 Worker 线程的
vm.Context中,全局 API 只有agent/parallel/pipeline/phase/log/args。 drive()永不 reject——所有故障映射为stopReason。- 并发用 FIFO 槽位控制,总量用
maxTotalAgents封顶。 - 取消三层:hook 边界拦截 → 子 agent abort → grace timer + terminate。
- Host 保证
agent-start/agent-end的 exactly-once 配对(endStrandedAgents兜底)。 - 结果通过
workflowRun.result获取,不通过事件。 - Fatal error(
WorkflowError)终止脚本,non-fatal 让单项返回null继续执行。