WebSocket API

📎 引用文件

本文引用的文件 - agent/src/channels/websocket.py - agent/src/channelsui/gateway_services.py - agent/src/channels/bus/events.py - agent/src/channels/bus/queue.py - agent/src/api/sessions_routes.py - agent/src/api/swarm_routes.py - frontend/src/lib/apiAuth.ts

目录

  1. 简介
  2. 项目结构
  3. 核心组件
  4. 架构总览
  5. 详细组件分析
  6. 依赖关系分析
  7. 性能考虑
  8. 故障排查指南
  9. 结论
  10. 附录

简介

本文件为 Vibe-Trading 的 WebSocket API 实时通信文档,覆盖连接建立、消息格式、事件类型与生命周期管理;说明实时数据流(市场数据推送、回测进度更新、系统状态变更等)的结构;解释订阅机制、过滤规则与批量处理;提供客户端实现要点(连接管理、重连策略、错误处理);总结性能优化技巧(压缩、连接池、内存优化);并给出在交易监控、策略执行和用户交互中的典型应用场景。

项目结构

Vibe-Trading 将 WebSocket 作为“通道”之一,接入统一的消息总线,并与会话/运行事件体系集成: - WebSocket 服务端通道:负责监听、鉴权、路由、订阅与广播。 - 网关服务:提供令牌签发、工作区范围控制、媒体与转录桥接等能力。 - 消息总线:解耦通道与 Agent 核心,承载入站/出站消息。 - HTTP/SSE 接口:会话事件流、Swarm 运行事件流,供前端或外部系统消费。 - 前端认证辅助:为浏览器 SSE 获取一次性票据,避免长链接暴露密钥。

graph TB Client["客户端"] --> WS["WebSocket 通道<br/>agent/src/channels/websocket.py"] WS --> Bus["消息总线<br/>agent/src/channels/bus/queue.py"] Bus --> Core["Agent 核心/会话服务"] WS --> GW["网关服务<br/>agent/src/channelsui/gateway_services.py"] Core --> SSE["SSE 事件流<br/>agent/src/api/sessions_routes.py"] Core --> SwarmSSE["Swarm SSE<br/>agent/src/api/swarm_routes.py"] FE["前端<br/>frontend/src/lib/apiAuth.ts"] --> SSE

图表来源 - agent/src/channels/websocket.py:264-528 - agent/src/channelsui/gateway_services.py:16-273 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800 - agent/src/api/swarm_routes.py:169-211 - frontend/src/lib/apiAuth.ts:18-53

章节来源 - agent/src/channels/websocket.py:264-528 - agent/src/channelsui/gateway_services.py:16-273 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800 - agent/src/api/swarm_routes.py:169-211 - frontend/src/lib/apiAuth.ts:18-53

核心组件

章节来源 - agent/src/channels/websocket.py:264-528 - agent/src/channelsui/gateway_services.py:16-273 - agent/src/channels/bus/events.py:20-55 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800 - agent/src/api/swarm_routes.py:169-211

架构总览

WebSocket 通道作为独立通道接入系统,既可作为 Web UI 的实时通道,也可被其他客户端复用。它通过消息总线与 Agent 核心交互,同时通过网关服务完成鉴权、工作区隔离与媒体处理。SSE 用于事件回放与进度推送,适合浏览器环境。

sequenceDiagram participant C as "客户端" participant W as "WebSocket 通道" participant G as "网关服务" participant B as "消息总线" participant A as "Agent/会话服务" participant S as "SSE 端点" C->>W : 建立连接(可选 token/client_id) W->>G : 校验/签发令牌, 解析工作区范围 W-->>C : ready(chat_id, client_id) C->>W : new_chat/attach/message(含媒体) W->>B : publish_inbound(InboundMessage) B->>A : 入站消息进入 Agent 处理 A-->>B : 出站消息(OutboundMessage) B-->>W : 出站消息分发 W-->>C : event/attached/session_updated/文本帧 C->>S : 订阅会话事件流(支持 Last-Event-ID) S-->>C : 事件增量/完成信号

图表来源 - agent/src/channels/websocket.py:386-443 - agent/src/channels/websocket.py:530-589 - agent/src/channelsui/gateway_services.py:16-150 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800

详细组件分析

WebSocket 通道与连接生命周期

