File size: 20,704 Bytes
2e658e7
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
"""Analysis-first Telegram webhook adapter for the existing Futures pipeline.

This module deliberately contains no order path. It translates authorized
Telegram updates into calls to ``run_futures_cycle(..., execute=False)`` and
returns compact Bot API responses. Direct delivery is best-effort; webhook
responses, a configured proxy, or an authenticated relay can be used without
bringing Telegram networking into the Futures engine.
"""
from __future__ import annotations

import hashlib
import hmac
import html
import json
import os
import re
import secrets
import threading
import time
from pathlib import Path
from typing import Any

import httpx
from fastapi import APIRouter, Request
from fastapi.responses import JSONResponse

from trading.futures_execution import get_futures_account, get_futures_positions
from trading.trade_cycle import run_futures_cycle
from trading.symbols import normalize_symbol
from tools.futures_dashboard_api import _load_symbols
from trading.dual_datasource_client import DS4_BASE, DS2_BASE

router = APIRouter()
_MAX_WEBHOOK_BYTES = 256 * 1024
_STORE_PATH = Path(os.environ.get("TELEGRAM_STATE_PATH", "/opt/data/telegram_state.json"))
_LOCK = threading.Lock()
_RATE: dict[str, list[float]] = {}
_NONCES: dict[str, dict[str, Any]] = {}
_LAST: dict[str, Any] = {"received": None, "reply": None, "error": None}


def _enabled() -> bool:
    return os.environ.get("TELEGRAM_ENABLED", "false").lower() in {"1", "true", "yes", "on"}


def _allowed_ids() -> set[str]:
    return {x.strip() for x in os.environ.get("TELEGRAM_ALLOWED_USER_IDS", "").split(",") if x.strip()}


def _load_store() -> dict:
    try:
        data = json.loads(_STORE_PATH.read_text(encoding="utf-8"))
        return data if isinstance(data, dict) else {}
    except (OSError, ValueError, TypeError):
        return {}


def _save_store(data: dict) -> None:
    try:
        _STORE_PATH.parent.mkdir(parents=True, exist_ok=True)
        tmp = _STORE_PATH.with_suffix(".tmp")
        tmp.write_text(json.dumps(data, ensure_ascii=False), encoding="utf-8")
        tmp.replace(_STORE_PATH)
    except OSError:
        pass


def _user(user_id: str) -> dict:
    with _LOCK:
        data = _load_store()
        users = data.setdefault("users", {})
        user = users.setdefault(user_id, {"watchlist": [], "settings": {"alerts": False, "direction": "both", "min_score": 0, "cooldown": 900}})
        _save_store(data)
        return user


def _authorized(user_id: str) -> bool:
    if not user_id:
        return False
    try:
        persisted = _load_store().get("owner_user_id")
    except Exception:
        persisted = None
    return user_id in _allowed_ids() or str(persisted or "") == user_id


def _claim_owner(user_id: str, secret: str, *, private_chat: bool) -> bool:
    """Atomically consume the one-time bootstrap secret."""
    if not private_chat or not user_id or not secret:
        return False
    expected = os.environ.get("TELEGRAM_BOOTSTRAP_SECRET", "").strip()
    if not expected:
        return False
    with _LOCK:
        data = _load_store()
        if data.get("bootstrap_consumed") or data.get("owner_user_id"):
            return False
        if not hmac.compare_digest(secret, expected):
            return False
        data["owner_user_id"] = user_id
        data["bootstrap_consumed"] = True
        _save_store(data)
        # Prevent reuse in this process; persistent bootstrap_consumed protects
        # restarts. The HF Secret can be removed separately by an operator.
        os.environ.pop("TELEGRAM_BOOTSTRAP_SECRET", None)
        return True


def _rate_ok(user_id: str) -> bool:
    now = time.time(); window = 60.0
    values = [x for x in _RATE.get(user_id, []) if now - x < window]
    limit = max(1, int(os.environ.get("TELEGRAM_COMMAND_RATE_LIMIT", "10")))
    if len(values) >= limit:
        _RATE[user_id] = values
        return False
    values.append(now); _RATE[user_id] = values
    return True


def _esc(value: Any) -> str:
    return html.escape(str(value if value not in (None, "") else "—"), quote=False)


