# 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) — 进化模块归档说明