四块互相交织的 benchmark 覆盖增量,统一提交: 1) 通信遥测(#23):orchestrator 路由 peer 消息时按 correlation_id 计请求/应答到 SwarmRun.collaboration(内部状态,不进 Manager 事件流);collector 算 s_communication。 治理计数由 run.approvals 派生(合规/总数)→ s_governance。 2) Q_quality 掩码归一(v2.1 裁定):metrics.quality_score 改为对 present 输入加权归一, 非编码任务自动忽略 TestPassRate,全缺 → NaN(不伪造)。 3) 质量插桩 / Group B:新增 Pod 内代码测试沙箱(orchestrator/sandbox.py,环境清洗 + 超时强杀 + 资源限额 + 路径越界校验,门控 ENABLE_QUALITY_EVAL)与 held-out fixture (benchmark/fixtures/);run 完成时用留出测试评分得 TestPassRate → Q_quality → collector 合成 reward。安全边界见 docs/integration/security-boundary.md §8.1。 4) 决策引擎 / Group A(#10,Option A score-at-pull):新增 orchestrator/decision_engine.py —— 信息素 τ(Redis 持久、(role,agent) 键控、冷启动 0.5、ρ 蒸发、夹紧、学习常开)+ η 启发式评分 + ε-greedy 概率采样;每次 dispatch 产一条 DecisionTrace → SwarmRun.decisions;collector 算 tau/eta/p_decision。概率选择门控 ENABLE_ACO_DISPATCH (默认关,CI 用 ACO_SEED 固定)。 覆盖:单次 run 真实可算字段由 4 提升至最多 10/15(新增 communication/reward/tau/eta/ p_decision,外加 governance 有条件)。 测试:新增 test-sandbox / test-quality / test-decision-engine;扩充 collector/metrics 用例; CI 纳入全部 benchmark 套件 + flag-on 的 ACO e2e。本地 11 项 gate 全绿。 诚实边界(未越界声称): - Group A 为单边匹配(Option B 待 Group C);概率派发优于贪心未证;默认关闭。 - reward 的 CodeReview/UserAcceptance 未采集(掩码忽略);P_risk 为审批派生低估。 - s_gain/s_swarm/g_e/g_e_cost/benchmark 仍 NaN —— 需基线(#21/#13),本 PR 不动验收。 影响范围:Swarm(orchestrator + benchmark + docs + CI)。不改 Manager↔Swarm 事件契约 (遥测均为运行时内部状态);不影响 Client/计费/密钥/发布链路。新增 ENABLE_QUALITY_EVAL / ENABLE_ACO_DISPATCH 两个开关,默认关闭。 Closes #10 Closes #23 Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
7.6 KiB
Peer Communication Logistics(对等通信机制与深化设计)
状态:现状 + 深化设计。描述专家 Agent 之间「就重叠领域相互沟通」的消息机制、路由逻辑、当前局限,以及深化方案与分阶段计划。
相关代码:
agent/main.py(request_peer_collaboration/handle_peer_message/answer_peer_query)、orchestrator/main.py(peer_message路由 /build_dispatch_context)、agent/task_executor.py(_execute_subtask触发)。相关文档:event-schema、frontend-event-api、benchmark/swarm-metrics-schema(S_communication)。
1. 当前实现(mechanics)
1.1 消息信封
{
"type": "peer_message",
"agent_id": "发送方",
"target_agent_id": "接收方",
"from_agent_id": "由编排器转发时补盖(接收方据此回复)",
"task_id": "...",
"content": "文本内容",
"correlation_id": "peer-<task>-<rand>", // 关联请求与回复
"is_reply": false, // true = 回复
"timestamp": 0.0
}
1.2 路由(编排器为 broker)
Agent 出站连编排器 WebSocket,彼此不直连。编排器收到 peer_message → 按 target_agent_id 转发,并补盖 from_agent_id(发送方身份)。无 target_agent_id 则记告警丢弃。
1.3 请求/回复(一问一答)
- 发起方
request_peer_collaboration(task_id, target_agent_id, content, timeout):生成correlation_id,挂一个asyncio.Future到peer_waiters[correlation_id],发出 query,await asyncio.wait_for(waiter, timeout)。 - 接收方
handle_peer_message:若correlation_id命中自身 waiter(或is_reply)→ 解析为「回复」并 resolve waiter;否则视为「入站查询」→answer_peer_query。 answer_peer_query:已实现实质性回复——对查询做一次受限 LLM 推理(TaskExecutor.peer_reply,受PEER_CONSULT_MAX_TOKENS,默认 500 约束),结合自身 workspace 产物,返回结构化{stance, content, evidence, refs};无模型/超时/失败时回退到last_summary摘要。回复在后台任务中完成,不阻塞消息循环(PEER_REPLY_TIMEOUT_SECONDS,默认 20s)。
1.4 触发与编排(logistics)
- 派发时
build_dispatch_context向任务上下文注入:peer_agents(本次 run 中已连接的其它 Agent:agent_id/role/capabilities)与dependency_artifacts(已完成依赖产物)。 task_executor._execute_subtask仅当specialist_role ∈ {testing, documentation}且存在peer_agents且提供了回调时发起咨询;按prefer implementation排序,取前max_peer_consults(默认 2)个,peer_timeout_seconds(默认 10s)。- 对齐原则(提示词内):实现产物为 API/异常语义的 source of truth;peer 输入与实现冲突时以实现为准并在变更摘要中说明。
2. 当前局限(为何"浅")
| 局限 | 说明 |
|---|---|
answer_peer_query 现做受限 LLM 推理 + 引用 workspace 产物 + 结构化字段(stance/content/evidence/refs);无模型时回退摘要(见 §3.3) |
|
| 触发窄 | 仅 testing/documentation 角色发起;实现角色不主动协作 |
| 单轮 | 一问一答,无多轮/澄清/协商线程 |
| 无广播 | 只能点对点 target_agent_id,无按角色/能力的群发或发现 |
orchestrator 路由 peer 消息时按 correlation_id 记请求/应答与投递失败到 SwarmRun.collaboration(内部计数,不进 Manager 事件流)→ collector 据此算 S_communication(请求→应答率);无 peer 通信的 run 仍 NaN(见 metric-coverage-gaps) |
|
| 弱冲突处理 | 仅"以实现为准"的提示约定,无结构化协商/升级协议 |
| 无前端可见 | peer 消息未作为事件暴露,前端协作时间线无数据 |
3. 深化设计(target)
3.1 消息类型(扩展 kind)
query(提问)·reply(回答)·clarify(追问)·proposal(提案)·critique(质疑)·broadcast(群发)·ack(确认)。
3.2 多轮线程
增加 conversation_id + turn;peer_waiters 升级为按会话的队列,支持 clarify 往返与超时续约,受 max_rounds 约束。
3.3 实质性回复 ✅ 已实现
answer_peer_query 已由「回 last_summary」升级为:针对 query 生成有据回复(TaskExecutor.peer_reply 做一次受限推理,引用自身 workspace 产物),返回结构化 {stance, content, evidence, refs},受 PEER_CONSULT_MAX_TOKENS(默认 500)约束;无模型/失败回退摘要;后台任务回复不阻塞消息循环。由 scripts/test-merge-smoke.py 覆盖(fallback + 实质回复两条断言)。
3.4 寻址与发起
- 点对点(
target_agent_id)+ 按角色/能力寻址(编排器按注册表解析)+ 受控广播。 - 任何角色均可发起咨询(不限 testing/doc),由预算与策略限制频次。
3.5 冲突解决协议
critique/proposal 往返;默认「实现为 source of truth」;僵局上交 Master Agent(master_agent 仲裁,纳入评审决策),而非各 Agent 私下定夺。
3.6 遥测(接入 benchmark)
编排器对每条 peer_message 计数:total / delivered / failed(接收方离线/超时)。
- 产出
S_communication = SuccessfulMessages / TotalMessages(见 swarm-metrics-schema)。 - 作为协作事件暴露(见 §4),供前端协作时间线与审计。
3.7 预算与可靠性
max_peer_consults、peer_timeout_seconds、max_rounds、peer_consult_max_tokens;接收方离线 → 立即回 failed 而非干等;correlation_id 幂等去重。
4. 事件与可见性
为满足前端与审计,建议把对等协作纳入事件流(与 event-schema 对齐):
- 新增/复用事件:
collaboration.message(含from/to/kind/conversation_id/correlation_id,内容脱敏)、collaboration.resolved/collaboration.escalated。 - 前端按
conversation_id渲染协作时间线(见 frontend-event-api §4)。 - HM 注册表暂无
collaboration.*事件 → 列为与 Workflow/Frontend Team 的待对齐项。
5. 与其它机制的关系
| 机制 | 作用 | 区别 |
|---|---|---|
dependency_artifacts |
顺序依赖:下游看到上游产物 | 单向、派发时注入,非交互 |
| handoff | 把子任务委派给更合适的专家 | 转移所有权,非对等沟通 |
| peer communication | 重叠领域的对等交流 | 双向、运行中、不转移所有权 |
| Master Agent 评审 | 全局裁决是否达标 | 中心化决策,非点对点 |
6. 分阶段计划
- 遥测先行(低成本):编排器对
peer_message计数 → 点亮S_communication;把消息作为collaboration.message事件暴露。 - 实质回复 ✅ 已实现:
answer_peer_query改为受限 LLM 推理 + 引用产物(替换 last_summary)。 - 多轮 + 寻址:
conversation_id/clarify、角色/能力寻址、广播。 - 冲突协议 + 升级 Master:
critique/proposal,僵局上交master_agent。 - 与 Frontend/Workflow Team 冻结
collaboration.*事件与协作时间线契约。
7. 待对齐
collaboration.*事件类型与载荷(Workflow/Frontend Team;当前 HM 注册表未含)。- 对等咨询的预算口径(
max_rounds/token)与频次策略(Governance)。 - 实质回复是否计入用量/计费(每次 peer consult 是一次模型调用 → 见 usage-billing §5)。