465 lines
17 KiB
Python
465 lines
17 KiB
Python
"""
|
|
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")
|
|
# detect (Chan BSP signals)
|
|
p_detect = sub.add_parser("detect", help="Detect Chan BSP signals and populate signal_features")
|
|
p_detect.add_argument("--from", dest="from_date", default="2024-01-01")
|
|
p_detect.add_argument("--to", dest="to_date")
|
|
# 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 == "detect":
|
|
from chan_integration import ChanSignalDetector
|
|
start = args.from_date
|
|
end = args.to_date or Date.today().isoformat()
|
|
detector = ChanSignalDetector()
|
|
count = detector.populate_signal_features(start, end)
|
|
logger.info(f"写入 {count} 条信号记录")
|
|
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()
|