def _keyboard(user_id: str = "") -> dict:
    nonce = _nonce(user_id)
    def cb(name: str) -> str: return f"menu:{name}|{nonce}"
    return {"inline_keyboard": [
        [{"text": "Promising Coins", "callback_data": cb("scan")}, {"text": "Analyze Symbol", "callback_data": cb("analyze")}],
        [{"text": "Current Signals", "callback_data": cb("signals")}, {"text": "Market Overview", "callback_data": cb("market")}],
        [{"text": "Watchlist", "callback_data": cb("watchlist")}, {"text": "Alerts", "callback_data": cb("alerts")}],
        [{"text": "Paper Account", "callback_data": cb("account")}, {"text": "Datasource Health", "callback_data": cb("sources")}],
        [{"text": "Settings", "callback_data": cb("settings")}],
    ]}


def _response(chat_id: Any, text: str, *, keyboard: dict | None = None) -> dict:
    result = {"method": "sendMessage", "chat_id": chat_id, "text": text[:4096], "parse_mode": "HTML"}
    if keyboard:
        result["reply_markup"] = keyboard
    return result


def _nonce(user_id: str) -> str:
    value = secrets.token_urlsafe(12)
    _NONCES[value] = {"expires": time.time() + 300, "user": user_id}
    return value


async def _direct_send(method: str, payload: dict) -> bool:
    token = os.environ.get("TELEGRAM_BOT_TOKEN", "").strip()
    if not token:
        return False
    proxy = os.environ.get("TELEGRAM_PROXY_URL", "").strip() or None
    url = f"https://api.telegram.org/bot{token}/{method}"
    try:
        async with httpx.AsyncClient(proxy=proxy, timeout=httpx.Timeout(5.0, connect=3.0)) as client:
            response = await client.post(url, json=payload)
            return response.status_code < 400
    except Exception:
        return False


async def _relay_send(method: str, payload: dict) -> bool:
    relay = os.environ.get("TELEGRAM_RELAY_URL", "").strip()
    secret = os.environ.get("TELEGRAM_RELAY_SECRET", "").strip()
    if not relay or not secret:
        return False
    body = json.dumps({"method": method, "payload": payload}, separators=(",", ":")).encode()
    signature = hmac.new(secret.encode(), body, hashlib.sha256).hexdigest()
    try:
        async with httpx.AsyncClient(timeout=5.0) as client:
            response = await client.post(relay, content=body, headers={"X-Relay-Signature": signature, "Content-Type": "application/json"})
            return response.status_code < 400
    except Exception:
        return False


async def deliver(method: str, payload: dict) -> bool:
    return await _relay_send(method, payload) or await _direct_send(method, payload)


async def _scan(user_id: str) -> str:
    items, _ = await _load_symbols()
    max_symbols = min(300, max(1, int(os.environ.get("TELEGRAM_SCAN_MAX_SYMBOLS", "300"))))
    shortlist_size = max(1, int(os.environ.get("TELEGRAM_SCAN_SHORTLIST_SIZE", "20")))
    result_count = max(1, int(os.environ.get("TELEGRAM_SCAN_RESULT_COUNT", "10")))
    candidates = [x for x in items[:max_symbols] if x.get("futuresVerified")]
    # Stage 1 is intentionally cheap and deterministic: verified contracts are
    # preferred, then catalog rank when supplied. No advisory calls occur here.
    candidates.sort(key=lambda x: (x.get("rank") is None, x.get("rank") or 999999, x.get("symbol", "")))
    import asyncio
    sem = asyncio.Semaphore(max(1, int(os.environ.get("TELEGRAM_SCAN_MAX_CONCURRENCY", "4"))))
    async def one(item: dict) -> dict:
        async with sem:
            try:
                plan = await run_futures_cycle(item["symbol"], risk_profile="moderate", execute=False, include_external_context=False)
                return {"symbol": item["symbol"], "decision": plan.get("decision", "NO_TRADE"), "score": plan.get("score"), "verified": True, "reason": (plan.get("rejection_reasons") or plan.get("core_reasons") or [""])[0]}
            except Exception as exc:
                return {"symbol": item["symbol"], "decision": "NO_TRADE", "verified": True, "reason": "analysis unavailable"}
    rows = await asyncio.gather(*(one(x) for x in candidates[:shortlist_size]))
    rows.sort(key=lambda x: (x.get("score") is None, -(float(x.get("score") or 0))))
    lines = ["<b>Promising Coins</b>", f"Evaluated {len(candidates[:max_symbols])} verified contracts; shortlist {len(rows)}."]
    for row in rows[:result_count]:
        lines.append(f"{_esc(row['symbol'])} · <b>{_esc(row['decision'])}</b> · score {_esc(row.get('score'))}")
    return "\n".join(lines)


