Files
gongzhiyongandClaude Opus 4.8 1288fd19d7
CI / tests (push) Failing after 15m3s
CI / guardrails (push) Failing after 15m3s
feat(swarm): 协作聚合收敛取代蜂后选优 + sandbox 用 pytest 验证
按 juejin 协作聚合模型重构收敛(取代 best-of-N 选优):
- 删蜂后选优(queen.py/test-queen.py)
- 新增聚合节点 result_aggregator.py:共享池收集→同文件 LLM/AST 整合→沙箱验证→单次落 main
- 质量驱动闭环:不达标打回迭代(AGGREGATE_ACCEPTANCE_THRESHOLD + MAX_REVIEW_CYCLES)
- agent 停 git 工作分支,产出走 task.result.files 共享池(AGENT_GIT_PUSH_ENABLED 默认 false)
- sandbox_runner 改用 pytest(原生支持 pytest 风格 class),修 stdlib runner 收集失败
- 文档同步重写为协作聚合模型

本地验证:产物仓单分支 main + 三函数完整 + pytest 12/12 pass_rate=100 一次达标。
影响:Swarm 收敛/聚合层;Manager/客户端契约不变(artifact字段/sequence/状态机;契约测试全过)。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-20 22:03:00 +08:00

199 lines
8.9 KiB
Python

