#!/usr/bin/env python3 """ main.py - 缠论买卖点监控主程序。 每整分钟: 1. 从 data_provider 拉取所有币对最新 1m K 线 2. 每个币对独立跑缠论管线 3. 检测新笔 → 检查上一笔终点是否 BSP → 推送 4. 检测中枢特征变化 → 推送 """ import asyncio import logging import sys import os import time from dataclasses import dataclass, field from datetime import datetime, timezone, timedelta from typing import Optional sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from fetcher import fetch_ohlcv, get_symbols from engine import ChanEngine from notify import send_bsp_alert, send_telegram_message, BOT_TOKEN, CHAT_ID _PARENT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) if _PARENT not in sys.path: sys.path.insert(0, _PARENT) from ChanPivotMonitor import ChanPivotMonitor from ChanEnum import Chan_BI_DIR logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) logger = logging.getLogger("bsp_monitor") def _short(symbol: str) -> str: """BTC/USDT:USDT → BTCUSDT""" return symbol.split(":")[0].replace("/", "") def _bi_id(bi) -> Optional[tuple]: """笔的稳定标识,基于首K线时间戳。""" if bi.start_klc is None: return None return (bi.start_klc.start_time,) @dataclass class SymbolState: symbol: str pivot_monitor: ChanPivotMonitor = field(default_factory=ChanPivotMonitor) last_bi_id: Optional[tuple] = None # 上次检查过的最后一笔 ID last_df_ts: object = None # 最新已处理的 K 线时间戳 first_run: bool = True class BSPMonitor: def __init__(self): symbols = get_symbols() self._states: dict[str, SymbolState] = { s: SymbolState(symbol=s) for s in symbols } logger.info(f"监控 {len(symbols)} 个币对: {', '.join(_short(s) for s in symbols)}") async def tick(self): tick_start = time.monotonic() logger.info("── tick 开始 ──") for symbol, st in self._states.items(): await self._tick_symbol(symbol, st) elapsed = (time.monotonic() - tick_start) * 1000 logger.info(f"── tick 结束 ({elapsed:.0f}ms) ──") async def _tick_symbol(self, symbol: str, st: SymbolState): name = _short(symbol) # 1. 拉取 K 线 try: df = fetch_ohlcv(symbol) except Exception as e: logger.error(f"[{name}] 拉取失败: {e}") return if df.empty: logger.warning(f"[{name}] DataFrame 为空") return # 2. 检查是否有新 K 线 latest_ts = df.iloc[-1]["timestamp"] if st.last_df_ts and latest_ts <= st.last_df_ts: return st.last_df_ts = latest_ts # 3. 运行缠论管线 try: engine = ChanEngine(df) except Exception as e: logger.error(f"[{name}] 缠论计算失败: {e}", exc_info=True) return # 4. 获取已确认的笔 confirmed = [b for b in engine.bi_list if b.is_sure] if len(confirmed) < 2: return last_confirmed = confirmed[-1] current_bi_id = _bi_id(last_confirmed) if current_bi_id is None: return # 5. 首轮:记录状态,不推送 if st.first_run: st.first_run = False st.last_bi_id = current_bi_id st.pivot_monitor.update(engine.bi_zs_list) logger.info( f"[{name}] 首次完成 — {len(confirmed)} 笔, " f"{len(engine.bsp_list)} BSP (不推送)" ) return # 6. 检测新笔确认 if current_bi_id == st.last_bi_id: return prev_bi_id = st.last_bi_id st.last_bi_id = current_bi_id bi_dir = "⬆️" if last_confirmed.dir == Chan_BI_DIR.UP else "⬇️" logger.info(f"[{name}] 新笔确认 — #{len(confirmed)} " f"{bi_dir} 高度: ${last_confirmed.height:.2f}") # 7. 找到刚被终结的那笔(上一轮的 confirmed[-1]),检查其终点是否为 BSP prev_bi = None for b in engine.bi_list: if b.is_sure and _bi_id(b) == prev_bi_id: prev_bi = b break if prev_bi: bsp = engine.get_bsp_for_bi(prev_bi) else: bsp = None if bsp: key = f"{symbol}_{bsp.type}_{bsp.klc.end_time}" msg = engine.format_bsp_detail(bsp, symbol) msg = _escape_html(msg) if send_bsp_alert(msg, bsp_key=key): logger.info(f"[{name}] ✅ BSP: {key}") # 8. 中枢特征更新 pivot_state = st.pivot_monitor.update(engine.bi_zs_list) if pivot_state: logger.info( f"[{name}] 中枢更新: bi_count={pivot_state['bi_count']} " f"contraction={pivot_state['contraction']:.4f} " f"shift={pivot_state['shift_norm']:+.4f}" ) self._push_pivot(pivot_state, name) def _push_pivot(self, state: dict, name: str): if state["is_sure"]: phase = "✅ 已确认" elif state["bi_count"] > 3: phase = "🔄 延伸中" else: phase = "🆕 刚形成" zs_dir = state["zs_dir"] dir_label = "⬆️ 向上" if "UP" in zs_dir else "⬇️ 向下" c = state["contraction"] if c < 0.85: contraction_note = "收敛(振幅缩小,可能快出方向)" elif c > 1.15: contraction_note = "扩张(振幅放大,波动加剧)" else: contraction_note = "稳定" s = state["shift_norm"] if s > 0.3: shift_note = "重心上移(偏多)" elif s < -0.3: shift_note = "重心下移(偏空)" else: shift_note = "重心居中" msg = ( f"🏠 中枢更新 — {name} 1m\n" f"\n" f"📐 笔数: {state['bi_count']} {dir_label} {phase}\n" f"📏 收敛率: {state['contraction']:.4f} → {contraction_note}\n" f"⚖️ 重心漂移: {state['shift_norm']:+.4f} → {shift_note}\n" f"⏱️ 持续: {state['duration_raw']}K " f"(norm: {state['duration_norm']:.2f})\n" f"📦 区间: {state['zd']:.2f} – {state['zg']:.2f} " f"(gg/dd: {state['gg']:.2f}/{state['dd']:.2f})" ) send_telegram_message(msg) async def run(self): logger.info("=" * 50) logger.info(f"bsp_monitor 启动 — {len(self._states)} 币对 1m 缠论监控") logger.info(f"Telegram: {'已配置' if BOT_TOKEN and CHAT_ID else '⚠️ 未配置'}") logger.info("=" * 50) logger.info("首次运行(初始化,不推送)...") await self.tick() while True: now = datetime.now(timezone.utc) next_minute = now.replace(second=0, microsecond=0) + timedelta(minutes=1) wait_seconds = max(0.1, (next_minute - now).total_seconds()) logger.info(f"等待 {wait_seconds:.0f}s 到 {next_minute.strftime('%H:%M:%S')}UTC") await asyncio.sleep(wait_seconds) try: await self.tick() except Exception as e: logger.error(f"tick 异常: {e}", exc_info=True) await asyncio.sleep(5) def _escape_html(msg: str) -> str: """HTML 转义,保留已有的 / 标签。""" msg = msg.replace("&", "&") msg = msg.replace("", "\x00B\x00").replace("", "\x00/B\x00") msg = msg.replace("", "\x00C\x00").replace("", "\x00/C\x00") msg = msg.replace("<", "<").replace(">", ">") msg = msg.replace("\x00B\x00", "").replace("\x00/B\x00", "") msg = msg.replace("\x00C\x00", "").replace("\x00/C\x00", "") return msg if __name__ == "__main__": monitor = BSPMonitor() try: asyncio.run(monitor.run()) except KeyboardInterrupt: logger.info("收到中断信号,退出")