"""Telegram alert sender for high-confidence trade cards. Routing: HF Spaces cannot reach api.telegram.org directly (IP blocked). Messages are routed via a Cloudflare Worker relay: HF Space → CLOUDFLARE_WORKER_URL → api.telegram.org → Telegram Env vars required (set in HF Spaces Secrets): CLOUDFLARE_WORKER_URL — your Worker URL, e.g. https://trade-alert.yourname.workers.dev (TELEGRAM_BOT_TOKEN and TELEGRAM_CHAT_ID are set as env vars on the Worker itself, not here) Dedup design (verified by 7-agent panel): - Background loop uses key prefix "bg:" — 4-hour cooldown - Manual /api/scan uses key prefix "manual:" — 60-second cooldown - Separate namespaces prevent background loop from suppressing manual scans - _sent dict is pruned every background loop cycle to prevent memory leak """ from __future__ import annotations import os import time import logging import aiohttp from datetime import datetime, timezone logger = logging.getLogger("telegram_alert") CLOUDFLARE_WORKER_URL = os.environ.get("CLOUDFLARE_WORKER_URL", "") # Keep these for backward compat / fallback reference TELEGRAM_BOT_TOKEN = os.environ.get("TELEGRAM_BOT_TOKEN", "") TELEGRAM_CHAT_ID = os.environ.get("TELEGRAM_CHAT_ID", "") _warned_missing = False # ── Dedup store ─────────────────────────────────────────────────────────────── # Keys: "bg:SYMBOL_DIRECTION" or "manual:SYMBOL_DIRECTION" # Values: epoch float of last alert time _sent: dict[str, float] = {} BG_DEDUP_WINDOW = 4 * 3600 # 4 hours for background loop MANUAL_DEDUP_WINDOW = 60 # 60 seconds for manual scan (prevents rapid-click spam) def _bg_key(card: dict) -> str: return f"bg:{card.get('symbol','')}_{card.get('direction','').upper()}" def _manual_key(card: dict) -> str: return f"manual:{card.get('symbol','')}_{card.get('direction','').upper()}" def _is_suppressed(card: dict, manual: bool) -> bool: key = _manual_key(card) if manual else _bg_key(card) window = MANUAL_DEDUP_WINDOW if manual else BG_DEDUP_WINDOW last = _sent.get(key) return last is not None and (time.time() - last) < window def _mark_sent(card: dict, manual: bool) -> None: key = _manual_key(card) if manual else _bg_key(card) _sent[key] = time.time() def prune_sent_cache() -> None: """Remove entries older than the background dedup window. Call from alert loop.""" cutoff = time.time() - BG_DEDUP_WINDOW stale = [k for k, v in _sent.items() if v < cutoff] for k in stale: del _sent[k] if stale: logger.debug("Pruned %d stale dedup entries", len(stale)) # ── Helpers ─────────────────────────────────────────────────────────────────── def _fmt_price(p) -> str: if p is None: return "—" p = float(p) if p >= 1000: return f"{p:,.2f}" if p >= 1: return f"{p:.4f}" return f"{p:.6f}" def _money(v, decimals: int = 2) -> str: if v is None: return "—" return f"${float(v):,.{decimals}f}" def _rr(card: dict) -> str: try: entry = float(card["close"]) sl = float(card["levels"]["sl"]) tp1 = float(card["levels"]["tp1"]) risk = abs(entry - sl) reward = abs(tp1 - entry) if risk == 0: return "—" return f"{reward / risk:.1f}R" except Exception: return "—" # ── Message builder ─────────────────────────────────────────────────────────── def _tier_line(conf: float) -> str: """Confidence tier prefix (Strategy Fix Plan F1a). Cosmetic label only — the actual send/no-send decision lives in main.py's _passes_alert_gate() (confidence >= 7.0 AND PA grade A/A+, changed 13 Jul 2026). This just tells Utkarsh at a glance where in the confidence range a delivered alert sits. """ if conf >= 9.0: return "💎 TRADE-GRADE" if conf >= 8.8: return "🔭 WATCH ONLY — below auto-trade bar" if conf >= 8.0: return "🟡 STRONG" return "🟠 WATCH" def _grade_line(card: dict) -> str: """PA confluence grade line — required by the alert gate itself now, so surface it on every message (grade is always A or A+ for delivered alerts, since B/C/D never pass _passes_alert_gate).""" grade = (card.get("price_action") or {}).get("confluence_grade") if not grade: return "" return f"📐 Grade {grade}" def _momentum_line(card: dict) -> str: """Entry-momentum line (Strategy Fix Plan F2b). Empty if data unavailable.""" mom = card.get("entry_momentum") or {} m15v = mom.get("m15_roc_pct") if m15v is None: return "" mark = "✓ aligned" if mom.get("aligned") else "⚠ AGAINST direction" sign = "+" if m15v >= 0 else "" return f"⚡ Momentum 15m: {sign}{m15v:.2f}% ({mark})" def _order_flow_line(card: dict) -> str: """Format the taker buy/sell ratio line for Telegram. Returns '' if unavailable.""" taker_ratio = ( card.get("evidence", {}) .get("positioning", {}) .get("taker_ratio") ) if taker_ratio is None: return "" direction = card.get("direction", "").upper() pct = int(round(taker_ratio * 100)) if taker_ratio > 0.55: emoji, label = "🟢", "Buy pressure" note = "confirms LONG" if direction == "LONG" else "conflicts SHORT" elif taker_ratio < 0.45: emoji, label = "🔴", "Sell pressure" note = "confirms SHORT" if direction == "SHORT" else "conflicts LONG" else: emoji, label = "⚪", "Neutral flow" note = "no directional bias" return ( f"{emoji} Order flow " f"{label} " f"{taker_ratio:.2f} " f"({note})" ) def build_message(card: dict, tv_sym: str) -> str: sym = card.get("symbol", "") direction = card.get("direction", "").upper() conf = card.get("confidence", 0) lvl = card.get("levels", {}) lev = card.get("leverage", {}) sz = card.get("sizing", {}) dur = card.get("duration", {}) hmm = card.get("market_state", {}).get("hmm", {}) dir_emoji = "🟢" if direction == "LONG" else "🔴" state = hmm.get("state", "—") hmm_conf = hmm.get("confidence") state_emoji = {"BULL": "🐂", "BEAR": "🐻", "RANGING": "↔️"}.get(state, "❓") bingx_url = f"https://bingx.com/en/perpetual/{sym.replace('-', '')}/" tv_url = f"https://www.tradingview.com/chart/?symbol={tv_sym}&interval=15" liq_dist = lev.get("liq_dist_pct") liq_str = f"{float(liq_dist):.1f}%" if liq_dist is not None else "—" hmm_line = ( f"{state_emoji} HMM {state} ({float(hmm_conf):.0%} conf)" if hmm_conf is not None else f"{state_emoji} HMM {state}" ) flow_line = _order_flow_line(card) tier_line = _tier_line(float(conf)) grade_line = _grade_line(card) momentum_line = _momentum_line(card) tp_dyn = lvl.get("tp_dyn") lines = [ *([tier_line] if tier_line else []), f"{dir_emoji} {sym} {direction} · {conf:.1f}/10", *([grade_line] if grade_line else []), "", f"📍 Entry {_fmt_price(card.get('close'))}", f"🛑 SL {_fmt_price(lvl.get('sl'))}", f"🎯 TP1 {_fmt_price(lvl.get('tp1'))} R/R {_rr(card)}", f" TP2 {_fmt_price(lvl.get('tp2'))}", f" TP3 {_fmt_price(lvl.get('tp3'))}", *([f" TP-dyn {_fmt_price(tp_dyn)} (0.7R)"] if tp_dyn is not None else []), *([momentum_line] if momentum_line else []), "", f"⚙️ Leverage {lev.get('user_leverage', '—')}x Margin {_money(lev.get('margin_to_post'))}", f"💸 Loss@SL {_money(lev.get('loss_at_stop'))} Risk {_money(sz.get('risk_usd'))}", f"💥 Liq dist {liq_str}", f"📦 Notional {_money(sz.get('notional'), 0)}", "", hmm_line, *([flow_line] if flow_line else []), f"⏱ Duration {dur.get('label', '—')}", "", f'📊 Chart (TradingView) 🔗 BingX', f"{datetime.now(timezone.utc).strftime('%H:%M UTC · %d %b')}", ] return "\n".join(lines) # ── Public API ──────────────────────────────────────────────────────────────── async def send_alert(card: dict, tv_sym: str, manual: bool = False) -> bool: """Send a Telegram alert for one card. Args: card: trade card dict from score_symbol() tv_sym: TradingView symbol string (e.g. "BINGX:AGLDUSDT.P") manual: True when triggered by user's /api/scan click (uses 60s cooldown) False when triggered by background loop (uses 4h cooldown) Returns True if delivered, False otherwise. Never raises. """ global _warned_missing if not CLOUDFLARE_WORKER_URL: if not _warned_missing: logger.warning( "Telegram alerts disabled — set CLOUDFLARE_WORKER_URL in HF Spaces Secrets." ) _warned_missing = True return False if _is_suppressed(card, manual): window = MANUAL_DEDUP_WINDOW if manual else BG_DEDUP_WINDOW logger.debug( "Dedup hit (%s): %s %s suppressed (window=%ds)", "manual" if manual else "bg", card.get("symbol"), card.get("direction"), window, ) return False if not CLOUDFLARE_WORKER_URL: logger.warning("CLOUDFLARE_WORKER_URL not set in HF Secrets — alerts disabled.") return False text = build_message(card, tv_sym) payload = {"text": text, "parse_mode": "HTML"} try: async with aiohttp.ClientSession() as session: async with session.post( CLOUDFLARE_WORKER_URL, json=payload, timeout=aiohttp.ClientTimeout(total=10) ) as resp: if resp.status == 200: _mark_sent(card, manual) logger.info( "Telegram alert sent via CF Worker [%s]: %s %s conf=%.1f", "manual" if manual else "bg", card.get("symbol"), card.get("direction"), float(card.get("confidence", 0)), ) return True body = await resp.text() logger.error( "CF Worker error %d for %s: %s", resp.status, card.get("symbol"), body[:300], ) except Exception as exc: logger.error( "CF Worker request failed for %s: %s", card.get("symbol"), exc, exc_info=True, ) return False async def send_raw_message(text: str) -> bool: """Send a plain text/HTML message via the CF Worker relay. No dedup, no card parsing — used for heartbeats and system status messages. Returns True if delivered, False otherwise. Never raises. """ if not CLOUDFLARE_WORKER_URL: return False payload = {"text": text, "parse_mode": "HTML"} try: async with aiohttp.ClientSession() as session: async with session.post( CLOUDFLARE_WORKER_URL, json=payload, timeout=aiohttp.ClientTimeout(total=10) ) as resp: if resp.status == 200: logger.info("Raw message sent via CF Worker") return True logger.error("CF Worker raw-message error %d", resp.status) except Exception as exc: logger.error("CF Worker raw-message failed: %s", exc) return False