"""北极星采集器 CLI — 从运行中的 orchestrator 采集一次 swarm run 的真实 metrics,持久化到本地。
设计:纯 stdlib + benchmark.metrics 公式,**不 import orchestrator 重依赖**,宿主直接跑、直接落桌面,
无需重构建镜像。数据源 = orchestrator HTTP API(tasks/events/summary)。
诚实(组织规则 #9):只采有真实数据源的 metrics;无数据的标 NaN + coverage False,绝不伪造 0/100。
持久化:① sqlite(~/Desktop/swarm-benchmark.db,可累积/可查询)② JSON(~/Desktop/)。
涌现增益 s_gain = Q_swarm - Q_base:给 --baseline 时用两个 run 的质量分(--q-swarm/--q-base 优先,
否则回退 completion 作代理并显式标注 proxy)。
用法:
python -m benchmark.collect_cli <deployment_id> \
[--baseline <deployment_id>] [--q-swarm 80] [--q-base 60] \
[--orchestrator http://localhost:8000] [--label humaneval-bon] [--db PATH]
"""
from __future__ import annotations
import argparse
import json
import math
import os
import sqlite3
import time
import urllib.request
from collections import Counter
from .metrics import (
completion_score, collaboration_score, robustness_score, communication_score,
cost_score, governance_score, emergence_gain,
)
def _get(base: str, path: str) -> dict:
req = urllib.request.Request(base.rstrip("/") + path)
with urllib.request.urlopen(req, timeout=20) as r:
return json.loads(r.read().decode() or "{}")
def _is_status(t: dict, name: str) -> bool:
return str(t.get("status", "")).lower() == name
def collect(base: str, dep_id: str) -> dict:
"""采集一次 run 的真实 metrics + coverage(诚实标注)。"""
tasks = (_get(base, f"/api/swarms/{dep_id}/tasks").get("data") or {}).get("tasks") or []
events = (_get(base, f"/api/swarms/{dep_id}/events").get("data") or {}).get("events") or []
summary = _get(base, f"/api/swarms/{dep_id}").get("data") or {}
et = Counter(e.get("event_type") for e in events)
cov: dict = {}
# completion(真实)
total = len(tasks)
done = sum(1 for t in tasks if _is_status(t, "completed"))
s_completion = completion_score(done, total) if total else math.nan
cov["s_completion"] = total > 0
# collaboration(真实):handoff 成功率 + 依赖解析率 + 负载均衡
req_h, comp_h = et.get("handoff.requested", 0), et.get("handoff.completed", 0)
handoff = (100.0 * comp_h / req_h) if req_h else 100.0
deps = [t for t in tasks if t.get("depends_on")]
done_ids = {t["task_id"] for t in tasks if _is_status(t, "completed")}
resolved = [t for t in deps if all(d in done_ids for d in t["depends_on"])]
dep_res = (100.0 * len(resolved) / len(deps)) if deps else 100.0
per_agent = Counter(t.get("assigned_agent_id") for t in tasks if t.get("assigned_agent_id"))
bal = (100.0 * min(per_agent.values()) / max(per_agent.values())) if per_agent else 100.0
s_collaboration = collaboration_score(handoff, dep_res, bal) if total else math.nan
cov["s_collaboration"] = total > 0
# robustness(真实):从失败中恢复
failures = [t for t in tasks if (t.get("retry_count") or 0) > 0 or _is_status(t, "failed")]
recovered = [t for t in failures if _is_status(t, "completed")]
s_robustness = robustness_score(len(recovered), len(failures))
cov["s_robustness"] = True
# communication(有 peer 消息才采)
msg_req = et.get("agent.message.request", 0) or et.get("message.request", 0)
msg_rep = et.get("agent.message.reply", 0) or et.get("message.reply", 0)
if msg_req:
s_communication = communication_score(min(msg_rep, msg_req), msg_req)
cov["s_communication"] = True
else:
s_communication = math.nan
cov["s_communication"] = False
# governance(有审批才采)
approvals = summary.get("approvals") or {}
appr_list = list(approvals.values()) if isinstance(approvals, dict) else (approvals or [])
if appr_list:
compliant = sum(1 for a in appr_list if a.get("decision") in ("approved", "rejected"))
s_governance = governance_score(compliant, len(appr_list))
cov["s_governance"] = True
else:
s_governance = math.nan
cov["s_governance"] = False
# cost(需 budget+usage,HTTP summary 一般不给 → 诚实 NaN)
s_cost = math.nan
cov["s_cost"] = False
quality = summary.get("quality") or {}
return {
"deployment_id": dep_id,
"swarm_id": summary.get("swarm_id"),
"status": summary.get("status") or summary.get("state"),
"n_tasks": total,
"n_completed": done,
"n_agents": len(per_agent),
"metrics": {
"s_completion": s_completion,
"s_collaboration": s_collaboration,
"s_robustness": s_robustness,
"s_communication": s_communication,
"s_governance": s_governance,
"s_cost": s_cost,
},
"coverage": cov,
"q_quality_run": quality.get("q_quality"),
"event_counts": dict(et),
}
def _persist_sqlite(db_path: str, row: dict):
os.makedirs(os.path.dirname(db_path), exist_ok=True)
con = sqlite3.connect(db_path)
con.execute("""CREATE TABLE IF NOT EXISTS runs(
ts INTEGER, label TEXT, deployment_id TEXT, swarm_id TEXT, status TEXT,
n_tasks INTEGER, n_completed INTEGER, n_agents INTEGER,
s_completion REAL, s_collaboration REAL, s_robustness REAL,
s_gain REAL, q_swarm REAL, q_base REAL,
coverage_json TEXT, raw_json TEXT)""")
m = row["metrics"]
con.execute("INSERT INTO runs VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)", (
row["ts"], row.get("label", ""), row["deployment_id"], row.get("swarm_id"),
row.get("status"), row["n_tasks"], row["n_completed"], row["n_agents"],
m["s_completion"], m["s_collaboration"], m["s_robustness"],
row.get("s_gain", math.nan), row.get("q_swarm", math.nan), row.get("q_base", math.nan),
json.dumps(row["coverage"]), json.dumps(row, default=str),
))
con.commit(); con.close()
def main():
ap = argparse.ArgumentParser(description="北极星采集器:采集 swarm run metrics 并持久化")
ap.add_argument("deployment_id")
ap.add_argument("--baseline", help="单 agent 基线 run 的 deployment_id(算涌现增益 s_gain)")
ap.add_argument("--q-swarm", type=float, help="蜂群质量分(0-100,如测试 pass rate);省略则用 completion 代理")
ap.add_argument("--q-base", type=float, help="基线质量分(0-100)")
ap.add_argument("--orchestrator", default=os.getenv("ORCH_URL", "http://localhost:8000"))
ap.add_argument("--label", default="")
ap.add_argument("--db", default=os.path.expanduser("~/Desktop/swarm-benchmark.db"))
a = ap.parse_args()
ts = int(time.time())
row = collect(a.orchestrator, a.deployment_id)
row["ts"] = ts
row["label"] = a.label
# 涌现增益 s_gain = Q_swarm - Q_base
s_gain = math.nan; q_s = a.q_swarm; q_b = a.q_base; proxy = False
if a.baseline:
base = collect(a.orchestrator, a.baseline)
if q_s is None:
q_s = row["metrics"]["s_completion"]; proxy = True
if q_b is None:
q_b = base["metrics"]["s_completion"]; proxy = True
if not (math.isnan(q_s) or math.isnan(q_b)):
s_gain = emergence_gain(q_s, q_b)
row["baseline"] = base
row["s_gain"] = s_gain; row["q_swarm"] = q_s if q_s is not None else math.nan
row["q_base"] = q_b if q_b is not None else math.nan
row["s_gain_proxy"] = proxy
_persist_sqlite(a.db, row)
json_path = os.path.join(os.path.dirname(a.db),
f"swarm-benchmark-{a.label or a.deployment_id}-{ts}.json")
with open(json_path, "w", encoding="utf-8") as f:
json.dump(row, f, ensure_ascii=False, indent=2, default=str)
def fmt(v):
return "NaN(未采)" if isinstance(v, float) and math.isnan(v) else f"{v:.1f}" if isinstance(v, float) else v
m = row["metrics"]
print(f"=== 采集 {a.deployment_id} (label={a.label}) ===")
print(f" 任务 {row['n_completed']}/{row['n_tasks']} 完成 | agent 数 {row['n_agents']} | 状态 {row['status']}")
print(f" s_completion = {fmt(m['s_completion'])} [cov={row['coverage']['s_completion']}]")
print(f" s_collaboration = {fmt(m['s_collaboration'])} [cov={row['coverage']['s_collaboration']}]")
print(f" s_robustness = {fmt(m['s_robustness'])} [cov={row['coverage']['s_robustness']}]")
print(f" s_communication = {fmt(m['s_communication'])} [cov={row['coverage']['s_communication']}]")
print(f" s_governance = {fmt(m['s_governance'])} [cov={row['coverage']['s_governance']}]")
print(f" s_cost = {fmt(m['s_cost'])} [cov={row['coverage']['s_cost']}]")
if a.baseline:
tag = "(completion 代理)" if proxy else "(质量分)"
print(f" 涌现增益 s_gain = {fmt(s_gain)} = Q_swarm({fmt(q_s)}) - Q_base({fmt(q_b)}) {tag}")
print(f"\n 持久化 → sqlite: {a.db}")
print(f" 持久化 → json: {json_path}")
if __name__ == "__main__":
main()