# 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` ## 配置(环境变量) ```bash # 存储 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 断连后其在执行的任务会被回收为可重派。 ## 本地运行 ```bash 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 ``` ## 测试 ```bash python scripts/test-runtime-contract.py # Manager 契约校验 python scripts/test-merge-smoke.py # 工作流冒烟测试(评审/协作/汇总等) ```