async def _analysis(symbol: str) -> str:
    try:
        norm = normalize_symbol(symbol).ds4
        items, _ = await _load_symbols()
        verified = next((x for x in items if x.get("symbol") == norm and x.get("futuresVerified")), None)
        if not verified:
            return f"<b>{_esc(norm)}</b> · Market-only / unverified\nDecision: <b>NO_TRADE</b>\nPaper analysis requires a verified Futures contract."
        plan = await run_futures_cycle(norm, risk_profile="moderate", execute=False, include_external_context=True)
        advisory = plan.get("external_advisory") or {}
        lines = [f"<b>{_esc(norm)}</b> · {'Verified Futures' if verified else 'Market-only / unverified'}",
                 f"Decision: <b>{_esc(plan.get('decision'))}</b> · Score: {_esc(plan.get('score'))}",
                 f"Entry {_esc(plan.get('entry'))} · SL {_esc(plan.get('stop_loss'))} · TP {_esc(plan.get('take_profit'))}",
                 f"Leverage {_esc(plan.get('effective_leverage') or plan.get('requested_leverage'))} · Quantity {_esc(plan.get('quantity'))}",
                 f"Risk approval: {_esc(plan.get('risk_approved'))}"]
        reasons = list(dict.fromkeys((plan.get("rejection_reasons") or []) + (plan.get("core_reasons") or [])))
        if reasons: lines.append("Reasons: " + "; ".join(_esc(x) for x in reasons[:3]))
        if advisory: lines.append(f"Advisory: {_esc(advisory.get('provider'))} · {_esc(advisory.get('model'))} · {_esc(advisory.get('market_bias'))}\n{_esc(advisory.get('summary'))}")
        return "\n".join(lines)[:4096]
    except Exception:
        return "Analysis unavailable; deterministic safety state was preserved."


async def _sources() -> str:
    async with httpx.AsyncClient(timeout=4.0) as client:
        async def probe(url: str) -> str:
            try:
                response = await client.get(url); return "ok" if response.status_code < 400 else f"degraded ({response.status_code})"
            except Exception: return "unavailable"
        ds4, ds2 = await __import__("asyncio").gather(probe(f"{DS4_BASE}/api/short-hunter/health"), probe(f"{DS2_BASE}/api/market"))
    return f"<b>Datasource Health</b>\nDS4 authoritative: {_esc(ds4)}\nDS2 complementary: {_esc(ds2)}\nBinance: optional market fallback (regional HTTP 451 may apply)."


