Files
MoFin/docs/message-architecture-20260821.md

5.6 KiB
Raw Permalink Blame History

MoFin 消息架构重构(2026-08-21

阶段总结:建立统一消息生产-消费架构,整合播报系统、异常浮窗、消息通道路由。

一、背景与问题

知微(hermes)通过 XMPP 发送的消息过多,真正有用的操作推荐被埋没。同时存在多套独立的消息系统,导致:

  • 交易推荐和系统异常混在一起
  • 部分消息(如 deploy_guard 回滚通知)只在 XMPP,系统界面看不到
  • 右上角异常浮窗显示几天前的旧异常
  • 各脚本消息输出不可控(直接 print 就被推 xmpp)

二、统一消息架构

2.1 核心原则

  • 生产者(脚本)与消费者(消息处理者)解耦:脚本只负责发消息,不关心去向
  • 消费者按"生产者是谁"(哪个 job)决定消息通道,不按消息内容判断
  • broadcast_messages 表是唯一消息源,所有消息都进这里

2.2 三层结构

[生产者] 77个 cron 脚本
    └─ 每个脚本开头: from messenger import install_stdio_hook; install_stdio_hook()
           (自动识别当前脚本/job, 拦截 stdout)
    ▼
[消费者] messenger.py (统一消息处理者)
    ├─ 读 jobs.json 该脚本的 delivery 配置
    ├─ broadcast 通道 → print 被吞, 内容写 broadcast_messages 表
    ├─ xmpp 通道 → print 输出给 hermes(推xmpp) + 写 broadcast 表
    └─ both 通道 → 两者
    ▼
[通道]
    ├─ broadcast_messages 表 (播报Tab/历史查询)
    └─ XMPP (操作推荐/交易触发)

2.3 通道配置(jobs.json delivery 字段)

  • broadcast:只写 broadcast_messages 表(不推 xmpp
  • xmpp:写广播表 + 推 xmpp
  • both:两者

三、组件清单

3.1 数据库表

用途
broadcast_messages 统一消息源(id/ts/category/title/content/source/archived/created_at

3.2 核心模块

模块 职责
messenger.py 统一消息处理者(分类+路由+hook拦截)
broadcast.py 播报查询(get_recent/search_history/archive_old
alert_logger.py 异常记录(写 alerts.json + broadcast
alert_helper.py 告警网关(分级推送 + broadcast 归档)

3.3 API 端点

端点 功能
GET /api/broadcast/recent?hours=72 最近72h消息
GET /api/broadcast/search 历史搜索(日期/关键字/类型)
POST /api/broadcast/archive 每日归档
POST /api/broadcast/toggle_delivery 切换 job 消息通道

3.4 前端

页面/组件 功能
播报 Tab (index.html) 72h 消息流,分类标签+颜色,点击浮动全文弹窗
broadcast.html 历史消息查询(日期范围+关键字+类型筛选)
健康 Tab cron 列表 新增"消息通道"列(点击切换 broadcast/xmpp/both
右上角异常浮窗 从 broadcast system_error 类读取(24h内)

四、消息分类(category

分类 关键词 去向
trading 买入/卖出/止损/止盈/加仓/持仓异动/推荐/区间触发 xmpp + 播报
system_error LLM端点故障/API错误/超时/异常/失败/告警 播报 + 异常浮窗
health 健康检查/系统体检/数据完整性 播报
market 大盘/市场/板块/行业/指数 播报
news 新闻/消息面/资讯/公告/政策 播报
strategy 策略/重评/评估/温区/12维 播报
general 其他 播报

五、异常浮窗行为规则

  • 24小时过滤:只有24小时内的异常才显示
  • 新异常自动展开:检测到从未展开过的异常 → 展开 → 3秒后自动折叠
  • 已展开不重复localStorage expandedAlertKeys 持久化,旧异常不再展开
  • 首次加载不弹旧异常alertsInitialized 标记,首次打开把当前异常标记为已展开
  • 24小时后算历史:超过24小时的异常自动从"最新异常"移除

六、Cron Job 清理

已禁用(30个失效/冗余 job

  • 引用已归档模块:branch_scanner/prune_branches/meta_growth/meta_watchdog/stale_push_wlin/strategy_evaluator/advice_reconciliation/ab_research/evolution_daily
  • 空脚本:市场精选推荐/宏观风险扫描系列
  • 功能重复:cron_to_xmpp(与 messenger 冲突)、stale_detector + strategy-staleness-check(每日重评覆盖)

最终架构(77 enabled

  • xmpp 通道(7个):价格监控/买入区提醒/持仓异动/知微洞察/策略复盘/重评跟进/自选退出
  • broadcast 通道(70个):采集/加工/健康/备份

七、修复记录

问题 修复
price_monitor ModuleNotFoundError alert_logger.py 复制到 deploy
realtime_indicators no such column id 移除 SELECT id
realtime_indicators 写错日期 改用 live_prices 实时价格
hermes cron 不执行 job jobs.json 缺 id 字段(知微修复)
hermes provider 429 ocg-new → ocg-router(多key路由)
daily_kline 600s超时 script_timeout_seconds: 1800 + sh统一 no_agent
promote_candidates DB锁 _exec_retry 自动重试
xmpp_logger score TypeError usage_percent None 兜底
health_monitor_daily 引用 evolution 禁用 job
messenger hook 递归 写表直接操作DB不经send
cycDelivery 非法正则 \x00-\x7F → [^ -~]
alert_float 旧异常 24h过滤+已展开记录+3秒折叠

八、设计文档关联