diff --git a/README.md b/README.md index 72d64a5..9388a6a 100644 --- a/README.md +++ b/README.md @@ -35,11 +35,19 @@
| 看板 Dashboard | -选股 Screener | +策略 Screener | ||||||||||||||
![]() |
- ![]() |
+ ![]() |
+ ||||||||||||||
| 监控中心 Monitor | +概念分析 Concept | +|||||||||||||||
![]() |
+ ![]() |
|||||||||||||||
| 回测 Backtest | @@ -112,11 +120,16 @@ - **SSE 流式进度**:长任务实时推送进度,支持刷新 / 切页后**重连恢复**(相同参数任务只启动一次) - **统计输出**:净值曲线 · 夏普 · 最大回撤 · 胜率 · 每笔交易明细 -### 📡 实时监控(Strategy Monitor) +### 📡 监控中心(Monitor) -- **盘中 SSE 推送**:行情刷新(`quotes_updated`)+ 策略告警(`strategy_alert`)双事件流,前端实时弹通知 -- **策略监控**:订阅策略的 entry / exit 信号 + 自定义提醒条件(如 `rsi_14 > 80`),命中即推送 -- **Webhook 告警**:命中规则后可选推送外部 webhook +**统一监控规则引擎** —— 一个页面管理所有类型的监控,实时推送 + 持久化触发记录: + +- **四类监控**:策略监控 · 个股信号监控(选信号即加) · 个股价格/涨跌监控 · 全市场异动监控 +- **灵活条件**:多条件 AND/OR 组合 + 冷却期去重(防刷屏) + 严重级别(info/warn/critical) +- **多入口配置**:监控中心页面新建规则 · 个股详情页「加监控」· 策略卡片一键开启 +- **实时 SSE 推送**:命中规则后右下角弹窗通知(可配声效) + 持久化到 `alerts.jsonl` +- **触发记录**:时间倒序展示,支持按来源过滤 · 单条删除 · 清空 · 点击查看个股日K +- **菜单未读徽标**:离开监控中心后有新触发,菜单显示未读数;进入页面后清零 ### 🤖 AI 策略生成(可选) @@ -207,7 +220,7 @@ pnpm dev # http://localhost:3011 3. **自选**页:添加跟踪标的;点代码进 **K 线**页看蜡烛图 + 买卖点 4. **选股**页:点任一内置策略卡片即时扫描;或用自定义信号组合条件 5. **回测**页:选策略 / 信号 + 时间区间 → 跑回测 → 看净值 / 夏普 / 交易明细(SSE 实时进度) -6. **监控**页:配置告警规则,盘中 SSE 推送行情与策略信号;命中后写入告警日志(可选 webhook) +6. **监控中心**页:配置监控规则(策略/个股信号/价格/市场异动),盘中 SSE 实时弹窗通知 + 持久化触发记录;或在个股详情页点「加监控」快速添加 --- @@ -273,7 +286,8 @@ DATA_DIR=./data # Parquet / DuckDB 数据存储目录 | **2** | Polars enriched 流水线 + Screener + 信号扫描 | ✅ | | **3** | vectorbt 回测 + T+1 + 手续费 + 止损 + max-hold | ✅ | | **4** | 监控引擎 + 告警规则 + Webhook + APScheduler 盘后定时 | ✅ | -| **v2** | 自定义信号 / 策略商店 / AI 策略生成 / 外部数据源插件 / 早晚报 / Onboarding | 🚧 | +| **5** | 统一监控中心 + 四类监控规则 + 实时推送 + 持久化触发记录 + 声效通知 | ✅ | +| **v2** | Webhook 推送(QMT/掘金下单) · 板块异动 · 早晚报 · 更多扩展 | 🚧 | --- diff --git a/VERSION b/VERSION index 02dc6f6..d516242 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -v0.1.19 +v0.1.28 diff --git a/backend/app/__init__.py b/backend/app/__init__.py index 688ef21..b5d5881 100644 --- a/backend/app/__init__.py +++ b/backend/app/__init__.py @@ -1,3 +1,3 @@ """TickFlow Stock Panel backend.""" -__version__ = "0.1.0" +__version__ = "0.1.28" diff --git a/backend/app/api/alerts.py b/backend/app/api/alerts.py new file mode 100644 index 0000000..ab27b10 --- /dev/null +++ b/backend/app/api/alerts.py @@ -0,0 +1,129 @@ +"""告警触发记录 API — 查询/清空/生成演示数据 alerts.jsonl。""" +from __future__ import annotations + +import random +import time +from pathlib import Path + +from fastapi import APIRouter, HTTPException, Request + +from app.services import alert_store + +router = APIRouter(prefix="/api/alerts", tags=["alerts"]) + + +def _data_dir(request: Request) -> Path: + return request.app.state.repo.store.data_dir + + +@router.get("") +def list_alerts( + request: Request, + days: int = 7, + limit: int = 5000, + source: str | None = None, + type: str | None = None, +): + """查询触发记录 (时间倒序)。""" + events = alert_store.list_recent( + _data_dir(request), days=days, limit=limit, source=source, type=type, + ) + total = alert_store.count(_data_dir(request)) + return {"alerts": events, "total": total} + + +@router.delete("") +def clear_alerts(request: Request): + """清空全部触发记录。""" + n = alert_store.clear(_data_dir(request)) + return {"ok": True, "cleared": n} + + +@router.delete("/{ts}") +def delete_alert(ts: int, request: Request): + """删除单条触发记录 (按 ts 毫秒时间戳)。""" + deleted = alert_store.delete_one(_data_dir(request), ts) + if not deleted: + raise HTTPException(status_code=404, detail="记录不存在") + return {"ok": True} + + +# ── 演示数据生成 (仅 Dev 页用) ───────────────────────── + +_DEMO_STOCKS = [ + ("600519.SH", "贵州茅台"), ("000001.SZ", "平安银行"), ("300750.SZ", "宁德时代"), + ("002594.SZ", "比亚迪"), ("000858.SZ", "五粮液"), ("601318.SH", "中国平安"), + ("002475.SZ", "立讯精密"), ("600036.SH", "招商银行"), ("000725.SZ", "京东方A"), + ("300059.SZ", "东方财富"), +] +_DEMO_TEMPLATES = [ + ("signal", "MA金叉触发", ["signal_ma_golden_5_20"], "info"), + ("signal", "放量突破新高", ["signal_volume_surge", "signal_n_day_high"], "warn"), + ("signal", "MACD金叉", ["signal_macd_golden"], "info"), + ("signal", "跌破MA20", ["signal_ma20_breakdown"], "info"), + ("price", "涨幅超 5%", [], "warn"), + ("price", "RSI 极度超卖", [], "warn"), + ("price", "跌幅超 3%", [], "info"), + ("market", "涨停封板", ["signal_limit_up"], "critical"), + ("market", "连板异动", ["signal_limit_up"], "warn"), + ("market", "炸板", ["signal_broken_limit_up"], "warn"), + ("strategy", "策略「趋势突破」买入信号", ["signal_n_day_high", "signal_volume_surge"], "info"), + ("strategy", "策略「趋势突破」卖出信号", ["signal_ma20_breakdown"], "info"), + ("strategy", "策略「新低反转」买入信号", ["signal_n_day_low"], "warn"), +] + + +@router.post("/seed") +def seed_demo_alerts(request: Request, count: int = 12, recent: bool = True): + """生成演示触发记录 (Dev 页用)。 + + Args: + count: 生成条数 (1-50) + recent: True=时间戳设为"刚刚"(用于测试闪烁效果); False=分散在近3天 + """ + count = max(1, min(50, count)) + now_ms = int(time.time() * 1000) + events = [] + for i in range(count): + source, message, signals, severity = _DEMO_TEMPLATES[i % len(_DEMO_TEMPLATES)] + sym, name = _DEMO_STOCKS[i % len(_DEMO_STOCKS)] + # recent 模式: 时间戳从现在往前每条错开 30 秒 (最新在前) + ts = now_ms - (i * 30000) if recent else now_ms - random.randint(60, 4320) * 60 * 1000 + events.append({ + "ts": ts, + "rule_id": f"demo_rule_{i}", + "rule_name": message, + "source": source, + "type": source, + "symbol": sym, + "name": name, + "message": message, + "price": round(random.uniform(8, 1800), 2), + "change_pct": round(random.uniform(-0.06, 0.098), 4), + "signals": signals, + "severity": severity, + }) + alert_store.append_many(_data_dir(request), events) + + # 同步推入 SSE 队列, 让所有连着 SSE 的客户端实时收到 (不依赖轮询) + qs = getattr(request.app.state, "quote_service", None) + if qs: + # 转成 SSE 推送格式 (和 _evaluate_monitors 一致) + sse_alerts = [{ + "source": ev["source"], + "type": ev["type"], + "rule_id": ev.get("rule_id"), + "symbol": ev["symbol"], + "name": ev["name"], + "message": ev["message"], + "price": ev["price"], + "change_pct": ev["change_pct"], + "signals": ev["signals"], + "severity": ev.get("severity", "info"), + } for ev in events] + with qs._lock: + qs._pending_alerts.extend(sse_alerts) + qs._alert_event.set() + + return {"ok": True, "generated": len(events)} + diff --git a/backend/app/api/data.py b/backend/app/api/data.py index b1feec7..2afe935 100644 --- a/backend/app/api/data.py +++ b/backend/app/api/data.py @@ -701,10 +701,29 @@ def table_schema(request: Request, table: str) -> list[dict]: @router.get("/version") def get_version(request: Request) -> dict: - """返回当前项目版本号(读取项目根目录 VERSION 文件)。""" + """返回当前项目版本号。 + + 优先从 pyproject.toml 读取 (项目权威版本源), + 回退到 VERSION 文件, 最后兜底 v0.0.0。 + """ from app.config import settings - version_file = Path(settings.data_dir).parent / "VERSION" - version = "v0.0.0" + project_root = Path(settings.data_dir).parent + + # 1. 优先读 pyproject.toml + pyproject = project_root / "pyproject.toml" + if pyproject.exists(): + for line in pyproject.read_text(encoding="utf-8").splitlines(): + if line.strip().startswith("version"): + # version = "0.1.28" + v = line.split("=", 1)[1].strip().strip('"').strip("'") + if v: + return {"version": f"v{v}" if not v.startswith("v") else v} + + # 2. 回退到 VERSION 文件 + version_file = project_root / "VERSION" if version_file.exists(): - version = version_file.read_text(encoding="utf-8").strip() or version - return {"version": version} + v = version_file.read_text(encoding="utf-8").strip() + if v: + return {"version": v} + + return {"version": "v0.0.0"} diff --git a/backend/app/api/kline.py b/backend/app/api/kline.py index 308282b..860ece8 100644 --- a/backend/app/api/kline.py +++ b/backend/app/api/kline.py @@ -55,6 +55,21 @@ def search_instruments( return {"results": rows} +@router.post("/instruments/names") +def instruments_names(request: Request, symbols: list[str]): + """批量查股票名称。传入 symbol 列表, 返回 {symbol: name}。""" + if not symbols: + return {"names": {}} + repo = request.app.state.repo + df = repo.get_instruments() + if df.is_empty(): + return {"names": {}} + import polars as pl + matched = df.filter(pl.col("symbol").is_in(symbols)).select(["symbol", "name"]) + names = {row["symbol"]: row["name"] for row in matched.iter_rows(named=True)} + return {"names": names} + + def _get_stock_info(repo, symbol: str) -> dict: """从 instruments 视图查标的名称 + 股本。""" try: diff --git a/backend/app/api/monitor_rules.py b/backend/app/api/monitor_rules.py new file mode 100644 index 0000000..4902231 --- /dev/null +++ b/backend/app/api/monitor_rules.py @@ -0,0 +1,212 @@ +"""监控规则 API 路由 — HTTP 请求 → 调用 monitor_rules 模块 → 同步引擎内存态。 + +只做胶水: 校验 → 持久化 → 失效引擎内存态。不含评估逻辑。 +""" +from __future__ import annotations + +from pathlib import Path + +from fastapi import APIRouter, HTTPException, Request +from pydantic import BaseModel + +from app.strategy import monitor_rules + +router = APIRouter(prefix="/api/monitor-rules", tags=["monitor-rules"]) + + +def _data_dir(request: Request) -> Path: + return request.app.state.repo.store.data_dir + + +def _sync_engine(request: Request) -> None: + """保存/删除后,把最新规则集 reload 到引擎内存态。""" + engine = getattr(request.app.state, "monitor_engine", None) + if engine is not None: + rules = monitor_rules.load_all(_data_dir(request)) + engine.set_rules(rules) + + +# ── Pydantic 模型 ─────────────────────────────────────── +class ConditionModel(BaseModel): + field: str + op: str # truth | > >= < <= == != + value: float | None = None # op 非 truth 时必填 + + +class RuleModel(BaseModel): + id: str + name: str + enabled: bool = True + type: str # strategy | signal | price | market + scope: str = "symbols" # symbols | all | sector + symbols: list[str] = [] + sector: str | None = None + strategy_id: str | None = None + direction: str = "entry" # entry | exit | both + conditions: list[ConditionModel] = [] + logic: str = "and" # and | or + cooldown_seconds: int = 3600 + severity: str = "info" # info | warn | critical + webhook_url: str = "" # Webhook 推送地址 (推送到 QMT 等外部软件, 开发中) + webhook_enabled: bool = False + message: str = "" + + +# ── 字段选项 ───────────────────────────────────────────── +@router.get("/options") +def get_options(request: Request): + """返回可选字段、信号列、运算符、枚举,供前端表单使用。""" + from app.indicators.pipeline import ENRICHED_COLUMNS + from app.strategy.custom_signals import ALLOWED_FIELDS, load_all as load_csg + + # 阈值字段 (带中文标签) + threshold_fields = [ + {"key": f, "label": ENRICHED_COLUMNS.get(f, f)} + for f in sorted(ALLOWED_FIELDS) + ] + # 内置信号列 (布尔, 用于 op=truth) + builtin_signals = [ + {"key": k, "label": v} + for k, v in ENRICHED_COLUMNS.items() + if k.startswith("signal_") + ] + # 自定义信号列 (csg_) + custom_sigs = [] + try: + for cs in load_csg(_data_dir(request)): + if cs.get("enabled") is not False: + custom_sigs.append({ + "key": f"csg_{cs['id']}", + "label": cs.get("name", cs["id"]), + }) + except Exception: + pass + + return { + "threshold_fields": threshold_fields, + "builtin_signals": builtin_signals, + "custom_signals": custom_sigs, + "operators": [">", ">=", "<", "<=", "==", "!="], + "types": [ + {"key": "signal", "label": "个股信号"}, + {"key": "price", "label": "价格/涨跌"}, + {"key": "market", "label": "市场异动"}, + {"key": "strategy", "label": "策略监控"}, + ], + "scopes": [ + {"key": "symbols", "label": "指定股票"}, + {"key": "all", "label": "全市场"}, + {"key": "sector", "label": "板块"}, + ], + "logics": [ + {"key": "and", "label": "全部满足 (AND)"}, + {"key": "or", "label": "任一满足 (OR)"}, + ], + "severities": [ + {"key": "info", "label": "普通"}, + {"key": "warn", "label": "警告"}, + {"key": "critical", "label": "重要"}, + ], + "directions": [ + {"key": "entry", "label": "买入"}, + {"key": "exit", "label": "卖出"}, + {"key": "both", "label": "买卖都报"}, + ], + } + + +# ── 列表 ─────────────────────────────────────────────── +@router.get("") +def list_rules(request: Request): + rules = monitor_rules.load_all(_data_dir(request)) + # 按 created_at 倒序 + rules.sort(key=lambda r: r.get("created_at", ""), reverse=True) + return {"rules": rules} + + +# ── 新建 / 更新 ──────────────────────────────────────── +@router.post("") +def save_rule(req: RuleModel, request: Request): + rule = monitor_rules.normalize(req.model_dump()) + # 编辑现有规则时, 保留原 created_at (避免按时间排序时位置跳动) + existing = monitor_rules.load_one(_data_dir(request), rule["id"]) + if existing and existing.get("created_at"): + rule["created_at"] = existing["created_at"] + try: + monitor_rules.validate(rule) + except ValueError as e: + raise HTTPException(status_code=400, detail=str(e)) + monitor_rules.save_one(_data_dir(request), rule) + _sync_engine(request) + return {"ok": True, "rule": rule} + + +# ── 删除 ─────────────────────────────────────────────── +@router.delete("/{rule_id}") +def delete_rule(rule_id: str, request: Request): + if not monitor_rules.ID_RE.match(rule_id): + raise HTTPException(status_code=400, detail="规则 id 非法") + deleted = monitor_rules.delete_one(_data_dir(request), rule_id) + if not deleted: + raise HTTPException(status_code=404, detail="规则不存在") + _sync_engine(request) + return {"ok": True} + + +# ── 演示数据生成 (仅 Dev 页用) ───────────────────────── + +import time as _time +from datetime import datetime, timezone + + +def _demo_rule(rule_id: str, name: str, rtype: str, scope: str, symbols: list[str], + conditions: list[dict], logic: str = "or", cooldown: int = 3600, + severity: str = "info", message: str = "") -> dict: + return monitor_rules.normalize({ + "id": rule_id, + "name": name, + "type": rtype, + "scope": scope, + "symbols": symbols, + "conditions": conditions, + "logic": logic, + "cooldown_seconds": cooldown, + "severity": severity, + "message": message, + "enabled": True, + }) + + +_DEMO_RULES_TEMPLATE = [ + ("个股信号 · 茅台放量突破", "signal", "symbols", ["600519.SH"], + [{"field": "signal_volume_surge", "op": "truth"}, + {"field": "signal_n_day_high", "op": "truth"}], "or", "info"), + ("个股信号 · 宁德金叉", "signal", "symbols", ["300750.SZ"], + [{"field": "signal_ma_golden_5_20", "op": "truth"}], "or", "info"), + ("价格 · 平安跌幅监控", "price", "symbols", ["000001.SZ"], + [{"field": "change_pct", "op": "<", "value": -0.03}], "or", "warn", "warn"), + ("价格 · 比亚迪RSI超卖", "price", "symbols", ["002594.SZ"], + [{"field": "rsi_14", "op": "<", "value": 30}], "and", "warn", "warn"), + ("市场异动 · 全市场涨停", "market", "all", [], + [{"field": "signal_limit_up", "op": "truth"}], "or", "critical", "critical"), + ("市场异动 · 全市场炸板", "market", "all", [], + [{"field": "signal_broken_limit_up", "op": "truth"}], "or", "warn", "warn"), + ("市场异动 · 跌幅超5%", "market", "all", [], + [{"field": "change_pct", "op": "<", "value": -0.05}], "or", "warn", "warn"), + ("个股信号 · 茅台跌破MA20", "signal", "symbols", ["600519.SH"], + [{"field": "signal_ma20_breakdown", "op": "truth"}], "or", "info"), +] + + +@router.post("/seed") +def seed_demo_rules(request: Request): + """生成演示监控规则 (Dev 页用)。覆盖 signal/price/market 三类。""" + ts = int(_time.time() * 1000) + created = [] + for i, (name, rtype, scope, symbols, conditions, logic, severity, sev) in enumerate(_DEMO_RULES_TEMPLATE): + rule_id = f"demo_{ts}_{i}" + rule = _demo_rule(rule_id, name, rtype, scope, symbols, conditions, logic, 3600, sev) + monitor_rules.save_one(_data_dir(request), rule) + created.append(rule_id) + _sync_engine(request) + return {"ok": True, "generated": len(created), "ids": created} diff --git a/backend/app/api/settings.py b/backend/app/api/settings.py index aa475ee..8ffa77b 100644 --- a/backend/app/api/settings.py +++ b/backend/app/api/settings.py @@ -352,34 +352,31 @@ class RealtimeMonitorConfigIn(BaseModel): @router.put("/preferences/realtime-monitor") def update_realtime_monitor_config(req: RealtimeMonitorConfigIn, request: Request) -> dict: - """更新实时监控配置。""" + """更新实时监控配置。策略监控统一迁移为 MonitorRule,由监控引擎评估。""" from app.services import preferences cfg = req.model_dump(exclude_none=True) result = preferences.set_realtime_monitor_config(cfg) - # 如果策略监控开关变化,更新 StrategyMonitorService 的监控池 + # 策略监控开关/池变化 → 同步迁移为 type=strategy 规则 + reload 引擎 if req.strategy_monitor_ids is not None or req.strategy_monitor_enabled is not None: - monitor = getattr(request.app.state, "strategy_monitor", None) - if monitor: - if preferences.get_strategy_monitor_enabled(): - # 从策略引擎加载监控配置 - engine = getattr(request.app.state, "strategy_engine", None) - ids = preferences.get_strategy_monitor_ids() - if engine and ids: - monitor.stop_all() - for sid in ids: - try: - s = engine.get(sid) - monitor.start(sid, { - "entry_signals": s.entry_signals, - "exit_signals": s.exit_signals, - "alerts": s.alerts, - }) - except ValueError: - pass - else: - monitor.stop_all() + monitor_engine = getattr(request.app.state, "monitor_engine", None) + strategy_engine = getattr(request.app.state, "strategy_engine", None) + data_dir = request.app.state.repo.store.data_dir + if monitor_engine is not None and strategy_engine is not None: + from app.strategy import monitor_rules as mr_store + try: + if preferences.get_strategy_monitor_enabled(): + ids = preferences.get_strategy_monitor_ids() + names = {s.id: s.name for s in strategy_engine.list_strategies()} + mr_store.migrate_strategy_monitors(data_dir, ids, names) + else: + # 关闭策略监控: 停用所有策略规则 + mr_store.migrate_strategy_monitors(data_dir, [], {}) + # reload 规则到引擎 + monitor_engine.set_rules(mr_store.load_all(data_dir)) + except Exception: + pass return result diff --git a/backend/app/api/strategy.py b/backend/app/api/strategy.py index a462e5e..dae6ddc 100644 --- a/backend/app/api/strategy.py +++ b/backend/app/api/strategy.py @@ -426,36 +426,8 @@ def delete_strategy(strategy_id: str, request: Request): # ── 监控 ───────────────────────────────────────────────────────────── - - -@router.post("/monitor/start") -def monitor_start(req: MonitorStartRequest, request: Request): - engine = _get_engine(request) - monitor = _get_monitor(request) - try: - s = engine.get(req.strategy_id) - except ValueError as e: - raise HTTPException(status_code=404, detail=str(e)) from e - - monitor.start(req.strategy_id, { - "entry_signals": s.entry_signals, - "exit_signals": s.exit_signals, - "alerts": s.alerts, - }) - return {"ok": True, "watching": list(monitor.watching.keys())} - - -@router.post("/monitor/stop/{strategy_id}") -def monitor_stop(strategy_id: str, request: Request): - monitor = _get_monitor(request) - monitor.stop(strategy_id) - return {"ok": True, "watching": list(monitor.watching.keys())} - - -@router.get("/monitor/status") -def monitor_status(request: Request): - monitor = _get_monitor(request) - return {"watching": list(monitor.watching.keys())} +# 注: 策略监控已统一迁移到 MonitorRuleEngine (监控通知页), 旧的 start/stop/status +# 路由已移除。StrategyMonitorService 类保留 (其 _check_signals 被 MonitorRuleEngine 复用)。 # ── 热重载 ─────────────────────────────────────────────────────────── diff --git a/backend/app/main.py b/backend/app/main.py index 0c6590f..219c4f2 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -11,7 +11,7 @@ from fastapi.responses import FileResponse from fastapi.staticfiles import StaticFiles from app import __version__ -from app.api import analysis, backtest, data, ext_data, financials, indices, intraday, kline, overview, pipeline, screener, settings as settings_api, signals, strategy, watchlist +from app.api import analysis, backtest, data, ext_data, financials, indices, intraday, kline, monitor_rules, alerts, overview, pipeline, screener, settings as settings_api, signals, strategy, watchlist from app.api.routes import router as core_router from app.config import settings from app.jobs import daily_pipeline @@ -113,6 +113,32 @@ async def lifespan(app: FastAPI): app.state.strategy_engine = strategy_engine logger.info("strategy engine loaded: %d strategies", len(strategy_engine.list_strategies())) + # 通用监控规则引擎: 启动时 reload 规则到内存态 (修复重启后告警失效) + from app.strategy.monitor import MonitorRuleEngine + from app.strategy import monitor_rules as mr_store + from app.services import preferences + monitor_engine = MonitorRuleEngine() + monitor_engine.set_strategy_engine(strategy_engine) + + # 自动迁移: 把旧 strategy_monitor_ids 同步为 type=strategy 规则 (统一到监控页) + try: + if preferences.get_strategy_monitor_enabled(): + ids = preferences.get_strategy_monitor_ids() + if ids: + names = {s.id: s.name for s in strategy_engine.list_strategies()} + mr_store.migrate_strategy_monitors(store.data_dir, ids, names) + logger.info("strategy monitor migrated: %d strategies", len(ids)) + except Exception as e: # noqa: BLE001 + logger.warning("strategy monitor migration failed: %s", e) + + try: + rules = mr_store.load_all(store.data_dir) + monitor_engine.set_rules(rules) + logger.info("monitor engine loaded: %d rules", monitor_engine.rule_count) + except Exception as e: # noqa: BLE001 + logger.warning("monitor engine load failed: %s", e) + app.state.monitor_engine = monitor_engine + yield if app.state.scheduler: @@ -139,11 +165,13 @@ app = FastAPI( lifespan=lifespan, ) -# 开发期 CORS 允许 Vite dev server +# CORS: 允许局域网访问 (自托管场景, 放开所有来源) +# 注: allow_credentials=True 与 allow_origins=['*'] 不能共存 (浏览器规范), +# 本项目认证走 header (API Key), 不依赖 cookie, 故关闭 credentials 换取通配来源。 app.add_middleware( CORSMiddleware, - allow_origins=["http://localhost:3011", "http://127.0.0.1:3011"], - allow_credentials=True, + allow_origins=["*"], + allow_credentials=False, allow_methods=["*"], allow_headers=["*"], ) @@ -165,6 +193,8 @@ app.include_router(financials.router) app.include_router(settings_api.router) app.include_router(strategy.router) app.include_router(signals.router) +app.include_router(monitor_rules.router) +app.include_router(alerts.router) # 生产期静态文件(前端 dist) _static = Path(settings.static_dir) diff --git a/backend/app/services/alert_store.py b/backend/app/services/alert_store.py new file mode 100644 index 0000000..961331a --- /dev/null +++ b/backend/app/services/alert_store.py @@ -0,0 +1,209 @@ +"""告警触发记录存储 — JSONL 追加写 + 滚动清理。 + +职责: + - 把每次触发的 AlertEvent 追加写入 data/user_data/alerts.jsonl + - 提供查询 (按来源/类型过滤、时间倒序、限量) + - 滚动清理: 保留近 N 天 + 上限 M 条 (取交集) + +设计: + - JSONL 每行一个 JSON 对象,便于增量追加和流式读取 + - 清理策略: 追加后按需 prune (按 ts 删旧),避免文件无限膨胀 + - 读时全量加载到内存过滤 (记录量受上限约束, 5000 条量级无压力) +""" +from __future__ import annotations + +import json +import logging +import threading +from pathlib import Path + +logger = logging.getLogger(__name__) + +# 保留策略 +MAX_DAYS = 7 +MAX_RECORDS = 5000 +# 每隔多少次写入触发一次清理 (避免每次写都 prune) +PRUNE_EVERY = 20 + +_lock = threading.Lock() +_write_count = 0 + + +def _path(data_dir: Path) -> Path: + p = data_dir / "user_data" / "alerts.jsonl" + p.parent.mkdir(parents=True, exist_ok=True) + return p + + +def append(data_dir: Path, event: dict) -> None: + """追加一条触发记录。event 应含 ts(毫秒)、rule_id、source 等字段。""" + line = json.dumps(event, ensure_ascii=False) + with _lock: + p = _path(data_dir) + with p.open("a", encoding="utf-8") as f: + f.write(line + "\n") + global _write_count + _write_count += 1 + if _write_count >= PRUNE_EVERY: + _write_count = 0 + _prune_locked(p) + + +def append_many(data_dir: Path, events: list[dict]) -> None: + """批量追加。""" + if not events: + return + with _lock: + p = _path(data_dir) + with p.open("a", encoding="utf-8") as f: + for ev in events: + f.write(json.dumps(ev, ensure_ascii=False) + "\n") + global _write_count + _write_count += len(events) + if _write_count >= PRUNE_EVERY: + _write_count = 0 + _prune_locked(p) + + +def list_recent( + data_dir: Path, + days: int = MAX_DAYS, + limit: int = MAX_RECORDS, + source: str | None = None, + type: str | None = None, +) -> list[dict]: + """读取近 N 天记录,按时间倒序,支持按 source/type 过滤。""" + import time + cutoff = (time.time() - days * 86400) * 1000 # 毫秒 + out: list[dict] = [] + p = _path(data_dir) + if not p.exists(): + return [] + try: + with p.open("r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + ev = json.loads(line) + except Exception: + continue + if ev.get("ts", 0) < cutoff: + continue + if source and ev.get("source") != source: + continue + if type and ev.get("type") != type: + continue + out.append(ev) + except Exception as e: + logger.warning("alert_store read failed: %s", e) + return [] + # 时间倒序 + 截断 + out.sort(key=lambda x: x.get("ts", 0), reverse=True) + return out[:limit] + + +def clear(data_dir: Path) -> int: + """清空全部记录,返回清除的条数。""" + with _lock: + p = _path(data_dir) + if not p.exists(): + return 0 + count = 0 + try: + with p.open("r", encoding="utf-8") as f: + count = sum(1 for line in f if line.strip()) + except Exception: + pass + p.write_text("", encoding="utf-8") + return count + + +def delete_one(data_dir: Path, ts: int) -> bool: + """删除指定 ts 的单条记录,返回是否删除成功。 + + JSONL 无主键, 用 ts(毫秒时间戳) 作为标识。 + 若存在多条同 ts, 只删第一条。 + """ + with _lock: + p = _path(data_dir) + if not p.exists(): + return False + kept: list[dict] = [] + deleted = False + try: + with p.open("r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + ev = json.loads(line) + except Exception: + continue + if not deleted and ev.get("ts") == ts: + deleted = True + continue + kept.append(ev) + except Exception as e: + logger.warning("alert_store delete_one read failed: %s", e) + return False + if not deleted: + return False + try: + with p.open("w", encoding="utf-8") as f: + for ev in kept: + f.write(json.dumps(ev, ensure_ascii=False) + "\n") + except Exception as e: + logger.warning("alert_store delete_one write failed: %s", e) + return False + return True + return count + + +def count(data_dir: Path) -> int: + """返回当前记录总数。""" + p = _path(data_dir) + if not p.exists(): + return 0 + try: + with p.open("r", encoding="utf-8") as f: + return sum(1 for line in f if line.strip()) + except Exception: + return 0 + + +def _prune_locked(p: Path) -> None: + """(调用方需持锁) 保留近 MAX_DAYS 天 + 上限 MAX_RECORDS 条。""" + import time + cutoff = (time.time() - MAX_DAYS * 86400) * 1000 + kept: list[dict] = [] + try: + with p.open("r", encoding="utf-8") as f: + for line in f: + line = line.strip() + if not line: + continue + try: + ev = json.loads(line) + except Exception: + continue + if ev.get("ts", 0) >= cutoff: + kept.append(ev) + except FileNotFoundError: + return + except Exception as e: + logger.warning("alert_store prune read failed: %s", e) + return + # 上限截断 (保留最新的) + if len(kept) > MAX_RECORDS: + kept.sort(key=lambda x: x.get("ts", 0)) + kept = kept[-MAX_RECORDS:] + # 重写文件 + try: + with p.open("w", encoding="utf-8") as f: + for ev in kept: + f.write(json.dumps(ev, ensure_ascii=False) + "\n") + except Exception as e: + logger.warning("alert_store prune write failed: %s", e) diff --git a/backend/app/services/quote_service.py b/backend/app/services/quote_service.py index c61ae66..fdd6392 100644 --- a/backend/app/services/quote_service.py +++ b/backend/app/services/quote_service.py @@ -57,6 +57,7 @@ class QuoteService: self._update_event = threading.Event() # SSE 通知: 行情更新后 set self._alert_event = threading.Event() # SSE 通知: 有告警时 set self._pending_alerts: list[dict] = [] # 待推送的告警 + self._max_pending_alerts: int = 1000 # 背压上限: 超出丢弃最旧 self._strategy_monitor = None # 延迟注入 self._app_state = None # 延迟注入 (FastAPI app.state) @@ -448,9 +449,7 @@ class QuoteService: # ================================================================ def _evaluate_monitors(self, daily_df: pl.DataFrame, quote_extra: pl.DataFrame | None) -> None: - """行情更新后评估策略监控,并刷新策略结果缓存。""" - from app.services import preferences - + """行情更新后评估统一监控规则引擎,并刷新策略结果缓存。""" try: # 获取 enriched 数据 (刚算好的) enriched_today, enriched_date = self.get_enriched_today() @@ -459,34 +458,50 @@ class QuoteService: all_alerts: list[dict] = [] - # 1. 策略监控评估 - if preferences.get_strategy_monitor_enabled(): - monitor = getattr(self._app_state, "strategy_monitor", None) if self._app_state else None - if monitor and monitor.watching: - strategy_alerts = monitor.on_quote_update(enriched_today) - for a in strategy_alerts: - all_alerts.append({ - "source": "strategy", - "type": a.type, - "strategy_id": a.strategy_id, - "symbol": a.symbol, - "name": a.name, - "message": a.message, - "price": a.price, - "change_pct": a.change_pct, - "signals": a.signals, - }) + # 通用监控规则评估 (统一引擎: signal/price/market/strategy) + if self._app_state: + engine = getattr(self._app_state, "monitor_engine", None) + if engine and engine.rule_count > 0: + rule_events = engine.evaluate(enriched_today) + if rule_events: + # 落盘到 alerts.jsonl + try: + from app.services import alert_store + alert_store.append_many( + self._app_state.repo.store.data_dir, rule_events, + ) + except Exception as e: # noqa: BLE001 + logger.warning("告警落盘失败: %s", e) + # 转为 SSE 推送格式 (兼容旧 alert schema) + for ev in rule_events: + all_alerts.append({ + "source": ev["source"], + "type": ev["type"], + "rule_id": ev.get("rule_id"), + "strategy_id": ev.get("rule_id") if ev["source"] == "strategy" else None, + "symbol": ev["symbol"], + "name": ev["name"], + "message": ev["message"], + "price": ev["price"], + "change_pct": ev["change_pct"], + "signals": ev["signals"], + "severity": ev.get("severity", "info"), + }) - # 2. 刷新策略结果缓存 (实时行情开启时,每轮行情更新后自动重算) + # 刷新策略结果缓存 (实时行情开启时,每轮行情更新后自动重算) if self._enabled and self._app_state: self._refresh_strategy_cache(enriched_today, enriched_date) - # 推入待推送队列 + 通知 SSE + # 推入待推送队列 + 通知 SSE (含背压保护) if all_alerts: with self._lock: self._pending_alerts.extend(all_alerts) + # 背压: 超出上限丢弃最旧 + if len(self._pending_alerts) > self._max_pending_alerts: + overflow = len(self._pending_alerts) - self._max_pending_alerts + self._pending_alerts = self._pending_alerts[overflow:] self._alert_event.set() - logger.info("策略监控评估完成: %d 条通知", len(all_alerts)) + logger.info("监控评估完成: %d 条通知", len(all_alerts)) except Exception as e: # noqa: BLE001 logger.warning("监控评估失败: %s", e) diff --git a/backend/app/strategy/monitor.py b/backend/app/strategy/monitor.py index 9b5b262..10f3b32 100644 --- a/backend/app/strategy/monitor.py +++ b/backend/app/strategy/monitor.py @@ -3,15 +3,23 @@ 职责: 接收实时行情 DataFrame → 检查监控中策略的信号/提醒 → 推送告警。 不知道: 策略加载逻辑、AI、API、配置持久化、回测。 依赖: 外部调用 on_quote_update() 传入实时数据。 + +本模块含两个评估器: + 1. StrategyMonitorService — 旧的策略监控 (type=strategy),第二步迁移到 MonitorRuleEngine + 2. MonitorRuleEngine — 通用规则引擎,覆盖 signal/price/market/strategy 四类, + 支持 scope (symbols/all/sector) + 多条件 AND/OR + cooldown 去重 """ from __future__ import annotations import logging +import time from dataclasses import dataclass, field from typing import Any, Callable import polars as pl +from app.strategy.custom_signals import _OP_BUILDERS # type: ignore # 复用运算符构造器 + logger = logging.getLogger(__name__) @@ -202,3 +210,259 @@ class StrategyMonitorService: row.get("change_pct"), )) return results + + +# ================================================================ +# 通用监控规则引擎 MonitorRuleEngine +# ================================================================ + +_SIGNAL_PREFIXES = ("signal_", "csg_") + + +def _is_signal_field(field: str) -> bool: + return any(field.startswith(p) for p in _SIGNAL_PREFIXES) + + +def _build_condition_mask(df: pl.DataFrame, conditions: list[dict], logic: str) -> pl.DataFrame: + """根据 conditions + logic 构建过滤后的命中 DataFrame。 + + conditions: [{"field","op","value"?}] — op=truth 为布尔信号, 否则阈值比较 + logic: "and" | "or" + 返回命中行 (含 symbol/name/close/change_pct + 各信号列) + """ + cols = set(df.columns) + parts: list[pl.Expr] = [] + for c in conditions: + field = c["field"] + if field not in cols: + return df.head(0) # 字段缺失,无法判定 → 空结果 + op = c["op"] + if op == "truth": + parts.append(pl.col(field).fill_null(False)) + elif op in _OP_BUILDERS: + parts.append(_OP_BUILDERS[op](pl.col(field), c["value"])) + else: + return df.head(0) + if not parts: + return df.head(0) + if logic == "or": + mask = pl.any_horizontal(parts) + else: + mask = pl.all_horizontal(parts) + return df.filter(mask) + + +class MonitorRuleEngine: + """通用监控规则引擎 — 接收实时行情 DataFrame,评估所有规则,返回 AlertEvent。 + + 与 StrategyMonitorService 的区别: + - 规则来自 monitor_rules 存储 (用户可配), 而非写死的 strategy config + - 支持 scope (symbols/all/sector) 过滤作用域 + - 支持 conditions + logic (AND/OR) 任意组合 + - ★ cooldown 去重: 同一 (rule_id, symbol) 在冷却期内不重复触发 + """ + + def __init__(self, alert_handler: Callable[[dict], None] | None = None): + self._alert_handler = alert_handler + self._rules: dict[str, dict] = {} # rule_id → rule + # (rule_id, symbol) → 上次触发时间戳(秒)。用于 cooldown 去重。 + self._last_fire: dict[tuple[str, str], float] = {} + self._strategy_engine = None # 延迟注入, type=strategy 规则用它读策略信号 + + def set_strategy_engine(self, engine) -> None: + """注入 StrategyEngine, type=strategy 规则据此读策略的 entry/exit_signals。""" + self._strategy_engine = engine + + # ── 规则管理 ─────────────────────────────────────── + def set_rules(self, rules: list[dict]) -> None: + """批量设置规则 (覆盖)。用于启动时 reload。""" + self._rules = {} + for r in rules: + if r.get("enabled") is not False: + self._rules[r["id"]] = r + logger.info("MonitorRuleEngine: 装载 %d 条规则", len(self._rules)) + + def add_rule(self, rule: dict) -> None: + if rule.get("enabled") is not False: + self._rules[rule["id"]] = rule + else: + self._rules.pop(rule["id"], None) + + def remove_rule(self, rule_id: str) -> None: + self._rules.pop(rule_id, None) + # 清理对应的 cooldown 记录 + self._last_fire = {k: v for k, v in self._last_fire.items() if k[0] != rule_id} + + def clear(self) -> None: + self._rules.clear() + self._last_fire.clear() + + @property + def rules(self) -> dict[str, dict]: + return dict(self._rules) + + @property + def rule_count(self) -> int: + return len(self._rules) + + # ── 评估 ─────────────────────────────────────────── + def evaluate(self, df: pl.DataFrame) -> list[dict]: + """行情更新后评估所有规则。 + + Args: + df: 实时 enriched 数据 (~5500行, 含 signal_/csg_/指标列) + Returns: + 触发的 AlertEvent dict 列表 (含 ts/rule_id/source/type/symbol/...) + """ + if not self._rules or df.is_empty(): + return [] + + now = time.time() + events: list[dict] = [] + + for rule_id, rule in self._rules.items(): + try: + events.extend(self._evaluate_rule(df, rule, now)) + except Exception as e: + logger.warning("规则评估失败 %s: %s", rule_id, e) + + return events + + def _evaluate_rule(self, df: pl.DataFrame, rule: dict, now: float) -> list[dict]: + """评估单条规则,返回触发的 events。""" + # 1. 按 scope 过滤作用域 + scoped = self._apply_scope(df, rule) + if scoped.is_empty(): + return [] + + # 2. 根据 type 构建命中集 + hit_rows: list[tuple[str, Any, Any, Any, list[str]]] = [] # (symbol,name,price,pct,signals) + + rtype = rule.get("type", "signal") + if rtype == "strategy": + # 策略类型: 从 StrategyEngine 读策略的 entry/exit_signals, 按 direction 评估 + hit_rows = self._match_strategy(scoped, rule) + else: + # signal / price / market: 通用条件匹配 + hit_rows = self._match_conditions(scoped, rule) + + if not hit_rows: + return [] + + # 3. cooldown 去重 + 生成 events + cooldown = rule.get("cooldown_seconds", 3600) + severity = rule.get("severity", "info") + message = rule.get("message", "") or self._default_message(rule) + source = rtype if rtype != "strategy" else "strategy" + ev_type = rule.get("direction", "entry") if rtype == "strategy" else rtype + + events: list[dict] = [] + for sym, name, price, pct, hit_sigs in hit_rows: + key = (rule["id"], sym) + last = self._last_fire.get(key) + if last is not None and (now - last) < cooldown: + continue # 冷却期内, 跳过 + self._last_fire[key] = now + ev = { + "ts": int(now * 1000), + "rule_id": rule["id"], + "rule_name": rule.get("name", ""), + "source": source, + "type": ev_type, + "symbol": sym, + "name": name, + "message": message, + "price": price, + "change_pct": pct, + "signals": hit_sigs, + "severity": severity, + } + events.append(ev) + if self._alert_handler: + try: + self._alert_handler(ev) + except Exception as e: + logger.warning("alert handler failed: %s", e) + + return events + + @staticmethod + def _apply_scope(df: pl.DataFrame, rule: dict) -> pl.DataFrame: + """按 scope 过滤 DataFrame。""" + scope = rule.get("scope", "symbols") + if scope == "all": + return df + if scope == "symbols": + syms = rule.get("symbols", []) + if not syms: + return df.head(0) + return df.filter(pl.col("symbol").is_in(syms)) + if scope == "sector": + # sector 过滤: 需 df 含板块列 (后续接入 ext_data JOIN) + # 当前先返回全量, sector 精确过滤第二步完善 + return df + return df + + def _match_strategy( + self, df: pl.DataFrame, rule: dict, + ) -> list[tuple[str, Any, Any, Any, list[str]]]: + """策略类型评估: 从 StrategyEngine 读策略信号, 按 direction 用 OR 匹配。 + + direction=entry → 策略 entry_signals + direction=exit → 策略 exit_signals + direction=both → entry + exit 合并 (命中信号名区分来源) + """ + if self._strategy_engine is None: + return [] + sid = rule.get("strategy_id") + if not sid: + return [] + try: + s = self._strategy_engine.get(sid) + except Exception: + return [] + if s is None: + return [] + + direction = rule.get("direction", "entry") + # 收集要评估的信号 (OR 组合), 与旧 StrategyMonitorService 行为一致 + sigs: list[str] = [] + if direction in ("entry", "both"): + sigs.extend(s.entry_signals or []) + if direction in ("exit", "both"): + sigs.extend(s.exit_signals or []) + if not sigs: + return [] + + # 复用旧的 _check_signals 静态方法 (已支持 signal_/csg_ 前缀) + return StrategyMonitorService._check_signals(df, sigs) + + @staticmethod + def _match_conditions( + df: pl.DataFrame, rule: dict, + ) -> list[tuple[str, Any, Any, Any, list[str]]]: + """按 conditions + logic 匹配,返回命中行 [(symbol,name,price,pct,signals)]。""" + conditions = rule.get("conditions", []) + logic = rule.get("logic", "and") + if not conditions: + return [] + hit_df = _build_condition_mask(df, conditions, logic) + results = [] + for row in hit_df.iter_rows(named=True): + sym = row.get("symbol", "") + name = row.get("name") + price = row.get("close") + pct = row.get("change_pct") + # 收集命中的信号列名 (仅 op=truth 且为真的) + hit_sigs = [ + c["field"] for c in conditions + if c.get("op") == "truth" and row.get(c["field"]) + ] + results.append((sym, name, price, pct, hit_sigs)) + return results + + @staticmethod + def _default_message(rule: dict) -> str: + rtype = rule.get("type", "signal") + name_map = {"signal": "信号触发", "price": "价格触发", "market": "市场异动", "strategy": "策略触发"} + return name_map.get(rtype, "监控触发") diff --git a/backend/app/strategy/monitor_rules.py b/backend/app/strategy/monitor_rules.py new file mode 100644 index 0000000..08392fb --- /dev/null +++ b/backend/app/strategy/monitor_rules.py @@ -0,0 +1,241 @@ +"""监控规则 — 统一的 MonitorRule 模型,覆盖策略/个股信号/个股价格/市场异动四类。 + +职责: + - 从 data/user_data/monitor_rules/*.json 加载规则定义 + - 校验规则字段合法性 + - 提供 CRUD (load_all / save_one / delete_one) + +不知道: 行情评估引擎、API、告警落盘。纯函数 + 文件存储。 + +设计 (镜像 custom_signals.py 的写法): + - 一对象一文件 + glob 全扫 + 全量重写 + - 字段白名单复用 custom_signals.ALLOWED_FIELDS (阈值条件) + 信号列清单 (布尔条件) + - id 正则与 custom_signals 一致,保证可纳入同一索引体系 +""" +from __future__ import annotations + +import json +import logging +import re +from datetime import datetime, timezone +from pathlib import Path + +from app.strategy.custom_signals import ALLOWED_FIELDS + +logger = logging.getLogger(__name__) + +# ── 常量 ──────────────────────────────────────────────── +ID_RE = re.compile(r"^[a-z0-9_]{1,40}$") +RULE_TYPES = {"strategy", "signal", "price", "market"} +SCOPES = {"symbols", "all", "sector"} +LOGICS = {"and", "or"} +DIRECTIONS = {"entry", "exit", "both"} +SEVERITIES = {"info", "warn", "critical"} +OPS = {">", ">=", "<", "<=", "==", "!="} + +# 布尔信号列前缀 (op=truth 时 field 取这些) +_SIGNAL_PREFIXES = ("signal_", "csg_") + + +# ── 持久化 (镜像 custom_signals.py) ───────────────────── +def _dir(data_dir: Path) -> Path: + d = data_dir / "user_data" / "monitor_rules" + d.mkdir(parents=True, exist_ok=True) + return d + + +def _path(data_dir: Path, rule_id: str) -> Path: + return _dir(data_dir) / f"{rule_id}.json" + + +def load_all(data_dir: Path) -> list[dict]: + """读取全部监控规则。损坏的文件被跳过。""" + d = _dir(data_dir) + out: list[dict] = [] + for f in sorted(d.glob("*.json")): + try: + out.append(json.loads(f.read_text(encoding="utf-8"))) + except Exception as e: + logger.warning("monitor rule load failed %s: %s", f.name, e) + return out + + +def load_one(data_dir: Path, rule_id: str) -> dict | None: + p = _path(data_dir, rule_id) + if not p.exists(): + return None + try: + return json.loads(p.read_text(encoding="utf-8")) + except Exception as e: + logger.warning("monitor rule load failed %s: %s", rule_id, e) + return None + + +def save_one(data_dir: Path, rule: dict) -> None: + p = _path(data_dir, rule["id"]) + p.parent.mkdir(parents=True, exist_ok=True) + p.write_text(json.dumps(rule, ensure_ascii=False, indent=2), encoding="utf-8") + + +def delete_one(data_dir: Path, rule_id: str) -> bool: + p = _path(data_dir, rule_id) + if p.exists(): + p.unlink() + return True + return False + + +# ── 校验 ──────────────────────────────────────────────── +def _is_signal_field(field: str) -> bool: + """判断 field 是否为布尔信号列 (signal_ / csg_ 前缀)。""" + return any(field.startswith(p) for p in _SIGNAL_PREFIXES) + + +def validate(rule: dict) -> None: + """校验一条监控规则,非法则抛 ValueError (含中文信息)。""" + rid = rule.get("id", "") + if not isinstance(rid, str) or not ID_RE.match(rid): + raise ValueError(f"规则 id 非法 (仅小写字母数字下划线, 1-40字符): {rid!r}") + if not isinstance(rule.get("name"), str) or not rule["name"].strip(): + raise ValueError("规则 name 不能为空") + if rule.get("type") not in RULE_TYPES: + raise ValueError(f"type 必须是 {RULE_TYPES} 之一") + + # 策略类型: 需要 strategy_id + direction,conditions 可空 + if rule.get("type") == "strategy": + if not rule.get("strategy_id"): + raise ValueError("策略类型规则必须指定 strategy_id") + if rule.get("direction", "entry") not in DIRECTIONS: + raise ValueError(f"direction 必须是 {DIRECTIONS} 之一") + else: + # 信号/价格/市场类型: 需要 conditions + conds = rule.get("conditions") + if not isinstance(conds, list) or len(conds) == 0: + raise ValueError("conditions 不能为空") + if len(conds) > 8: + raise ValueError("conditions 最多 8 条") + if rule.get("logic", "and") not in LOGICS: + raise ValueError(f"logic 必须是 {LOGICS} 之一") + for i, c in enumerate(conds): + if not isinstance(c, dict): + raise ValueError(f"第 {i+1} 个条件格式错误") + field = c.get("field", "") + op = c.get("op", "") + if op == "truth": + # 布尔信号: field 必须是 signal_/csg_ 前缀 + if not _is_signal_field(field): + raise ValueError(f"第 {i+1} 个条件: op=truth 时 field 必须是信号列 (signal_/csg_ 前缀): {field!r}") + elif op in OPS: + # 阈值比较: field 必须在白名单, 需要 value + if field not in ALLOWED_FIELDS: + raise ValueError(f"第 {i+1} 个条件: 阈值字段 {field!r} 不在白名单") + if not isinstance(c.get("value"), (int, float)): + raise ValueError(f"第 {i+1} 个条件: value 必须是数字") + else: + raise ValueError(f"第 {i+1} 个条件: op {op!r} 非法 (应为 truth 或 {OPS})") + + # scope 校验 + if rule.get("scope", "symbols") not in SCOPES: + raise ValueError(f"scope 必须是 {SCOPES} 之一") + if rule.get("scope") == "symbols": + syms = rule.get("symbols") + if not isinstance(syms, list) or len(syms) == 0: + raise ValueError("scope=symbols 时 symbols 不能为空") + + # 其余枚举 + if rule.get("severity", "info") not in SEVERITIES: + raise ValueError(f"severity 必须是 {SEVERITIES} 之一") + cd = rule.get("cooldown_seconds", 3600) + if not isinstance(cd, int) or cd < 0: + raise ValueError("cooldown_seconds 必须是非负整数") + + +def normalize(rule: dict) -> dict: + """补全默认字段,返回规范化后的规则 (不校验)。""" + r = dict(rule) + r.setdefault("enabled", True) + r.setdefault("scope", "symbols") + r.setdefault("symbols", []) + r.setdefault("sector", None) + r.setdefault("strategy_id", None) + r.setdefault("direction", "entry") + r.setdefault("conditions", []) + r.setdefault("logic", "and") + r.setdefault("cooldown_seconds", 3600) + r.setdefault("severity", "info") + r.setdefault("message", "") + r.setdefault("webhook_url", "") + r.setdefault("webhook_enabled", False) + r.setdefault("created_at", datetime.now(timezone.utc).isoformat()) + return r + + +# 策略监控自动迁移的规则 id 前缀 (固定, 保证幂等) +STRATEGY_RULE_PREFIX = "mr_strategy_" + + +def strategy_rule_id(strategy_id: str) -> str: + """策略监控规则 id = mr_strategy_{strategy_id}。""" + return f"{STRATEGY_RULE_PREFIX}{strategy_id}" + + +def migrate_strategy_monitors(data_dir: Path, strategy_ids: list[str], strategy_names: dict[str, str]) -> list[dict]: + """把 preferences.strategy_monitor_ids 里的策略,同步生成/更新 type=strategy 规则。 + + 幂等: 已存在的策略规则会被更新 (方向/名称),不会重复创建。 + 已从 strategy_ids 移除的策略, 其规则会被停用 (enabled=False) 而非删除 (保留历史触发记录的关联)。 + + Args: + data_dir: 数据目录 + strategy_ids: 当前监控池中的策略 id 列表 + strategy_names: {strategy_id: 策略名} 用于规则显示名 + Returns: + 本次生成/更新的规则列表 + """ + desired = set(strategy_ids) + existing = load_all(data_dir) + # 已存在的策略规则 {strategy_id: rule} + existing_strategy_rules: dict[str, dict] = {} + for r in existing: + rid = r.get("id", "") + if rid.startswith(STRATEGY_RULE_PREFIX): + sid = rid[len(STRATEGY_RULE_PREFIX):] + if sid: + existing_strategy_rules[sid] = r + + touched: list[dict] = [] + # 1. 为当前监控池的策略 upsert 规则 + for sid in desired: + rule_id = strategy_rule_id(sid) + name = strategy_names.get(sid, sid) + rule = existing_strategy_rules.get(sid) + if rule is None: + rule = normalize({ + "id": rule_id, + "name": f"策略监控 · {name}", + "type": "strategy", + "scope": "all", + "strategy_id": sid, + "direction": "entry", + "conditions": [], + "cooldown_seconds": 3600, + "enabled": True, + }) + else: + rule = dict(rule) + rule["enabled"] = True + rule["strategy_id"] = sid + rule["name"] = f"策略监控 · {name}" + rule.setdefault("scope", "all") + rule.setdefault("direction", "entry") + save_one(data_dir, rule) + touched.append(rule) + + # 2. 不在监控池的策略 → 停用其规则 (不删除) + for sid, rule in existing_strategy_rules.items(): + if sid not in desired and rule.get("enabled") is not False: + rule = dict(rule) + rule["enabled"] = False + save_one(data_dir, rule) + + return touched diff --git a/dev.ps1 b/dev.ps1 index d9ca3be..c9bc16a 100644 --- a/dev.ps1 +++ b/dev.ps1 @@ -148,14 +148,14 @@ $backendJob = Start-Job -Name 'backend' -ScriptBlock { $PID | Out-File -FilePath $pidFile -Encoding ascii -Force $env:PYTHONUNBUFFERED = '1' Set-Location $dir - & .\.venv\Scripts\python.exe -m uvicorn app.main:app --reload --port $port 2>&1 + & .\.venv\Scripts\python.exe -m uvicorn app.main:app --reload --host 0.0.0.0 --port $port 2>&1 } -ArgumentList $backendPidFile, $BackendDir, $BackendPort $frontendJob = Start-Job -Name 'frontend' -ScriptBlock { param($pidFile, $dir, $port) $PID | Out-File -FilePath $pidFile -Encoding ascii -Force Set-Location $dir - & pnpm dev --port $port 2>&1 + & pnpm dev --host 0.0.0.0 --port $port 2>&1 } -ArgumentList $frontendPidFile, $FrontendDir, $FrontendPort # Wait up to 5 seconds for the PID files to materialise diff --git a/dev.sh b/dev.sh index 987a86e..4304b36 100755 --- a/dev.sh +++ b/dev.sh @@ -120,14 +120,14 @@ echo ( cd "$BACKEND_DIR" - uv run uvicorn app.main:app --reload --port "$BACKEND_PORT" 2>&1 \ + uv run uvicorn app.main:app --reload --host 0.0.0.0 --port "$BACKEND_PORT" 2>&1 \ | prefix_awk "$(printf "${BLUE}[backend ]${NC} ")" ) & PIDS+=("$!") ( cd "$FRONTEND_DIR" - pnpm dev --port "$FRONTEND_PORT" 2>&1 \ + pnpm dev --host 0.0.0.0 --port "$FRONTEND_PORT" 2>&1 \ | prefix_awk "$(printf "${GREEN}[frontend]${NC} ")" ) & PIDS+=("$!") diff --git a/docs/screenshots/concept-analysis.png b/docs/screenshots/concept-analysis.png new file mode 100644 index 0000000..13255c7 Binary files /dev/null and b/docs/screenshots/concept-analysis.png differ diff --git a/docs/screenshots/dashboard.png b/docs/screenshots/dashboard.png index 7d4af68..bba7a1e 100644 Binary files a/docs/screenshots/dashboard.png and b/docs/screenshots/dashboard.png differ diff --git a/docs/screenshots/monitor.png b/docs/screenshots/monitor.png new file mode 100644 index 0000000..cf4a9bb Binary files /dev/null and b/docs/screenshots/monitor.png differ diff --git a/docs/screenshots/screener.png b/docs/screenshots/screener.png index c4f6c09..f1d1804 100644 Binary files a/docs/screenshots/screener.png and b/docs/screenshots/screener.png differ diff --git a/frontend/src/components/AlertToast.tsx b/frontend/src/components/AlertToast.tsx new file mode 100644 index 0000000..b8ec877 --- /dev/null +++ b/frontend/src/components/AlertToast.tsx @@ -0,0 +1,142 @@ +import { useCallback, useEffect, useState } from 'react' +import { useNavigate } from 'react-router-dom' +import { motion, AnimatePresence } from 'framer-motion' +import { Bell, TrendingUp, TrendingDown, X } from 'lucide-react' +import type { AlertEvent } from '@/lib/api' +import { fmtPct, fmtPrice } from '@/lib/format' +import { cn } from '@/lib/cn' +import { playNotificationSound } from '@/lib/notificationSound' + +// ===== 全局状态 (模块级, 仿 Toast.tsx 模式) ===== +type Item = { id: number; alert: AlertEvent } +let _id = 0 +let _queue: Item[] = [] +const AUTO_DISMISS = 5000 // 5 秒自动消失 +const _listeners: Set<(items: Item[]) => void> = new Set() + +/** 从 localStorage 读取配置 */ +function getEnabled(): boolean { + try { + const v = localStorage.getItem('alert_toast_enabled') + return v === null ? true : v === '1' // 默认开启 + } catch { return true } +} + +function getMaxVisible(): number { + try { + const v = parseInt(localStorage.getItem('alert_toast_max') || '', 10) + return v >= 1 && v <= 10 ? v : 3 // 默认 3, 范围 1-10 + } catch { return 3 } +} + +/** 通知外部配置变更后刷新 (设置页改了配置后调用) */ +export function refreshAlertToastConfig() { + _emit() +} + +function _emit() { _listeners.forEach(fn => fn([..._queue])) } + +/** 推入监控告警通知 (外部调用) */ +export function pushAlertToast(alert: AlertEvent) { + if (!getEnabled()) return // 开关关闭: 不弹 + const maxVisible = getMaxVisible() + const item = { id: ++_id, alert } + _queue = [..._queue, item] + // 超出上限: 丢弃最旧的 + if (_queue.length > maxVisible) { + _queue = _queue.slice(-maxVisible) + } + _emit() + setTimeout(() => dismiss(item.id), AUTO_DISMISS) + playNotificationSound() // 播放声效 +} + +/** 手动关闭 */ +export function dismiss(id: number) { + _queue = _queue.filter(t => t.id !== id) + _emit() +} + +// ===== 配色 ===== +const SEVERITY_BAR: Record
| 日期 | +分钟K条数 | +数据来源 | +状态 | +
|---|---|---|---|
| {r.date} | +{r.rows} | ++ {r.source} + | +
+ {r.ok ? (
+
+ |
+
+ 缺失日期若为周末/节假日属正常; + 若为停牌日(成交量为 0)也属正常; + 若为正常交易日(日K有成交量)却缺失分钟K, + 则是 TickFlow 数据源未提供该日分钟数据。 +
++ 生成模拟的触发记录,用于测试监控中心页面的展示效果、未读徽标、新增闪烁等功能。生成的数据可随时清空。 +
++ 生成多种类型的演示监控规则 (个股信号/价格/市场异动),用于测试监控中心规则列表展示。 +
+
- 逐日调用 /api/kline/minute 接口,
- 检测每只股票最近若干天的分钟K数据是否齐全。本地无数据时会自动走 TickFlow 实时拉取。
-
| 日期 | -分钟K条数 | -数据来源 | -状态 | -
|---|---|---|---|
| {r.date} | -{r.rows} | -- - {r.source} - - | -
- {r.ok ? (
-
- |
-
- 缺失日期若为周末/节假日属正常; - 若为停牌日(成交量为 0)也属正常; - 若为正常交易日(日K有成交量)却缺失分钟K, - 则是 TickFlow 数据源未提供该日分钟数据(常见于停牌后复牌首日,存在补数延迟)。 -
-{item.title}
-{item.desc}
-{message}
+- 每次行情刷新时自动评估监控池中的策略。命中买入/卖出信号或阈值条件时弹通知。 - 与当前打开的页面无关 — 后端始终在评估。 +
+ 策略监控、个股信号监控、价格监控已统一到「监控中心」页面,支持灵活配置触发条件、冷却期和作用范围。
-