async def _command(user_id: str, chat_id: Any, text: str) -> dict:
    parts = text.strip().split(maxsplit=1); command = parts[0].lower().split("@", 1)[0] if parts else "/help"; arg = parts[1].strip() if len(parts) > 1 else ""
    if command in {"/start", "/help", "/menu"}:
        return _response(chat_id, "<b>Hermes Futures Desk</b>\nAnalysis-only Telegram access. Paper/Testnet/Live execution is not available here.", keyboard=_keyboard(user_id))
    if command in {"/scan", "/top"}:
        return _response(chat_id, await _scan(user_id), keyboard=_keyboard(user_id))
    if command in {"/analyze", "/signal", "/market"}:
        return _response(chat_id, await _analysis(arg or "BTCUSDT"), keyboard=_keyboard(user_id))
    user = _user(user_id)
    if command == "/watch":
        try: symbol = normalize_symbol(arg).ds4
        except ValueError: return _response(chat_id, "Use a valid symbol such as BTCUSDT.")
        if symbol not in user["watchlist"]: user["watchlist"].append(symbol); _user_save(user_id, user)
        return _response(chat_id, f"Watching {_esc(symbol)}.")
    if command == "/unwatch":
        try: symbol = normalize_symbol(arg).ds4 if arg else ""
        except ValueError: return _response(chat_id, "Use a valid symbol such as BTCUSDT.")
        user["watchlist"] = [x for x in user["watchlist"] if x != symbol]; _user_save(user_id, user); return _response(chat_id, "Watchlist updated.")
    if command == "/watchlist": return _response(chat_id, "<b>Watchlist</b>\n" + ("\n".join(_esc(x) for x in user["watchlist"]) or "Empty."))
    if command == "/alerts":
        setting = arg.lower().split()
        if setting and setting[0] in {"on", "off"}:
            user["settings"]["alerts"] = setting[0] == "on"; _user_save(user_id, user)
        elif len(setting) >= 1 and setting[0] in {"long", "short", "both"}:
            user["settings"]["direction"] = setting[0]; _user_save(user_id, user)
        elif len(setting) >= 2 and setting[0] == "score":
            try: user["settings"]["min_score"] = float(setting[1]); _user_save(user_id, user)
            except ValueError: pass
        elif len(setting) >= 2 and setting[0] == "cooldown":
            try: user["settings"]["cooldown"] = max(60, int(setting[1])); _user_save(user_id, user)
            except ValueError: pass
        return _response(chat_id, f"Alerts: {'enabled' if user['settings'].get('alerts') else 'disabled'} · direction {user['settings'].get('direction','both')} · min score {user['settings'].get('min_score',0)} · cooldown {user['settings'].get('cooldown',900)}s")
    if command == "/account": return _response(chat_id, _format_account(await get_futures_account()))
    if command == "/positions": return _response(chat_id, _format_positions(await get_futures_positions()))
    if command == "/sources" or command == "/health": return _response(chat_id, await _sources())
    if command in {"/risk", "/settings"}: return _response(chat_id, "Use the authenticated Futures dashboard for this read-only status panel.")
    return _response(chat_id, "Unknown command. Use /help.")


def _user_save(user_id: str, user: dict) -> None:
    with _LOCK:
        data = _load_store(); data.setdefault("users", {})[user_id] = user; _save_store(data)


def _format_account(account: dict) -> str: return f"<b>Paper Account</b>\nMode: {_esc(account.get('mode'))}\nEquity: {_esc(account.get('equity'))}\nDaily PnL: {_esc(account.get('realizedPnlToday'))}"
def _format_positions(data: dict) -> str:
    rows = data.get("positions") or []; return "<b>Paper Positions</b>\n" + ("\n".join(f"{_esc(x.get('symbol'))} · {_esc(x.get('side'))} · size {_esc(x.get('size'))}" for x in rows) or "No open positions.")


async def evaluate_alerts(user_id: str) -> list[dict]:
    """Evaluate a user's watchlist with cooldown and state-change suppression.

    Delivery is optional; failures are returned as unavailable rather than
    retried in a loop. This function is safe to call from an external scheduler.
    """
    user = _user(user_id); settings = user.get("settings", {})
    if not settings.get("alerts"):
        return []
    data = _load_store(); states = data.setdefault("alert_state", {}); now = time.time(); results = []
    for symbol in list(user.get("watchlist", [])):
        try: plan = await run_futures_cycle(symbol, risk_profile="moderate", execute=False, include_external_context=False)
        except Exception: continue
        decision = str(plan.get("decision") or "NO_TRADE"); score = float(plan.get("score") or 0)
        direction = str(settings.get("direction", "both"))
        if decision not in {"LONG", "SHORT"} or (direction == "long" and decision != "LONG") or (direction == "short" and decision != "SHORT") or score < float(settings.get("min_score", 0)):
            continue
        key = f"{user_id}:{symbol}"; previous = states.get(key) or {}
        if previous.get("decision") == decision and now - float(previous.get("sentAt", 0)) < int(settings.get("cooldown", 900)):
            continue
        payload = {"chat_id": user_id, "text": f"<b>Watchlist alert</b>\n{_esc(symbol)} · {decision} · score {_esc(score)}", "parse_mode": "HTML"}
        delivered = await deliver("sendMessage", payload)
        states[key] = {"decision": decision, "score": score, "sentAt": now, "delivery": "ok" if delivered else "unavailable"}
        results.append({"symbol": symbol, "decision": decision, "delivery": "ok" if delivered else "unavailable"})
    _save_store(data)
    return results


