diff --git a/ChanMacro/cli.py b/ChanMacro/cli.py index bcfed84..39cb9f2 100644 --- a/ChanMacro/cli.py +++ b/ChanMacro/cli.py @@ -5,6 +5,7 @@ cli.py — Command-line interface for ChanMacro. import argparse import json import logging +import time from datetime import date as Date, datetime, timedelta logging.basicConfig( @@ -331,6 +332,8 @@ def main(): # validate p_validate = sub.add_parser("validate", help="Run validation framework") + # cron + p_cron = sub.add_parser("cron", help="Run scheduled fetch+score loop") # serve p_serve = sub.add_parser("serve", help="Start web dashboard") @@ -353,7 +356,22 @@ def main(): report = ValidationReporter().run_all() print(report) elif args.command == "serve": - logger.info("Web dashboard not yet implemented (Phase 7)") + from scheduler import get_scheduler + get_scheduler().start() + logger.info("启动 Web Dashboard: http://127.0.0.1:8124") + from web.app import app + app.run(host="0.0.0.0", port=8124, debug=False) + elif args.command == "cron": + from scheduler import get_scheduler + logger.info("启动后台调度器 (Ctrl+C 停止)") + s = get_scheduler() + s.start() + try: + while True: + time.sleep(60) + except KeyboardInterrupt: + s.stop() + logger.info("调度器已停止") else: parser.print_help() diff --git a/ChanMacro/scheduler.py b/ChanMacro/scheduler.py new file mode 100644 index 0000000..c6b3646 --- /dev/null +++ b/ChanMacro/scheduler.py @@ -0,0 +1,137 @@ +""" +scheduler.py — 后台自动调度:定时拉取数据 + 计算因子 + 制度判定。 + +Python main.py cron → 前台阻塞运行,每 N 分钟一个 tick +Web app 启动时自动启动调度器 → 后台线程,不阻塞 Web 请求 +""" + +import threading +import logging +import time +from datetime import datetime, timezone, timedelta +from typing import Optional + +logger = logging.getLogger("chanmacro.scheduler") + + +class MacroScheduler: + """后台调度器:定时 fetch + score。""" + + def __init__(self, interval_minutes: int = 60): + self.interval = interval_minutes + self._thread: Optional[threading.Thread] = None + self._stop = threading.Event() + self._last_run: Optional[datetime] = None + self._running = False + + def start(self) -> None: + """启动后台线程。""" + if self._running: + return + self._stop.clear() + self._thread = threading.Thread(target=self._loop, name="macro-scheduler", daemon=True) + self._thread.start() + self._running = True + logger.info(f"调度器已启动, 每 {self.interval} 分钟执行一次") + + def stop(self) -> None: + """停止后台线程。""" + self._stop.set() + self._running = False + logger.info("调度器已停止") + + @property + def last_run(self) -> Optional[datetime]: + return self._last_run + + def _loop(self) -> None: + """后台循环。""" + # 首次启动立即跑一次 + self._tick() + + while not self._stop.wait(self.interval * 60): + self._tick() + + def _tick(self) -> None: + """执行一次:fetch → score。""" + try: + from fetchers.ohlcv import OHLCVFetcher + from fetchers.breadth import BreadthFetcher + from fetchers.derivatives import DerivativesFetcher + from database import init_db + from datetime import date as Date + + init_db() + today = Date.today() + + # Fetch + ohlcv = OHLCVFetcher() + df = ohlcv.fetch() + if not df.empty: + ohlcv.store_df(df) + + breadth = BreadthFetcher() + record = breadth.fetch() + if record: + breadth.store(record=record) + + deriv = DerivativesFetcher() + records = deriv.fetch(today) + if records: + deriv.store(records=records) + + # Score + Regime (also persisted inside _build_state) + from scoring.price_structure import PriceStructureScorer + from scoring.breadth_scorer import BreadthScorer + from scoring.oi_matrix import OIMatrixScorer + from scoring.volatility_regime import VolatilityRegimeScorer + from regime_detector import RegimeDetector + from models import MarketStateVector + from config import config + import json + from database import get_connection + + ps = PriceStructureScorer().compute(today) + br = BreadthScorer().compute(today) + oi = OIMatrixScorer().compute(today) + vol = VolatilityRegimeScorer().compute(today) + + detector = RegimeDetector() + detector.load_state(config.db_path) + r = detector.detect(ps.score, br.breadth_top50, vol.vol_regime.value, today) + + conn = get_connection() + conn.execute(""" + INSERT OR REPLACE INTO regime_history + (date, regime, confidence, regime_version, maturity_score, + all_scores_json, prior_regime, confirmation_days) + VALUES (?, ?, ?, ?, ?, ?, ?, ?) + """, ( + str(today), r.regime.value, r.confidence, r.regime_version, + r.maturity_score, json.dumps(r.all_scores), + r.prior_regime.value if r.prior_regime else None, + r.confirmation_days, + )) + conn.commit() + conn.close() + + self._last_run = datetime.now(timezone.utc) + logger.info( + f"Tick 完成: regime={r.regime.value} conf={r.confidence:.2f} " + f"breadth={br.score:.0f}({br.breadth_bucket.value}) " + f"price={ps.score:.0f} oi={oi.oi_state.value} vol={vol.vol_regime.value}" + ) + + except Exception as e: + logger.error(f"Tick 失败: {e}", exc_info=True) + + +# 单例 +_scheduler: Optional[MacroScheduler] = None + + +def get_scheduler() -> MacroScheduler: + global _scheduler + if _scheduler is None: + _scheduler = MacroScheduler(interval_minutes=60) + return _scheduler diff --git a/ChanMacro/web/app.py b/ChanMacro/web/app.py index feeb5e3..0fc152f 100644 --- a/ChanMacro/web/app.py +++ b/ChanMacro/web/app.py @@ -173,4 +173,6 @@ def api_expectancy(): if __name__ == "__main__": + from scheduler import get_scheduler + get_scheduler().start() app.run(host="0.0.0.0", port=8124, debug=True)