363 lines
16 KiB
Python
363 lines
16 KiB
Python
#!/usr/bin/env python3
|
||
"""
|
||
serial-4agent-pipeline.py — 4 Agent 串行流水线(全 2B,一次一个 agent)
|
||
=====================================================================
|
||
牧尘 2026-09-10 定案:复杂任务分 4 阶段串行,每个阶段一个 agent,全用 2B:
|
||
|
||
阶段1 计划 → hermes profile(planner 角色, 2B) 输出 plan.json
|
||
阶段2 执行 → opencode CLI + 2B(harness 补 2B 短板)产出文件
|
||
阶段3 验收 → hermes profile(verifier 角色, 2B) 跑验证 → verdict
|
||
阶段4 总结 → hermes profile(reflector 角色, 2B) 复盘沉淀
|
||
|
||
状态流转:
|
||
plan.json → workdir/ → exec → verify_report.json
|
||
→ pass ? reflect : (rework → 回到 exec 重跑, 最多 N 次)
|
||
|
||
用法:
|
||
python3 serial-4agent-pipeline.py "任务描述" [--workdir /tmp/xxx]
|
||
python3 serial-4agent-pipeline.py "任务" --max-rework 3
|
||
|
||
一次只能有一个 agent 工作(串行),用状态文件防并发。
|
||
"""
|
||
import argparse
|
||
import json
|
||
import os
|
||
import subprocess
|
||
import sys
|
||
import time
|
||
|
||
HERMES = "/home/muc/bin/hermes"
|
||
OPENCODE = "/home/muc/nodejs/node-v24.16.0-linux-x64/bin/opencode"
|
||
HERMES_PROFILE = "default" # 计划/验收/总结跑在哪个 hermes profile(模型 2B 覆盖)
|
||
MODEL_2B = "llama-local/minicpm5-2b"
|
||
ROLES_DIR = os.path.expanduser("~/.hermes/scripts/pipeline-roles")
|
||
LOCK_FILE = "/tmp/serial-4agent.lock"
|
||
|
||
# 模型链:执行阶段用 opencode + 2B 主力 → 智谱 → 商汤(全免费)
|
||
EXEC_MODEL_CHAIN = [
|
||
"llama-local/minicpm5-2b",
|
||
"zhipu/glm-4-flash",
|
||
"sensenova/deepseek-v4-flash",
|
||
]
|
||
|
||
|
||
def sh(cmd, timeout=300, cwd=None, env_extra=None):
|
||
env = dict(os.environ)
|
||
env["PATH"] = "/home/muc/nodejs/node-v24.16.0-linux-x64/bin:" + env.get("PATH", "")
|
||
if env_extra:
|
||
env.update(env_extra)
|
||
try:
|
||
r = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout,
|
||
cwd=cwd, env=env)
|
||
return r.returncode, (r.stdout or "") + ("\n[STDERR] " + r.stderr if r.stderr else "")
|
||
except Exception as e:
|
||
return -1, str(e)[:500]
|
||
|
||
|
||
def acquire_lock():
|
||
if os.path.exists(LOCK_FILE):
|
||
age = time.time() - os.path.getmtime(LOCK_FILE)
|
||
if age < 3600:
|
||
return False, f"已有流水线在跑 (锁文件 {int(age)}s 前创建)"
|
||
os.unlink(LOCK_FILE)
|
||
open(LOCK_FILE, "w").write(str(time.time()))
|
||
return True, ""
|
||
|
||
|
||
def release_lock():
|
||
try:
|
||
os.unlink(LOCK_FILE)
|
||
except Exception:
|
||
pass
|
||
|
||
|
||
def load_role(name):
|
||
path = os.path.join(ROLES_DIR, name)
|
||
return open(path).read() if os.path.exists(path) else ""
|
||
|
||
|
||
def call_hermes_role(role_name, instruction, workdir, timeout=280):
|
||
"""调 hermes profile 扮演角色(2B 模型),返回 stdout"""
|
||
role_prompt = load_role(f"{role_name}.md")
|
||
full = f"{role_prompt}\n\n===== 任务输入 =====\n{instruction}"
|
||
# -m 覆盖模型为 2B(hermes -z 单轮, 不带历史累积)
|
||
rc, out = sh([HERMES, "-p", HERMES_PROFILE, "-m", MODEL_2B, "-z", full],
|
||
timeout=timeout, cwd=workdir)
|
||
return rc, out
|
||
|
||
|
||
def extract_json(text):
|
||
"""从输出剥离 JSON。
|
||
优先取「最后一个完整 JSON 对象」——LLM 常先贴参考内容再输出真正的 JSON。
|
||
"""
|
||
if not text:
|
||
return None
|
||
# 从后往前找平衡的 { }
|
||
depth = 0
|
||
for i in range(len(text) - 1, -1, -1):
|
||
if text[i] == "}":
|
||
depth += 1
|
||
elif text[i] == "{":
|
||
depth -= 1
|
||
if depth == 0:
|
||
cand = text[i:]
|
||
try:
|
||
return json.loads(cand)
|
||
except Exception:
|
||
pass
|
||
# fallback: 找 code fence 里的 JSON
|
||
import re
|
||
m = re.search(r"```(?:json)?\s*(\{.*?\})\s*```", text, re.DOTALL)
|
||
if m:
|
||
try:
|
||
return json.loads(m.group(1))
|
||
except Exception:
|
||
pass
|
||
# fallback: 简单首尾截取
|
||
s, e = text.find("{"), text.rfind("}")
|
||
if s >= 0 and e > s:
|
||
try:
|
||
return json.loads(text[s:e + 1])
|
||
except Exception:
|
||
return None
|
||
return None
|
||
|
||
|
||
def stage_plan(task_desc, workdir):
|
||
"""阶段1: 计划 agent → plan.json"""
|
||
print("\n[1/4] 计划 agent (2B)...", flush=True)
|
||
rc, out = call_hermes_role("01-planner", f"任务: {task_desc}", workdir, timeout=280)
|
||
if rc != 0:
|
||
return False, {"error": f"planner failed: {out[:300]}"}
|
||
plan = extract_json(out)
|
||
if not plan or "steps" not in plan:
|
||
return False, {"error": f"planner no valid JSON: {out[:400]}"}
|
||
with open(os.path.join(workdir, "plan.json"), "w") as f:
|
||
json.dump(plan, f, ensure_ascii=False, indent=2)
|
||
print(f" ✅ 计划产出 {len(plan['steps'])} 步", flush=True)
|
||
for st in plan["steps"]:
|
||
print(f" [{st.get('id')}] {st.get('action','')[:70]}", flush=True)
|
||
return True, plan
|
||
|
||
|
||
def stage_exec(task_desc, plan, workdir, step_override=None):
|
||
"""阶段2: 执行 agent (opencode + 2B)"""
|
||
print("\n[2/4] 执行 agent (opencode+2B)...", flush=True)
|
||
steps = plan.get("steps", [])
|
||
if step_override is not None:
|
||
steps = [s for s in steps if s.get("id") in step_override]
|
||
if not steps:
|
||
return False, {"error": "no steps to execute"}
|
||
# 组装执行指令:给 opencode 完整上下文(任务 + 步骤 + 验收标准)
|
||
exec_prompt = f"任务: {task_desc}\n\n按以下计划执行(每步做完要实际验证):\n"
|
||
for st in steps:
|
||
exec_prompt += (f"\n步骤{st.get('id')}: {st.get('action','')}\n"
|
||
f" 期望产出: {st.get('expected_output','')}\n"
|
||
f" 验收: {st.get('verify','')}\n")
|
||
exec_prompt += ("\n全部完成后输出 JSON: {\"done_steps\": [id...], \"files\": [...], "
|
||
"\"verified\": true/false, \"notes\": \"说明\"}")
|
||
last_err = "no model attempted"
|
||
for model in EXEC_MODEL_CHAIN:
|
||
t0 = time.time()
|
||
# --auto 自动批准权限;--dir 强制 opencode 工作目录(它默认会跑到家目录,忽略 subprocess cwd)
|
||
rc, out = sh([OPENCODE, "run", "--model", model, "--auto", "--dir", workdir,
|
||
exec_prompt],
|
||
timeout=400, cwd=workdir)
|
||
if rc == 0 and out.strip():
|
||
result = {"model": model, "elapsed_s": round(time.time() - t0),
|
||
"raw": out[-3000:]}
|
||
with open(os.path.join(workdir, "exec_result.json"), "w") as f:
|
||
json.dump(result, f, ensure_ascii=False, indent=2)
|
||
print(f" ✅ 执行完成 ({model}, {result['elapsed_s']}s)", flush=True)
|
||
return True, result
|
||
last_err = out[-300:]
|
||
return False, {"error": f"exec all failed: {last_err}"}
|
||
|
||
|
||
def stage_verify(task_desc, plan, workdir):
|
||
"""阶段3: 验收 agent (2B)
|
||
确定性硬检查由脚本做(文件存在/可运行),2B 只判断语义质量。
|
||
🔴 防 2B 幻觉:exec_result.raw 的文本描述不可信,以实际文件检查为准。
|
||
"""
|
||
print("\n[3/4] 验收 agent (2B)...", flush=True)
|
||
|
||
# ===== 硬检查层(确定性,代码执行,不靠 LLM)=====
|
||
hard_issues = []
|
||
steps = plan.get("steps", []) or []
|
||
# 1. 每步 expected_output 声明的文件必须真实存在
|
||
for st in steps:
|
||
exp = st.get("expected_output", "") or ""
|
||
# 提取期望文件路径(粗匹配:找含 .py/.sh/.json/.txt/.md 的路径片段)
|
||
import re
|
||
candidates = re.findall(r'[\w\-./]+\.(?:py|sh|json|txt|md|xlsx?)', exp)
|
||
for c in candidates:
|
||
if not c.startswith("/"):
|
||
c = os.path.join(workdir, c)
|
||
if not os.path.exists(c):
|
||
hard_issues.append({"step_id": st.get("id"),
|
||
"severity": "critical",
|
||
"description": f"期望产出文件不存在: {c}",
|
||
"expected": exp, "actual": "文件缺失"})
|
||
# 2. workdir 里应有非元数据产物(排除 plan/exec/verify/reflect json 和隐藏文件)
|
||
real_files = []
|
||
for f in os.listdir(workdir):
|
||
fp = os.path.join(workdir, f)
|
||
if os.path.isfile(fp) and not f.startswith(".") and f not in (
|
||
"plan.json", "exec_result.json", "verify_report.json",
|
||
"reflect_report.json"):
|
||
real_files.append(f)
|
||
if not real_files and not hard_issues:
|
||
hard_issues.append({"step_id": None, "severity": "critical",
|
||
"description": "工作目录无任何产物文件",
|
||
"expected": "至少一个产物", "actual": "空目录"})
|
||
# 3. 如果有 .py 产物且是脚本任务,跑 py_compile 验证语法
|
||
for f in real_files:
|
||
if f.endswith(".py"):
|
||
fp = os.path.join(workdir, f)
|
||
rc, _ = sh(["python3", "-m", "py_compile", fp], timeout=30, cwd=workdir)
|
||
if rc != 0:
|
||
hard_issues.append({"step_id": None, "severity": "critical",
|
||
"description": f"{f} 语法错误",
|
||
"expected": "py_compile 通过", "actual": f"exit {rc}"})
|
||
|
||
if hard_issues:
|
||
report = {"verdict": "rework", "checked_steps": [],
|
||
"issues": hard_issues,
|
||
"rework_instructions":
|
||
"硬检查未通过:期望产物缺失/语法错误。请真实创建文件并确保可运行,"
|
||
"不要只在文本里声称已写。"}
|
||
with open(os.path.join(workdir, "verify_report.json"), "w") as f:
|
||
json.dump(report, f, ensure_ascii=False, indent=2)
|
||
print(" verdict: rework (硬检查失败, 无需 LLM)", flush=True)
|
||
for i in hard_issues[:5]:
|
||
print(f" 🔴 {i['description'][:80]}", flush=True)
|
||
return True, report
|
||
|
||
# ===== 语义检查层(2B 判断质量)=====
|
||
files_ctx = []
|
||
for f in real_files:
|
||
fp = os.path.join(workdir, f)
|
||
size = os.path.getsize(fp)
|
||
files_ctx.append({"name": f, "size": size,
|
||
"head": open(fp).read()[:800] if size < 5000 else "(大文件略)"})
|
||
verify_ctx = json.dumps({
|
||
"task": task_desc,
|
||
"plan_steps": steps,
|
||
"hard_check": "✅ 全部通过(文件存在/语法 OK)",
|
||
"produced_files": files_ctx,
|
||
"note": "硬检查已过,你只需判断语义质量:逻辑是否正确、是否满足任务意图。"
|
||
"exec_result.raw 是执行 agent 的自述,不可全信,以 produced_files 内容为准。",
|
||
}, ensure_ascii=False)
|
||
rc, out = call_hermes_role("03-verifier", verify_ctx, workdir, timeout=280)
|
||
if rc != 0:
|
||
return False, {"error": f"verifier failed: {out[:300]}"}
|
||
report = extract_json(out)
|
||
if not report or "verdict" not in report:
|
||
return False, {"error": f"verifier no valid JSON: {out[:400]}"}
|
||
with open(os.path.join(workdir, "verify_report.json"), "w") as f:
|
||
json.dump(report, f, ensure_ascii=False, indent=2)
|
||
verdict = report.get("verdict", "rework")
|
||
print(f" verdict: {verdict} (语义检查)", flush=True)
|
||
if verdict == "rework":
|
||
for i in report.get("issues", [])[:5]:
|
||
print(f" ⚠️ step{i.get('step_id')} [{i.get('severity')}] "
|
||
f"{i.get('description','')[:60]}", flush=True)
|
||
return True, report
|
||
|
||
|
||
def stage_reflect(task_desc, plan, workdir):
|
||
"""阶段4: 总结 agent (2B) → reflect_report.json"""
|
||
print("\n[4/4] 总结提升 agent (2B)...", flush=True)
|
||
ctx = {"task": task_desc, "plan": plan}
|
||
for fname in ("exec_result.json", "verify_report.json"):
|
||
fp = os.path.join(workdir, fname)
|
||
if os.path.exists(fp):
|
||
ctx[fname.replace(".json", "")] = json.load(open(fp))
|
||
rc, out = call_hermes_role("04-reflector", json.dumps(ctx, ensure_ascii=False),
|
||
workdir, timeout=280)
|
||
if rc != 0:
|
||
return False, {"error": f"reflector failed: {out[:300]}"}
|
||
report = extract_json(out)
|
||
if not report:
|
||
return False, {"error": f"reflector no valid JSON: {out[:400]}"}
|
||
with open(os.path.join(workdir, "reflect_report.json"), "w") as f:
|
||
json.dump(report, f, ensure_ascii=False, indent=2)
|
||
print(f" ✅ 总结完成: {str(report.get('task_result',''))[:80]}", flush=True)
|
||
for l in report.get("lessons", [])[:5]:
|
||
print(f" 📌 {l[:70]}", flush=True)
|
||
return True, report
|
||
|
||
|
||
def main():
|
||
ap = argparse.ArgumentParser()
|
||
ap.add_argument("task", help="任务描述")
|
||
ap.add_argument("--workdir", default="/tmp/serial-4agent-work")
|
||
ap.add_argument("--max-rework", type=int, default=2)
|
||
ap.add_argument("--skip-exec", action="store_true", help="调试: 跳过执行阶段")
|
||
args = ap.parse_args()
|
||
|
||
ok, msg = acquire_lock()
|
||
if not ok:
|
||
print(json.dumps({"ok": False, "error": msg}))
|
||
return
|
||
os.makedirs(args.workdir, exist_ok=True)
|
||
|
||
try:
|
||
# 阶段1: 计划
|
||
ok1, plan = stage_plan(args.task, args.workdir)
|
||
if not ok1:
|
||
print(json.dumps({"ok": False, "stage": "plan", **plan}))
|
||
return
|
||
if not args.skip_exec:
|
||
# 阶段2: 执行 (rework 循环)
|
||
rework_count = 0
|
||
exec_ok, exec_res = stage_exec(args.task, plan, args.workdir)
|
||
if not exec_ok:
|
||
print(json.dumps({"ok": False, "stage": "exec", **exec_res}))
|
||
return
|
||
# 阶段3: 验收 (rework 循环)
|
||
while True:
|
||
ver_ok, ver_res = stage_verify(args.task, plan, args.workdir)
|
||
if not ver_ok:
|
||
print(json.dumps({"ok": False, "stage": "verify", **ver_res}))
|
||
return
|
||
if ver_res.get("verdict") == "pass":
|
||
break
|
||
rework_count += 1
|
||
if rework_count > args.max_rework:
|
||
print(json.dumps({"ok": False, "stage": "rework-exhausted",
|
||
"msg": f"超过 {args.max_rework} 次返修上限"}))
|
||
return
|
||
# 返修:提取 rework 指令 → 只重跑失败步骤
|
||
instr = ver_res.get("rework_instructions", "") or ""
|
||
issues = ver_res.get("issues", []) or []
|
||
failed_steps = [i.get("step_id") for i in issues
|
||
if isinstance(i, dict) and i.get("severity") == "critical"]
|
||
print(f" ↻ 返修 #{rework_count}: 重跑步骤 {failed_steps or '全部'}", flush=True)
|
||
# 更新 plan 的 action 附返修指令
|
||
plan2 = dict(plan)
|
||
steps2 = plan2.get("steps", []) or []
|
||
for st in steps2:
|
||
if isinstance(st, dict) and st.get("id") in failed_steps:
|
||
st["action"] = f"{st.get('action','')} [返修要求: {instr}]"
|
||
exec_ok, exec_res = stage_exec(args.task, plan2, args.workdir)
|
||
if not exec_ok:
|
||
print(json.dumps({"ok": False, "stage": "exec-rework", **exec_res}))
|
||
return
|
||
# 阶段4: 总结
|
||
ok4, ref = stage_reflect(args.task, plan, args.workdir)
|
||
if not ok4:
|
||
print(json.dumps({"ok": False, "stage": "reflect", **ref}))
|
||
return
|
||
print("\n=== 流水线完成 ===")
|
||
print(json.dumps({"ok": True, "workdir": args.workdir,
|
||
"verdict": "pass",
|
||
"reflect": ref.get("task_result", "")[:200]},
|
||
ensure_ascii=False))
|
||
finally:
|
||
release_lock()
|
||
|
||
|
||
if __name__ == "__main__":
|
||
main()
|