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

126 lines
5.6 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# 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秒折叠 |
## 八、设计文档关联
- [holding-strategies-multewriter-analysis.md](holding-strategies-multewriter-analysis.md) — 多写方分析(A-E组)
- [llm-execution-protocol.md](llm-execution-protocol.md) — LLM执行协议
- [evolution-archive-readme.md](evolution-archive-readme.md) — 进化模块归档说明