From 5eeb31668d3326260b867a4a7537570ae1693f82 Mon Sep 17 00:00:00 2001 From: xxm Date: Fri, 21 Aug 2026 15:09:48 +0800 Subject: [PATCH] =?UTF-8?q?feat(broadcast):=20=E6=92=AD=E6=8A=A5=E7=B3=BB?= =?UTF-8?q?=E7=BB=9FTab+API+=E5=8E=86=E5=8F=B2=E6=9F=A5=E8=AF=A2+=E6=B6=88?= =?UTF-8?q?=E6=81=AF=E5=88=86=E7=B1=BB?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- deploy/profile-scripts/broadcast.py | 93 ++++++++++++++++++++++++++++ server.py | 27 ++++++++ static/broadcast.html | 96 +++++++++++++++++++++++++++++ 3 files changed, 216 insertions(+) create mode 100644 deploy/profile-scripts/broadcast.py create mode 100644 static/broadcast.html diff --git a/deploy/profile-scripts/broadcast.py b/deploy/profile-scripts/broadcast.py new file mode 100644 index 00000000..7f08e320 --- /dev/null +++ b/deploy/profile-scripts/broadcast.py @@ -0,0 +1,93 @@ +#!/usr/bin/env python3 +"""broadcast.py — 播报系统核心模块 +消息分类 + DB写入 + API查询接口 +""" +import os, re, sqlite3, json +from datetime import datetime, timedelta + +DB = "/home/hmo/MoFin/data/mofin.db" + +# === 消息分类规则 === +CATEGORIES = { + "trading": ["买入", "卖出", "止损", "止盈", "加仓", "减仓", "调仓", "持仓异动", "推荐", "操作建议", "区间触发", "价格监控"], + "system_error": ["LLM端点故障", "API错误", "连接失败", "超时", "异常", "error", "失败", "Error", "Exception"], + "health": ["健康检查", "系统体检", "数据完整性", "cron", "部署"], + "market": ["大盘", "市场", "板块", "行业", "指数", "行情", "涨跌"], + "news": ["新闻", "消息面", "资讯", "公告", "政策"], + "strategy": ["策略", "重评", "评估", "温区", "regime"], + "general": [], +} + +def classify(title, content): + """根据标题和内容分类消息""" + text = (title + " " + content).lower() + for cat, keywords in CATEGORIES.items(): + if cat == "general": + continue + for kw in keywords: + if kw.lower() in text: + return cat + return "general" + +def save_message(ts, title, content, source="xmpp"): + """保存消息到DB""" + category = classify(title, content) + conn = sqlite3.connect(DB, timeout=30) + conn.execute( + "INSERT INTO broadcast_messages (ts, category, title, content, source) VALUES (?,?,?,?,?)", + (ts, category, title, content, source)) + conn.commit() + conn.close() + return category + +def get_recent(hours=72, category=None, limit=200): + """获取最近N小时的消息""" + conn = sqlite3.connect(DB, timeout=30) + conn.row_factory = sqlite3.Row + since = (datetime.now() - timedelta(hours=hours)).isoformat() + query = "SELECT * FROM broadcast_messages WHERE ts >= ?" + params = [since] + if category: + query += " AND category=?" + params.append(category) + query += " ORDER BY ts DESC LIMIT ?" + params.append(limit) + rows = conn.execute(query, params).fetchall() + conn.close() + return [dict(r) for r in rows] + +def search_history(start_date=None, end_date=None, keyword=None, category=None, limit=100): + """历史消息搜索""" + conn = sqlite3.connect(DB, timeout=30) + conn.row_factory = sqlite3.Row + query = "SELECT * FROM broadcast_messages WHERE 1=1" + params = [] + if start_date: + query += " AND ts >= ?" + params.append(start_date) + if end_date: + query += " AND ts <= ?" + params.append(end_date + " 23:59:59") + if keyword: + query += " AND (title LIKE ? OR content LIKE ?)" + params.extend([f"%{keyword}%", f"%{keyword}%"]) + if category: + query += " AND category=?" + params.append(category) + query += " ORDER BY ts DESC LIMIT ?" + params.append(limit) + rows = conn.execute(query, params).fetchall() + conn.close() + return [dict(r) for r in rows] + +def archive_old(days=7): + """归档超过N天的消息""" + cutoff = (datetime.now() - timedelta(days=days)).isoformat() + conn = sqlite3.connect(DB, timeout=30) + conn.execute("UPDATE broadcast_messages SET archived=1 WHERE ts < ? AND archived=0", (cutoff,)) + conn.commit() + conn.close() + +if __name__ == "__main__": + print("broadcast 模块已加载") + print(f" 当前消息数: {len(get_recent(hours=9999))}") diff --git a/server.py b/server.py index 10340396..904be14e 100644 --- a/server.py +++ b/server.py @@ -2445,6 +2445,33 @@ def api_research_execution_log(): return jsonify([dict(r) for r in rows]) +# ── Broadcast System API ── +@app.route("/api/broadcast/recent") +def api_broadcast_recent(): + from broadcast import get_recent + hours = int(request.args.get("hours", "72")) + category = request.args.get("category") + limit = int(request.args.get("limit", "200")) + return jsonify(get_recent(hours=hours, category=category, limit=limit)) + +@app.route("/api/broadcast/search") +def api_broadcast_search(): + from broadcast import search_history + start = request.args.get("start") + end = request.args.get("end") + keyword = request.args.get("keyword") + category = request.args.get("category") + limit = int(request.args.get("limit", "100")) + return jsonify(search_history(start_date=start, end_date=end, keyword=keyword, category=category, limit=limit)) + +@app.route("/api/broadcast/archive", methods=["POST"]) +def api_broadcast_archive(): + from broadcast import archive_old + days = int(request.args.get("days", "7")) + archive_old(days=days) + return jsonify({"ok": True}) + + if __name__ == "__main__": port = int(os.environ.get("PORT", 8899)) print(f"🚀 MoFin Dashboard → http://0.0.0.0:{port}") diff --git a/static/broadcast.html b/static/broadcast.html new file mode 100644 index 00000000..bb02cba3 --- /dev/null +++ b/static/broadcast.html @@ -0,0 +1,96 @@ + + + + + +Broadcast History - MoFin + + + +← Back to Dashboard +

Broadcast History

+
+ + + + + +
+
+ + + +