流式传输实现

📎 引用文件

本文引用的文件 - agent/src/channels/base.py - agent/src/channels/telegram.py - agent/src/channels/discord.py - agent/src/channels/runtime.py - agent/src/channels/manager.py - frontend/src/stores/agent.ts

目录

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

简介

本文件聚焦 Vibe-Trading 的“流式传输”能力,围绕 send_delta() 与 send_reasoning_delta() 的设计与实现要求,系统阐述增量更新、状态管理、错误重试、会话隔离与并发控制。同时记录不同渠道(Telegram、Discord)的适配策略,给出长文本、富媒体与交互式内容的处理要点,并提供性能优化、内存管理与网络异常处理的实践建议,以及调试与性能分析指引。

项目结构

graph TB subgraph "通道抽象" Base["BaseChannel<br/>定义流式接口"] end subgraph "平台适配器" TG["TelegramChannel<br/>edit_message_text 渐进更新"] DC["DiscordChannel<br/>_StreamBuf 增量缓冲"] end subgraph "运行时" RT["ChannelRuntime<br/>启动/停止/消费循环"] CM["ChannelManager<br/>通道发现/配置/开关"] end subgraph "前端" FE["Agent Store<br/>appendDelta/status"] end Base --> TG Base --> DC RT --> CM RT --> TG RT --> DC FE --> RT

图表来源 - agent/src/channels/base.py:22-163 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/discord.py:40-48 - agent/src/channels/runtime.py:34-102 - agent/src/channels/manager.py:36-137 - frontend/src/stores/agent.ts:136-165

章节来源 - agent/src/channels/base.py:22-163 - agent/src/channels/runtime.py:34-102 - agent/src/channels/manager.py:36-137

核心组件

章节来源 - agent/src/channels/base.py:85-163 - agent/src/channels/telegram.py:349-356 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/discord.py:40-48 - agent/src/channels/discord.py:246-319 - agent/src/channels/runtime.py:34-102 - agent/src/channels/manager.py:36-137 - frontend/src/stores/agent.ts:136-165

架构总览

下图展示从入站到流式输出的端到端流程,包含通道抽象、平台适配、运行时调度与前端消费。

sequenceDiagram participant Client as "客户端/用户" participant Bus as "消息总线" participant RT as "ChannelRuntime" participant CM as "ChannelManager" participant CH as "通道(BaseChannel)" participant TG as "TelegramChannel" participant DC as "DiscordChannel" participant FE as "前端Store" Client->>Bus : 入站消息 Bus->>RT : publish_inbound RT->>CM : 路由到对应通道 CM->>CH : 选择具体通道实现 CH->>TG : send_delta / send_reasoning_delta TG-->>Client : edit_message_text(渐进更新) CH->>DC : send_delta / send_reasoning_delta DC-->>Client : 分片/嵌入更新 FE->>FE : appendDelta/status(增量拼接/状态切换)

图表来源 - agent/src/channels/base.py:85-163 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/discord.py:246-319 - agent/src/channels/runtime.py:34-102 - agent/src/channels/manager.py:36-137 - frontend/src/stores/agent.ts:136-165

详细组件分析

BaseChannel 流式契约与默认行为

classDiagram class BaseChannel { +name +display_name +send_progress +send_tool_hints +show_reasoning +supports_streaming bool +send_delta(chat_id, delta, metadata) +send_reasoning_delta(chat_id, delta, metadata) +send_reasoning_end(chat_id, metadata) +send_reasoning(msg) +_handle_message(...) }

图表来源 - agent/src/channels/base.py:22-163

章节来源 - agent/src/channels/base.py:85-163

Telegram 流式实现:编辑消息与富文本回退

