diff --git a/.gitignore b/.gitignore index ecdd432..01c06f4 100644 --- a/.gitignore +++ b/.gitignore @@ -46,3 +46,6 @@ research/out/*.jsonl.gz research/out/penetration.csv research/out/shadow_*.csv research/out/run_meta_*.json + +# Telegram 凭据。**不要提交** +research/live/deploy/tg.env diff --git a/research/live/deploy/start.sh b/research/live/deploy/start.sh index 7a9f351..7dfc1a1 100755 --- a/research/live/deploy/start.sh +++ b/research/live/deploy/start.sh @@ -113,6 +113,11 @@ docker run -d --name "$NAME" -w /home/hummingbot \ -e PYTHONPATH=/home/hummingbot:/repo/research:/repo/research/live:/repo \ -e SHADOW_SITE="$SHADOW_SITE" \ -e SHADOW_LEAN="$SHADOW_LEAN" \ + -e SHADOW_INCR="${SHADOW_INCR:-1}" \ + -e TG_TOKEN="${TG_TOKEN:-}" \ + -e TG_CHAT="${TG_CHAT:-}" \ + -e TG_NOTIONAL="${TG_NOTIONAL:-200}" \ + -e TG_STALE_S="${TG_STALE_S:-90}" \ -v "$REPO_ROOT:/repo:ro" \ -v "$OUT:/out" \ --entrypoint /opt/conda/envs/hummingbot/bin/python \ diff --git a/research/live/deploy/tg.env.example b/research/live/deploy/tg.env.example new file mode 100644 index 0000000..d6054b3 --- /dev/null +++ b/research/live/deploy/tg.env.example @@ -0,0 +1,13 @@ +# 复制成 tg.env 再填。tg.env 已在 .gitignore 里,不会被提交。 +# +# 拿 token:Telegram 里找 @BotFather → /newbot → 按提示起名 +# 拿 chat id:给你的 bot 随便发一句,然后打开 +# https://api.telegram.org/bot/getUpdates +# 返回的 result[0].message.chat.id 就是 +export TG_TOKEN="" +export TG_CHAT="" +# 每笔名义额(USDT)。小额实盘先用小的,它只影响推送里的下单数量 +export TG_NOTIONAL="200" +# 距「参考价成立」超过这么多秒就标为已失效。参考价是次根开盘价, +# 过了就不是回测那个成交价了 +export TG_STALE_S="90" diff --git a/research/live/shadow_hb.py b/research/live/shadow_hb.py index d344b91..122c13e 100644 --- a/research/live/shadow_hb.py +++ b/research/live/shadow_hb.py @@ -71,6 +71,8 @@ import pandas as pd from lib.shadow_budget import LAG_ALARM_MS, LAG_WINDOW, lag_healthy +import tg_notify + # 站点标识。跨地对比时两台机器的 CSV 要能合起来读,没有这一列就分不清哪行 # 来自哪台。默认取主机名,部署脚本会显式传 SHADOW_SITE(如 sg-hetzner) SITE = os.environ.get("SHADOW_SITE") or socket.gethostname() @@ -646,6 +648,18 @@ class Shadow: self._record_later(sym, kline_ts, h, t_data, t_signal, baseline, atr_pct, lag_ok), f"record {sym}") + # 手工执行的推送。只推过全部滤网的,且 lag 退化时不推——那与 + # 「停开新仓」是同一条规则,不能只在自动化里执行 + if h["pass_all"] and atr_pct: + if not lag_ok: + print(f" [TG] {sym} lag 退化,按停开新仓规则不推", + flush=True) + else: + self._spawn( + tg_notify.push_signal( + sym, h["direction"], float(baseline), + float(atr_pct), kline_ts, t_data - kline_ts), + f"tg {sym}") def _probe_lag(self, sym: str, lag_ms: int) -> tuple[float, bool]: """记一根的到达延迟,返回 (滚动中位数, 该币是否健康)。 diff --git a/research/live/tg_notify.py b/research/live/tg_notify.py new file mode 100644 index 0000000..9eb5086 --- /dev/null +++ b/research/live/tg_notify.py @@ -0,0 +1,144 @@ +"""把过全部滤网的信号推到 Telegram,供手工执行。 + +## 为什么要这个 + +自动执行链一行都还没写(下单 / 持仓状态 / 跨重启持久化 / 对账 / 熔断),而 +过全部滤网的信号只有约 5.3 笔/天——低到人手能接。先手工跑一批,就能在写 +自动化**之前**拿到真实费率档、真实成交价、真实出场行为,让执行链的每个假设 +都有实测对照,而不是写完再发现出场模型不对。 + +## 时效是这条路最大的风险 + +回测的成交价是**信号根的次根开盘价**。信号在收盘瞬间产生,人看到推送、解锁 +手机、下单,几十秒就过去了,成交价已经不是那个开盘价。所以推送里必须带: + + - 参考开盘价(回测口径的成交价) + - 该币的滑点预算(还能容忍多少偏离) + - 距信号产生已过多久 + +并且**超过 TG_STALE_S 就直接标记为已失效**,而不是让人自己判断。宁可漏做, +不要在偏离预算之外入场——那等于在负期望上开仓。 + +## 环境变量 + + TG_TOKEN BotFather 给的 token(缺失则整个推送静默关闭) + TG_CHAT chat id + TG_NOTIONAL 每笔名义额,默认 200 USDT(小额实盘) + TG_STALE_S 超过多少秒算失效,默认 90 +""" +from __future__ import annotations + +import os +import time + +TOKEN = os.environ.get("TG_TOKEN", "") +CHAT = os.environ.get("TG_CHAT", "") +NOTIONAL = float(os.environ.get("TG_NOTIONAL", "200")) +STALE_S = float(os.environ.get("TG_STALE_S", "90")) +ENABLED = bool(TOKEN and CHAT) + +# 出场参数。必须与 step43_fill_aware_budget 的口径一致,否则推的价位和 +# 预算所依据的收益结构不是一回事 +SL_ATR, SCALE_ATR, RUNNER_ATR, MAXB = 2.0, 3.0, 8.0, 48 + +_sent: set = set() + + +def levels(entry: float, atr: float, direction: int) -> dict: + """按 2/3/8 ATR 算出绝对价位。 + + direction=+1 做多、-1 做空。剩余半仓的止损**保持在 2ATR**、不移到成本, + 这是回测参数(RUNNER_STOP=2.0),移了就不是同一个收益结构。 + """ + s = 1.0 if direction > 0 else -1.0 + return {"entry": entry, + "stop": entry - s * SL_ATR * atr, + "scale": entry + s * SCALE_ATR * atr, + "runner": entry + s * RUNNER_ATR * atr} + + +def _fmt(px: float) -> str: + # 币价跨度从 DOGE 的 0.2 到 BTC 的 10 万,固定小数位会把小价币截成 0 + if px >= 1000: + return f"{px:,.1f}" + if px >= 10: + return f"{px:,.3f}" + return f"{px:.6f}" + + +def build(sym: str, direction: int, entry: float, atr_pct: float, + kline_ts: int, lag_ms: float, budget_bp: float, + age_s: float) -> str: + atr = entry * atr_pct + lv = levels(entry, atr, direction) + side = "做多 LONG" if direction > 0 else "做空 SHORT" + qty = NOTIONAL / entry + stale = age_s > STALE_S + + head = f"⛔ 已失效({age_s:.0f}s > {STALE_S:.0f}s)· 不要入场" if stale \ + else f"✅ {side} {sym}" + lines = [ + head, + "", + f"参考成交价 {_fmt(lv['entry'])} ← 回测口径(次根开盘)", + f"数量 {qty:.6f}(名义 {NOTIONAL:,.0f} USDT)", + f"距参考价成立 {age_s:.1f}s(含数据延迟 {lag_ms:.0f}ms,不可压缩)", + "", + f"止损 {_fmt(lv['stop'])} (2 ATR,stop-market)", + f"减半 {_fmt(lv['scale'])} (3 ATR,限价 maker)", + f"目标 {_fmt(lv['runner'])} (8 ATR,限价 maker)", + f"超时 {MAXB} 分钟后市价平(剩余半仓止损仍在 2 ATR,不移成本)", + "", + f"ATR {atr_pct * 1e4:.1f}bp · 滑点预算 {budget_bp:.1f}bp", + f"→ 实际成交偏离参考价超过 {budget_bp:.1f}bp 就不值得做", + ] + if stale: + lines.append("") + lines.append("时效已过:成交价已不是回测那个价,宁可漏做。") + return "\n".join(lines) + + +async def send(text: str) -> None: + """推一条。任何失败都只打日志——推送挂了不能连坐采集。""" + if not ENABLED: + return + try: + import aiohttp + url = f"https://api.telegram.org/bot{TOKEN}/sendMessage" + async with aiohttp.ClientSession() as s: + async with s.post(url, json={"chat_id": CHAT, "text": text}, + timeout=aiohttp.ClientTimeout(total=10)) as r: + if r.status != 200: + print(f" [TG] 推送失败 HTTP {r.status} " + f"{(await r.text())[:200]}", flush=True) + except Exception as e: + print(f" [TG] 推送异常 {type(e).__name__}: {e}", flush=True) + + +async def push_signal(sym: str, direction: int, entry: float, atr_pct: float, + kline_ts: int, lag_ms: float) -> None: + """去重后推一条信号。 + + 去重键取 (币, K线时刻, 方向):同一根被重复处理(补根、池重建后重放)不该 + 推两次,否则人会开两次仓。 + + 时效的起点是 `kline_ts` 而不是信号产生时刻——参考成交价(次根开盘)就是 + 在 kline_ts 那一刻存在的。从信号时刻起算会漏掉数据延迟加计算那 0.5~1.5s, + 而那段是无法压缩的固定成本,必须计入。 + """ + if not ENABLED: + return + key = (sym, int(kline_ts), int(direction)) + if key in _sent: + return + _sent.add(key) + if len(_sent) > 5000: + _sent.clear() + + from lib.shadow_budget import budget_of + b = budget_of(sym) + if not (b == b): # nan:该币当前环境不可做(如 TRX) + print(f" [TG] {sym} 无预算(当前环境不可做),不推", flush=True) + return + age = time.time() - kline_ts / 1000.0 + await send(build(sym, direction, entry, atr_pct, kline_ts, lag_ms, b, age))