oh-my-openagent senpi-task DAG 调度器修复实录:复活节点导致的 run 误终态(dag_2d12c2f7 事故复盘与回归验证)
【免费下载链接】oh-my-openagentOmO: Just type "mass ulw" keyword with your prompt. Now you are the master of graph engineering.项目地址: https://gitcode.com/gh_mirrors/oh/oh-my-openagent
导读
本文基于 .omo/evidence/20260821-dag-revive-terminalization/README.md 及 senpi-task 的 DAG 调度器源码,完整复盘一次真实线上事故:一次三节点 DAG run 在节点 B 被dag send复活、节点 C 仍处于 pending 时,调度器错误地发出了dag.run.completed终态事件。文章逐层还原事故时间线、根因分析、修复方案(共享结算注册表 + 终态断言),并给出可复现的 RED/GREEN 回归测试与完整验证命令,帮助你理解 oh-my-openagent 中 DAG 波次调度、节点复活(revive)与 run 终态判定之间的微妙交互。
事故背景:一次真实的 DAG run 提前终态
2026-08-21 的线上运行dag_2d12c2f7-38f3-487e-bead-2ceafd0ab87e在一个三节点、两边的 DAG 上执行,其定义与最终状态快照保存在 incident-run.json 中:
| 节点 | 依赖 | 角色 |
|---|---|---|
omo-sdk-error-surfacing(记为 A) | 无 | 修复 omo 侧 dag SDK 错误透出 |
senpi-js-localroots(记为 B) | 无 | 修复 senpi 侧 JS kernellocal://路径 |
verify-merges(记为 C) | A、B | 验证两个 PR 均已合并 |
波次结构为:wave 0 = [A, B],wave 1 = [C],C 依赖 A、B 都完成后才能运行。事故发生时,B 已被dag send复活且处于 running 状态,而调度器却在 A 完成的那一刻错误地终结了整个 run。事件日志 incident-events.jsonl 完整记录了当时的执行序列:
seq 1 dag.run.created runId=dag_2d12c2f7-... nodeCount=3 edgeCount=2 seq 2 dag.run.started generation=1 seq 5 dag.wave.started waveIndex=0 [A, B] seq 10 dag.node.transitioned B: running -> failed seq 11 dag.node.steered B delivery=revive seq 12 dag.node.transitioned B: failed -> running (revived) seq 13 dag.node.transitioned A: running -> completed seq 14 dag.wave.completed waveIndex=0 # B 仍在 running! seq 15 dag.run.completed counts={total:3, pending:1, running:1, completed:1}seq 15 的counts字段就是最直观的"病征":total:3而pending:1(C 还等着调度)、running:1(B 刚被复活还在跑)、completed:1(只有 A 完成),这样一个明显未完成的 run 却被标记为completed。事故时刻的快照 incident-run.json 中,verify-merges节点状态为pending、status却已是completed,形成了一副不可能被任何恢复路径继续认领的终态检查点。
根因分析:波次局部的结算表与复活后的"身份丢失"
复盘文档给出的根因指向调度器内部的一个状态所有权问题:
The scheduler's
admitAndSettleWavekept the wave's attached settlement map local to the wave. After B failed, its settlement was removed.dag sendrevived B throughnode-send.ts, which armed a separate watcher but did not register B's new settlement in the active wave map. When A settled, the wave appeared empty andrunWavesemitteddag.run.completeddespite B running and C pending.
拆解来看,问题由三个环节叠加而成:
- 结算表是波次局部的:调度器把"已挂载任务的结算(settlement)"映射放在波次内部。B 失败后,其结算从波次映射中被移除。
- 复活没有回到主结算循环:
dag send复活 B 走的是 node-send.ts 的sendToDagNode路径,它为复活的子任务单独挂了一个 watcher(watchRevivedTask),但没有把 B 的新结算注册进活跃波次的结算映射。 - 空波次触发误终态:A 结算后,波次内部结算集合已空,调度器据此判定"没有活跃任务",进而发出
dag.run.completed——尽管 B 正在运行、C 尚在等待。
值得注意的是,事故文档中提到的admitAndSettleWave/runWaves这些命名,在当前 scheduler.ts 中已被**依赖前沿(dependency-frontier)**执行循环取代:runFrontier不再以波次为执行屏障,编译出的 wave 仅作为信息性事件分组存在(见emitWaveAdmissions/emitCompletedWaves,scheduler.ts)。这正是本次修复重构的底层背景——波次已不再是执行单元,却仍在结算归属上残留着"波次局部"的旧语义,才让空波次钻了空子。
修复方案:调度器生命周期共享的结算注册表与终态断言
修复的核心是把结算注册表的所有权从"波次"提升到"调度器"本身。当前源码中的SchedulerContext持有两个贯穿调度器生命周期的映射:
// scheduler.ts type SchedulerContext = { ... readonly attachedTaskIds: Map<DagNodeId, string> readonly attachedTasks: Map<DagNodeId, AttachedTask> readonly settlementChanged: () => Promise<void> readonly resolveSettlementChanged: () => void ... }attachedTasks记录每个已挂载任务节点的结算 Promise(settled)与折叠完成信号(folded),其生命周期等于调度器实例本身,而不是某个波次。关键路径如下:
1. 复活结算重新进入共享注册表
sendToNode在投递 revive 时,通过watchRevived回调判断当前 run 是否仍处于running:scheduler.ts
sendToNode: (runId, nodeId, message) => sendToDagNode( options, runId, nodeId, message, (revivedNodeId, taskId) => context.journal.refresh().status === "running" ? watchRevivedInScheduler(context, revivedNodeId, taskId) : undefined, ),若 run 仍在运行,则走watchRevivedInScheduler(scheduler.ts)——它把复活的 taskId 写回共享的attachedTaskIds/attachedTasks,并等待结算被主循环折叠:
async function watchRevivedInScheduler( context: SchedulerContext, nodeId: DagNodeId, taskId: string, ): Promise<void> { context.attachedTaskIds.set(nodeId, taskId) const task = attachTaskSettlement(context, nodeId, taskId) await task.folded }attachTaskSettlement(scheduler.ts)在注册结算的同时调用context.resolveSettlementChanged()——这正好唤醒了settleOne中等待的"波次等待者"(wave waiter),使主结算循环立刻重新感知到 B 的存在。
2. 复活结果在折叠完成后才 resolve
settleOne(scheduler.ts)等待attachedTasks中任一结算落地,并在finally中调用task?.resolveFolded()。因此sendToDagNode返回的settledPromise 只有在调度器真正把复活节点的新终态折叠(foldTaskOutcome)之后才 resolve——调用方可以据此知道"复活结果已被 run 主循环吸收"。
3. 终态发射前的全节点终态断言
runFrontier的主循环(scheduler.ts)在发射dag.run.completed/dag.run.failed之前,先要求所有节点处于终态:
const current = context.journal.snapshot() if (current.nodes.every((node) => TERMINAL_NODE_STATES.has(node.state))) break if (context.attachedTasks.size === 0) { if (hasCascadableDependent(current)) continue throw new Error(`DAG run "${current.runId}" cannot terminalize while nodes are active`) } if (!await settleOne(context)) return cancelledSnapshot(context)其中TERMINAL_NODE_STATES定义为completed / failed / cancelled / skipped四态(scheduler.ts)。换句话说:只要还有节点非终态(B 复活后 running、C pending),循环就不会 break,自然也不会进入终态事件发射;即使出现"无挂载任务但节点仍活跃"的反常局面,也会直接抛错而不是静默完成。
4. 已终态 run 的复活保留独立 watcher
若 run 已经进入终态(completed/cancelled)后仍有人对节点send,watchRevived回调返回undefined,sendToDagNode会退回到 node-send.ts 的watchRevivedTask独立 watcher——它通过自己的 control journal 折叠复活结果并持久化产物,与主结算循环解耦。此外 node-send.ts 对已completed/cancelled的 run 直接抛出node_not_continuable,只有failed状态的 run 保留可 revive 语义(复活正是失败节点的文档化补救手段)。
回归测试:A/B/C 场景的完整复现
修复对应的回归用例位于 e2e-failure.test.ts,测试名直接标注了事故编号:
test("#given a live wave with one failed node revived before its sibling settles \ #when the sibling completes #then the run awaits the revival and schedules their dependent", async () => { // given - mirrors dag_2d12c2f7: A and B share wave 0; C depends on both. ... })测试路径与事故完全一致:A、B 同属 wave 0,C 依赖 A、B——B 先失败(runner.settle("B", failed("B"))),再通过与生产相同的sendToNode路径复活(fixture.send(runId, "B", ...),断言sent.delivery === "revive");随后 A 完成,此时必须满足:
- run 状态仍为
running(修复前此处断言失败); - 事件流中不存在任何
dag.run.completed/dag.run.failed终态事件。
接着 B 完成(runner.settle("B", completed("B-revived")))、await sent.settled确认复活已被折叠、C 被调度并完成,最后才允许出现且仅出现一次dag.run.completed,终态节点列表严格为A:completed / B:completed / C:completed。
同文件的第二个用例 L1109-L1141 覆盖"run 已失败后复活驻留子任务"的独立 watcher 路径:failedResult.status === "failed"时仍可 revive,复活任务完成只折叠一次completed转换,且结果产物落盘(output:stuck-revived)。
验证流程:RED → GREEN 与全套证据
事故文档将整个修复过程固化为标准的 TDD 证据链,全部捕获在 .omo/evidence/20260821-dag-revive-terminalization/ 目录下:
| 文件 | 阶段 | 内容 |
|---|---|---|
| red.txt | 修复前 | 聚焦回归测试按预期失败 |
| green.txt | 修复后 | 聚焦回归测试 17 pass / 0 fail |
| focused-final.txt | 聚焦 | scheduler 专项测试全部通过 |
| package-suite-final.txt | 全量 | senpi-task 包级全套测试通过 |
| typecheck-final.txt | 类型检查 | tsgo --noEmit -p tsconfig.json无错误 |
| incident-events.jsonl | 事故证据 | 事故 run 的完整事件日志 |
| incident-run.json | 事故证据 | 事故时刻的检查点快照 |
RED 阶段的关键失败信息位于 red.txt 中:
error: expect(received).toBe(expected) Expected: "running" Received: "completed" at e2e-failure.test.ts:1078:67 (fail) #given a live wave with one failed node revived before its sibling settles ... 16 pass 1 fail即修复前,A 完成后 run 状态被错误置为completed;修复后同一断言通过,且终态事件在全 run 中只出现一次。验证命令为:在仓库根目录运行bun test packages/senpi-task/src/dag/e2e-failure.test.ts(聚焦)、完整包套件(bun test对应文件)以及tsgo --noEmit -p tsconfig.json类型检查。
相关源码与事件模型索引
- 调度器主循环与结算注册表:packages/senpi-task/src/dag/scheduler.ts(
runFrontierL512、settleOneL839、attachTaskSettlementL809、watchRevivedInSchedulerL825、终态集合 L47) - 复活/投递路径:packages/senpi-task/src/dag/node-send.ts(
sendToDagNodeL32、revive 折叠 L102、独立 watcherwatchRevivedTaskL164) - 回归测试:packages/senpi-task/src/dag/e2e-failure.test.ts(L1052、L1109)
- 事件构造器:packages/senpi-task/src/dag/events.ts(
dag.run.completedL75、dag.node.steeredL168) - 节点/状态类型:packages/senpi-task/src/dag/types.ts(
DagNodeState、DagNodeTransitionReason) - 事故证据目录:.omo/evidence/20260821-dag-revive-terminalization/(RED/GREEN 输出、事件日志与快照)
从源码结构看,本次事故的价值在于揭示了 DAG 调度器的一条不变式:run 的终态必须以全部节点的终态为准,而不是以"当前波次是否还有已注册结算"为准。复活(revive)作为失败节点的事实补救手段,必须被纳入主结算循环的可见范围,否则任何"局部结算表 + 外部 watcher"的组合都可能让 run 在仍有活跃节点的前提下被误判完成。这一修复同时为dag.send的调用方明确了语义:返回的settledPromise 即"复活结果已被 run 吸收"的确认信号。
【免费下载链接】oh-my-openagentOmO: Just type "mass ulw" keyword with your prompt. Now you are the master of graph engineering.项目地址: https://gitcode.com/gh_mirrors/oh/oh-my-openagent
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考