Swarm运行状态管理

📎 引用文件

本文引用的文件 - agent/src/swarm/models.py - agent/src/swarm/store.py - agent/src/swarm/runtime.py - agent/src/swarm/task_store.py - agent/src/swarm/worker.py - agent/src/api/swarm_routes.py - frontend/src/stores/agent.ts

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能与并发特性
  8. 故障排除指南
  9. 结论
  10. 附录:API与事件速查

简介

本文件面向 Vibe-Trading 研究页面的“Swarm运行状态管理系统”,系统性说明多Swarm实例的状态隔离、状态更新机制、持久化策略,以及与消息系统(SSE)的集成方式。重点覆盖以下主题: - SwarmRunStatus 数据结构与生命周期 - 多Swarm实例的状态隔离(按run_id分目录、任务独立存储) - 状态更新机制(reconcile_run、hydrate_run、原子写入、心跳检测) - upsertSwarmStatus/updateSwarmStatus 在前端状态管理中的使用场景 - 状态持久化策略(run.json、events.jsonl、tasks/*.json) - 运行监控、状态同步、冲突解决(Windows重命名重试、读取重试、幂等恢复) - 调试与故障排除(SSE断线、僵尸运行、任务阻塞、超时与内容过滤)

项目结构

Swarm运行状态管理由后端运行时、持久化层、工作进程以及前端状态管理共同构成: - 数据模型:定义运行、任务、事件、工作者结果等核心结构 - 持久化:基于文件系统,按run_id隔离;run.json为聚合快照,events.jsonl为追加日志,tasks/*.json为任务实时状态 - 运行时:编排DAG执行、并发调度、取消与重试、最终状态收敛 - 工作进程:轻量ReAct循环、工具调用、心跳上报、摘要与产物输出 - API路由:提供REST接口与SSE事件流,供前端查询与订阅 - 前端状态:维护swarmRuns映射,支持upsert/update以增量更新UI

graph TB subgraph "后端" A["SwarmRuntime<br/>编排与调度"] --> B["SwarmStore<br/>run.json + events.jsonl"] A --> C["TaskStore<br/>tasks/*.json"] A --> D["Worker<br/>ReAct循环"] E["HTTP/SSE路由<br/>/swarm/runs*"] --> A E --> B E --> C end subgraph "前端" F["Zustand Store<br/>swarmRuns 映射"] G["页面组件<br/>列表/详情/事件流"] end F --> G E --> F

图表来源 - agent/src/swarm/runtime.py:49-148 - agent/src/swarm/store.py:115-218 - agent/src/swarm/task_store.py:16-110 - agent/src/swarm/worker.py:297-758 - agent/src/api/swarm_routes.py:91-211 - frontend/src/stores/agent.ts:251-280

章节来源 - agent/src/swarm/models.py:14-217 - agent/src/swarm/store.py:115-566 - agent/src/swarm/runtime.py:49-751 - agent/src/swarm/task_store.py:16-249 - agent/src/swarm/worker.py:297-758 - agent/src/api/swarm_routes.py:91-211 - frontend/src/stores/agent.ts:251-280

核心组件

章节来源 - agent/src/swarm/models.py:14-217 - agent/src/swarm/store.py:115-566 - agent/src/swarm/task_store.py:16-249 - agent/src/swarm/runtime.py:49-751 - agent/src/swarm/worker.py:297-758 - agent/src/api/swarm_routes.py:91-211 - frontend/src/stores/agent.ts:251-280

架构总览

Swarm运行状态管理的整体流程如下: - 启动:API接收请求,构建SwarmRun并持久化,标记为running,后台线程执行 - 执行:按拓扑分层并行执行任务,任务间通过依赖图与blocked_by协调 - 事件:每个关键阶段写入events.jsonl,并通过SSE推送给前端 - 收敛:层边界将任务快照回写run.json;结束时根据任务状态推导运行终态 - 恢复:reconcile_run在读取时合并任务实时状态、修复僵尸运行、回收超时无心跳的运行

sequenceDiagram participant Client as "客户端" participant API as "HTTP/SSE路由" participant RT as "SwarmRuntime" participant ST as "SwarmStore" participant TS as "TaskStore" participant W as "Worker" Client->>API : POST /swarm/runs API->>RT : start_run(preset, user_vars) RT->>ST : create_run(run) RT->>ST : update_run(status=running) RT->>TS : save_task(task_i) loop 每层 RT->>W : 并行执行任务 W-->>RT : WorkerResult RT->>TS : update_status(completed/failed) RT->>ST : append_event(...) RT->>ST : _sync_run_tasks_snapshot() end RT->>ST : update_run(final status) API-->>Client : SSE事件流(任务/运行事件)

图表来源 - agent/src/api/swarm_routes.py:91-211 - agent/src/swarm/runtime.py:211-391 - agent/src/swarm/store.py:159-218 - agent/src/swarm/task_store.py:47-110 - agent/src/swarm/worker.py:297-758

详细组件分析

数据模型与状态机

stateDiagram-v2 [*] --> pending pending --> running : "start_run" running --> completed : "所有任务完成" running --> failed : "任一任务失败/异常" running --> cancelled : "用户取消/层超时" note right of running : "reconcile_run会合并任务状态并回收僵尸运行"

图表来源 - agent/src/swarm/models.py:14-41 - agent/src/swarm/store.py:364-423

章节来源 - agent/src/swarm/models.py:14-217

持久化与状态隔离

flowchart TD Start(["写入run.json"]) --> Tmp["写入.tmp临时文件"] Tmp --> Rename{"os.replace成功?"} Rename -- 否(Windows共享冲突) --> Retry["指数退避重试"] Retry --> Rename Rename -- 是 --> Done(["完成"])

图表来源 - agent/src/swarm/store.py:94-113 - agent/src/swarm/store.py:159-218 - agent/src/swarm/store.py:555-566

章节来源 - agent/src/swarm/store.py:115-566 - agent/src/swarm/task_store.py:16-110

状态更新与收敛(reconcile_run/hydrate_run)

flowchart TD RStart(["reconcile_run(run)"]) --> Hydrate["hydrate_run(合并任务)"] Hydrate --> CheckTerm{"run已是终态?"} CheckTerm -- 是 --> ReturnHyd["返回hydrated"] CheckTerm -- 否 --> AllTerm{"所有任务均为终态?"} AllTerm -- 是 --> Recover["_recover_terminal(推导run终态)"] Recover --> PersistR{"write=True?"} PersistR -- 是 --> SaveR["持久化恢复事件"] PersistR -- 否 --> EndR["返回"] AllTerm -- 否 --> Stale{"是否过期?"} Stale -- 是 --> Reap["_reap_stale(标记失败)"] Reap --> PersistS{"write=True?"} PersistS -- 是 --> SaveS["持久化回收事件"] PersistS -- 否 --> EndS["返回"] Stale -- 否 --> ReturnHyd

图表来源 - agent/src/swarm/store.py:286-423

章节来源 - agent/src/swarm/store.py:286-423

运行编排与并发(SwarmRuntime)

sequenceDiagram participant RT as "SwarmRuntime" participant TS as "TaskStore" participant ST as "SwarmStore" participant W as "Worker" RT->>ST : reap_stale_running_runs() RT->>ST : create_run(run) RT->>ST : update_run(running) loop 拓扑层 RT->>TS : load_all() RT->>W : 并行执行层内任务 W-->>RT : WorkerResult RT->>TS : update_status(...) RT->>ST : append_event(...) RT->>ST : _sync_run_tasks_snapshot() end RT->>ST : update_run(final)

图表来源 - agent/src/swarm/runtime.py:86-148 - agent/src/swarm/runtime.py:211-391 - agent/src/swarm/runtime.py:465-634

章节来源 - agent/src/swarm/runtime.py:49-751

工作者与心跳(Worker)

flowchart TD WStart["run_worker入口"] --> Build["构建工具注册表/LLM/提示词"] Build --> Loop{"迭代预算"} Loop -- 继续 --> LLM["stream_chat(带心跳)"] LLM --> Tools{"有工具调用?"} Tools -- 否 --> Finalize["生成摘要/产物/返回completed或incomplete"] Tools -- 是 --> Exec["执行工具(带心跳)"] Exec --> Append["追加消息/结果"] Append --> Loop Loop -- 超时/超限/熔断 --> Fail["返回timeout/token_limit/failed"]

图表来源 - agent/src/swarm/worker.py:297-758

章节来源 - agent/src/swarm/worker.py:297-758

API与SSE事件流

sequenceDiagram participant FE as "前端" participant API as "swarm_routes" participant RT as "SwarmRuntime" participant ST as "SwarmStore" FE->>API : GET /swarm/runs/{run_id}/events?last_index=N API->>ST : read_events(after_index=N) loop 轮询 ST-->>API : 事件片段 API-->>FE : event : task_heartbeat/tool_result/... API->>ST : load_run + reconcile_run alt run终态 API-->>FE : event : done end end

图表来源 - agent/src/api/swarm_routes.py:169-211 - agent/src/swarm/store.py:246-284

章节来源 - agent/src/api/swarm_routes.py:91-260

前端状态管理:upsertSwarmStatus与updateSwarmStatus

flowchart TD Msg["收到swarm_status消息"] --> Find{"是否存在同runId?"} Find -- 是 --> Upsert["覆盖swarmRuns[runId]"] Find -- 否 --> NewMsg["新增swarm_status消息"] --> Insert["插入swarmRuns[runId]"] Update["updateSwarmStatus(runId, updater)"] --> Patch["局部更新现有状态"]

图表来源 - frontend/src/stores/agent.ts:251-280

章节来源 - frontend/src/stores/agent.ts:251-280

依赖关系分析

graph LR models["models.py"] --> store["store.py"] models --> runtime["runtime.py"] models --> worker["worker.py"] store --> runtime task_store["task_store.py"] --> runtime task_store --> store worker --> runtime routes["swarm_routes.py"] --> runtime routes --> store frontend["agent.ts"] --> routes

图表来源 - agent/src/swarm/runtime.py:25-44 - agent/src/swarm/store.py:21-24 - agent/src/swarm/worker.py:16-36 - agent/src/api/swarm_routes.py:18-35 - frontend/src/stores/agent.ts:251-280

章节来源 - agent/src/swarm/runtime.py:25-44 - agent/src/swarm/store.py:21-24 - agent/src/swarm/worker.py:16-36 - agent/src/api/swarm_routes.py:18-35 - frontend/src/stores/agent.ts:251-280

性能与并发特性

[本节为通用指导,无需特定文件引用]

故障排除指南

章节来源 - agent/src/swarm/store.py:315-423 - agent/src/api/swarm_routes.py:169-211 - agent/src/swarm/runtime.py:465-634 - agent/src/swarm/worker.py:416-758 - frontend/src/stores/agent.ts:251-280

结论

Swarm运行状态管理系统通过清晰的数据模型、严格的持久化策略、健壮的收敛与恢复机制,实现了多Swarm实例的安全隔离与可靠执行。结合SSE事件流与前端状态管理,提供了实时、一致、可观测的运行体验。建议在生产环境中: - 保持心跳配置合理,避免误判僵尸运行 - 关注reconcile_run的恢复事件,及时排查异常 - 利用SSE断点续传与任务状态字段进行前端健壮性处理 - 定期巡检events.jsonl与tasks/*.json,辅助定位复杂问题

[本节为总结,无需特定文件引用]

附录:API与事件速查

章节来源 - agent/src/api/swarm_routes.py:91-260 - agent/src/swarm/runtime.py:166-209 - agent/src/swarm/worker.py:80-109