#!/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))}")