Compare commits
10
Commits
f391020f78
...
2e905e7238
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
2e905e7238 | ||
|
|
02a52c04dd | ||
|
|
19c8f86862 | ||
|
|
ffe7074fef | ||
|
|
7b91f459d7 | ||
|
|
3c72aa1310 | ||
|
|
8d916371e2 | ||
|
|
7e19c9858e | ||
|
|
efb721b39f | ||
|
|
7813e319b4 |
@@ -0,0 +1,224 @@
|
||||
"""
|
||||
chan_integration.py — 缠论引擎集成:检测 BSP 信号并写入 signal_features。
|
||||
|
||||
复用 bsp_monitor/engine.py 的 ChanEngine 管线,对历史日线数据批量跑缠论,
|
||||
提取 B1/B2/B3/S1/S2/S3 信号,通过 SignalTracker 记录到 signal_features。
|
||||
"""
|
||||
|
||||
import sys
|
||||
import os
|
||||
from datetime import date as Date, timedelta
|
||||
from typing import List, Optional
|
||||
import logging
|
||||
|
||||
# 确保 Chan 引擎在路径上(与 bsp_monitor/engine.py 相同的路径设置)
|
||||
_PARENT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
|
||||
if _PARENT not in sys.path:
|
||||
sys.path.insert(0, _PARENT)
|
||||
|
||||
import pandas as pd
|
||||
|
||||
from ChanEnum import Chan_BSP_TYPE, Chan_BSP_DIR
|
||||
from ChanBSP import ChanBSP
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
class ChanSignalDetector:
|
||||
"""
|
||||
对历史日线数据运行缠论管线,提取所有 BSP 信号。
|
||||
|
||||
Usage:
|
||||
detector = ChanSignalDetector()
|
||||
signals = detector.detect_from_db("2026-01-01", "2026-06-24")
|
||||
# → [{"date": Date, "signal_type": "B3", "entry_price": 96500, ...}, ...]
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
from TF_DF import TF_DF as _TF_DF_Class
|
||||
self._TF_DF_Class = _TF_DF_Class
|
||||
|
||||
def detect_from_db(self, start_date: str, end_date: str) -> list[dict]:
|
||||
"""从数据库加载日线数据,跑缠论管线,提取信号。"""
|
||||
from database import get_connection
|
||||
conn = get_connection()
|
||||
df = pd.read_sql_query(
|
||||
"SELECT date, open, high, low, close, volume "
|
||||
"FROM ohlcv_daily WHERE symbol='BTC/USDT:USDT' "
|
||||
"AND date BETWEEN ? AND ? ORDER BY date",
|
||||
conn, params=(start_date, end_date)
|
||||
)
|
||||
conn.close()
|
||||
|
||||
if df.empty or len(df) < 50:
|
||||
logger.warning(f"日线数据不足: {len(df)} 根")
|
||||
return []
|
||||
|
||||
return self.detect_from_df(df)
|
||||
|
||||
def detect_from_df(self, df: pd.DataFrame) -> list[dict]:
|
||||
"""从 DataFrame 运行缠论管线,提取 BSP 信号。"""
|
||||
# 需要 datetime 列才能跑 TF_DF
|
||||
df = df.copy()
|
||||
df["timestamp"] = pd.to_datetime(df["date"])
|
||||
df["date"] = df["timestamp"]
|
||||
|
||||
try:
|
||||
engine = self._build_engine(df)
|
||||
except Exception as e:
|
||||
logger.error(f"缠论管线失败: {e}")
|
||||
return []
|
||||
|
||||
return self._extract_signals(engine)
|
||||
|
||||
def _build_engine(self, df: pd.DataFrame):
|
||||
"""构建缠论管线(对齐 bsp_monitor/engine.py 的 ChanEngine)。"""
|
||||
from TF_DF import TF_DF as _TF_DF_Class
|
||||
|
||||
if df.empty or len(df) < 50:
|
||||
raise ValueError(f"数据不足: {len(df)} 根 K 线")
|
||||
|
||||
if "date" not in df.columns and "timestamp" in df.columns:
|
||||
df["date"] = df["timestamp"]
|
||||
|
||||
# 使用 __new__ 避免触发 TF_DF.__init__
|
||||
engine = type('ChanEngine', (), {})() # 简单容器
|
||||
tf = _TF_DF_Class.__new__(_TF_DF_Class)
|
||||
|
||||
df_with_indicators = tf.add_indicators(df.copy())
|
||||
engine.klu_list = tf.get_klu_list(df_with_indicators)
|
||||
engine.klc_list = tf.get_klc_list(engine.klu_list)
|
||||
engine.bi_list = tf.cal_bi_list(engine.klc_list)
|
||||
engine.seg_list = tf.get_seg_list(engine.bi_list)
|
||||
engine.bi_zs_list = tf.cal_bi_zs(engine.seg_list)
|
||||
engine.bsp_list = tf.find_all_bsp(engine.bi_list, engine.bi_zs_list)
|
||||
|
||||
return engine
|
||||
|
||||
def _extract_signals(self, engine) -> list[dict]:
|
||||
"""从 ChanEngine 输出中提取所有 BSP 信号。"""
|
||||
signals = []
|
||||
for bsp in engine.bsp_list:
|
||||
if bsp.type == Chan_BSP_TYPE.NONE:
|
||||
continue
|
||||
if bsp.klc is None:
|
||||
continue
|
||||
|
||||
signal_type = self._bsp_type_str(bsp.type)
|
||||
entry_price = bsp.klc.close
|
||||
signal_date = self._klc_date(bsp.klc)
|
||||
|
||||
if signal_date is None:
|
||||
continue
|
||||
|
||||
# 信号质量:根据分型强度判断
|
||||
strength = self._calc_strength(bsp)
|
||||
grade = "A" if strength >= 70 else "B" if strength >= 50 else "C"
|
||||
|
||||
signals.append({
|
||||
"date": signal_date,
|
||||
"signal_type": signal_type,
|
||||
"entry_price": float(entry_price),
|
||||
"signal_grade": grade,
|
||||
"signal_strength": float(strength),
|
||||
})
|
||||
|
||||
details = ", ".join(f"{s['signal_type']}({s['date']})" for s in signals)
|
||||
logger.info(f"检测到 {len(signals)} 个信号: {details}")
|
||||
return signals
|
||||
|
||||
def populate_signal_features(self, start_date: str = "2024-01-01",
|
||||
end_date: Optional[str] = None) -> int:
|
||||
"""
|
||||
完整流程:检测信号 → 计算市场状态 → 写入 signal_features。
|
||||
|
||||
Returns: 写入的信号数量。
|
||||
"""
|
||||
if end_date is None:
|
||||
end_date = Date.today().isoformat()
|
||||
|
||||
logger.info(f"开始信号检测: {start_date} → {end_date}")
|
||||
|
||||
# Step 1: 检测缠论信号
|
||||
signals = self.detect_from_db(start_date, end_date)
|
||||
if not signals:
|
||||
logger.warning("未检测到任何 BSP 信号")
|
||||
return 0
|
||||
|
||||
# Step 2: 去重 — 跳过已存在的信号
|
||||
from database import get_connection
|
||||
conn = get_connection()
|
||||
existing = set()
|
||||
for row in conn.execute(
|
||||
"SELECT date, signal_type FROM signal_features"
|
||||
).fetchall():
|
||||
existing.add((row[0], row[1]))
|
||||
conn.close()
|
||||
|
||||
new_signals = [s for s in signals
|
||||
if (str(s["date"]), s["signal_type"]) not in existing]
|
||||
if not new_signals:
|
||||
logger.info("所有信号已存在,跳过")
|
||||
return 0
|
||||
|
||||
# Step 3: 写入 signal_features
|
||||
from expectancy.tracker import SignalTracker
|
||||
tracker = SignalTracker()
|
||||
count = tracker.backfill_signals(new_signals)
|
||||
|
||||
logger.info(f"信号入库完成: {count}/{len(signals)}")
|
||||
return count
|
||||
|
||||
@staticmethod
|
||||
def _bsp_type_str(t: Chan_BSP_TYPE) -> str:
|
||||
mapping = {
|
||||
Chan_BSP_TYPE.B1: "B1", Chan_BSP_TYPE.B2: "B2", Chan_BSP_TYPE.B3: "B3",
|
||||
Chan_BSP_TYPE.S1: "S1", Chan_BSP_TYPE.S2: "S2", Chan_BSP_TYPE.S3: "S3",
|
||||
}
|
||||
return mapping.get(t, "UNKNOWN")
|
||||
|
||||
@staticmethod
|
||||
def _klc_date(klc) -> Optional[Date]:
|
||||
"""从 KLC 提取信号确认日期。"""
|
||||
end_time = getattr(klc, "end_time", None)
|
||||
if end_time is None:
|
||||
start_time = getattr(klc, "start_time", None)
|
||||
if start_time is None:
|
||||
return None
|
||||
end_time = start_time
|
||||
if hasattr(end_time, "date"):
|
||||
return end_time.date()
|
||||
if isinstance(end_time, str):
|
||||
return Date.fromisoformat(end_time[:10])
|
||||
return None
|
||||
|
||||
@staticmethod
|
||||
def _calc_strength(bsp: ChanBSP) -> float:
|
||||
"""根据 BSP 特征计算信号强度 0-100。"""
|
||||
score = 50.0
|
||||
klc = bsp.klc
|
||||
if klc is None:
|
||||
return score
|
||||
|
||||
# 分型强度
|
||||
from ChanEnum import Chan_KLC_FX
|
||||
fx = getattr(klc, "klc_fx_type", None)
|
||||
if fx is not None:
|
||||
strong_fxs = {Chan_KLC_FX.TOP2, Chan_KLC_FX.TOP3, Chan_KLC_FX.BOTTOM2, Chan_KLC_FX.BOTTOM3}
|
||||
medium_fxs = {Chan_KLC_FX.TOP1, Chan_KLC_FX.BOTTOM1, Chan_KLC_FX.TOP4, Chan_KLC_FX.BOTTOM4}
|
||||
if fx in strong_fxs:
|
||||
score += 25
|
||||
elif fx in medium_fxs:
|
||||
score += 10
|
||||
|
||||
# BSP 类型
|
||||
if bsp.type in (Chan_BSP_TYPE.B1, Chan_BSP_TYPE.S1):
|
||||
score += 10 # 一类买卖点: 背驰确认, 额外加分
|
||||
|
||||
# 笔特征
|
||||
bi = getattr(bsp, "bi", None)
|
||||
if bi and hasattr(bi, "height") and hasattr(bi, "width"):
|
||||
if bi.width > 3 and abs(bi.height) > 100:
|
||||
score += 10
|
||||
|
||||
return min(score, 100.0)
|
||||
+150
-50
@@ -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(
|
||||
@@ -49,6 +50,24 @@ def _build_market_state(target: Date) -> tuple:
|
||||
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
|
||||
|
||||
|
||||
@@ -93,13 +112,13 @@ def cmd_fetch(args):
|
||||
|
||||
def cmd_score(args):
|
||||
"""Compute all factor scores and regime for a date."""
|
||||
from database import init_db, get_connection
|
||||
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, regime_result = _build_market_state(target)
|
||||
state, _ = _build_market_state(target)
|
||||
|
||||
# Output
|
||||
ps = state.price_structure_score
|
||||
@@ -128,26 +147,6 @@ def cmd_score(args):
|
||||
print(f" Market State Hash: {state.market_state_hash}")
|
||||
print()
|
||||
|
||||
# Store regime to DB
|
||||
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),
|
||||
state.regime.value,
|
||||
state.regime_confidence,
|
||||
state.regime_version,
|
||||
state.regime_maturity_score,
|
||||
json.dumps(regime_result.all_scores),
|
||||
regime_result.prior_regime.value if regime_result.prior_regime else None,
|
||||
regime_result.confirmation_days,
|
||||
))
|
||||
conn.commit()
|
||||
conn.close()
|
||||
|
||||
return state
|
||||
|
||||
|
||||
@@ -193,23 +192,104 @@ def cmd_track(args):
|
||||
|
||||
|
||||
def cmd_backfill(args):
|
||||
"""Backfill historical scores and/or signals."""
|
||||
"""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()
|
||||
|
||||
# First, backfill OHLCV data
|
||||
logger.info(f"Backfilling OHLCV from {start} to {end}...")
|
||||
fetcher = OHLCVFetcher()
|
||||
df = fetcher.fetch()
|
||||
if not df.empty:
|
||||
fetcher.store_df(df)
|
||||
# Step 1: Ensure OHLCV data exists for the range
|
||||
logger.info(f"Step 1/3: Fetching BTC OHLCV...")
|
||||
OHLCVFetcher().store_df(OHLCVFetcher().fetch())
|
||||
|
||||
# Then compute scores for each date
|
||||
# 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
|
||||
@@ -217,10 +297,8 @@ def cmd_backfill(args):
|
||||
from regime_detector import RegimeDetector
|
||||
|
||||
detector = RegimeDetector()
|
||||
conn = get_connection()
|
||||
|
||||
current = start
|
||||
count = 0
|
||||
score_count = 0
|
||||
while current <= end:
|
||||
try:
|
||||
ps = PriceStructureScorer().compute(current)
|
||||
@@ -228,33 +306,27 @@ def cmd_backfill(args):
|
||||
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)
|
||||
r = detector.detect(ps.score, br.breadth_top50, vol.vol_regime.value, current)
|
||||
|
||||
conn.execute("""
|
||||
INSERT OR REPLACE INTO regime_history
|
||||
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,
|
||||
))
|
||||
count += 1
|
||||
if count % 30 == 0:
|
||||
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" Backfilled {count} days... ({current})")
|
||||
logger.info(f" Scored {score_count} days... ({current})")
|
||||
except Exception as e:
|
||||
logger.debug(f" Skip {current}: {e}")
|
||||
logger.debug(f" Score skip {current}: {e}")
|
||||
current += timedelta(days=1)
|
||||
|
||||
conn.commit()
|
||||
conn.close()
|
||||
logger.info(f"Backfill complete: {count} days scored")
|
||||
logger.info(f"Backfill complete: {breadth_count} breadth + {score_count} regime days")
|
||||
|
||||
|
||||
def cmd_expectancy(args):
|
||||
@@ -333,6 +405,12 @@ 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")
|
||||
# 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")
|
||||
|
||||
@@ -354,8 +432,30 @@ def main():
|
||||
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":
|
||||
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()
|
||||
|
||||
|
||||
@@ -0,0 +1,153 @@
|
||||
"""
|
||||
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()
|
||||
|
||||
# 检测新信号(每天运行一次,UTC 0 点后首次触发)
|
||||
now = datetime.now(timezone.utc)
|
||||
if self._last_run is None or now.date() > self._last_run.date():
|
||||
try:
|
||||
from chan_integration import ChanSignalDetector
|
||||
detector = ChanSignalDetector()
|
||||
# 检测最近 90 天的 4h 信号
|
||||
count = detector.populate_signal_features(
|
||||
start_date=(today - __import__('datetime').timedelta(days=90)).isoformat(),
|
||||
end_date=today.isoformat(),
|
||||
)
|
||||
if count > 0:
|
||||
logger.info(f"新增 {count} 条信号记录")
|
||||
except Exception as e:
|
||||
logger.debug(f"信号检测跳过: {e}")
|
||||
|
||||
self._last_run = now
|
||||
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
|
||||
+22
-2
@@ -6,7 +6,8 @@ import sys
|
||||
import os
|
||||
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
|
||||
|
||||
from datetime import date as Date, timedelta
|
||||
import json
|
||||
from datetime import date as Date
|
||||
from flask import Flask, render_template, jsonify, request
|
||||
|
||||
from database import get_connection
|
||||
@@ -23,7 +24,7 @@ app = Flask(__name__)
|
||||
|
||||
|
||||
def _build_state(target: Date):
|
||||
"""Shared: build MarketStateVector for a date."""
|
||||
"""Build MarketStateVector and persist regime to DB."""
|
||||
ps = PriceStructureScorer().compute(target)
|
||||
br = BreadthScorer().compute(target)
|
||||
oi = OIMatrixScorer().compute(target)
|
||||
@@ -44,6 +45,23 @@ def _build_state(target: Date):
|
||||
oi_matrix_score=oi, volatility_regime_score=vol,
|
||||
)
|
||||
state.market_state_hash = state.compute_hash()
|
||||
|
||||
# Persist regime to DB so load_state() works across requests
|
||||
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
|
||||
|
||||
|
||||
@@ -155,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)
|
||||
|
||||
@@ -1,60 +1,75 @@
|
||||
// dashboard.js — ChanMacro dashboard
|
||||
// dashboard.js — ChanMacro
|
||||
|
||||
const C = { TREND: "#3fb950", RANGE: "#d29922", PANIC: "#f85149" };
|
||||
let regimeChart = null, breadthChart = null;
|
||||
|
||||
const REGIME_COLORS = { TREND: "#3fb950", RANGE: "#d29922", PANIC: "#f85149" };
|
||||
const BUCKET_CLASS = { EXTREME: "bucket-EXTREME", STRONG: "bucket-STRONG",
|
||||
NORMAL: "bucket-NORMAL", WEAK: "bucket-WEAK", PANIC: "bucket-PANIC" };
|
||||
|
||||
async function loadState() {
|
||||
try {
|
||||
const r = await fetch("/api/state");
|
||||
const d = await r.json();
|
||||
if (d.error) { document.getElementById("db-status").textContent = d.error; return; }
|
||||
if (d.error) { document.getElementById("update-time").textContent = d.error; return; }
|
||||
|
||||
document.getElementById("db-status").textContent = "✓ " + d.date;
|
||||
document.getElementById("update-time").textContent = "更新于 " + new Date().toLocaleTimeString();
|
||||
document.getElementById("update-time").textContent = d.date;
|
||||
|
||||
// Hero
|
||||
const regime = d.regime;
|
||||
document.getElementById("hero-regime").textContent = regime === "TREND" ? "趋势" : regime === "RANGE" ? "震荡" : "恐慌";
|
||||
document.getElementById("hero-regime").className = "hero-regime regime-" + regime;
|
||||
const names = { TREND: "TREND", RANGE: "RANGE", PANIC: "PANIC" };
|
||||
document.getElementById("hero-regime").textContent = names[regime] || regime;
|
||||
document.getElementById("hero-regime").className = "regime-name " + regime.toLowerCase();
|
||||
document.getElementById("hero-badge").textContent = regime;
|
||||
document.getElementById("hero-badge").className = "badge-regime badge-" + regime;
|
||||
document.getElementById("hero-badge").className = "regime-badge " + regime.toLowerCase();
|
||||
document.getElementById("hero-conf").textContent = (d.regime_confidence * 100).toFixed(0) + "%";
|
||||
document.getElementById("hero-maturity").textContent = d.regime_maturity.toFixed(0) + "/100";
|
||||
|
||||
// Factors
|
||||
document.getElementById("f-price").textContent = d.price_structure.score.toFixed(0);
|
||||
document.getElementById("f-price").style.color =
|
||||
document.getElementById("hero-maturity").textContent = d.regime_maturity.toFixed(0);
|
||||
document.getElementById("hero-ps").textContent = d.price_structure.score.toFixed(0);
|
||||
document.getElementById("hero-ps").style.color =
|
||||
d.price_structure.score >= 60 ? "#3fb950" : d.price_structure.score >= 40 ? "#d29922" : "#f85149";
|
||||
document.getElementById("f-price-sub").textContent = d.price_structure.label;
|
||||
document.getElementById("f-price-narr").textContent = d.price_structure.narrative;
|
||||
document.getElementById("hero-br").textContent = d.breadth.score.toFixed(0);
|
||||
document.getElementById("hero-br").style.color =
|
||||
d.breadth.bucket === "EXTREME" || d.breadth.bucket === "STRONG" ? "#3fb950" :
|
||||
d.breadth.bucket === "WEAK" || d.breadth.bucket === "PANIC" ? "#f85149" : "#d29922";
|
||||
|
||||
const b = d.breadth;
|
||||
document.getElementById("f-breadth").textContent = b.score.toFixed(0);
|
||||
document.getElementById("f-breadth").setAttribute("style",
|
||||
"color: " + (b.bucket === "EXTREME" || b.bucket === "STRONG" ? "#3fb950" :
|
||||
b.bucket === "WEAK" || b.bucket === "PANIC" ? "#f85149" :
|
||||
b.bucket === "NORMAL" ? "#d29922" : "#e6edf3"));
|
||||
// Factor cards
|
||||
const ps = d.price_structure;
|
||||
document.getElementById("f-price").textContent = ps.score.toFixed(0);
|
||||
document.getElementById("f-price").style.color =
|
||||
ps.score >= 60 ? "#3fb950" : ps.score >= 40 ? "#d29922" : "#f85149";
|
||||
document.getElementById("f-price-sub").textContent =
|
||||
`趋势 ${ps.trend.toFixed(0)} · 波动 ${ps.vol_comp.toFixed(0)} · 动量 ${ps.momentum.toFixed(0)}`;
|
||||
document.getElementById("bar-price").style.width = ps.score + "%";
|
||||
document.getElementById("bar-price").style.background =
|
||||
ps.score >= 60 ? "#3fb950" : ps.score >= 40 ? "#d29922" : "#f85149";
|
||||
|
||||
const br = d.breadth;
|
||||
document.getElementById("f-breadth").textContent = br.score.toFixed(0);
|
||||
document.getElementById("f-breadth").style.color =
|
||||
br.bucket === "EXTREME" || br.bucket === "STRONG" ? "#3fb950" :
|
||||
br.bucket === "WEAK" || br.bucket === "PANIC" ? "#f85149" : "#d29922";
|
||||
document.getElementById("f-breadth-sub").textContent =
|
||||
`${b.bucket} · T20=${b.top20.toFixed(0)} T50=${b.top50.toFixed(0)} div=${b.divergence > 0 ? "+" : ""}${b.divergence.toFixed(0)}`;
|
||||
document.getElementById("f-breadth-narr").textContent = b.narrative;
|
||||
`${br.bucket} · T20=${br.top20.toFixed(0)} T50=${br.top50.toFixed(0)}`;
|
||||
document.getElementById("bar-breadth").style.width = br.score + "%";
|
||||
document.getElementById("bar-breadth").style.background =
|
||||
br.bucket === "EXTREME" || br.bucket === "STRONG" ? "#3fb950" :
|
||||
br.bucket === "WEAK" || br.bucket === "PANIC" ? "#f85149" : "#d29922";
|
||||
|
||||
document.getElementById("f-oi").textContent = d.oi_state;
|
||||
document.getElementById("f-oi").textContent = d.oi_state.toUpperCase().replace(" ", "\n");
|
||||
document.getElementById("f-oi").style.color =
|
||||
d.oi_state === "New Longs" ? "#3fb950" : d.oi_state === "New Shorts" ? "#f85149" :
|
||||
d.oi_state === "Short Covering" ? "#d29922" : d.oi_state === "Long Exit" ? "#f85149" : "#e6edf3";
|
||||
document.getElementById("f-oi-sub").textContent = `分数: ${d.oi_score.toFixed(0)}`;
|
||||
document.getElementById("f-oi-narr").textContent = d.oi_narrative;
|
||||
d.oi_state === "New Longs" ? "#3fb950" : d.oi_state.includes("Short") || d.oi_state === "Long Exit" ? "#f85149" : "#8b949e";
|
||||
document.getElementById("f-oi-sub").textContent = d.oi_narrative;
|
||||
|
||||
document.getElementById("f-vol").textContent = d.volatility;
|
||||
const vm = { LOW_VOL: "低波动", NORMAL_VOL: "正常", HIGH_VOL: "高波动", EXPLOSIVE_VOL: "极端" };
|
||||
document.getElementById("f-vol").textContent = vm[d.volatility] || d.volatility;
|
||||
document.getElementById("f-vol").style.color =
|
||||
d.volatility === "LOW_VOL" ? "#58a6ff" : d.volatility === "NORMAL_VOL" ? "#e6edf3" :
|
||||
d.volatility === "LOW_VOL" ? "#58a6ff" : d.volatility === "NORMAL_VOL" ? "#8b949e" :
|
||||
d.volatility === "HIGH_VOL" ? "#d29922" : "#f85149";
|
||||
document.getElementById("f-vol-sub").textContent = `分数: ${d.price_structure.score.toFixed(0)}`;
|
||||
document.getElementById("f-vol-sub").textContent = d.volatility;
|
||||
document.getElementById("bar-vol").style.width =
|
||||
(d.volatility === "EXPLOSIVE_VOL" ? 95 : d.volatility === "HIGH_VOL" ? 70 :
|
||||
d.volatility === "NORMAL_VOL" ? 40 : 20) + "%";
|
||||
document.getElementById("bar-vol").style.background =
|
||||
d.volatility === "EXPLOSIVE_VOL" ? "#f85149" : d.volatility === "HIGH_VOL" ? "#d29922" :
|
||||
d.volatility === "NORMAL_VOL" ? "#8b949e" : "#58a6ff";
|
||||
} catch (e) {
|
||||
document.getElementById("db-status").textContent = "连接失败";
|
||||
document.getElementById("update-time").textContent = "连接失败";
|
||||
}
|
||||
}
|
||||
|
||||
@@ -63,72 +78,50 @@ async function loadHistory() {
|
||||
const r = await fetch("/api/history?days=60");
|
||||
const d = await r.json();
|
||||
|
||||
// Regime chart
|
||||
const dates = d.regimes.map(x => x.date);
|
||||
const regimes = d.regimes.map(x => x.regime);
|
||||
const colors = regimes.map(r => REGIME_COLORS[r] || "#8b949e");
|
||||
const colors = d.regimes.map(x => C[x.regime] || "#5c6675");
|
||||
|
||||
if (regimeChart) regimeChart.destroy();
|
||||
const ctx1 = document.getElementById("chart-regime").getContext("2d");
|
||||
regimeChart = new Chart(ctx1, {
|
||||
regimeChart = new Chart(document.getElementById("chart-regime").getContext("2d"), {
|
||||
type: "bar",
|
||||
data: {
|
||||
labels: dates,
|
||||
datasets: [{
|
||||
label: "置信度",
|
||||
data: d.regimes.map(x => x.confidence * 100),
|
||||
backgroundColor: colors,
|
||||
borderWidth: 0,
|
||||
borderRadius: 2,
|
||||
}]
|
||||
},
|
||||
data: { labels: dates, datasets: [{ data: d.regimes.map(x => x.confidence * 100),
|
||||
backgroundColor: colors, borderWidth: 0, borderRadius: 2 }] },
|
||||
options: {
|
||||
responsive: true,
|
||||
maintainAspectRatio: false,
|
||||
plugins: {
|
||||
legend: { display: false },
|
||||
tooltip: {
|
||||
callbacks: {
|
||||
label: ctx => `${d.regimes[ctx.dataIndex].regime} · ${ctx.raw.toFixed(0)}%`
|
||||
}
|
||||
}
|
||||
},
|
||||
responsive: true, maintainAspectRatio: false,
|
||||
plugins: { legend: { display: false } },
|
||||
scales: {
|
||||
x: { ticks: { color: "#8b949e", maxTicksLimit: 15, maxRotation: 45 } },
|
||||
y: { max: 100, ticks: { color: "#8b949e" } }
|
||||
x: { ticks: { color: "#5c6675", maxTicksLimit: 15, maxRotation: 45, font: { size: 10 } },
|
||||
grid: { color: "#151a23" } },
|
||||
y: { max: 100, ticks: { color: "#5c6675", font: { size: 10 } }, grid: { color: "#151a23" } }
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Breadth chart
|
||||
if (breadthChart) breadthChart.destroy();
|
||||
const ctx2 = document.getElementById("chart-breadth").getContext("2d");
|
||||
breadthChart = new Chart(ctx2, {
|
||||
breadthChart = new Chart(document.getElementById("chart-breadth").getContext("2d"), {
|
||||
type: "line",
|
||||
data: {
|
||||
labels: d.breadth.map(x => x.date),
|
||||
datasets: [
|
||||
{ label: "上涨", data: d.breadth.map(x => x.advance), borderColor: "#3fb950",
|
||||
backgroundColor: "rgba(63,185,80,0.1)", fill: true, tension: 0.3, pointRadius: 0 },
|
||||
backgroundColor: "rgba(63,185,80,0.08)", fill: true, tension: 0.3, pointRadius: 0 },
|
||||
{ label: "下跌", data: d.breadth.map(x => x.decline), borderColor: "#f85149",
|
||||
backgroundColor: "rgba(248,81,73,0.1)", fill: true, tension: 0.3, pointRadius: 0 },
|
||||
backgroundColor: "rgba(248,81,73,0.06)", fill: true, tension: 0.3, pointRadius: 0 },
|
||||
{ label: ">EMA20", data: d.breadth.map(x => x.above_ema20), borderColor: "#58a6ff",
|
||||
borderDash: [4, 2], tension: 0.3, pointRadius: 0 },
|
||||
borderDash: [3, 3], tension: 0.3, pointRadius: 0 },
|
||||
]
|
||||
},
|
||||
options: {
|
||||
responsive: true,
|
||||
maintainAspectRatio: false,
|
||||
plugins: { legend: { labels: { color: "#8b949e", usePointStyle: true, boxWidth: 8 } } },
|
||||
responsive: true, maintainAspectRatio: false,
|
||||
plugins: { legend: { labels: { color: "#5c6675", usePointStyle: true, boxWidth: 6, font: { size: 10 } } } },
|
||||
scales: {
|
||||
x: { ticks: { color: "#8b949e", maxTicksLimit: 15, maxRotation: 45 } },
|
||||
y: { max: 50, ticks: { color: "#8b949e" } }
|
||||
x: { ticks: { color: "#5c6675", maxTicksLimit: 15, maxRotation: 45, font: { size: 10 } },
|
||||
grid: { color: "#151a23" } },
|
||||
y: { ticks: { color: "#5c6675", font: { size: 10 } }, grid: { color: "#151a23" } }
|
||||
}
|
||||
}
|
||||
});
|
||||
} catch (e) {
|
||||
console.error("History load failed:", e);
|
||||
}
|
||||
} catch (e) { console.error(e); }
|
||||
}
|
||||
|
||||
async function loadExpectancy() {
|
||||
@@ -136,36 +129,32 @@ async function loadExpectancy() {
|
||||
try {
|
||||
const r = await fetch(`/api/expectancy?signal=${signal}`);
|
||||
const d = await r.json();
|
||||
if (d.error) { document.getElementById("exp-layers").innerHTML = `<tr><td colspan="6">${d.error}</td></tr>`; return; }
|
||||
if (d.error) { document.getElementById("exp-layers").innerHTML =
|
||||
`<tr><td colspan="6" style="color:#f85149">${d.error}</td></tr>`; return; }
|
||||
|
||||
document.getElementById("exp-sufficiency").textContent = d.sufficiency;
|
||||
document.getElementById("exp-sufficiency").className =
|
||||
"badge " + (d.sufficiency === "HIGH" ? "bg-success" : d.sufficiency === "MEDIUM" ? "bg-warning" :
|
||||
d.sufficiency === "LOW" ? "bg-danger" : "bg-secondary");
|
||||
const el = document.getElementById("exp-sufficiency");
|
||||
el.textContent = d.sufficiency;
|
||||
el.className = "suff suff-" + d.sufficiency;
|
||||
|
||||
let html = "";
|
||||
for (const l of d.layers) {
|
||||
html += `<tr>
|
||||
<td>${l.name}</td>
|
||||
<td>${l.samples}</td>
|
||||
<td>${l.effective_samples.toFixed(0)}</td>
|
||||
<td>${l.name}</td><td>${l.samples}</td><td>${l.effective_samples.toFixed(0)}</td>
|
||||
<td>${l.raw_winrate ? (l.raw_winrate * 100).toFixed(1) + "%" : "—"}</td>
|
||||
<td><strong>${(l.posterior_winrate * 100).toFixed(1)}%</strong></td>
|
||||
<td>${l.avg_return ? (l.avg_return > 0 ? "+" : "") + l.avg_return.toFixed(1) + "%" : "—"}</td>
|
||||
<td style="color:${l.avg_return > 0 ? '#3fb950' : l.avg_return < 0 ? '#f85149' : '#8b949e'}">${l.avg_return ? (l.avg_return > 0 ? "+" : "") + l.avg_return.toFixed(2) + "%" : "—"}</td>
|
||||
</tr>`;
|
||||
}
|
||||
document.getElementById("exp-layers").innerHTML = html;
|
||||
|
||||
let summary = `最终估计: <strong>${(d.final_estimate * 100).toFixed(1)}%</strong>`;
|
||||
if (d.avg_return_7d) summary += ` · 平均收益: <strong>${d.avg_return_7d > 0 ? "+" : ""}${d.avg_return_7d.toFixed(1)}%</strong>`;
|
||||
if (d.profit_factor) summary += ` · 盈亏比: <strong>${d.profit_factor}</strong>`;
|
||||
document.getElementById("exp-summary").innerHTML = summary;
|
||||
} catch (e) {
|
||||
console.error("Expectancy load failed:", e);
|
||||
}
|
||||
let s = `后验胜率 <strong style="color:#58a6ff">${(d.final_estimate * 100).toFixed(1)}%</strong>`;
|
||||
if (d.avg_return_7d) s += ` · 平均收益 <strong>${d.avg_return_7d > 0 ? "+" : ""}${d.avg_return_7d.toFixed(2)}%</strong>`;
|
||||
if (d.profit_factor) s += ` · 盈亏比 <strong>${d.profit_factor}</strong>`;
|
||||
if (d.max_adverse) s += ` · MAE <strong>${d.max_adverse.toFixed(1)}%</strong>`;
|
||||
document.getElementById("exp-summary").innerHTML = s;
|
||||
} catch (e) { console.error(e); }
|
||||
}
|
||||
|
||||
// Init
|
||||
loadState();
|
||||
loadHistory();
|
||||
loadExpectancy();
|
||||
|
||||
+126
-107
@@ -4,138 +4,157 @@
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>ChanMacro — 市场状态</title>
|
||||
<link href="https://cdn.jsdelivr.net/npm/bootstrap@5.3.0/dist/css/bootstrap.min.css" rel="stylesheet">
|
||||
<script src="https://cdn.jsdelivr.net/npm/chart.js@4.4.0/dist/chart.umd.min.js"></script>
|
||||
<style>
|
||||
:root { --bg: #0d1117; --card: #161b22; --border: #30363d; --text: #e6edf3; --muted: #8b949e;
|
||||
--green: #3fb950; --red: #f85149; --orange: #d29922; --blue: #58a6ff; }
|
||||
body { background: var(--bg); color: var(--text); font-family: -apple-system, BlinkMacSystemFont, sans-serif; }
|
||||
.card { background: var(--card); border: 1px solid var(--border); border-radius: 10px; }
|
||||
.hero-regime { font-size: 3rem; font-weight: 700; }
|
||||
.hero-conf { font-size: 1.2rem; color: var(--muted); }
|
||||
.factor-value { font-size: 2.2rem; font-weight: 700; color: var(--text); }
|
||||
.factor-label { color: var(--muted); font-size: 0.85rem; }
|
||||
.regime-TREND { color: var(--green); }
|
||||
.regime-RANGE { color: var(--orange); }
|
||||
.regime-PANIC { color: var(--red); }
|
||||
.bucket-EXTREME, .bucket-STRONG { color: var(--green); }
|
||||
.bucket-NORMAL { color: var(--orange); }
|
||||
.bucket-WEAK, .bucket-PANIC { color: var(--red); }
|
||||
.badge-regime { font-size: 0.85rem; padding: 4px 12px; border-radius: 20px; }
|
||||
.badge-TREND { background: #1a3a1a; color: var(--green); }
|
||||
.badge-RANGE { background: #3a2a0a; color: var(--orange); }
|
||||
.badge-PANIC { background: #3a0a0a; color: var(--red); }
|
||||
.narrative { color: var(--muted); font-size: 0.9rem; }
|
||||
canvas { max-height: 300px; }
|
||||
.text-muted { color: var(--muted) !important; }
|
||||
.table { color: var(--text); }
|
||||
.table-dark { --bs-table-color: var(--text); --bs-table-bg: var(--card); }
|
||||
.form-select { background-color: var(--card); color: var(--text); border-color: var(--border); }
|
||||
.btn-primary { background-color: var(--blue); border-color: var(--blue); }
|
||||
strong { color: var(--text); }
|
||||
small { color: var(--muted); }
|
||||
* { margin: 0; padding: 0; box-sizing: border-box; }
|
||||
body { background: #0a0e14; color: #c9d1d9; font-family: -apple-system, BlinkMacSystemFont, "SF Mono", monospace; }
|
||||
.app { max-width: 1200px; margin: 0 auto; padding: 20px 24px; }
|
||||
|
||||
/* Header */
|
||||
.header { display: flex; justify-content: space-between; align-items: flex-end; padding: 20px 0 28px;
|
||||
border-bottom: 1px solid #1c2333; margin-bottom: 24px; }
|
||||
.header h1 { font-size: 22px; font-weight: 600; letter-spacing: 1px; }
|
||||
.header h1 span { color: #58a6ff; }
|
||||
.header .time { color: #5c6675; font-size: 13px; }
|
||||
.dot { display: inline-block; width: 7px; height: 7px; border-radius: 50%; background: #3fb950;
|
||||
margin-right: 6px; animation: pulse 2s infinite; }
|
||||
@keyframes pulse { 0%,100%{opacity:1} 50%{opacity:0.4} }
|
||||
|
||||
/* Regime Hero */
|
||||
.hero { display: flex; gap: 16px; margin-bottom: 24px; }
|
||||
.hero-card { flex: 1; background: #11161e; border: 1px solid #1c2333; border-radius: 8px; padding: 20px 24px; }
|
||||
.hero-card.main { flex: 2; display: flex; align-items: center; gap: 28px; }
|
||||
.regime-badge { display: inline-block; padding: 5px 16px; border-radius: 4px; font-size: 13px;
|
||||
font-weight: 600; letter-spacing: 2px; }
|
||||
.regime-badge.trend { background: rgba(63,185,80,0.12); color: #3fb950; border: 1px solid rgba(63,185,80,0.3); }
|
||||
.regime-badge.range { background: rgba(210,153,34,0.12); color: #d29922; border: 1px solid rgba(210,153,34,0.3); }
|
||||
.regime-badge.panic { background: rgba(248,81,73,0.12); color: #f85149; border: 1px solid rgba(248,81,73,0.3); }
|
||||
.regime-name { font-size: 42px; font-weight: 700; letter-spacing: 2px; }
|
||||
.regime-name.trend { color: #3fb950; }
|
||||
.regime-name.range { color: #d29922; }
|
||||
.regime-name.panic { color: #f85149; }
|
||||
.hero-stat { text-align: center; }
|
||||
.hero-stat .val { font-size: 28px; font-weight: 600; color: #e6edf3; }
|
||||
.hero-stat .lbl { font-size: 11px; color: #5c6675; letter-spacing: 1px; margin-top: 4px; }
|
||||
|
||||
/* Factor Grid */
|
||||
.grid { display: grid; grid-template-columns: repeat(4, 1fr); gap: 12px; margin-bottom: 24px; }
|
||||
.fcard { background: #11161e; border: 1px solid #1c2333; border-radius: 8px; padding: 18px 20px; }
|
||||
.fcard .title { font-size: 11px; color: #5c6675; letter-spacing: 1.5px; margin-bottom: 10px; }
|
||||
.fcard .score { font-size: 38px; font-weight: 700; margin-bottom: 4px; }
|
||||
.fcard .sub { font-size: 12px; color: #5c6675; }
|
||||
.fcard .bar-wrap { height: 3px; background: #1c2333; border-radius: 2px; margin-top: 12px; }
|
||||
.fcard .bar { height: 100%; border-radius: 2px; transition: width 0.6s; }
|
||||
|
||||
/* Charts */
|
||||
.charts { display: grid; grid-template-columns: 1fr 1fr; gap: 12px; margin-bottom: 24px; }
|
||||
.chart-box { background: #11161e; border: 1px solid #1c2333; border-radius: 8px; padding: 18px 20px; }
|
||||
.chart-box h3 { font-size: 12px; color: #5c6675; letter-spacing: 1.5px; margin-bottom: 14px; }
|
||||
.chart-box canvas { max-height: 260px; }
|
||||
|
||||
/* Expectancy */
|
||||
.exp { background: #11161e; border: 1px solid #1c2333; border-radius: 8px; padding: 18px 20px; }
|
||||
.exp h3 { font-size: 12px; color: #5c6675; letter-spacing: 1.5px; margin-bottom: 14px; }
|
||||
.exp-row { display: flex; gap: 12px; align-items: center; margin-bottom: 14px; }
|
||||
.exp select { background: #0a0e14; color: #c9d1d9; border: 1px solid #1c2333; padding: 6px 12px;
|
||||
border-radius: 4px; font-size: 13px; }
|
||||
.exp button { background: #1c3a5c; color: #58a6ff; border: 1px solid #2d4f7c; padding: 6px 18px;
|
||||
border-radius: 4px; cursor: pointer; font-size: 13px; }
|
||||
.exp button:hover { background: #254d7a; }
|
||||
.exp .suff { font-size: 11px; padding: 3px 10px; border-radius: 3px; }
|
||||
.suff-HIGH { background: rgba(63,185,80,0.12); color: #3fb950; }
|
||||
.suff-MEDIUM { background: rgba(210,153,34,0.12); color: #d29922; }
|
||||
.suff-LOW { background: rgba(248,81,73,0.12); color: #f85149; }
|
||||
.suff-INSUFFICIENT { background: rgba(92,102,117,0.12); color: #5c6675; }
|
||||
table { width: 100%; border-collapse: collapse; font-size: 13px; }
|
||||
th { text-align: left; color: #5c6675; font-weight: 500; padding: 8px 10px; border-bottom: 1px solid #1c2333; }
|
||||
td { padding: 7px 10px; border-bottom: 1px solid #0e1219; color: #8b949e; }
|
||||
td strong { color: #e6edf3; }
|
||||
.exp-summary { margin-top: 14px; font-size: 13px; color: #8b949e; padding: 10px 14px;
|
||||
background: #0d1117; border-radius: 6px; border-left: 3px solid #58a6ff; }
|
||||
.exp-summary strong { color: #e6edf3; }
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<div class="container-fluid py-3 px-4">
|
||||
<div class="app">
|
||||
|
||||
<!-- Header -->
|
||||
<div class="d-flex justify-content-between align-items-center mb-4">
|
||||
<div class="header">
|
||||
<div>
|
||||
<h4 class="mb-0">ChanMacro <span class="text-muted fs-6">市场状态</span></h4>
|
||||
<small class="text-muted" id="update-time"></small>
|
||||
</div>
|
||||
<div>
|
||||
<span class="badge bg-secondary" id="db-status">加载中...</span>
|
||||
<h1><span>Chan</span>Macro</h1>
|
||||
</div>
|
||||
<div class="time"><span class="dot"></span> <span id="update-time">加载中...</span></div>
|
||||
</div>
|
||||
|
||||
<!-- Hero: Regime -->
|
||||
<div class="card p-4 mb-3 text-center">
|
||||
<div class="hero-conf mb-1">当前制度</div>
|
||||
<div class="hero-regime" id="hero-regime">—</div>
|
||||
<div>
|
||||
<span class="badge-regime" id="hero-badge">—</span>
|
||||
<span class="ms-2" style="color:#8b949e">置信度 <strong id="hero-conf" style="color:#e6edf3">—</strong></span>
|
||||
<span class="ms-2" style="color:#8b949e">成熟度 <strong id="hero-maturity" style="color:#e6edf3">—</strong></span>
|
||||
<!-- Regime Hero -->
|
||||
<div class="hero">
|
||||
<div class="hero-card main">
|
||||
<div>
|
||||
<div class="regime-badge" id="hero-badge">—</div>
|
||||
<div class="regime-name" id="hero-regime">—</div>
|
||||
</div>
|
||||
<div style="display:flex; gap:32px; margin-left:auto;">
|
||||
<div class="hero-stat"><div class="val" id="hero-conf">—</div><div class="lbl">置信度</div></div>
|
||||
<div class="hero-stat"><div class="val" id="hero-maturity">—</div><div class="lbl">成熟度</div></div>
|
||||
</div>
|
||||
</div>
|
||||
<div class="hero-card" style="flex:1">
|
||||
<div class="hero-stat"><div class="val" id="hero-ps">—</div><div class="lbl">价格结构</div></div>
|
||||
</div>
|
||||
<div class="hero-card" style="flex:1">
|
||||
<div class="hero-stat"><div class="val" id="hero-br">—</div><div class="lbl">市场广度</div></div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- 4 Factor Cards -->
|
||||
<div class="row g-3 mb-3">
|
||||
<div class="col-md-3">
|
||||
<div class="card p-3 h-100">
|
||||
<div class="factor-label">价格结构</div>
|
||||
<div class="factor-value" id="f-price">—</div>
|
||||
<div class="text-muted small" id="f-price-sub"></div>
|
||||
<div class="narrative mt-1" id="f-price-narr"></div>
|
||||
</div>
|
||||
<div class="grid">
|
||||
<div class="fcard">
|
||||
<div class="title">价格结构 PRICE STRUCTURE</div>
|
||||
<div class="score" id="f-price">—</div>
|
||||
<div class="sub" id="f-price-sub"></div>
|
||||
<div class="bar-wrap"><div class="bar" id="bar-price"></div></div>
|
||||
</div>
|
||||
<div class="col-md-3">
|
||||
<div class="card p-3 h-100">
|
||||
<div class="factor-label">市场广度</div>
|
||||
<div class="factor-value" id="f-breadth">—</div>
|
||||
<div class="text-muted small" id="f-breadth-sub"></div>
|
||||
<div class="narrative mt-1" id="f-breadth-narr"></div>
|
||||
</div>
|
||||
<div class="fcard">
|
||||
<div class="title">市场广度 BREADTH</div>
|
||||
<div class="score" id="f-breadth">—</div>
|
||||
<div class="sub" id="f-breadth-sub"></div>
|
||||
<div class="bar-wrap"><div class="bar" id="bar-breadth"></div></div>
|
||||
</div>
|
||||
<div class="col-md-3">
|
||||
<div class="card p-3 h-100">
|
||||
<div class="factor-label">OI 状态</div>
|
||||
<div class="factor-value fs-4" id="f-oi">—</div>
|
||||
<div class="text-muted small" id="f-oi-sub"></div>
|
||||
<div class="narrative mt-1" id="f-oi-narr"></div>
|
||||
</div>
|
||||
<div class="fcard">
|
||||
<div class="title">持仓状态 OI MATRIX</div>
|
||||
<div class="score" id="f-oi" style="font-size:24px">—</div>
|
||||
<div class="sub" id="f-oi-sub"></div>
|
||||
</div>
|
||||
<div class="col-md-3">
|
||||
<div class="card p-3 h-100">
|
||||
<div class="factor-label">波动率</div>
|
||||
<div class="factor-value" id="f-vol">—</div>
|
||||
<div class="text-muted small" id="f-vol-sub"></div>
|
||||
</div>
|
||||
<div class="fcard">
|
||||
<div class="title">波动率 VOLATILITY</div>
|
||||
<div class="score" id="f-vol">—</div>
|
||||
<div class="sub" id="f-vol-sub"></div>
|
||||
<div class="bar-wrap"><div class="bar" id="bar-vol"></div></div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- Charts Row -->
|
||||
<div class="row g-3 mb-3">
|
||||
<div class="col-md-6">
|
||||
<div class="card p-3">
|
||||
<h6 class="mb-3">制度历史</h6>
|
||||
<canvas id="chart-regime"></canvas>
|
||||
</div>
|
||||
</div>
|
||||
<div class="col-md-6">
|
||||
<div class="card p-3">
|
||||
<h6 class="mb-3">市场广度</h6>
|
||||
<canvas id="chart-breadth"></canvas>
|
||||
</div>
|
||||
</div>
|
||||
<!-- Charts -->
|
||||
<div class="charts">
|
||||
<div class="chart-box"><h3>制度历史 REGIME HISTORY</h3><canvas id="chart-regime"></canvas></div>
|
||||
<div class="chart-box"><h3>市场广度 BREADTH</h3><canvas id="chart-breadth"></canvas></div>
|
||||
</div>
|
||||
|
||||
<!-- Expectancy -->
|
||||
<div class="card p-3">
|
||||
<h6 class="mb-3">信号期望查询</h6>
|
||||
<div class="row g-2 align-items-end">
|
||||
<div class="col-auto">
|
||||
<select class="form-select form-select-sm" id="exp-signal">
|
||||
<option value="B3">B3 (三买)</option><option value="B2">B2 (二买)</option><option value="B1">B1 (一买)</option>
|
||||
<option value="S3">S3 (三卖)</option><option value="S2">S2 (二卖)</option><option value="S1">S1 (一卖)</option>
|
||||
</select>
|
||||
</div>
|
||||
<div class="col-auto">
|
||||
<button class="btn btn-sm btn-primary" onclick="loadExpectancy()">查询</button>
|
||||
</div>
|
||||
<div class="col-auto">
|
||||
<span class="badge bg-secondary" id="exp-sufficiency">—</span>
|
||||
</div>
|
||||
<div class="exp">
|
||||
<h3>信号期望 SIGNAL EXPECTANCY</h3>
|
||||
<div class="exp-row">
|
||||
<select id="exp-signal">
|
||||
<option value="B3">B3 · 三买</option><option value="B2">B2 · 二买</option><option value="B1">B1 · 一买</option>
|
||||
<option value="S3">S3 · 三卖</option><option value="S2">S2 · 二卖</option><option value="S1">S1 · 一卖</option>
|
||||
</select>
|
||||
<button onclick="loadExpectancy()">查询</button>
|
||||
<span class="suff" id="exp-sufficiency">—</span>
|
||||
</div>
|
||||
<div class="table-responsive mt-2">
|
||||
<table class="table table-sm table-dark mb-0" style="--bs-table-bg:#161b22">
|
||||
<thead><tr><th>层级</th><th>样本</th><th>有效样本</th><th>原始胜率</th><th>后验胜率</th><th>平均收益</th></tr></thead>
|
||||
<tbody id="exp-layers"></tbody>
|
||||
</table>
|
||||
</div>
|
||||
<div class="mt-2 text-muted small" id="exp-summary"></div>
|
||||
<table>
|
||||
<thead><tr><th>层级</th><th>样本</th><th>有效样本</th><th>原始胜率</th><th>后验胜率</th><th>平均收益</th></tr></thead>
|
||||
<tbody id="exp-layers"></tbody>
|
||||
</table>
|
||||
<div class="exp-summary" id="exp-summary"></div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
|
||||
@@ -0,0 +1,380 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
twitter_web.py — Twitter 监控账号管理 Web 界面。
|
||||
单文件,零依赖,只用到 Python 标准库。
|
||||
"""
|
||||
|
||||
import json
|
||||
import os
|
||||
import sys
|
||||
import re
|
||||
from datetime import datetime, timezone
|
||||
from http.server import HTTPServer, BaseHTTPRequestHandler
|
||||
from urllib.parse import urlparse, parse_qs
|
||||
|
||||
SCRIPT_DIR = os.path.dirname(os.path.abspath(__file__))
|
||||
WATCHLIST_PATH = os.path.join(SCRIPT_DIR, "twitter_watchlist.json")
|
||||
STATE_PATH = os.path.join(SCRIPT_DIR, "twitter_state.json")
|
||||
|
||||
PORT = int(sys.argv[1]) if len(sys.argv) > 1 else 8010
|
||||
|
||||
|
||||
def extract_username(value: str) -> str:
|
||||
value = value.strip().rstrip("/")
|
||||
if value.startswith("@"):
|
||||
return value[1:]
|
||||
for pattern in [r"(?:twitter\.com|x\.com)/(\w+)(?:/|$)", r"/(\w+)$"]:
|
||||
m = re.search(pattern, value)
|
||||
if m:
|
||||
return m.group(1)
|
||||
if re.match(r"^\w+$", value):
|
||||
return value
|
||||
raise ValueError(f"无法提取用户名: {value}")
|
||||
|
||||
|
||||
def load_json(path):
|
||||
if os.path.exists(path):
|
||||
with open(path) as f:
|
||||
return json.load(f)
|
||||
return {}
|
||||
|
||||
|
||||
def save_json(path, data):
|
||||
with open(path, "w") as f:
|
||||
json.dump(data, f, indent=2, ensure_ascii=False)
|
||||
|
||||
|
||||
def get_watchlist():
|
||||
return load_json(WATCHLIST_PATH).get("users", [])
|
||||
|
||||
|
||||
def save_watchlist(users):
|
||||
save_json(WATCHLIST_PATH, {"users": users})
|
||||
|
||||
|
||||
def get_state():
|
||||
return load_json(STATE_PATH)
|
||||
|
||||
|
||||
HTML = """<!DOCTYPE html>
|
||||
<html lang="zh">
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){dataLayer.push(arguments);}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>Twitter 监控管理</title>
|
||||
<style>
|
||||
:root {
|
||||
--bg: #0d1117; --card: #161b22; --border: #30363d;
|
||||
--text: #c9d1d9; --muted: #8b949e; --accent: #58a6ff;
|
||||
--green: #3fb950; --red: #f85149; --yellow: #d2991d;
|
||||
}
|
||||
* { margin:0; padding:0; box-sizing:border-box; }
|
||||
body { font:14px/1.6 -apple-system,BlinkMacSystemFont,"Segoe UI",sans-serif;
|
||||
background:var(--bg); color:var(--text); padding:24px; max-width:680px; margin:auto; }
|
||||
h1 { font-size:20px; margin-bottom:4px; }
|
||||
.sub { color:var(--muted); font-size:12px; margin-bottom:20px; }
|
||||
.add-bar { display:flex; gap:8px; margin-bottom:20px; }
|
||||
.add-bar input { flex:1; padding:8px 12px; border:1px solid var(--border);
|
||||
border-radius:6px; background:var(--card); color:var(--text); font-size:14px; outline:none; }
|
||||
.add-bar input:focus { border-color:var(--accent); }
|
||||
.add-bar input::placeholder { color:var(--muted); }
|
||||
button { padding:8px 16px; border:none; border-radius:6px; cursor:pointer; font-size:13px;
|
||||
font-weight:500; transition:opacity .15s; }
|
||||
button:hover { opacity:0.85; }
|
||||
.btn-add { background:var(--accent); color:#fff; }
|
||||
.btn-edit, .btn-save { background:var(--yellow); color:#000; }
|
||||
.btn-del { background:var(--red); color:#fff; }
|
||||
.btn-cancel { background:var(--border); color:var(--text); }
|
||||
.account { background:var(--card); border:1px solid var(--border); border-radius:8px;
|
||||
padding:12px 16px; margin-bottom:8px; display:flex; align-items:center; gap:12px; }
|
||||
.account .name { font-weight:600; min-width:160px; }
|
||||
.account .name a { color:var(--accent); text-decoration:none; }
|
||||
.account .name a:hover { text-decoration:underline; }
|
||||
.account .meta { font-size:12px; color:var(--muted); flex:1; }
|
||||
.account .actions { display:flex; gap:6px; flex-shrink:0; }
|
||||
.edit-row { display:flex; gap:6px; align-items:center; width:100%; }
|
||||
.edit-row input { flex:1; padding:6px 10px; border:1px solid var(--accent);
|
||||
border-radius:4px; background:var(--bg); color:var(--text); font-size:13px; outline:none; }
|
||||
.badge { display:inline-block; font-size:11px; padding:2px 8px; border-radius:10px;
|
||||
background:var(--green); color:#000; margin-left:6px; }
|
||||
.empty { text-align:center; padding:60px 20px; color:var(--muted); }
|
||||
.empty p { margin-bottom:8px; }
|
||||
.toast { position:fixed; bottom:20px; right:20px; padding:10px 20px; border-radius:6px;
|
||||
font-size:13px; color:#fff; opacity:0; transition:opacity .3s; z-index:100; }
|
||||
.toast.show { opacity:1; }
|
||||
.toast.ok { background:var(--green); }
|
||||
.toast.err { background:var(--red); }
|
||||
</style>
|
||||
</head>
|
||||
<body>
|
||||
<h1>🐦 Twitter 账号监控</h1>
|
||||
<p class="sub">管理 twitterapi.io 监控账号 · 增删改查</p>
|
||||
|
||||
<div class="add-bar">
|
||||
<input id="urlInput" type="text" placeholder="输入 Twitter/X 链接或用户名..." autofocus>
|
||||
<button class="btn-add" onclick="addAccount()">➕ 添加</button>
|
||||
</div>
|
||||
|
||||
<div id="list"></div>
|
||||
|
||||
<div class="toast" id="toast"></div>
|
||||
|
||||
<script>
|
||||
const API = '/twitter/api/accounts';
|
||||
let editing = null;
|
||||
|
||||
async function api(method, path='', body=null) {
|
||||
const opts = { method, headers:{} };
|
||||
if (body) { opts.headers['Content-Type']='application/json'; opts.body=JSON.stringify(body); }
|
||||
const r = await fetch(API + path, opts);
|
||||
const data = await r.json();
|
||||
if (!r.ok) throw new Error(data.error || '请求失败');
|
||||
return data;
|
||||
}
|
||||
|
||||
function toast(msg, ok=true) {
|
||||
const t = document.getElementById('toast');
|
||||
t.textContent = msg; t.className = 'toast ' + (ok?'ok':'err') + ' show';
|
||||
setTimeout(() => t.classList.remove('show'), 2500);
|
||||
}
|
||||
|
||||
async function load() {
|
||||
const data = await api('GET');
|
||||
const div = document.getElementById('list');
|
||||
if (!data.accounts.length) {
|
||||
div.innerHTML = '<div class="empty"><p>📭 暂无监控账号</p><p style="font-size:12px;color:var(--muted)">在上方输入 Twitter/X 链接或用户名添加</p></div>';
|
||||
return;
|
||||
}
|
||||
div.innerHTML = data.accounts.map(a => `
|
||||
<div class="account" id="row-${a.username}">
|
||||
${editing===a.username ? `
|
||||
<div class="edit-row">
|
||||
<input id="editInput" value="${esc(a.display_name || a.username)}" placeholder="备注名称">
|
||||
<button class="btn-save" onclick="saveEdit('${esc(a.username)}')">保存</button>
|
||||
<button class="btn-cancel" onclick="cancelEdit()">取消</button>
|
||||
</div>
|
||||
` : `
|
||||
<div class="name">
|
||||
<a href="https://x.com/${esc(a.username)}" target="_blank">@${esc(a.username)}</a>
|
||||
${a.display_name && a.display_name !== a.username ? `<span style="color:var(--text)">(${esc(a.display_name)})</span>` : ''}
|
||||
</div>
|
||||
<div class="meta">
|
||||
添加: ${a.added_at?.slice(0,10) || '?'}
|
||||
${a.last_check ? ` · 上次检查: ${a.last_check}` : ''}
|
||||
</div>
|
||||
<div class="actions">
|
||||
<button class="btn-edit" onclick="startEdit('${esc(a.username)}','${esc(a.display_name||a.username)}')">✏️</button>
|
||||
<button class="btn-del" onclick="removeAccount('${esc(a.username)}')">🗑</button>
|
||||
</div>
|
||||
`}
|
||||
</div>
|
||||
`).join('');
|
||||
}
|
||||
|
||||
function esc(s) { return s.replace(/&/g,'&').replace(/"/g,'"').replace(/</g,'<').replace(/>/g,'>').replace(/'/g,'''); }
|
||||
|
||||
async function addAccount() {
|
||||
const inp = document.getElementById('urlInput');
|
||||
const val = inp.value.trim();
|
||||
if (!val) { toast('请输入链接或用户名', false); return; }
|
||||
try {
|
||||
const r = await api('POST', '', {url: val});
|
||||
toast(r.message || '添加成功');
|
||||
inp.value = '';
|
||||
load();
|
||||
} catch(e) { toast(e.message, false); }
|
||||
}
|
||||
|
||||
async function removeAccount(username) {
|
||||
if (!confirm(`确定删除 @${username}?`)) return;
|
||||
try {
|
||||
const r = await api('DELETE', '/' + username);
|
||||
toast(r.message || '已删除');
|
||||
load();
|
||||
} catch(e) { toast(e.message, false); }
|
||||
}
|
||||
|
||||
function startEdit(username, name) {
|
||||
editing = username;
|
||||
load();
|
||||
setTimeout(() => {
|
||||
const inp = document.getElementById('editInput');
|
||||
if (inp) { inp.focus(); inp.select(); }
|
||||
}, 50);
|
||||
}
|
||||
|
||||
function cancelEdit() { editing = null; load(); }
|
||||
|
||||
async function saveEdit(username) {
|
||||
const val = document.getElementById('editInput').value.trim();
|
||||
editing = null;
|
||||
try {
|
||||
const r = await api('PUT', '/' + username, {display_name: val});
|
||||
toast(r.message || '已更新');
|
||||
load();
|
||||
} catch(e) { toast(e.message, false); }
|
||||
}
|
||||
|
||||
document.getElementById('urlInput').addEventListener('keydown', e => {
|
||||
if (e.key === 'Enter') addAccount();
|
||||
});
|
||||
|
||||
load();
|
||||
</script>
|
||||
</body>
|
||||
</html>"""
|
||||
|
||||
|
||||
class Handler(BaseHTTPRequestHandler):
|
||||
def log_message(self, format, *args):
|
||||
pass # silent
|
||||
|
||||
def _send(self, code, body, content_type="application/json"):
|
||||
body = body.encode() if isinstance(body, str) else json.dumps(body, ensure_ascii=False).encode()
|
||||
self.send_response(code)
|
||||
self.send_header("Content-Type", content_type + "; charset=utf-8")
|
||||
self.send_header("Content-Length", str(len(body)))
|
||||
self.send_header("Access-Control-Allow-Origin", "*")
|
||||
self.end_headers()
|
||||
self.wfile.write(body)
|
||||
|
||||
def _json(self, code, data):
|
||||
self._send(code, data)
|
||||
|
||||
def _error(self, code, msg):
|
||||
self._json(code, {"error": msg})
|
||||
|
||||
def do_OPTIONS(self):
|
||||
self.send_response(204)
|
||||
self.send_header("Access-Control-Allow-Origin", "*")
|
||||
self.send_header("Access-Control-Allow-Methods", "GET,POST,PUT,DELETE,OPTIONS")
|
||||
self.send_header("Access-Control-Allow-Headers", "Content-Type")
|
||||
self.end_headers()
|
||||
|
||||
def do_GET(self):
|
||||
path = urlparse(self.path).path
|
||||
if path == "/" or path == "/index.html":
|
||||
self._send(200, HTML, "text/html")
|
||||
return
|
||||
if path.startswith("/api/accounts"):
|
||||
username = path[len("/api/accounts"):].strip("/")
|
||||
if username:
|
||||
# GET /api/accounts/<username> — single account
|
||||
users = get_watchlist()
|
||||
state = get_state()
|
||||
for u in users:
|
||||
if u["username"].lower() == username.lower():
|
||||
entry = dict(u)
|
||||
entry["last_check"] = state.get(u["username"], {}).get("last_check")
|
||||
self._json(200, entry)
|
||||
return
|
||||
self._error(404, "账号不存在")
|
||||
return
|
||||
# GET /api/accounts — list all
|
||||
users = get_watchlist()
|
||||
state = get_state()
|
||||
accounts = []
|
||||
for u in users:
|
||||
entry = dict(u)
|
||||
sc = state.get(u["username"], {})
|
||||
ts = sc.get("last_check")
|
||||
if ts:
|
||||
try:
|
||||
ts = datetime.fromisoformat(ts).strftime("%m-%d %H:%M")
|
||||
except Exception:
|
||||
pass
|
||||
else:
|
||||
ts = "从未"
|
||||
entry["last_check"] = ts
|
||||
accounts.append(entry)
|
||||
self._json(200, {"accounts": accounts})
|
||||
else:
|
||||
self._error(404, "Not Found")
|
||||
|
||||
def do_POST(self):
|
||||
path = urlparse(self.path).path
|
||||
if path != "/api/accounts":
|
||||
self._error(404, "Not Found")
|
||||
return
|
||||
length = int(self.headers.get("Content-Length", 0))
|
||||
body = json.loads(self.rfile.read(length)) if length else {}
|
||||
url = body.get("url", "").strip()
|
||||
if not url:
|
||||
self._error(400, "缺少 url 参数")
|
||||
return
|
||||
try:
|
||||
username = extract_username(url)
|
||||
except ValueError:
|
||||
self._error(400, "无法从输入中提取用户名,请输入 Twitter/X 链接或 @用户名")
|
||||
return
|
||||
|
||||
users = get_watchlist()
|
||||
if any(u["username"].lower() == username.lower() for u in users):
|
||||
self._error(409, f"@{username} 已在监控列表中")
|
||||
return
|
||||
|
||||
display_name = body.get("display_name", "").strip() or username
|
||||
users.append({
|
||||
"username": username,
|
||||
"display_name": display_name,
|
||||
"added_at": datetime.now(timezone.utc).isoformat(),
|
||||
})
|
||||
save_watchlist(users)
|
||||
self._json(201, {"message": f"✅ 已添加 @{username}", "username": username})
|
||||
|
||||
def do_PUT(self):
|
||||
path = urlparse(self.path).path
|
||||
username = path[len("/api/accounts"):].strip("/")
|
||||
if not username:
|
||||
self._error(400, "缺少用户名")
|
||||
return
|
||||
length = int(self.headers.get("Content-Length", 0))
|
||||
body = json.loads(self.rfile.read(length)) if length else {}
|
||||
display_name = body.get("display_name", "").strip()
|
||||
|
||||
users = get_watchlist()
|
||||
for u in users:
|
||||
if u["username"].lower() == username.lower():
|
||||
if display_name:
|
||||
u["display_name"] = display_name
|
||||
save_watchlist(users)
|
||||
self._json(200, {"message": f"✅ @{username} 已更新"})
|
||||
return
|
||||
self._error(404, "账号不存在")
|
||||
|
||||
def do_DELETE(self):
|
||||
path = urlparse(self.path).path
|
||||
username = path[len("/api/accounts"):].strip("/")
|
||||
if not username:
|
||||
self._error(400, "缺少用户名")
|
||||
return
|
||||
users = get_watchlist()
|
||||
before = len(users)
|
||||
users = [u for u in users if u["username"].lower() != username.lower()]
|
||||
if len(users) < before:
|
||||
save_watchlist(users)
|
||||
self._json(200, {"message": f"🗑 已移除 @{username}"})
|
||||
else:
|
||||
self._error(404, "账号不存在")
|
||||
|
||||
|
||||
def main():
|
||||
print(f"🐦 Twitter 监控管理: http://0.0.0.0:{PORT}")
|
||||
server = HTTPServer(("0.0.0.0", PORT), Handler)
|
||||
try:
|
||||
server.serve_forever()
|
||||
except KeyboardInterrupt:
|
||||
print("\n已停止")
|
||||
server.server_close()
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -1,6 +1,14 @@
|
||||
<!DOCTYPE html>
|
||||
<html lang="zh-CN">
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){dataLayer.push(arguments);}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>Chan 数据提供商 - API 文档</title>
|
||||
|
||||
@@ -1,6 +1,14 @@
|
||||
<!DOCTYPE html>
|
||||
<html lang="zh-CN">
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){dataLayer.push(arguments);}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
<meta charset="UTF-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>Chan 数据提供商</title>
|
||||
|
||||
@@ -24,6 +24,12 @@ from fastapi.responses import HTMLResponse
|
||||
import uvicorn
|
||||
from technical.util import resample_to_interval
|
||||
|
||||
from onchain_metrics import (
|
||||
OnchainMetricsManager,
|
||||
api_keys_from_env,
|
||||
create_onchain_router,
|
||||
)
|
||||
|
||||
# docker compose logs --tail=200
|
||||
# docker compose down && docker compose build --no-cache && docker compose up -d
|
||||
|
||||
@@ -1132,6 +1138,13 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
"""构造 FastAPI 应用:lifespan 内同步 initialize 并启动后台拉数;WebSocket 实时推送。"""
|
||||
ws_manager = WebSocketManager()
|
||||
|
||||
# 初始化链上指标模块
|
||||
onchain_manager = OnchainMetricsManager(
|
||||
data_dir=provider.data_dir,
|
||||
api_keys=api_keys_from_env(),
|
||||
)
|
||||
onchain_router = create_onchain_router(onchain_manager)
|
||||
|
||||
def _on_data_update(symbol: str, base_tf: str) -> None:
|
||||
"""后台刷新线程回调:广播基础及衍生周期更新给 WebSocket 订阅者。"""
|
||||
with provider._lock:
|
||||
@@ -1165,9 +1178,12 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
provider.start_background_workers()
|
||||
# 启动 WebSocket 实时监听(asyncio 后台任务)
|
||||
provider.start_watch_tasks()
|
||||
# 启动链上指标后台刷新
|
||||
onchain_manager.start()
|
||||
try:
|
||||
yield
|
||||
finally:
|
||||
onchain_manager.stop()
|
||||
provider.stop()
|
||||
|
||||
app = FastAPI(title="Chan 数据提供商", version="1.0.0", lifespan=lifespan)
|
||||
@@ -1180,6 +1196,9 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
allow_headers=["*"],
|
||||
)
|
||||
|
||||
# 注册链上指标 API
|
||||
app.include_router(onchain_router)
|
||||
|
||||
@app.get("/health")
|
||||
async def health() -> Dict[str, object]:
|
||||
"""存活检查:交易所、交易对、基础/衍生周期、是否已完成冷启动。"""
|
||||
@@ -1291,6 +1310,14 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
return f"""<!DOCTYPE html>
|
||||
<html lang="zh-CN">
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){{dataLayer.push(arguments);}}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>Data Provider — API Manual</title>
|
||||
@@ -1347,6 +1374,7 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
<ul class="nav nav-pills mb-4" id="manualTabs" role="tablist">
|
||||
<li class="nav-item"><button class="nav-link active" data-section="overview">Overview</button></li>
|
||||
<li class="nav-item"><button class="nav-link" data-section="rest">REST API</button></li>
|
||||
<li class="nav-item"><button class="nav-link" data-section="onchain">Onchain Metrics</button></li>
|
||||
<li class="nav-item"><button class="nav-link" data-section="ws">WebSocket</button></li>
|
||||
<li class="nav-item"><button class="nav-link" data-section="timeframes">Timeframes</button></li>
|
||||
<li class="nav-item"><button class="nav-link" data-section="config">Configuration</button></li>
|
||||
@@ -1421,6 +1449,9 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/health</td><td>Structured health check (JSON)</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/timeframes</td><td>List available base & derived timeframes</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/api/candles</td><td>Fetch OHLCV candles</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/api/derivatives</td><td>Funding rate, OI, basis</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/api/onchain/metrics</td><td>On-chain metric time series</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/api/onchain/latest</td><td>Latest on-chain snapshot</td></tr>
|
||||
<tr><td><span class="method-badge method-WS">WS</span></td><td class="endpoint-path">/ws</td><td>Real-time K-line streaming</td></tr>
|
||||
<tr><td><span class="method-badge method-GET">GET</span></td><td class="endpoint-path">/api-manual</td><td>This page</td></tr>
|
||||
</tbody>
|
||||
@@ -1520,6 +1551,117 @@ def create_app(provider: DataProvider) -> FastAPI:
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- ==================== ONCHAIN METRICS ==================== -->
|
||||
<div class="section" id="section-onchain">
|
||||
|
||||
<h5 class="mb-3"><span class="method-badge method-GET">GET</span> /api/onchain/latest</h5>
|
||||
<div class="card endpoint-card">
|
||||
<div class="card-body">
|
||||
<p class="small">获取全部链上指标最新快照。</p>
|
||||
<p><strong>Response</strong> <span class="text-muted small">200 OK</span></p>
|
||||
<pre><code>{{
|
||||
"timestamp": 1704067200000,
|
||||
"btc_netflow": {{
|
||||
"timestamp": 1704067200000,
|
||||
"datetime": "2024-01-01T00:00:00Z",
|
||||
"value": -33361122.64,
|
||||
"sub_value": 1366105343.72,
|
||||
"extra": {{ "inflow_usd": 1366105343.72, "outflow_usd": 1399466466.36 }}
|
||||
}},
|
||||
"stablecoin_supply": {{
|
||||
"timestamp": 1704067200000,
|
||||
"datetime": "2024-01-01T00:00:00Z",
|
||||
"value": 276932114785,
|
||||
"sub_value": 276.93,
|
||||
"extra": {{ "USDT": 184419658464, "USDC": 73258860594, "DAI": 4625356465 }}
|
||||
}},
|
||||
"etf_flow": {{
|
||||
"timestamp": 1704067200000,
|
||||
"datetime": "2024-01-01T00:00:00Z",
|
||||
"value": -296.0,
|
||||
"sub_value": -296000000.0,
|
||||
"extra": {{ "total_million_usd": -296.0, "breakdown": {{ "IBIT": -219.4, "GBTC": -62.8 }} }}
|
||||
}},
|
||||
"mvrv_zscore": {{
|
||||
"timestamp": 1704067200000,
|
||||
"datetime": "2024-01-01T00:00:00Z",
|
||||
"value": -1.3895,
|
||||
"sub_value": 1.1321,
|
||||
"extra": {{ "mvrv_ratio": 1.1321, "rolling_mean": 1.6787, "rolling_stddev": 0.3934 }}
|
||||
}},
|
||||
"sopr": null
|
||||
}}</code></pre>
|
||||
<p class="small text-muted">4个免费指标 (netflow/stablecoin/etf/mvrv) 每5分钟自动刷新。SOPR 需 Glassnode API key。</p>
|
||||
<button class="btn btn-sm btn-outline-secondary" onclick="fetch('/api/onchain/latest').then(r=>r.json()).then(d=>alert(JSON.stringify(d,null,2)))">Try it</button>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<h5 class="mb-3 mt-4"><span class="method-badge method-GET">GET</span> /api/onchain/metrics</h5>
|
||||
<div class="card endpoint-card">
|
||||
<div class="card-body">
|
||||
<p class="small">获取单个链上指标的时间序列。</p>
|
||||
<p><strong>Query Parameters</strong></p>
|
||||
<table class="table table-bordered table-sm">
|
||||
<thead class="table-light"><tr><th>Param</th><th>Type</th><th>Required</th><th>Default</th><th>Description</th></tr></thead>
|
||||
<tbody>
|
||||
<tr><td><code>metric</code></td><td>string</td><td><span class="param-required">Yes</span></td><td>—</td><td>指标: btc_netflow, stablecoin_supply, etf_flow, mvrv_zscore, sopr</td></tr>
|
||||
<tr><td><code>limit</code></td><td>int</td><td>No</td><td><code>100</code></td><td>返回条数上限 (1-5000)</td></tr>
|
||||
<tr><td><code>start</code></td><td>int</td><td>No</td><td>—</td><td>起始时间戳(ms)</td></tr>
|
||||
<tr><td><code>end</code></td><td>int</td><td>No</td><td>—</td><td>结束时间戳(ms)</td></tr>
|
||||
</tbody>
|
||||
</table>
|
||||
<p><strong>cURL Example</strong></p>
|
||||
<pre><code>curl "http://localhost:80/api/onchain/metrics?metric=btc_netflow&limit=10"
|
||||
curl "http://localhost:80/api/onchain/metrics?metric=mvrv_zscore&start=1704067200000&end=1711929600000"</code></pre>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<h5 class="mb-3 mt-4"><span class="method-badge method-GET">GET</span> /api/onchain/available</h5>
|
||||
<div class="card endpoint-card">
|
||||
<div class="card-body">
|
||||
<p class="small">列出所有指标及其数据量、时间范围。</p>
|
||||
<pre><code>{{
|
||||
"btc_netflow": {{"count": 4015, "has_data": true, "first_ts": "2024-01-01T00:00:00Z", "last_ts": "2026-07-01T00:00:00Z"}},
|
||||
"stablecoin_supply": {{"count": 2, "has_data": true, ...}},
|
||||
"etf_flow": {{"count": 12, "has_data": true, ...}},
|
||||
"mvrv_zscore": {{"count": 4036, "has_data": true, ...}},
|
||||
"sopr": {{"count": 0, "has_data": false, ...}}
|
||||
}}</code></pre>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<h5 class="mb-3 mt-4">指标说明</h5>
|
||||
<div class="card">
|
||||
<div class="card-body p-0">
|
||||
<table class="table table-hover mb-0">
|
||||
<thead class="table-light"><tr><th>Metric</th><th>来源</th><th>value 含义</th><th>sub_value</th><th>费用</th></tr></thead>
|
||||
<tbody>
|
||||
<tr><td><code>btc_netflow</code></td><td>CoinMetrics</td><td>净流量 (USD)</td><td>流入 (USD)</td><td><span class="badge bg-success">免费</span></td></tr>
|
||||
<tr><td><code>stablecoin_supply</code></td><td>CoinGecko</td><td>总市值 (USD)</td><td>总市值 (B)</td><td><span class="badge bg-success">免费</span></td></tr>
|
||||
<tr><td><code>etf_flow</code></td><td>Farside</td><td>净流入 (M USD)</td><td>净流入 (USD)</td><td><span class="badge bg-success">免费</span></td></tr>
|
||||
<tr><td><code>mvrv_zscore</code></td><td>CoinMetrics</td><td>Z-Score</td><td>MVRV Ratio</td><td><span class="badge bg-success">免费</span></td></tr>
|
||||
<tr><td><code>sopr</code></td><td>Glassnode</td><td>SOPR 值</td><td>—</td><td><span class="badge bg-warning text-dark">需 key</span></td></tr>
|
||||
</tbody>
|
||||
</table>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<h5 class="mb-3 mt-4">CSV 存储</h5>
|
||||
<div class="card">
|
||||
<div class="card-body">
|
||||
<p class="small">数据落盘在 <code>data/onchain/</code> 下,每个指标一个 CSV 文件:</p>
|
||||
<ul class="small">
|
||||
<li><code>btc_netflow.csv</code> — 交易所净流量 (timestamp, datetime, value, sub_value, extra)</li>
|
||||
<li><code>stablecoin_supply.csv</code> — 稳定币供应</li>
|
||||
<li><code>etf_flow.csv</code> — ETF 净流入</li>
|
||||
<li><code>mvrv_zscore.csv</code> — MVRV Z-Score</li>
|
||||
<li><code>sopr.csv</code> — SOPR</li>
|
||||
</ul>
|
||||
<p class="small text-muted">每5分钟自动刷新并写盘。extra 字段存 JSON (明细/分解数据)。</p>
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
|
||||
<!-- ==================== WEBSOCKET ==================== -->
|
||||
<div class="section" id="section-ws">
|
||||
|
||||
|
||||
@@ -0,0 +1,597 @@
|
||||
"""
|
||||
链上指标数据模块:拉取 BTC 交易所净流量、稳定币供应、ETF 净流入、MVRV Z-Score、SOPR。
|
||||
本地 CSV 落盘 + 内存缓存,参照 derivatives 模式。
|
||||
|
||||
免费数据源:
|
||||
- CoinMetrics Community API: 交易所流入/流出(USD)、MVRV Ratio
|
||||
- CoinGecko: 稳定币市值
|
||||
需要 API key 的指标(可配置):
|
||||
- ETF 净流入: Coinglass / Farside / Glassnode
|
||||
- SOPR: Glassnode / CoinMetrics Pro
|
||||
"""
|
||||
import csv
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import re
|
||||
import threading
|
||||
import time
|
||||
from datetime import datetime, timezone
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
import requests
|
||||
from fastapi import APIRouter, HTTPException, Query
|
||||
|
||||
logger = logging.getLogger("onchain_metrics")
|
||||
|
||||
# ── 常量 ─────────────────────────────────────────────────────────────
|
||||
METRIC_NAMES = [
|
||||
"btc_netflow",
|
||||
"stablecoin_supply",
|
||||
"etf_flow",
|
||||
"mvrv_zscore",
|
||||
"sopr",
|
||||
]
|
||||
REFRESH_INTERVAL = 300 # 后台刷新间隔(秒),免费 API 限频较严
|
||||
COINMETRICS_BASE = "https://community-api.coinmetrics.io/v4"
|
||||
COINGECKO_BASE = "https://api.coingecko.com/api/v3"
|
||||
REQUEST_TIMEOUT = 30
|
||||
HTTP_HEADERS = {"User-Agent": "Mozilla/5.0 (compatible; chan-data-provider/1.0)"}
|
||||
|
||||
|
||||
def to_utc_iso(ts_ms: int) -> str:
|
||||
dt = datetime.fromtimestamp(ts_ms / 1000, tz=timezone.utc)
|
||||
return dt.isoformat().replace("+00:00", "Z")
|
||||
|
||||
|
||||
# ── OnchainMetricsManager ────────────────────────────────────────────
|
||||
class OnchainMetricsManager:
|
||||
"""管理链上指标的拉取、缓存和持久化。每个指标独立存储。"""
|
||||
|
||||
def __init__(self, data_dir: Path, api_keys: Optional[Dict[str, str]] = None) -> None:
|
||||
self.data_dir = data_dir / "onchain"
|
||||
self.data_dir.mkdir(parents=True, exist_ok=True)
|
||||
self.api_keys = api_keys or {}
|
||||
|
||||
# 内存缓存: metric_name -> List[Dict]
|
||||
self._metrics: Dict[str, List[Dict]] = {name: [] for name in METRIC_NAMES}
|
||||
self._lock = threading.RLock()
|
||||
self._stop_event = threading.Event()
|
||||
self._thread: Optional[threading.Thread] = None
|
||||
|
||||
# 加载本地历史
|
||||
for name in METRIC_NAMES:
|
||||
local = self._load_local(name)
|
||||
if local:
|
||||
self._metrics[name] = local
|
||||
logger.info("链上指标 %s 加载本地记录: %d 条", name, len(local))
|
||||
|
||||
# ── CSV 路径 ─────────────────────────────────────────────────────
|
||||
def _path(self, metric: str) -> Path:
|
||||
return self.data_dir / f"{metric}.csv"
|
||||
|
||||
# ── CSV 读写 ─────────────────────────────────────────────────────
|
||||
def _load_local(self, metric: str) -> List[Dict]:
|
||||
path = self._path(metric)
|
||||
if not path.exists():
|
||||
return []
|
||||
records = []
|
||||
with path.open("r", encoding="utf-8", newline="") as fp:
|
||||
for row in csv.DictReader(fp):
|
||||
try:
|
||||
records.append({
|
||||
"timestamp": int(row["timestamp"]),
|
||||
"datetime": row.get("datetime", ""),
|
||||
"value": float(row.get("value", 0)),
|
||||
"sub_value": float(row.get("sub_value", 0)) if row.get("sub_value") else None,
|
||||
"extra": json.loads(row.get("extra", "{}")) if row.get("extra") else {},
|
||||
})
|
||||
except (KeyError, ValueError):
|
||||
continue
|
||||
records.sort(key=lambda r: r["timestamp"])
|
||||
return records
|
||||
|
||||
def _write_local(self, metric: str, records: List[Dict]) -> None:
|
||||
path = self._path(metric)
|
||||
fieldnames = ["timestamp", "datetime", "value", "sub_value", "extra"]
|
||||
tmp = path.with_suffix(".tmp")
|
||||
try:
|
||||
with tmp.open("w", encoding="utf-8", newline="") as fp:
|
||||
w = csv.DictWriter(fp, fieldnames=fieldnames, extrasaction="ignore")
|
||||
w.writeheader()
|
||||
for r in records:
|
||||
row = {
|
||||
"timestamp": r.get("timestamp", 0),
|
||||
"datetime": r.get("datetime", ""),
|
||||
"value": r.get("value", 0),
|
||||
"sub_value": "" if r.get("sub_value") is None else r["sub_value"],
|
||||
"extra": json.dumps(r.get("extra", {}), ensure_ascii=False),
|
||||
}
|
||||
w.writerow(row)
|
||||
tmp.replace(path)
|
||||
logger.info("链上指标 %s 已写入磁盘: %d 条", metric, len(records))
|
||||
except Exception as e:
|
||||
logger.error("链上指标 %s 落盘失败: %s", metric, e)
|
||||
if tmp.exists():
|
||||
tmp.unlink()
|
||||
|
||||
def _merge(self, base: List[Dict], new: List[Dict]) -> List[Dict]:
|
||||
"""按 timestamp 去重合并,新覆盖旧。"""
|
||||
merged = {r["timestamp"]: r for r in base}
|
||||
for r in new:
|
||||
merged[r["timestamp"]] = r
|
||||
return sorted(merged.values(), key=lambda r: r["timestamp"])
|
||||
|
||||
# ── 数据拉取 ─────────────────────────────────────────────────────
|
||||
def _fetch_cm_metrics(
|
||||
self, metrics: str, days: int = 365
|
||||
) -> Optional[List[Dict]]:
|
||||
"""从 CoinMetrics Community API 拉取指标。返回 [{time, metric: value}, ...]"""
|
||||
page_size = min(days, 10000)
|
||||
url = (
|
||||
f"{COINMETRICS_BASE}/timeseries/asset-metrics"
|
||||
f"?assets=btc&metrics={metrics}&frequency=1d&page_size={page_size}"
|
||||
)
|
||||
try:
|
||||
r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS)
|
||||
if r.status_code != 200:
|
||||
logger.warning("CoinMetrics %s 返回 %d: %s", metrics, r.status_code, r.text[:200])
|
||||
return None
|
||||
data = r.json().get("data", [])
|
||||
if not data:
|
||||
return None
|
||||
# 翻页取更多历史
|
||||
all_data = list(data)
|
||||
next_url = r.json().get("next_page_url")
|
||||
pages = 0
|
||||
while next_url and pages < 10:
|
||||
pages += 1
|
||||
time.sleep(0.5)
|
||||
r2 = requests.get(next_url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS)
|
||||
if r2.status_code != 200:
|
||||
break
|
||||
batch = r2.json()
|
||||
all_data.extend(batch.get("data", []))
|
||||
next_url = batch.get("next_page_url")
|
||||
if not next_url or next_url == r.json().get("next_page_url"):
|
||||
break
|
||||
return all_data
|
||||
except Exception as e:
|
||||
logger.error("CoinMetrics %s 拉取失败: %s", metrics, e)
|
||||
return None
|
||||
|
||||
def _fetch_btc_netflow(self) -> List[Dict]:
|
||||
"""BTC 交易所净流量 (USD) = 流入 - 流出。"""
|
||||
data = self._fetch_cm_metrics("FlowInExUSD,FlowOutExUSD", days=365)
|
||||
if not data:
|
||||
return []
|
||||
|
||||
records = []
|
||||
for d in data:
|
||||
ts_str = d.get("time", "")
|
||||
inflow = float(d.get("FlowInExUSD") or 0)
|
||||
outflow = float(d.get("FlowOutExUSD") or 0)
|
||||
netflow = inflow - outflow
|
||||
try:
|
||||
ts = int(datetime.fromisoformat(ts_str.replace("Z", "+00:00")).timestamp() * 1000)
|
||||
except Exception:
|
||||
continue
|
||||
records.append({
|
||||
"timestamp": ts,
|
||||
"datetime": to_utc_iso(ts),
|
||||
"value": round(netflow, 2),
|
||||
"sub_value": round(inflow, 2),
|
||||
"extra": {"inflow_usd": round(inflow, 2), "outflow_usd": round(outflow, 2)},
|
||||
})
|
||||
return sorted(records, key=lambda r: r["timestamp"])
|
||||
|
||||
def _fetch_stablecoin_supply(self) -> List[Dict]:
|
||||
"""稳定币总供应(USDT + USDC + DAI + FDUSD + TUSD 市值之和)。"""
|
||||
try:
|
||||
url = (
|
||||
f"{COINGECKO_BASE}/coins/markets"
|
||||
f"?vs_currency=usd&category=stablecoins&order=market_cap_desc"
|
||||
f"&per_page=5&page=1&sparkline=false"
|
||||
)
|
||||
r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS)
|
||||
if r.status_code != 200:
|
||||
logger.warning("CoinGecko stablecoins 返回 %d", r.status_code)
|
||||
return []
|
||||
coins = r.json()
|
||||
if not isinstance(coins, list):
|
||||
return []
|
||||
|
||||
total_mcap = sum(c.get("market_cap", 0) or 0 for c in coins)
|
||||
now_ms = int(time.time() * 1000)
|
||||
breakdown = {
|
||||
c.get("symbol", "?").upper(): c.get("market_cap", 0) or 0
|
||||
for c in coins[:5]
|
||||
}
|
||||
return [{
|
||||
"timestamp": now_ms,
|
||||
"datetime": to_utc_iso(now_ms),
|
||||
"value": round(total_mcap, 2),
|
||||
"sub_value": round(total_mcap / 1e9, 2), # Billions
|
||||
"extra": breakdown,
|
||||
}]
|
||||
except Exception as e:
|
||||
logger.error("稳定币供应拉取失败: %s", e)
|
||||
return []
|
||||
|
||||
def _fetch_etf_flow(self) -> List[Dict]:
|
||||
"""BTC ETF 净流入/流出。优先从 Farside 免费爬取,Coinglass 为备用。
|
||||
|
||||
数据源优先级:
|
||||
1. Farside (免费, HTML 爬取)
|
||||
2. Coinglass (需 coinglass_key)
|
||||
3. Glassnode (需 glassnode_key)
|
||||
"""
|
||||
# ── 优先:Farside 免费爬取 ──
|
||||
try:
|
||||
records = self._fetch_etf_farside()
|
||||
if records:
|
||||
logger.info("ETF 数据从 Farside 拉取: %d 条", len(records))
|
||||
return records
|
||||
except Exception as e:
|
||||
logger.warning("Farside ETF 爬取失败: %s", e)
|
||||
|
||||
# ── 备用:Coinglass ──
|
||||
cg_key = self.api_keys.get("coinglass_key")
|
||||
if cg_key:
|
||||
try:
|
||||
url = "https://open-api-v3.coinglass.com/api/bitcoin/etf/net-inflow?interval=30"
|
||||
r = requests.get(url, timeout=REQUEST_TIMEOUT, headers={
|
||||
**HTTP_HEADERS, "coinglassSecret": cg_key,
|
||||
})
|
||||
if r.status_code == 200:
|
||||
data = r.json()
|
||||
records = []
|
||||
if isinstance(data, dict) and "data" in data:
|
||||
for item in data["data"]:
|
||||
ts = int(item.get("date", 0)) * 1000 if item.get("date") else 0
|
||||
if ts:
|
||||
records.append({
|
||||
"timestamp": ts,
|
||||
"datetime": to_utc_iso(ts),
|
||||
"value": float(item.get("netInflow", 0)),
|
||||
"sub_value": None,
|
||||
"extra": item,
|
||||
})
|
||||
return records
|
||||
logger.warning("Coinglass ETF 返回 %d", r.status_code)
|
||||
except Exception as e:
|
||||
logger.error("Coinglass ETF 拉取失败: %s", e)
|
||||
|
||||
logger.warning("ETF 数据源均不可用(Farside/Coinglass/Glassnode)")
|
||||
return []
|
||||
|
||||
def _fetch_etf_farside(self) -> List[Dict]:
|
||||
"""从 Farside 网站爬取 BTC ETF 每日净流量。
|
||||
|
||||
https://farside.co.uk/btc/ 页面包含一个 HTML 表格,
|
||||
每行包含日期和各 ETF 的当日流量(百万美元),最后一列为总计。
|
||||
"""
|
||||
url = "https://farside.co.uk/btc/"
|
||||
r = requests.get(url, timeout=REQUEST_TIMEOUT, headers={
|
||||
**HTTP_HEADERS,
|
||||
"Accept": "text/html,application/xhtml+xml,*/*",
|
||||
})
|
||||
if r.status_code != 200:
|
||||
logger.warning("Farside 返回 %d", r.status_code)
|
||||
return []
|
||||
|
||||
html = r.text
|
||||
|
||||
# 查找表格
|
||||
table_match = re.search(r'<table[^>]*>(.*?)</table>', html, re.DOTALL)
|
||||
if not table_match:
|
||||
logger.warning("Farside 页面未找到表格")
|
||||
return []
|
||||
|
||||
rows_html = re.findall(r'<tr[^>]*>(.*?)</tr>', table_match.group(1), re.DOTALL)
|
||||
month_map = {
|
||||
"jan": 1, "feb": 2, "mar": 3, "apr": 4, "may": 5, "jun": 6,
|
||||
"jul": 7, "aug": 8, "sep": 9, "oct": 10, "nov": 11, "dec": 12,
|
||||
}
|
||||
|
||||
records = []
|
||||
for row_html in rows_html:
|
||||
cells = re.findall(r'<t[dh][^>]*>(.*?)</t[dh]>', row_html, re.DOTALL)
|
||||
# 清理 HTML 标签和空白
|
||||
clean = []
|
||||
for c in cells:
|
||||
t = re.sub(r'<[^>]+>', '', c).strip()
|
||||
t = t.replace('\xa0', ' ').replace(' ', ' ').strip()
|
||||
clean.append(t)
|
||||
|
||||
if not clean:
|
||||
continue
|
||||
|
||||
# 第一列应为日期格式 "15 Jun 2026"
|
||||
date_str = clean[0]
|
||||
parts = date_str.split()
|
||||
if len(parts) != 3:
|
||||
continue
|
||||
day_str, mon_str, year_str = parts
|
||||
mon = month_map.get(mon_str.lower()[:3])
|
||||
if mon is None:
|
||||
continue
|
||||
try:
|
||||
day = int(day_str)
|
||||
year = int(year_str)
|
||||
except ValueError:
|
||||
continue
|
||||
|
||||
# 最后一列为总计(可能带括号表示负值)
|
||||
total_str = clean[-1] if len(clean) > 1 else ""
|
||||
if not total_str or total_str in ("", "Total", "-"):
|
||||
continue
|
||||
# 解析 "(123.4)" → -123.4, "123.4" → 123.4
|
||||
total_str = total_str.replace(",", "")
|
||||
is_negative = total_str.startswith("(") and total_str.endswith(")")
|
||||
if is_negative:
|
||||
total_str = total_str[1:-1]
|
||||
try:
|
||||
total_m = float(total_str)
|
||||
except ValueError:
|
||||
continue
|
||||
if is_negative:
|
||||
total_m = -total_m
|
||||
|
||||
# 构建时间戳(UTC 午夜)
|
||||
from datetime import datetime, timezone as tz
|
||||
dt = datetime(year, mon, day, tzinfo=tz.utc)
|
||||
ts = int(dt.timestamp() * 1000)
|
||||
|
||||
# 分解各 ETF 明细
|
||||
etf_breakdown = {}
|
||||
if len(clean) > 2:
|
||||
etf_names = ["IBIT", "FBTC", "BITB", "ARKB", "BTCO", "EZBC",
|
||||
"BRRR", "HODL", "BTCW", "MSBT", "GBTC", "BTC"]
|
||||
for i, name in enumerate(etf_names):
|
||||
idx = i + 1
|
||||
if idx < len(clean) - 1:
|
||||
val_str = clean[idx].replace(",", "").replace("(", "").replace(")", "")
|
||||
neg = clean[idx].startswith("(")
|
||||
try:
|
||||
val = float(val_str)
|
||||
etf_breakdown[name] = -val if neg else val
|
||||
except ValueError:
|
||||
pass
|
||||
|
||||
records.append({
|
||||
"timestamp": ts,
|
||||
"datetime": dt.isoformat().replace("+00:00", "Z"),
|
||||
"value": round(total_m, 2), # 净流入 (百万 USD)
|
||||
"sub_value": round(total_m * 1_000_000, 2), # 净流入 (USD)
|
||||
"extra": {"total_million_usd": round(total_m, 2), "breakdown": etf_breakdown},
|
||||
})
|
||||
|
||||
return sorted(records, key=lambda r: r["timestamp"])
|
||||
|
||||
def _fetch_mvrv_zscore(self) -> List[Dict]:
|
||||
"""MVRV Z-Score = (当前 MVRV - 滚动均值) / 滚动标准差。
|
||||
|
||||
从 CoinMetrics 拉取 CapMVRVCur(免费),计算 365 天滚动 Z-Score。
|
||||
"""
|
||||
data = self._fetch_cm_metrics("CapMVRVCur", days=400)
|
||||
if not data:
|
||||
return []
|
||||
|
||||
# 解析时间序列
|
||||
mvrv_series = []
|
||||
for d in data:
|
||||
ts_str = d.get("time", "")
|
||||
mvrv_val = d.get("CapMVRVCur")
|
||||
if mvrv_val is None:
|
||||
continue
|
||||
try:
|
||||
ts = int(datetime.fromisoformat(ts_str.replace("Z", "+00:00")).timestamp() * 1000)
|
||||
except Exception:
|
||||
continue
|
||||
mvrv_series.append((ts, float(mvrv_val)))
|
||||
mvrv_series.sort(key=lambda x: x[0])
|
||||
|
||||
# 计算滚动 Z-Score (365 天窗口,即约 365 个数据点)
|
||||
window = 365
|
||||
records = []
|
||||
values = []
|
||||
timestamps = []
|
||||
for ts, val in mvrv_series:
|
||||
values.append(val)
|
||||
timestamps.append(ts)
|
||||
if len(values) < window:
|
||||
continue
|
||||
window_vals = values[-window:]
|
||||
mean = sum(window_vals) / window
|
||||
variance = sum((v - mean) ** 2 for v in window_vals) / window
|
||||
stddev = variance ** 0.5
|
||||
zscore = (val - mean) / stddev if stddev > 0 else 0.0
|
||||
records.append({
|
||||
"timestamp": ts,
|
||||
"datetime": to_utc_iso(ts),
|
||||
"value": round(zscore, 4),
|
||||
"sub_value": round(val, 4),
|
||||
"extra": {"mvrv_ratio": round(val, 4), "rolling_mean": round(mean, 4), "rolling_stddev": round(stddev, 4)},
|
||||
})
|
||||
return records
|
||||
|
||||
def _fetch_sopr(self) -> List[Dict]:
|
||||
"""SOPR (Spent Output Profit Ratio)。需要 API key。
|
||||
|
||||
支持:
|
||||
- glassnode_key: Glassnode v1/metrics/indicators/sopr
|
||||
"""
|
||||
gn_key = self.api_keys.get("glassnode_key")
|
||||
if gn_key:
|
||||
try:
|
||||
url = (
|
||||
f"https://api.glassnode.com/v1/metrics/indicators/sopr"
|
||||
f"?a=btc&f=json&api_key={gn_key}"
|
||||
)
|
||||
r = requests.get(url, timeout=REQUEST_TIMEOUT, headers=HTTP_HEADERS)
|
||||
if r.status_code == 200:
|
||||
data = r.json()
|
||||
records = []
|
||||
for item in data if isinstance(data, list) else []:
|
||||
ts = int(item.get("t", 0)) * 1000
|
||||
if ts:
|
||||
records.append({
|
||||
"timestamp": ts,
|
||||
"datetime": to_utc_iso(ts),
|
||||
"value": float(item.get("v", 0)),
|
||||
"sub_value": None,
|
||||
"extra": {"raw": item},
|
||||
})
|
||||
return records
|
||||
logger.warning("Glassnode SOPR 返回 %d", r.status_code)
|
||||
except Exception as e:
|
||||
logger.error("Glassnode SOPR 拉取失败: %s", e)
|
||||
|
||||
logger.warning(
|
||||
"SOPR 数据未配置 API key。请设置环境变量 GLASSNODE_API_KEY"
|
||||
)
|
||||
return []
|
||||
|
||||
# ── 刷新入口 ─────────────────────────────────────────────────────
|
||||
def refresh_all(self) -> None:
|
||||
"""拉取全部 5 个指标,合并到内存并写盘。"""
|
||||
fetchers = {
|
||||
"btc_netflow": self._fetch_btc_netflow,
|
||||
"stablecoin_supply": self._fetch_stablecoin_supply,
|
||||
"etf_flow": self._fetch_etf_flow,
|
||||
"mvrv_zscore": self._fetch_mvrv_zscore,
|
||||
"sopr": self._fetch_sopr,
|
||||
}
|
||||
|
||||
for name, fetcher in fetchers.items():
|
||||
try:
|
||||
new_records = fetcher()
|
||||
if not new_records:
|
||||
continue
|
||||
with self._lock:
|
||||
base = list(self._metrics.get(name, []))
|
||||
merged = self._merge(base, new_records)
|
||||
self._metrics[name] = merged
|
||||
self._write_local(name, merged)
|
||||
logger.info("链上指标 %s 刷新: +%d 条, 总计 %d 条",
|
||||
name, len(new_records), len(merged))
|
||||
except Exception:
|
||||
logger.exception("链上指标 %s 刷新异常", name)
|
||||
|
||||
# ── 后台线程 ─────────────────────────────────────────────────────
|
||||
def _refresh_loop(self) -> None:
|
||||
"""后台线程:启动立即拉取一次,之后每 REFRESH_INTERVAL 秒刷新。"""
|
||||
logger.info("链上指标后台线程启动,间隔 %ds", REFRESH_INTERVAL)
|
||||
# 启动立即拉取
|
||||
self.refresh_all()
|
||||
while not self._stop_event.wait(REFRESH_INTERVAL):
|
||||
self.refresh_all()
|
||||
logger.info("链上指标后台线程已退出")
|
||||
|
||||
def start(self) -> None:
|
||||
self._stop_event.clear()
|
||||
self._thread = threading.Thread(
|
||||
target=self._refresh_loop, name="onchain-loop", daemon=True,
|
||||
)
|
||||
self._thread.start()
|
||||
logger.info("链上指标模块已启动")
|
||||
|
||||
def stop(self) -> None:
|
||||
self._stop_event.set()
|
||||
if self._thread:
|
||||
self._thread.join(timeout=5)
|
||||
|
||||
# ── 查询接口 ─────────────────────────────────────────────────────
|
||||
def get_metric(
|
||||
self,
|
||||
metric: str,
|
||||
limit: int = 100,
|
||||
start_ms: Optional[int] = None,
|
||||
end_ms: Optional[int] = None,
|
||||
) -> List[Dict]:
|
||||
if metric not in METRIC_NAMES:
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail=f"不支持的指标: {metric}。可选: {METRIC_NAMES}",
|
||||
)
|
||||
with self._lock:
|
||||
records = list(self._metrics.get(metric, []))
|
||||
if start_ms is not None:
|
||||
records = [r for r in records if r["timestamp"] >= start_ms]
|
||||
if end_ms is not None:
|
||||
records = [r for r in records if r["timestamp"] <= end_ms]
|
||||
return records[-limit:] if limit else records
|
||||
|
||||
def get_latest(self) -> Dict[str, Any]:
|
||||
"""获取所有指标的最新值。"""
|
||||
result = {"timestamp": int(time.time() * 1000)}
|
||||
for name in METRIC_NAMES:
|
||||
with self._lock:
|
||||
records = self._metrics.get(name, [])
|
||||
if records:
|
||||
latest = records[-1]
|
||||
result[name] = {
|
||||
"timestamp": latest["timestamp"],
|
||||
"datetime": latest["datetime"],
|
||||
"value": latest["value"],
|
||||
"sub_value": latest.get("sub_value"),
|
||||
"extra": latest.get("extra", {}),
|
||||
}
|
||||
else:
|
||||
result[name] = None
|
||||
return result
|
||||
|
||||
|
||||
# ── FastAPI Router ────────────────────────────────────────────────────
|
||||
def create_onchain_router(manager: OnchainMetricsManager) -> APIRouter:
|
||||
router = APIRouter(prefix="/api/onchain", tags=["onchain"])
|
||||
|
||||
@router.get("/metrics")
|
||||
async def get_metrics(
|
||||
metric: str = Query(..., description=f"指标名称: {', '.join(METRIC_NAMES)}"),
|
||||
limit: int = Query(100, ge=1, le=5000, description="返回条数上限"),
|
||||
start: Optional[int] = Query(None, description="起始时间戳(ms)"),
|
||||
end: Optional[int] = Query(None, description="结束时间戳(ms)"),
|
||||
):
|
||||
"""获取单个链上指标的时间序列。"""
|
||||
data = manager.get_metric(metric, limit=limit, start_ms=start, end_ms=end)
|
||||
return {"metric": metric, "count": len(data), "data": data}
|
||||
|
||||
@router.get("/latest")
|
||||
async def get_latest():
|
||||
"""获取全部 5 个指标的最新快照。"""
|
||||
return manager.get_latest()
|
||||
|
||||
@router.get("/available")
|
||||
async def get_available():
|
||||
"""列出可用指标及当前数据量。"""
|
||||
info = {}
|
||||
for name in METRIC_NAMES:
|
||||
records = manager.get_metric(name, limit=0)
|
||||
info[name] = {
|
||||
"count": len(records),
|
||||
"has_data": len(records) > 0,
|
||||
"first_ts": records[0]["datetime"] if records else None,
|
||||
"last_ts": records[-1]["datetime"] if records else None,
|
||||
}
|
||||
return info
|
||||
|
||||
return router
|
||||
|
||||
|
||||
# ── 环境变量辅助 ─────────────────────────────────────────────────────
|
||||
def api_keys_from_env() -> Dict[str, str]:
|
||||
"""从环境变量读取 API keys。"""
|
||||
keys = {}
|
||||
for env_var, key_name in [
|
||||
("COINGLASS_API_KEY", "coinglass_key"),
|
||||
("GLASSNODE_API_KEY", "glassnode_key"),
|
||||
("FARSIDE_API_KEY", "farside_key"),
|
||||
("COINMETRICS_API_KEY", "coinmetrics_key"),
|
||||
]:
|
||||
val = os.getenv(env_var, "").strip()
|
||||
if val and "***" not in val: # 忽略占位符
|
||||
keys[key_name] = val
|
||||
return keys
|
||||
@@ -3,3 +3,4 @@ fastapi>=0.110.0,<1.0.0
|
||||
uvicorn[standard]>=0.23.0,<1.0.0
|
||||
pandas>=2.0.0,<3.0.0
|
||||
technical==1.5.0
|
||||
requests>=2.28.0
|
||||
|
||||
@@ -1,6 +1,15 @@
|
||||
<!DOCTYPE html>
|
||||
<html lang="zh-CN">
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){dataLayer.push(arguments);}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
<title>缠论 Chart — TradingView</title>
|
||||
|
||||
@@ -1,6 +1,15 @@
|
||||
<!DOCTYPE html>
|
||||
<html>
|
||||
<head>
|
||||
<!-- Google tag (gtag.js) -->
|
||||
<script async src="https://www.googletagmanager.com/gtag/js?id=G-LVVXH3TL04"></script>
|
||||
<script>
|
||||
window.dataLayer = window.dataLayer || [];
|
||||
function gtag(){dataLayer.push(arguments);}
|
||||
gtag('js', new Date());
|
||||
gtag('config', 'G-LVVXH3TL04');
|
||||
</script>
|
||||
|
||||
<title>缠论分析系统</title>
|
||||
<meta charset="utf-8">
|
||||
<meta name="viewport" content="width=device-width, initial-scale=1.0">
|
||||
|
||||
Reference in New Issue
Block a user