/** * Workflow Extension — lifecycle * * Workflow run 生命周期 free functions(D-12)。 * * 5 个导出函数: * - runWorkflow(spec, deps, signal?) → Promise * - abortRun(runId, deps, reason?, doneReason?) → Promise(done no-op) * - terminateRunningRuns(deps, reason) → Promise(session 切换/关闭终止) * - evictDoneRunsBeyondCap(runs, keepDone) → number(done run 内存淘汰) * - scheduleTimeBudget(runId, deps, budgetTimeMs) → timer(C.7 时间预算) * * 私有 makeHandlers(run, deps) → WorkerHandlers: * - onMessage → handleWorkerMessage(run, raw, deps, handlers) * - onError → handleWorkerError(run, err, deps, handlers) + workerErrorCount++ * - onExit(code, handle) → handleWorkerExit(run, code, handle, deps, handlers) * (G-025:handle.isCurrent 检查内化在 handleWorkerExit 内) * * **A4 原子性**:abort/terminate 内部 transition 先 releaseRuntime(cleanup before * mutate),失败时 status 不变。transition("done") 在 WorkflowRun.transition 内已实现 * 「releaseRuntime → 改 status」原子顺序。 * * **G3-001**(run 一次性生命周期):AbortController 一次性无法复用,runtime 释放后 * 只有两类重建——rebuildRuntime(error-recovery,崩溃重试路径,run 保持 running、 * replaceRuntime 原子换新)与 abort/terminate 的终态释放(transition("done") 内 * releaseRuntime,run 不再恢复)。 * * **D-13**:maxConcurrency=4(ConcurrencyGate 默认值)。 * * 层归属:Engine。依赖 LifecycleDeps + ConcurrencyGate + WorkerHost via port + * WorkflowRun + handleWorker* 函数。 * * 参考:domain-models.md §1(聚合根状态机)。 */ import { getLogger } from "@zhushanwen/pi-extension-logger"; import { validateRunArgs } from "./args-validator.ts"; import { ConcurrencyGate, DEFAULT_CONCURRENCY } from "./concurrency-gate.ts"; import { handleWorkerError, handleWorkerExit, handleWorkerMessage, } from "./error-recovery.ts"; import { Budget } from "./models/budget.ts"; import type { LifecycleDeps, WorkerHandlers } from "./models/ports.ts"; import { RunRuntime } from "./models/run-runtime.ts"; import type { RunSpec } from "./models/run-spec.ts"; import { Trace } from "./models/trace.ts"; import type { DoneReason } from "./models/types.ts"; import { WorkflowRun } from "./models/workflow-run.ts"; import type { WorkerHandle } from "./worker-handle.ts"; const logger = getLogger("subagents"); // ── 常量 ───────────────────────────────────────────────────── /** runId 生成:wf--。 */ const RUNID_RADIX = 36; const RUNID_SLICE_START = 2; const RUNID_SLICE_END = 8; /** * done run 内存保留窗口(K=20)。 * * 本淘汰是 done run 内存有界性的唯一来源:calls.result 不裁、单聚合大小不随 wave1 * 裁剪缩小,故内存上限 = K × 实际聚合大小。同时定义 actionStatus 可查刚完成 run * 的窗口(超出窗口的 done run 不再出现在列表中——已接受的用户可见变化)。 */ export const MAX_RETAINED_DONE_RUNS = 20; function generateRunId(): string { return `wf-${Date.now()}-${Math.random().toString(RUNID_RADIX).slice(RUNID_SLICE_START, RUNID_SLICE_END)}`; } // ── makeHandlers(路由 worker 事件到 error-recovery handle* 函数) ────── /** * 构造 WorkerHandlers——将 worker 的 onMessage/onError/onExit 事件路由到 * error-recovery 的 handleWorker* 函数。 * * 闭包捕获 run + deps。runtime 重建(replaceRuntime)后 run 实例不变、deps 不变, * 故 handlers 对新 worker 仍有效(lifecycle 与 error-recovery 共用 handlers)。 * * **onExit G-025**:handleWorkerExit 内部检查 handle.isCurrent(stale exit 丢弃)。 * 本函数不在 onExit 里重复检查——error-recovery.handleWorkerExit 是单一守卫点。 * * **workerErrorCount**:onError 触发时递增(C.5 跨 runtime 存活的重试计数载体)。 * 注意 handleWorkerError 内部也会递增——这里 onError 递增是 worker 事件层面的 * 「error 事件到达」计数,handleWorkerError 内的是「错误处理决策」计数。 * 实际 handleWorkerError 会做最终计数(含重试上限判断),onError 不重复递增。 */ function makeHandlers(run: WorkflowRun, deps: LifecycleDeps): WorkerHandlers { // 自引用——error-recovery rebuildRuntime 需要 handlers 参数(handlers 引用自身) const handlers: WorkerHandlers = { async onMessage(raw: unknown): Promise { await handleWorkerMessage(run, raw, deps, handlers); }, async onError(err: Error): Promise { await handleWorkerError(run, err, deps, handlers); }, async onExit(code: number, handle: WorkerHandle): Promise { // H-2:用 worker-host 传入的 handle(即真正触发 exit 的那个 handle),而非 // run.runtime?.worker——重试竞态下 runtime.worker 可能已被 replaceRuntime 替换 // 为新 handle,导致 handleWorkerExit 内的 isCurrent 检查误判(漏判 stale exit 或 // 误杀新 worker)。G-025 检查仍在 handleWorkerExit 内(handle.isCurrent)。 await handleWorkerExit(run, code, handle, deps, handlers); }, }; return handlers; } // ── scheduleTimeBudget(C.7 Run 级时间预算调度) ────────── /** * 启动 run 级墙钟时间预算计时器:到期后 abortRun(doneReason="time_limited")。 * * 恢复旧 orchestrator-budget.ts 的 scheduleTimeBudgetCheck 语义——runWorkflow 启动 * 一个 setTimeout(maxTimeMs),到期若 run 仍未终态则转 done,time_limited。 * 计时器存入 RunRuntime.timeBudgetTimer,release(abort/replaceRuntime)时 * 自动清理,避免孤儿触发。worker/script 错误重试经 rebuildRuntime 重排新计时器。 * * @returns 计时器句柄(未设预算时 undefined) */ export function scheduleTimeBudget( runId: string, deps: LifecycleDeps, budgetTimeMs: number, ): ReturnType { const timer = setTimeout(() => { void abortRun(runId, deps, "Time budget exceeded", "time_limited").catch( (err: unknown) => { const msg = err instanceof Error ? err.message : String(err); logger.error(`[workflow] time budget abort failed: ${msg}`); }, ); }, budgetTimeMs); // unref:不阻止 Node 退出(workflow 是后台任务,不应因计时器持有事件循环)。 timer.unref(); return timer; } // ── runWorkflow ────────────────────────────────────────────── /** * 启动一个 workflow run。 * * 流程:创建 WorkflowRun(running,I1 构造期跳过)+ makeHandlers + 构建 RunRuntime * (worker+gate+controller)+ assignRuntime(注入 runtime,恢复 I1)+ 注册到 * deps.runs + store.save。 * * @param spec RunSpec(scriptSource 只读;args 会被原地注入 _runId——rfl C2 契约, * worker 启动与崩溃重建共用同一 args 对象) * @param deps LifecycleDeps(store/workerHost/runner/runs) * @param signal 外部 abort signal(可选;abort 时调 abortRun) * @returns runId(wf--) * @throws signal 已 abort(pre-abort fail fast) */ export async function runWorkflow( spec: RunSpec, deps: LifecycleDeps, signal?: AbortSignal, ): Promise { // m3 E9:参数校验单一 chokepoint,钉在所有副作用前(generateRunId/log/signal // listener/runs.set/workerHost.start/store.save/pending:register)。校验失败时 // zero side effects。coerceTypes 原地规范化 spec.args——worker 启动与崩溃重建 // 共用同一对象(run.spec === spec),恢复路径参数一致。 validateRunArgs(spec); const runId = generateRunId(); // rfl 仪表(tier-1 §7.1):注入稳定 _runId。runAndWait 与 executeNestedWorkflow // 两个 args 入口都经本 choke point;rebuildRuntime 复用 run.spec.args 同一对象 // (error-recovery.ts),worker rebuild 后脚本侧 $ARGS._runId 不漂移——修复 // 「rebuild 回退 run- 导致同一逻辑 run 碎裂到多个 state 目录」。 // 注入在 validateRunArgs 之后,不参与脚本参数 schema 校验(引擎内部字段)。 if (spec.args && typeof spec.args === "object") { spec.args._runId = runId; } deps.log?.("debug", "workflow:lifecycle", "runWorkflow start", { runId, scriptName: spec.scriptName }); // P1-2: pre-aborted signal → fail fast if (signal?.aborted) { throw new Error("Workflow run aborted before start"); } const run = new WorkflowRun( runId, spec, { status: "running", budget: spec.budgetRef ?? new Budget({ maxTokens: spec.budgetTokens, maxTimeMs: spec.budgetTimeMs, }), calls: new Map(), trace: new Trace(), errorLogs: [], }, { startedAt: new Date().toISOString() }, ); // signal abort → abortRun(一次性监听) if (signal) { signal.addEventListener( "abort", () => { void abortRun(runId, deps, "External signal aborted").catch((err: unknown) => { const msg = err instanceof Error ? err.message : String(err); logger.error(`[workflow] abortRun on signal failed: ${msg}`); }); }, { once: true }, ); } // 构造 handlers + runtime(worker + gate + controller) const handlers = makeHandlers(run, deps); const controller = new AbortController(); const gate = new ConcurrencyGate({ maxConcurrency: DEFAULT_CONCURRENCY }); const worker = deps.workerHost.start(spec, spec.args, handlers); // C.7:run 级时间预算计时器(spec.budgetTimeMs > 0 时启用,到期 abortRun time_limited)。 const timeBudgetTimer = spec.budgetTimeMs && spec.budgetTimeMs > 0 ? scheduleTimeBudget(runId, deps, spec.budgetTimeMs) : undefined; const runtime = new RunRuntime(worker, gate, controller, timeBudgetTimer); // assignRuntime(注入 runtime,恢复 I1:running ⟺ runtime!==undefined) run.assignRuntime(runtime); // 注册到 deps.runs(assignRuntime 之后——构造到 assignRuntime 之间 run 处于 // I1 跳过窗口(running 而 runtime undefined),后移保证窗口对外不可见; // worker.start 抛错时 run 未注册,无孤儿 run 残留) deps.runs.set(runId, run); await deps.store.save(run); deps.log?.("debug", "workflow:lifecycle", "run saved", { runId, status: run.state.status }); // pending-notifications: run 启动 → 注册(所有 workflow 启动路径的单一汇聚点: // runAndWait / actionRun / 未来入口全覆盖) deps.log?.("debug", "workflow:lifecycle", "emit pending:register", { runId }); deps.eventBus?.emit("pending:register", { id: runId, type: "workflow", name: spec.slug || spec.scriptName || runId, }); deps.log?.("debug", "workflow:lifecycle", "emit pending:register done", { runId }); return runId; } // ── abortRun ───────────────────────────────────────────────── /** * 中止 workflow(running)。 * * **done 状态 no-op**:已终态的 run 不重复 abort。 * **A4 原子性**:transition("done", doneReason) 内部先 releaseRuntime。 * * @param runId * @param deps * @param reason 可选中止原因(存 run.state.error) * @param doneReason 终态原因(默认 "aborted";超时场景传 "time_limited",C.7) * @throws runId 不存在 */ export async function abortRun( runId: string, deps: LifecycleDeps, reason?: string, doneReason: DoneReason = "aborted", ): Promise { const run = deps.runs.get(runId); if (!run) { throw new Error(`Workflow '${runId}' not found`); } deps.log?.("debug", "workflow:lifecycle", "abortRun", { runId, status: run.state.status, reason, doneReason }); // done 状态 no-op if (run.state.status === "done") { deps.log?.("debug", "workflow:lifecycle", "abortRun no-op: already done", { runId }); return; } // 记录中止原因 if (reason) { run.state.error = reason; } // A4: transition 内部 releaseRuntime(cleanup before mutate) run.transition("done", doneReason); await deps.store.save(run); deps.log?.("debug", "workflow:lifecycle", "abortRun transition done", { runId, reason: run.state.reason }); // C-4: run 到达 done 终态 → 注销 pending-notification + 通知 Interface 层 deps.log?.("debug", "workflow:lifecycle", "emit pending:unregister", { runId, reason: run.state.reason }); deps.eventBus?.emit("pending:unregister", { id: run.runId, reason: run.state.reason ?? "completed" }); deps.log?.("debug", "workflow:lifecycle", "emit pending:unregister done", { runId }); deps.onRunDone?.(run); } // ── terminateRunningRuns(session 切换/关闭:终止全部 running run) ──────── /** * 终止 deps.runs 中全部 running run(session 切换 / session 关闭时调用)。 * * 一次性生命周期(D-2):session 离开当刻,running run 的 token 投入作废,转 * done,failed 持久化落盘——重启后 kill-9 恢复不误判,也不再存在「挂起待恢复」 * 的中间态。 * * per-run 行为:`state.error = reason` → `transition("done","failed")`(内部先 * releaseRuntime,A4)→ `await store.save(run)` → `eventBus.emit("pending:unregister", * {reason:"failed"})`。 * * **不调 deps.onRunDone**:对齐 session_start 恢复先例(index.ts kill-9 恢复只发 * unregister、不发 onRunDone)——session 切换/关闭语境下主 agent 已离开本 session, * 注入完成通知只会把消息发给已离开的 session。 * * **不调 discardInFlightCalls**:run 已转终态不再 replay(无恢复路径),在飞 call * 缓存清不清都不影响结果;该清理仅 rebuildRuntime 需要(崩溃重试会重放脚本, * 假失败结果会污染重跑输出)。 * * 单 run 失败(try/catch + log 带 runId/reason)不中断其余 run——终止是批量收尾, * 一个 run 落盘失败不应放走其余 run 的 failed 状态。 * * @param deps LifecycleDeps(runs/store/eventBus/log) * @param reason 终止原因(写入 run.state.error,如 "Session switched: run terminated") */ export async function terminateRunningRuns( deps: LifecycleDeps, reason: string, ): Promise { for (const run of deps.runs.values()) { if (run.state.status !== "running") continue; try { deps.log?.("debug", "workflow:lifecycle", "terminateRunningRuns", { runId: run.runId, reason }); run.state.error = reason; // A4: transition 内部先 releaseRuntime(cleanup before mutate) run.transition("done", "failed"); await deps.store.save(run); // C-4: run 到达 done 终态 → 注销 pending-notification(reason 固定 "failed") deps.eventBus?.emit("pending:unregister", { id: run.runId, reason: "failed" }); deps.log?.("debug", "workflow:lifecycle", "run terminated", { runId: run.runId, reason: run.state.reason }); } catch (err) { const msg = err instanceof Error ? err.message : String(err); logger.error( `[workflow] terminateRunningRuns failed for run ${run.runId}: ${msg} (reason: ${reason})`, ); } } } // ── evictDoneRunsBeyondCap(done run 内存淘汰,原地裁剪函数) ──────── /** * 淘汰 runs Map 中超出保留窗口的 done run,返回本次淘汰数量。 * * 规则(契约 W3C1): * 1. **状态白名单**:仅 `state.status === "done"` 可淘汰(RunStatus 封闭两态, * 显式白名单而非「非 running」——未来新增状态不落淘汰端)。running * (活跃执行,isScriptRunning 遍历依赖)永不淘汰,即使 completedAt 缺失也 * 绝不参与排序淘汰。 * 2. **排序**:done 项按 `meta.completedAt` ISO 字符串字典序升序(toISOString 恒 * UTC 毫秒格式,字典序=时间序)。completedAt 缺失(防御旧格式/异常快照)fallback * 排序键为空串——字典序最小=最旧,先被淘汰。 * 3. **tie 稳定排序**:比较器三态返回(相等返回 0),Array#sort 稳定性(Node≥12) * 保持元素原序——原序 = Map 插入序 = 创建序,tie 组内先创建者视为更旧先被淘汰 * (kill-9 批量恢复同 ms completedAt 场景的确定性保证)。 * 4. **淘汰执行**:超限数 excess = doneCount - keepDone(<=0 时 no-op 返回 0), * 对升序前 excess 项逐个 `runs.delete(runId)`。 * 5. **边界不变式**:禁止按 Map 插入序直接淘汰——嵌套 workflow 父 run 创建最早、 * 完成最晚,插入序淘汰会在其自身 onRunDone 同步裁剪中淘汰它,runAndWait 轮询 * 窗口内 get 不到 → 误返 "Run not found"。 * 6. **副作用边界**:只清内存 runs Map,不动磁盘 state 文件、不删 * workflow-state-link 指针条目、不发任何事件。 * * @param runs per-session 的 run 注册表(原地裁剪) * @param keepDone done run 保留数(生产传 MAX_RETAINED_DONE_RUNS) * @returns 本次淘汰的 run 数量 */ export function evictDoneRunsBeyondCap( runs: Map, keepDone: number, ): number { // 显式白名单:仅 done 参与淘汰(running 误删即功能破坏) const done = Array.from(runs.values()).filter((r) => r.state.status === "done"); const excess = done.length - keepDone; if (excess <= 0) return 0; // ISO 字典序=时间序;缺失 fallback 空串(ISO 串恒以 '2' 开头非空,空串严格最小=最旧) const keyOf = (r: WorkflowRun): string => r.meta.completedAt ?? ""; // 三态比较器 + sort 稳定性:tie 保持 Array.from 的 Map 插入序(=创建序)—— // 先插入者更旧先淘汰,禁止按插入序直接 slice 淘汰(边界不变式 5) done.sort((a, b) => (keyOf(a) < keyOf(b) ? -1 : keyOf(a) > keyOf(b) ? 1 : 0)); for (const r of done.slice(0, excess)) { runs.delete(r.runId); } return excess; }