按 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>
199 lines
8.9 KiB
Python
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()
|