#!/usr/bin/env python3 """ main.py - 缠论买卖点监控主程序。 每整分钟: 1. 从 Binance 拉取最新 1m K 线 2. 跑完整缠论管线 3. 检测新出现的买卖点(BSP) 4. 推送到 Telegram(同一 BSP 只推一次) """ import asyncio import logging import sys import os import time from datetime import datetime, timezone, timedelta sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) from fetcher import fetch_ohlcv from engine import ChanEngine from notify import send_bsp_alert, BOT_TOKEN, CHAT_ID logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(levelname)s] %(name)s: %(message)s", ) logger = logging.getLogger("bsp_monitor") def _bsp_stable_key(bsp) -> str: """生成 BSP 的稳定唯一键(不依赖笔边界的微小变化)。""" return f"{bsp.type}_{bsp.klc.end_time}" class BSPMonitor: def __init__(self): self._first_run = True self._last_df_ts = None self._known_bsp_keys: set = set() # 已见过的 BSP 键(含已推送和历史的) async def tick(self): """单次 tick。""" tick_start = time.monotonic() logger.info("── tick 开始 ──") # 1. 拉取 K 线 try: df = fetch_ohlcv() except Exception as e: logger.error(f"拉取 K 线失败: {e}") return if df.empty: logger.warning("DataFrame 为空,跳过") return latest_ts = df.iloc[-1]["timestamp"] if self._last_df_ts and latest_ts <= self._last_df_ts: logger.info("无新 K 线,跳过") return self._last_df_ts = latest_ts # 2. 运行缠论管线 try: engine = ChanEngine(df) except Exception as e: logger.error(f"缠论计算失败: {e}", exc_info=True) return # 3. 检测新 BSP(用 stable key 去重) current_bsps = engine.bsp_list current_keys = {_bsp_stable_key(b) for b in current_bsps} if self._first_run: self._first_run = False self._known_bsp_keys = current_keys confirmed = [b for b in engine.bi_list if b.is_sure] logger.info( f"首次运行完成 — {len(confirmed)} 笔已确认, " f"{len(current_bsps)} 个买卖点 (不推送历史)" ) elapsed = (time.monotonic() - tick_start) * 1000 logger.info(f"── tick 结束 ({elapsed:.0f}ms) ──") return new_keys = current_keys - self._known_bsp_keys self._known_bsp_keys |= current_keys # 永增,BSP 一旦见过就不会再"新" if not new_keys: logger.debug("无新买卖点") elapsed = (time.monotonic() - tick_start) * 1000 logger.info(f"── tick 结束 ({elapsed:.0f}ms) ──") return logger.info(f"检测到 {len(new_keys)} 个新买卖点") # 4. 推送新 BSP for bsp in current_bsps: key = _bsp_stable_key(bsp) if key not in new_keys: continue msg = engine.format_bsp_detail(bsp) # HTML 转义(Telegram parse_mode=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", "") ok = send_bsp_alert(msg, bsp_key=key) if ok: logger.info(f"✅ 推送: {key}") else: logger.warning(f"❌ 推送失败: {key}") elapsed = (time.monotonic() - tick_start) * 1000 logger.info(f"── tick 结束 ({elapsed:.0f}ms) ──") async def run(self): logger.info("=" * 50) logger.info("bsp_monitor 启动 — BTC/USDT 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) if __name__ == "__main__": monitor = BSPMonitor() try: asyncio.run(monitor.run()) except KeyboardInterrupt: logger.info("收到中断信号,退出")