已关闭
feat(harness): add workflow-orchestrator skill #2
sunday创建于  1 天前关闭于  2 小时前
sunday
1 天前 创建

设计方案:workflow-orchestrator 技能

关联:Issue #2 / PR #6。本文档补充 harness/workflow-orchestrator 的完整设计,供评审与后续迭代参考。

1. 背景与目标

harness 需要一个可落盘、可恢复、可并发的开发工作流执行器:把一条开发任务(如"实现某算子")拆成若干有依赖关系的子任务,由 Coding Agent 子代理接力完成,并在中途崩溃 / 重启后从断点继续,而不是从头再来。

核心诉求(按优先级):

  1. 状态持久化:所有进度落盘,重启后进度不丢。
  2. 执行/验证分离:每个任务由 executor 子代理执行、verifier 子代理对照验收标准裁决,二者接力。
  3. DAG 调度:执行顺序由 depends_on 边决定,支持并发派发(max_parallel)。
  4. 运行时动态扇出:子任务在运行期才确定(task_type: subgraph)。
  5. 多 CLI 可插拔:底层 Coding Agent CLI(opencode / claude / codex / pi / cannbot)可替换。
  6. 黑盒化:编排引擎只通过 CLI 调用、只凭输出和退出码评判,不暴露内部实现。

2. 总体架构

用户 / 上层 Agent
   │  /workflow-orchestrator  (SKILL.md —— 收集四项输入 → tmux 后台运行 → 回传日志)
   ▼
tmux 会话 (orchestrator)
   │  python3 orchestrator.py --yaml … --work-dir … --provider … --prompt …
   │
   ├─ init_status.py    首次运行:校验 yaml → 生成 .workflow/status.json + user_prompt.md + log.jsonl
   ├─ get_task.py       每轮:校验 → 物化子图 → 重试裁决 → 产出派发批次
   ├─ update_status.py  状态迁移:按 agent 回复推进 pending→running→…→pass/fail
   │
   ▼
Provider 抽象层 (opencode / cannbot / claude / pi / codex)
   │  build_command + parse_reply,无头调用 CLI
   ▼
executor / verifier 子代理  ── 在 work_dir 内产出真实工件
   │
   ▼