flowchart TD Start(["连接建立"]) --> Auth{"是否启用令牌?"} Auth --> |是| CheckToken["校验静态/issued token"] Auth --> |否| Allow["允许无令牌(受 allow_from 控制)"] CheckToken --> Ready["发送 ready + 默认 chat_id"] Allow --> Ready Ready --> Envelope{"是否为新式信封?"} Envelope --> |是| Dispatch["new_chat/attach/message 等"] Envelope --> |否| Legacy["解析 content/text/message"] Dispatch --> Media{"是否包含媒体?"} Media --> |是| ValidateMedia["MIME/数量/大小校验"] Media --> |否| Publish["发布到消息总线"] ValidateMedia --> Publish Publish --> Fanout["按 chat_id 广播"] Fanout --> End(["结束/等待下一帧"])

图表来源 - agent/src/channels/websocket.py:386-443 - agent/src/channels/websocket.py:530-589 - agent/src/channels/websocket.py:593-799

章节来源 - agent/src/channels/websocket.py:53-143 - agent/src/channels/websocket.py:386-443 - agent/src/channels/websocket.py:530-589 - agent/src/channels/websocket.py:593-799

网关服务与令牌签发

章节来源 - agent/src/channelsui/gateway_services.py:16-150 - agent/src/channelsui/gateway_services.py:153-273

消息总线与消息模型

章节来源 - agent/src/channels/bus/events.py:20-55 - agent/src/channels/bus/queue.py:8-44

SSE 事件流与会话/运行事件

sequenceDiagram participant FE as "前端" participant API as "FastAPI" participant SRV as "会话/运行服务" FE->>API : GET /sessions/{id}/events?Last-Event-ID=... API->>SRV : subscribe(session_id, last_event_id, replay?) loop 事件循环 SRV-->>API : 事件对象 API-->>FE : id/event/data SSE 帧 end API-->>FE : event : done (完成/缺失)

图表来源 - agent/src/api/sessions_routes.py:752-800 - agent/src/api/swarm_routes.py:169-211 - frontend/src/lib/apiAuth.ts:18-53

章节来源 - agent/src/api/sessions_routes.py:752-800 - agent/src/api/swarm_routes.py:169-211 - frontend/src/lib/apiAuth.ts:18-53

实时数据流结构与事件类型

章节来源 - agent/src/channels/websocket.py:321-349 - agent/src/channels/websocket.py:661-799 - agent/src/channelsui/transcription_ws.py:8-12 - agent/src/channels/bus/events.py:20-55

订阅机制、过滤与批量处理

章节来源 - agent/src/channels/websocket.py:304-349 - agent/src/channels/websocket.py:593-799 - agent/src/channelsui/gateway_services.py:74-134

客户端实现示例(要点)

章节来源 - agent/src/channels/websocket.py:386-443 - agent/src/channels/websocket.py:593-799 - agent/src/api/sessions_routes.py:752-800 - frontend/src/lib/apiAuth.ts:18-53

依赖关系分析

graph LR WS["WebSocketChannel"] --> GW["GatewayServices"] WS --> MB["MessageBus"] MB --> AG["Agent/会话服务"] AG --> SSE["SSE 端点"] SSE --> FE["前端"]

图表来源 - agent/src/channels/websocket.py:264-528 - agent/src/channelsui/gateway_services.py:16-273 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800

章节来源 - agent/src/channels/websocket.py:264-528 - agent/src/channelsui/gateway_services.py:16-273 - agent/src/channels/bus/queue.py:8-44 - agent/src/api/sessions_routes.py:752-800

性能考虑

[本节为通用指导,不直接分析具体文件]

故障排查指南

章节来源 - agent/src/channels/websocket.py:386-443 - agent/src/channels/websocket.py:593-799 - agent/src/api/sessions_routes.py:752-800 - frontend/src/lib/apiAuth.ts:18-53

结论

Vibe-Trading 的 WebSocket API 提供了高内聚、低耦合的实时通信能力:通过统一的通道与消息总线,既能支撑 Web UI 的即时交互,也能服务于外部系统集成。结合 SSE 的事件回放与完成信号,可实现可靠的进度追踪与状态同步。在生产环境中,建议启用 WSS、配置短效令牌与工作区范围,并结合客户端重连与错误处理策略,以获得稳定高效的实时体验。

[本节为总结性内容,不直接分析具体文件]

附录

章节来源 - agent/src/channels/websocket.py:53-143 - agent/src/channels/websocket.py:593-799 - agent/src/channels/websocket.py:321-349