Files
Chan/bsp_monitor/main.py
T
jackyu66git f0ea6a6065 添加 bsp_monitor: BTC/USDT 1m 缠论买卖点实时监控
- 每整分钟拉取 Binance 永续合约 1m K 线
- 运行完整缠论管线检测买卖点 (BSP)
- 新 BSP 推送到 Telegram
- fix: fetcher 用 limit=1000 替代固定 since,避免 API 500 根限制截断新数据
2026-05-18 08:43:23 +08:00

150 lines
5.0 KiB
Python

#!/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("&", "&amp;")
msg = msg.replace("<b>", "\x00B\x00").replace("</b>", "\x00/B\x00")
msg = msg.replace("<code>", "\x00C\x00").replace("</code>", "\x00/C\x00")
msg = msg.replace("<", "&lt;").replace(">", "&gt;")
msg = msg.replace("\x00B\x00", "<b>").replace("\x00/B\x00", "</b>")
msg = msg.replace("\x00C\x00", "<code>").replace("\x00/C\x00", "</code>")
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("收到中断信号,退出")