@router.post("/api/telegram/webhook")
async def telegram_webhook(request: Request) -> JSONResponse:
    if not _enabled(): return JSONResponse({"ok": False, "detail": "Telegram disabled"}, status_code=404)
    expected = os.environ.get("TELEGRAM_WEBHOOK_SECRET", "").strip()
    supplied = request.headers.get("X-Telegram-Bot-Api-Secret-Token", "")
    if not expected or not hmac.compare_digest(supplied, expected): return JSONResponse({"ok": False}, status_code=401)
    try:
        declared_length = int(request.headers.get("content-length", "0") or 0)
        if declared_length > _MAX_WEBHOOK_BYTES:
            return JSONResponse({"ok": False, "detail": "request too large"}, status_code=413)
        raw_body = await request.body()
        if len(raw_body) > _MAX_WEBHOOK_BYTES:
            return JSONResponse({"ok": False, "detail": "request too large"}, status_code=413)
        update = json.loads(raw_body)
    except Exception: return JSONResponse({"ok": False}, status_code=400)
    message = update.get("message") or update.get("edited_message")
    callback = update.get("callback_query")
    user = (message or callback or {}).get("from") or {}; user_id = str(user.get("id", "")); chat_id = ((message or {}).get("chat") or {}).get("id") or (((callback or {}).get("message") or {}).get("chat") or {}).get("id")
    if not _rate_ok(user_id): return JSONResponse({"ok": False, "detail": "rate limited"}, status_code=429)
    _LAST.update(received=time.time(), error=None)
    text = (message or {}).get("text") or ""
    if text.strip().lower().startswith("/claim "):
        private_chat = str(((message or {}).get("chat") or {}).get("type") or "") == "private"
        claim_secret = text.strip().split(None, 1)[1]
        if _claim_owner(user_id, claim_secret, private_chat=private_chat):
            return JSONResponse(_response(chat_id, "Owner claimed. Use /start to open the Futures menu."))
        return JSONResponse({"ok": False}, status_code=403)
    if not _authorized(user_id): return JSONResponse({"ok": False}, status_code=403)
    if callback:
        raw_callback = str(callback.get("data") or "")
        callback_name, _, callback_nonce = raw_callback.partition("|")
        nonce_info = _NONCES.get(callback_nonce) or {}
        if not callback_nonce or nonce_info.get("expires", 0.0) < time.time() or nonce_info.get("user") != user_id: return JSONResponse({"ok": False, "detail": "expired or unauthorized callback"}, status_code=409)
        _NONCES.pop(callback_nonce, None)
        text = "/" + callback_name.replace("menu:", "")
        # Acknowledge the callback so the Telegram client clears the button's
        # loading spinner. This is best-effort and must not block the reply.
        callback_id = callback.get("id")
        if callback_id:
            await deliver("answerCallbackQuery", {"callback_query_id": callback_id})
    reply = await _command(user_id, chat_id, text)
    _LAST["reply"] = time.time()
    return JSONResponse(reply)


@router.get("/api/telegram/status")
async def telegram_status() -> JSONResponse:
    return JSONResponse({"enabled": _enabled(), "mode": os.environ.get("TELEGRAM_MODE", "webhook"), "webhookEndpointReady": _enabled() and bool(os.environ.get("TELEGRAM_WEBHOOK_SECRET")), "proxyConfigured": bool(os.environ.get("TELEGRAM_PROXY_URL")), "relayConfigured": bool(os.environ.get("TELEGRAM_RELAY_URL") and os.environ.get("TELEGRAM_RELAY_SECRET")), "lastReceivedUpdate": _LAST.get("received"), "lastSuccessfulReply": _LAST.get("reply"), "authorizedUserCount": len(_allowed_ids()), "alertScheduler": "enabled" if os.environ.get("TELEGRAM_ALERTS_ENABLED", "false").lower() == "true" else "disabled", "lastError": _LAST.get("error")})


@router.get("/api/telegram/bootstrap/status")
async def telegram_bootstrap_status(request: Request) -> JSONResponse:
    """Return claim state to the activation workstation without exposing owner data."""
    expected = os.environ.get("TELEGRAM_WEBHOOK_SECRET", "").strip()
    supplied = request.headers.get("X-Telegram-Bot-Api-Secret-Token", "")
    if not expected or not hmac.compare_digest(supplied, expected):
        return JSONResponse({"ok": False}, status_code=401)
    data = _load_store()
    return JSONResponse({"ok": True, "ownerClaimed": bool(data.get("owner_user_id")), "bootstrapConsumed": bool(data.get("bootstrap_consumed"))})