work_dir/.workflow/
   ├─ status.json      任务状态表(进度的唯一权威来源)
   ├─ log.jsonl        只追加事件日志
   ├─ sessions/*.jsonl 各阶段子代理原始对话存档
   └─ orchestrator.lock 单实例锁

职责划分orchestrator.py 只负责"主循环 + 进程编排 + 并发 + 崩溃恢复",是唯一的进程级协调者;get_task.py 是"纯函数式派发器"(读状态、算批次、无副作用地物化子图);update_status.py 是"状态机执行器"(单任务、单次迁移)。三者通过 status.json / log.jsonl 文件解耦,各自可由脚本独立调用、独立测试。

3. 核心概念与名词

概念 含义
workflow 一个 DAG:顶层 workflow(名)/ max_parallel(并发上限)/ nodes(任务节点列表)
任务节点 DAG 顶点,两种类型:normal(执行+验证两阶段)与 subgraph(动态扇出容器)
任务 id 节点全局唯一标识;是 status.json 的主键,也是 depends_on 引用的对象
阶段 phase execute(执行器)与 verify(验证器),一个任务依次经历两个阶段
子代理 agent 每阶段各指定一个 Coding Agent(executor / verifier 字段),由 Provider 无头调用
瞬态 / 终态 running/verifying 是进程内瞬态(崩溃后可重置);pass/fail 是终态(不可回退)

4. 状态机设计

每个任务在 status.json 中维护 statusretries,状态集为:

pending ──占坑──▶ running ──执行器回复──▶ executed ──占坑──▶ verifying
   ▲                                                        │
   │                                      pass/fail 回复     │ 其他回复
   │      ┌────────────── pass ─────────────────────────────┤
   │      │                                                 │
   └──────┴────── fail(预算内,重试) ◀── fail              │(原地不动,触发空轮自愈)
                                (预算外 → exhausted → STUCK)

迁移规则(update_status.py,单任务单次迁移):

当前态 回复 新态 说明
pending 任意 running 占坑,开始执行
running 任意 executed 不校验回复关键词,由 verifier 把关
executed 任意 verifying 占坑,进入验证
verifying pass pass 终态,解锁下游
verifying fail fail + retries+1 是否重试由调度方按 max_retries 决定
verifying 其他 verifying(不写回) 触发"空轮自愈"重验
pass / fail 任意 拒绝更新 终态

回复分类(关键防误判):取最后一条非空行,忽略大小写与结尾标点后必须恰好等于 pass/fail/executed 关键词——宁误判为"其他"(可重试),绝不把 "the fix did not pass" 之类误判成 pass

5. 数据模型与持久化

work_dir/.workflow/ 下的产物契约:

文件 内容 写入方式
status.json {workflow, work_dir, user_prompt, tasks}tasks 以任务 id 为键,值为 {status, retries} 原子写(临时文件 + os.replace),全部由主线程串行写
user_prompt.md 本次开发任务的 prompt 原文;status.user_prompt 指向其绝对路径 初始化时写一次
log.jsonl 只追加事件日志(init / update / reset / finish / stuck / crash append-only
sessions/<task>.<phase>.<时间戳>.jsonl 每阶段子代理 stdout 流式存档;失败时追加 exit= 与 stderr 进程启动前创建
orchestrator.lock 单实例锁(pid + 存活探测) 启动写入、退出释放

设计要点:

  • 原子写:所有 status.json 写入都是"写临时文件 → os.replace",写半截崩溃不会损坏旧文件,旧状态仍可恢复。
  • append-only 日志log.jsonl 只追加,天然免疫并发写坏;是审计与排障的旁路记录。
  • 会话存档:子代理原始 stdout 先落盘再解析,失败信息也保留,供排障基础设施故障。

6. 调度主循环(轮次制)

orchestrator.pyloop() 是核心,采用轮次制(每轮:取批次 → 占坑 → 并发执行 → 回收):

  1. 取批次:调 get_task.py,拿到本批可派发任务([{task_id, agent, prompt}])。
  2. 占坑(主线程串行):对每个任务调 update_status.py(回复传空串),pending→running / executed→verifying。占坑在主线程串行,避免并发写坏 status.json
  3. 并发执行ThreadPoolExecutor 按批次大小并发拉起子代理(cwd=work_dir),stdout 流式写入会话存档。
  4. 回收(主线程串行):逐个调 update_status.py 用真实回复做状态迁移。单任务更新失败不停机——记录后继续回收其余任务,交下一轮重试/告警。

get_task.py 的输出 → 主循环的映射:

get_task 输出 主循环动作
[{task_id, agent, prompt}, …] 进入"占坑→并发→回收"
{"result":"[ALL TASK FINISHED]"} 全部 pass,打印任务表,exit 0
{"result":"[STUCK]", "reason"} 卡死(重试耗尽/依赖无法满足),打印原因,exit 1
[](在飞) 空轮:原地重置瞬态再试;连续 3 轮仍空 → exit 2

并发上限:每批派发量 ≤ max_parallel − 在飞数,由 get_task.py 计算,保证任意时刻在飞任务数不超过上限。

7. 节点提示词拼装

get_task.py 把节点字段拼成子代理提示词,顺序固定:

$WORK_DIR=<work_dir>
$USER_PROMPT=<user_prompt 路径>
[id] title
Goal:            (list → "- " 每项)
Approach:        (仅 execute 阶段)
Acceptance:
Out-of-Scope:
After finishing, reply with only: executed        (execute 阶段)
After verifying, reply with only: pass or fail    (verify 阶段)

设计要点:goal/acceptance/out_of_scope 两个阶段都收到(让执行器提前知道验收线),approach 只在执行阶段出现(仅供参考);最后一行是协议约定,子代理须原样回复关键词,供状态机分类。

8. 子图动态扇出

subgraph 节点解决"子任务数量运行时才知道"的场景(典型:一个普通节点先生成子工作流文件,子图节点消费它)。

  • 容器不入账:subgraph 本身不注册进 status.json("分组不是任务"),其有效状态由子任务实时聚合:无子任务 = pending,子任务全 pass = pass,否则 = running;不占 max_parallel 名额。
  • 幂等物化:subgraph 的 depends_on 全 pass 时,读其 file(相对 work_dir,顶层恰好一个 nodes 键,子节点均 normal),子任务以 父id/子id 注册进 status.json(幂等、只增),按普通节点同规则派发。
  • 聚合解锁:子任务全 pass → subgraph 聚合为 pass,解锁依赖它的下游节点。
  • 限制:不支持嵌套 subgraph(子图内只允许 normal 节点)。

9. 多 CLI 适配层(Provider)

orchestrator.py 内置 Provider 抽象(build_command + parse_reply),已注册 5 个实现:

Provider 命令形态 回复提取
opencode opencode run --format json --agent <agent> <prompt> NDJSON 事件流,text 事件的 part.text 依序拼接
cannbot cannbot run --format json --agent <agent> <prompt> 同上(继承 opencode)
claude claude -p <prompt> --output-format json --dangerously-skip-permissions --agent <agent> 单 JSON 对象的 result 字段
pi pi -p --mode json --no-session <prompt> NDJSON,message_end(assistant) 的 text 拼接
codex codex exec --json --ephemeral --skip-git-repo-check --dangerously-bypass-approvals-and-sandbox <prompt> JSONL,item.completed(agent_message) 的 text 拼接

接入新 CLI 只需继承 Provider 并注册进 PROVIDERS,主循环与状态机完全复用。

10. Schema 校验(显式、无默认值)

校验在 get_task.pyvalidate_workflow / validate_node_set 完成,全量、首个错误即失败

  • 顶层恰好 workflow / max_parallel / nodes 三键,缺/多均报错。
  • normal 节点恰好 12 键(id/task_type/title/goal/approach/acceptance/out_of_scope/depends_on/executor/verifier/max_retries/on_exhaust);subgraph 恰好 4 键(id/task_type/file/depends_on)。
  • id 全局唯一;depends_on 引用存在;依赖成环检测(递归 DFS)。
  • 类型约束:字符串键非空、列表键非空 list[非空 str]max_retries ≥ 0on_exhaust ∈ {exit, continue}max_parallel ≥ 1

设计动机:schema 完全显式、没有默认值——键名写错、漏键都会在派发前暴露,而不是把含糊默认值埋进运行时。

11. 崩溃恢复与并发安全

这是设计上投入最多的部分,靠三个机制保证"随便什么时候杀进程都能恢复":

  1. 瞬态重置:启动时(或空轮自愈时)reset_transientrunning→pendingverifying→executed——瞬态是进程内的,崩溃滞留会造成"永久在飞 → 死锁",必须复位。
  2. 单实例锁orchestrator.lock 记录 pid + os.kill(pid,0) 存活探测;陈旧锁自动接管;正常退出释放;拒绝双跑(防止两个实例并发写坏状态)。
  3. 原子写:见 §5,写半截崩溃不损坏旧文件。

主循环中"占坑 / 回收"两次 update_status.py 都在主线程串行执行,从根上避免并发写坏 status.json;子代理进程只读状态、只写自己的会话存档,不触碰 status.json

12. 错误处理与退出码

退出码 含义 触发
0 工作流完成 全部任务 pass
1 工作流卡死 [STUCK](重试耗尽 / 依赖无法满足 / 空节点等)
2 基础设施故障 / 配置错误 agent 进程非零退出、校验失败、状态读写失败、连续空轮

agent 进程非零退出 = 基础设施故障:立即停机 exit 2,任务留在瞬态,重跑本命令自动恢复(§11 的瞬态重置)。

13. 测试设计

tests/test_orchestrator.py黑盒 UT:只通过 --dry-run 模式调用 orchestrator.py,不读源码,断言退出码、stdout、落盘产物(status.json / log.jsonl)。--dry-run 不拉起真 CLI,回复取自 dry_replies.json{task_id: [回复…]},FIFO 跨 execute/verify 消费;"$CRASH" 模拟进程崩溃)。

覆盖的关键路径:单节点 happy path、依赖链完成、独立节点并行、空节点拒绝、缺键拒绝、非法 task_type 拒绝、重试后通过、预算耗尽失败。

14. 关键设计决策与取舍

决策 理由 代价
轮次制 + 每轮"占坑→并发→回收" 主循环单一、易推理;状态写入集中主线程 单批并发粒度,批次间有串行边界
状态机拆成独立脚本(get_task/update_status 各脚本可独立测试、独立排障;文件解耦 多进程调用,进程启动有开销
原子写(tmp + replace)而非锁内写 崩溃安全,代价极低
回复关键词精确匹配、宁误判 others 防止"误判 pass"造成的静默错误 偶发"其他"回复触发重验
执行器回复不校验、由 verifier 把关 执行器只需"做完",验证器判定质量 执行器回复不当也推进到 executed
黑盒化编排引擎 调用方只依赖 CLI 契约,内部可自由演进 排障只能靠日志与产物

15. 已知限制与待办

  • on_exhaust: continue 语义未完全落地:README 描述"预算耗尽后由 on_exhaust 裁决(exit 终止 / continue 让其他分支照常派发)",但当前 get_task.py 对任何预算耗尽都直接 [STUCK] 短路(等价于 exit),未区分 continue 分支。若需 continue 语义,应在调度层区分:耗尽节点保持 fail 终态、只阻塞其下游,其余分支继续派发。
  • 不支持嵌套 subgraph:子图内只允许 normal 节点(§8)。
  • 依赖成环检测用递归 DFS:超深依赖链(>1000 层)需改迭代(代码注释已标注)。
  • 双启动竞窗:单实例锁未用 O_EXCL,同微秒级别的双启竞窗未覆盖(双启间隔足够大即安全,代码注释已标注)。
likedislike
Ssunday
1 天前 关联了pull request:feat(harness): add workflow-orchestrator skill
Ssunday
1 天前 修改了issue 的描述
sunday
2 小时前 评论:

需求基础版本已提供

likedislike
Ssunday
2 小时前 issue类型由 任务 改变为 需求
Ssunday
2 小时前 issue状态由 待办的 改变为 已解决
Ssunday
2 小时前 关闭了 issue
CANN-robotCANN-robot成员
2 小时前 添加了label:resolved