diff --git a/docs/harness-review-2026-08-16.md b/docs/harness-review-2026-08-16.md new file mode 100644 index 0000000000..4ab307934e --- /dev/null +++ b/docs/harness-review-2026-08-16.md @@ -0,0 +1,327 @@ +# Harness 评审:REJECT 通路与假并行(2026-08-16) + +讨论背景:用户的 harness 底层思想——用并发空间换时间、用总编排与反复监督(loop)换最终代码质量。 +目标:必须使用 DAG 的触发场景、模板优先、ask-matt 路由价值植入检查点、无需 skill 也能模拟 matt skill 工作。 + +## 一、根因确认:为什么从未见过检查点后的 replan + +**是 harness 问题,不是模型一次通过。** 三个叠加的结构性原因: + +1. **积木编译器无法表达纠正波。** `packages/opencode/src/dag/blocks.ts:141` 对任何依赖 review 的积木只生成 + `condition: ${review}.output.verdict == "ACCEPT"`。REJECT 时下游被跳过(`runtime/loop.ts:148` + `condition_false`),不是回环。积木 DSL 里没有「REJECT 回到 coding」的语法。 +2. **REJECT 把工作流终结成 FAILED,而 FAILED 不可变。** review 节点带 REJECT 正常完成 + (`runtime/capture.ts:142-150` 只在 payload/fingerprint 无效时才失败),但 `runtime/loop.ts:328-334` + 的 `unresolvedReviewOutcomes` 触发 `dag.fail("unresolved review outcome(s)")`。`dag.ts:535` 拒绝 + 对终结工作流 replan;`dag.ts:755-763` reopen 只允许 `completed`——「failed and cancelled workflows + are immutable」。结果:`orchestration-policy.md:257-270` Verdict Disposal Contract 的选项 1(extend) + 和选项 2(pause→replan→resume)在 REJECT 之后物理不可达,只剩开新工作流。 +3. **运行时零自动 replan。** `loop.ts` 从不调用 replan;`max_node_replan_attempts` 只是手动 replan 上限。 + +生命周期层**支持**纠正波(`review-lifecycle.ts:173` `isCorrectionReview`;policy 写明 +`REJECT → corrected implementation → verification(PASS) → new diff review`),但只能用手写低层 +`nodes` 表达。`~/.config/opencode/workflows/` 下 14 个 spec 没有一个带纠正波,全部 `review → synthesize` 收尾。 + +### 意外发现:reopen 机制就是为 REJECT 设计的,但够不着 + +`dag.ts:727-733` 注释:被跳过的 dependent(`condition_false`)不算 executed,因此图在检查点处「实质终结」, +自然完成的工作流仍可通过 additive extend 重新打开。REJECT 场景恰好满足全部 reopen 条件: +review 带 REJECT 完成(checkpoint 候选,`report_to_parent` 默认 true)+ synthesize 被跳过 → +`executedDependents = 0` → `hasReportingLeafCheckpoint = true`。唯一阻塞是 `loop.ts` 抢先判 `failed`, +而 `reopenDenial`(`dag.ts:756`)只放行 `completed`。 + +**这不是缺失功能,是一条被自己堵死的既有通路。** + +### harness 内部自相矛盾 + +`orchestration-policy.md:250` 明文:「A gate that successfully returns `REVISE` or `REJECT` is a +**completed node, not a failed node**」。运行时在节点层遵守(`capture.ts`),在工作流层违反 +(`loop.ts:328-334`)。同一条规则两层不一致——建模错误,非取舍。 + +## 二、用户的两个已确认决策 + +1. **REJECT 通路**:选「让 review REJECT 保持 completed」。 +2. **写集**:并行的任务不涉及同文件修改,因此必须并行且同一 worktree;有的「并行」实际上是 + 并行 review / 并行探索——那些本来就是并行的。 + +## 三、方向 A:REJECT 通路 + +### A-1 隐藏耦合:`dag.complete` 自己会抛 + +自然完成调 `dag.complete`,而 `dag.ts:509-512` 在其内部也做 `unresolvedReviewOutcomes` 检查并抛 +`ReviewGateError`(HTTP 处理器 `server/routes/instance/httpapi/handlers/dag.ts:22` 也识别它)。 +只改 loop 不够。守卫有正当用途:拦 agent 用 `control(complete)` 抹平未接受的 review。 + +**方向**:自然完成放行,显式抄近路(agent 的 `control(complete)` 与 HTTP complete)仍被拦。 +守卫语义是「谁在调用」而非「图的状态」。 +**落地实况(#300)**:守卫保留在 `dag.complete` 内,新增 `skipReviewGate` 选项仅由调度循环 +的自然完成路径传入。与「上移工具层」语义等价,且同时守住 HTTP 显式完成面。 + +### A-2 变体:completed vs paused(A1 vs A4) + +| | A1 · complete-on-REJECT | A4 · pause-on-REJECT | +|---|---|---| +| 可用处置 | 只有 extend(reopen 要求 `addsNewNode`,天然只能追加) | 全部四种:extend / pause→replan→resume / 新工作流 / 有理由地停 | +| replan | 不可用(`dag.ts:535` 拒绝终结工作流) | 可用(工作流仍 live) | +| 父对话失职 | 无兜底 | `orchestrator_unresponsive` 看门狗生效(`loop.ts:1229`) | +| 改动量 | 小:两处守卫 | 中:需设计 pause 幂等(全节点已终结的图 resume 会立刻再撞 checkCompletion,须「先注入新节点才允许 resume」) | + +A4 直接服务「反复监督 loop」诉求,且给非 ACCEPT 装上强制处置牙齿——`orchestration-policy.md:270` +自己承认终结检查点逃出了看门狗,A4 补上。 +**建议**:A1 是通往 A4 的安全第一步;先跑通,有真实 loop 手感后再决定是否需要 A4。不要现在决定。 + +### A-3 响度(改成 completed 不会静默) + +review 积木默认 `report_to_parent: true`(`blocks.ts:227`)→ 检查点唤醒父对话,携带 +verdict + findings + required_actions。响度载体是 wake 而非终态状态字。 +需补:`status` 输出里显式列出未解决 verdict(现在只能从终态原因串读到)。 + +## 四、方向 B:并行写集(假并行) + +用户前提成立,但通往它的路上有三个互相咬合的阻塞,只拆一个会直接编译失败。 + +### B-1 序列化本身 + +`serializeWorkspaceWriters`(`blocks.ts:353-373`)无条件把所有 `coding`/`prototype` 串成链。最易改。 + +### B-2 真阻塞:canonical writer 假设 + +`implementationReviewRoute`(`blocks.ts:396-406`)要求唯一「被所有其他 writer 传递依赖」的 writer。 +writer 真并行后该查找返回 `undefined` → 抛 `no canonical serialized implementation writer`。 +**只删 B-1,带 review 的图全部编译不过。** 必须同时解。 + +### B-3 fingerprint 绑定 + +`review.implementation_node_id` 与 `DIFF_REVIEW_SCHEMA` 都假设单个实现节点。N 个并行 writer 时 +「实现指纹」指谁?三方案: + +| 方案 | 做法 | 评价 | +|---|---|---| +| **B3-a 聚合节点**(推荐) | 编译器为 review 路线注入只读汇聚节点,依赖全部 writer,输出 changed_files 并集 + 该点 HEAD 指纹;`implementation_node_id` 指它 | 保持 implementation 与 verification 独立,`validateDiffReview` 不动;多一个廉价节点,聚合过程可观测 | +| B3-b 复用 verify | verify 已被强制依赖每个 writer(`blocks.ts:394-400`),扩展 `VERIFICATION_SCHEMA` 带 fingerprint | 不加节点,但 implementation 与 verification 塌成同一个 | +| B3-c 指纹列表 | `implementation_node_id` 改数组 | 波及 schema、`validateDiffReview`、`reviewEvidenceKeys`、`unresolvedReviewOutcomes`、`isCorrectionReview`、review 契约文案——最大改动最小收益 | + +旁证:`project-development-full.yaml:62` verify 指令已写「Bind results to the exact HEAD +fingerprint」——作者早就按「汇聚点取指纹」思考。B3-a 只是形式化。 + +### B-4 几乎免费的安全网 + +`IMPLEMENTATION_SCHEMA`(`blocks.ts:66-74`)已强制 `changed_files`。运行时可在汇聚点比对各 writer +申报文件集,实际重叠即大声失败。可先「默认并行 + 事后检测」,真踩到再考虑编译期 `write_set` 声明。 + +### B-5 边界:文件不重叠 ≠ 可以并行 + +同 worktree 并行写的共享状态: +- **生成物 / codegen**:`project-development-full.yaml:39-45` 的 `integration-slice` 明确负责 + generated-artifact wiring,即使源文件不重叠也不安全并行。 +- **lockfile / 包管理器**:并发 `bun install` 直接互相破坏。 +- **git HEAD 与暂存区**:writer 若自算指纹会读到别人写到一半的树——B3-a 的独立理由:指纹只在 + 唯一汇聚点计算一次,writer 只管写文件、不碰 git。 + +结论:并行判据是「源文件、生成物、锁文件三者不重叠,且不触发共享构建」,要写进 plan 积木产出契约。 + +### B-6 config 仓库已过时的手写方案 + +`project-development-full.yaml:79-84` DELIVERY CONTRACT(写临时 md 只返回路径)是手工绕开上下文 +膨胀的土办法。运行时 v1.0.15 Train B 已内建 output-by-reference(`runtime/output-ref.ts`:提交 +绝对路径 → 记录 `content_ref`/`size`/`sha256`,报告区 `.opencode/workflow-reports/` 并自动 +gitignore)。该 YAML 段落可删改原生能力;顺带修掉它写 `${TMPDIR}` 而非报告区、绕过完整性收据的问题。 + +## 五、方向 C:tier 判据与检查点处置 + +### C-1 tier 判据保持风险维度 + +不要换成「>2 / <=2 角色」。块数 full 6–9、lite 5–6,区分度低;块数 ≠ 并发数(review 展开 3 节点、 +debug 展开 2);角色数是风险的**结果**而非原因。`workflow-routing.md:50-57` 现有判据(可逆性、 +跨模块、公共契约/并发/持久化/迁移/身份/授权/CI)是原因维度,保留。 + +### C-2 ask-matt 的正确植入位置 + +「路由作为 replan 判据」概念错位:ask-matt 是开局按情境选路;检查点问的是「给定 finding,什么最小波 +能修掉」。选路 ≠ 处置,直接套用会导出「开新工作流」(`orchestration-domains.md:25` 明令禁止)。 +落点:Verdict Disposal Contract(`orchestration-policy.md:257-270`)已有四个选项,缺的是选哪个的判据。 +ask-matt 的情境分类当**处置分类器**: + +| finding 性质 | 处置 | 契约选项 | +|---|---|---| +| 当前 scope 内、有界 | 同图追加纠正波,tier 不变 | 1 extend | +| 揭示跨模块/触边界,lite 形状覆盖不了 | 按 full 形状升级追加保障 lane | 1 extend(升级) | +| 情境判断错了(以为是变更,实为未知成因缺陷) | 唯一该重选路线的场景:debug backbone 新图 | 3 新工作流 | +| 无界 / 迷雾 | 见方向 D | — | + +tier 升级 `workflow-routing.md:59-61` 已在暗示,提升为显式处置表即可。 + +## 六、方向 D:两个结构性缺口 + +### D-1 缺「决策而非交付物」路线 + +现有 7 条路线全是 deliverable 导向,缺「路还看不见」的情形(ask-matt 的 wayfinder 位置,即「深挖 +设计方案」的真实缺口)。`technical-design-full` 假设一个图内产出 implementation-ready design;迷雾期 +需要逐个消解决策点、跨 session 累积决策记录。产物是决策记录,跨图累积,机制需单独设计。 + +### D-2 父对话上下文卫生——直接威胁 loop 主张 + +每次检查点 wake 都往父对话堆上下文;反复 loop 的图会把父对话推出 smart zone(~150k),监督质量 +随轮次单调下降,与「用反复监督换质量」正相反。harness 目前没有这个概念。ask-matt 的 +phase boundaries(continue / clear / handoff / subagent / compact)是现成词汇表。 +原语已有:output-by-reference 让父对话只拿 `content_ref` + 200 字摘要。缺的是把「wave 之间 +父对话该压缩什么」写成契约,不是缺机制。 + +## 七、落地顺序(按解锁关系排序) + +1. **恢复 REJECT 可达性**(本仓库,最小)。守卫安排见 A-1 落地实况(`dag.complete` 保留守卫 + + `skipReviewGate` 自然路径放行);`loop.ts:328-334` 不再 `dag.fail`。验收:REJECT 图断言 + `completed`、reopen-extend 被接受。单独可验证,立刻第一次看到检查点后的 replan。 +2. **决定要不要 A4**。第 1 步跑通后凭真实 loop 手感判断「只能 extend」够不够。唯一应等经验数据的选择。 +3. **并行写集**(本仓库,原子改动):B-1 解序列化 + B-2 汇聚节点替代 canonical writer + B3-a + 汇聚点指纹 + B-4 `changed_files` 重叠检测。拆开编译失败。 +4. **提示词与 spec**(config 仓库,须等 1–3 落地并推进 `runtime-compat.json`): + 处置分类表入 Verdict Disposal Contract;tier 升级显式化;plan 契约加三重不重叠判据; + 删 `project-development-full.yaml:79-84` 手写 DELIVERY CONTRACT;full spec「parallel slices」 + 措辞与第 3 步同步发布(此前不诚实)。 +5. **缺失路线与上下文卫生**(主要 config 仓库):D-1、D-2,独立于前四步。 + +跨仓库铁律:1–3 是运行时(`blocks.ts` / `loop.ts` / `dag.ts` / `review-lifecycle.ts`),4–5 主要 +config 仓库。按 `runtime-compat.json` 策略运行时先合并,同一 PR 推进 SHA 并更新模板。 +第 3 步改变积木编译规则,属于「incompatible block rule」。 + +## 八、关键文件索引 + +| 主题 | 文件 | +|---|---| +| 积木 DSL / 编译器 / 序列化 | `packages/opencode/src/dag/blocks.ts` | +| review 生命周期 / 裁决结算 | `packages/opencode/src/dag/review-lifecycle.ts` | +| 输出捕获与结算 | `packages/opencode/src/dag/runtime/capture.ts` | +| 调度循环 / checkCompletion | `packages/opencode/src/dag/runtime/loop.ts` | +| 工作流状态机 / extend / reopen / replan | `packages/opencode/src/dag/dag.ts` | +| output-by-reference | `packages/opencode/src/dag/runtime/output-ref.ts` | +| 路由器提示词 | `packages/core/src/plugin/command/workflow-routing.md` | +| 积木指南 | `packages/core/src/plugin/command/workflow-blocks.md` | +| 策略 / 裁决处置契约 | `packages/core/src/plugin/command/orchestration-policy.md` | +| 域冲突组合 | `packages/core/src/plugin/command/orchestration-domains.md` | +| 参考 spec 库 | `~/.config/opencode/workflows/*.yaml`(权威源 LeXwDeX/opencode-dag-config) | + +## 九、用户已决策 / 待决策 + +已决策:REJECT 保持 completed(A1 起步);写集必须真并行且同一 worktree(写集不重叠为前提)。 +待决策(有数据后再定):A1 之后是否上 A4(pause + 看门狗)。 +明确不做:把角色数当 tier 判据;把 ask-matt 路由器当 replan 判据(只取情境分类做处置分类器)。 + +--- + +# 附篇:假并行深挖与修复设计(同日) + +## A1 证据盘点 + +1. **假并行纯粹发生在编译期。** 运行时天生支持并行:每个工作流一个信号量 + (`runtime/loop.ts:432-434`、`dag.ts:462-463`,`max_concurrency` 默认 5),就绪节点并行 spawn。 + 唯一拦截者是 `serializeWorkspaceWriters` 在编译期注入的依赖边。**删掉它并行就真实生效, + 运行时零改动。** +2. **所有子会话同一 worktree。** `spawn.ts:425-431` `sessions.create({parentID, title, agent, + model, permission})` 不带 directory 覆盖——子会话继承工作流目录。符合「同一 worktree」前提。 +3. **fingerprint 是 writer 自报的,下游只校验 echo 一致性。** 链路:writer 通过 + `IMPLEMENTATION_SCHEMA` 自报 `fingerprint`(无格式约束,纯提示词契约)→ 编译器用 + `input_mapping` 注入 reviewer 提示词 → reviewer 必须原样 echo + (`review-lifecycle.ts:126-138` `validateReviewResult`)→ `settleCapturedOutput` + (`capture.ts:142-150`)比对相等。它**不是**对工作树的密码学绑定。 + 推论:「聚合指纹」在汇聚点计算一次即可,语义与现状完全兼容。 +4. **`implementationReviewRoute` 的「no canonical serialized implementation writer」抛错 + (`blocks.ts:407-408`)目前是死代码。** 序列化总先构造全序,canonical 查找必然成功。 + 删除序列化后它才可达——改造点正是这里:把 throw 换成聚合器注入。 +5. **review 生命周期全部经由 `review.implementation_node_id` + `input_mapping` 工作** + (`validateDiffReview` `review-lifecycle.ts:197-257`、`reviewEvidenceKeys`、 + `reviewImplementationFingerprint`、`isCorrectionReview`、recovery 路径 + `runtime/recovery.ts:188-190`)。只要聚合节点以 `IMPLEMENTATION_SCHEMA` 形状输出且 + mapping 指向它,**review 生命周期全部零改动**。 +6. **explore agent 有 bash 权限、无写权限**(`agent.ts:200-211`:grep/glob/list/bash/read + allow,`*` deny)——是聚合器 worker 的现成最佳人选(标准 tier,便宜)。 +7. **只有 1 个测试钉死序列化**:`test/dag/blocks.test.ts:252`「serializes unordered workspace + writers...」。`blocks.test.ts:229`(链式 prototype 案例)在保留链行为的设计下不受影响。 + `workflow-authoring.test.ts` 与 `dag-review-lifecycle.test.ts` 无钉死。 + +## A2 修复设计:两种模式,最小差异 + +在 `compileWorkflowBlocks`(`blocks.ts:118-133`)中: + +**模式一:writers 全序(链)→ 行为不变。** 存在唯一「被所有其他 writer 传递依赖」的 writer +时,canonical 逻辑原样保留。链式图(如 `blocks.test.ts:229` 的 coding→prototype→verify→review) +编译结果一字不改。 + +**模式二:writers 无全序(真并行)→ 注入聚合节点。** 把现有的 throw 替换为: + +- 新节点 `${review.id}--aggregate`(沿用 `--` 子节点命名惯例,碰撞时按现有 + `uniqueDuplicates` 抛错);`depends_on` = 全部 writer。 +- `worker_type: explore`(只读 + bash,可跑 `git rev-parse` / 内容哈希);`required: true`; + `report_to_parent: false`。 +- `input_mapping`:逐一绑定每个 writer 的 `output.changed_files`(顺带 `summary`)。 +- `output_schema` 复用 `IMPLEMENTATION_SCHEMA`(union changed_files + 单一 fingerprint)→ + `validateDiffReview` 等所有下游检查零改动。 +- 契约文案:比对各 writer 申报写集,**有交集则不提交、直接失败**(大声报错);无交集则 + 提交并集 + 在汇聚点计算的一次性 fingerprint(对 union 文件内容做哈希,报告所用命令)。 +- **改写 verify 的 writer 依赖为依赖聚合节点**(preserve 非-writer 依赖); + review 的 `implementation_node_id` 指向聚合节点;review `input_mapping` 的三个键 + (`implementation_changed_files` / `implementation_fingerprint` / `verification`)不变, + 仅源节点 ID 变化。 + +## A3 重叠门禁的位置与残余风险 + +- 门禁放在聚合点:intersection of `changed_files` == ∅ 是 writer 契约的机械验证。 + 能抓住:同文件双写、`package.json`/`bun.lock` 双改、同一 codegen 输出双写。 +- 抓不住的残余:不同生成文件但共享构建缓存目录、并发 `bun install` 的锁争用。 + → 写进 plan 积木产出契约与 `workflow-blocks.md`:并行判据 =「源文件、生成物、锁文件 + 三者不重叠,且不触发共享构建」。这是计划纪律,不是引擎能单独兜底的。 + +## A4 附带缺陷(同 PR 可顺手修) + +1. **verify 的指纹契约是空头支票。** `BLOCK_CONTRACTS.verify`(`blocks.ts:111`)写「Bind + evidence to the supplied implementation fingerprint」,但编译器从不给 verify 注入任何 + `input_mapping`——「supplied」并不存在。修复:把聚合/canonical 节点的 + `changed_files`+`fingerprint` 绑进 verify 的 `input_mapping`(verify 已依赖该节点, + 数据依赖与图依赖一致)。 + **落地实况(#299)**:绑定只落在被聚合器改写的 verify 上。canonical 链式路线的 verify + 刻意不绑——绑了会破坏「链式 routes 字节级不变」这一更强的验收项;两条文档主张冲突时 + 字节级不变优先,空头支票在并行路线上已兑现、链式路线留档为已知残余。 +2. **verify FAIL 的终态表现混乱。** verify FAIL → review 被条件跳过(`condition_false`)→ + required review 的 skip 不算 failure → 工作流最终死于 + 「unresolved review outcome(s)」(`loop.ts:328-334`)——原因串误导。属 Direction A + 的 REJECT 通路改造的邻域,一并考虑。 + **落地实况(#300)**:已修复——此类图在检查点自然 completed(不再 fail), + `workflow(action="status")` 以 `unresolved_reviews` 字段显式列出未解决 verdict。 + +## A5 改动清单(原子)与波及面 + +| 文件 | 改动 | +|---|---| +| `src/dag/blocks.ts` | 删 `serializeWorkspaceWriters`;新增聚合预处理 pass(在 `requireValidReviewRoutes` 之前);`implementationReviewRoute` 的 throw 点改为聚合器注入;verify 指纹绑定 | +| `test/dag/blocks.test.ts` | 重写 :252 为并行断言;新增:聚合器形状、用户块名碰撞、部分序(A→B 且 C 独立)、链式图零变化 | +| `packages/core/src/plugin/command/workflow-blocks.md:150-151` | 「compiler serializes unordered writers」改写为并行 + 聚合器 + 三重不重叠判据 | +| config 仓库(后置 PR) | `project-development-full` 的「parallel slices」措辞自此真实;plan 指令加三重不重叠契约;`runtime-compat.json` 推进(属 incompatible block rule) | + +**运行时文件零改动**:`loop.ts` / `dag.ts` / `spawn.ts` / `review-lifecycle.ts` / `capture.ts` / +`recovery.ts` 均不需要动。已持久化的工作流不受影响(编译只发生在 create 时)。 + +## A6 关键判断记录 + +- 为什么聚合器而不是复用 verify 当指纹源:`validateDiffReview:236-238` 要求 verification + 传递依赖 implementation——二者同一节点会直接违反;聚合器保持两者独立,零检查改动。 +- 为什么聚合器用 explore:只读 + bash 恰好覆盖「读树、算哈希、跑 git rev-parse」, + 标准 tier 成本低,无写权限杜绝聚合器自己改工作树。 +- 为什么不做编译期 `write_set` 声明字段:`IMPLEMENTATION_SCHEMA` 已强制 `changed_files`, + 事后机械检测覆盖了大部分场景;声明字段是积木 DSL 的 schema 变更,等真实踩坑再加。 +- 为什么保留链式 canonical:最小差异原则——当前能编译的链式图编译结果完全不变, + 行为变化只发生在「曾经被强行串行」与「曾经编译失败」两类图上。 + +## 附篇 B:落地对照(2026-08-16 review 回填) + +| 文档决策 | 落地载体 | 状态与偏差 | +|---|---|---| +| A-1 守卫安排 | PR #300(merge 6221fdf71) | 机制与 A-1 落地实况一致;语义等价 | +| A-3 status 显式化 | PR #300 | `unresolved_reviews` 字段已实现 | +| A2 并行两模式 + 聚合器 | PR #299(merge ebc5ad089)+ ADR-0002 | 按 A2 全部落地;运行时六文件零改动物证:#299 diff 不触及 runtime/review-lifecycle/dag.ts | +| A4.1 verify 绑定 | PR #299 | 仅聚合路线兑现;链式路线为保字节级不变刻意留白(见 A4 落地实况) | +| A4.2 verify FAIL 终态 | PR #300 | 已修复 | +| B-1~B-4 | PR #299 | 全部落地(重叠门禁为聚合器节点级失败,非编译期检测) | +| B-6 config 仓库 | opencode-dag-config PR #9 | 措辞 + 原生交付 + SHA→ebc5ad089 | +| C-1/C-2 | PR #301(merge 773330b14) | 分类表与 tier 显式化按票 #296 原文落地 | +| D-1/D-2 | — | 刻意未排期,保持结构性缺口记录 | +| (新增)qwen union 双重编码 | PR #298(merge e3e9d76f1)、issue #297 | 调研中新发现的运行时缺陷,不在本文档原始范围 | diff --git a/packages/core/src/plugin/command/orchestration-policy.md b/packages/core/src/plugin/command/orchestration-policy.md index becee83c62..648c4a0c74 100644 --- a/packages/core/src/plugin/command/orchestration-policy.md +++ b/packages/core/src/plugin/command/orchestration-policy.md @@ -267,6 +267,22 @@ of it in the same wake turn with exactly one of: 4. A reasoned stop — tell the user, finding by finding, why no further wave is warranted. Silence is not a stop decision. +Classify the findings first; the class selects the option: + +| Finding nature | Disposal | +| --- | --- | +| Bounded within the current scope | Option 1, same graph, tier unchanged | +| Reveals cross-module, contract, persistence, or boundary risk the current shape cannot cover | Option 1 escalated: append full-shaped assurance lanes (broader review axes, extra verification) instead of the lite correction alone | +| The situation itself was misclassified (a change assumed, an unknown-cause defect found; a repair assumed, a design gap found) | Option 3 only — the single legitimate route switch; start the workflow whose backbone matches the real deliverable and name which prior evidence carries over | +| Unbounded or foggy — findings that no bounded wave can discharge | Option 4, or escalate to the user with a decision request; do not launder fog into a speculative wave | + +Route reselection is never the default: same-objective work stays in one +workflow, and escalation keeps additive semantics — full-shaped assurance +lanes are appended waves, not replacement graphs. The tier discriminator is +risk only (reversibility, module span, public contracts, concurrency, +persistence, migration, identity, authorization, upstream executables, +CI/release). Role or block count never selects a tier. + Merely summarizing a non-ACCEPT verdict and ending the turn is an orchestration failure. The runtime's `orchestrator_unresponsive` guard only fires for workflows that are still live; a checkpoint that terminalizes its diff --git a/packages/core/src/plugin/command/workflow-blocks.md b/packages/core/src/plugin/command/workflow-blocks.md index 59e29148ce..9e5ddfb654 100644 --- a/packages/core/src/plugin/command/workflow-blocks.md +++ b/packages/core/src/plugin/command/workflow-blocks.md @@ -147,9 +147,16 @@ stay quiet. A block immediately after a review gate is conditioned on its accepted verdict. Because the condition language handles one verdict reference, fan multiple review lanes into one review block before continuing. -All block workers share one workspace. The compiler serializes -otherwise-unordered `coding` and `prototype` writers, while read-only lanes may -remain parallel. The resident Orchestration Router owns route selection and +All block workers share one workspace. Unordered `coding` and `prototype` +writers run in parallel, so their work packages must be triple-disjoint: +source files, generated artifacts, and lockfiles must not overlap, and no +shared build may be triggered. The plan block owns this partition. For an +implementation review over parallel writers, the compiler injects one +read-only aggregation node between the writers and the verification gate; it +fails loudly when the declared write sets overlap and otherwise publishes the +union with a single implementation fingerprint computed at the convergence +point. Total-ordered writer chains compile unchanged. Read-only lanes remain +parallel throughout. The resident Orchestration Router owns route selection and phase pruning; this guide owns block fields, contracts, and graph mechanics. ## When to use low-level nodes diff --git a/packages/core/src/plugin/command/workflow-routing.md b/packages/core/src/plugin/command/workflow-routing.md index bdd6fda2ad..878aea801b 100644 --- a/packages/core/src/plugin/command/workflow-routing.md +++ b/packages/core/src/plugin/command/workflow-routing.md @@ -53,12 +53,18 @@ suffice, work is reversible, and no high-risk boundary is involved. Use `full` when any are true: requirements or design are uncertain; work crosses modules or write owners; a public contract, concurrency, persistence, migration, identity, authorization, upstream executable dependencies, CI/release, or -production behavior is in scope. A single matching custom workflow has no tier -to infer: read and retarget it directly. - -If a lite reporting gate returns non-`ACCEPT`, let that graph finish and use -additive `extend` with new node IDs after reassessing the live library. Do not -pause or replan a completed workflow; the parent owns this control decision. +production behavior is in scope. Only these risk dimensions select a tier — +never role count or block count, which are consequences of risk, not causes. +A single matching custom workflow has no tier to infer: read and retarget it +directly. + +If a lite reporting gate returns non-`ACCEPT`, let that graph finish and +dispose of the verdict under the Verdict Disposal Contract. When the findings +cross any `full` criterion above, the correction wave MUST be full-shaped: +escalate by appending full-shaped assurance lanes with new node IDs in the +same workflow — a tier escalation is an additive wave, never a replacement +workflow. Do not pause or replan a completed workflow; the parent owns this +control decision. The primary reference follows the final artifact, not every concern. For code or repairs, review, security, and performance are secondary assurance in that diff --git a/packages/core/test/plugin/command.test.ts b/packages/core/test/plugin/command.test.ts index e49d44aede..253aac387e 100644 --- a/packages/core/test/plugin/command.test.ts +++ b/packages/core/test/plugin/command.test.ts @@ -118,8 +118,10 @@ describe("CommandPlugin.Plugin", () => { expect(CommandPlugin.WorkflowContent).toContain("upstream executable dependencies") expect(CommandPlugin.WorkflowContent).toContain("single matching custom workflow") expect(CommandPlugin.WorkflowContent).toContain("Do not concatenate two complete references") - expect(CommandPlugin.WorkflowContent).toContain("additive `extend` with new node IDs") - expect(CommandPlugin.WorkflowContent).toContain("Do not\npause or replan a completed workflow") + expect(CommandPlugin.WorkflowContent).toContain("never role count or block count, which are consequences of risk") + expect(CommandPlugin.WorkflowContent).toContain("dispose of the verdict under the Verdict Disposal Contract") + expect(CommandPlugin.WorkflowContent).toContain("the correction wave MUST be full-shaped") + expect(CommandPlugin.WorkflowContent).toContain("Do not pause or replan a completed workflow") }), ) @@ -249,6 +251,9 @@ describe("CommandPlugin.Plugin", () => { ) expect(CommandPlugin.OrchestrationPolicyContent).toContain("escapes that guard") expect(CommandPlugin.OrchestrationPolicyContent).toContain("Silence is not a stop decision") + expect(CommandPlugin.OrchestrationPolicyContent).toContain("Classify the findings first; the class selects the option") + expect(CommandPlugin.OrchestrationPolicyContent).toContain("The situation itself was misclassified") + expect(CommandPlugin.OrchestrationPolicyContent).toContain("Role or block count never selects a tier") expect(CommandPlugin.WorkflowFactsContent).toContain("Verdict Disposal Contract") }), ) diff --git a/packages/opencode/src/dag/blocks.ts b/packages/opencode/src/dag/blocks.ts index f1c7989892..db36f4cd53 100644 --- a/packages/opencode/src/dag/blocks.ts +++ b/packages/opencode/src/dag/blocks.ts @@ -120,9 +120,9 @@ export function compileWorkflowBlocks( options: WorkflowBlockCompileOptions = {}, ): NodeConfig[] { requireValidBlockGraph(graph, options) - const blocks = serializeWorkspaceWriters(graph.blocks) - requireValidReviewRoutes(blocks) - const nodes = blocks.flatMap((block) => compileBlock(graph.objective, block, blocks)) + requireValidReviewRoutes(graph.blocks) + const { blocks, aggregations, verifyAggregators } = aggregateParallelWriters(graph.blocks) + const nodes = blocks.flatMap((block) => compileBlock(graph.objective, block, blocks, aggregations, verifyAggregators)) const duplicateNodeIDs = uniqueDuplicates(nodes.map((node) => node.id)) if (duplicateNodeIDs.length > 0) { throw new Error( @@ -132,7 +132,13 @@ export function compileWorkflowBlocks( return nodes } -function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowBlock[]): NodeConfig[] { +function compileBlock( + objective: string, + block: WorkflowBlock, + blocks: WorkflowBlock[], + aggregations: Map, + verifyAggregators: Map, +): NodeConfig[] { const dependencies = block.depends_on ?? [] const required = block.required ?? (block.kind === "plan" || block.kind === "verify" || block.kind === "synthesize") const reviewDependency = dependencies.find( @@ -171,20 +177,22 @@ function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowB } if (block.kind === "review") { - const standardsID = `${block.id}--standards` - const intentID = `${block.id}--intent` - const route = implementationReviewRoute(block, blocks) - const reviewCondition = route ? `${route.verification.id}.output.verdict == "PASS"` : condition + const aggregation = aggregations.get(block.id) + const legacyRoute = aggregation ? undefined : implementationReviewRoute(block, blocks) + const implementationID = aggregation ? aggregation.aggregatorID : legacyRoute?.implementation.id + const verificationID = aggregation ? aggregation.verificationID : legacyRoute?.verification.id + const route = implementationID && verificationID ? { implementationID, verificationID } : undefined + const reviewCondition = route ? `${route.verificationID}.output.verdict == "PASS"` : condition const reviewEvidence = route ? { - implementation_changed_files: `${route.implementation.id}.output.changed_files`, - implementation_fingerprint: `${route.implementation.id}.output.fingerprint`, - verification: `${route.verification.id}.output`, + implementation_changed_files: `${route.implementationID}.output.changed_files`, + implementation_fingerprint: `${route.implementationID}.output.fingerprint`, + verification: `${route.verificationID}.output`, } : undefined - return [ + const lanes = [ node({ - id: standardsID, + id: `${block.id}--standards`, name: `${block.id}: standards review`, workerType: block.worker_type ?? "general", dependencies, @@ -197,7 +205,7 @@ function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowB inputMapping: reviewEvidence, }), node({ - id: intentID, + id: `${block.id}--intent`, name: `${block.id}: intent review`, workerType: block.worker_type ?? "general", dependencies, @@ -213,7 +221,7 @@ function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowB id: block.id, name: `${block.id}: review decision`, workerType: block.worker_type ?? "general", - dependencies: [standardsID, intentID, ...(route ? [route.verification.id] : [])], + dependencies: [`${block.id}--standards`, `${block.id}--intent`, ...(route ? [route.verificationID] : [])], objective, instruction: block.instruction, contract: [ @@ -229,22 +237,45 @@ function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowB inputMapping: route ? { ...reviewEvidence, - standards_review: `${standardsID}.output`, - intent_review: `${intentID}.output`, + standards_review: `${block.id}--standards.output`, + intent_review: `${block.id}--intent.output`, } : undefined, review: route ? { phase: "diff", - implementation_node_id: route.implementation.id, - verification_node_id: route.verification.id, + implementation_node_id: route.implementationID, + verification_node_id: route.verificationID, } : undefined, outputSchema: route ? DIFF_REVIEW_SCHEMA : GENERAL_VERDICT_SCHEMA, }), ] + if (!aggregation) return lanes + return [ + node({ + id: aggregation.aggregatorID, + name: `${block.id}: aggregate parallel implementation evidence`, + workerType: "explore", + dependencies: aggregation.writerIDs, + objective, + contract: AGGREGATOR_CONTRACT, + required: true, + reportToParent: false, + inputMapping: Object.fromEntries( + aggregation.writerIDs.flatMap((writerID: string) => [ + [`${writerID.replace(/-/g, "_")}_changed_files`, `${writerID}.output.changed_files`], + [`${writerID.replace(/-/g, "_")}_summary`, `${writerID}.output.summary`], + ]), + ), + outputSchema: IMPLEMENTATION_SCHEMA, + }), + ...lanes, + ] } + const verifyAggregatorIDs = verifyAggregators.get(block.id) + const verifyAggregator = verifyAggregatorIDs && verifyAggregatorIDs.length > 0 ? verifyAggregatorIDs[0] : undefined return [ node({ id: block.id, @@ -257,6 +288,12 @@ function compileBlock(objective: string, block: WorkflowBlock, blocks: WorkflowB required, reportToParent: block.report_to_parent ?? block.kind === "synthesize", condition, + inputMapping: verifyAggregator + ? { + implementation_changed_files: `${verifyAggregator}.output.changed_files`, + implementation_fingerprint: `${verifyAggregator}.output.fingerprint`, + } + : undefined, outputSchema: WRITER_KINDS.has(block.kind) ? IMPLEMENTATION_SCHEMA : block.kind === "verify" @@ -350,33 +387,75 @@ function requireValidBlockGraph(graph: WorkflowBlockGraph, options: WorkflowBloc topologicalBlocks(graph.blocks) } -function serializeWorkspaceWriters(blocks: WorkflowBlock[]) { - const writers = topologicalBlocks(blocks).filter((block) => WRITER_KINDS.has(block.kind)) - const previousWriter = new Map( - writers.slice(1).map((block, index) => [block.id, writers[index]?.id ?? block.id] as const), - ) - const serialized = blocks.map((block) => { - const previous = previousWriter.get(block.id) - if (!previous || dependsTransitively(blocks, block.id, previous)) return block +// Injected between parallel implementation writers and their verification +// gate: mechanically detects declared write-set overlap (loud node failure) +// and publishes the union with one fingerprint computed at the convergence +// point, so diff review binds to a single post-merge state. +const AGGREGATOR_CONTRACT = + "Collect the supplied changed-file lists and summaries from each parallel implementation writer. If any file path appears in more than one list, do not submit; fail the node naming the exact overlapping paths. Otherwise submit the union of all changed files and one stable fingerprint computed at this convergence point (for example a sha256 over the sorted union of current file contents, reporting the exact commands used). Do not modify any file." + +interface WriterAggregation { + aggregatorID: string + writerIDs: string[] + verificationID: string +} + +function aggregateParallelWriters(blocks: WorkflowBlock[]) { + const aggregations = new Map() + for (const block of blocks) { + if (block.kind !== "review") continue + const topology = reviewWriterTopology(block, blocks) + if (!topology) continue + if (canonicalWriter(topology, blocks)) continue + aggregations.set(block.id, { + aggregatorID: `${block.id}--aggregate`, + writerIDs: topology.implementations.map((writer) => writer.id), + verificationID: topology.verification.id, + }) + } + if (aggregations.size === 0) { + return { blocks, aggregations, verifyAggregators: new Map() } + } + const writerToAggregators = new Map() + for (const aggregation of aggregations.values()) { + for (const writerID of aggregation.writerIDs) { + writerToAggregators.set(writerID, [...(writerToAggregators.get(writerID) ?? []), aggregation.aggregatorID]) + } + } + const aggregatorIDs = new Set([...aggregations.values()].map((aggregation) => aggregation.aggregatorID)) + const verifyAggregators = new Map() + const rewired = blocks.map((block) => { + if (block.kind !== "verify") return block + const original = block.depends_on ?? [] + const replaced = [...new Set(original.flatMap((dependency) => writerToAggregators.get(dependency) ?? [dependency]))] + if (replaced.length === original.length && replaced.every((dependency, index) => dependency === original[index])) { + return block + } + verifyAggregators.set(block.id, replaced.filter((dependency) => aggregatorIDs.has(dependency))) return new WorkflowBlock({ id: block.id, kind: block.kind, - depends_on: [...(block.depends_on ?? []), previous], + depends_on: replaced, instruction: block.instruction, worker_type: block.worker_type, required: block.required, report_to_parent: block.report_to_parent, }) }) - topologicalBlocks(serialized) - return serialized + topologicalBlocks(rewired) + return { blocks: rewired, aggregations, verifyAggregators } } function requireValidReviewRoutes(blocks: WorkflowBlock[]) { - blocks.filter((block) => block.kind === "review").forEach((block) => implementationReviewRoute(block, blocks)) + blocks.filter((block) => block.kind === "review").forEach((block) => reviewWriterTopology(block, blocks)) } -function implementationReviewRoute(block: WorkflowBlock, blocks: WorkflowBlock[]) { +interface ReviewWriterTopology { + implementations: WorkflowBlock[] + verification: WorkflowBlock +} + +function reviewWriterTopology(block: WorkflowBlock, blocks: WorkflowBlock[]): ReviewWriterTopology | undefined { const implementations = blocks.filter( (candidate) => WRITER_KINDS.has(candidate.kind) && dependsTransitively(blocks, block.id, candidate.id), ) @@ -389,8 +468,7 @@ function implementationReviewRoute(block: WorkflowBlock, blocks: WorkflowBlock[] `Implementation review "${block.id}" requires exactly one verification ancestor; found ${verifications.length}`, ) } - const verification = verifications[0] - if (!verification) throw new Error(`Implementation review "${block.id}" has no verification ancestor`) + const verification = verifications[0]! const verifiedImplementations = implementations.filter((candidate) => dependsTransitively(blocks, verification.id, candidate.id), ) @@ -399,15 +477,27 @@ function implementationReviewRoute(block: WorkflowBlock, blocks: WorkflowBlock[] `Implementation review "${block.id}" requires its verification ancestor to depend on every implementation writer`, ) } - const implementation = verifiedImplementations.find((candidate) => - verifiedImplementations.every( + return { implementations: verifiedImplementations, verification } +} + +function canonicalWriter(topology: ReviewWriterTopology, blocks: WorkflowBlock[]): WorkflowBlock | undefined { + return topology.implementations.find((candidate) => + topology.implementations.every( (other) => other.id === candidate.id || dependsTransitively(blocks, candidate.id, other.id), ), ) +} + +function implementationReviewRoute(block: WorkflowBlock, blocks: WorkflowBlock[]) { + const topology = reviewWriterTopology(block, blocks) + if (!topology) return undefined + const implementation = canonicalWriter(topology, blocks) if (!implementation) { + // Unreachable for compiled graphs: aggregateParallelWriters injects an + // aggregator whenever no canonical writer exists. throw new Error(`Implementation review "${block.id}" has no canonical serialized implementation writer`) } - return { implementation, verification } + return { implementation, verification: topology.verification } } function dependsTransitively( diff --git a/packages/opencode/src/dag/dag.ts b/packages/opencode/src/dag/dag.ts index bb0b980538..1c83a45092 100644 --- a/packages/opencode/src/dag/dag.ts +++ b/packages/opencode/src/dag/dag.ts @@ -262,7 +262,7 @@ export interface Interface { readonly resume: (dagID: string) => Effect.Effect readonly step: (dagID: string) => Effect.Effect<{ status: "stepping"; nodeID?: string } | { status: "no_ready_nodes" }, Error> readonly cancel: (dagID: string) => Effect.Effect - readonly complete: (dagID: string) => Effect.Effect + readonly complete: (dagID: string, options?: { readonly skipReviewGate?: boolean }) => Effect.Effect readonly fail: (dagID: string, reason: string) => Effect.Effect readonly replan: (dagID: string, fragment: { nodes: NodeConfig[] }) => Effect.Effect< { cancel: string[]; restart: string[]; replace: string[]; add: string[]; ignore: string[] }, @@ -502,11 +502,19 @@ export const layer = Layer.effect( yield* events.publish(DagEvent.WorkflowCancelled, { dagID: dagID as ID, timestamp: yield* DateTime.now }) yield* terminateNonTerminalNodes(lock, dagID, "workflow_cancelled", "workflow_cancelled", false) }) - const complete = Effect.fn("Dag.complete")(function* (lock: WorkflowLock, dagID: string) { + const complete = Effect.fn("Dag.complete")(function* ( + lock: WorkflowLock, + dagID: string, + options?: { readonly skipReviewGate?: boolean }, + ) { yield* guardWorkflow(dagID, WorkflowStatus.COMPLETED) const workflow = yield* store.getWorkflow(dagID).pipe(Effect.orDie) const config = workflow ? parseWorkflowConfig(workflow.config) : undefined - const unresolvedReviews = config + // The review gate guards EXPLICIT completion (tool/HTTP shortcuts). A + // natural completion at a REJECT checkpoint must stay reachable so the + // parent can dispose of the verdict via reopen-extend (issue #294); the + // scheduling loop passes skipReviewGate for that path. + const unresolvedReviews = !options?.skipReviewGate && config ? unresolvedReviewOutcomes(config, yield* store.getNodes(dagID)) : [] if (unresolvedReviews.length > 0) yield* Effect.fail(new ReviewGateError(dagID, unresolvedReviews)) @@ -898,7 +906,7 @@ export const layer = Layer.effect( resume: (dagID) => withWorkflowLock(dagID)((lock) => resume(lock, dagID)), step: (dagID) => withWorkflowLock(dagID)((lock) => step(lock, dagID)), cancel: (dagID) => withWorkflowLock(dagID)((lock) => cancel(lock, dagID)), - complete: (dagID) => withWorkflowLock(dagID)((lock) => complete(lock, dagID)), + complete: (dagID, options) => withWorkflowLock(dagID)((lock) => complete(lock, dagID, options)), fail: (dagID, reason) => withWorkflowLock(dagID)((lock) => fail(lock, dagID, reason)), replan: (dagID, fragment) => withWorkflowLock(dagID)((lock) => _replan(lock, dagID, fragment)), extend: (dagID, nodes) => withWorkflowLock(dagID)((lock) => _extend(lock, dagID, nodes)), diff --git a/packages/opencode/src/dag/docs/adr/0002-parallel-writers-aggregator.md b/packages/opencode/src/dag/docs/adr/0002-parallel-writers-aggregator.md new file mode 100644 index 0000000000..29f0f50c0f --- /dev/null +++ b/packages/opencode/src/dag/docs/adr/0002-parallel-writers-aggregator.md @@ -0,0 +1,60 @@ +# ADR-0002: Parallel workspace writers with an implementation aggregator + +- Status: Accepted +- Date: 2026-08-16 + +## Context + +The block compiler used to chain every unordered `coding` and `prototype` +writer into one serial lane before the graph reached the runtime. Saved +routes advertised parallel implementation slices that never overlapped in +execution, so the runtime's concurrency budget bought nothing for the phase +where time pressure is highest. The serialization existed because the diff +review gate binds one implementation reference and one fingerprint: with +several independent writers there is no single canonical source of +implementation evidence. + +The fingerprint contract is a worker-reported value verified by echo +(implementation report → reviewer echo → settlement equality), not a +cryptographic binding to the worktree, so centralizing where the fingerprint +is produced does not change what it is. + +## Decision + +Unordered writers compile to truly parallel nodes; the runtime semaphore +schedules them within the workflow's `max_concurrency`. Writer ordering only +exists where the author declared it. + +When an implementation review covers writers with no total order, the +compiler injects one aggregation node per review route: + +- The aggregator depends on every writer of the route, runs read-only with + shell access, is required, and reuses the implementation output schema. +- It receives each writer's declared `changed_files` and fails its node + loudly on any non-empty write-set intersection; otherwise it publishes the + union plus one fingerprint computed at the convergence point. +- The verify block's writer dependencies are re-pointed to the aggregator, + the diff review's implementation reference points at the aggregator, and + the verify node receives the implementation fingerprint binding. + +Writer chains that already have a total order keep the canonical-writer +behavior and compile byte-identically. Overlap of an author-defined block id +with a generated aggregator id is rejected by the existing duplicate-node +check. + +Author discipline for parallel writers is the triple-disjoint rule — source +files, generated artifacts, and lockfiles disjoint, and no shared build — +owned by the plan block's work packages. Mechanical enforcement is the +aggregator's changed-file intersection check; shared-cache and lock-contention +races remain plan discipline. + +## Consequences + +- The runtime, review lifecycle, settlement, and recovery paths are untouched; + they observe ordinary durable nodes with an ordinary implementation schema. +- The empty fingerprint promise in the verify contract is filled by the + verify binding for aggregated routes. +- Compiled graphs for parallel implementation routes gain one node per + review route; node ceilings must account for it. +- Block guide wording and saved-route wording must describe parallel writers + truthfully; the serialization claim is removed. diff --git a/packages/opencode/src/dag/runtime/loop.ts b/packages/opencode/src/dag/runtime/loop.ts index 9076807b7c..a725cdd563 100644 --- a/packages/opencode/src/dag/runtime/loop.ts +++ b/packages/opencode/src/dag/runtime/loop.ts @@ -22,7 +22,6 @@ import { reviewContractForNode, validateReviewExecutionInput, reviewEvidenceKeys, - unresolvedReviewOutcomes, } from "../review-lifecycle" import { Agent } from "@/agent/agent" import { Session } from "@/session/session" @@ -325,14 +324,12 @@ const serviceLayer = Layer.effect( yield* dag.fail(dagID, `required node(s) failed: ${entry.runtime.getRequiredFailures().join(", ")}`) return } - const unresolvedReviews = entry.config - ? unresolvedReviewOutcomes(entry.config, nodes) - : [] - if (unresolvedReviews.length > 0) { - yield* dag.fail(dagID, `unresolved review outcome(s): ${unresolvedReviews.join(", ")}`) - return - } - yield* dag.complete(dagID) + // Unresolved review outcomes terminalize the graph as COMPLETED at + // the checkpoint, not as a failure: the REJECT shape (skipped + // dependents, reporting leaf) is exactly what reopen-extend is + // designed to pick up, and a failed workflow is immutable + // (issue #294). Explicit completion shortcuts keep the review gate. + yield* dag.complete(dagID, { skipReviewGate: true }) }) const checkSessionStatus = makeSessionStatusChecker(sessionSvc) diff --git a/packages/opencode/src/tool/tool.ts b/packages/opencode/src/tool/tool.ts index 2e631b216e..8dafc6ed1d 100644 --- a/packages/opencode/src/tool/tool.ts +++ b/packages/opencode/src/tool/tool.ts @@ -1,5 +1,5 @@ import { PermissionV1 } from "@opencode-ai/core/v1/permission" -import { Effect, Schema } from "effect" +import { Effect, Option, Schema } from "effect" import { SessionV1 } from "@opencode-ai/core/v1/session" import type { JSONSchema7 } from "@ai-sdk/provider" import type { SessionID, MessageID } from "../session/schema" @@ -107,10 +107,13 @@ export type InferDef = /** * The OpenAI tools contract requires `parameters` to be a JSON Schema object. * A root-level combinator (anyOf/oneOf/allOf) is outside that contract: - * OpenAI tolerates it, DeepSeek rejects it with a schema error, and GLM - * silently emits empty tool arguments. Tools that need a discriminated union - * must nest it under a property (e.g. `{ params: }`). Violations fail - * at construction time here instead of degrading at provider runtime. + * OpenAI tolerates it, DeepSeek rejects it with a schema error, GLM silently + * emits empty tool arguments, and qwen-family models string-encode property + * values whose schema is a nested union (issue #297 — repaired by the + * retry pass in the execute wrapper below). Tools that need a discriminated + * union must nest it under a property (e.g. `{ params: }`). + * Violations fail at construction time here instead of degrading at provider + * runtime. */ function assertObjectRootedParameters(id: string, toolInfo: DefWithoutID | { parameters: unknown; jsonSchema?: unknown }) { const root = toolInfo.jsonSchema ?? ToolJsonSchema.fromSchema(toolInfo.parameters as Schema.Top) @@ -162,16 +165,19 @@ function wrap, Result extends Metadat "message.id": ctx.messageID, ...(ctx.callID ? { "tool.call_id": ctx.callID } : {}), } + const invalidArguments = (error: unknown) => + new InvalidArgumentsError({ + tool: id, + detail: toolInfo.formatValidationError ? toolInfo.formatValidationError(error) : String(error), + }) return Effect.gen(function* () { - const decoded = yield* decode(args).pipe( - Effect.mapError( - (error) => - new InvalidArgumentsError({ - tool: id, - detail: toolInfo.formatValidationError ? toolInfo.formatValidationError(error) : String(error), - }), - ), - ) + // Strict decode first; the lenient retry only runs once strict + // decoding already failed, so legitimate string arguments that look + // like JSON are never re-parsed. + const strict = yield* decode(args).pipe(Effect.option) + const decoded = Option.isSome(strict) + ? strict.value + : yield* decode(repairStringifiedContainers(args)).pipe(Effect.mapError(invalidArguments)) const result = yield* execute(decoded as Schema.Schema.Type, ctx) if (result.metadata.truncated !== undefined) { return result @@ -193,6 +199,33 @@ function wrap, Result extends Metadat }) } +// Some models string-encode a tool-argument container whose schema is a +// nested union (qwen family, issue #297): the wire carries {"params": +// "{\"action\": \"list\"}"} instead of a nested object. Re-parse strings that +// sit where a container is expected and let the strict decode judge the +// result. Runs only after the strict decode failed, so plain string +// parameters are never touched. +function repairStringifiedContainers(value: unknown): unknown { + if (typeof value === "string") { + const trimmed = value.trim() + if (!trimmed.startsWith("{") && !trimmed.startsWith("[")) return value + try { + const parsed: unknown = JSON.parse(trimmed) + if (typeof parsed === "object" && parsed !== null) return repairStringifiedContainers(parsed) + } catch { + return value + } + return value + } + if (Array.isArray(value)) return value.map(repairStringifiedContainers) + if (typeof value === "object" && value !== null) { + return Object.fromEntries( + Object.entries(value).map(([key, item]) => [key, repairStringifiedContainers(item)]), + ) + } + return value +} + export function define< Parameters extends Schema.Decoder, Result extends Metadata, diff --git a/packages/opencode/src/tool/workflow.ts b/packages/opencode/src/tool/workflow.ts index 58c42b86de..a592593889 100644 --- a/packages/opencode/src/tool/workflow.ts +++ b/packages/opencode/src/tool/workflow.ts @@ -5,6 +5,7 @@ import { Tool } from "./tool" import { CommandPlugin } from "@opencode-ai/core/plugin/command" import { Effect, Option, Schema } from "effect" import { Dag } from "@/dag/dag" +import { DagReviewLifecycle } from "@/dag/review-lifecycle" import { DagConfig } from "@/dag/config" import { DagWorkflows } from "@/dag/workflows" import { DagModel } from "@/dag/model" @@ -451,6 +452,12 @@ export const WorkflowTool = Tool.define< // stay reachable via the result seam (getNode is unfiltered). const nodes = yield* dag.store.getCurrentNodes(params.workflow_id).pipe(Effect.orDie) const config = Dag.parseWorkflowConfig(workflow.config) + // A completed graph can still carry an unresolved review verdict + // (REJECT checkpoint, issue #294); surface it explicitly instead + // of burying it in a terminal reason string. + const unresolvedReviews = config + ? DagReviewLifecycle.unresolvedReviewOutcomes(config, nodes) + : [] return { title: `Workflow status: ${workflow.title}`, output: JSON.stringify( @@ -460,6 +467,7 @@ export const WorkflowTool = Tool.define< status: workflow.status, session_id: workflow.sessionId, mode: config?.mode ?? "standard", + ...(unresolvedReviews.length > 0 ? { unresolved_reviews: unresolvedReviews } : {}), ...(config?.admission ? { admission: { diff --git a/packages/opencode/test/dag/blocks-parallel-writers.test.ts b/packages/opencode/test/dag/blocks-parallel-writers.test.ts new file mode 100644 index 0000000000..1630a96e43 --- /dev/null +++ b/packages/opencode/test/dag/blocks-parallel-writers.test.ts @@ -0,0 +1,144 @@ +import { describe, expect, it } from "bun:test" +import { DagBlocks } from "@/dag/blocks" + +describe("parallel workspace writers (issue #293)", () => { + const parallelRoute = [ + { id: "plan", kind: "plan" }, + { id: "slice-a", kind: "coding", depends_on: ["plan"] }, + { id: "slice-b", kind: "coding", depends_on: ["plan"] }, + { id: "slice-c", kind: "coding", depends_on: ["plan"] }, + { id: "gates", kind: "verify", depends_on: ["slice-a", "slice-b", "slice-c"] }, + { id: "decision", kind: "review", depends_on: ["gates"] }, + ] as const + + it("keeps unordered writers parallel and injects one aggregation node", () => { + const nodes = DagBlocks.compileWorkflowBlocks({ + objective: "Deliver three disjoint slices", + blocks: [...parallelRoute], + }) + const byID = new Map(nodes.map((node) => [node.id, node])) + + expect(byID.get("slice-a")?.depends_on).toEqual(["plan"]) + expect(byID.get("slice-b")?.depends_on).toEqual(["plan"]) + expect(byID.get("slice-c")?.depends_on).toEqual(["plan"]) + + const aggregate = byID.get("decision--aggregate") + expect(aggregate).toMatchObject({ + worker_type: "explore", + required: true, + report_to_parent: false, + depends_on: ["slice-a", "slice-b", "slice-c"], + }) + expect(aggregate?.output_schema).toEqual( + expect.objectContaining({ required: expect.arrayContaining(["changed_files", "fingerprint"]) }), + ) + expect(aggregate?.input_mapping).toEqual({ + slice_a_changed_files: "slice-a.output.changed_files", + slice_a_summary: "slice-a.output.summary", + slice_b_changed_files: "slice-b.output.changed_files", + slice_b_summary: "slice-b.output.summary", + slice_c_changed_files: "slice-c.output.changed_files", + slice_c_summary: "slice-c.output.summary", + }) + expect(aggregate?.prompt_template.inline).toContain("overlapping paths") + + expect(byID.get("gates")?.depends_on).toEqual(["decision--aggregate"]) + const decision = byID.get("decision") + expect(decision?.review).toEqual({ + phase: "diff", + implementation_node_id: "decision--aggregate", + verification_node_id: "gates", + }) + expect(decision?.input_mapping).toMatchObject({ + implementation_changed_files: "decision--aggregate.output.changed_files", + implementation_fingerprint: "decision--aggregate.output.fingerprint", + verification: "gates.output", + }) + }) + + it("aggregates a partially ordered writer set without losing any writer", () => { + const nodes = DagBlocks.compileWorkflowBlocks({ + objective: "Deliver a partially ordered change", + blocks: [ + { id: "inner", kind: "coding" }, + { id: "outer", kind: "coding", depends_on: ["inner"] }, + { id: "side", kind: "coding" }, + { id: "gates", kind: "verify", depends_on: ["outer", "side"] }, + { id: "decision", kind: "review", depends_on: ["gates"] }, + ], + }) + const byID = new Map(nodes.map((node) => [node.id, node])) + + expect(byID.get("outer")?.depends_on).toEqual(["inner"]) + expect(byID.get("decision--aggregate")?.depends_on).toEqual(["inner", "outer", "side"]) + expect(byID.get("gates")?.depends_on).toEqual(["decision--aggregate"]) + }) + + it("keeps total-ordered writer chains byte-identical without an aggregator", () => { + const nodes = DagBlocks.compileWorkflowBlocks({ + objective: "Deliver a chained change", + blocks: [ + { id: "first", kind: "coding" }, + { id: "second", kind: "coding", depends_on: ["first"] }, + { id: "gates", kind: "verify", depends_on: ["second"] }, + { id: "decision", kind: "review", depends_on: ["gates"] }, + ], + }) + const byID = new Map(nodes.map((node) => [node.id, node])) + + expect(nodes.some((node) => node.id.endsWith("--aggregate"))).toBe(false) + expect(byID.get("gates")?.depends_on).toEqual(["second"]) + expect(byID.get("gates")?.input_mapping).toBeUndefined() + const decision = byID.get("decision") + expect(decision?.review).toEqual({ + phase: "diff", + implementation_node_id: "second", + verification_node_id: "gates", + }) + expect(decision?.input_mapping).toMatchObject({ + implementation_changed_files: "second.output.changed_files", + implementation_fingerprint: "second.output.fingerprint", + }) + }) + + it("rewires verify onto the aggregator while preserving non-writer dependencies", () => { + const nodes = DagBlocks.compileWorkflowBlocks({ + objective: "Deliver with prior evidence", + blocks: [ + { id: "map", kind: "explore" }, + { id: "slice-a", kind: "coding" }, + { id: "slice-b", kind: "coding" }, + { id: "gates", kind: "verify", depends_on: ["map", "slice-a", "slice-b"] }, + { id: "decision", kind: "review", depends_on: ["gates"] }, + ], + }) + const gates = nodes.find((node) => node.id === "gates") + expect(gates?.depends_on).toEqual(["map", "decision--aggregate"]) + }) + + it("binds the implementation fingerprint into the rewired verify node", () => { + const nodes = DagBlocks.compileWorkflowBlocks({ + objective: "Deliver three disjoint slices", + blocks: [...parallelRoute], + }) + expect(nodes.find((node) => node.id === "gates")?.input_mapping).toEqual({ + implementation_changed_files: "decision--aggregate.output.changed_files", + implementation_fingerprint: "decision--aggregate.output.fingerprint", + }) + }) + + it("rejects an author block that collides with an aggregator id", () => { + expect(() => + DagBlocks.compileWorkflowBlocks({ + objective: "Colliding aggregator id", + blocks: [ + { id: "slice-a", kind: "coding" }, + { id: "slice-b", kind: "coding" }, + { id: "decision--aggregate", kind: "verify" }, + { id: "gates", kind: "verify", depends_on: ["slice-a", "slice-b"] }, + { id: "decision", kind: "review", depends_on: ["gates"] }, + ], + }), + ).toThrow("duplicate node ids: decision--aggregate") + }) +}) diff --git a/packages/opencode/test/dag/blocks.test.ts b/packages/opencode/test/dag/blocks.test.ts index 108d20074d..9404fa8423 100644 --- a/packages/opencode/test/dag/blocks.test.ts +++ b/packages/opencode/test/dag/blocks.test.ts @@ -249,7 +249,7 @@ describe("workflow blocks", () => { }) }) - it("serializes unordered workspace writers while leaving read-only lanes parallel", () => { + it("keeps unordered workspace writers parallel and read-only lanes parallel", () => { const nodes = DagBlocks.compileWorkflowBlocks({ objective: "Build two packages from independent evidence", blocks: [ @@ -264,8 +264,8 @@ describe("workflow blocks", () => { expect(nodes.find((node) => node.id === "map-a")?.depends_on).toEqual([]) expect(nodes.find((node) => node.id === "map-b")?.depends_on).toEqual([]) expect(nodes.find((node) => node.id === "package-a")?.depends_on).toEqual(["map-a"]) - expect(nodes.find((node) => node.id === "experiment")?.depends_on).toEqual(["map-b", "package-a"]) - expect(nodes.find((node) => node.id === "package-b")?.depends_on).toEqual(["map-b", "experiment"]) + expect(nodes.find((node) => node.id === "experiment")?.depends_on).toEqual(["map-b"]) + expect(nodes.find((node) => node.id === "package-b")?.depends_on).toEqual(["map-b"]) }) it("routes volume blocks to the standard tier and decision blocks to the advanced tier", () => { @@ -299,6 +299,7 @@ describe("workflow blocks", () => { "diagnose--evidence": "standard", diagnose: "advanced", build: "standard", + "decision--aggregate": "advanced", verify: "advanced", "decision--standards": "standard", "decision--intent": "standard", @@ -312,6 +313,7 @@ describe("workflow blocks", () => { "diagnose--evidence": false, diagnose: true, build: false, + "decision--aggregate": true, verify: true, "decision--standards": false, "decision--intent": false, diff --git a/packages/opencode/test/dag/dag-wake-integration.test.ts b/packages/opencode/test/dag/dag-wake-integration.test.ts index 451d9574e9..fe72bf4a76 100644 --- a/packages/opencode/test/dag/dag-wake-integration.test.ts +++ b/packages/opencode/test/dag/dag-wake-integration.test.ts @@ -1260,7 +1260,10 @@ describe("DagLoop atomic wake integration", () => { ) }) - it.each(["deep", "standard"] as const)("fails a recovered %s workflow when verification skips every diff review", async (mode) => { + // Issue #294: a workflow whose verification skips every diff review settles + // as COMPLETED at the checkpoint (not failed) so the parent can dispose of + // the unresolved review verdict via reopen-extend; failed stays immutable. + it.each(["deep", "standard"] as const)("completes a recovered %s workflow when verification skips every diff review", async (mode) => { await Effect.runPromise( runWakeTest( ({ store, parentPrompts }) => @@ -1271,13 +1274,13 @@ describe("DagLoop atomic wake integration", () => { ), "recovered workflow without an accepted review did not settle", ) - expect(workflow.status).toBe("failed") + expect(workflow.status).toBe("completed") expect((yield* store.getNode(workflow.id, "review-diff"))?.status).toBe("skipped") expect((yield* store.getNode(workflow.id, "final-audit"))?.status).toBe("skipped") - const parent = yield* takeWithin(parentPrompts, "review rejection failure did not wake the parent") + const parent = yield* takeWithin(parentPrompts, "review rejection completion did not wake the parent") expect(promptText(parent.input)).toContain( - '[DAG Workflow failed] Workflow "Recovered review rejection" has reached terminal status.', + '[DAG Workflow completed] Workflow "Recovered review rejection" has reached terminal status.', ) yield* Deferred.succeed(parent.release, "success") }), diff --git a/packages/opencode/test/tool/tool-define.test.ts b/packages/opencode/test/tool/tool-define.test.ts index 8a6afe39de..16c4c90d47 100644 --- a/packages/opencode/test/tool/tool-define.test.ts +++ b/packages/opencode/test/tool/tool-define.test.ts @@ -106,6 +106,96 @@ describe("Tool.define", () => { }), ) + // Regression for #297: qwen-family models string-encode a nested-union + // property value ({"params": "{\"action\": \"list\"}"}). The execute wrap + // must retry with the container re-parsed so the tool still runs. + it.effect("stringified container arguments decode through the lenient retry", () => + Effect.gen(function* () { + const parameters = Schema.Struct({ + params: Schema.Union([ + Schema.Struct({ action: Schema.Literal("list") }), + Schema.Struct({ action: Schema.Literal("validate"), spec_path: Schema.String }), + ]), + }) + const calls: Array> = [] + const info = yield* Tool.define( + "test-repair", + Effect.succeed({ + description: "test tool", + parameters, + parseOptions: { onExcessProperty: "error" }, + execute(args: Schema.Schema.Type) { + calls.push(args) + return Effect.succeed({ title: "test", output: "ok", metadata: { truncated: false } }) + }, + }), + ) + const ctx = makeCtx() + const tool = yield* info.init() + const execute = tool.execute as unknown as (args: unknown, ctx: Tool.Context) => ReturnType + + yield* execute({ params: '{"action": "list"}' }, ctx) + yield* execute({ params: '{"action": "validate", "spec_path": "spec.yaml"}' }, ctx) + yield* execute({ params: { action: "list" } }, ctx) + + expect(calls).toEqual([ + { params: { action: "list" } }, + { params: { action: "validate", spec_path: "spec.yaml" } }, + { params: { action: "list" } }, + ]) + }), + ) + + it.effect("unrepairable arguments still surface as InvalidArgumentsError", () => + Effect.gen(function* () { + const parameters = Schema.Struct({ + params: Schema.Union([Schema.Struct({ action: Schema.Literal("list") })]), + }) + const info = yield* Tool.define( + "test-repair-fail", + Effect.succeed({ + description: "test tool", + parameters, + execute() { + return Effect.succeed({ title: "test", output: "ok", metadata: { truncated: false } }) + }, + }), + ) + const tool = yield* info.init() + const execute = tool.execute as unknown as (args: unknown, ctx: Tool.Context) => ReturnType + + const exit = yield* execute({ params: "not a container" }, makeCtx()).pipe(Effect.exit) + expect(Exit.isFailure(exit)).toBe(true) + if (!Exit.isFailure(exit)) return + const die = exit.cause.reasons.find(Cause.isDieReason) + expect(die?.defect).toBeInstanceOf(Tool.InvalidArgumentsError) + }), + ) + + it.effect("plain string parameters that look like JSON are not re-parsed", () => + Effect.gen(function* () { + const parameters = Schema.Struct({ note: Schema.String }) + const calls: Array> = [] + const info = yield* Tool.define( + "test-string-passthrough", + Effect.succeed({ + description: "test tool", + parameters, + execute(args: Schema.Schema.Type) { + calls.push(args) + return Effect.succeed({ title: "test", output: "ok", metadata: { truncated: false } }) + }, + }), + ) + const tool = yield* info.init() + const execute = tool.execute as unknown as (args: unknown, ctx: Tool.Context) => ReturnType + + yield* execute({ note: '{"kept": "string"}' }, makeCtx()) + + expect(calls).toEqual([{ note: '{"kept": "string"}' }]) + }), + ) + // Regression for #28438: the wrap is the canonical "untyped → typed" boundary. // When the LLM emits a tool call with a payload that fails the parameter // schema, the wrap must surface a typed `Tool.InvalidArgumentsError` whose diff --git a/packages/opencode/test/tool/workflow-authoring.test.ts b/packages/opencode/test/tool/workflow-authoring.test.ts index aea2b2b5b7..431558fee2 100644 --- a/packages/opencode/test/tool/workflow-authoring.test.ts +++ b/packages/opencode/test/tool/workflow-authoring.test.ts @@ -134,24 +134,28 @@ describe("worktree-lifecycle regression fixtures", () => { blocks: [...worktreeLifecycleStartInput.spec.config.blocks], }) const byID = new Map(nodes.map((node) => [node.id, node])) - // Workspace-writer serialization: the second coding writer waits for the - // first even though both only declared the plan dependency. - const writerOrder = nodes.filter((node) => node.worker_type === "build").map((node) => node.id) - expect(writerOrder).toEqual(["coding-worktree-core", "coding-callers-and-fixture"]) - expect(byID.get("coding-callers-and-fixture")?.depends_on).toContain("coding-worktree-core") - // Verification depends on every writer; the review binds to the - // canonical implementation fingerprint. - expect(byID.get("verify")?.depends_on).toEqual( - expect.arrayContaining(["coding-worktree-core", "coding-callers-and-fixture"]), + // Parallel writers: both coding blocks only declared the plan dependency, + // and compilation keeps them unordered — no injected writer-to-writer edge. + expect(byID.get("coding-callers-and-fixture")?.depends_on).toEqual(["plan"]) + expect(byID.get("coding-worktree-core")?.depends_on).toEqual(["plan"]) + // The compiler injects one aggregation node between the parallel writers + // and the verification gate; verify is rewired onto it. + expect(byID.get("verify")?.depends_on).toEqual(["review--aggregate"]) + expect(byID.get("review--aggregate")).toMatchObject({ + worker_type: "explore", + required: true, + depends_on: ["coding-worktree-core", "coding-callers-and-fixture"], + }) + expect(byID.get("review--aggregate")?.output_schema).toEqual( + expect.objectContaining({ required: expect.arrayContaining(["changed_files", "fingerprint"]) }), ) const reviewDecision = byID.get("review") expect(reviewDecision?.review?.phase).toBe("diff") - // Canonical writer is the serialized one that transitively depends on - // every other writer — the second package after serialization. - expect(reviewDecision?.review?.implementation_node_id).toBe("coding-callers-and-fixture") + // The diff review binds to the aggregator's post-merge fingerprint. + expect(reviewDecision?.review?.implementation_node_id).toBe("review--aggregate") expect(reviewDecision?.review?.verification_node_id).toBe("verify") expect(reviewDecision?.input_mapping?.["implementation_fingerprint"]).toBe( - "coding-callers-and-fixture.output.fingerprint", + "review--aggregate.output.fingerprint", ) })