青雲的博客
深入浅出 DeepSeek Harness 第五部:Goal与多Agent——谁在掌舵 第 31 章

Workflow:Worker 线程里的脚本编排

追踪一段 workflow 脚本从 start() 校验到 Worker 线程内 vm.Context 沙箱执行,再到子 agent RPC、并发槽位、取消/grace/terminate 的完整流程。

源码版本
47f943859bef60e4160492346772ded9b24f765a
验证日期
Commit
47f943859bef60e4160492346772ded9b24f765a

引言: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) 时,引擎做三件事:

  1. Meta 校验validateMeta()request.meta 做形状检查(name、description 必须非空字符串,phases 数组可选),违规抛 META_INVALID
  2. Body 预编译assertBodyParses() 用与 Worker 侧相同的包装器 (async () => {\n${body}\n})() 做一次 new vm.Script(),确保语法正确。如果脚本以 export const meta 开头,直接报错——meta 通过 request 数据传递,不由脚本导出。
  3. 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 只有三行:导入 parentPortworkerData,调用 runWorkerSession()


第三帧:Session 握手与 Go 信号

Worker 侧启动后,runWorkerSession() 做以下事情:

  1. 构建 ChildRpcBridge(子 agent 的 RPC 管道)和 ExecutionObserver(通过 port.postMessage 把 phase/log/agentStart/agentEnd 发回 host)。
  2. 创建 WorkflowExecution——编译脚本、构建 vm.Context
  3. 发送 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
}

执行流程:

  1. 如果在 body 执行前已 cancelled → 直接返回 cancelled result。
  2. this.compiled.runInContext(this.context, { timeout: syncTimeoutMs }) — 启动脚本的同步前缀(超时保护)。
  3. await 脚本返回的 Promise。
  4. 如果 await 结束后发现已 cancelled → 返回 cancelled(防止 cancel 信号到达时脚本恰好自然结束)。
  5. 对返回值做 materializeFromRealm() 序列化为 plain JSON → 返回 completed。
  6. 任何 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(跨线程)

获得槽位后,执行路径跨越线程边界:

  1. Worker 侧 ChildRpcBridge.startAgent() 分配一个 callId,发送 ChildStart 消息。
  2. Host 侧 onChildStart() 接收消息,调用 this.subagents.start(provider, {...}) 启动真正的子 agent。
  3. 启动成功 → host 注册 ChildRecord,发回 ChildStarted { callId, childId }
  4. 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(如 CANCELLEDAGENT_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 强制:

  1. 合成缺失的 agent-end 事件(endStrandedAgents())。
  2. settleResult(cancelledResult)
  3. 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 的消息:gocancelchild-startedchild-start-errorchild-settledchild-failedchild-disposed


第十一帧:Terminal 竞争与 Exactly-Once 配对

一次 workflow 运行有三种终止来源:

  1. Worker Result 消息 — 脚本正常完成或出错。
  2. Worker Death — 线程崩溃(error/exit 事件)。
  3. Grace Timer 到期 — cancel 后脚本不响应。

它们通过 terminalClaimed 标志做先到先得竞争。第一个 claim 成功的来源决定最终 result。

Agent 配对保证

Host 维护 liveAgents: Map<seq, WorkflowAgentInfo>。每个 workflow/agent-start 事件必须恰好配一个 workflow/agent-end

  • 正常路径:Worker 发 agent-end,host endAgent() 从 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 的路径:

  1. Worker 发 child-dispose RPC(正常路径)。
  2. reapChildren() — cancel/death 时批量 abort + dispose。
  3. 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 事件的处理:

  1. 第一个死亡信号关闭消息准入(workerDeathObserved = true)——之后收到的 queued message 被静默丢弃。
  2. reapChildren() — abort + dispose 所有注册子 agent。
  3. endStrandedAgents() — 合成缺失的 agent-end。
  4. 如果 terminal 尚未 claim → 以 error result 结算(或如果之前已发起 cancel,则以 cancelled 结算)。

exit 事件额外做一轮 disposal sweep——确保所有 ChildRecord 都启动了 dispose。


第十六帧:配置全景

配置项默认值作用
provider'spawn'子 agent 使用的 subagent provider
maxConcurrentAgents0(自动)并发 agent 上限
maxTotalAgents1000单次 run 总 agent 上限(防死循环)
maxItemsPerCall4096parallel/pipeline 单次调用的 items 上限
syncTimeoutMs5000vm 同步切片超时
disposeGraceMs5000cancel 后的宽限期,到期 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 线程不是安全边界,但提供三层隔离:

  1. 线程隔离 — 同步阻塞不影响 host 事件循环;可 terminate 强杀。
  2. 环境隔离 — scrubbed env 不暴露凭证;空 execArgv 不暴露调试端口。
  3. 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 继续执行。