diff --git a/deploy/profile-scripts/strategy_executor.py b/deploy/profile-scripts/strategy_executor.py new file mode 100644 index 00000000..800c1eb3 --- /dev/null +++ b/deploy/profile-scripts/strategy_executor.py @@ -0,0 +1,149 @@ +#!/usr/bin/env python3 +# -*- coding: utf-8 -*- +"""strategy_executor.py — 策略执行调度器(2026-08-13 温区自适应核心) + +核心(老莫要求):根据当前温区激活策略,调度对应的选股/买卖/重评任务。 + +温区自适应数据流: + regime_tracker(K=5 平滑温区) + temp_band(温度) + regime_perf(实测) + → strategy_router(当前温区激活策略 strategy_weights.json) + → strategy_executor(按激活策略调度对应 scanner/监控/重评) + +调度逻辑: + - 读 strategy_weights.json 的 active(当前温区激活策略集合) + - 对每个激活策略,调度其对应的选股扫描器(scanner) + - 非激活策略的 scanner 不跑(温区休眠) + +策略→scanner 映射: + v_weak → mr_scanner.py (choppy/trend_down 超跌确认) + v_oversold → predictive_oversold_scanner.py(trend_down 预测超跌反弹) + s2_panic → s2_scanner.py (trend_down 恐慌买超跌) + v_mr/v_mr_sel → mr_scanner.py (复用 v_weak 扫描器,同家族) + v_next4/v8.1 → leader_scanner.py (trend_up 龙头识别) + +用法:cron 每30分钟盘中跑(替代各 scanner 分散 cron) + python3 strategy_executor.py +""" +import sys +import json +import subprocess +from pathlib import Path +from datetime import datetime + +_SCRIPT_DIR = Path(__file__).resolve().parent +sys.path.insert(0, str(_SCRIPT_DIR)) +sys.path.insert(0, "/home/hmo/MoFin") + +WEIGHTS_FILE = "/home/hmo/MoFin/data/strategy_weights.json" +LOG = "/home/hmo/MoFin/gateway/logs/strategy_executor.log" + +# 策略 → 选股扫描器映射 +STRATEGY_SCANNER = { + "v_weak": "mr_scanner.py", + "v_oversold": "predictive_oversold_scanner.py", + "p_oversold": "predictive_oversold_scanner.py", + "s2_panic": "s2_scanner.py", + "v_mr": "mr_scanner.py", + "v_mr_sel": "mr_scanner.py", + "v_mr2": "mr_scanner.py", + "v_mr3": "mr_scanner.py", + "v_mr4": "mr_scanner.py", + "v_lurk_v1": "mr_scanner.py", + "v_lurk_v2": "mr_scanner.py", + "v_lurk_v3": "mr_scanner.py", + "v_next4": "leader_scanner.py", + "v8.1": "leader_scanner.py", + "v_next5": "leader_scanner.py", + "v_combo": "leader_scanner.py", +} + + +def log(msg): + line = f"[{datetime.now().isoformat(timespec='seconds')}] {msg}" + print(line, flush=True) + try: + Path(LOG).parent.mkdir(parents=True, exist_ok=True) + with open(LOG, "a", encoding="utf-8") as f: + f.write(line + "\n") + except Exception: + pass + + +def load_active(): + """读当前温区激活策略""" + try: + p = Path(WEIGHTS_FILE) + if p.exists(): + d = json.loads(p.read_text(encoding="utf-8")) + return d.get("active") or [], d.get("state", "unknown") + except Exception as e: + log(f"加载 weights 失败: {e}") + return [], "unknown" + + +def is_trading_time(): + """盘中时段检查(9:30-11:30 / 13:00-15:00 连续竞价)""" + now = datetime.now() + t = now.strftime("%H:%M") + wd = now.weekday() + if wd >= 5: + return False + return ("09:30" <= t <= "11:30") or ("13:00" <= t <= "15:00") + + +def run_scanner(scanner_name): + """执行一个扫描器(带超时防挂)""" + script = _SCRIPT_DIR / scanner_name + if not script.exists(): + log(f" ⚠️ {scanner_name} 不存在,跳过") + return False + log(f" ▶ 执行 {scanner_name}") + try: + r = subprocess.run( + [sys.executable, str(script)], + capture_output=True, text=True, timeout=600, + cwd=str(_SCRIPT_DIR), + ) + out = (r.stdout or "").strip()[:300] + err = (r.stderr or "").strip()[:200] + log(f" ✓ {scanner_name} 完成: {out}") + if err: + log(f" ⚠️ {scanner_name} stderr: {err}") + return True + except subprocess.TimeoutExpired: + log(f" ⏱ {scanner_name} 超时(600s)跳过") + return False + except Exception as e: + log(f" ✗ {scanner_name} 异常: {e}") + return False + + +def main(): + if not is_trading_time(): + print(f"[SKIP] {datetime.now().strftime('%H:%M')} 非连续竞价时段", flush=True) + return 0 + + active, regime = load_active() + log(f"温区: {regime} | 激活策略: {active}") + + # 去重:同一 scanner 只跑一次(多策略映射同一扫描器) + scanners = set() + for strat in active: + sc = STRATEGY_SCANNER.get(strat) + if sc: + scanners.add(sc) + log(f"调度扫描器: {sorted(scanners)}") + + if not scanners: + log(" 无激活策略的扫描器(当前温区无对应选股任务)") + return 0 + + for sc in sorted(scanners): + run_scanner(sc) + + log("调度完成") + return 0 + + +if __name__ == "__main__": + sys.exit(main())