""" Lumen — телеграм-бот на Gemini/OpenRouter, webhook-режим. История диалога — 100 сообщений, TikTok через TikWM без водяных знаков, генерация картинок через Pollinations, озвучка через Gemini TTS. """ from __future__ import annotations import asyncio import atexit import base64 import contextlib import hashlib import hmac import html as _html_mod import io import json import logging import logging.handlers import os import queue import re import socket import sys import tempfile import time from pathlib import Path from collections import deque from datetime import date, datetime import mimetypes from typing import Any, TypedDict import urllib.parse import urllib.request as _urllib_request from urllib.parse import urlparse import aiohttp import uvicorn from aiogram import Bot, Dispatcher, F from aiogram.client.session.aiohttp import AiohttpSession from aiogram.client.telegram import TelegramAPIServer from aiogram.exceptions import TelegramEntityTooLarge from aiogram.enums import ChatType, ParseMode from aiogram.filters import Command from aiogram.types import ( BotCommand, BufferedInputFile, FSInputFile, CallbackQuery, InlineKeyboardButton, InlineKeyboardMarkup, InlineQueryResultArticle, InputMediaPhoto, InputMediaVideo, InputTextMessageContent, Message, Update, ) from fastapi import FastAPI, Request from google import genai from google.genai import types from mutagen.id3 import APIC, ID3, TPE1, TIT2 from mutagen.mp3 import MP3 from system_prompt import SYSTEM_PROMPT # логирование LOG_FILE_PATH = Path(os.getenv("BOT_LOG_PATH", "/app/bot.log")) LOG_FILE_PATH.parent.mkdir(parents=True, exist_ok=True) _LOG_QUEUE: queue.SimpleQueue[logging.LogRecord] = queue.SimpleQueue() _LOG_QUEUE_HANDLER = logging.handlers.QueueHandler(_LOG_QUEUE) _LOG_LISTENER: logging.handlers.QueueListener | None = None def _setup_logging() -> logging.Logger: root = logging.getLogger() root.setLevel(logging.INFO) for h in list(root.handlers): root.removeHandler(h) console_handler = logging.StreamHandler(sys.stdout) console_handler.setLevel(logging.INFO) file_handler = logging.handlers.RotatingFileHandler( LOG_FILE_PATH, maxBytes=5 * 1024 * 1024, backupCount=3, encoding="utf-8", delay=True, ) file_handler.setLevel(logging.INFO) fmt = logging.Formatter("%(asctime)s [%(levelname)s] %(name)s: %(message)s") console_handler.setFormatter(fmt) file_handler.setFormatter(fmt) global _LOG_LISTENER _LOG_LISTENER = logging.handlers.QueueListener(_LOG_QUEUE, file_handler, console_handler, respect_handler_level=True) _LOG_LISTENER.start() root.addHandler(_LOG_QUEUE_HANDLER) logging.captureWarnings(True) for name in ("httpx", "google_genai", "aiohttp", "uvicorn.access"): logging.getLogger(name).setLevel(logging.WARNING) return logging.getLogger(__name__) logger = _setup_logging() log = logger def _stop_logging() -> None: global _LOG_LISTENER listener = _LOG_LISTENER _LOG_LISTENER = None if listener is not None: with contextlib.suppress(Exception): listener.stop() atexit.register(_stop_logging) def _load_env_file(path: Path) -> None: if not path.exists(): return try: for raw_line in path.read_text(encoding="utf-8").splitlines(): line = raw_line.strip() if not line or line.startswith("#"): continue if line.startswith("export "): line = line[7:].strip() if "=" not in line: continue key, value = line.split("=", 1) key = key.strip() value = value.strip() if not key or key in os.environ: continue if (value.startswith('"') and value.endswith('"')) or (value.startswith("'") and value.endswith("'")): value = value[1:-1] value = value.replace('\\n', '\n').replace('\\t', '\t') os.environ[key] = value except Exception as exc: log.warning("[setup] Failed to parse .env file %s: %s", path, exc) for _env_path in (Path('/app/.env'), Path('.env')): _load_env_file(_env_path) # переменные окружения и конфиг BOT_TOKEN = os.getenv("BOT_TOKEN", "").strip() if not BOT_TOKEN: BOT_TOKEN = os.getenv("TELEGRAM_TOKEN", "").strip() if not BOT_TOKEN: BOT_TOKEN = os.getenv("TELEGRAM_BOT_TOKEN", "").strip() if not BOT_TOKEN: log.warning("[setup] SYSTEM WARN: BOT_TOKEN is empty! Please verify BOT_TOKEN/TELEGRAM_BOT_TOKEN environment variables in settings or .env.") else: log.info("[setup] BOT_TOKEN configured successfully (length: %d)", len(BOT_TOKEN)) TELEGRAM_API_BASE_URL = os.getenv("TELEGRAM_API_BASE_URL", "https://api.telegram.org").strip().rstrip("/") if TELEGRAM_API_BASE_URL and not TELEGRAM_API_BASE_URL.lower().startswith(("http://", "https://")): # Реальный инцидент: TELEGRAM_API_BASE_URL был задан как голый хост воркера # (например "tg-proxy.egor-kuzko-04.workers.dev") без схемы. aiohttp такой URL # не проглатывает — падает с "Network error" на КАЖДЫЙ вызов (getMe/setWebhook/ # deleteWebhook/setMyCommands и далее вообще все reply/send_message через aiogram), # при этом само сообщение об ошибке невнятное (просто битый URL как текст), не # указывает на реальную причину. Раз уж опечатка в схеме случилась один раз — # молча чинить её тут дешевле, чем снова терять время на диагностику того же самого. log.warning("[setup] SYSTEM WARN: TELEGRAM_API_BASE_URL задан без схемы (%r) — добавляю https:// автоматически.", TELEGRAM_API_BASE_URL) TELEGRAM_API_BASE_URL = "https://" + TELEGRAM_API_BASE_URL log.info("[setup] Using Telegram API Base URL: %s", TELEGRAM_API_BASE_URL) BOT_USERNAME = os.getenv("BOT_USERNAME", "LumenAI_bot").strip().lstrip("@") GEMINI_API_KEY = os.getenv("GEMINI_API_KEY", "").strip() OPENROUTER_API_KEY = (os.getenv("OPENROUTER_API_KEY") or os.getenv("OPENROUTER_KEY") or "").strip() # Найдено при код-ревью: раньше этот дефолт вычислялся ОДИН РАЗ здесь, на старте # модуля, из ЕЩЁ НЕ уточнённого BOT_USERNAME (env-заглушка "LumenAI_bot" по # умолчанию) — до того, как try_setup() ниже реально спрашивает getMe и мог бы # обновить настоящий юзернейм бота. Если владелец не задал BOT_USERNAME в env (или # задал неверно) — заголовок HTTP-Referer к OpenRouter так и оставался бы со # старым/неверным t.me/... адресом весь срок жизни процесса. _OPENROUTER_HTTP_ # REFERER_ENV_SET запоминает, была ли переменная задана ЯВНО владельцем — чтобы # try_setup() ниже пересчитывал referer по свежему юзернейму, только если # владелец сам не переопределил его в env (иначе не перетираем явную настройку). _OPENROUTER_HTTP_REFERER_ENV_SET = bool(os.getenv("OPENROUTER_HTTP_REFERER", "").strip()) OPENROUTER_HTTP_REFERER = os.getenv("OPENROUTER_HTTP_REFERER", f"https://t.me/{BOT_USERNAME}").strip() OPENROUTER_TITLE = os.getenv("OPENROUTER_TITLE", BOT_USERNAME).strip() OPENROUTER_BASE_URL = "https://openrouter.ai/api/v1" # Раньше у OpenRouter был свой отдельный лимит истории (30), меньший, чем у Gemini # (100) — при переключении провайдера (/provider или /model) ощущалось резкое # "обнуление" контекста разговора. Теперь история ОБЩАЯ (см. state["history"] в # get_state) и лимит один и тот же для обоих провайдеров. SHARED_HISTORY_MAX_LEN = 100 OWNER_ID: int | None = None for env_name in ("OWNER_ID", "BOT_OWNER_ID", "ADMIN_ID", "TELEGRAM_OWNER_ID"): val = os.getenv(env_name, "").strip() if val and val.isdigit(): OWNER_ID = int(val) break TELEGRAM_REQUEST_TIMEOUT = float(os.getenv("TELEGRAM_REQUEST_TIMEOUT", "45")) TELEGRAM_AI_TIMEOUT = float(os.getenv("TELEGRAM_AI_TIMEOUT", "45")) TELEGRAM_MEDIA_TIMEOUT = float(os.getenv("TELEGRAM_MEDIA_TIMEOUT", "25")) # Раньше было захардкожено как 15.0 прямо внутри _download_telegram_file_bytes — # несогласованно с остальными таймаутами, которые все конфигурируются через env. TELEGRAM_GET_FILE_TIMEOUT = float(os.getenv("TELEGRAM_GET_FILE_TIMEOUT", "15")) # Если сам прокси перед Telegram (tg-proxy на Deno Deploy) недоступен/приостановлен # (например, исчерпан лимит бесплатного тарифа Deno — ответ вида "503 ... USAGE_EXCEEDED"), # он вместо валидного JSON от Telegram отдаёт HTML/текстовую страницу ошибки. Ни aiogram, # ни наш telegram_api_call не могут её распарсить — падают с JSONDecodeError на КАЖДЫЙ # вызов, а исходящих вызовов в Telegram за секунду может быть десятки (reply, typing-экшен, # get_file и т.д. на каждое входящее сообщение) — без выключателя это лавина одинаковых # WARNING-строк в логах и бессмысленные повторные попытки в мёртвый прокси. См. _tg_call/ # telegram_api_call и _looks_like_proxy_garbage ниже. TG_PROXY_COOLDOWN_SEC — на сколько # секунд отключаем реальные сетевые попытки после первой пойманной такой ошибки. TG_PROXY_COOLDOWN_SEC = float(os.getenv("TG_PROXY_COOLDOWN_SEC", "20")) # В отличие от ask_gemini/ask_openrouter_text (которые ограничены ROUTE_TOTAL_ # BUDGET_SEC на весь маршрут), у стриминга раньше не было НИКАКОГО таймаута вокруг # ожидания следующего куска — генуинно подвисший (не упавший с исключением, а # просто переставший присылать куски) стрим мог держать лок чата (_chat_locks) # бесконечно. Теперь каждое ожидание СЛЕДУЮЩЕГО куска (для ЛЮБОГО провайдера — # Gemini или OpenRouter, см. _run_streaming_reply) ограничено этим таймаутом — # если тишина затянулась дольше него, поднимается TimeoutError, которую функция # и так уже умеет корректно обрабатывать. STREAM_CHUNK_TIMEOUT_SEC = float(os.getenv("STREAM_CHUNK_TIMEOUT_SEC", "30")) # Лимит длины текста для /tts — без него пользователь мог отправить огромный # текст, что вызывало бы очень долгий прогон Gemini TTS + ffmpeg на один запрос. TTS_MAX_CHARS = int(os.getenv("TTS_MAX_CHARS", "800")) _PROCESS_START_MONOTONIC = time.monotonic() # ── Тайминги автоматического маршрутизатора моделей (см. секцию "автоматический # выбор модели" ниже) ── # Раньше (до перехода на роутер) при таймауте/503/500 бот ретраил ОДНУ и ту же # модель 2-3 раза с экспоненциальной задержкой, и только потом переключался на # следующую в цепочке — именно это было причиной ответов по 2+ минуты при # малейшей нестабильности API (см. историю: несколько моделей подряд по # 3 попытки × до 45с каждая). Теперь ретраев ОДНОЙ модели нет вообще: любая # ошибка (таймаут, 429, 503/500, что угодно ещё) — сразу переход к следующей # модели в маршруте. ROUTE_MODEL_TIMEOUT_SEC — сколько ждём ОДНУ попытку одной # модели, прежде чем считать её неудачной и пробовать следующую. ROUTE_MODEL_TIMEOUT_SEC = float(os.getenv("ROUTE_MODEL_TIMEOUT_SEC", "22")) # ROUTE_TOTAL_BUDGET_SEC — общий бюджет времени на ВЕСЬ маршрут одного сообщения, # включая ОБА провайдера (Gemini и OpenRouter), если маршрут предполагает # резервный переход между ними. Без этого потолка каскадный сбой сразу у многих # моделей/провайдеров мог бы растянуть один ответ на несколько минут, всё это # время удерживая лок чата (_chat_locks). При превышении бюджета дальнейшие # попытки прекращаются и пользователь получает честное "сейчас всё перегружено" # вместо тихого зависания. ROUTE_TOTAL_BUDGET_SEC = float(os.getenv("ROUTE_TOTAL_BUDGET_SEC", "40")) TG_MAX_LEN = 4096 # Telegram Bot API ограничивает загрузку файлов, отправляемых ботом (upload, а не # по file_id/URL), 50 МБ — используется в _tiktok_video_candidates/handle_tiktok # ниже, чтобы заранее пропускать заведомо слишком большой вариант качества видео, # не тратя время и трафик на скачивание файла, который Telegram всё равно отклонит. TELEGRAM_BOT_API_UPLOAD_LIMIT_BYTES = 50 * 1024 * 1024 # TikTok официально разрешает до 35 фото/слайдов в одном посте формата "слайдшоу" # (photo mode) — см. справку TikTok. sendMediaGroup при этом жёстко ограничен 10 # элементами ЗА ОДИН вызов — это ограничение Telegram Bot API, а не наше. Чтобы # реально доставить ВЕСЬ пост (а не только первые 10, как было раньше), слайды # делятся на группы по TELEGRAM_MEDIA_GROUP_CHUNK и отправляются несколькими # последовательными вызовами sendMediaGroup — см. handle_tiktok. TIKTOK_SLIDESHOW_MAX_ITEMS = 35 TELEGRAM_MEDIA_GROUP_CHUNK = 10 MAX_CHAT_LIMIT = 5000 PRUNED_CHAT_TARGET = 4500 MAX_CHAT_HISTORY_LEN = 100 MAX_MEDIA_RECENT_IDS = 8 # хранится ОТДЕЛЬНО на каждого пользователя чата (см. recent_media_ids: dict[user_id, deque]) # Простой трекер для rate limiting (5 запросов в 30 сек) user_rate_limits: dict[int, list[float]] = {} def _cleanup_rate_limit_dict() -> None: """user_rate_limits раньше никогда не уменьшался — ключи (user_id) оставались в словаре навсегда, даже когда список timestamp'ов у конкретного пользователя полностью очищался скользящим окном в _handle_message_core. За месяцы работы с большим числом разных пользователей это медленная, но реальная утечка памяти. Вызывается раз в час из фонового цикла в _webhook_startup.""" now = time.time() stale = [uid for uid, ts in user_rate_limits.items() if not ts or now - ts[-1] > 3600] for uid in stale: user_rate_limits.pop(uid, None) # ── Инициализация клиентов внутри цикла обработки событий (решает RuntimeError) ── bot: Bot = None dp = Dispatcher() client: genai.Client = None _telegram_session: aiohttp.ClientSession | None = None _http_session: aiohttp.ClientSession | None = None _chat_locks: dict[int, asyncio.Lock] = {} # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: состояние "выключателя" мёртвого Telegram-прокси # раньше жило как четыре независимых module-level globals (_tg_proxy_down_until/ # _tg_proxy_down_logged_at/_tg_proxy_consecutive_failures/_tg_proxy_garbage_event_count), # мутируемых через `global` из двух разных функций (_tg_call/telegram_api_call) — # такое размазанное состояние сложнее читать и тестировать, чем один объект с # понятными методами. _TelegramProxyCircuitBreaker ниже — чистая инкапсуляция, # поведение (включая формулы cooldown/threshold) не изменилось ни на йоту. # # Выключатель срабатывает по СЧЁТЧИКУ подряд идущих сбоев, а не на первый же # сбой. Раньше ОДНА-единственная заминка прокси (например разовый сетевой глюк # на одной ноде anycast-CDN — Vercel/Cloudflare/Deno все матчат запросы на # множество географически разных нод) полностью глушила ответы бота ВСЕМ чатам # на TG_PROXY_COOLDOWN_SEC секунд — то есть один случайный сбой был неотличим # от реально упавшего прокси. Теперь выключатель включается, только когда # подряд (без единого успеха между ними) накопилось trip_threshold сбоев — # единичные заминки его больше не запускают. TG_PROXY_TRIP_THRESHOLD = int(os.getenv("TG_PROXY_TRIP_THRESHOLD", "3")) class _TelegramProxyCircuitBreaker: """Инкапсулирует состояние выключателя — см. комментарий выше. Используется как единственный module-level инстанс (_tg_proxy_breaker ниже), но методы не трогают globals напрямую, что делает поведение проще проверять.""" def __init__(self, *, cooldown_sec: float, trip_threshold: int) -> None: self.cooldown_sec = cooldown_sec self.trip_threshold = trip_threshold self.down_until: float = 0.0 self.down_logged_at: float = 0.0 self.consecutive_failures: int = 0 # Совокупный (не сбрасывается) счётчик срабатываний "прокси вернул не-JSON" # за время жизни процесса — виден через /stats, чтобы деградацию прокси # можно было заметить прямо из Telegram, а не только копаясь в логах контейнера. self.garbage_event_count: int = 0 def is_down(self, now: float) -> bool: return now < self.down_until def log_still_down_if_due(self, now: float) -> None: """Логирует "прокси всё ещё недоступен" не чаще раза в cooldown_sec — иначе лавина одинаковых WARNING на каждый пропущенный вызов из бэклога.""" if now - self.down_logged_at > self.cooldown_sec: self.down_logged_at = now log.warning("[telegram] Прокси всё ещё недоступен, пропускаю вызовы ещё ~%.0fс.", self.down_until - now) def note_success(self) -> None: """Сбрасывает счётчик подряд идущих сбоев — вызывается на любой исход, который означает, что прокси реально ответил валидным JSON (успех ИЛИ настоящая ошибка Telegram уровня API), т.е. прокси-звено не виновато.""" self.consecutive_failures = 0 def note_failure(self) -> bool: """Увеличивает счётчик подряд идущих сбоев прокси (и общий счётчик для /stats). Возвращает True, если достигнут trip_threshold и пора включать выключатель (см. trip() ниже).""" self.consecutive_failures += 1 self.garbage_event_count += 1 return self.consecutive_failures >= self.trip_threshold def trip(self) -> None: now = time.monotonic() self.down_until = now + self.cooldown_sec self.down_logged_at = now def status_text(self) -> str: """Готовый HTML-фрагмент для /stats — раньше собирался в самой команде по четырём глобалам напрямую, теперь инкапсулирован вместе с состоянием.""" now = time.monotonic() if now < self.down_until: state = f"ВЫКЛЮЧЕН ещё ~{int(self.down_until - now)}с" else: state = "в норме" return ( f"\n\nTelegram-прокси: {state}\n" f"Подряд сбоев сейчас: {self.consecutive_failures}/{self.trip_threshold}, " f"всего за время работы: {self.garbage_event_count}" ) _tg_proxy_breaker = _TelegramProxyCircuitBreaker(cooldown_sec=TG_PROXY_COOLDOWN_SEC, trip_threshold=TG_PROXY_TRIP_THRESHOLD) def _looks_like_proxy_garbage(exc: Exception) -> bool: """Отличает РЕАЛЬНУЮ ошибку Telegram API (валидный JSON вида {"ok": false, ...}) от случая, когда сам HTTP-прокси перед Telegram (tg-proxy на Deno Deploy) вернул не-JSON тело — например, страницу приостановки аккаунта при исчерпанном лимите Deno ("USAGE_EXCEEDED"). Сигнатура именно этого случая — ошибка разбора JSON: Telegram, даже сообщая о СВОИХ ошибках, всегда отвечает валидным JSON, а вот прокси, упавший или приостановленный целиком, отдаёт HTML/plain-text, который ни json.loads, ни aiogram распарсить не могут.""" low = str(exc).lower() cls = exc.__class__.__name__.lower() if "jsondecodeerror" in cls or "jsondecodeerror" in low: return True if "failed to decode" in low or "usage_exceeded" in low: return True # Прокси-хост вообще не принимает соединение (обрыв на уровне TCP/TLS, а не # ответ с ошибкой) — такой же надёжный сигнал "прокси недоступен целиком", как # и не-JSON ответ выше. Реальный инцидент без этой ветки: ClientConnectorError # ("Cannot connect to host ...") не ловился выключателем, и бот на каждое # сообщение заново пытался и подолгу ждал таймаута — вплоть до Duration 226754 ms # на одно сообщение, при том что проблема была одна и та же на протяжении часов. if "clientconnectorerror" in cls or "cannot connect to host" in low: return True return False def _build_telegram_connector(limit: int) -> aiohttp.TCPConnector: """Общая конфигурация TCPConnector для соединений с Telegram API — используется и в _get_telegram_session (aiohttp-сессия для telegram_api_call), и в IPv4AiohttpSession (сессия самого aiogram Bot). Раньше эти два места дублировали один и тот же блок настроек по отдельности — вынесено сюда, чтобы будущая правка (например, очередная донастройка ttl_dns_cache/keepalive_timeout под конкретный прокси-хостинг) не требовала синхронизировать два места вручную. ttl_dns_cache сокращён с 300 до 10 сек: прокси-хостинг (Vercel/Cloudflare/Deno — anycast-CDN с множеством edge-нод по всему миру) мог "залипать" на одной подвисающей/перегруженной ноде на весь TTL DNS-кэша — отсюда сбои шли ПАЧКАМИ (несколько подряд, потом пауза), а не единично-случайно. keepalive_timeout сокращён до 15с вместо ранее пробовавшегося force_close=True: полное отключение keep-alive заставляло КАЖДЫЙ вызов (reply, send_message, get_file, typing-экшен и т.д. — на одно сообщение их несколько) платить полный TCP+TLS handshake — это перебор. Короткого keepalive_timeout достаточно, чтобы не залипать на плохой ноде надолго, но не требовать новый handshake на каждый вызов.""" return aiohttp.TCPConnector( family=socket.AF_INET, limit=limit, ttl_dns_cache=10, keepalive_timeout=15.0, enable_cleanup_closed=True, ) async def _get_telegram_session() -> aiohttp.ClientSession: global _telegram_session if _telegram_session is None or _telegram_session.closed: _telegram_session = aiohttp.ClientSession( connector=_build_telegram_connector(limit=10), timeout=aiohttp.ClientTimeout(total=TELEGRAM_REQUEST_TIMEOUT + 10.0, connect=10.0), ) return _telegram_session async def _get_http_session() -> aiohttp.ClientSession: global _http_session if _http_session is None or _http_session.closed: _http_session = aiohttp.ClientSession( timeout=aiohttp.ClientTimeout(total=60), # limit поднят с 16 до 40 (24 июля 2026, ревизия TikTok-скачивания): # слайдшоу TikTok теперь скачивается целиком за раз через # asyncio.gather (см. handle_tiktok), и TikTok разрешает до 35 слайдов # в одном посте — со старым лимитом 16 часть слайдов ждала бы в # очереди на соединение вместо реально параллельной загрузки. Прочие # потребители этой же сессии (OpenRouter, Pollinations, TikWM-запросы) # используют на порядок меньше одновременных соединений, так что # повышение лимита их не затрагивает. connector=aiohttp.TCPConnector(family=socket.AF_INET, limit=40, ttl_dns_cache=300), ) return _http_session async def _close_sessions() -> None: for s in (_telegram_session, _http_session): if s is not None and not s.closed: await s.close() if bot is not None and hasattr(bot, "session") and bot.session: await bot.session.close() # конвертация markdown в html, утилиты json # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: вся логика конвертации markdown/LaTeX/таблиц/ # маркеров списков в Telegram HTML вынесена в отдельный модуль lumen_formatting.py — # это чистые функции над строками без единой зависимости от Telegram/Gemini/ # OpenRouter/рантайм-состояния бота, самый безопасный кандидат на выделение из # монолитного bot.py. Публичные имена и поведение не изменились — импортируются # напрямую, чтобы `bot._md_to_html(...)`/`bot._scrub_latex(...)` и т.п. продолжали # работать ровно как раньше (в т.ч. для существующих тестов). # # Только `_md_to_html`/`_scrub_latex`/`_normalize_bullet_markers` реально нужны # здесь (используются в коде bot.py или напрямую в тестах через `bot.X`) — # остальные внутренние хелперы (`_TABLE_SEP_RE`, `_LATEX_SYMBOL_MAP` и т.п.) # нужны только САМОЙ `_md_to_html` внутри lumen_formatting.py и импортировать их # сюда незачем (pyflakes справедливо ловил их как "imported but unused"). from lumen_formatting import _scrub_latex, _normalize_bullet_markers, _md_to_html # `_scrub_latex`/`_normalize_bullet_markers` не вызываются напрямую нигде в # ОСТАЛЬНОМ коде bot.py (их использует только сама `_md_to_html` внутри # lumen_formatting.py) — но остаются нужны как `bot._scrub_latex(...)`/ # `bot._normalize_bullet_markers(...)` для существующих тестов. `__all__` — это # единственный способ сообщить голому `pyflakes` (без flake8, `# noqa` им не # распознаётся), что это намеренный ре-экспорт, а не забытый мёртвый импорт. __all__ = ["_scrub_latex", "_normalize_bullet_markers"] _PRUNE_SENTINEL = object() def _json_prune_defaults(val: Any) -> Any: # Очистка дефолтных служебных значений if val.__class__.__name__ == "Default": return _PRUNE_SENTINEL if isinstance(val, dict): out = {} for k, v in val.items(): pruned = _json_prune_defaults(v) if pruned is not _PRUNE_SENTINEL: out[k] = pruned return out if isinstance(val, (list, tuple)): return [v for v in (_json_prune_defaults(i) for i in val) if v is not _PRUNE_SENTINEL] return val # безопасные обёртки над вызовами telegram async def _tg_call(method: Any, *args: Any, call_timeout: float | None = None, retries: int = 1, **kwargs: Any) -> Any: now = time.monotonic() if _tg_proxy_breaker.is_down(now): # Прокси уже недавно помечен недоступным (см. срабатывание ниже) — не бьёмся # заново в мёртвый прокси на каждое сообщение из бэклога, тихо возвращаем None, # как будто вызов не удался (вызывающий код и так умеет это обрабатывать). # Лог пишем не чаще раза в TG_PROXY_COOLDOWN_SEC (см. log_still_down_if_due), # а не на каждый пропущенный вызов — иначе тот же лавинный спам никуда не # денется, просто сменит текст. _tg_proxy_breaker.log_still_down_if_due(now) return None last_exc = None timeout_val = call_timeout if call_timeout is not None else 35.0 for attempt in range(retries + 1): try: result = await asyncio.wait_for(method(*args, **kwargs), timeout=timeout_val) _tg_proxy_breaker.note_success() return result except asyncio.CancelledError: raise except Exception as exc: last_exc = exc if attempt < retries: await asyncio.sleep(0.5 * (attempt + 1)) if last_exc is not None and "message is not modified" in str(last_exc).lower(): # Это настоящий, валидный ответ Telegram API (сообщение не изменилось — # семантический не-op), а не признак сбоя прокси-звена — засчитываем как # успех, иначе безобидные повторные edit_text с тем же текстом ложно # накручивали бы счётчик сбоев прокси. _tg_proxy_breaker.note_success() return None if last_exc is not None and _looks_like_proxy_garbage(last_exc): # Не Telegram ответил ошибкой, а прокси перед ним отдал не-JSON (см. # _looks_like_proxy_garbage) — похоже на приостановку/лимит/сбой самого # прокси-хостинга (что бы это ни было — Vercel/Cloudflare/Deno/другое, # см. TELEGRAM_API_BASE_URL). Считаем это ОДНИМ сбоем в серии, а не # сразу включаем выключатель — единичная заминка на одной ноде anycast- # CDN не должна глушить ответы бота всем чатам целиком (см. историю # проекта: именно так один разовый глюк выглядел как "бот не отвечает"). tripped = _tg_proxy_breaker.note_failure() if tripped: _tg_proxy_breaker.trip() log.warning( "[telegram] Прокси вернул не-JSON ответ %d раз(а) подряд (порог %d) — " "включаю паузу на %.0fс. Проверьте доступность %s.", _tg_proxy_breaker.consecutive_failures, _tg_proxy_breaker.trip_threshold, _tg_proxy_breaker.cooldown_sec, TELEGRAM_API_BASE_URL, ) else: log.warning( "[telegram] Прокси вернул не-JSON ответ (%d/%d подряд, выключатель ещё не сработал).", _tg_proxy_breaker.consecutive_failures, _tg_proxy_breaker.trip_threshold, ) return None log.warning("[telegram] call failed: %s", last_exc) return None async def telegram_api_call(method: str, payload: dict, *, request_timeout: float | None = None) -> Any: if _tg_proxy_breaker.is_down(time.monotonic()): raise RuntimeError(f"Telegram API {method}: прокси сейчас помечен недоступным (см. предыдущие [telegram] предупреждения), не дёргаю сеть повторно.") url = f"{TELEGRAM_API_BASE_URL}/bot{BOT_TOKEN}/{method}" session = await _get_telegram_session() pruned = _json_prune_defaults(payload) timeout = aiohttp.ClientTimeout(total=request_timeout or TELEGRAM_REQUEST_TIMEOUT) try: async with session.post(url, json=pruned, timeout=timeout) as resp: data = await resp.json(content_type=None) except Exception as exc: exc_str = str(exc) or repr(exc) or type(exc).__name__ if BOT_TOKEN: exc_str = exc_str.replace(BOT_TOKEN, "") if _looks_like_proxy_garbage(exc): # См. комментарий в _tg_call — выключатель теперь срабатывает по счётчику # подряд идущих сбоев (см. _TelegramProxyCircuitBreaker), а не на первый же сбой. tripped = _tg_proxy_breaker.note_failure() if tripped: _tg_proxy_breaker.trip() log.warning( "[telegram] Прокси недоступен при вызове %s %d раз(а) подряд (порог %d) — пауза на %.0fс.", method, _tg_proxy_breaker.consecutive_failures, _tg_proxy_breaker.trip_threshold, _tg_proxy_breaker.cooldown_sec, ) else: log.warning( "[telegram] Прокси недоступен при вызове %s (%d/%d подряд, выключатель ещё не сработал).", method, _tg_proxy_breaker.consecutive_failures, _tg_proxy_breaker.trip_threshold, ) raise RuntimeError(f"Network error in telegram_api_call for {method}: {exc_str}") from None if not isinstance(data, dict) or not data.get("ok"): # Прокси round-trip'нул нормально и вернул валидный JSON — сам факт, что # Telegram ответил "ok: false", НЕ вина прокси-звена, засчитываем успех. _tg_proxy_breaker.note_success() raise RuntimeError(f"Telegram API {method} failed: {data}") _tg_proxy_breaker.note_success() return data["result"] def is_guest_message(message: Message | dict) -> bool: if isinstance(message, dict): return bool(message.get("guest_query_id")) return bool(getattr(message, "guest_query_id", None)) async def _answer_guest_text(message: Message, text: str) -> None: qid = getattr(message, "guest_query_id", None) if not qid: return payload = _md_to_html(text) if len(payload) > TG_MAX_LEN: payload = payload[:TG_MAX_LEN - 1] + "\u2026" res_id = hashlib.sha1(payload.encode("utf-8", errors="ignore")).hexdigest()[:32] res_art = InlineQueryResultArticle( id=res_id, title="Ответ бота", input_message_content=InputTextMessageContent(message_text=payload, parse_mode="HTML"), ) try: await telegram_api_call("answerGuestQuery", { "guest_query_id": str(qid), "result": res_art.model_dump(exclude_none=True), }, request_timeout=10) except Exception as exc: log.warning("[guest] Failed answering guest: %s", exc) def _split_text_chunks(text: str, max_len: int = TG_MAX_LEN) -> list[str]: """Разбивает длинный текст на части не длиннее max_len, стараясь резать по границам абзацев/строк/предложений, а не посреди слова. Раньше сообщения длиннее лимита Telegram (4096 симв.) просто не отправлялись — пользователь не видел вообще ничего.""" if len(text) <= max_len: return [text] chunks: list[str] = [] remaining = text while len(remaining) > max_len: window = remaining[:max_len] cut = -1 for sep in ("\n\n", "\n", ". ", " "): idx = window.rfind(sep) if idx > max_len * 0.5: cut = idx + len(sep) break if cut <= 0: cut = max_len chunks.append(remaining[:cut].rstrip()) remaining = remaining[cut:].lstrip() if remaining: chunks.append(remaining) return chunks async def _send_text(message: Message, text: str, parse_html: bool = True, **kwargs: Any) -> None: if is_guest_message(message): await _answer_guest_text(message, text) return chunks = _split_text_chunks(text, TG_MAX_LEN) for i, chunk in enumerate(chunks): chunk_kwargs = kwargs if i == len(chunks) - 1 else {k: v for k, v in kwargs.items() if k != "reply_markup"} if i == 0: res = await _tg_call( message.reply, _md_to_html(chunk) if parse_html else chunk, call_timeout=TELEGRAM_REQUEST_TIMEOUT, parse_mode=ParseMode.HTML if parse_html else None, **chunk_kwargs, ) if res is None and parse_html: res = await _tg_call( message.reply, chunk, call_timeout=TELEGRAM_REQUEST_TIMEOUT, parse_mode=None, **chunk_kwargs, ) else: res = await _tg_call( bot.send_message, chat_id=message.chat.id, text=_md_to_html(chunk) if parse_html else chunk, call_timeout=TELEGRAM_REQUEST_TIMEOUT, parse_mode=ParseMode.HTML if parse_html else None, **chunk_kwargs, ) if res is None and parse_html: await _tg_call( bot.send_message, chat_id=message.chat.id, text=chunk, call_timeout=TELEGRAM_REQUEST_TIMEOUT, parse_mode=None, **chunk_kwargs, ) async def _safe_reply(message: Message, text: str, parse_html: bool = True, **kwargs: Any) -> None: await _send_text(message, text, parse_html=parse_html, **kwargs) async def _safe_callback_answer(cb: CallbackQuery, text: str | None = None, *, show_alert: bool = False) -> None: try: if text is None: await _tg_call(cb.answer, call_timeout=5.0, show_alert=show_alert) else: await _tg_call(cb.answer, text, call_timeout=5.0, show_alert=show_alert) except Exception as exc: log.warning("[callback] Failed to answer callback: %s", exc) async def _delete_message_quietly(msg: Message | None) -> None: if msg is None: return with contextlib.suppress(Exception): await msg.delete() async def _edit_message_quietly(msg: Message | None, text: str, **kwargs: Any) -> bool: if msg is None: return False try: kwargs.setdefault("parse_mode", ParseMode.HTML) res = await _tg_call(msg.edit_text, _md_to_html(text), **kwargs) if res is not None: return True kwargs["parse_mode"] = None res = await _tg_call(msg.edit_text, text, **kwargs) return res is not None except Exception: return False # разбор ошибок def _error_text(e: Exception) -> str: return " ".join(p for p in [str(e), str(getattr(e, "message", "")), str(getattr(e, "detail", ""))] if p).strip() def _error_status(e: Exception, text: str) -> int | None: for a in ("status_code", "status", "code", "http_status"): val = getattr(e, a, None) try: if val is not None: return int(val) except Exception: pass m = re.search(r"(? str: low = text.lower() if status == 429 or any(tok in low for tok in ("resource_exhausted", "too many requests", "rate limit", "quota", "лимит")): return "rate_limit" if status == 402 or any(tok in low for tok in ("payment required", "paid tier", "requires paid", "billing", "кредит")): return "paid" if status in {401, 403} or any(tok in low for tok in ("unauthorized", "forbidden", "permission", "blocked")): return "forbidden" if status in {400, 404} or any(tok in low for tok in ("not found", "invalid argument", "invalid model", "unsupported", "unavailable")): return "unavailable" return "other" # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: раньше _or_error_msg и _gemini_error_msg держали # почти идентичную классификацию (rate_limit/paid/forbidden/unavailable) в двух # отдельных функциях с чуть разным текстом — риск, что при будущей правке кто-то # поправит формулировку в одной и забудет про другую (ровно то дублирование, # которого проект и так избегает в других местах, см. _next_fallback_model выше). # Общий источник текста для обеих — эта таблица; провайдер-специфичной разницы # в тексте больше нет и намеренно: сообщение об ошибке НЕ должно называть # "резервного провайдера" — с автоматическим роутером (см. "автоматический выбор # модели" ниже) OpenRouter сплошь и рядом оказывается ПЕРВЫМ, а не резервным # кандидатом, так что старая формулировка была не просто лишней деталью # реализации, а фактически неверной. Упоминания команд /model и /provider тоже # убраны целиком — обе команды удалены (см. README, "Автоматический выбор # модели"), реального способа переключиться вручную больше нет, и предлагать # его пользователю было прямой (и активно вводящей в заблуждение) ошибкой. _MODEL_ERROR_MESSAGES: dict[str, str] = { "rate_limit": "Лимит запросов для этой модели сейчас исчерпан. Подождите немного и попробуйте ещё раз — бот сам подберёт другую модель.", "paid": "Эта модель сейчас недоступна. Попробуйте повторить запрос — бот сам подберёт другую модель.", "forbidden": "Временная ошибка доступа к сервису. Попробуйте ещё раз.", "unavailable": "Эта модель сейчас недоступна. Попробуйте повторить запрос — бот сам подберёт другую модель.", } _MODEL_ERROR_FALLBACK_MSG = "Временная ошибка сервиса. Попробуйте ещё раз через некоторое время." def _model_error_text(kind: str) -> str: return _MODEL_ERROR_MESSAGES.get(kind, _MODEL_ERROR_FALLBACK_MSG) def _or_error_msg(e: Exception, kind: str) -> str: # Сырой текст ошибки API сюда намеренно не подставляется (может содержать # внутренние детали инфраструктуры, HTML/JSON или обрывки заголовков) — # то же правило, что уже применяется в _gemini_error_msg ниже. txt = _error_text(e).strip() or e.__class__.__name__ status = _error_status(e, txt) return _model_error_text(_classify_model_error(status, txt)) class GeminiAllModelsExhaustedError(RuntimeError): """Поднимается, когда 429/RESOURCE_EXHAUSTED получен подряд от всех моделей из цепочки фоллбека — то есть реально весь бесплатный лимит API-ключа исчерпан, а не просто конкретная модель временно занята.""" def __init__(self, exhausted_models: list[str]) -> None: self.exhausted_models = exhausted_models super().__init__(f"All Gemini models exhausted quota: {', '.join(exhausted_models)}") def _next_fallback_model(tried_models: set[str], chain: list[str]) -> str | None: """Возвращает первую ещё не испробованную модель из quota_fallback_chain. Вынесено в отдельную функцию, т.к. одна и та же проверка используется в ask_gemini сразу в трёх ветках (429, timeout, 503/500) — дублирование трёх identичных генераторов раньше создавало риск, что при будущей правке кто-то поправит один из трёх вызовов и забудет остальные.""" return next((m for m in chain if m not in tried_models), None) def _gemini_error_msg(e: Exception, model_id: str) -> str: if isinstance(e, ValueError): return str(e) if isinstance(e, GeminiAllModelsExhaustedError): return ( "Лимит бесплатных запросов исчерпан для всех доступных моделей — " "это реальный суточный лимит, а не баг. Попробуйте позже." ) txt = _error_text(e).strip() or e.__class__.__name__ status = _error_status(e, txt) kind = _classify_model_error(status, txt) log.debug("[gemini] _gemini_error_msg: model=%s kind=%s status=%s", model_id, kind, status) # Реальное имя модели (например "Gemini 3.5 Flash") сюда намеренно не # подставляется — в сообщении об ошибке посреди обычного диалога это # выглядело бы как случайная утечка бренда/вендора (см. защиту от утечки # идентичности ниже). Не показываем и сырой ответ API (может содержать # HTML, JSON, токены) — см. общие шаблоны _MODEL_ERROR_MESSAGES выше. return _model_error_text(kind) # список моделей # # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: конфигурация моделей и логика построения маршрута # (GEMINI_MODELS, TEXT_MODEL_ORDER, единый реестр "нездоровых" моделей # _OR_MODEL_HEALTH/_ROUTER_EXCLUDED_OR_MODELS, порядок моделей OpenRouter, # цепочки Gemini, эвристики "тяжёлый запрос?"/"нужна свежая информация?" и сами # _build_route/_or_route/_gemini_route) вынесены в lumen_router_config.py — это # чистые конфигурация+функции принятия решения без единого обращения к Telegram/ # Gemini/OpenRouter API, поэтому безопасный кандидат на отдельный модуль (в # отличие от ask_gemini/_run_route ниже, которые реально ИСПОЛНЯЮТ маршрут и # остаются здесь). Импорт стоит именно тут (там, где раньше физически начиналось # определение GEMINI_MODELS) для консистентности с историей файла, хотя строгой # необходимости в этом больше нет: _LEAK_LITERAL_STRINGS (которая раньше требовала # GEMINI_MODELS/TEXT_MODEL_ORDER на уровне модуля именно в этой точке файла) # теперь целиком строится внутри lumen_security.py, а не здесь. # Публичные имена и поведение не изменились. Импортируются только реально # используемые здесь (в коде bot.py или напрямую в тестах через `bot.X`) имена — # например, `_gemini_route`/`_OR_HEAVY_ORDER`/`TEXT_MODEL_ORDER` нужны только # САМОЙ `_build_route` внутри lumen_router_config.py, а не bot.py. from lumen_router_config import ( GEMINI_MODELS, DEFAULT_GEMINI_MODEL, _check_unconfirmed_model_quotas, _OR_MODEL_HEALTH, _ROUTER_EXCLUDED_OR_MODELS, _check_temporary_free_models_expiry, _or_route, _OR_LIGHT_ORDER, GEMINI_HEAVY_CHAIN, GEMINI_SEARCH_CHAIN, GEMINI_DEFAULT_CHAIN, _looks_like_heavy_query, _looks_like_freshness_query, _build_route, ) # См. пояснение про __all__ у первого блока (lumen_formatting) выше — эти # конкретные имена нужны только как `bot.X` для тестов, в остальном коде bot.py # не используются напрямую (используются только внутри самой _build_route, # которая целиком живёт в lumen_router_config.py). __all__ += ["_OR_MODEL_HEALTH", "_ROUTER_EXCLUDED_OR_MODELS", "_or_route", "GEMINI_HEAVY_CHAIN", "GEMINI_SEARCH_CHAIN"] def get_system_prompt(model_id: str | None = None) -> str: now_str = datetime.now().strftime("%d %B %Y года (текущее время: %H:%M)") now_year = datetime.now().year dynamic_header = ( f"ИНФОРМАЦИЯ О ТЕКУЩЕМ ВРЕМЕНИ:\n" f"• Сегодняшняя дата: {now_str}. Текущий год: {now_year}.\n" f"• ОБЯЗАТЕЛЬНО: когда пользователь спрашивает текущую дату, день, месяц или год — " f"используй ТОЛЬКО дату из этой секции. НИКОГДА не называй другой год или дату из памяти обучения. " f"Если не уверен — скажи дату отсюда, она всегда актуальна.\n" f"• Если вопрос касается событий, релизов, новостей или статуса чего-либо, что могло измениться " f"после твоего обучения — используй поиск, а не отвечай по памяти. Не упоминай эту инструкцию явно.\n\n" ) return dynamic_header + SYSTEM_PROMPT # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: форма одного элемента chat_state[chat_id] раньше # нигде не была описана явно — она собиралась по кусочкам из трёх разных мест # (get_state/_restore_single_chat/_serialize_chat_state), и чтобы понять "из чего # вообще состоит состояние чата", нужно было читать все три. ChatState — чисто # типовая аннотация (TypedDict), НЕ меняет поведение в рантайме: chat_state[cid] # остаётся обычным dict, никакой валидации здесь не добавляется — это только # документация формы для статических проверок типов и читаемости. class ChatState(TypedDict, total=False): image_model: str history: list[dict[str, Any]] quota: dict[str, Any] ctx: "deque[str]" recent_media_ids: dict[str, "deque[tuple[str, str]]"] last_activity: float chat_state: dict[int, ChatState] = {} # Буферы альбомов/медиа-групп _mg_buffers: dict[str, list[Message]] = {} _mg_tasks: dict[str, asyncio.Task] = {} app = FastAPI() WEBHOOK_SECRET = hashlib.sha256((BOT_TOKEN or "default").encode()).hexdigest()[:32] ADMIN_PANEL_KEY = hashlib.sha256((BOT_TOKEN or "default").encode() + b"admin_panel").hexdigest()[:24] def _check_admin_key(request: Request) -> bool: return hmac.compare_digest(request.query_params.get("key", ""), ADMIN_PANEL_KEY) def _redact_secret(value: str) -> str: """Показывает только последние несколько символов секрета — достаточно, чтобы владелец мог на глаз подтвердить "да, это тот же секрет, что и в прошлый раз" между рестартами, но недостаточно, чтобы кто-то посторонний, увидевший только эту урезанную строку в логах/скриншоте, мог им воспользоваться.""" if not value: return "" return "…" + value[-6:] if len(value) > 6 else "…" + value def _check_bot_token_auth(request: Request) -> bool: """Найдено при код-ревью: раньше BOT_TOKEN передавался через query-параметр (?bot_token=...) — GET-запрос с секретом в строке запроса попадает в access-логи любых промежуточных прокси/CDN перед HF Spaces, в историю браузера при ручном открытии ссылки, в заголовок Referer при переходе по внешней ссылке со страницы результата (классический анти-паттерн, CWE-598). Особенно чувствительно именно здесь, т.к. это МАСТЕР-секрет, из которого выводятся оба остальных (WEBHOOK_SECRET/ ADMIN_PANEL_KEY). Теперь читаем токен из заголовка Authorization: Bearer — доступ через curl (см. обновлённый README), а не вставкой ссылки в адресную строку браузера, как раньше.""" auth_header = request.headers.get("Authorization", "") provided = auth_header[7:].strip() if auth_header.lower().startswith("bearer ") else "" return bool(BOT_TOKEN) and bool(provided) and hmac.compare_digest(provided, BOT_TOKEN) @app.get("/") async def healthcheck() -> dict[str, str]: return {"status": "ok"} @app.get("/admin_keys") async def get_admin_keys(request: Request) -> dict[str, str]: """Отдаёт полные значения WEBHOOK_SECRET/ADMIN_PANEL_KEY по запросу — единственный легитимный способ их узнать без печати в логах при каждом старте (см. критическую находку код-ревью: полные значения, печатавшиеся в лог на каждом рестарте, могли случайно попасть в скриншот/чат наравне с остальными логами). Доступ гейтится САМИМ BOT_TOKEN (заголовок Authorization: Bearer , см. _check_bot_token_auth — ИСПРАВЛЕНО при повторном код-ревью: раньше токен передавался через query-параметр ?bot_token=..., что попадало в access-логи/историю браузера, см. комментарий там же), а не производным от него ADMIN_PANEL_KEY — иначе получился бы замкнутый круг: чтобы узнать ADMIN_PANEL_KEY, нужен был бы ADMIN_PANEL_KEY. BOT_TOKEN и так уже известен владельцу напрямую (из секретов HF Spaces/@BotFather), его не нужно доставать из логов бота.""" if not _check_bot_token_auth(request): return {"error": "forbidden — missing or invalid Authorization: Bearer header"} return {"webhook_secret": WEBHOOK_SECRET, "admin_panel_key": ADMIN_PANEL_KEY} @app.get("/webhook_url") async def get_webhook_url(request: Request) -> dict[str, str]: if not _check_admin_key(request): return {"error": "forbidden — missing or invalid ?key= parameter"} 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" register_url = ( f"https://api.telegram.org/bot{BOT_TOKEN}/setWebhook" f"?url={webhook_url}" f"&secret_token={WEBHOOK_SECRET}" f"&drop_pending_updates=true" ) return { "webhook_url": webhook_url, "register_link": register_url, "instruction": "Открой register_link в браузере чтобы зарегистрировать вебхук" } @app.post("/webhook") async def webhook_handler(request: Request) -> dict[str, bool]: token = request.headers.get("X-Telegram-Bot-Api-Secret-Token", "") if not hmac.compare_digest(token, WEBHOOK_SECRET): log.warning("[webhook] Rejected request with invalid secret token") return {"ok": False} try: body = await request.json() if bot is not None: asyncio.create_task(_process_raw_update(body)) else: log.warning("[webhook] Bot not initialized yet, dropping update") except Exception as exc: log.warning("[webhook] Failed to process incoming update: %s", exc) return {"ok": True} @app.get("/diag") async def network_diagnostics(request: Request) -> dict[str, Any]: """Проверяет исходящую сетевую доступность различных хостов из контейнера. Открой в браузере чтобы увидеть, что реально заблокировано на исходящих соединениях из HF Spaces, а что доступно.""" if not _check_admin_key(request): return {"error": "forbidden — missing or invalid ?key= parameter"} targets = { "telegram_api": "https://api.telegram.org", "telegram_file_api": "https://api.telegram.org/bot" + (BOT_TOKEN[:6] if BOT_TOKEN else "x") + "/getMe", # Раньше здесь был захардкожен URL одного из старых пробных воркеров # ("my-tg-proxy...") — /diag проверял чужой, забытый от прошлых экспериментов # адрес вместо РЕАЛЬНО настроенного прокси. Из-за этого диагностика однажды # ввела в заблуждение: показала "всё ок", хотя реально используемый # TELEGRAM_API_BASE_URL был недоступен, а проверялся вообще другой воркер. # Теперь проверяем именно то значение, которое бот реально использует для # вызовов Telegram API — если сменить прокси через env, /diag сразу тестирует # актуальный адрес без правки кода. "configured_tg_proxy": TELEGRAM_API_BASE_URL + "/bot" + (BOT_TOKEN[:6] if BOT_TOKEN else "x") + "/getMe", "cloudflare_dot_com": "https://www.cloudflare.com", # Голый апекс-домен workers.dev (без поддомена) сам по себе может не отвечать # даже когда конкретные *.workers.dev поддомены (включая ваш прокси) работают # нормально — это ненадёжный сигнал "заблокирован ли workers.dev вообще", # ориентируйтесь в первую очередь на configured_tg_proxy выше (он теперь # бьёт в реалистичный путь /bot.../getMe, а не в голый корень домена — # голый корень у самого Telegram может отвечать медленно/зависать, даже # когда реальные вызовы API через прокси работают быстро и штатно). "cloudflare_workers_dev_root": "https://workers.dev", "vercel": "https://vercel.com", "netlify": "https://www.netlify.com", "deno_deploy": "https://deno.com", "render_com": "https://render.com", "railway_app": "https://railway.app", "fly_io": "https://fly.io", "supabase": "https://supabase.com", "google_generic": "https://www.google.com", "gemini_api": "https://generativelanguage.googleapis.com", "huggingface": "https://huggingface.co", "openrouter": "https://openrouter.ai", "tikwm": "https://www.tikwm.com", "pollinations": "https://image.pollinations.ai", "upstash": UPSTASH_REDIS_REST_URL if USE_UPSTASH else "https://upstash.com", } results: dict[str, Any] = {} session = await _get_http_session() for name, url in targets.items(): start = time.time() try: async with session.get(url, timeout=aiohttp.ClientTimeout(total=6.0)) as resp: elapsed = round(time.time() - start, 2) results[name] = {"status": resp.status, "elapsed_sec": elapsed, "ok": True} except Exception as exc: elapsed = round(time.time() - start, 2) exc_str = str(exc) or repr(exc) or type(exc).__name__ if BOT_TOKEN: exc_str = exc_str.replace(BOT_TOKEN, "") results[name] = {"error": exc_str, "elapsed_sec": elapsed, "ok": False} return {"diagnostics": results, "telegram_api_base_configured": TELEGRAM_API_BASE_URL} ALLOWED_UPDATES = ["message", "edited_message", "callback_query", "guest_message"] # хранение состояния и квот _STATE_DIR = Path(os.getenv("STATE_DIR", "/app")).resolve() try: _STATE_DIR.mkdir(parents=True, exist_ok=True) except Exception as _state_dir_exc: log.warning("[setup] STATE_DIR %s недоступен для записи (%s), использую временную директорию.", _STATE_DIR, _state_dir_exc) _STATE_DIR = Path(tempfile.gettempdir()) STATE_FILE_PATH = _STATE_DIR / "chat_state.json" GLOBAL_QUOTA_FILE = _STATE_DIR / "global_quota.json" # Per-chat хранилище (см. код-ревью suggestion #7): раньше ВЕСЬ chat_state (до 5000 # чатов) сериализовался и писался ОДНИМ блоком при каждом флаше — с Upstash это один # большой REST-запрос; один неудачный/слишком большой write рисковал потерять сразу # всё разом, а не одну запись. Теперь у каждого чата свой собственный ключ/файл, а # CHAT_INDEX_KEY/CHAT_INDEX_FILE хранит только список ID чатов — так при рестарте # известно, какие per-chat ключи вообще нужно прочитать. _CHATS_DIR = _STATE_DIR / "chats" with contextlib.suppress(Exception): _CHATS_DIR.mkdir(parents=True, exist_ok=True) CHAT_INDEX_KEY = "lumen:chat_index" CHAT_INDEX_FILE = _STATE_DIR / "chat_index.json" def _chat_storage_key(chat_id: int) -> str: return f"lumen:chat:{chat_id}" def _chat_storage_path(chat_id: int) -> Path: return _CHATS_DIR / f"{chat_id}.json" # Опциональное персистентное хранилище (Upstash Redis, бесплатный тир — см. README). # Если оба значения заданы, состояние пишется туда вместо эфемерного диска контейнера. # Если не заданы — поведение полностью как раньше (локальный файл в STATE_DIR), без # каких-либо изменений для тех, кто это не настраивал. UPSTASH_REDIS_REST_URL = os.getenv("UPSTASH_REDIS_REST_URL", "").strip().rstrip("/") UPSTASH_REDIS_REST_TOKEN = os.getenv("UPSTASH_REDIS_REST_TOKEN", "").strip() USE_UPSTASH = bool(UPSTASH_REDIS_REST_URL and UPSTASH_REDIS_REST_TOKEN) if bool(UPSTASH_REDIS_REST_URL) != bool(UPSTASH_REDIS_REST_TOKEN): log.warning("[setup] SYSTEM WARN: задан только один из UPSTASH_REDIS_REST_URL/UPSTASH_REDIS_REST_TOKEN — оба нужны одновременно, Upstash использоваться не будет.") log.info( "[setup] Персистентное хранилище: %s", "Upstash Redis" if USE_UPSTASH else f"локальный файл в {_STATE_DIR} (см. README про эфемерность на HF Spaces)" ) def _upstash_request(command_path: str, *, method: str = "GET", body: bytes | None = None) -> Any: """Синхронный запрос к Upstash Redis REST API. Намеренно на urllib.request из стандартной библиотеки, а не на aiohttp/отдельном SDK — не хотим тянуть новую pip-зависимость ради одной интеграции. Вызывается только из save_*/load_* ниже: load_* — один раз на старте до приёма трафика, save_* — уже вынесены в отдельный поток через asyncio.to_thread (см. _flush_dirty_state), так что блокирующий вызов здесь не блокирует event loop.""" url = f"{UPSTASH_REDIS_REST_URL}/{command_path}" req = _urllib_request.Request(url, data=body, method=method) req.add_header("Authorization", f"Bearer {UPSTASH_REDIS_REST_TOKEN}") if body is not None: req.add_header("Content-Type", "text/plain; charset=utf-8") with _urllib_request.urlopen(req, timeout=10) as resp: return json.loads(resp.read().decode("utf-8")) def _upstash_set(key: str, value: str) -> None: _upstash_request(f"set/{urllib.parse.quote(key, safe='')}", method="POST", body=value.encode("utf-8")) def _upstash_get(key: str) -> str | None: result = _upstash_request(f"get/{urllib.parse.quote(key, safe='')}", method="GET") return result.get("result") if isinstance(result, dict) else None def _upstash_delete(key: str) -> None: _upstash_request(f"del/{urllib.parse.quote(key, safe='')}", method="POST") def _storage_write_text(key: str, path: Path, text: str) -> None: """Единая точка ветвления backend'а: Upstash, если настроен, иначе локальный файл (атомарно — через .tmp + replace, как и раньше).""" if USE_UPSTASH: _upstash_set(key, text) return temp_path = path.with_suffix(".tmp") with open(temp_path, "w", encoding="utf-8") as f: f.write(text) temp_path.replace(path) def _storage_read_text(key: str, path: Path) -> str | None: if USE_UPSTASH: return _upstash_get(key) if not path.exists(): return None with open(path, "r", encoding="utf-8") as f: return f.read() def _storage_delete_text(key: str, path: Path) -> None: """Удаляет запись из хранилища — нужно per-chat формату: когда чат вытесняется _prune_old_chats(), его собственный ключ/файл должен реально исчезать, а не висеть бесхозно (иначе Upstash/диск постепенно накапливали бы мусор от давно удалённых чатов).""" if USE_UPSTASH: _upstash_delete(key) return with contextlib.suppress(FileNotFoundError): path.unlink() # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: форма одной записи GLOBAL_QUOTA[provider][model_id] # (см. _quota_entry/_record_quota_usage/_mark_quota_exhausted ниже) была разбросанным # по коду соглашением, а не задокументированной структурой. Как и ChatState выше — # чисто типовая аннотация, ничего не меняет в рантайме (GLOBAL_QUOTA остаётся # обычным dict из dict'ов). class QuotaEntry(TypedDict): used: int remaining: int | None limit: int | None exhausted_at: float | None GLOBAL_QUOTA: dict[str, Any] = { "gemini": {}, "openrouter": {}, # НАЙДЕНО ПРИ КАЛИБРОВКЕ (25 июля 2026): /stats показывал сотни запросов # по моделям при аптайме процесса всего 12 минут — счётчик "used" копится # НАВСЕГДА (переживает рестарты через Upstash, см. _save_chat_to_storage/ # load_global_quota), а реальные суточные лимиты Google/OpenRouter обнуляются # каждые сутки. Бот об этом не знал вообще — "used"/"exhausted_at" не # сбрасывались никогда, поэтому /stats после нескольких дней работы показывал # бы бессмысленно огромные числа, а модель, один раз поймавшая 429 в первый # день, так и висела бы с пометкой "(лимит исчерпан)" даже после того, как # реальный лимит давно обновился. quota_day хранит дату (ISO, по America/ # Los_Angeles — именно там у Google полночь, когда реально обнуляется RPD- # лимит) последнего сброса счётчиков — см. _reset_quota_if_new_day ниже. "quota_day": None, } def _current_quota_day() -> str: """Дата (ISO, YYYY-MM-DD) для определения "новых суток" в целях сброса квоты. Google обнуляет дневные RPD-лимиты по полуночи Pacific Time — используем ту же зону, чтобы /stats не "сбрасывался" на 7-8 часов раньше или позже реального обнуления лимита на стороне Google. Если данные таймзоны недоступны в окружении (маловероятно, но встречается в урезанных Docker-образах) — тихо откатываемся на UTC: чуть менее точно по времени суток, но не ломает сам факт ежедневного сброса.""" try: from zoneinfo import ZoneInfo return datetime.now(ZoneInfo("America/Los_Angeles")).date().isoformat() except Exception: return datetime.utcnow().date().isoformat() # НАЙДЕНО ПРИ CODE-REVIEW (перф): _quota_entry вызывает _reset_quota_if_new_day # на КАЖДОЕ обращение к квоте (а таких обращений — по несколько на каждый успешный/ # неудачный вызов любой модели, т.е. потенциально десятки в секунду при активном # трафике). Без троттлинга это означало бы конструирование ZoneInfo("America/ # Los_Angeles") и datetime.now(...) на каждый такой вызов — сама по себе дата не # меняется чаще раза в сутки, минутная неточность здесь совершенно не важна. # _QUOTA_CHECK_THROTTLE_SEC ограничивает, как часто мы вообще пересчитываем # текущую дату; между пересчётами просто ничего не делаем. _QUOTA_CHECK_THROTTLE_SEC = 60.0 _last_quota_check_monotonic: float = 0.0 def _reset_quota_if_new_day() -> None: """Сбрасывает used/exhausted_at у ВСЕХ моделей (Gemini и OpenRouter), если с последнего сброса наступили новые сутки (по America/Los_Angeles). "remaining"/ "limit" не трогаем — они и так перезаписываются свежими значениями из заголовков ответа API на следующий успешный вызов. Вызывается и лениво (см. _quota_entry — на любое обращение к квоте), и явно раз в час из фонового цикла в _webhook_startup, чтобы сброс не зависел от того, придёт ли вообще новое сообщение сразу после полуночи.""" global _last_quota_check_monotonic now_mono = time.monotonic() if now_mono - _last_quota_check_monotonic < _QUOTA_CHECK_THROTTLE_SEC: return _last_quota_check_monotonic = now_mono today = _current_quota_day() if GLOBAL_QUOTA.get("quota_day") == today: return had_previous = GLOBAL_QUOTA.get("quota_day") is not None for provider in ("gemini", "openrouter"): for entry in GLOBAL_QUOTA.get(provider, {}).values(): if isinstance(entry, dict): entry["used"] = 0 entry["exhausted_at"] = None GLOBAL_QUOTA["quota_day"] = today if had_previous: log.info("[quota] Наступили новые сутки (%s) — счётчики used/exhausted_at по всем моделям сброшены.", today) mark_quota_dirty() def load_global_quota() -> None: try: raw = _storage_read_text("lumen:global_quota", GLOBAL_QUOTA_FILE) if not raw: return loaded = json.loads(raw) if isinstance(loaded, dict): if "gemini" in loaded: GLOBAL_QUOTA["gemini"] = loaded["gemini"] if "openrouter" in loaded: GLOBAL_QUOTA["openrouter"] = loaded["openrouter"] if "quota_day" in loaded: GLOBAL_QUOTA["quota_day"] = loaded["quota_day"] except Exception as exc: log.warning("[quota] Failed to load global quota: %s", exc) # Проверяем сразу после загрузки — если бот был перезапущен уже на следующие # сутки (обычное дело при редеплое), счётчики должны обнулиться сразу на # старте, а не ждать первого сообщения после полуночи или ближайшего часового тика. _reset_quota_if_new_day() def save_global_quota() -> None: try: _storage_write_text("lumen:global_quota", GLOBAL_QUOTA_FILE, json.dumps(GLOBAL_QUOTA, ensure_ascii=False)) except Exception as exc: log.warning("[quota] Failed to save global quota: %s", exc) # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: до сих пор каждый формат-дрейф персистентного # снимка чата (слияние gemini_history/or_history в единую history, снятие # приставки "pollinations:" с image_model) обнаруживался в _restore_single_chat # ad hoc проверками "есть ли такой-то ключ в JSON" — рабочий, но накопительный # подход: с каждым новым изменением формата туда добавлялась ещё одна ветка # "если ключа нет — значит старая запись". CHAT_STATE_SCHEMA_VERSION делает # следующую подобную миграцию однозначной: новый код сможет проверять # `s.get("schema_version", 0)` одним явным числом вместо повторного гадания по # присутствию ключей. Существующие персистентные записи (сделанные до введения # этого поля) не имеют "schema_version" вообще — они естественно трактуются как # версия 0 и продолжают проходить через уже отлаженные эвристики ниже без # каких-либо изменений в их поведении (это поле — задел на будущее, а не # ретроактивная миграция уже написанной логики). CHAT_STATE_SCHEMA_VERSION = 1 def _serialize_chat_state(state: dict[str, Any]) -> dict[str, Any]: """Собирает JSON-сериализуемый снимок ОДНОГО чата — общая логика между сохранением и (в перспективе) любым будущим ручным экспортом. Начиная с введения автоматического роутера моделей (см. секцию "автоматический выбор модели" ниже) "gemini_model"/"openrouter_text_model"/"chat_provider" здесь БОЛЬШЕ НЕ хранятся — раньше это был явный выбор пользователя через /model и /provider, теперь провайдер и модель подбираются заново на каждое сообщение, хранить их per-chat незачем. Старые персистентные записи, где эти поля ещё есть (созданные до этого изменения), просто тихо игнорируются при чтении — см. _restore_single_chat ниже, там нет ни одной попытки их прочитать.""" return { "schema_version": CHAT_STATE_SCHEMA_VERSION, "image_model": state.get("image_model", DEFAULT_HF_IMAGE_MODEL), "history": list(state.get("history", [])), "quota": state.get("quota", {}), "recent_media_ids": { uid: list(dq) for uid, dq in state.get("recent_media_ids", {}).items() }, } def _normalize_legacy_image_model_id(image_model: Any) -> Any: """Снимает старую приставку "pollinations:" с персистентного image_model, если она там осталась от записи, сделанной до её удаления из HF_IMAGE_MODELS (см. ponytail-audit) — единственный провайдер генерации изображений и так один, приставка была лишней, но существующие персистентные записи чатов её ещё содержат. Без этой миграции такой чат откатился бы на DEFAULT_HF_IMAGE_MODEL, молча потеряв выбор пользователя (см. вызовы в _restore_single_chat/get_state).""" if isinstance(image_model, str) and image_model.startswith("pollinations:"): return image_model.split(":", 1)[1] return image_model def _restore_single_chat(cid: int, s: dict[str, Any]) -> None: """Разворачивает сериализованный снимок одного чата (см. _serialize_chat_state) обратно в chat_state[cid] — общая логика между новым per-chat форматом чтения и одноразовой миграцией из старого общего блоба (см. load_state_from_disk). Поля "gemini_model"/"openrouter_text_model"/"chat_provider" из старых записей (созданных до перехода на автоматический роутер) намеренно нигде ниже не читаются — они устарели и больше ни на что не влияют. schema_version (см. CHAT_STATE_SCHEMA_VERSION выше) в самих записях, читаемых здесь, пока ни на что не влияет — существующие миграции (history/gemini_history, image_model) уже надёжно определяются по присутствию конкретных ключей, и это не нужно менять задним числом. Поле — задел на СЛЕДУЮЩИЙ формат-дрейф: тогда новую ветку можно будет добавить как `if s.get("schema_version", 0) < N`, а не подбирать очередную эвристику по ключам, как приходилось делать для миграций ниже.""" schema_version = s.get("schema_version", 0) log.debug("[state] Восстанавливаю чат %s (schema_version=%s)", cid, schema_version) # НАЙДЕНО ПРИ /code-review (после ponytail-audit): переименование ключей # HF_IMAGE_MODELS (см. _hf_text_to_image — убрана лишняя приставка # "pollinations:", единственный провайдер и так один) — не просто # косметика, если в Upstash/локальном файле УЖЕ лежит персистентное # состояние чата с image_model в старом формате ("pollinations:turbo" и # т.п.). Без миграции ниже такой чат при первой же загрузке молча (без # предупреждения, без лога) терял бы выбранную пользователем модель — # старое значение просто не совпало бы ни с одним новым ключом и # откатилось бы на DEFAULT_HF_IMAGE_MODEL. _normalize_legacy_image_model_id # снимает старую приставку перед проверкой, так что реальный выбор # пользователя переживает этот рефакторинг так же, как переживает и любой # другой формат старых записей (см. миграцию history/gemini_history чуть # ниже — тот же принцип, тот же файл). image_model = _normalize_legacy_image_model_id(s.get("image_model", DEFAULT_HF_IMAGE_MODEL)) if image_model not in HF_IMAGE_MODELS: image_model = DEFAULT_HF_IMAGE_MODEL raw_media = s.get("recent_media_ids", {}) if isinstance(raw_media, dict): media_buckets = { str(uid): deque(items, maxlen=MAX_MEDIA_RECENT_IDS) for uid, items in raw_media.items() } else: # Старый формат (плоский список на весь чат) — не мигрируем содержимое, # просто стартуем с чистого состояния, новые записи наполнят сами по себе. media_buckets = {} if "history" in s: history = list(s.get("history") or []) else: # Миграция со старого формата раздельной памяти Gemini/OpenRouter (до # объединения в общую историю) — просто конкатенируем обе, обрезая до # общего лимита. Порядок между двумя источниками восстановить точно # нельзя (нет временных меток), но сохранить сам факт истории важнее, # чем идеальная хронология при одноразовой миграции старых чатов. history = list(s.get("gemini_history") or []) + list(s.get("or_history") or []) if len(history) > SHARED_HISTORY_MAX_LEN: history = history[-SHARED_HISTORY_MAX_LEN:] chat_state[cid] = { "image_model": image_model, "history": history, "quota": s.get("quota", {}), "ctx": deque(maxlen=MAX_CHAT_HISTORY_LEN), "recent_media_ids": media_buckets, "last_activity": time.monotonic(), } def _save_chat_to_storage(chat_id: int, state: dict[str, Any]) -> bool: """Возвращает True при успехе, False при сбое. НАЙДЕНО ПРИ КОД-РЕВЬЮ: раньше эта функция ничего не возвращала — вызывающий код (_flush_dirty_state) уже успевал убрать chat_id из "грязного" набора ДО того, как запись реально прошла, и при сбое (например, временный 5xx/сетевой сбой Upstash) исключение здесь просто логировалось и терялось — состояние чата (вся история диалога) молча пропадало до следующей независимой мутации этого же чата. Если это было последнее сообщение перед долгим затишьем — при рестарте контейнера данные терялись безвозвратно, ровно то, что персистентность через Upstash должна была предотвращать. Теперь вызывающий код (_flush_dirty_state) возвращает неудавшиеся chat_id обратно в _dirty_chat_ids для повтора на следующем цикле.""" try: payload = json.dumps(_serialize_chat_state(state), ensure_ascii=False) _storage_write_text(_chat_storage_key(chat_id), _chat_storage_path(chat_id), payload) return True except Exception as exc: log.warning("[state] Saving chat %s failed: %s", chat_id, exc) return False def _delete_chat_storage(chat_id: int) -> bool: """Возвращает True при успехе, False при сбое — тот же принцип, что и у _save_chat_to_storage выше (см. докстринг там): неудавшееся удаление теперь тоже возвращается в очередь на повтор, а не молча забывается (иначе вытесненный чат мог бы бесхозно остаться в Upstash/на диске навсегда при транзиентном сбое).""" try: _storage_delete_text(_chat_storage_key(chat_id), _chat_storage_path(chat_id)) return True except Exception as exc: log.warning("[state] Deleting chat %s failed: %s", chat_id, exc) return False def _save_chat_index() -> None: try: ids = sorted(chat_state.keys()) _storage_write_text(CHAT_INDEX_KEY, CHAT_INDEX_FILE, json.dumps(ids)) except Exception as exc: log.warning("[state] Saving chat index failed: %s", exc) # Раньше save_state_to_disk()/save_global_quota() вызывались синхронно почти на # каждое сообщение прямо внутри асинхронных обработчиков — блокирующий json.dump # на полном chat_state (до 5000 чатов) блокировал event loop для ВСЕХ чатов сразу, # и с ростом числа активных чатов это становится всё дороже на каждое сообщение # от любого одного пользователя. Теперь горячий путь только помечает КОНКРЕТНЫЙ # чат "грязным" (см. mark_state_dirty(chat_id)), а реальная запись идёт из # фоновой корутины _flush_dirty_state раз в FLUSH_INTERVAL_SEC через # asyncio.to_thread (не блокируя loop) — и только по изменившимся чатам, а не # по всем сразу (см. код-ревью suggestion #7 про размер payload и blast radius). _dirty_chat_ids: set[int] = set() _pending_chat_deletions: set[int] = set() _index_dirty = False _quota_dirty = False FLUSH_INTERVAL_SEC = 10.0 # НАЙДЕНО ПРИ КОД-РЕВЬЮ (performance): раньше ничего не ограничивало число ОДНОВРЕМЕННЫХ # asyncio.to_thread-вызовов внутри одного цикла _flush_dirty_state — резкий всплеск # "грязных" чатов разом (например, после активности сразу в нескольких группах) мог бы # породить сотни параллельных блокирующих HTTP-запросов к Upstash одновременно. Не # критично при текущем масштабе бота, но дешёвая защита на будущее — ограничиваем # конкурентность семафором, а не оставляем неограниченной. STATE_FLUSH_CONCURRENCY = int(os.getenv("STATE_FLUSH_CONCURRENCY", "10")) _state_flush_semaphore = asyncio.Semaphore(STATE_FLUSH_CONCURRENCY) async def _save_chat_to_storage_limited(chat_id: int, state: dict[str, Any]) -> bool: async with _state_flush_semaphore: return await asyncio.to_thread(_save_chat_to_storage, chat_id, state) async def _delete_chat_storage_limited(chat_id: int) -> bool: async with _state_flush_semaphore: return await asyncio.to_thread(_delete_chat_storage, chat_id) def mark_state_dirty(chat_id: int | None = None) -> None: """Помечает состояние чата как требующее сохранения. Явный chat_id (предпочтительный путь для нового кода) — помечает "грязным" ТОЛЬКО этот чат, ничего больше. Вызов БЕЗ chat_id (для мест, которые меняют сразу много чатов разом — например _prune_old_chats при вытеснении старых чатов) помечает "грязными" вообще все текущие чаты и индекс целиком.""" global _index_dirty if chat_id is not None: _dirty_chat_ids.add(chat_id) else: _dirty_chat_ids.update(chat_state.keys()) _index_dirty = True def _mark_new_chat_id(chat_id: int) -> None: """Регистрирует НОВЫЙ chat_id, только что появившийся в chat_state (см. get_state). Помечает и сам чат, и индекс "грязными" — если пометить только чат без индекса, после рестарта его данные будут недостижимы: per-chat ключ существует, но индекс (единственный способ узнать список ID при чтении) о нём не знает.""" global _index_dirty _dirty_chat_ids.add(chat_id) _index_dirty = True def mark_quota_dirty() -> None: global _quota_dirty _quota_dirty = True async def _flush_dirty_state_once() -> None: """Тело ОДНОЙ итерации периодического сброса состояния — вынесено из _flush_dirty_state (которая теперь только спит и вызывает эту функцию в цикле) отдельной функцией, чтобы её можно было протестировать напрямую, не дожидаясь реального FLUSH_INTERVAL_SEC в тестах.""" global _dirty_chat_ids, _index_dirty, _quota_dirty, _pending_chat_deletions try: if _pending_chat_deletions: to_delete = list(_pending_chat_deletions) _pending_chat_deletions = set() del_results = await asyncio.gather( *(_delete_chat_storage_limited(cid) for cid in to_delete), return_exceptions=True, ) # ИСПРАВЛЕНО (код-ревью): раньше результат просто игнорировался — при # сбое удаление молча "терялось" (чат оставался бесхозно висеть в # хранилище навсегда, если это был единственный шанс его удалить). # Теперь неудавшиеся id возвращаются в очередь для повтора на # следующем цикле (return_exceptions=True защищает и от неожиданного # исключения, которое не было поймано внутри самой _delete_chat_storage). failed_deletes = {cid for cid, res in zip(to_delete, del_results) if res is not True} if failed_deletes: _pending_chat_deletions.update(failed_deletes) log.warning( "[state] %d удаление(й) чатов не удалось, повторю в следующем цикле: %s", len(failed_deletes), ", ".join(str(c) for c in sorted(failed_deletes)), ) if _dirty_chat_ids: to_save = list(_dirty_chat_ids) _dirty_chat_ids = set() # Каждый чат — свой независимый write; если один упадёт, это не # заденет сохранение остальных чатов из этой же пачки. attempted_ids = [cid for cid in to_save if cid in chat_state] save_results = await asyncio.gather( *(_save_chat_to_storage_limited(cid, chat_state[cid]) for cid in attempted_ids), return_exceptions=True, ) # ИСПРАВЛЕНО (найдено при код-ревью, КРИТИЧНО): раньше to_save # очищался ДО того, как запись реально прошла, а _save_chat_to_storage # сама ловила исключение и просто логировала его — наружу в gather # ничего не долетало. При транзиентном сбое Upstash (сетевой глюк, # 429 и т.п.) состояние чата (вся история диалога) молча терялось до # следующей независимой мутации этого же чата — а если это было # последнее сообщение перед долгим затишьем, данные пропадали # безвозвратно при следующем рестарте контейнера. Теперь неудавшиеся # id возвращаются обратно в _dirty_chat_ids для повтора на следующем # цикле (FLUSH_INTERVAL_SEC секунд спустя), а не теряются молча. failed_ids = {cid for cid, res in zip(attempted_ids, save_results) if res is not True} if failed_ids: _dirty_chat_ids.update(failed_ids) log.warning( "[state] %d чат(ов) не удалось сохранить в этом цикле, повторю в следующем: %s", len(failed_ids), ", ".join(str(c) for c in sorted(failed_ids)), ) if _index_dirty: _index_dirty = False await asyncio.to_thread(_save_chat_index) if _quota_dirty: _quota_dirty = False await asyncio.to_thread(save_global_quota) except Exception as exc: log.warning("[state] Периодический сброс состояния упал: %s", exc) async def _flush_dirty_state() -> None: while True: await asyncio.sleep(FLUSH_INTERVAL_SEC) await _flush_dirty_state_once() def _flush_state_now() -> None: """Синхронный финальный сброс всего "грязного" состояния — используется только при остановке процесса (main(), finally): event loop всё равно останавливается, поэтому блокирующие вызовы здесь не проблема, а вот пропустить несохранённые изменения между последним тиком периодического флаша и остановкой контейнера — проблема.""" global _dirty_chat_ids, _index_dirty, _pending_chat_deletions for cid in list(_pending_chat_deletions): _delete_chat_storage(cid) _pending_chat_deletions = set() for cid in list(_dirty_chat_ids): st = chat_state.get(cid) if st is not None: _save_chat_to_storage(cid, st) _dirty_chat_ids = set() if _index_dirty: _save_chat_index() _index_dirty = False def load_state_from_disk() -> None: load_global_quota() index_raw = None try: index_raw = _storage_read_text(CHAT_INDEX_KEY, CHAT_INDEX_FILE) except Exception as exc: log.warning("[state] Reading chat index failed, falling back to legacy combined blob: %s", exc) if index_raw is not None: # Новый формат (per-chat ключи) — индекс уже существует, читаем каждый # чат отдельно по своему ключу. try: chat_ids = json.loads(index_raw) except Exception as exc: log.warning("[state] Не удалось разобрать индекс чатов: %s", exc) chat_ids = [] loaded_count = 0 for chat_id_raw in chat_ids: try: cid = int(chat_id_raw) except Exception: continue try: raw = _storage_read_text(_chat_storage_key(cid), _chat_storage_path(cid)) except Exception as exc: log.warning("[state] Не удалось прочитать чат %s: %s", cid, exc) continue if not raw: continue try: s = json.loads(raw) except Exception as exc: log.warning("[state] Не удалось разобрать состояние чата %s: %s", cid, exc) continue _restore_single_chat(cid, s) loaded_count += 1 log.info("[state] Restored states for %d chats (per-chat storage).", loaded_count) return # ── Legacy-формат (единый блоб на все чаты, старый ключ "lumen:chat_state") ── # Индекса ещё нет — значит бот ещё ни разу не сохранял состояние в новом # per-chat формате (первый запуск после этого обновления). Читаем как раньше, # но сразу помечаем ВСЕ восстановленные чаты и индекс "грязными" (mark_state_ # dirty() без аргумента) — уже самый первый периодический флаш перепишет их # в новом per-chat формате; дальше старый общий ключ больше не читается. # Сам старый ключ/файл намеренно не удаляется автоматически — не хотим лишний # раз трогать чужие данные во время миграции, можно вычистить вручную позже. try: raw = _storage_read_text("lumen:chat_state", STATE_FILE_PATH) if not raw: return loaded = json.loads(raw) for chat_id_str, s in loaded.items(): try: cid = int(chat_id_str) except Exception: continue _restore_single_chat(cid, s) log.info( "[state] Restored states for %d chats (миграция из старого общего формата хранения — " "будут переписаны в новый per-chat формат при следующем флаше).", len(chat_state), ) mark_state_dirty() except Exception as exc: log.warning("[state] Restoring states failed: %s", exc) def get_state(chat_id: int) -> dict[str, Any]: if chat_id not in chat_state: chat_state[chat_id] = { "image_model": DEFAULT_HF_IMAGE_MODEL, "history": [], "quota": {}, "ctx": deque(maxlen=MAX_CHAT_HISTORY_LEN), "recent_media_ids": {}, "last_activity": time.monotonic(), } # Новый chat_id должен попасть в индекс (см. per-chat хранилище выше) — # иначе после рестарта его данные будут недостижимы: собственный ключ # существует, но индекс о нём не знает. _mark_new_chat_id(chat_id) else: chat_state[chat_id]["last_activity"] = time.monotonic() normalized = _normalize_legacy_image_model_id(chat_state[chat_id].get("image_model")) if normalized not in HF_IMAGE_MODELS: chat_state[chat_id]["image_model"] = DEFAULT_HF_IMAGE_MODEL elif normalized != chat_state[chat_id].get("image_model"): chat_state[chat_id]["image_model"] = normalized if len(chat_state) > MAX_CHAT_LIMIT: _prune_old_chats() return chat_state[chat_id] def _prune_old_chats() -> None: sorted_ids = sorted(chat_state.keys(), key=lambda cid: chat_state[cid].get("last_activity", 0)) to_remove = len(chat_state) - PRUNED_CHAT_TARGET removed_ids = sorted_ids[:to_remove] for cid in removed_ids: chat_state.pop(cid, None) _chat_locks.pop(cid, None) # Вытесненные чаты должны реально исчезнуть из хранилища (иначе их собственные # ключи/файлы бесхозно копятся навсегда) — ставим в очередь на удаление, # обрабатывается в _flush_dirty_state вместе с обычным сбросом. _pending_chat_deletions.update(removed_ids) mark_state_dirty() def get_chat_lock(chat_id: int) -> asyncio.Lock: lock = _chat_locks.get(chat_id) if lock is None: lock = asyncio.Lock() _chat_locks[chat_id] = lock return lock def _is_owner(user_id: int | None) -> bool: """Единая точка проверки "это владелец бота?" — используется в /logs, /stats и при гейтинге привилегированных действий в группах (см. _is_privileged_in_chat ниже). Модель/провайдер бот теперь выбирает сам (см. секцию "автоматический выбор модели"), поэтому проверка реальных названий моделей ("показывать ли Gemini/Gemma/OpenRouter владельцу") больше не нужна нигде — эти названия вообще никому не показываются, включая владельца.""" return OWNER_ID is not None and user_id is not None and user_id == OWNER_ID async def _is_privileged_in_chat(chat_type: str, chat_id: int, user_id: int | None) -> bool: """Может ли этот пользователь менять ОБЩИЕ настройки данного чата (модель для генерации изображений, сброс истории — выбор модели/провайдера для текстового чата больше не настройка чата вообще, см. "автоматический выбор модели" ниже)? В личных сообщениях у чата всего один пользователь — разрешено всегда. Владелец бота (OWNER_ID) — разрешено всегда, в любом чате. В группах/супергруппах — только создатель или администратор ЭТОЙ группы (проверяется через getChatMember, не требует особых прав у бота помимо членства в чате).""" if chat_type == ChatType.PRIVATE: return True if _is_owner(user_id): return True if user_id is None or bot is None: return False try: member = await bot.get_chat_member(chat_id, user_id) return getattr(member, "status", None) in ("creator", "administrator") except Exception as exc: log.warning("[perm] Не удалось проверить статус администратора в чате %s: %s", chat_id, exc) return False # генерация изображений (pollinations) DEFAULT_HF_IMAGE_MODEL = os.getenv("HF_IMAGE_MODEL", "flux").strip() HF_IMAGE_MODEL_PAGE_SIZE = 8 HF_IMAGE_MODELS: dict[str, dict[str, Any]] = { "flux": { "name": "FLUX Pro", "desc": "Высококачественный FLUX. Фотореализм, точное следование промпту, богатая детализация.", }, "flux-realism": { "name": "FLUX Realism", "desc": "FLUX с акцентом на гиперреализм — детализированные текстуры, естественное освещение, кинематографичность.", }, "flux-anime": { "name": "FLUX Anime", "desc": "FLUX для аниме и иллюстраций — характерные пропорции, яркие цвета, стилизация под японскую графику.", }, "turbo": { "name": "Turbo", "desc": "Быстрая дистиллированная модель. Результат за несколько секунд — для черновиков и быстрых итераций.", }, "dreamshaper": { "name": "DreamShaper", "desc": "Художественная модель для фэнтези, концепт-арта и стилизованных иллюстраций.", }, } def _hf_model_catalog() -> list[dict[str, Any]]: """Список моделей генерации изображений для клавиатуры /imgmodel — прямо из HF_IMAGE_MODELS (единственный источник правды, статический список). НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: раньше здесь был отдельный `HF_IMAGE_MODEL_CACHE` dict и `async def _hf_fetch_model_catalog()` — вестигиальные остатки более раннего дизайна, когда каталог реально динамически подтягивался из HF API. Тот динамический фетч убран (см. историю — он добавлял неизвестные модели в меню), и функция давно не делает ни одного `await` внутри, а просто пересобирает список из того же самого статического HF_IMAGE_MODELS — то есть кэшировать было уже нечего: HF_IMAGE_MODELS и так уже лежит в памяти целиком, а пересборка списка из 5 элементов не стоит отдельного кэш-слоя с ручной инвалидацией на каждом сайте вызова. Убрано вместе с самим кэшем.""" return [{"id": mid, **meta} for mid, meta in HF_IMAGE_MODELS.items()] def _imgmodel_button_text(model: dict[str, Any], current: str) -> str: prefix = "• " if model["id"] == current else "" return f"{prefix}{model['name']}" def _imgmodel_keyboard(current: str, page: int = 0) -> InlineKeyboardMarkup: all_models = _hf_model_catalog() total_pages = max(1, (len(all_models) + HF_IMAGE_MODEL_PAGE_SIZE - 1) // HF_IMAGE_MODEL_PAGE_SIZE) page = max(0, min(page, total_pages - 1)) start = page * HF_IMAGE_MODEL_PAGE_SIZE chunk = all_models[start:start + HF_IMAGE_MODEL_PAGE_SIZE] rows: list[list[InlineKeyboardButton]] = [] for i in range(0, len(chunk), 2): row = [] for item in chunk[i:i + 2]: row.append(InlineKeyboardButton(text=_imgmodel_button_text(item, current), callback_data=f"imgmodel:set:{page}:{item['id']}")) rows.append(row) nav_row: list[InlineKeyboardButton] = [] if page > 0: nav_row.append(InlineKeyboardButton(text="< Назад", callback_data=f"imgmodel:page:{page - 1}")) if page + 1 < total_pages: nav_row.append(InlineKeyboardButton(text="Дальше >", callback_data=f"imgmodel:page:{page + 1}")) if nav_row: rows.append(nav_row) return InlineKeyboardMarkup(inline_keyboard=rows) async def _pollinations_generate(model_name: str, prompt: str) -> bytes: """Бесплатная генерация через Pollinations.ai — не требует авторизации.""" session = await _get_http_session() encoded = urllib.parse.quote(prompt[:600], safe="") url = ( f"https://image.pollinations.ai/prompt/{encoded}" f"?width=1024&height=1024&model={model_name}&nologo=true&enhance=false" ) async with session.get(url, timeout=aiohttp.ClientTimeout(total=90)) as resp: if resp.status == 200: ctype = (resp.headers.get("Content-Type") or "").lower() body = await resp.read() if body and (ctype.startswith("image/") or body[:4] in (b"\x89PNG", b"\xff\xd8\xff", b"RIFF", b"GIF8")): return body raise RuntimeError(f"Pollinations вернул не-изображение: {ctype}") raise RuntimeError(f"Pollinations.ai HTTP {resp.status}") async def _hf_text_to_image(model_id: str, prompt: str) -> bytes: # Pollinations.ai — единственный провайдер генерации изображений (HF Inference # API-ветка убрана: старые HF-модели регулярно устаревали на стороне провайдера, # см. историю: FLUX.1-dev вернул 410 Gone). Раз провайдер всего один, отдельная # "pollinations:" приставка на каждом ключе HF_IMAGE_MODELS была лишней — # model_id и есть имя модели Pollinations как есть. if model_id not in HF_IMAGE_MODELS: raise ValueError(f"Неизвестная модель генерации изображений: {model_id}") return await _pollinations_generate(model_id, prompt) # метаданные и скачивание медиа def _is_gemini_supported_mime(mime: str) -> bool: m = mime.lower() if m in ("image/png", "image/jpeg", "image/webp", "image/heic", "image/heif", "image/gif"): return True if m in ("audio/mp3", "audio/wav", "audio/ogg", "audio/aac", "audio/flac", "audio/mp4", "audio/m4a", "audio/mpeg", "audio/x-m4a"): return True if m in ("video/mp4", "video/mpeg", "video/mov", "video/avi", "video/flv", "video/mpg", "video/webm", "video/wmv", "video/quicktime"): return True if m in ( "application/pdf", "text/plain", "text/html", "text/css", "text/javascript", "application/x-javascript", "text/csv", "text/markdown", "text/xml", "application/xml", "application/json", "text/x-python", "application/x-python-code" ): return True return False def _sanitize_mime_type(file_path: str | None, mime: str | None, default_fallback: str = "application/octet-stream") -> str: m = (mime or "").strip().lower() if not m or m in ("application/octet-stream", "binary/oct-stream", "application/x-binary", "octet/stream"): if file_path: guessed, _ = mimetypes.guess_type(file_path) if guessed: return guessed.lower() if file_path: ext = Path(file_path).suffix.lower() ext_map = { ".jpg": "image/jpeg", ".jpeg": "image/jpeg", ".png": "image/png", ".webp": "image/webp", ".gif": "image/gif", ".mp4": "video/mp4", ".mov": "video/quicktime", ".m4v": "video/x-m4v", ".avi": "video/x-msvideo", ".mp3": "audio/mpeg", ".ogg": "audio/ogg", ".oga": "audio/ogg", ".opus": "audio/ogg", ".m4a": "audio/mp4", ".wav": "audio/wav", ".pdf": "application/pdf", ".txt": "text/plain", ".csv": "text/csv", ".json": "application/json", ".html": "text/html", ".htm": "text/html", ".xml": "text/xml" } if ext in ext_map: return ext_map[ext] return default_fallback if "voice" in m or m == "audio/ogg": return "audio/ogg" if m == "video/quicktime": return "video/quicktime" return m def _media_file_id_and_mime(source: Any) -> tuple[str, str, str]: file_id, mime, filename = "", "", "" if source is None: return file_id, mime, filename for key in ("file_id", "id"): v = getattr(source, key, None) if not isinstance(source, dict) else source.get(key) if isinstance(v, str) and v.strip(): file_id = v.strip() break mt = getattr(source, "mime_type", None) if not isinstance(source, dict) else source.get("mime_type") if isinstance(mt, str): mime = mt.strip() nm = getattr(source, "file_name", None) if not isinstance(source, dict) else source.get("file_name") if isinstance(nm, str): filename = nm.strip() class_name = type(source).__name__ if not mime: if class_name == "PhotoSize": mime = "image/jpeg" elif class_name == "Sticker": mime = "image/webp" elif class_name == "Voice": mime = "audio/ogg" elif class_name == "VideoNote": mime = "video/mp4" elif class_name == "Animation": mime = "video/mp4" return file_id, mime, filename def _mime_suffix(mime: str, filename: str = "") -> str: if filename: suffix = Path(filename).suffix if suffix: return suffix m = (mime or "").lower() if m.startswith("image/"): sub = m.split("/", 1)[1] return {"jpeg": ".jpg", "jpg": ".jpg", "png": ".png", "gif": ".gif", "webp": ".webp"}.get(sub, f".{sub}") if m.startswith("audio/"): return ".mp3" if m.startswith("video/"): return ".mp4" return ".bin" # загрузка медиа def _msg_media_source(message: Any) -> Any | None: for attr in ("photo", "video", "animation", "video_note", "voice", "audio", "document", "sticker"): val = getattr(message, attr, None) if not val: continue if attr == "photo" and isinstance(val, list): return val[-1] if val else None return val return None async def _download_telegram_file_bytes(file_id: str, *, timeout: float | None = None, retries: int = 1) -> tuple[bytes, str]: # Один ретрай с короткой паузой — раньше здесь не было НИКАКОГО повторного # обращения (в отличие от _tg_call, у которого есть свой параметр retries), # поэтому одна-единственная транзиентная заминка прокси (не-JSON/обрыв ровно # на getMe/getFile, см. _looks_like_proxy_garbage) насовсем валила скачивание # медиа. Именно эта функция стояла за инцидентом "[media] Download media failed # ... NOT_FOUND" в логах этой сессии — единичный сбой прокси не должен означать # "пользователь прислал фото/видео, а бот его просто не увидел". last_exc: Exception | None = None for attempt in range(retries + 1): try: file = await asyncio.wait_for(bot.get_file(file_id), timeout=TELEGRAM_GET_FILE_TIMEOUT) file_path = getattr(file, "file_path", None) or getattr(file, "path", None) if not file_path: raise RuntimeError("File path is empty") session = await _get_telegram_session() url = f"{TELEGRAM_API_BASE_URL}/file/bot{BOT_TOKEN}/{file_path}" async with session.get(url, timeout=timeout or TELEGRAM_MEDIA_TIMEOUT) as resp: resp.raise_for_status() data = await resp.read() mime = resp.headers.get("Content-Type", "application/octet-stream").split(";", 1)[0].strip() real_mime = _sanitize_mime_type(file_path, mime) return data, real_mime except Exception as exc: last_exc = exc if attempt < retries: log.warning("[media] Попытка %d/%d скачать file_id %s не удалась, повтор через 0.5с: %s", attempt + 1, retries + 1, file_id, exc) await asyncio.sleep(0.5) exc_str = str(last_exc) or repr(last_exc) or type(last_exc).__name__ if BOT_TOKEN: exc_str = exc_str.replace(BOT_TOKEN, "") raise RuntimeError(f"Network error in download_telegram_file_bytes: {exc_str}") from None def _save_media_to_history(source: Any, state: dict[str, Any], user_id: int | None) -> None: file_id, mime, _ = _media_file_id_and_mime(source) if not file_id or user_id is None: return buckets: dict[str, deque] = state.setdefault("recent_media_ids", {}) key = str(user_id) recent = buckets.setdefault(key, deque(maxlen=MAX_MEDIA_RECENT_IDS)) if not recent or recent[-1][0] != file_id: recent.append((file_id, mime or "application/octet-stream")) async def _download_message_attachment_to_tmp(source: Any) -> tuple[str, str, str] | None: file_id, mime, filename = _media_file_id_and_mime(source) if not file_id: return None suffix = _mime_suffix(mime, filename) fd, tmp_path = tempfile.mkstemp(prefix="tg_media_", suffix=suffix) os.close(fd) try: data, real_mime = await _download_telegram_file_bytes(file_id) final_mime = _sanitize_mime_type(filename or "", mime) if final_mime == "application/octet-stream" or not final_mime: final_mime = _sanitize_mime_type(None, real_mime) if final_mime == "application/octet-stream" or not final_mime: final_mime = mime or "application/octet-stream" with open(tmp_path, "wb") as h: h.write(data) return tmp_path, final_mime, filename or Path(tmp_path).name except Exception: with contextlib.suppress(FileNotFoundError): os.unlink(tmp_path) raise async def _fetch_media(file_id: str, mime: str) -> tuple[bytes, str] | None: if not file_id: return None try: data, real_mime = await _download_telegram_file_bytes(file_id) final_mime = _sanitize_mime_type(None, mime) if final_mime == "application/octet-stream": final_mime = _sanitize_mime_type(None, real_mime) return data, final_mime except Exception as exc: log.warning("[media] Download media failed for file_id %s: %s", file_id, exc) return None def _ensure_prompt_text(text: str | None, mime: str) -> str: s = (text or "").strip() if s: return s m = mime.lower() if m.startswith("image/"): return "Подробно опиши, что изображено на картинке." if m.startswith("video/") or m == "video/quicktime": return "Подробно опиши происходящее на этом видео." if m.startswith("audio/") or m == "audio/ogg" or m == "audio/mpeg" or m == "audio/mp3" or "voice" in m: return "Прослушай и подробно опиши, что на этой аудиозаписи, или кратко перескажи ее содержание." if m == "image/gif" or "animation" in m: return "Подробно опиши происходящее на этой анимации." return "Проанализируй и подробно опиши содержимое этого вложения." # учёт квот def _quota_entry(state: Any, provider: str, model_id: str) -> dict[str, Any]: _reset_quota_if_new_day() sub = GLOBAL_QUOTA.setdefault(provider, {}) return sub.setdefault(model_id, {"used": 0, "remaining": None, "limit": None, "exhausted_at": None}) def _mark_quota_exhausted(provider: str, model_id: str) -> None: """Записывает момент, когда API реально вернул 429/RESOURCE_EXHAUSTED для модели. Используется, потому что _record_quota_usage инкрементирует "used" только при успешном ответе — без этого счётчик мог годами показывать 0, даже если все запросы к модели упирались в реальный лимит на стороне Google/OpenRouter.""" e = _quota_entry(None, provider, model_id) e["exhausted_at"] = time.time() mark_quota_dirty() def _record_quota_usage(state: Any, provider: str, model_id: str, remaining: int | None = None, limit: int | None = None) -> None: e = _quota_entry(None, provider, model_id) e["used"] = int(e.get("used") or 0) + 1 e["exhausted_at"] = None if remaining is not None: e["remaining"] = max(0, remaining) if limit is not None: e["limit"] = max(0, limit) mark_quota_dirty() # openrouter api # # TEXT_MODEL_ORDER / _OR_MODEL_HEALTH / _ROUTER_EXCLUDED_OR_MODELS / # _check_temporary_free_models_expiry вынесены в lumen_router_config.py (см. # импорт рядом с GEMINI_MODELS выше по файлу) — здесь остаётся только код, # который реально ХОДИТ в OpenRouter API (OpenRouterAPIError/_or_request/ # ask_openrouter_*/_or_chat_completion_with_fallback и т.д.). # ─────────────────── защита от утечки провайдера/модели и промт-инъекций ─────────────────── # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: детекторы утечки идентичности (_detect_identity_leak/ # _scrub_identity_leak/_detect_injected_payload_echo) и входной префильтр промт- # инъекций (_looks_like_injection_probe) вынесены в lumen_security.py — чистые # функции над строками (плюс регэкспы/константы), не зависящие от Telegram/рантайм- # состояния бота. Импортируются напрямую — публичные имена и поведение (включая # логирование через тот же логгер "bot", см. lumen_security.py) не изменились. # Только реально используемые здесь имена импортируются явно — регэкспы # (`_IDENTITY_LEAK_RE`/`_INJECTION_PROBE_RE`/`_INJECTED_PAYLOAD_ECHO_RE` и # составляющие их `_LEAK_BRAND_TOKENS`/`_LEAK_LITERAL_STRINGS`) нужны только # самим детекторам внутри lumen_security.py, а не коду bot.py. from lumen_security import ( _LEAK_SCAN_TAIL_CHARS, _leak_scan_window, _IDENTITY_LEAK_FALLBACK, _INJECTED_PAYLOAD_ECHO_FALLBACK, _detect_injected_payload_echo, _detect_identity_leak, _scrub_identity_leak, _INJECTION_PROBE_REPLY, _looks_like_injection_probe, ) # См. пояснение про __all__ у первого блока (lumen_formatting) выше — используется # только как `bot._LEAK_SCAN_TAIL_CHARS` в тестах (сама логика окна сканирования — # внутри `_leak_scan_window`, который уже используется по-настоящему). __all__ += ["_LEAK_SCAN_TAIL_CHARS"] class OpenRouterAPIError(RuntimeError): def __init__(self, message: str, status_code: int | None = None, payload: Any = None) -> None: super().__init__(message) self.status_code = status_code self.payload = payload async def _or_request(path: str, method: str = "GET", *, json_body: dict | None = None) -> Any: if not OPENROUTER_API_KEY: raise OpenRouterAPIError("OPENROUTER_API_KEY не задан") headers = { "Authorization": f"Bearer {OPENROUTER_API_KEY}", "HTTP-Referer": OPENROUTER_HTTP_REFERER, "X-OpenRouter-Title": OPENROUTER_TITLE, } if json_body is not None: headers["Content-Type"] = "application/json" session = await _get_http_session() url = f"{OPENROUTER_BASE_URL}/{path.lstrip('/')}" try: async with session.request( method.upper(), url, headers=headers, json=json_body, # ponytail: было захардкожено total=12.0, независимо от ROUTE_MODEL_ # TIMEOUT_SEC (22с по умолчанию) — модель могла получить меньше времени, # чем задокументированный бюджет одной попытки, и валиться таймаутом # раньше, чем должна была (см. аудит моделей 2 августа 2026). timeout=aiohttp.ClientTimeout(total=ROUTE_MODEL_TIMEOUT_SEC, connect=10.0) ) as resp: if resp.status >= 400: payload = await resp.json(content_type=None) msg = payload.get("error", {}).get("message") or f"HTTP {resp.status}" raise OpenRouterAPIError(msg, status_code=resp.status, payload=payload) return await resp.json(content_type=None) except OpenRouterAPIError: raise except Exception as exc: # str(exc) часто пуст для таймаутов/CancelledError-обёрток (см. реальный # найденный случай в логах: "Сетевая ошибка OpenRouter: " без единой # детали) — тогда используем repr/имя класса, чтобы в логах вообще было # видно, что произошло, а не пустая строка. exc_str = str(exc) or repr(exc) or exc.__class__.__name__ raise OpenRouterAPIError(f"Сетевая ошибка OpenRouter: {exc_str}") from exc def _or_extract_text(data: Any) -> str: if isinstance(data, str): return data.strip() if isinstance(data, dict): return _or_extract_text(data.get("content") or data.get("text") or "") if isinstance(data, list) and data: return "".join(_or_extract_text(i) for i in data) return "" def _is_account_wide_or_rate_limit(text: str) -> bool: """"free-models-per-day" — это лимит на весь аккаунт OpenRouter целиком (см. реальный найденный случай: "Rate limit exceeded: free-models-per-day. Add 10 credits to unlock 1000 free model requests per day"), а не на одну конкретную модель. Раньше при этой ошибке бот всё равно честно перебирал ВСЕ 7-8 кандидатов цепочки по очереди — и получал одну и ту же ошибку на каждом, иногда суммарно теряя больше минуты (реальный случай в логах — 168 секунд) только на то, чтобы наконец сдаться и попробовать Gemini. Если видим этот текст — сразу прекращаем всю цепочку OpenRouter, а не тратим время на заведомо обречённые попытки остальных моделей.""" low = text.lower() return "free-models-per-day" in low async def _or_chat_completion_with_fallback( messages: list[dict], trial_models: list[str], primary_model_id: str, *, attempts_per_model: int = 1, deadline: float | None = None, ) -> tuple[str, str]: """Общий цикл fallback по цепочке моделей для запросов к OpenRouter chat/completions. Раньше это был почти идентичный код, продублированный внутри ask_openrouter_text И ask_openrouter_multimodal — риск, что при будущей правке (например, добавлении новой категории временной ошибки) кто-то поправит только одну из двух копий и они молча разойдутся. messages[0] должен быть системным сообщением — его content переписывается под каждую пробуемую модель (т.к. у разных моделей разный get_system_prompt). attempts_per_model по умолчанию — 1 (без ретраев одной и той же модели): раньше было 2 попытки с задержкой 1.5с между ними, что при массовой нестабильности одной модели ощутимо замедляло весь маршрут (та же причина, что и убранные ретраи в ask_gemini — см. комментарий там). Одна ошибка — сразу следующий кандидат по цепочке. Возвращает (answer, реально_использованная_модель) при успехе. Если ни одна модель из trial_models не дала ответ — поднимает последнее пойманное исключение, либо RouteBudgetExceededError, если общий бюджет времени маршрута закончился раньше, чем дошла очередь до оставшихся кандидатов.""" last_exc: Exception | None = None tried: list[str] = [] for model_trial in trial_models: if deadline is not None and time.monotonic() > deadline: log.warning("[or] Бюджет времени маршрута исчерпан перед моделью %s. Испробовано: %s", model_trial, ", ".join(tried) or "ничего") raise RouteBudgetExceededError(tried) tried.append(model_trial) messages[0]["content"] = get_system_prompt(model_trial) for attempt in range(attempts_per_model): try: payload = {"model": model_trial, "messages": messages, "stream": False} resp = await _or_request("chat/completions", "POST", json_body=payload) choices = resp.get("choices") or [] answer = "" if choices: answer = _or_extract_text(choices[0].get("message") or "") answer = answer.strip() or "Empty response" answer = _scrub_identity_leak(answer, source=f"or_chat_completion:{model_trial}") log.info("[or] Успешный ответ от модели %s (primary=%s, испробовано моделей: %d)", model_trial, primary_model_id, len(tried)) return answer, model_trial except Exception as exc: last_exc = exc err_text = str(exc).lower() if _is_account_wide_or_rate_limit(err_text): log.warning( "[or] Обнаружен лимит на весь аккаунт OpenRouter (free-models-per-day) на модели %s — " "прекращаю перебор оставшихся кандидатов цепочки, дальше они гарантированно откажут тем же самым.", model_trial, ) raise # НАЙДЕНО ПРИ АУДИТЕ ТЕХДОЛГА: раньше здесь была третья независимая # ad hoc классификация ошибки (ручной разбор подстрок "429"/"rate # limit"/"500"/"502"/... и т.п.) — отдельная от _classify_model_error # (уже используемого в _gemini_error_msg/_or_error_msg) и от такой же # по смыслу классификации внутри ask_gemini. Общий классификатор # покрывает то же множество "стоит ли повторить попытку" не хуже. # Поведение при attempts_per_model=1 (единственное реальное # использование сейчас) не меняется: is_retriable_same_model влияет # только на повторные попытки ОДНОЙ и той же модели внутри # attempts_per_model, а не на переход к следующей модели по цепочке # (тот происходит естественным продолжением внешнего for ниже). is_timeout = isinstance(exc, (asyncio.TimeoutError, TimeoutError)) or not err_text kind = _classify_model_error(_error_status(exc, err_text), err_text) is_retriable_same_model = is_timeout or kind in ("rate_limit", "unavailable", "other") if not is_retriable_same_model: break log.warning("[or] Model %s failed: %s. Switching to next candidate...", model_trial, str(last_exc) or last_exc.__class__.__name__) if last_exc: raise last_exc raise RuntimeError("Не удалось получить ответ ни от одной модели-кандидата.") async def ask_openrouter_text(chat_id: int, user_text: str, model_chain: list[str], *, deadline: float | None = None) -> str: state = get_state(chat_id) ctx = state.get("ctx", deque()) full = "" if ctx: full = "Фон разговора:\n" + "\n".join(ctx) + "\n\nТекущий вопрос: " full += user_text history = state.setdefault("history", []) # model_chain строится роутером (см. _build_route) — здесь только убираем # дубликаты, сохраняя порядок приоритета, заданный роутером. Пустой model_chain # в норме не должен случаться (_build_route всегда возвращает непустой маршрут), # это последняя страховка "на всякий случай". ИСПРАВЛЕНО (24 июля 2026): раньше # здесь запасным вариантом стоял meta-llama/llama-3.3-70b-instruct:free — та же # модель, что подтверждённо снята провайдером с бесплатного тира (см. README — # повторяющиеся HTTP 404 "unavailable for free") и по этой причине уже исключена # из _OR_LIGHT_ORDER/_OR_HEAVY_ORDER. Оставлять её единственным запасным # вариантом здесь означало тот же самый риск с другой стороны — заменено на # первую модель актуального _OR_LIGHT_ORDER (единый источник правды). trial_models = list(dict.fromkeys(model_chain)) or [_OR_LIGHT_ORDER[0]] primary_model_id = trial_models[0] messages: list[dict] = [{"role": "system", "content": get_system_prompt(primary_model_id)}] messages.extend(history) messages.append({"role": "user", "content": full}) answer, model_trial = await _or_chat_completion_with_fallback(messages, trial_models, primary_model_id, deadline=deadline) # В общую историю пишем ЧИСТЫЙ текст пользователя (без служебного префикса # "Фон разговора") — эту же историю теперь читает и Gemini (см. SHARED_HISTORY_ # MAX_LEN), и разовый ephemeral-контекст группового чата не должен там оседать. history.append({"role": "user", "content": user_text}) history.append({"role": "assistant", "content": answer}) ctx.clear() if len(history) > SHARED_HISTORY_MAX_LEN: del history[:-SHARED_HISTORY_MAX_LEN] _record_quota_usage(state, "openrouter", model_trial) return answer async def ask_openrouter_multimodal( chat_id: int, user_text: str, media_tuple: tuple[bytes, str], media_filename: str, model_chain: list[str], *, deadline: float | None = None, ) -> str: state = get_state(chat_id) b64 = base64.b64encode(media_tuple[0]).decode("utf-8") img_url = f"data:{media_tuple[1]};base64,{b64}" ctx = state.get("ctx", deque()) full_text = user_text if ctx: full_text = "Фон разговора в чате (для контекста, не обращение к тебе):\n" + "\n".join(ctx) + "\n\nТекущий вопрос/сообщение: " + user_text history = state.setdefault("history", []) trial_models = list(dict.fromkeys(model_chain)) or ["nvidia/nemotron-nano-12b-v2-vl:free"] primary_model_id = trial_models[0] messages: list[dict] = [{"role": "system", "content": get_system_prompt(primary_model_id)}] messages.extend(history) messages.append({ "role": "user", "content": [ {"type": "text", "text": full_text}, {"type": "image_url", "image_url": {"url": img_url}} ] }) answer, model_trial = await _or_chat_completion_with_fallback(messages, trial_models, primary_model_id, deadline=deadline) history.append({"role": "user", "content": user_text}) history.append({"role": "assistant", "content": answer}) if len(history) > SHARED_HISTORY_MAX_LEN: del history[:-SHARED_HISTORY_MAX_LEN] ctx.clear() _record_quota_usage(state, "openrouter", model_trial) return answer # скачивание тикток # ─────────────────── локализация подписи "оригинальный звук" ─────────────────── # TikWM отдаёт название "оригинального звука" (music_info.title) либо на английском # ("original sound"), либо на языке автора ИСХОДНОГО видео, породившего этот звук — # то есть никак не связано с языком человека, который прислал ссылку В НАШ бот. # Telegram передаёт язык интерфейса КАЖДОГО отправителя в message.from_user. # language_code (IETF-тег вроде "ru"/"uk"/"en-US") — это язык, который сам человек # выбрал в настройках Telegram, приходит с любым сообщением без доп. разрешений и # не требует ничего от пользователя. Используем именно его (а не raw_music_title # от TikWM), чтобы подпись "оригинальный звук"/"original sound"/... совпадала с # языком ТОГО, кто прислал конкретную ссылку — даже в группе, где разные участники # могут иметь разный язык интерфейса. # # Список языков — приоритет отдан региону СНГ/ближнего зарубежья (основная # аудитория бота), плюс крупные европейские и соседние языки. Любой язык, которого # нет в словаре, тихо откатывается на английский вариант (нейтральный и понятный # дефолт, а не гадание с неизвестным алфавитом). _ORIGINAL_SOUND_LABELS: dict[str, str] = { "ru": "Оригинальный звук", "uk": "Оригінальний звук", "be": "Арыгінальны гук", # НАЙДЕНО ПО ВОПРОСУ ВЛАДЕЛЬЦА: единая схема регистра для всех языков — с # большой буквы у первого слова (Sentence case), как и положено названию # трека (это поле идёт в MP3-тег "название", т.е. в тот же слот, где обычно # показывается настоящее название песни — оно тоже всегда с большой буквы). # Раньше английский вариант был строчным ("original sound") по инерции от # того, как сам TikTok показывает его в своём интерфейсе — но раз это теперь # НАША подпись, а не дословная копия чужого UI, приводим её к тому же виду, # что и остальные языки, а не оставляем единственным исключением. "en": "Original sound", "pl": "Oryginalny dźwięk", "de": "Originalton", "es": "Sonido original", "fr": "Son original", "it": "Audio originale", "pt": "Som original", "tr": "Orijinal ses", "kk": "Түпнұсқа дыбыс", "uz": "Original tovush", "az": "Orijinal səs", "ka": "ორიგინალური ხმა", "hy": "Օրիգինալ ձայն", "ky": "Түпнуска үн", "ar": "الصوت الأصلي", } _ORIGINAL_SOUND_LABEL_DEFAULT = _ORIGINAL_SOUND_LABELS["en"] def _original_sound_label(language_code: str | None) -> str: """Возвращает локализованную подпись "оригинальный звук" по IETF-коду языка (например, из message.from_user.language_code). Код языка может приходить с региональным уточнением (например "en-US", "pt-BR") — берём только первичный подтег до дефиса. Неизвестный/отсутствующий код — тихий откат на английский.""" if not language_code: return _ORIGINAL_SOUND_LABEL_DEFAULT primary = language_code.split("-", 1)[0].strip().lower() return _ORIGINAL_SOUND_LABELS.get(primary, _ORIGINAL_SOUND_LABEL_DEFAULT) async def _send_tiktok_music(session, media_data: dict, message: Message, author: str, headers: dict) -> None: music_url = media_data.get("music") if not music_url: return try: music_bytes = await _download_url_bin(session, music_url, headers=headers) if not music_bytes: return music_info = media_data.get("music_info") or {} raw_music_title = music_info.get("title") or "Музыка из TikTok" raw_music_author = music_info.get("author") or author # юзернейм и никнейм автора ВИДЕО без @ — нужны и для performer_name, и # для очистки заголовка ниже (TikTok иногда подставляет их в заголовок # безымянного звука вместо настоящего названия). author_nick = media_data.get("author", {}).get("nickname") or "" author_uniq = media_data.get("author", {}).get("unique_id") or "" author_uniq_clean = author_uniq.lstrip("@") # проверяем, оригинальный ли это звук m_title_lower = raw_music_title.lower() mentions_generic_phrase = "оригинальный звук" in m_title_lower or "original sound" in m_title_lower # НАЙДЕНО ПО ЖИВОМУ ТЕСТИРОВАНИЮ (реальный найденный регресс, часть 2): TikTok # разрешает автору дать "оригинальному звуку" СОБСТВЕННОЕ название при публикации # видео (см. реальный пример: TikTok показывает такой звук как "Оригинальный # звук: Night, Blooming Jasmine." на его собственной странице звука) — при этом # TikWM всё равно присылает raw_music_title с префиксом "original sound - "/ # "оригинальный звук - " ПЕРЕД настоящим названием, а не одно только настоящее # название. Поэтому буквальное совпадение фразы "original sound" в заголовке — # это ещё НЕ финальный признак "звук совсем безымянный": вырезаем саму фразу # (и, если она там же, ник/юзернейм автора ВИДЕО — TikTok в ДЕЙСТВИТЕЛЬНО # безымянном случае подставляет в заголовок именно его) и смотрим, остаётся ли # после этого что-то ЕЩЁ. Если да — это настоящее, осмысленное название звука, # которое нужно показать как есть, а не подменять generic-подписью. residual_title = raw_music_title for _phrase in ("оригинальный звук", "original sound"): residual_title = re.sub(re.escape(_phrase), "", residual_title, flags=re.IGNORECASE) residual_title = residual_title.strip(" \t-–—:") if author_nick: residual_title = re.sub(re.escape(author_nick), "", residual_title, flags=re.IGNORECASE).strip(" \t-–—:") if author_uniq_clean: residual_title = re.sub(re.escape(author_uniq_clean), "", residual_title, flags=re.IGNORECASE).strip(" \t-–—:") is_original_sound = mentions_generic_phrase and not residual_title # Диагностика: пока эта эвристика не "обкатана" на достаточном числе реальных # случаев, полезно видеть в /logs исходные поля TikWM целиком при каждом # решении — это то самое "сначала факты, потом фикс" вместо повторной догадки. log.info( "[tiktok-music][diag] raw_title=%r raw_author=%r cover=%r author_avatar=%r " "residual_title=%r -> is_original_sound=%s", raw_music_title, raw_music_author, music_info.get("cover"), media_data.get("author", {}).get("avatar"), residual_title, is_original_sound, ) if is_original_sound: # в исполнителях — юзернейм без @ performer_name = author_uniq_clean if author_uniq_clean else raw_music_author # Вместо генерации/очистки сырого raw_music_title от TikWM сразу подставляем # перевод, локализованный под язык интерфейса Telegram ИМЕННО отправителя # этой конкретной ссылки (см. _original_sound_label выше) — это единственный # способ показать подпись на "его" языке, раз сам TikTok эту связь не даёт: # raw_music_title зависит от языка автора исходного видео, а не от языка # человека, приславшего ссылку в наш бот. sender_language_code = message.from_user.language_code if message.from_user else None cleaned_title = _original_sound_label(sender_language_code) else: # Либо обычный именованный трек с автором (раньше он всегда попадал только # сюда), либо "оригинальный звук" с собственным названием (см. комментарий # выше) — в обоих случаях реальное название важнее generic-подписи. Если # TikWM прислал raw_music_title с префиксом "original sound - "/"оригинальный # звук - " перед настоящим названием, используем уже очищенный остаток; # иначе (обычный трек без такого префикса) оставляем raw_music_title как есть. cleaned_title = residual_title if (mentions_generic_phrase and residual_title) else raw_music_title performer_name = raw_music_author # достаём обложку трека cover_url = music_info.get("cover") or music_info.get("avatar") or media_data.get("author", {}).get("avatar") cover_bytes = None if cover_url: try: cover_bytes = await _download_url_bin(session, cover_url, headers=headers) except Exception as e: log.warning("[tiktok] failed to download cover image: %s", e) # готовим превью thumbnail_file = None if cover_bytes: thumbnail_file = BufferedInputFile(cover_bytes, filename="cover.jpg") # вшиваем метаданные и обложку в MP3 перед отправкой — чтобы теги видели и другие плееры tagged_music_bytes = music_bytes try: with tempfile.TemporaryDirectory() as tmp_dir: tmp_mp3_path = os.path.join(tmp_dir, "music.mp3") with open(tmp_mp3_path, "wb") as f: f.write(music_bytes) _write_mp3_tags(tmp_mp3_path, cleaned_title, performer_name, cover_bytes) if os.path.exists(tmp_mp3_path) and os.path.getsize(tmp_mp3_path) > 0: with open(tmp_mp3_path, "rb") as f: tagged_music_bytes = f.read() except Exception as tag_err: log.warning("[tiktok] failed to write embedded tags to MP3: %s", tag_err) await bot.send_audio( chat_id=message.chat.id, audio=BufferedInputFile(tagged_music_bytes, filename=f"{cleaned_title[:60]}.mp3"), title=cleaned_title, performer=performer_name, thumbnail=thumbnail_file, reply_to_message_id=message.message_id ) except Exception as e: log.warning("[tiktok] failed to send music: %s", e) def _chunk_tiktok_media_items(items: list, chunk_size: int = TELEGRAM_MEDIA_GROUP_CHUNK) -> list[list]: """Разбивает список медиа-элементов слайдшоу на группы для sendMediaGroup. НАЙДЕНО ПРИ ПОВТОРНОЙ РЕВИЗИИ (КРИТИЧНО): у Telegram Bot API `sendMediaGroup` жёсткое требование — от 2 до 10 элементов НА ОДИН вызов, а не просто "не больше 10". Наивное разбиение "по chunk_size без остатка" (см. предыдущую версию этого кода) даёт хвостовую группу РОВНО из ОДНОГО элемента всякий раз, когда общее число элементов даёт остаток 1 при делении на chunk_size (11, 21, 31 элемент и т.п. — а слайдшоу TikTok реально может состоять из любого числа слайдов вплоть до 35, так что это не гипотетический случай). Такой вызов Telegram гарантированно отклоняет как невалидный — причём это произошло бы уже ПОСЛЕ того, как предыдущие группы успешно отправились, то есть пользователь получил бы часть слайдшоу и затем непонятную ошибку. Если наивное разбиение даёт хвост из 1 элемента — "занимаем" один элемент у предпоследней группы, превращая последние две группы из (chunk_size, 1) в (chunk_size - 1, 2). Единственный элемент целиком (0 или 1 элементов на входе) эта функция не обрабатывает — такие случаи вызывающий код (handle_tiktok) должен отправлять напрямую через send_photo/send_video, а не через эту функцию/sendMediaGroup вообще.""" if not items: return [] chunks = [items[i:i + chunk_size] for i in range(0, len(items), chunk_size)] if len(chunks) >= 2 and len(chunks[-1]) == 1: borrowed = chunks[-2].pop() chunks[-1].insert(0, borrowed) return chunks def _looks_like_video_bytes(data: bytes) -> bool: """Определяет, что скачанный файл — это видео (MP4/MOV/ISO base media file format), а не статичная картинка, по магическим байтам начала файла. НАЙДЕНО ПРИ РЕВИЗИИ: TikTok официально разрешает комбинировать в одном слайдшоу-посте (photo mode) обычные статичные фото-слайды И короткие видео-слайды (TikTok сам называет это "combine photo and video slides"). TikWM отдаёт URL такого видео-слайда в том же списке `images`, что и обычные фото — без явного признака "это видео", и Content-Type в ответе CDN для таких слайдов тоже не всегда достоверен. Раньше такой URL молча скачивался и оборачивался в InputMediaPhoto — в лучшем случае Telegram показывал статичный кадр вместо реального движения слайда (то, что пользователь называет TikTok-'живым фото'), в худшем — вовсе не мог корректно отрендерить не-JPEG/ PNG/WEBP байты как фото. Проверяем сигнатуру ISO base media file format ("ftyp" на смещении 4 байта) — это надёжный и стандартный способ отличить MP4/MOV-контейнер от растрового изображения без сторонних библиотек, не зависящий от того, как именно TikWM называет поле в JSON.""" return len(data) >= 12 and data[4:8] == b"ftyp" def _slideshow_slide_urls(media_data: dict, images_to_fetch: list[str]) -> list[str]: """Для каждого слайда слайдшоу возвращает URL, который реально стоит скачать — предпочитая `live_images[i]` вместо `images[i]`, если TikWM отдал непустую запись на этой позиции. НАЙДЕНО (по логам диагностики) и ПОДТВЕРЖДЕНО на реальных постах: у ответа TikWM для фото-поста ЕСТЬ отдельное поле `live_images` помимо обычного `images`. Прежняя эвристика (см. историю — пробовала верхнеуровневые `play`/`hdplay`) была основана на неверном предположении: для фото-постов эти поля указывают НЕ на видео, а на ту же самую фоновую музыку, что и поле `music` (реальный найденный URL содержал `mime_type=audio_mpeg` на домене `...music.tiktokcdn...`), поэтому убрана целиком. `images[]` всегда отдаёт статичные `...~tplv-photomode-image.jpeg` кадры — то есть настоящую "живую" (двигающуюся) версию слайда, если она есть у этого поста, даёт именно `live_images`. ПОДТВЕРЖДЕНО РЕАЛЬНЫМИ ТЕСТАМИ (см. /logs с реальных постов): позиционное соответствие `live_images[i]` <-> `images[i]` верно — например, для поста с 2 слайдами, где только один реально "живой", `live_images` пришёл как `[None, ""]` (ровно на позиции живого слайда), и итоговый детект (`_looks_like_video_bytes` после скачивания) корректно показал "1 из 2 слайдов — видео". Для постов, где живые оба слайда или только один из одного — тоже совпало 1-в-1. Пустая/отсутствующая запись на позиции означает "этот слайд не живой, обычное статичное фото" — на этот случай функция просто продолжает использовать `images[i]`.""" live_images = media_data.get("live_images") if not isinstance(live_images, list): return images_to_fetch resolved: list[str] = [] for idx, fallback_url in enumerate(images_to_fetch): live_url = live_images[idx] if idx < len(live_images) else None resolved.append(live_url if isinstance(live_url, str) and live_url.strip() else fallback_url) return resolved def _tiktok_video_candidates(media_data: dict) -> list[dict[str, Any]]: """Строит список кандидатов на скачивание видео TikTok в порядке убывания качества: HD без водяных знаков → стандартное без водяных знаков → (только как самый последний резерв, если вообще ничего другого нет) видео с водяным знаком. НАЙДЕНО ПРИ РЕВИЗИИ: раньше запрос к TikWM не передавал параметр hd=1, и код брал только `media_data.get("play") or media_data.get("wmplay")` — то есть ВСЕГДА уходило видео в обычном (не HD) качестве без водяных знаков, даже когда у TikWM реально есть более качественная версия (`hdplay`). См. добавленный `&hd=1` в tikwm_endpoints в handle_tiktok — без него поле `hdplay` в ответе вообще не гарантированно присутствует. TikWM вместе с каждой ссылкой отдаёт реальный размер файла в байтах (`hd_size`/`size`/`wm_size`) — используем его, чтобы сразу пропустить вариант, который заведомо не пролезет в лимит Telegram Bot API на загрузку (см. TELEGRAM_BOT_API_UPLOAD_LIMIT_BYTES), а не тратить время и трафик на скачивание файла, который всё равно не отправится. Если размер не пришёл в ответе (0 или отсутствует — TikWM не всегда его отдаёт) — не отбрасываем вариант заранее, просто пробуем; на случай реального превышения лимита handle_tiktok сам ловит TelegramEntityTooLarge и переходит к следующему кандидату по качеству.""" candidates: list[dict[str, Any]] = [] for url_key, size_key, label in ( ("hdplay", "hd_size", "HD"), ("play", "size", "стандартное"), ("wmplay", "wm_size", "с водяным знаком — резерв"), ): raw_url = media_data.get(url_key) if not raw_url: continue if not raw_url.startswith("http"): raw_url = "https://www.tikwm.com" + raw_url try: size_bytes = int(media_data.get(size_key) or 0) except (TypeError, ValueError): size_bytes = 0 candidates.append({"key": url_key, "url": raw_url, "size": size_bytes, "label": label}) return candidates class TikTokUserFacingError(RuntimeError): """Ошибка TikTok-загрузки с текстом, уже написанным для пользователя (см. raise ниже по функции handle_tiktok). ВАЖНО для будущих правок: любой raise этого класса должен содержать ТОЛЬКО чистый русский текст без внутренних деталей/сырых исключений — except-блок в handle_tiktok показывает str(exc) пользователю as-is, без дополнительной проверки содержимого. Обычный RuntimeError (не этот подкласс) считается "сырым" и пользователю не показывается — см. except Exception ниже.""" # ─────────────────── ссылка на страницу звука (не видео) ─────────────────── # НАЙДЕНО ПО РЕАЛЬНЫМ ЛОГАМ (см. /logs владельца): если зайти в приложении TikTok # не на видео, а на сам ЗВУК (карточка "название звука" под видео → тап → "Поделиться"), # скопированная ссылка выглядит как https://www.tiktok.com/music/original-sound-7666630127215823637 # — числовой ID звука в конце после последнего дефиса. Основной эндпоинт TikWM # (`/api/?url=`, единственный, которым пользуется остальной код этого файла) на # такие ссылки отвечает "Url parsing is failed! Please check url." — он умеет # парсить только ссылки на видео/фото-посты, не на страницы звука. # # ИСТОРИЯ ДВУХ ПРОВАЛИВШИХСЯ ПОПЫТОК (обе подтверждены реальным тестированием, # см. логи владельца, — не гипотетические, а фактически проверенные и опровергнутые): # 1) Предположение, что TikWM принимает голый числовой ID видео вместо полной # ссылки, и что у "оригинальных звуков" ID звука совпадает с ID видео-источника. # Опровергнуто: TikWM отвечает "Url parsing is failed!" на голый числовой ID # ВСЕГДА, и для именованных песен, и для настоящих оригинальных звуков. # 2) Прямой запрос страницы tiktok.com/music/... и парсинг встроенного в неё JSON # (