Files
Agentswarm/orchestrator
FastheiandGitHub 5e29618bf5 Merge pull request #32 from xmindlab-heicode/feat/max-agents-per-user
每用户并发 Agent 上限:MAX_AGENTS_PER_USER 默认 10
2026-06-10 22:44: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         # 工作流冒烟测试(评审/协作/汇总等)