关联:Issue #2 / PR #6。本文档补充 harness/workflow-orchestrator 的完整设计,供评审与后续迭代参考。
harness/workflow-orchestrator
harness 需要一个可落盘、可恢复、可并发的开发工作流执行器:把一条开发任务(如"实现某算子")拆成若干有依赖关系的子任务,由 Coding Agent 子代理接力完成,并在中途崩溃 / 重启后从断点继续,而不是从头再来。
核心诉求(按优先级):
executor
verifier
depends_on
max_parallel
task_type: subgraph
用户 / 上层 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 文件解耦,各自可由脚本独立调用、独立测试。
orchestrator.py
get_task.py
update_status.py
status.json
log.jsonl
workflow
nodes
normal
subgraph
execute
verify
running
verifying
pass
fail
每个任务在 status.json 中维护 status 与 retries,状态集为:
status
retries
pending ──占坑──▶ running ──执行器回复──▶ executed ──占坑──▶ verifying ▲ │ │ pass/fail 回复 │ 其他回复 │ ┌────────────── pass ─────────────────────────────┤ │ │ │ └──────┴────── fail(预算内,重试) ◀── fail │(原地不动,触发空轮自愈) (预算外 → exhausted → STUCK)
迁移规则(update_status.py,单任务单次迁移):
pending
executed
retries+1
max_retries
回复分类(关键防误判):取最后一条非空行,忽略大小写与结尾标点后必须恰好等于 pass/fail/executed 关键词——宁误判为"其他"(可重试),绝不把 "the fix did not pass" 之类误判成 pass。
"the fix did not pass"
work_dir/.workflow/ 下的产物契约:
work_dir/.workflow/
{workflow, work_dir, user_prompt, tasks}
tasks
{status, retries}
os.replace
user_prompt.md
status.user_prompt
init
update
reset
finish
stuck
crash
sessions/<task>.<phase>.<时间戳>.jsonl
exit=
orchestrator.lock
设计要点:
orchestrator.py 的 loop() 是核心,采用轮次制(每轮:取批次 → 占坑 → 并发执行 → 回收):
loop()
[{task_id, agent, prompt}]
pending→running
executed→verifying
ThreadPoolExecutor
cwd=work_dir
get_task.py 的输出 → 主循环的映射:
[{task_id, agent, prompt}, …]
{"result":"[ALL TASK FINISHED]"}
{"result":"[STUCK]", "reason"}
[]
并发上限:每批派发量 ≤ max_parallel − 在飞数,由 get_task.py 计算,保证任意时刻在飞任务数不超过上限。
max_parallel − 在飞数
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 只在执行阶段出现(仅供参考);最后一行是协议约定,子代理须原样回复关键词,供状态机分类。
goal
acceptance
out_of_scope
approach
subgraph 节点解决"子任务数量运行时才知道"的场景(典型:一个普通节点先生成子工作流文件,子图节点消费它)。
file
work_dir
父id/子id
orchestrator.py 内置 Provider 抽象(build_command + parse_reply),已注册 5 个实现:
Provider
build_command
parse_reply
opencode
opencode run --format json --agent <agent> <prompt>
text
part.text
cannbot
cannbot run --format json --agent <agent> <prompt>
claude
claude -p <prompt> --output-format json --dangerously-skip-permissions --agent <agent>
result
pi
pi -p --mode json --no-session <prompt>
message_end(assistant)
codex
codex exec --json --ephemeral --skip-git-repo-check --dangerously-bypass-approvals-and-sandbox <prompt>
item.completed(agent_message)
接入新 CLI 只需继承 Provider 并注册进 PROVIDERS,主循环与状态机完全复用。
PROVIDERS
校验在 get_task.py 的 validate_workflow / validate_node_set 完成,全量、首个错误即失败:
validate_workflow
validate_node_set
workflow / max_parallel / nodes
id/task_type/title/goal/approach/acceptance/out_of_scope/depends_on/executor/verifier/max_retries/on_exhaust
id/task_type/file/depends_on
id
list[非空 str]
max_retries ≥ 0
on_exhaust ∈ {exit, continue}
max_parallel ≥ 1
设计动机:schema 完全显式、没有默认值——键名写错、漏键都会在派发前暴露,而不是把含糊默认值埋进运行时。
这是设计上投入最多的部分,靠三个机制保证"随便什么时候杀进程都能恢复":
reset_transient
running→pending
verifying→executed
os.kill(pid,0)
主循环中"占坑 / 回收"两次 update_status.py 都在主线程串行执行,从根上避免并发写坏 status.json;子代理进程只读状态、只写自己的会话存档,不触碰 status.json。
0
1
[STUCK]
2
agent 进程非零退出 = 基础设施故障:立即停机 exit 2,任务留在瞬态,重跑本命令自动恢复(§11 的瞬态重置)。
exit 2
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" 模拟进程崩溃)。
tests/test_orchestrator.py
--dry-run
dry_replies.json
{task_id: [回复…]}
"$CRASH"
覆盖的关键路径:单节点 happy path、依赖链完成、独立节点并行、空节点拒绝、缺键拒绝、非法 task_type 拒绝、重试后通过、预算耗尽失败。
task_type
get_task
update_status
on_exhaust: continue
on_exhaust
exit
continue
O_EXCL
需求基础版本已提供
设计方案:workflow-orchestrator 技能
1. 背景与目标
harness 需要一个可落盘、可恢复、可并发的开发工作流执行器:把一条开发任务(如"实现某算子")拆成若干有依赖关系的子任务,由 Coding Agent 子代理接力完成,并在中途崩溃 / 重启后从断点继续,而不是从头再来。
核心诉求(按优先级):
executor子代理执行、verifier子代理对照验收标准裁决,二者接力。depends_on边决定,支持并发派发(max_parallel)。task_type: subgraph)。2. 总体架构
职责划分:
orchestrator.py只负责"主循环 + 进程编排 + 并发 + 崩溃恢复",是唯一的进程级协调者;get_task.py是"纯函数式派发器"(读状态、算批次、无副作用地物化子图);update_status.py是"状态机执行器"(单任务、单次迁移)。三者通过status.json/log.jsonl文件解耦,各自可由脚本独立调用、独立测试。3. 核心概念与名词
workflow(名)/max_parallel(并发上限)/nodes(任务节点列表)normal(执行+验证两阶段)与subgraph(动态扇出容器)status.json的主键,也是depends_on引用的对象execute(执行器)与verify(验证器),一个任务依次经历两个阶段executor/verifier字段),由 Provider 无头调用running/verifying是进程内瞬态(崩溃后可重置);pass/fail是终态(不可回退)4. 状态机设计
每个任务在
status.json中维护status与retries,状态集为:迁移规则(
update_status.py,单任务单次迁移):pendingrunningrunningexecutedexecutedverifyingverifyingpasspassverifyingfailfail+retries+1max_retries决定verifyingverifying(不写回)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.mdstatus.user_prompt指向其绝对路径log.jsonlinit/update/reset/finish/stuck/crash)sessions/<task>.<phase>.<时间戳>.jsonlexit=与 stderrorchestrator.lock设计要点:
status.json写入都是"写临时文件 →os.replace",写半截崩溃不会损坏旧文件,旧状态仍可恢复。log.jsonl只追加,天然免疫并发写坏;是审计与排障的旁路记录。6. 调度主循环(轮次制)
orchestrator.py的loop()是核心,采用轮次制(每轮:取批次 → 占坑 → 并发执行 → 回收):get_task.py,拿到本批可派发任务([{task_id, agent, prompt}])。update_status.py(回复传空串),pending→running/executed→verifying。占坑在主线程串行,避免并发写坏status.json。ThreadPoolExecutor按批次大小并发拉起子代理(cwd=work_dir),stdout 流式写入会话存档。update_status.py用真实回复做状态迁移。单任务更新失败不停机——记录后继续回收其余任务,交下一轮重试/告警。get_task.py的输出 → 主循环的映射:[{task_id, agent, prompt}, …]{"result":"[ALL TASK FINISHED]"}{"result":"[STUCK]", "reason"}[](在飞)并发上限:每批派发量 ≤
max_parallel − 在飞数,由get_task.py计算,保证任意时刻在飞任务数不超过上限。7. 节点提示词拼装
get_task.py把节点字段拼成子代理提示词,顺序固定:设计要点:
goal/acceptance/out_of_scope两个阶段都收到(让执行器提前知道验收线),approach只在执行阶段出现(仅供参考);最后一行是协议约定,子代理须原样回复关键词,供状态机分类。8. 子图动态扇出
subgraph节点解决"子任务数量运行时才知道"的场景(典型:一个普通节点先生成子工作流文件,子图节点消费它)。status.json("分组不是任务"),其有效状态由子任务实时聚合:无子任务 =pending,子任务全pass=pass,否则 =running;不占max_parallel名额。depends_on全 pass 时,读其file(相对work_dir,顶层恰好一个nodes键,子节点均normal),子任务以父id/子id注册进status.json(幂等、只增),按普通节点同规则派发。normal节点)。9. 多 CLI 适配层(Provider)
orchestrator.py内置Provider抽象(build_command+parse_reply),已注册 5 个实现:opencodeopencode run --format json --agent <agent> <prompt>text事件的part.text依序拼接cannbotcannbot run --format json --agent <agent> <prompt>claudeclaude -p <prompt> --output-format json --dangerously-skip-permissions --agent <agent>result字段pipi -p --mode json --no-session <prompt>message_end(assistant)的 text 拼接codexcodex exec --json --ephemeral --skip-git-repo-check --dangerously-bypass-approvals-and-sandbox <prompt>item.completed(agent_message)的 text 拼接接入新 CLI 只需继承
Provider并注册进PROVIDERS,主循环与状态机完全复用。10. Schema 校验(显式、无默认值)
校验在
get_task.py的validate_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 ≥ 0、on_exhaust ∈ {exit, continue}、max_parallel ≥ 1。设计动机:schema 完全显式、没有默认值——键名写错、漏键都会在派发前暴露,而不是把含糊默认值埋进运行时。
11. 崩溃恢复与并发安全
这是设计上投入最多的部分,靠三个机制保证"随便什么时候杀进程都能恢复":
reset_transient把running→pending、verifying→executed——瞬态是进程内的,崩溃滞留会造成"永久在飞 → 死锁",必须复位。orchestrator.lock记录 pid +os.kill(pid,0)存活探测;陈旧锁自动接管;正常退出释放;拒绝双跑(防止两个实例并发写坏状态)。主循环中"占坑 / 回收"两次
update_status.py都在主线程串行执行,从根上避免并发写坏status.json;子代理进程只读状态、只写自己的会话存档,不触碰status.json。12. 错误处理与退出码
0pass1[STUCK](重试耗尽 / 依赖无法满足 / 空节点等)2agent 进程非零退出 = 基础设施故障:立即停机
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)15. 已知限制与待办
on_exhaust: continue语义未完全落地:README 描述"预算耗尽后由on_exhaust裁决(exit终止 /continue让其他分支照常派发)",但当前get_task.py对任何预算耗尽都直接[STUCK]短路(等价于exit),未区分continue分支。若需continue语义,应在调度层区分:耗尽节点保持fail终态、只阻塞其下游,其余分支继续派发。normal节点(§8)。O_EXCL,同微秒级别的双启竞窗未覆盖(双启间隔足够大即安全,代码注释已标注)。