flowchart TD Start(["send_delta 入口"]) --> CheckEnd{"是否 _stream_end?"} CheckEnd --> |是| Finalize["最终输出: rich/HTML/纯文本"] Finalize --> Cleanup["清理缓冲区/停止打字/移除反应"] CheckEnd --> |否| GetBuf["获取/创建 _StreamBuf"] GetBuf --> Append["追加 delta 到 buf.text"] Append --> FirstMsg{"是否首条消息?"} FirstMsg --> |是| SendPreview["发送预览消息"] FirstMsg --> |否| Interval{"是否达到编辑间隔?"} Interval --> |否| End(["返回"]) Interval --> |是| Overflow{"是否超长?"} Overflow --> |是| Flush["分片: 编辑主片段/发送中间片段/新建尾部消息"] Overflow --> |否| Edit["edit_message_text 渐进更新"] Flush --> End Edit --> End Cleanup --> End

图表来源 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/telegram.py:1061-1098 - agent/src/channels/telegram.py:860-880

章节来源 - agent/src/channels/telegram.py:349-356 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/telegram.py:1061-1098 - agent/src/channels/telegram.py:860-880

Discord 流式实现:嵌入更新与分片发送

classDiagram class DiscordChannel { +name +_STREAM_EDIT_INTERVAL +_stream_bufs dict +_known_channels dict +send_outbound(msg) +_build_chunks(content, failed_media, sent_media) list } class _StreamBuf { +text str +message Any +last_edit float +stream_id str } DiscordChannel --> _StreamBuf : "按聊天聚合"

图表来源 - agent/src/channels/discord.py:40-48 - agent/src/channels/discord.py:246-319

章节来源 - agent/src/channels/discord.py:40-48 - agent/src/channels/discord.py:246-319

运行时与通道管理器

sequenceDiagram participant App as "应用" participant RT as "ChannelRuntime" participant CM as "ChannelManager" participant CH as "通道实例" App->>RT : start(start_manager=True) RT->>CM : start_all() CM->>CM : 发现/初始化通道 CM-->>RT : 通道就绪 RT->>RT : 启动消费者任务 App->>RT : stop() RT->>RT : 取消任务/清理 RT->>CM : stop_all()

图表来源 - agent/src/channels/runtime.py:72-102 - agent/src/channels/manager.py:36-137

章节来源 - agent/src/channels/runtime.py:72-102 - agent/src/channels/manager.py:36-137

前端流式消费:增量拼接与会话隔离

flowchart TD S(["收到SSE/WS增量"]) --> A["appendDelta(delta)"] A --> U["UI渲染更新"] S --> ST["setStatus('streaming'|'idle')"] ST --> SI{"是否streaming且sessionId存在?"} SI --> |是| SetSID["设置streamingSessionId=sessionId"] SI --> |否| ClearSID["清理streamingSessionId"]

图表来源 - frontend/src/stores/agent.ts:136-165

章节来源 - frontend/src/stores/agent.ts:136-165

依赖关系分析

graph LR Base["BaseChannel"] --> TG["TelegramChannel"] Base --> DC["DiscordChannel"] CM["ChannelManager"] --> TG CM --> DC RT["ChannelRuntime"] --> CM FE["前端Store"] --> RT

图表来源 - agent/src/channels/base.py:22-163 - agent/src/channels/manager.py:36-137 - agent/src/channels/runtime.py:34-102 - frontend/src/stores/agent.ts:136-165

章节来源 - agent/src/channels/base.py:22-163 - agent/src/channels/manager.py:36-137 - agent/src/channels/runtime.py:34-102 - frontend/src/stores/agent.ts:136-165

性能考虑

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

故障排查指南

章节来源 - agent/src/channels/telegram.py:860-880 - agent/src/channels/telegram.py:920-1060 - agent/src/channels/discord.py:246-319 - frontend/src/stores/agent.ts:136-165

结论

Vibe-Trading 的流式传输以 BaseChannel 为契约,结合 Telegram 与 Discord 的平台特性,实现了稳定、高效、可扩展的增量更新机制。通过编辑节流、分片溢出处理、重试与限流、会话隔离与前端状态管理,系统在长文本、富媒体与交互式内容场景下具备良好表现。建议在新增渠道时遵循该契约,复用重试与节流模式,并确保正确的状态清理与会话隔离。

[本节总结性内容,无需特定文件引用]

附录

[本节为路径索引,无需额外分析]