Files
Agentswarm/orchestrator
FastheiandClaude Opus 4.8 1edf73aba4 fix(swarm/#70): 任务超时透传+提默认 / 事件时间线降噪 / 失败 termination_reason 准确
真机端到端实测(#70)暴露三问题,本 PR 全部修复(仅 agent + orchestrator,不跨仓):

问题1【阻断】单任务超时只有 60s,生成类任务必挂
- agent/main.py: TASK_TIMEOUT_SECONDS 默认 60→300(仅对外部/独立启动 agent 生效)。
- agent_launcher.py: 新增 DEFAULT_TASK_TIMEOUT_SECONDS=300、_budget_duration_seconds、
  resolve_task_timeout(base=env 默认 300,与 run budget.duration_seconds 取较小);
  plan_launch_specs 把 TASK_TIMEOUT_SECONDS 透传进每个 agent env(非敏感,inline,
  k8s 不进 Secret)。

问题2【体验】事件时间线全是内部噪音(纯附加,未碰冻结契约)
- swarm_runtime.py: is_client_visible(=event_type∈FROZEN_CLIENT_EVENT_TYPES,单一真源);
  emit_event 给 envelope 加 metadata.client_visible 布尔 + 关键客户端事件回填可选
  payload.message(人话进度,仅取已有字段,不伪造)。task.heartbeat/retried/
  deployment.status_changed/timeline/budget 标 client_visible=false,仍持久化+回调
  但客户端据此过滤出时间线。冻结事件集/类型/sequence/artifact 形状一字未动。
- event-schema.md: 文档化两个附加字段 + 新增 §6.1,明确未解冻。

问题3【正确性】失败/超时 termination_reason 仍报 "tasks_completed"
- convergence.py: 新增 TIMEOUT/MAX_RETRIES_EXCEEDED/TASK_FAILED;classify_failure_reason
  按 timeout→max_retries→task_failed 取最具体(仅凭真实 per-task 信号);FAILED 分支
  再不会返回 tasks_completed(该 reason 仅用于成功),budget/rounds 仅在通用失败时才覆盖。
- task_queue.py: fail_task 永久失败时把 reason 落到 task.result({"success":false,"error":reason}),
  不覆盖已有结果,供 convergence 读取。
- main.py: compute_convergence_report 快照补 retry_count/max_retries。

测试:新增 test_resolve_task_timeout、扩 test-convergence(failed_timeout/max_retries/
generic + "FAILED 永不报 tasks_completed"不变量)。本地全过:test-agent-launcher /
test-convergence / test-runtime-contract / test-contract-freeze / test-merge-smoke /
test-workflow-e2e / test-security-boundary。

影响:agent + orchestrator + 文档;不动 Manager↔Swarm 冻结契约字段(问题2 纯附加)。
栈在 #64(agent_swarm git 注入)之上,#64 合并后本 PR base 自动转 main。

Closes #70

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-06-16 18:29:39 +08:00
..

Orchestrator(蜂群编排器)

基于 FastAPI 的编排器,是蜂群系统的控制核心。它对上对接 Heicode Manager(控制面),对下通过 WebSocket 调度多个 Agent 执行单元,负责任务分解、派发、专家协作、质量评审与结果汇总。

能力概览

  • Manager 对接:提供面向 Manager 的部署/任务/日志/指标/审批等 REST 接口,事件通过带签名的回调上报。
  • WebSocket 调度:与 Agent 实时双向通信,按能力与容量派发任务。
  • 任务图与依赖:支持任务间依赖(depends_on),依赖完成后才派发后继任务。
  • 任务规划(可选):在 Manager 未给出明确分工时,用大模型把目标拆解为「实现 → 测试 → 文档」等专家子任务。
  • 主控评审循环(可选):所有任务完成后由「评审者」判断结果是否达标,不达标则把相关任务退回重做,直至通过或达到上限,并对结果做汇总。
  • 专家协作:派发时注入依赖产物与同伴信息,支持 Agent 间消息路由。
  • 持久化:以 Redis 为权威存储;提供仅限开发/CI 的内存回退。
  • 故障恢复:心跳超时检测、任务重派与重试。

目录结构

orchestrator/
├── main.py             # FastAPI 应用:REST 接口、WebSocket、调度循环、评审与汇总
├── swarm_runtime.py    # Manager 桥接:运行持久化、事件/回调、审批、任务图构建
├── planner.py          # 任务规划 + 主控评审(critic)+ 结果汇总
├── task_queue.py       # 任务队列:创建、派发、依赖、重试、重开
├── agent_registry.py   # Agent 注册表:注册、心跳、状态
├── handoff_manager.py  # 移交协调与历史
├── redis_client.py     # 异步 Redis 封装(含受控的内存回退)
├── checkpoint_manager.py / model_tuner.py / tracing.py  # 检查点、模型调优、链路追踪
└── requirements.txt

工作流程

分解(Plan) → 派发(Dispatch) → 专家执行(Execute) → 协作/移交(Handoff)
            → 评审(Review,可选循环重做) → 汇总交付(Deliver)
  • 分解:优先采用 Manager 提供的编排方案;若未提供且开启了规划回退,则由 planner 生成专家子任务图。
  • 派发:调度循环将就绪任务按「能力匹配 + 剩余容量」派发给已连接的空闲 Agent;派发时把已完成依赖的产物与同伴 Agent 信息注入任务上下文。
  • 评审循环:任务全部完成后,planner.review 判定是否达标;不达标则 reopen 指定任务重做,受 MAX_REVIEW_CYCLES 约束;通过后由 planner.synthesize 生成统一的最终回答。

WebSocket 协议(/ws/{agent_id})

Agent → 编排器:register(含 available_slots)、heartbeat、task_start、task_complete、task_failed、task_accepted、task_rejected、blocked_on_handoff、handoff_request、peer_message、status_update。

编排器 → Agent:task_assignment、cancel_task、peer_message(路由转发)、各类确认(registered、heartbeat_ack 等)。

REST 接口

健康与指标

  • GET /、GET /health、GET /metrics(Prometheus)

Manager 面(含 /api/swarms、/api/agent/swarm/deployments、/api/agnet/deployments 等别名)

  • POST .../:创建蜂群部署
  • GET .../{id}:部署概要
  • GET .../{id}/tasks | /logs | /events | /metrics | /workflow | /diagnostics
  • POST .../{id}/stop:停止部署(会向相关 Agent 发送 cancel_task)
  • POST .../{id}/approvals/{approval_id}:接收 Manager 审批决定

内部/调试

  • GET /agents、/agents/idle、/agents/{id}
  • POST /tasks、POST /tasks/assign、GET /tasks、GET /tasks/{id}、GET /handoffs

配置(环境变量)

# 存储
REDIS_HOST=redis-service
REDIS_PORT=6379
REDIS_DB=0
REDIS_FAKE=1            # 仅开发/CI:使用内存版 fakeredis
ALLOW_MEMORY_STORE=1    # 仅开发/CI:Redis 不可用时回退到内存(生产请勿开启)

# 工作流开关
ENABLE_PLANNER_FALLBACK=1   # 无 Manager 分工时启用规划回退
ENABLE_REVIEW_LOOP=1        # 启用主控评审/重做循环 + 结果汇总
MAX_REVIEW_CYCLES=2         # 评审最大重做轮数
ENABLE_SUBTASK_HANDOFF=false

# 规划/评审用模型(OpenAI 兼容;未配置时退化为静态分解 + 启发式评审)
OPENAI_API_KEY / OPENAI_API_BASE / OPENAI_MODEL
MASTER_REVIEW_MODEL / MAX_SUBTASKS / PLANNER_TIMEOUT_SECONDS

# 安全与回调(与 Manager 契约相关)
AGENT_RUNTIME_SERVICE_TOKEN / AGNET_RUNTIME_SERVICE_TOKEN     # 服务间鉴权
AGENT_CALLBACK_SERVICE_TOKEN / AGENT_CALLBACK_SIGNING_SECRET  # 回调令牌与 HMAC 签名

# 可观测性
OTEL_EXPORTER_OTLP_ENDPOINT / SWARM_RUNTIME_SOURCE / SWARM_RUNTIME_PLATFORM

说明:未配置运行时服务令牌时,Manager 面接口进入「非安全开发模式」(不校验鉴权),仅供本地调试。

持久化与回退

Redis 为权威存储。仅当显式设置 REDIS_FAKE=1 或 ALLOW_MEMORY_STORE=1 时,才会使用进程内的 fakeredis 作为开发/CI 回退;生产环境在 Redis 不可用时会快速失败,避免静默丢失持久化与 Manager 状态。

故障恢复

  • Agent 需每 15 秒心跳;超过 30 秒无心跳判定为失败。
  • 失败 Agent 的任务自动退回队列重派;任务失败按 max_retries 重试。
  • Agent 断连后其在执行的任务会被回收为可重派。

本地运行

pip install -r orchestrator/requirements.txt

# 开发模式(内存回退 + 工作流开关)
set "REDIS_FAKE=1"
set "ENABLE_PLANNER_FALLBACK=1"
set "ENABLE_REVIEW_LOOP=1"
python -m uvicorn orchestrator.main:app --host 0.0.0.0 --port 8000

测试

python scripts/test-runtime-contract.py   # Manager 契约校验
python scripts/test-merge-smoke.py         # 工作流冒烟测试(评审/协作/汇总等)