""" 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( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) logger = logging.getLogger("chanmacro") def parse_date(date_str: str) -> Date: """Parse YYYY-MM-DD string to Date.""" return datetime.strptime(date_str, "%Y-%m-%d").date() def _build_market_state(target: Date) -> tuple: """Shared helper: compute all scores → (MarketStateVector, RegimeResult).""" from config import config 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 ps = PriceStructureScorer().compute(target) br = BreadthScorer().compute(target) oi = OIMatrixScorer().compute(target) vol = VolatilityRegimeScorer().compute(target) detector = RegimeDetector() detector.load_state(config.db_path) r = detector.detect(ps.score, br.breadth_top50, vol.vol_regime.value, target) state = MarketStateVector( date=target, regime=r.regime, regime_confidence=r.confidence, regime_version=r.regime_version, regime_maturity_score=r.maturity_score, breadth_top20=br.breadth_top20, breadth_top30=br.breadth_top30, breadth_top50=br.breadth_top50, breadth_bucket=br.breadth_bucket, breadth_divergence=br.breadth_divergence, oi_state=oi.oi_state, volatility_regime=vol.vol_regime, price_structure_score=ps, breadth_score=br, oi_matrix_score=oi, volatility_regime_score=vol, ) state.market_state_hash = state.compute_hash() # Persist regime to DB so subsequent calls have correct state from database import get_connection 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(target), 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() return state, r def cmd_fetch(args): """Fetch raw data and store to DB.""" from database import init_db from fetchers.ohlcv import OHLCVFetcher from fetchers.breadth import BreadthFetcher target = parse_date(args.date) if args.date else Date.today() init_db() module = args.module or "all" if module in ("ohlcv", "all"): logger.info(f"Fetching OHLCV for {target}...") fetcher = OHLCVFetcher() df = fetcher.fetch(target) if not df.empty: n = fetcher.store_df(df) logger.info(f"OHLCV: stored {n} rows") if module in ("breadth", "all"): logger.info(f"Fetching Breadth for {target}...") fetcher = BreadthFetcher() record = fetcher.fetch(target) if record: fetcher.store(record=record) logger.info(f"Breadth: stored (adv={record.get('advance_top50')}, " f"dec={record.get('decline_top50')}, " f"ema20={record.get('above_ema20_top50')})") if module in ("derivatives", "all"): logger.info(f"Fetching Derivatives for {target}...") from fetchers.derivatives import DerivativesFetcher fetcher = DerivativesFetcher() records = fetcher.fetch(target) if records: n = fetcher.store(records=records) logger.info(f"Derivatives: stored {n} records") def cmd_score(args): """Compute all factor scores and regime for a date.""" from database import init_db target = parse_date(args.date) if args.date else Date.today() init_db() logger.info(f"Computing scores for {target}...") state, _ = _build_market_state(target) # Output ps = state.price_structure_score br = state.breadth_score oi = state.oi_matrix_score vol = state.volatility_regime_score print(f"\n{'='*60}") print(f" {target} Market State") print(f"{'='*60}") print(f" Regime: {state.regime.value} (conf={state.regime_confidence:.2f}, " f"v={state.regime_version})") print(f" Maturity: {state.regime_maturity_score:.0f}/100") print(f" Breadth: {state.breadth_bucket.value} " f"(T20={state.breadth_top20:.0f} T30={state.breadth_top30:.0f} " f"T50={state.breadth_top50:.0f} div={state.breadth_divergence:+.0f})") print(f" OI State: {state.oi_state.value}") print(f" Volatility: {state.volatility_regime.value}") print(f"{'='*60}") print(f" Scores:") print(f" Price Structure: {ps.score:.0f} {ps.label}") print(f" Breadth: {br.score:.0f} {br.breadth_bucket.value}") print(f" OI Matrix: {oi.score:.0f} {oi.oi_state.value}") print(f" Volatility: {vol.score:.0f} {vol.vol_regime.value}") print(f"{'='*60}") print(f" Market State Hash: {state.market_state_hash}") print() return state def cmd_regime(args): """Show regime history.""" from database import get_connection days = args.days or 30 conn = get_connection() rows = conn.execute( "SELECT date, regime, confidence, maturity_score, confirmation_days " "FROM regime_history ORDER BY date DESC LIMIT ?", (days,) ).fetchall() conn.close() print(f"\n{'='*50}") print(f" Regime History (last {days} days)") print(f"{'='*50}") for r in rows: print(f" {r['date']} {r['regime']:7s} conf={r['confidence']:.2f} " f"mat={r['maturity_score']:.0f} days={r['confirmation_days']}") print() def cmd_track(args): """Record a trading signal with current market state.""" from database import init_db from expectancy.tracker import SignalTracker target = parse_date(args.date) if args.date else Date.today() init_db() logger.info(f"Recording {args.signal} on {target} @ {args.price}") state, _ = _build_market_state(target) tracker = SignalTracker() rid = tracker.record( date=target, signal_type=args.signal, entry_price=args.price, state=state, signal_grade=args.grade, signal_strength=args.strength, ) logger.info(f"Signal recorded: id={rid}") def cmd_backfill(args): """Backfill historical breadth + regime scores.""" from datetime import date as Date, timedelta from database import init_db, get_connection from fetchers.ohlcv import OHLCVFetcher from fetchers.breadth import BreadthFetcher from config import config import pandas as pd import requests start = parse_date(args.from_date) end = parse_date(args.to_date) if args.to_date else Date.today() init_db() # Step 1: Ensure OHLCV data exists for the range logger.info(f"Step 1/3: Fetching BTC OHLCV...") OHLCVFetcher().store_df(OHLCVFetcher().fetch()) # Step 2: Backfill breadth — fetch TOP50 daily data and compute per date logger.info(f"Step 2/3: Backfilling breadth {start} → {end}...") provider_url = config.provider_url all_symbol_data = {} for sym in config.top50_symbols: try: df = pd.DataFrame(requests.get( f"{provider_url}/api/candles", params={"symbol": sym, "tf": "1d", "limit": 400}, timeout=30 ).json()) if not df.empty and "timestamp" in df.columns: df["date"] = pd.to_datetime(df["timestamp"], unit="ms").dt.date df["close"] = df["close"].astype(float) df["high"] = df["high"].astype(float) df["ema20"] = df["close"].ewm(20).mean() all_symbol_data[sym] = df except Exception as e: logger.debug(f" Skip {sym}: {e}") logger.info(f" Fetched {len(all_symbol_data)}/{len(config.top50_symbols)} symbols") # Compute breadth for each date conn = get_connection() current = start breadth_count = 0 while current <= end: target_str = str(current) try: advances_50 = declines_50 = above_ema20_50 = new_highs_50 = 0 advances_30 = advances_20 = above_ema20_30 = above_ema20_20 = 0 new_highs_30 = new_highs_20 = 0 for rank, (sym, df) in enumerate(all_symbol_data.items()): rows = df[df["date"] == current] if rows.empty: continue row = rows.iloc[0] prev_rows = df[df["date"] < current] if prev_rows.empty: continue prev = prev_rows.iloc[-1] if row["close"] > prev["close"]: if rank < 50: advances_50 += 1 if rank < 30: advances_30 += 1 if rank < 20: advances_20 += 1 elif row["close"] < prev["close"]: if rank < 50: declines_50 += 1 if not pd.isna(row.get("ema20")) and row["close"] > row["ema20"]: if rank < 50: above_ema20_50 += 1 if rank < 30: above_ema20_30 += 1 if rank < 20: above_ema20_20 += 1 recent_highs = df[(df["date"] < current) & (df["date"] >= current - timedelta(days=20))] if not recent_highs.empty and row["high"] > recent_highs["high"].max(): if rank < 50: new_highs_50 += 1 if rank < 30: new_highs_30 += 1 if rank < 20: new_highs_20 += 1 conn.execute("""INSERT OR REPLACE INTO breadth_daily (date, total_tracked, advance_top50, decline_top50, above_ema20_top50, new_highs_20d_top50, advance_top30, advance_top20, above_ema20_top30, above_ema20_top20, new_highs_20d_top30, new_highs_20d_top20) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)""", (target_str, len(all_symbol_data), advances_50, declines_50, above_ema20_50, new_highs_50, advances_30, advances_20, above_ema20_30, above_ema20_20, new_highs_30, new_highs_20)) breadth_count += 1 except Exception as e: logger.debug(f" Breadth skip {current}: {e}") current += timedelta(days=1) conn.commit() logger.info(f" Breadth backfill: {breadth_count} days") # Step 3: Compute regime scores for each date logger.info(f"Step 3/3: Computing regime scores {start} → {end}...") 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 detector = RegimeDetector() current = start score_count = 0 while current <= end: try: ps = PriceStructureScorer().compute(current) br = BreadthScorer().compute(current) if br.score == 50.0 and br.label == "No Data": current += timedelta(days=1) continue oi = OIMatrixScorer().compute(current) vol = VolatilityRegimeScorer().compute(current) r = detector.detect(ps.score, br.breadth_top50, vol.vol_regime.value, current) conn.execute("""INSERT OR REPLACE INTO regime_history (date, regime, confidence, regime_version, maturity_score, all_scores_json, confirmation_days) VALUES (?, ?, ?, ?, ?, ?, ?)""", (str(current), r.regime.value, r.confidence, r.regime_version, r.maturity_score, json.dumps(r.all_scores), r.confirmation_days)) score_count += 1 if score_count % 30 == 0: conn.commit() logger.info(f" Scored {score_count} days... ({current})") except Exception as e: logger.debug(f" Score skip {current}: {e}") current += timedelta(days=1) conn.commit() conn.close() logger.info(f"Backfill complete: {breadth_count} breadth + {score_count} regime days") def cmd_expectancy(args): """Query signal expectancy for current market state.""" from database import init_db from expectancy.engine import BayesianExpectancyEngine target = parse_date(args.date) if args.date else Date.today() init_db() state, _ = _build_market_state(target) engine = BayesianExpectancyEngine() signal = args.signal or "B3" report = engine.estimate(state, signal_type=signal, target_date=target) print(f"\n{'='*60}") print(f" {target} Signal Expectancy: {signal}") print(f"{'='*60}") print(f" Regime: {state.regime.value} (conf={state.regime_confidence:.2f})") print(f" Breadth: {state.breadth_bucket.value} (T50={state.breadth_top50:.0f})") print(f" OI State: {state.oi_state.value}") print(f" Volatility: {state.volatility_regime.value}") print(f"{'='*60}") for layer in report.layers: print(f" {layer.name:15s} N={layer.samples:4d} eff={layer.effective_samples:.0f} " f"raw={layer.raw_winrate or 0:.1%} post={layer.posterior_winrate:.1%} " f"ret={layer.avg_return or 0:+.1f}%") print(f"{'='*60}") print(f" Final: {report.final_estimate:.1%} " f"(sufficiency={report.sufficiency.value}, source={report.source})") if report.profit_factor: print(f" PF={report.profit_factor} MAE={report.max_adverse_excursion}%") print() def main(): parser = argparse.ArgumentParser( description="ChanMacro — Crypto Market Memory System" ) sub = parser.add_subparsers(dest="command", help="Commands") # fetch p_fetch = sub.add_parser("fetch", help="Fetch raw data") p_fetch.add_argument("--date", help="Target date (YYYY-MM-DD)") p_fetch.add_argument("--module", choices=["ohlcv", "breadth", "derivatives", "all"]) # score p_score = sub.add_parser("score", help="Compute scores and regime") p_score.add_argument("--date", help="Target date (YYYY-MM-DD)") # regime p_regime = sub.add_parser("regime", help="Show regime history") p_regime.add_argument("--days", type=int, default=30) # track p_track = sub.add_parser("track", help="Record a trading signal") p_track.add_argument("--date", help="Signal date (YYYY-MM-DD)") p_track.add_argument("--signal", required=True, help="Signal type (B1/B2/B3/S1/S2/S3)") p_track.add_argument("--price", type=float, required=True, help="Entry price") p_track.add_argument("--grade", choices=["A", "B", "C"], help="Signal quality grade") p_track.add_argument("--strength", type=float, help="Signal strength 0-100") # backfill p_backfill = sub.add_parser("backfill", help="Backfill historical scores") p_backfill.add_argument("--from", dest="from_date", required=True) p_backfill.add_argument("--to", dest="to_date") # expectancy p_expectancy = sub.add_parser("expectancy", help="Query signal expectancy") p_expectancy.add_argument("--date", help="Target date (YYYY-MM-DD)") p_expectancy.add_argument("--signal", default="B3", help="Signal type") # 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") args = parser.parse_args() if args.command == "fetch": cmd_fetch(args) elif args.command == "score": cmd_score(args) elif args.command == "regime": cmd_regime(args) elif args.command == "track": cmd_track(args) elif args.command == "backfill": cmd_backfill(args) elif args.command == "expectancy": cmd_expectancy(args) elif args.command == "validate": from validation.reporter import ValidationReporter report = ValidationReporter().run_all() print(report) elif args.command == "serve": 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() if __name__ == "__main__": main()