Spaces:
Running
Running
File size: 13,572 Bytes
9aa928f 5f80ef1 9aa928f 350f2ff 9aa928f 1c833d2 9aa928f 0b65b92 9aa928f b0fd5ea 9aa928f 1c833d2 9aa928f 0b65b92 9aa928f 0b65b92 9aa928f 0b65b92 9aa928f 1c833d2 0b65b92 9aa928f 1c833d2 9aa928f 1c833d2 9aa928f 1c833d2 9aa928f 0562341 9aa928f 0562341 9aa928f b0fd5ea 0562341 9aa928f 0562341 9aa928f 0562341 9aa928f 0e3aacf a5fc83f 9aa928f 1c833d2 9aa928f 0b65b92 9aa928f 0b65b92 9aa928f 0b65b92 9aa928f 9b9fc79 9aa928f a5fc83f 9aa928f 0e3aacf 9aa928f dd4c939 9aa928f 0b65b92 9aa928f 1c833d2 9aa928f 350f2ff 9aa928f 350f2ff 9aa928f | 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 | """
lumen_admin.py — HTTP-слой: FastAPI-приложение, секрет-гейты, admin/diag/export.
Связь с bot.py — только через отложенный `import bot` внутри функций.
"""
from __future__ import annotations
import asyncio
import contextlib
import copy
import hmac
import logging
import os
import time
from datetime import datetime
from typing import Any
import aiohttp
from fastapi import FastAPI, Request
from fastapi.responses import JSONResponse
from lumen_state_storage import _serialize_chat_state
log = logging.getLogger("bot")
# Схемы FastAPI выключены намеренно: /docs, /redoc и /openapi.json на публичном
# Space описывали все эндпоинты и их параметры бесплатно (аудит 26.09.2026).
app = FastAPI(docs_url=None, redoc_url=None, openapi_url=None)
def _check_bearer_token(request: Request, expected: str) -> bool:
"""Только Bearer-заголовок: секрет в URL светится в логах/истории/Referer (CWE-598)."""
auth_header = request.headers.get("Authorization", "")
provided = auth_header[7:].strip() if auth_header.lower().startswith("bearer ") else ""
# Сравнение в байтах: str-вариант compare_digest падает TypeError на не-ASCII.
return bool(expected) and bool(provided) and hmac.compare_digest(provided.encode(), expected.encode())
def _check_admin_key(request: Request) -> bool:
import bot
return _check_bearer_token(request, bot.ADMIN_PANEL_KEY)
def _log_denied(request: Request, what: str) -> None:
"""Попытка без верного ключа видна в логах: раньше брутфорс /admin_keys или
/export_state не оставлял следов, а ответ был 200 с телом {"error": ...}
(аудит 26.09.2026)."""
log.warning(
'[admin] Denied %s: no valid Authorization: Bearer <key> from %s',
what, getattr(getattr(request, "client", None), "host", "?"),
)
def _forbidden() -> JSONResponse:
return JSONResponse(
status_code=401,
content={"error": "forbidden — missing or invalid Authorization: Bearer <key> header"},
)
def _redact_secret(value: str) -> str:
"""Отпечаток: последние символы для сверки между рестартами, воспользоваться нельзя."""
if not value:
return "<empty>"
return "…" + value[-6:] if len(value) > 6 else "…" + value
def _check_bot_token_auth(request: Request) -> bool:
"""Как _check_admin_key, но мастер-секрет BOT_TOKEN; только для /admin_keys (иначе круг)."""
import bot
return _check_bearer_token(request, bot.BOT_TOKEN)
@app.get("/")
async def healthcheck() -> dict[str, Any]:
# Только in-memory готовность (bot/client), без сети — healthcheck обязан быть дешёвым.
import bot
ready = bot.bot is not None and bot.client is not None
return {"status": "ok" if ready else "starting", "ready": ready}
@app.get("/admin_keys")
async def get_admin_keys(request: Request) -> Any:
"""Полные секреты только по Bearer BOT_TOKEN (не ADMIN_PANEL_KEY — иначе круг); раньше светились в логах."""
if not _check_bot_token_auth(request):
_log_denied(request, "GET /admin_keys")
return _forbidden()
import bot
return {"webhook_secret": bot.WEBHOOK_SECRET, "admin_panel_key": bot.ADMIN_PANEL_KEY}
@app.get("/webhook_url")
async def get_webhook_url(request: Request) -> Any:
if not _check_admin_key(request):
_log_denied(request, "GET /webhook_url")
return _forbidden()
space_host = os.getenv("SPACE_HOST", "").strip()
if not space_host:
author = os.getenv("SPACE_AUTHOR_NAME", "silverelixir").lower()
repo = os.getenv("SPACE_REPO_NAME", "lumen").lower()
space_host = f"{author}-{repo}.hf.space"
webhook_url = f"https://{space_host}/webhook"
# Токен в URL не отдаём: ссылка с botTOKEN светилась бы в истории браузера
# и логах прокси (AUD-D-001). Токен владелец берёт у @BotFather.
return {
"webhook_url": webhook_url,
"register_url_template": (
"https://api.telegram.org/bot<TOKEN>/setWebhook"
f"?url={webhook_url}"
"&secret_token=<ADMIN_SECRET>"
"&drop_pending_updates=true"
),
"instruction": "Подставь BOT_TOKEN и WEBHOOK_SECRET вручную и вызови ссылку curl, а не браузером",
}
# Апдейт Telegram маленький: всё сверх капа отклоняем до разбора.
WEBHOOK_MAX_BODY_BYTES = 512 * 1024
# Фоновых задач апдейтов держим ограниченно: переполнение просим повторить.
WEBHOOK_MAX_INFLIGHT_TASKS = 100
# Отклонения по секрету/капу гроздьями логируем суммарно, а не по одному.
_WEBHOOK_DENIED_LOG_INTERVAL_SEC = 60.0
_webhook_denied_count = 0
_webhook_denied_last_log = 0.0
def _log_webhook_denied(reason: str) -> None:
"""Троттлинг отказов: каждая строка в SimpleQueue-лог без края, гроздья — одной."""
global _webhook_denied_count, _webhook_denied_last_log
now = time.monotonic()
_webhook_denied_count += 1
if now - _webhook_denied_last_log >= _WEBHOOK_DENIED_LOG_INTERVAL_SEC:
log.warning(
"[webhook] Rejected %d request(s): %s",
_webhook_denied_count, reason,
)
_webhook_denied_count = 0
_webhook_denied_last_log = now
@app.post("/webhook")
async def webhook_handler(request: Request) -> Any:
import bot
import json as _json
token = request.headers.get("X-Telegram-Bot-Api-Secret-Token", "")
# Пустой секрет невалиден, сравнение в байтах держит не-ASCII без TypeError.
if not bot.WEBHOOK_SECRET or not token or not hmac.compare_digest(token.encode(), bot.WEBHOOK_SECRET.encode()):
_log_webhook_denied("invalid secret token")
return {"ok": False}
# Тело апдейта маленькое: отсекаем мусор по заголовку до чтения.
content_length = request.headers.get("Content-Length")
if content_length is not None:
try:
if int(content_length) > WEBHOOK_MAX_BODY_BYTES:
_log_webhook_denied(f"body larger than the {WEBHOOK_MAX_BODY_BYTES} byte cap")
return {"ok": False}
except ValueError:
pass
# Очередь фоновых задач ограничена: переполнение просим повторить позже.
if len(bot._inflight_tasks) >= WEBHOOK_MAX_INFLIGHT_TASKS:
log.warning("[webhook] Too many in-flight updates (%d), asking Telegram to retry", len(bot._inflight_tasks))
return JSONResponse(status_code=503, content={"ok": False, "retry": True})
try:
raw_body = await request.body()
if len(raw_body) > WEBHOOK_MAX_BODY_BYTES:
# Заголовок соврал или его не было: сырые байты сверх капа не разбираем.
_log_webhook_denied(f"body larger than the {WEBHOOK_MAX_BODY_BYTES} byte cap")
return {"ok": False}
body = _json.loads(raw_body)
if bot.bot is not None:
bot._track_inflight_task(asyncio.create_task(bot._process_raw_update(body)))
else:
# 503 вместо 200: пусть Telegram повторит апдейт, а не считает дроп успехом (AUD-E-003).
log.warning("[webhook] Bot not initialized yet, asking Telegram to retry")
return JSONResponse(status_code=503, content={"ok": False, "retry": True})
except Exception as exc:
log.warning("[webhook] Failed to process incoming update: %s", exc)
return {"ok": True}
async def probe_url(session: Any, url: str, *, timeout_sec: float = 6.0, redact: str = "") -> dict[str, Any]:
"""Один GET-зонд: та же логика и тот же формат результата для /diag и /selftest."""
start = time.time()
try:
async with session.get(url, timeout=aiohttp.ClientTimeout(total=timeout_sec)) as resp:
return {"status": resp.status, "elapsed_sec": round(time.time() - start, 2), "ok": True}
except Exception as exc:
elapsed = round(time.time() - start, 2)
exc_str = str(exc) or repr(exc) or type(exc).__name__
if redact:
exc_str = exc_str.replace(redact, "<TOKEN>")
return {"error": exc_str, "elapsed_sec": elapsed, "ok": False}
def _diag_budget_sec() -> float:
"""Лениво через bot._env_number: голый float() давал 500 на мусоре в env."""
import bot
return bot._env_number("DIAG_TOTAL_BUDGET_SEC", 25, min_value=1)
@app.get("/diag")
async def network_diagnostics(request: Request) -> dict[str, Any]:
"""Проверяет исходящую сетевую доступность различных хостов из контейнера.
Открой в браузере чтобы увидеть, что реально заблокировано на исходящих
соединениях из HF Spaces, а что доступно."""
if not _check_admin_key(request):
_log_denied(request, "GET /diag")
return _forbidden()
import bot
bot_token = bot.BOT_TOKEN
targets = {
"telegram_api": "https://api.telegram.org",
"telegram_file_api": "https://api.telegram.org/bot" + (bot_token[:6] if bot_token else "x") + "/getMe",
# Проверяем реально настроенный прокси, а не забытый хардкод: иначе "всё ок" при мёртвом адресе.
"configured_tg_proxy": bot.TELEGRAM_API_BASE_URL + "/bot" + (bot_token[:6] if bot_token else "x") + "/getMe",
"cloudflare_dot_com": "https://www.cloudflare.com",
# Голый workers.dev — ненадёжный сигнал; смотреть на configured_tg_proxy.
"cloudflare_workers_dev_root": "https://workers.dev",
"deno_deploy": "https://deno.com",
# Лишние хостинги убраны: сигнал нужен только по реально используемым (workers/deno).
"google_generic": "https://www.google.com",
"gemini_api": "https://generativelanguage.googleapis.com",
"huggingface": "https://huggingface.co",
"openrouter": "https://openrouter.ai",
"groq_api": "https://api.groq.com",
"tikwm": "https://www.tikwm.com",
"pollinations": "https://image.pollinations.ai",
"upstash": bot.UPSTASH_REDIS_REST_URL if bot.USE_UPSTASH else "https://upstash.com",
}
results: dict[str, Any] = {}
session = await bot._get_http_session()
# Общий бюджет вместо последовательных 14×6с (~84с висящей диагностики):
# зонды идут параллельно, хвост обрезается бюджетом (AUD-F-001).
diag_budget = _diag_budget_sec()
async def _probe(name: str, url: str) -> tuple[str, dict[str, Any]]:
return name, await probe_url(session, url, timeout_sec=6.0, redact=bot_token or "")
tasks = {asyncio.create_task(_probe(name, url)): name for name, url in targets.items()}
done, pending = await asyncio.wait(tasks, timeout=diag_budget)
for task in done:
try:
name, res = task.result()
except Exception as exc:
name, res = tasks[task], {"error": str(exc), "elapsed_sec": diag_budget, "ok": False}
results[name] = res
for task in pending:
task.cancel()
results[tasks[task]] = {"error": "diag budget exceeded", "elapsed_sec": diag_budget, "ok": False}
with contextlib.suppress(Exception):
await asyncio.gather(*pending, return_exceptions=True)
return {"diagnostics": results, "telegram_api_base_configured": bot.TELEGRAM_API_BASE_URL}
@app.get("/export_state")
async def export_state(request: Request) -> dict[str, Any]:
"""Ручной бэкап для cron: эфемерный диск и тир Upstash без реплики; гейт ADMIN_PANEL_KEY."""
if not _check_admin_key(request):
_log_denied(request, "GET /export_state")
return _forbidden()
import bot
# Снимок живых структур: экспорт отдавал ссылки на мутабельные объекты —
# запись между возвратом и сериализацией ответа давала несогласованный бэкап.
chats = {str(cid): _serialize_chat_state(state) for cid, state in list(bot.chat_state.items())}
# Квота вложенная (счётчики per-model) — только deepcopy отцепляет её целиком.
quota = copy.deepcopy(dict(bot.GLOBAL_QUOTA))
return {
"exported_at": datetime.now().isoformat(),
"chats": chats,
"global_quota": quota,
}
|