diff --git "a/bot.py" "b/bot.py" deleted file mode 100644--- "a/bot.py" +++ /dev/null @@ -1,2853 +0,0 @@ -# >>> START OF BLOCK 0: GLOBAL NETWORK PATCH AND SYSTEM IMPORTS <<< - -import bot_modules.network_patch # noqa: F401 -# noqa: F401 — применяет SSL, proxy, trust_env патчи - -import os -import ssl -import certifi -import aiohttp -import asyncio -import logging -import sys -import traceback -import traceback -import re -import random -import html -import uuid -import ipaddress -import asyncpg -import httpx -import time -import json -from datetime import datetime, date -from io import BytesIO -from typing import Callable, Dict, Any, Awaitable - -# <<< END OF BLOCK 0 >>> -# >>> START OF BLOCK 1: LIBRARY IMPORTS, CONFIG AND CLIENT INITIALIZATION <<< - -from aiohttp import web -from dotenv import load_dotenv -from openai import AsyncOpenAI -from docx import Document -import logging -logger = logging.getLogger(__name__) - -from aiogram import Bot, Dispatcher, F, types -from aiogram.filters import CommandStart, Command -from aiogram.fsm.context import FSMContext -from aiogram.fsm.state import State, StatesGroup -from aiogram.types import ( - ReplyKeyboardRemove, - BufferedInputFile, - InlineKeyboardMarkup, - InlineKeyboardButton, - CallbackQuery, - Message -) -from aiogram import BaseMiddleware -from aiogram.client.session.aiohttp import AiohttpSession -from aiogram.webhook.aiohttp_server import SimpleRequestHandler, setup_application - -logging.basicConfig(level=logging.INFO, format="%(asctime)s - [%(levelname)s] - %(message)s") - -load_dotenv() -BOT_TOKEN = os.getenv("BOT_TOKEN", "").strip().replace('"', '').replace("'", "") -OPENAI_API_KEY = os.getenv("OPENAI_API_KEY", "").strip().replace('"', '').replace("'", "") -ADMIN_USERNAME = "Pteneev" -ADMIN_ID = int(os.getenv("ADMIN_ID", 0)) - -WEBHOOK_PATH = "/webhook" - -# Нативная интеграция Cloudflare Worker для обхода блокировки РКН -from aiogram.client.telegram import TelegramAPIServer - -proxy_server = os.getenv("TELEGRAM_API_SERVER") -if proxy_server: - proxy_server = proxy_server.strip().rstrip('/') - # Превращаем твой Worker в "официальный сервер Telegram" для aiogram - custom_api_server = TelegramAPIServer.from_base(proxy_server) - session = AiohttpSession(api=custom_api_server) - logging.info(f"✅ [PROXY] Telegram API перенаправлен через Cloudflare: {proxy_server}") -else: - session = AiohttpSession() - logging.warning("⚠️ [PROXY] TELEGRAM_API_SERVER не найден! Прямое подключение (может быть заблокировано).") - -bot = Bot(token=BOT_TOKEN, session=session) -dp = Dispatcher() - - -# OpenAI клиент с строгой SSL-валидацией через certifi -client = AsyncOpenAI( - api_key=OPENAI_API_KEY, - base_url="https://models.inference.ai.azure.com", - http_client=httpx.AsyncClient(timeout=60.0, verify=certifi.where(), trust_env=True) -) - -SUPPORT_API_KEY = os.getenv("SUPPORT_API_KEY", OPENAI_API_KEY).strip().replace('"', '').replace("'", "") - -client_support = AsyncOpenAI( - api_key=SUPPORT_API_KEY, - base_url="https://models.inference.ai.azure.com", - http_client=httpx.AsyncClient(timeout=45.0, verify=certifi.where(), trust_env=True) -) - # ═══════════════════════════════════════════════════════════════════ -# ИМПОРТЫ ИЗ МОДУЛЕЙ -# ═══════════════════════════════════════════════════════════════════ -from bot_modules.utils import ( - UserLockManager, - GlobalSpamProtector, - openai_semaphore, - _sanitize_memory_key, - _sanitize_memory_value, - safe_background_task, -) - -# ═══════════════════════════════════════════════════════════════════ -# ЭКЗЕМПЛЯРЫ БЛОКИРОВОК (создаются здесь, используются в хендлерах) -# ═══════════════════════════════════════════════════════════════════ -user_lock_manager = UserLockManager() -knowledge_lock_manager = UserLockManager() -summary_lock_manager = UserLockManager() -free_tier_protector = GlobalSpamProtector(daily_limit=1000) - -# ═══════════════════════════════════════════════════════════════════ -# B4: RAG LITE SYSTEM — Константы и промпт -# ═════════════════════════════════════════════════════════════���═════ - -# Максимальное количество знаний на пользователя -KNOWLEDGE_MAX_PER_USER = 100 - -# Порог важности для удаления при превышении лимита -KNOWLEDGE_CLEANUP_IMPORTANCE_THRESHOLD = 3 - -# Максимальное количество извлечённых знаний за один вызов -KNOWLEDGE_MAX_EXTRACTED_PER_CALL = 10 - -# Максимальное количество знаний для retrieval в контекст -KNOWLEDGE_RETRIEVAL_LIMIT = 10 - -# Системный промпт для извлечения знаний из сообщений пользователя -KNOWLEDGE_EXTRACTION_PROMPT = """Ты — аналитик знаний пользователя. Твоя задача: извлечь из сообщения пользователя полезные долгосрочные знания и вернуть их в строгом JSON-формате. - -ИЗВЛЕКАЙ только следующие типы знаний: -- Проекты (что пользователь делает или планирует) -- Технологии (с какими инструментами, фреймворками, языками работает) -- Бизнес-цели (чего хочет достичь) -- Архитектурные решения (как строит системы) -- Предпочтения пользователя (что нравится/не нравится) -- Долгосрочные задачи (планы на месяцы/годы) - -НЕ извлекай: -- Приветствия, прощания, случайный чат -- Временные или контекстно-зависимые факты -- Одноразовые вопросы без долгосрочной ценности -- Шум, воду, эмоциональные реакции без смысловой нагрузки -- Ответы модели или предположения модели -- Системные промпты или технические инструкции - -Формат ответа (строго JSON, без markdown-блоков): -{ - "knowledge": [ - {"content": "Пользователь делает SaaS на aiogram", "importance": 8}, - {"content": "Предпочитает PostgreSQL вместо MongoDB", "importance": 6} - ] -} - -Правила: -1. content — конкретное, самодостаточное знание (до 500 символов) -2. importance — целое число 1-10, где 10 = критически важное для будущих ответов -3. Если знаний нет — верни {"knowledge": []} -4. Извлекай только из сообщений ПОЛЬЗОВАТЕЛЯ, никогда из ответов AI""" - -# ═══════════════════════════════════════════════════════════════════ -# B3: CONVERSATION SUMMARY SYSTEM — Константы и промпт -# ═══════════════════════════════════════════════════════════════════ - -# Порог срабатывания summary: каждые N сообщений пользователя -SUMMARY_MESSAGE_THRESHOLD = 30 - -# Максимальная длина summary в символах (жёсткий лимит) -SUMMARY_MAX_LENGTH = 5000 - -# Целевая длина summary, генерируемого GPT -SUMMARY_TARGET_LENGTH = 1000 - -# Системный промпт для генерации conversation summary -SUMMARY_SYSTEM_PROMPT = """Ты — аналитик диалогов. Твоя задача: создать краткий, структурированный summary диалога между пользователем и AI-ассистентом. - -ИЗВЛЕКАЙ только следующее: -- Цели пользователя (что он хочет достичь) -- Проекты пользователя (над чем работает) -- Текущие задачи (что решает прямо сейчас) -- Принятые решения (что уже согласовано/решено) -- Незавершённые вопросы (что осталось открытым) -- Важный контекст диалога (ключевые факты, договорённости) - -НЕ включай: -- Приветствия и прощания -- Шум, воду, повторения -- Технические детали форматирования -- Эмоциональные реакции без смысловой нагрузки - -Формат вывода: -Краткий структурированный текст до 1000 символов. Без markdown, без списков с эмодзи, без таблиц. - -Если предоставлен предыдущий summary + новые сообщения — объедини их в единый, обновлённый summary, сохраняя всю важную информацию из старого summary и добавляя новое из сообщений.""" - -# <<< END OF BLOCK 1 >>> -# >>> START OF BLOCK 2: ANTI-SPAM MIDDLEWARE AND USER LOCKS <<< - -class ThrottlingMiddleware(BaseMiddleware): - """Rate-limiting per user: защита от флуда и случайного дублирования кнопок. - - Бизнес-контекст: снижает нагрузку на БД и OpenAI API, улучшает UX - (пользователь не запускает 5 параллельных генераций случайно). - """ - def __init__(self, rate_limit: float = 3.0): - self.rate_limit = rate_limit - self.users_timers: Dict[int, float] = {} - - async def __call__(self, handler: Callable, event: types.Message, data: Dict[str, Any]) -> Any: - if not event.from_user: - return await handler(event, data) - - user_id = event.from_user.id - now = datetime.now().timestamp() - - # Очистка старых таймеров для предотвращения утечки памяти (OOM) - if len(self.users_timers) > 5000: - self.users_timers = {k: v for k, v in self.users_timers.items() if now - v < 300} - - if user_id in self.users_timers: - time_passed = now - self.users_timers[user_id] - if time_passed < self.rate_limit: - logging.warning(f"🛡 Anti-Flood: Молча заблокирован спам от пользователя {user_id}") - return - - self.users_timers[user_id] = now - return await handler(event, data) - -class LengthLimitMiddleware(BaseMiddleware): - """Ограничение длины входящих сообщений — защита от prompt injection и OOM.""" - async def __call__( - self, - handler: Callable[[Message, Dict[str, Any]], Awaitable[Any]], - event: Message, - data: Dict[str, Any] - ) -> Any: - if event.text and len(event.text) > 1000: - await event.answer("⚠️ Ваше сообщение слишком длинное (максимум 1000 символов). Пожалуйста, сделайте его короче.") - return - return await handler(event, data) - -dp.message.middleware(ThrottlingMiddleware(rate_limit=3.0)) - -# <<< END OF BLOCK 2 >>> -# >>> START OF BLOCK 3: ASYNC POSTGRESQL CLOUD DATABASE MANAGEMENT <<< - -db_pool = None -_cleanup_task: asyncio.Task | None = None - -async def init_db(): - global db_pool - db_url = os.getenv("DATABASE_URL") - if not db_url: - logging.critical("🚨 КРИТИЧЕСКАЯ ОШИБКА: Переменная DATABASE_URL не найдена в Secrets!") - return - - try: - # SSL-контекст для PostgreSQL — отдельный от внешнего custom_ssl - # Используем certifi для публичных БД, но НЕ подменяем глобальный контекст - pg_ssl = ssl.create_default_context(cafile=certifi.where()) - pg_ssl.check_hostname = True - pg_ssl.verify_mode = ssl.CERT_REQUIRED - - # connection_init_hook: валидация коннекта перед выдачей из пула - # Бизнес-контекст: предотвращает InterfaceError при рестарте Serverless БД - async def init_conn(conn): - await conn.execute("SELECT 1") - - db_pool = await asyncpg.create_pool( - db_url, - ssl=pg_ssl, - min_size=2, - max_size=15, - command_timeout=10.0, - max_inactive_connection_lifetime=300.0, - init=init_conn - ) - - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute(''' - CREATE TABLE IF NOT EXISTS users ( - user_id BIGINT PRIMARY KEY, - username TEXT, - generation_count INTEGER DEFAULT 0, - created_at TEXT, - tier TEXT DEFAULT 'Free', - max_limit INTEGER DEFAULT 1 - ) - ''') - await conn.execute(''' - CREATE TABLE IF NOT EXISTS tickets ( - ticket_id BIGINT PRIMARY KEY, - user_id BIGINT, - created_at TEXT - ) - ''') - # СТРУКТУРА: Новая таблица логирования ИИ-запросов - await conn.execute(''' - CREATE TABLE IF NOT EXISTS ai_logs ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL, - model TEXT NOT NULL, - role TEXT NOT NULL, - prompt_tokens INTEGER DEFAULT 0, - completion_tokens INTEGER DEFAULT 0, - latency_ms INTEGER DEFAULT 0, - status_code INTEGER, - success BOOLEAN DEFAULT TRUE, - error TEXT, - created_at TIMESTAMP DEFAULT NOW() - ) - ''') - - # ═══════════════════════════════════════════════════════════════ - # ИНДЕКСЫ ai_logs: Оптимизация для роста пользовательской базы - # Безопасная идемпотентная миграция — выполняется один раз - # ═══════════════════════════════════════════════════════════════ - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_ai_logs_user - ON ai_logs(user_id) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_ai_logs_created - ON ai_logs(created_at) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_ai_logs_success - ON ai_logs(success) - ''') - await conn.execute(''' - CREATE TABLE IF NOT EXISTS user_memory ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL, - memory_key TEXT NOT NULL, - memory_value TEXT NOT NULL, - importance INTEGER DEFAULT 5, - created_at TIMESTAMP DEFAULT NOW(), - updated_at TIMESTAMP DEFAULT NOW(), - UNIQUE(user_id, memory_key), - CHECK (LENGTH(memory_key) <= 50), - CHECK (LENGTH(memory_value) <= 200), - CHECK (importance >= 1 AND importance <= 10) - ) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_user_memory_user_id - ON user_memory(user_id) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_user_memory_updated - ON user_memory(updated_at) - ''') - # B2.1: Таблица для хранения истории диалогов (conversation_history) - await conn.execute(''' - CREATE TABLE IF NOT EXISTS conversation_history ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL, - role VARCHAR(20) NOT NULL, - content TEXT NOT NULL, - created_at TIMESTAMP DEFAULT NOW() - ) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_conversation_history_user - ON conversation_history(user_id, created_at DESC) - ''') - # Legacy: таблица message_history оставлена для обратной совместимости - await conn.execute(''' - CREATE TABLE IF NOT EXISTS message_history ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL, - role TEXT NOT NULL, - content TEXT NOT NULL, - message_type TEXT NOT NULL DEFAULT 'user', - created_at TIMESTAMP DEFAULT NOW() - ) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_message_history_user_created - ON message_history(user_id, created_at) - ''') - # ═══════════════════════════════════════════════════════════════ - # B3: Таблица conversation_summaries (Conversation Summary System) - # ═══════════════════════════════════════════════════════════════ - await conn.execute(''' - CREATE TABLE IF NOT EXISTS conversation_summaries ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL UNIQUE, - summary TEXT NOT NULL DEFAULT '', - message_count INTEGER NOT NULL DEFAULT 0, - updated_at TIMESTAMP DEFAULT NOW() - ) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_conversation_summaries_user - ON conversation_summaries(user_id) - ''') - # ═══════════════════════════════════════════════════════════════ - # B4: Таблица knowledge_base (RAG Lite System) - # ═══════════════════════════════════════════════════════════════ - await conn.execute(''' - CREATE TABLE IF NOT EXISTS knowledge_base ( - id BIGSERIAL PRIMARY KEY, - user_id BIGINT NOT NULL, - content TEXT NOT NULL, - source TEXT DEFAULT 'user_message', - importance INTEGER NOT NULL DEFAULT 5, - created_at TIMESTAMP DEFAULT NOW(), - CHECK (LENGTH(content) <= 500), - CHECK (importance >= 1 AND importance <= 10) - ) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_knowledge_base_user_id - ON knowledge_base(user_id) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_knowledge_base_created - ON knowledge_base(created_at) - ''') - await conn.execute(''' - CREATE INDEX IF NOT EXISTS idx_knowledge_base_importance - ON knowledge_base(user_id, importance DESC, created_at DESC) - ''') - - # Миграция: строго INTEGER для max_limit, TEXT для tier - - try: - await conn.execute("ALTER TABLE users ADD COLUMN IF NOT EXISTS tier TEXT DEFAULT 'Free'") - await conn.execute("ALTER TABLE users ADD COLUMN IF NOT EXISTS max_limit INTEGER DEFAULT 1") - except Exception as e: - logging.warning(f"Миграция колонок пропущена или выполнена ранее: {e}") - - logging.info("🗄 Облачная база данных (asyncpg) успешно инициализирована с SSL!") - except Exception as e: - logging.error(f"❌ Ошибка инициализации базы данных: {e}") - -async def cleanup_old_ai_logs(): - """Автоматическая очистка логов ИИ-запросов старше 90 дней. - - Бизнес-контекст: предотвращает бесконтрольный рост таблицы ai_logs - в облачной PostgreSQL. Логи старше 6 месяцев неактуальны для - операционного мониторинга latency и cost per user. - """ - if not db_pool: - logging.warning("⚠️ cleanup_old_ai_logs: пул соединений не инициализирован, пропуск.") - return - - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM ai_logs WHERE created_at < NOW() - INTERVAL '90 days'" - ) - # asyncpg возвращает строку вида "DELETE 0" — извлекаем количество - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🧹 cleanup_old_ai_logs: удалено {deleted_count} записей старше 90 дней.") - else: - logging.info("🧹 cleanup_old_ai_logs: старых записей не найдено.") - except Exception as e: - logging.error(f"❌ Ошибка в cleanup_old_ai_logs: {e}") - -async def cleanup_old_message_history(days: int = 7): - """Автоматическая очистка истории сообщений старше N дней. - - Бизнес-контекст: message_history — оперативный контекст, не архив. - Храним только активное окно (7 дней) для снижения объёма БД в HF Spaces. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM message_history WHERE created_at < NOW() - $1 * INTERVAL '1 day'", - days - ) - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🧹 cleanup_old_message_history: удалено {deleted_count} записей старше {days} дней.") - except Exception as e: - logging.error(f"❌ Ошибка в cleanup_old_message_history: {e}") - -async def cleanup_old_conversation_history(days: int = 7): - """Автоматическая очистка истории диалогов conversation_history старше N дней. - - Бизнес-контекст: conversation_history — оперативный контекст, не архив. - Храним только активное окно (7 дней) для снижения объёма БД в HF Spaces. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM conversation_history WHERE created_at < NOW() - $1 * INTERVAL '1 day'", - days - ) - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🧹 cleanup_old_conversation_history: удалено {deleted_count} записей старше {days} дней.") - except Exception as e: - logging.error(f"❌ Ошибка в cleanup_old_conversation_history: {e}") - -async def cleanup_old_conversation_summaries(days: int = 90): - """Автоматическая очистка summary неактивных пользователей старше N дней. - - Бизнес-контекст: summary неактивных пользователей неактуальны. - Храним 90 дней для баланса между памятью и объёмом БД. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM conversation_summaries WHERE updated_at < NOW() - $1 * INTERVAL '1 day'", - days - ) - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info( - f"🧹 cleanup_old_conversation_summaries: " - f"удалено {deleted_count} записей старше {days} дней." - ) - except Exception as e: - logging.error(f"❌ Ошибка в cleanup_old_conversation_summaries: {e}") - -async def _cleanup_loop(): - """Фоновый цикл очистки ai_logs и message_history каждые 12 часов. - - Бизнес-контекст: предотвращает бесконтрольный рост таблиц между - перезапусками контейнера HF Spaces. Корректно обрабатывает ошибки - внутри цикла — одна неудачная итерация не убивает задачу навсегда. - """ - while True: - # A: Очистка ai_logs (самая частая таблица) - try: - logging.info("🧹 [Cleanup Scheduler] Запуск периодической очистки ai_logs...") - await cleanup_old_ai_logs() - logging.info("✅ [Cleanup Scheduler] Очистка ai_logs завершена.") - except asyncio.CancelledError: - logging.info("🛑 [Cleanup Scheduler] Задача отменена (shutdown).") - raise - except Exception as e: - logging.error(f"❌ [Cleanup Scheduler] Ошибка в цикле очистки ai_logs: {e}") - - # B: Очистка conversation_history (оперативный контекст, не архив) - try: - logging.info("🧹 [Cleanup Scheduler] Запуск очистки conversation_history...") - await cleanup_old_conversation_history(days=7) - logging.info("✅ [Cleanup Scheduler] Очистка conversation_history завершена.") - except Exception as e: - logging.error(f"❌ [Cleanup Scheduler] Ошибка очистки conversation_history: {e}") - - # C: Очистка message_history (legacy) - try: - logging.info("🧹 [Cleanup Scheduler] Запуск очистки message_history...") - await cleanup_old_message_history(days=7) - logging.info("✅ [Cleanup Scheduler] Очистка message_history завершена.") - except Exception as e: - logging.error(f"❌ [Cleanup Scheduler] Ошибка очистки message_history: {e}") - - # D: Очистка conversation_summaries (неактивные пользователи) - try: - logging.info("🧹 [Cleanup Scheduler] Запуск очистки conversation_summaries...") - await cleanup_old_conversation_summaries(days=90) - logging.info("✅ [Cleanup Scheduler] Очистка conversation_summaries завершена.") - except Exception as e: - logging.error(f"❌ [Cleanup Scheduler] Ошибка очистки conversation_summaries: {e}") - - # E: Очистка knowledge_base (устаревшие знания) - try: - logging.info("🧹 [Cleanup Scheduler] ��апуск очистки knowledge_base...") - await cleanup_old_knowledge(days=180) - logging.info("✅ [Cleanup Scheduler] Очистка knowledge_base завершена.") - except Exception as e: - logging.error(f"❌ [Cleanup Scheduler] Ошибка очистки knowledge_base: {e}") - - # Все таблицы очищены — спим до следующей итерации - await asyncio.sleep(43200) # 12 часов - -async def start_cleanup_scheduler(): - """Запускает фоновую задачу очистки ai_logs с защитой от дубликатов. - - Гарантирует, что при повторном вызове on_startup() (например, тесты - или горячий перезапуск) не создаётся вторая задача. - """ - global _cleanup_task - if _cleanup_task is not None and not _cleanup_task.done(): - logging.warning("⚠️ [Cleanup Scheduler] Задача уже запущена, пропуск создания дубликата.") - return - _cleanup_task = asyncio.create_task(_cleanup_loop()) - logging.info("🚀 [Cleanup Scheduler] Фоновая задача очистки ai_logs запущена (интервал: 12ч).") - -# ═══════════════════════════════════════════════════════════════════ -# AI LOGS (Observability) -# ═══════════════════════════════════════════════════════════════════ -async def save_ai_log( - user_id: int, - model: str, - role: str, - prompt_tokens: int, - completion_tokens: int, - latency_ms: int, - status_code: int, - success: bool, - error: str -): - """Структурированное логирование LLM-запросов для Grafana/дэшбордов. - - Бизнес-контекст: позволяет отслеживать cost per user, anomaly detection - (всплески latency), и отправлять тикеты в OpenAI support с request_id. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute(''' - INSERT INTO ai_logs ( - user_id, model, role, prompt_tokens, completion_tokens, - latency_ms, status_code, success, error - ) VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9) - ''', user_id, model, role, prompt_tokens, completion_tokens, latency_ms, status_code, success, error) - except Exception as e: - logging.error(f"❌ Ошибка в save_ai_log: {e}") - -async def save_user_memory(user_id: int, memory_key: str, memory_value: str, importance: int = 5): - """Сохранение или обновление факта о пользователе с upsert-логикой и валидацией.""" - if not db_pool: - return - key = _sanitize_memory_key(memory_key) - value = _sanitize_memory_value(memory_value) - if not key or not value: - return - try: - importance = min(10, max(1, int(importance))) - except (TypeError, ValueError): - importance = 5 - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute(''' - INSERT INTO user_memory (user_id, memory_key, memory_value, importance, created_at, updated_at) - VALUES ($1, $2, $3, $4, NOW(), NOW()) - ON CONFLICT (user_id, memory_key) - DO UPDATE SET - memory_value = EXCLUDED.memory_value, - importance = EXCLUDED.importance, - updated_at = NOW() - ''', user_id, key, value, importance) - except Exception as e: - logging.error(f"❌ Ошибка в save_user_memory: {e}") - -async def get_user_memory(user_id: int) -> list: - """Получение фактов о пользователе с защитой от OOM (LIMIT 100).""" - if not db_pool: - return [] - try: - async with db_pool.acquire(timeout=5.0) as conn: - rows = await conn.fetch( - "SELECT memory_key, memory_value, importance FROM user_memory WHERE user_id = $1 ORDER BY updated_at DESC LIMIT 100", - user_id - ) - return [(row['memory_key'], row['memory_value'], row['importance']) for row in rows] - except Exception as e: - logging.error(f"❌ Ошибка в get_user_memory: {e}") - return [] - -async def get_conversation_summary(user_id: int) -> dict | None: - """Получение conversation summary для пользователя. - - Возвращает dict с ключами summary, message_count, updated_at - или None если записи нет. - """ - if not db_pool: - return None - try: - async with db_pool.acquire(timeout=5.0) as conn: - row = await conn.fetchrow( - "SELECT summary, message_count, updated_at FROM conversation_summaries WHERE user_id = $1", - user_id - ) - if row: - return { - 'summary': row['summary'], - 'message_count': row['message_count'], - 'updated_at': row['updated_at'] - } - return None - except Exception as e: - logging.error(f"❌ [B3] Ошибка в get_conversation_summary: {e}") - return None - -async def save_conversation_summary(user_id: int, summary: str, message_count: int): - """Сохранение или создание conversation summary с upsert-логикой.""" - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute(''' - INSERT INTO conversation_summaries (user_id, summary, message_count, updated_at) - VALUES ($1, $2, $3, NOW()) - ON CONFLICT (user_id) - DO UPDATE SET - summary = EXCLUDED.summary, - message_count = EXCLUDED.message_count, - updated_at = NOW() - ''', user_id, summary, message_count) - except Exception as e: - logging.error(f"❌ [B3] Ошибка в save_conversation_summary: {e}") - -async def update_conversation_summary(user_id: int, summary: str, message_count: int | None = None): - """Явное обновление существующего conversation summary. - - Используется для атомарных обновлений summary без сброса счётчика - или со сбросом счётчика при полном обновлении. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - if message_count is not None: - await conn.execute(''' - UPDATE conversation_summaries - SET summary = $1, message_count = $2, updated_at = NOW() - WHERE user_id = $3 - ''', summary, message_count, user_id) - else: - await conn.execute(''' - UPDATE conversation_summaries - SET summary = $1, updated_at = NOW() - WHERE user_id = $2 - ''', summary, user_id) - except Exception as e: - logging.error(f"❌ [B3] Ошибка в update_conversation_summary: {e}") - -async def generate_conversation_summary(messages: list, existing_summary: str | None = None) -> str: - """Генерация summary диалога через GPT с инкрементальным сжатием. - - Аргументы: - messages: список dict'ов [{'role': 'user'|'assistant', 'content': str}, ...] - existing_summary: предыдущий summary для инкрементального обновления. - - Возвращает: - Строку summary, обрезанную до SUMMARY_MAX_LENGTH при необходимости. - """ - if not messages: - return "" - - # Формируем текст сообщений - messages_text = [] - for msg in messages: - role_label = "Пользователь" if msg.get('role') == 'user' else "Ассистент" - content = msg.get('content', '') - if content: - messages_text.append(f"{role_label}: {content}") - - conversation_text = "\n\n".join(messages_text) - - # Формируем user prompt в зависимости от наличия старого summary - if existing_summary: - user_prompt = f"""Существующий summary диалога: -{existing_summary} - -Новые сообщения для интеграции: -{conversation_text} - -Обнови summary, объединив старую информацию с новыми сообщениями. -Сохрани всё важное из старого summary. Удали устаревшее. Результат — до 1000 символов.""" - else: - user_prompt = f"""Сообщения диалога: -{conversation_text} - -Создай краткий summary диалога. Результат — до 1000 символов.""" - - model_name = "gpt-4o-mini" - role = "conversation_summarizer" - - try: - async with openai_semaphore: - response = await client.chat.completions.create( - model=model_name, - messages=[ - {"role": "system", "content": SUMMARY_SYSTEM_PROMPT}, - {"role": "user", "content": user_prompt} - ], - temperature=0.3, - max_tokens=800 - ) - - summary = response.choices[0].message.content or "" - summary = summary.strip() - - # Логирование токенов - prompt_tokens = response.usage.prompt_tokens if response.usage else 0 - completion_tokens = response.usage.completion_tokens if response.usage else 0 - logging.info( - f"🤖 [B3 Summary] Generated | model={model_name} | " - f"tokens={prompt_tokens}p+{completion_tokens}c | " - f"length={len(summary)} | existing={bool(existing_summary)}" - ) - - return summary - - except Exception as e: - logging.error(f"❌ [B3 Summary] Ошибка генерации summary: {e}") - raise - -async def _maybe_update_conversation_summary(user_id: int): - """Проверяет необходимость обновления summary и выполняет инкрементальное сжатие. - - Логика: - - Инкрементирует счётчик сообщений пользователя. - - Если счётчик >= SUMMARY_MESSAGE_THRESHOLD: генерирует новый summary - из старого summary + последних сообщений, сбрасывает счётчик. - - Если счётчик < порога: просто обновляет счётчик. - - Защита от гонок: per-user lock через summary_lock_manager. - """ - if not db_pool: - logging.info(f"📊 [B3 Summary] user_id={user_id} | Summary skipped: db_pool unavailable") - return - - lock = await summary_lock_manager.get(user_id) - try: - async with lock: - # 1. Получаем текущий summary - summary_row = await get_conversation_summary(user_id) - - if summary_row: - current_summary = summary_row['summary'] - message_count = summary_row['message_count'] - else: - current_summary = "" - message_count = 0 - - # 2. Инкрементируем счётчик (только сообщения пользователя) - message_count += 1 - - # 3. Проверяем порог - if message_count < SUMMARY_MESSAGE_THRESHOLD: - # Порог не достигнут — просто обновляем счётчик - await save_conversation_summary(user_id, current_summary, message_count) - logging.info( - f"📊 [B3 Summary] user_id={user_id} | " - f"message_count={message_count}/{SUMMARY_MESSAGE_THRESHOLD} | " - f"threshold not reached" - ) - return - - # 4. Порог достигнут — генерируем новый summary - logging.info( - f"📝 [B3 Summary] user_id={user_id} | " - f"Threshold reached ({SUMMARY_MESSAGE_THRESHOLD}) | " - f"Generating incremental summary..." - ) - - # Получаем последние N сообщений для обработки - recent_messages = await get_recent_messages(user_id, limit=SUMMARY_MESSAGE_THRESHOLD) - - if not recent_messages: - await save_conversation_summary(user_id, current_summary, 0) - logging.info( - f"📊 [B3 Summary] user_id={user_id} | " - f"No messages found, resetting counter" - ) - return - - # 5. Генерируем summary (инкрементально: старый + новые) - try: - new_summary = await generate_conversation_summary( - recent_messages, - existing_summary=current_summary if current_summary else None - ) - - # 6. Безопасная обрезка до жёсткого лимита - original_len = len(new_summary) - if original_len > SUMMARY_MAX_LENGTH: - # Обрезаем по последнему пробелу/переносу для безопасности - truncated = new_summary[:SUMMARY_MAX_LENGTH] - last_break = max(truncated.rfind(' '), truncated.rfind('\n')) - if last_break > SUMMARY_MAX_LENGTH * 0.8: - new_summary = truncated[:last_break] - else: - new_summary = truncated - logging.warning( - f"⚠️ [B3 Summary] user_id={user_id} | " - f"Summary truncated from {original_len} to {len(new_summary)} chars" - ) - - # 7. Сохраняем новый summary и сбрасываем счётчик - await save_conversation_summary(user_id, new_summary, 0) - - logging.info( - f"✅ [B3 Summary] user_id={user_id} | " - f"Summary updated | length={len(new_summary)} | " - f"messages_processed={len(recent_messages)}" - ) - - except Exception as e: - logging.error( - f"❌ [B3 Summary] user_id={user_id} | " - f"Summary generation failed: {e}", - exc_info=True - ) - # Не сбрасываем счётчик — повторим при следующем сообщении - await save_conversation_summary(user_id, current_summary, message_count) - - finally: - await summary_lock_manager.release(user_id) - -# ═══════════════════════════════════════════════════════════════════ -# B2.1 CONTEXT BUILDER: CONVERSATION HISTORY AND CONTEXT ASSEMBLY -# ═══════════════════════════════════════════════════════════════════ - -async def save_conversation_history(user_id: int, role: str, content: str): - """Сохранение сообщения в историю диалогов + триггер B3 summary. - - Бизнес-контекст: conversation_history хранит полный текст сообщений - пользователя и AI для формирования контекста при последующих запросах. - Автоочистка через cleanup_old_conversation_history предотвращает рост таблицы. - - B3: На каждом сообщении пользователя (role='user') запускается фоновая - проверка порога summary. Инкрементальное сжатие выполняется каждые - SUMMARY_MESSAGE_THRESHOLD сообщений. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute(''' - INSERT INTO conversation_history (user_id, role, content) - VALUES ($1, $2, $3) - ''', user_id, role, content) - except Exception as e: - logging.error(f"❌ Ошибка в save_conversation_history: {e}") - return - - # B3: Триггер обновления summary только на сообщениях пользователя - if role == 'user': - safe_background_task( - _maybe_update_conversation_summary(user_id), - name=f"summary_update_{user_id}" - ) - -async def get_recent_messages(user_id: int, limit: int = 10) -> list: - """Загрузка последних сообщений пользователя и AI. - - Возвращает список словарей в хронологическом порядке: - [{'role': 'user'|'assistant', 'content': str}, ...] - Если сообщений меньше limit — используются доступные. - """ - if not db_pool: - return [] - try: - async with db_pool.acquire(timeout=5.0) as conn: - rows = await conn.fetch( - """ - SELECT role, content - FROM conversation_history - WHERE user_id = $1 - ORDER BY created_at DESC - LIMIT $2 - """, - user_id, limit - ) - # Реверс для хронологического порядка (старые → новые) - history = [] - for row in reversed(rows): - history.append({ - 'role': row['role'], - 'content': row['content'] - }) - return history - except Exception as e: - logging.error(f"❌ Ошибка в get_recent_messages: {e}") - return [] - -async def build_user_context(user_id: int, current_message: str) -> str: - """Сборка единого контекста из памяти пользователя, summary диалога, RAG знаний и истории. - - Структура контекста (B4): - ======================== - USER MEMORY - - {memory} - - ======================== - CONVERSATION SUMMARY - - {summary} - - ======================== - RAG KNOWLEDGE - - {knowledge} - - ======================== - RECENT HISTORY - - {history} - - ======================== - CURRENT MESSAGE - - {current_message} - - При ошибке загрузки любого блока — используется доступное + current_message. - """ - memory_count = 0 - history_count = 0 - summary_length = 0 - knowledge_count = 0 - context_parts = [] - - try: - # 1. Загружаем память пользователя (B1) - memory_facts = [] - try: - memory_facts = await get_user_memory(user_id) - memory_count = len(memory_facts) - except Exception as e: - logging.warning(f"⚠️ [Context Build] Не удалось загрузить память для user_id={user_id}: {e}") - - # 2. Загружаем conversation summary (B3) - summary_text = "" - try: - summary_row = await get_conversation_summary(user_id) - if summary_row and summary_row.get('summary'): - summary_text = summary_row['summary'] - summary_length = len(summary_text) - except Exception as e: - logging.warning(f"⚠️ [Context Build] Не удалось загрузить summary для user_id={user_id}: {e}") - - # 3. Загружаем RAG знания (B4) - knowledge_items = [] - try: - knowledge_items = await search_knowledge(user_id, current_message, limit=KNOWLEDGE_RETRIEVAL_LIMIT) - knowledge_count = len(knowledge_items) - except Exception as e: - logging.warning(f"⚠️ [Context Build] Не удалось выполнить поиск знаний для user_id={user_id}: {e}") - - # 4. Загружаем последние сообщения (B2.1) - history_messages = [] - try: - history_messages = await get_recent_messages(user_id, limit=10) - history_count = len(history_messages) - except Exception as e: - logging.warning(f"⚠️ [Context Build] Не удалось загрузить историю для user_id={user_id}: {e}") - - # 5. Формируем блок памяти (B1) - if memory_facts: - memory_lines = [] - for key, value, importance in memory_facts: - if importance >= 6: - memory_lines.append(f"- {key}: {value}") - if memory_lines: - context_parts.append("========================") - context_parts.append("USER MEMORY") - context_parts.append("") - context_parts.append("\n".join(memory_lines)) - - # 6. Формируем блок summary (B3) - if summary_text: - context_parts.append("========================") - context_parts.append("CONVERSATION SUMMARY") - context_parts.append("") - context_parts.append(summary_text) - - # 7. Формируем блок RAG знаний (B4) - if knowledge_items: - knowledge_lines = [] - for item in knowledge_items: - importance_marker = "★" if item.get('importance', 5) >= 8 else "" - knowledge_lines.append(f"- {importance_marker}{item['content']}") - if knowledge_lines: - context_parts.append("========================") - context_parts.append("RAG KNOWLEDGE") - context_parts.append("") - context_parts.append("\n".join(knowledge_lines)) - - # 8. Формируем блок истории (B2.1) - if history_messages: - history_lines = [] - for msg in history_messages: - role_label = "User" if msg['role'] == 'user' else "Assistant" - history_lines.append(f"{role_label}: {msg['content']}") - context_parts.append("========================") - context_parts.append("RECENT HISTORY") - context_parts.append("") - context_parts.append("\n\n".join(history_lines)) - - # 9. Добавляем текущее сообщение - context_parts.append("========================") - context_parts.append("CURRENT MESSAGE") - context_parts.append("") - context_parts.append(current_message) - - # 10. Собираем финальный контекст - context = "\n".join(context_parts) - context_length = len(context) - - # 11. Логирование - logging.info( - f"📋 [Context Build] user_id={user_id} | " - f"memory={memory_count} | summary={summary_length}ch | " - f"knowledge={knowledge_count} | history={history_count} | " - f"total_length={context_length}" - ) - - return context - - except Exception as e: - logging.error(f"❌ [Context Build] Критическая ошибка сборки контекста для user_id={user_id}: {e}") - # Graceful degradation: возвращаем только текущее сообщение - return current_message - -# Legacy: сохраняем обратную совместимость ��ля существующих вызовов -async def save_message_history(user_id: int, role: str, content: str, message_type: str = "user"): - """Legacy-обёртка: сохраняет в conversation_history с маппингом ролей. - - message_type='user' → role='user' - message_type='assistant' → role='assistant' - """ - mapped_role = "assistant" if message_type == "assistant" else "user" - await save_conversation_history(user_id, mapped_role, content) - - -async def get_message_history_full(user_id: int, limit: int = 10) -> list: - """Legacy-обёртка: делегирует get_recent_messages для обратной совместимости.""" - return await get_recent_messages(user_id, limit) - - -async def build_context_for_user( - user_id: int, - system_prompt: str, - current_user_message: str, - history_limit: int = 10, - max_context_tokens: int = 3000 -) -> list: - """Legacy-обёртка: собирает messages в формате OpenAI API из нового контекста. - - Использует build_user_context для формирования контекста и сохраняет - структуру возвращаемого значения для обратной совместимости. - """ - start_time = time.perf_counter() - - # 1. Загружаем память пользователя - memory_facts = [] - try: - memory_facts = await get_user_memory(user_id) - except Exception as e: - logging.warning(f"⚠️ [Context] Не удалось загрузить память для user_id={user_id}: {e}") - - # 2. Загружаем историю сообщений - history_messages = [] - try: - history_messages = await get_recent_messages(user_id, limit=history_limit) - except Exception as e: - logging.warning(f"⚠️ [Context] Не удалось загрузить историю для user_id={user_id}: {e}") - - # 3. Формируем блок памяти для system prompt - memory_block = "" - if memory_facts: - memory_lines = [] - for key, value, importance in memory_facts: - if importance >= 6: - memory_lines.append(f"- {key}: {value}") - if memory_lines: - memory_block = "\n\n[ИЗВЕСТНЫЕ ФАКТЫ О ПОЛЬЗОВАТЕЛЕ]\n" + "\n".join(memory_lines) + "\n[КОНЕЦ ФАКТОВ]" - - # 4. Собираем system message - system_content = system_prompt - if memory_block: - system_content = system_prompt + memory_block - - # 5. Оценка токенов и truncation при необходимости - system_tokens = _estimate_tokens(system_content) - current_tokens = _estimate_tokens(current_user_message) - history_tokens = sum(_estimate_tokens(m['content']) for m in history_messages) - total_estimate = system_tokens + current_tokens + history_tokens - - # Если превышаем лимит — урезаем историю с начала (оставляем более свежие в конце) - if total_estimate > max_context_tokens and history_messages: - while history_messages and total_estimate > max_context_tokens: - removed = history_messages.pop(0) - total_estimate -= _estimate_tokens(removed['content']) - if not history_messages: - break - - # 6. Формируем финальный messages список - messages = [{"role": "system", "content": system_content}] - messages.extend(history_messages) - messages.append({"role": "user", "content": current_user_message}) - - latency_ms = int((time.perf_counter() - start_time) * 1000) - - # 7. Логирование контекста - logging.info( - f"📋 [Context Build] user_id={user_id} | " - f"latency={latency_ms}ms | " - f"context_messages_count={len(history_messages)} | " - f"context_memory_count={len([f for f in memory_facts if f[2] >= 6])} | " - f"context_tokens_estimate={total_estimate} | " - f"total_messages={len(messages)}" - ) - - return messages - -async def extract_and_save_memory(user_id: int, user_text: str, ai_response: str): - """Извлечение долгосрочных фактов из диалога и сохранение в БД. - - Бизнес-контекст: после каждого успешного AI-ответа анализируем диалог - на предмет персональных фактов (цели, предпочтения, проекты) и сохраняем - для персонализации будущих взаимодействий. - """ - if not db_pool: - return - - memory_system_prompt = """Ты — аналитик персональных данных. Извлеки из диалога факты о пользователе, которые будут полезны в б��дущем. - -Правила: -1. Извлекай ТОЛЬКО долгосрочные факты: профессия, цели, проекты, предпочтения, технологии, ниша. -2. НЕ извлекай временные или контекстно-зависимые факты. -3. Каждый факт — это пара ключ=значение. -4. Ключ — короткий snake_case идентификатор (например: project, goal, niche, experience). -5. Значение — конкретная информация, максимум 200 символов. -6. Важность от 1 до 10, где 10 — критически важный факт для персонализации. - -Ответ строго в JSON-формате: -{ - "facts": [ - {"key": "project", "value": "telegram ai saas", "importance": 8}, - {"key": "goal", "value": "создать ai бизнес", "importance": 9} - ] -} - -Если фактов нет — верни {"facts": []}.""" - - combined_text = f"Вопрос пользователя:\n{user_text}\n\nОтвет AI:\n{ai_response}" - model_name = "gpt-4o-mini" - role = "memory_extractor" - - try: - async with openai_semaphore: - response = await client.chat.completions.create( - model=model_name, - messages=[ - {"role": "system", "content": memory_system_prompt}, - {"role": "user", "content": combined_text} - ], - temperature=0.1, - max_tokens=300, - response_format={"type": "json_object"} - ) - - raw_content = response.choices[0].message.content - if not raw_content or not raw_content.strip(): - logging.info(f"🧠 [Memory] Пустой ответ от модели, user_id={user_id}") - return - - try: - parsed = json.loads(raw_content) - except json.JSONDecodeError as je: - logging.warning(f"⚠️ [Memory] Невалидный JSON от модели: {je} | raw={raw_content[:200]}") - return - - if not isinstance(parsed, dict): - logging.warning(f"⚠️ [Memory] Ответ не является dict: {type(parsed)}") - return - - facts = parsed.get("facts") - if not isinstance(facts, list): - logging.info(f"🧠 [Memory] Фактов для извлечения не найдено (не список), user_id={user_id}") - return - - if not facts: - logging.info(f"🧠 [Memory] Фактов для извлечения не найдено, user_id={user_id}") - return - - # Защита от спама фактами: максимум 10 за один вызов - facts = facts[:10] - - valid_facts = [] - for fact in facts: - if not isinstance(fact, dict): - continue - - key = fact.get("key", "") - value = fact.get("value", "") - - importance = fact.get("importance", 5) - - if not isinstance(key, str) or not isinstance(value, str): - continue - - key = _sanitize_memory_key(key) - value = _sanitize_memory_value(value) - - if not key or not value: - continue - - try: - importance = int(importance) - except (TypeError, ValueError): - importance = 5 - - importance = min(10, max(1, importance)) - - valid_facts.append((user_id, key, value, importance)) - - if not valid_facts: - return - - # Batch insert в одной транзакции для снижения нагрузки на БД - async with db_pool.acquire(timeout=5.0) as conn: - async with conn.transaction(): - for uid, key, value, importance in valid_facts: - await conn.execute(''' - INSERT INTO user_memory (user_id, memory_key, memory_value, importance, created_at, updated_at) - VALUES ($1, $2, $3, $4, NOW(), NOW()) - ON CONFLICT (user_id, memory_key) - DO UPDATE SET - memory_value = EXCLUDED.memory_value, - importance = EXCLUDED.importance, - updated_at = NOW() - ''', uid, key, value, importance) - - logging.info(f"🧠 [Memory] Сохранен факт: {key}={value} (importance={importance}) для user_id={user_id}") - - except Exception as e: - logging.error(f"❌ [Memory] Ошибка извлечения памяти для user_id={user_id}: {e}") - -# ═══════════════════════════════════════════════════════════════════ -# B4: RAG LITE SYSTEM — API функции -# ═══════════════════════════════════════════════════════════════════ - -def _normalize_knowledge_text(text: str) -> str: - """Нормализация текста знания для дедупликации. - - Удаляет лишние пробелы, приводит к нижнему регистру, - удаляет пунктуацию для сравнения. - """ - if not text: - return "" - text = text.lower().strip() - text = re.sub(r'\s+', ' ', text) - text = re.sub(r'[^\w\s]', '', text) - return text[:500] - - -async def add_knowledge(user_id: int, content: str, source: str = "user_message", importance: int = 5) -> bool: - """Добавление знания в knowledge_base с проверкой дедупликации и лимитов. - - Возвращает True если знание сохранено, False если пропущено (дубликат или ошибка). - """ - if not db_pool: - logging.warning("⚠️ [B4 Knowledge] add_knowledge skipped: db_pool unavailable") - return False - - if not content or not content.strip(): - logging.info(f"📝 [B4 Knowledge] add_knowledge skipped: empty content for user_id={user_id}") - return False - - content = content.strip() - if len(content) > 500: - content = content[:500] - - try: - importance = min(10, max(1, int(importance))) - except (TypeError, ValueError): - importance = 5 - - normalized = _normalize_knowledge_text(content) - if not normalized: - logging.info(f"📝 [B4 Knowledge] add_knowledge skipped: normalized empty for user_id={user_id}") - return False - - lock = await knowledge_lock_manager.get(user_id) - try: - async with lock: - async with db_pool.acquire(timeout=5.0) as conn: - # Проверка дедупликации: ищем по нормализованному content - existing = await conn.fetchval( - """ - SELECT COUNT(*) FROM knowledge_base - WHERE user_id = $1 - AND LOWER(REGEXP_REPLACE(content, '[^\\w\\s]', '', 'g')) = $2 - """, - user_id, normalized - ) - - if existing and existing > 0: - logging.info(f"📝 [B4 Knowledge] add_knowledge skipped: duplicate for user_id={user_id} | content_hash={normalized[:50]}") - return False - - # Проверка лимита - count = await conn.fetchval( - "SELECT COUNT(*) FROM knowledge_base WHERE user_id = $1", - user_id - ) - - if count and count >= KNOWLEDGE_MAX_PER_USER: - # Пытаемся удалить старые записи с importance <= 3 - deleted = await conn.execute( - """ - DELETE FROM knowledge_base - WHERE user_id = $1 - AND importance <= $2 - AND id IN ( - SELECT id FROM knowledge_base - WHERE user_id = $1 AND importance <= $2 - ORDER BY created_at ASC - LIMIT 1 - ) - """, - user_id, KNOWLEDGE_CLEANUP_IMPORTANCE_THRESHOLD - ) - deleted_count = int(deleted.split()[-1]) if deleted and deleted.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🧹 [B4 Knowledge] Cleanup: removed {deleted_count} old low-importance records for user_id={user_id}") - else: - logging.warning(f"⚠️ [B4 Knowledge] Limit reached ({count}/{KNOWLEDGE_MAX_PER_USER}) for user_id={user_id}, all records are high-importance. Keeping as is.") - return False - - # Вставка нового знания - await conn.execute( - ''' - INSERT INTO knowledge_base (user_id, content, source, importance, created_at) - VALUES ($1, $2, $3, $4, NOW()) - ''', - user_id, content, source, importance - ) - - logging.info(f"💾 [B4 Knowledge] Saved: user_id={user_id} | importance={importance} | content={content[:80]}...") - return True - - except Exception as e: - logging.error(f"❌ [B4 Knowledge] add_knowledge error for user_id={user_id}: {e}") - return False - finally: - await knowledge_lock_manager.release(user_id) - - -async def get_knowledge(user_id: int, limit: int = 100) -> list: - """Получение всех знаний пользователя, отсортированных по важности и дате. - - Возвращает список dict'ов: [{'id', 'content', 'source', 'importance', 'created_at'}, ...] - """ - if not db_pool: - return [] - try: - async with db_pool.acquire(timeout=5.0) as conn: - rows = await conn.fetch( - """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - ORDER BY importance DESC, created_at DESC - LIMIT $2 - """, - user_id, limit - ) - return [ - { - 'id': row['id'], - 'content': row['content'], - 'source': row['source'], - 'importance': row['importance'], - 'created_at': row['created_at'] - } - for row in rows - ] - except Exception as e: - logging.error(f"❌ [B4 Knowledge] get_knowledge error for user_id={user_id}: {e}") - return [] - - -async def delete_knowledge(user_id: int, knowledge_id: int) -> bool: - """Удаление конкретного знания пользователя.""" - if not db_pool: - return False - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM knowledge_base WHERE id = $1 AND user_id = $2", - knowledge_id, user_id - ) - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🗑 [B4 Knowledge] Deleted: knowledge_id={knowledge_id} for user_id={user_id}") - return True - return False - except Exception as e: - logging.error(f"❌ [B4 Knowledge] delete_knowledge error for user_id={user_id}, id={knowledge_id}: {e}") - return False - - -async def search_knowledge(user_id: int, query: str, limit: int = KNOWLEDGE_RETRIEVAL_LIMIT) -> list: - """Простой keyword-based поиск по knowledge_base. - - B5-ready: интерфейс сохраняется, внутренняя реализация заменится на vector search. - - Возвращает список dict'ов: [{'id', 'content', 'source', 'importance', 'created_at'}, ...] - """ - if not db_pool: - logging.warning("⚠️ [B4 Knowledge] search_knowledge failed: db_pool unavailable") - return [] - - if not query or not query.strip(): - return [] - - try: - # Извлекаем ключевые слова (убираем стоп-слова, короткие слова) - words = re.findall(r'\b[a-zA-Zа-яА-ЯёЁ]{3,}\b', query.lower()) - if not words: - return [] - - # Используем статический SQL с параметризованными LIKE — безопасно от SQL injection - # Поддерживаем до 5 ключевых слов (фиксированное количество параметров) - async with db_pool.acquire(timeout=5.0) as conn: - # Берём максимум 5 ключевых слов - search_words = words[:5] - - # Формируем паттерны для LIKE - patterns = [f"%{word}%" for word in search_words] - - # Выбираем SQL в зависимости от количества слов - if len(search_words) == 1: - sql = """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - AND LOWER(content) LIKE $2 - ORDER BY importance DESC, created_at DESC - LIMIT $3 - """ - params = [user_id, patterns[0], limit] - elif len(search_words) == 2: - sql = """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - AND (LOWER(content) LIKE $2 OR LOWER(content) LIKE $3) - ORDER BY importance DESC, created_at DESC - LIMIT $4 - """ - params = [user_id, patterns[0], patterns[1], limit] - elif len(search_words) == 3: - sql = """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - AND (LOWER(content) LIKE $2 OR LOWER(content) LIKE $3 OR LOWER(content) LIKE $4) - ORDER BY importance DESC, created_at DESC - LIMIT $5 - """ - params = [user_id, patterns[0], patterns[1], patterns[2], limit] - elif len(search_words) == 4: - sql = """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - AND (LOWER(content) LIKE $2 OR LOWER(content) LIKE $3 OR LOWER(content) LIKE $4 OR LOWER(content) LIKE $5) - ORDER BY importance DESC, created_at DESC - LIMIT $6 - """ - params = [user_id, patterns[0], patterns[1], patterns[2], patterns[3], limit] - else: # 5 слов - sql = """ - SELECT id, content, source, importance, created_at - FROM knowledge_base - WHERE user_id = $1 - AND (LOWER(content) LIKE $2 OR LOWER(content) LIKE $3 OR LOWER(content) LIKE $4 OR LOWER(content) LIKE $5 OR LOWER(content) LIKE $6) - ORDER BY importance DESC, created_at DESC - LIMIT $7 - """ - params = [user_id, patterns[0], patterns[1], patterns[2], patterns[3], patterns[4], limit] - - rows = await conn.fetch(sql, *params) - - results = [ - { - 'id': row['id'], - 'content': row['content'], - 'source': row['source'], - 'importance': row['importance'], - 'created_at': row['created_at'] - } - for row in rows - ] - - logging.info(f"🔍 [B4 Knowledge] Search: user_id={user_id} | query_words={search_words} | found={len(results)} | limit={limit}") - return results - - except Exception as e: - logging.error(f"❌ [B4 Knowledge] search_knowledge error for user_id={user_id}: {e}") - return [] - - -async def cleanup_old_knowledge(days: int = 180): - """Автоматическая очистка знаний старше N дней. - - Бизнес-контекст: knowledge_base — долгосрочная память, но устаревшие - знания (старше 6 месяцев) могут потерять актуальность. - """ - if not db_pool: - return - try: - async with db_pool.acquire(timeout=5.0) as conn: - result = await conn.execute( - "DELETE FROM knowledge_base WHERE created_at < NOW() - $1 * INTERVAL '1 day'", - days - ) - deleted_count = int(result.split()[-1]) if result and result.startswith("DELETE") else 0 - if deleted_count > 0: - logging.info(f"🧹 [B4 Knowledge] cleanup_old_knowledge: удалено {deleted_count} записей старше {days} дней.") - except Exception as e: - logging.error(f"❌ [B4 Knowledge] Ошибка в cleanup_old_knowledge: {e}") - - -async def get_knowledge_count(user_id: int) -> int: - """Получение количества знаний пользователя.""" - if not db_pool: - return 0 - try: - async with db_pool.acquire(timeout=5.0) as conn: - return await conn.fetchval( - "SELECT COUNT(*) FROM knowledge_base WHERE user_id = $1", - user_id - ) or 0 - except Exception as e: - logging.error(f"❌ [B4 Knowledge] get_knowledge_count error for user_id={user_id}: {e}") - return 0 - - -# ═══════════════════════════════════════════════════════════════════ -# B4: RAG LITE SYSTEM — Извлечение и сохранение знаний -# ═══════════════════════════════════════════════════════════════════ - -async def extract_knowledge(user_text: str) -> list: - """Извлечение знаний из сообщения пользователя через GPT. - - Анализирует сообщение пользователя и возвращает список знаний. - Каждый элемент: {'content': str, 'importance': int} - - Защита от галлюцинаций: извлекает ТОЛЬКО из user_text, - никогда не сохраняет ответы модели или системные промпты. - """ - if not user_text or not user_text.strip(): - logging.info("📝 [B4 Knowledge] extract_knowledge: empty user_text") - return [] - - model_name = "gpt-4o-mini" - role = "knowledge_extractor" - - try: - async with openai_semaphore: - response = await client.chat.completions.create( - model=model_name, - messages=[ - {"role": "system", "content": KNOWLEDGE_EXTRACTION_PROMPT}, - {"role": "user", "content": user_text} - ], - temperature=0.1, - max_tokens=800, - response_format={"type": "json_object"} - ) - - raw_content = response.choices[0].message.content - if not raw_content or not raw_content.strip(): - logging.info("📝 [B4 Knowledge] extract_knowledge: empty response from model") - return [] - - try: - parsed = json.loads(raw_content) - except json.JSONDecodeError as je: - logging.warning(f"⚠️ [B4 Knowledge] extract_knowledge: invalid JSON from model: {je} | raw={raw_content[:200]}") - return [] - - if not isinstance(parsed, dict): - logging.warning(f"⚠️ [B4 Knowledge] extract_knowledge: response is not dict: {type(parsed)}") - return [] - - knowledge_items = parsed.get("knowledge") - if not isinstance(knowledge_items, list): - logging.info("📝 [B4 Knowledge] extract_knowledge: no knowledge list found") - return [] - - if not knowledge_items: - logging.info("📝 [B4 Knowledge] extract_knowledge: no knowledge extracted") - return [] - - # Защита от спама: максимум KNOWLEDGE_MAX_EXTRACTED_PER_CALL за один вызов - knowledge_items = knowledge_items[:KNOWLEDGE_MAX_EXTRACTED_PER_CALL] - - valid_items = [] - for item in knowledge_items: - if not isinstance(item, dict): - continue - - content = item.get("content", "") - importance = item.get("importance", 5) - - if not isinstance(content, str) or not content.strip(): - continue - - content = content.strip() - if len(content) > 500: - content = content[:500] - - try: - importance = int(importance) - except (TypeError, ValueError): - importance = 5 - - importance = min(10, max(1, importance)) - - valid_items.append({ - "content": content, - "importance": importance - }) - - logging.info(f"🧠 [B4 Knowledge] Extracted: {len(valid_items)} items from text") - return valid_items - - except Exception as e: - logging.error(f"❌ [B4 Knowledge] extract_knowledge error: {e}") - return [] - - -async def save_extracted_knowledge(user_id: int, user_text: str) -> int: - """Извлечение и сохранение знаний из сообщения пользователя. - - Вызывает extract_knowledge и сохраняет результаты в knowledge_base. - Возвращает количество успешно сохранённых знаний. - - Защита от галлюцинаций: сохраняет только подтвержденные данные - из сообщений пользователя. - """ - if not db_pool: - logging.warning("⚠️ [B4 Knowledge] save_extracted_knowledge skipped: db_pool unavailable") - return 0 - - if not user_text or not user_text.strip(): - logging.info("📝 [B4 Knowledge] save_extracted_knowledge skipped: empty text") - return 0 - - # Извлекаем знания - knowledge_items = await extract_knowledge(user_text) - - if not knowledge_items: - logging.info(f"📝 [B4 Knowledge] save_extracted_knowledge: no knowledge to save for user_id={user_id}") - return 0 - - saved_count = 0 - skipped_count = 0 - - for item in knowledge_items: - result = await add_knowledge( - user_id=user_id, - content=item["content"], - source="user_message", - importance=item["importance"] - ) - if result: - saved_count += 1 - else: - skipped_count += 1 - - logging.info( - f"💾 [B4 Knowledge] save_extracted_knowledge: user_id={user_id} | " - f"extracted={len(knowledge_items)} | saved={saved_count} | skipped={skipped_count}" - ) - return saved_count - -async def get_user_data(user_id: int, username: str) -> tuple: - if not db_pool: return (0, 1, "Free") - try: - async with db_pool.acquire(timeout=5.0) as conn: - row = await conn.fetchrow("SELECT generation_count, max_limit, tier FROM users WHERE user_id = $1", user_id) - if row is None: - await conn.execute( - "INSERT INTO users (user_id, username, created_at, tier, max_limit) VALUES ($1, $2, $3, 'Free', 1)", - user_id, username, datetime.now().isoformat() - ) - return (0, 1, "Free") - return (row['generation_count'], row['max_limit'], row['tier']) - except asyncio.TimeoutError: - logging.error("❌ Превышено время ожидания от БД (Timeout) в get_user_data") - return (0, 1, "Free") - except Exception as e: - logging.error(f"❌ Ошибка в get_user_data: {e}") - return (0, 1, "Free") - - -async def increment_limit(user_id: int): - if not db_pool: return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute("UPDATE users SET generation_count = generation_count + 1 WHERE user_id = $1", user_id) - except Exception as e: - logging.error(f"❌ Ошибка в increment_limit: {e}") - - -async def update_user_subscription(user_id: int, tier: str, max_limit: int): - if not db_pool: return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute( - "UPDATE users SET tier = $1, max_limit = $2, generation_count = 0 WHERE user_id = $3", - tier, max_limit, user_id - ) - logging.info(f"💎 Подписка успешно обновлена в облаке для пользователя {user_id}") - except Exception as e: - logging.error(f"❌ Ошибка в update_user_subscription: {e}") - - -async def get_all_users_count() -> int: - if not db_pool: return 0 - try: - async with db_pool.acquire(timeout=5.0) as conn: - return await conn.fetchval("SELECT COUNT(*) FROM users") - except Exception as e: - logging.error(f"❌ Ошибка в get_all_users_count: {e}") - return 0 - - -async def reset_user_limit(user_id: int): - if not db_pool: return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute("UPDATE users SET generation_count = 0 WHERE user_id = $1", user_id) - logging.info(f"🔄 Счётчик лимитов успешно обнулен для пользователя {user_id}") - except Exception as e: - logging.error(f"❌ Ошибка в reset_user_limit: {e}") - - -async def save_ticket(ticket_id: int, user_id: int): - if not db_pool: return - try: - async with db_pool.acquire(timeout=5.0) as conn: - await conn.execute( - "INSERT INTO tickets (ticket_id, user_id, created_at) VALUES ($1, $2, $3)", - ticket_id, user_id, datetime.now().isoformat() - ) - except Exception as e: - logging.error(f"❌ Ошибка в save_ticket: {e}") - - -async def get_ticket_user(ticket_id: int) -> int: - if not db_pool: return None - try: - async with db_pool.acquire(timeout=5.0) as conn: - return await conn.fetchval("SELECT user_id FROM tickets WHERE ticket_id = $1", ticket_id) - except Exception as e: - logging.error(f"❌ Ошибка в get_ticket_user: {e}") - return None - - -# <<< END OF BLOCK 3 >>> -# >>> START OF BLOCK 4: FSM STATES, AI GENERATION AND DOCX CREATION <<< - -class SprintForm(StatesGroup): - niche = State() - expert = State() - audience = State() - funnel = State() - -class SupportForm(StatesGroup): - subject = State() - description = State() - -class SupportStates(StatesGroup): - in_support = State() - -SUPPORT_SYSTEM_PROMPT = """ -Ты — умный ИИ-ассистент тех поддержки Telegram-бота "AI-Sprint Bot". -Твой контекст: этот бот написан на Python (aiogram) и развернут в облаке. Он помогает маркетологам и экспертам создавать стратегии запусков (AI Sprint) и выгружает их в .docx. -У бота есть тарифы: Plus (15 генераций — 250 ₽), Pro (25 генераций — 350 ₽), Ultra (40 генераций — 450 ₽). Есть 1 бесплатная тест-генерация. -Твоя задача — вежливо отвечать на вопросы пользователей. скажи оператор рассмотрит эту заявку. Если вопрос технический, дай базовый совет и подтверди, что передал тикет человеку. -Если вопрос касается сотрудничества скажи то что оператор напишет вам в лс -""" - -# ═══════════════════════════════════════════════════════════════════ -# AI RESPONSE GENERATOR (с логированием, retry-логикой, latency, B2 context) -# ═══════════════════════════════════════════════════════════════════ -async def generate_ai_response(answers: dict, user_id: int) -> str: - """Генерация стратегии запуска с полной observability и контекстом B2. - - Бизнес-контекст: каждый вызов — это списание токенов = деньги. - Мы логируем latency, tokens, errors для оптимизации модели и cost control. - B2: теперь с персонализацией через память и continuity через историю. - """ - system_prompt = """Ты — Senior AI-упаковщик запусков для продюсеров, контент-команд и launch-менеджеров рынка онлайн-образования СНГ. Твоя задача — скоростная упаковка сырых смыслов в продающий каркас, готовый к передаче команде без участия стратега и копирайтера. Ты даёшь не глубокий аудит, а рабочую основу: аватары, офферы, hooks, микро-воронку и первые шаги внедрения. - -Твой ответ будет автоматически сконвертирован в чистый DOCX-документ. -КРИТИЧЕСКИ ВАЖНО соблюдать правила оформления: - -1. Никакого Markdown (никаких ```, блоков кода и т.д.). -2. Никаких таблиц, графиков и спецсимволов. -3. Никаких эмодзи и смайлов — строго запрещено. -4. Только обычный текст, абзацы, нумерованные и маркированные списки (дефисы). -5. Заголовки выделяй ЗАГЛАВНЫМИ БУКВАМИ — так они визуально читаются в документе. -6. Пиши исключительно на русском языке. - -ПРАВИЛА КОНТЕНТА И СТИЛЯ: - -· Сразу к делу. Никаких приветствий, воды, фраз вроде «Вот ваш анализ» или «Надеюсь, это поможет». Первое слово — первый заголовок. Последнее слово — последний пункт плана внедрения. -· Не пересказывай вводные. Сразу выдавай решения. -· Стиль: циничный, деловой, конверсионный, без инфоцыганских клише и академических рассуждений. Каждая фраза должна быть готова к использованию в контенте или ТЗ. -· Работай как безжалостный упаковщик: сними тревогу перед запуском, дай команде готовые формулировки. -· Если данных от пользователя мало — достраивай гипотезы уверенно, опираясь на логику рынка СНГ (баннерная слепота, недоверие, жажда быстрых результатов, хроническая нехватка времени у продюсеров). - -СТРУКТУРА ДОКУМЕНТА (СТРОГО 5 БЛОКОВ И ДЕТАЛЬНЫЕ ПОДБЛОКИ В КАЖДОМ): - -БЛОК 1. АВАТАРЫ ЦЕЛЕВОЙ АУДИТОРИИ (3 САМЫХ ПЛАТЕЖЕСПОСОБНЫХ СЕГМЕНТА) -Для каждого сегмента обязательно раскрой все 9 подпунктов: - -· Название сегмента (ёмко, как называют себя сами люди). -· Главная боль (что реально не даёт спать, самая острая формулировка). -· Глубинное желание (истинный мотив купить, часто неосознаваемый). -· Язык сегмента (2–3 точных фразы, которыми они описывают проблему). -· Текущая альтернатива (чем они решают проблему сейчас, почему это не работает). -· Триггерное событие (конкретная ситуация, после которой они начинают искать решение). -· Ключевое возражение против покупки именно у этого эксперта (скрытое или явное). -· Триггер к немедленной покупке (на что надавить, чтобы заплатили сейчас). -· Идеальный канал касания (где конкретно они увидят первый контакт: рилс, сторис, Telegram, таргет). - -БЛОК 2. УБОЙНЫЕ ОФФЕРЫ (3 УГЛА ПРОДАЖ) -Для каждого оффера обязательно выдай 6 подпунктов: - -· Тип оффера (рациональный, эмоциональный, быстрого результата). -· Формулировка по формуле «Результат + Срок + Снятие главного страха» (1–2 предложения, сразу готовые в текст). -· Конкретное наполнение (3–4 пункта, что именно получит клиент: модули, шаги, инструменты). -· Механика снятия риска (гарантия, пробный период, возврат, демо-доступ). -· Ценностный якорь (с чем сравнить цену, чтобы она казалась оправданной: стоимость часа эксперта, цена ошибки, стоимость аналога). -· Дедлайн или ограничитель (почему нужно решаться сейчас: набор мест, повышение цены, старт группы). - -БЛОК 3. ИДЕИ ДЛЯ КОНТЕНТА И REELS (5 ГОТОВЫХ СЦЕНАРИЕВ) -Каждый сценарий расширь до 5 обязательных элементов: - -· Формат (Reels, Stories, пост в Telegram — выбрать один оптимальный). -· Хук (фраза для первых 2 секунд, которая режет скролл). -· Вскрытие боли (2–3 предложения, попадающих в реальную ситуацию зрителя). -· Микро-решение/инсайт (короткая ценность, которая формирует доверие и показывает экспертность). -· Призыв к действию (CTA) — конкретный шаг в воронку: «напиши слово…», «переходи по ссылке», «смотри следующее видео». - -БЛОК 4. БЫСТРАЯ МИКРО-ВОРОНКА -Опиши путь клиента детально, по шагам, добавив пропущенные элементы: - -· Шаг 0. Источник трафика (откуда берём людей: таргет, блог эксперта, партнёрская рассылка). -· Шаг 1. Точка входа и лид-магнит (конкретный формат: чек-лист, мини-урок, тест-драйв, вебинар; что именно отдаём). -· Шаг 1.1 Квалификация (как отсекаем нецелевых: вопрос в директ, анкета, бот). -· Шаг 2. Прогрев (какую одну ключевую ценность и какой контент выдать, чтобы снять основное возражение и создать желание). -· Шаг 2.1 Вовлекающий элемент (квиз, голосование, задание — то, что заставляет взаимодействовать). -· Шаг 3. Конверсионный этап (где и как предлагаем оплатить: ссылка на закрытый канал, прямой эфир с триггером, чат-бот с кассой; точная формулировка оффера в момент предложения). - -БЛОК 5. ПЛАН ВНЕДРЕНИЯ (QUICK WINS) -Ровно 5 конкретных, физических действий, которые владелец проекта или продюсер может выполнить в ближайшие 24 часа, чтобы запустить связку: - -· Действие 1: Быстрый тест оффера (кому отправить 3 вопроса для проверки спроса). -· Действие 2: Запись первого контента (какое именно видео/пост сделать прямо сейчас, с готовым скриптом). -· Действие 3: Технический минимум (что настроить: чат-бот, сбор контактов, ссылку на оплату). -· Действие 4: Сбор возражений (где взять 3–5 живых возражений ЦА для доработки прогрева). -· Действие 5: Запуск микро-трафика (конкретная сумма и площадка для тестового таргета или рассылки).""" - - current_user_message = f"Ниша: {answers.get('niche')}\nЭксперт: {answers.get('expert')}\nЦА: {answers.get('audience')}\nВоронка: {answers.get('funnel')}" - model_name = "gpt-4o-mini" - role = "sprint_generator" - - # B2.1: Сохраняем текущий user message в историю ПЕРЕД запросом - await save_conversation_history(user_id, "user", current_user_message) - - # B2.1: Собираем контекст с памятью и историей - user_context = await build_user_context(user_id, current_user_message) - - # B2.1: Формируем messages с контекстом в system prompt - context_enhanced_prompt = system_prompt + "\n\n" + user_context - messages = [ - {"role": "system", "content": context_enhanced_prompt}, - {"role": "user", "content": current_user_message} - ] - - - max_retries = 3 - backoff_delays = [1.0, 2.0, 4.0] - - for attempt in range(max_retries): - start_time = time.perf_counter() - prompt_tokens = 0 - completion_tokens = 0 - logging.info(f"📡 [GPT Запрос] Юзер: {user_id} | Модель: {model_name} | Попытка: {attempt + 1}/{max_retries}") - - try: - async with openai_semaphore: - response = await client.chat.completions.create( - model=model_name, - messages=messages, - temperature=0.75, - max_tokens=2500 - ) - latency_ms = int((time.perf_counter() - start_time) * 1000) - prompt_tokens = response.usage.prompt_tokens if response.usage else 0 - completion_tokens = response.usage.completion_tokens if response.usage else 0 - - logging.info(f"✅ [GPT Успех] Юзер: {user_id} | Latency: {latency_ms}ms | Tokens: {prompt_tokens}p + {completion_tokens}c") - - await save_ai_log( - user_id=user_id, model=model_name, role=role, - prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, - latency_ms=latency_ms, status_code=200, success=True, error=None - ) - - ai_response = response.choices[0].message.content - - # B2.1: Сохраняем ответ AI в историю - await save_conversation_history(user_id, "assistant", ai_response) - - # B4: Извлекаем знания из сообщения пользователя (фоново, не блокирует) - safe_background_task( - save_extracted_knowledge(user_id, current_user_message), - name=f"knowledge_extract_{user_id}" - ) - - # Асинхронное извлечение памяти без блокировки ответа пользователю - safe_background_task(extract_and_save_memory(user_id, current_user_message, ai_response), name="memory_extract") - - return ai_response - - - except Exception as e: - latency_ms = int((time.perf_counter() - start_time) * 1000) - status_code = getattr(e, "status_code", 500) - error_msg = str(e) - - logging.error(f"❌ [GPT Ошибка] Юзер: {user_id} | Попытка {attempt + 1} из {max_retries} провалена: {error_msg}") - - await save_ai_log( - user_id=user_id, model=model_name, role=role, - prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, - latency_ms=latency_ms, status_code=status_code, success=False, error=error_msg - ) - - if attempt < max_retries - 1: - sleep_time = backoff_delays[attempt] - logging.info(f"⏳ Ожидание {sleep_time}s перед повторным запросом...") - await asyncio.sleep(sleep_time) - else: - logging.critical(f"🚨 [GPT Фатал] Достигнут лимит попыток для пользователя {user_id}") - raise - -def create_docx_sync(text: str, username: str) -> BufferedInputFile: - """Синхронная функция, генерирующая DOCX, изолированная от Event Loop. - - Бизнес-контекст: python-docx — CPU-bound операция. Запуск в отдельном - потоке через asyncio.to_thread() предотвращает блокировку Event Loop - на 100-500ms, критично при 100+ параллельных пользователях. - """ - clean_text = re.sub(r'[*_#`~]', '', text) - - doc = Document() - doc.add_heading('AI Sprint: Премиальная Упаковка Запуска', 0) - doc.add_paragraph(f'Сгенерировано для: @{username}\nДата: {datetime.now().strftime("%d.%m.%Y")}\n') - doc.add_paragraph(clean_text) - - file_stream = BytesIO() - doc.save(file_stream) - file_stream.seek(0) - return BufferedInputFile(file_stream.read(), filename="AI_Sprint_Premium_Result.docx") - -# <<< END OF BLOCK 4 >>> -# >>> START OF BLOCK 5: ADMIN PANEL AND ANONYMOUS TICKET REPLIES <<< - -@dp.message(Command("admin")) -async def admin_panel(message: types.Message): - if message.from_user.username != ADMIN_USERNAME and message.from_user.id != ADMIN_ID: - return - stats = await get_all_users_count() - text = ( - f"👑 Панель управления (Admin)\n\n" - f"Всего пользователей в базе: {stats}\n\n" - f"🔹 /give_sub <user_id> <Тариф> — выдать тариф (Plus, Pro, Ultra)\n" - f"🔹 /reset_limit <user_id> — обнулить израсходованные лимиты (сохранив тариф)\n" - f"🔹 /reset_free <user_id> — обнулить и жестко вернуть лимит в 1 генерацию (Free)" - ) - await message.answer(text, parse_mode="HTML") - -@dp.message(Command("give_sub")) -async def admin_give_sub(message: types.Message): - if message.from_user.username != ADMIN_USERNAME and message.from_user.id != ADMIN_ID: - return - args = message.text.split() - if len(args) < 3 or not args[1].isdigit(): - await message.answer("Формат: /give_sub 123456789 Pro", parse_mode="HTML") - return - target_id, tier_name = args[1], args[2].capitalize() - limits = {"Free": 1, "Plus": 15, "Pro": 25, "Ultra": 40} - if tier_name not in limits: - await message.answer("❌ Нет такого тарифа!") - return - await update_user_subscription(int(target_id), tier_name, limits[tier_name]) - await message.answer(f"✅ Пользователю {target_id} установлен тариф {tier_name}") - -@dp.message(Command("reset_limit")) -async def admin_reset_limit_cmd(message: types.Message): - if message.from_user.username != ADMIN_USERNAME and message.from_user.id != ADMIN_ID: - return - args = message.text.split() - if len(args) < 2 or not args[1].isdigit(): - await message.answer("Формат: /reset_limit 123456789", parse_mode="HTML") - return - target_id = int(args[1]) - await reset_user_limit(target_id) - await message.answer(f"✅ Счетчик израсходованных генераций для {target_id} успешно обнулен. Текущий тариф сохранен!") - -@dp.message(Command("reset_free")) -async def admin_reset_free_cmd(message: types.Message): - if message.from_user.username != ADMIN_USERNAME and message.from_user.id != ADMIN_ID: - return - args = message.text.split() - if len(args) < 2 or not args[1].isdigit(): - await message.answer("Формат: /reset_free 123456789", parse_mode="HTML") - return - target_id = int(args[1]) - await update_user_subscription(target_id, "Free", 1) - await message.answer(f"✅ Пользователю {target_id} возвращен базовый тариф Free.") - -@dp.message(F.reply_to_message & ((F.from_user.id == ADMIN_ID) | (F.from_user.username == ADMIN_USERNAME))) -async def admin_reply_handler(message: types.Message): - original_text = message.reply_to_message.text or message.reply_to_message.caption - if not original_text or "[TICKET_ID:" not in original_text: - return - - try: - ticket_id_str = original_text.split("[TICKET_ID:")[1].split("]")[0] - ticket_id = int(ticket_id_str) - user_chat_id = await get_ticket_user(ticket_id) - - if user_chat_id: - await bot.send_message( - chat_id=user_chat_id, - text=f"👨‍💻 Ответ от оператора техподдержки:\n\n{html.escape(message.text)}", - parse_mode="HTML" - ) - await message.answer("✅ Ответ успешно доставлен пользователю анонимно.") - else: - await message.answer("❌ Ошибка: тикет не найден в базе данных.") - except Exception as e: - await message.answer(f"❌ Ошибка отправки: {e}") - -# <<< END OF BLOCK 5 >>> -# >>> START OF BLOCK 6: NAVIGATION COMMANDS AND SUBSCRIPTION MENU <<< - -@dp.message(Command("profile")) -async def cmd_profile(message: types.Message): - user_id = message.from_user.id - username = message.from_user.username or str(user_id) - count, max_limit, tier = await get_user_data(user_id, username) - text = ( - f"👤 Ваш профиль: @{username}\n" - f"🆔 ID: {user_id}\n\n" - f"💎 Текущий тариф: {tier}\n" - f"📊 Израсходовано лимитов: {count} из {max_limit}\n\n" - f"💡 Когда лимит генераций закончится, вы можете приобрести любой тариф заново." - ) - await message.answer(text, parse_mode="HTML") - -@dp.message(Command("help")) -async def cmd_help(message: types.Message): - text = ( - "🛠 Доступные команды:\n\n" - "🔸 /start — Главное меню\n" - "🔸 /profile — Узнать свой тариф\n" - "🔸 /help — Показать это сообщение\n" - ) - await message.answer(text, parse_mode="HTML") - -@dp.message(CommandStart()) -async def cmd_start(message: types.Message, state: FSMContext): - await state.clear() - text = ( - "👋 Добро пожаловать в AI-Sprint Bot!\n\n" - "Я — ваш личный AI-маркетолог. Помогу разработать упакованную концепцию.\n\n" - "💳 Тарифные планы:\n" - "🟢 Plus (15 генераций) — 250 ₽\n" - "🟢 Pro (25 генераций) — 350 ₽\n" - "🟢 Ultra (40 генераций) — 450 ₽\n\n" - "Доступна 1 бесплатная генерация для теста!\n" - ) - keyboard = InlineKeyboardMarkup( - inline_keyboard=[ - [InlineKeyboardButton(text="🚀 Начать AI Sprint", callback_data="start_sprint")], - [InlineKeyboardButton(text="🔥 Купить подписку", callback_data="choose_tier")], - [InlineKeyboardButton(text="💡 О боте", callback_data="about_bot"), InlineKeyboardButton(text="🆘 Поддержка", callback_data="open_support")] - ] - ) - await message.answer(text, reply_markup=keyboard, parse_mode="HTML") - - -@dp.callback_query(F.data == "about_bot") -async def process_about_bot(callback: CallbackQuery): - await callback.answer() - text = ( - "🤖 Об AI-Sprint Bot — Твоем AI-Упаковщике Запусков\n\n" - "Это скоростной инструмент для продюсеров, маркетологов и экспертов, который за 10–15 минут превращает сырые смыслы в готовый продающий каркас проекта, экономя недели работы и десятки тысяч рублей на копирайтерах.\n\n" - "📦 ЧТО ВЫ ПОЛУЧАЕТЕ НА ВЫХОДЕ (В СТРУКТУРИРОВАННОМ DOCX):\n" - "Вы отвечаете на 4 вопроса — ИИ выдает плотный бизнес-документ без воды:\n" - "• 3 аватара ЦА: истинные боли, скрытые мотивы, триггеры и дословный язык рынка.\n" - "• 3 убойных оффера: рациональный, эмоциональный, быстрого результата + снятие рисков.\n" - "• 5 контент-сценариев (Reels/Stories/ТГ): цепляющие хуки, вскрытие боли и четкий CTA.\n" - "• Микро-воронка: пошаговый путь клиента от источника трафика до кассы.\n" - "• План Quick Wins: ровно 5 физических действий на ближайшие 24 часа.\n\n" - "🎯 КОМУ ЭТО СЭКОНОМИТ ВРЕМЯ И ДЕНЬГИ:\n" - "• Продюсерам онлайн-школ: когда сроки горят, а эксперт не может внятно описать продукт.\n" - "• Контент-командам: готовый фундамент прогревов и четкие ТЗ для сценаристов.\n" - "• Экспертам на самозапуске: быстрый способ вытащить из себя смыслы и убрать кашу из головы.\n\n" - "⚡ Бот думает как циничный и жесткий маркетолог рынка СНГ. Если вводных данных мало — он сам уверенно достроит рабочие гипотезы!\n\n" - "Запустите сессию прямо сейчас, разгрузите команду и заберите готовую стратегию за 10 минут!" - ) - await bot.send_message(chat_id=callback.from_user.id, text=text, parse_mode="HTML") - - -@dp.callback_query(F.data == "choose_tier") -async def process_choose_tier(callback: CallbackQuery): - await callback.answer() - text = "📋 Выберите подходящий тариф для покупки:\n\n" - keyboard = InlineKeyboardMarkup( - inline_keyboard=[ - [InlineKeyboardButton(text="➕ Тариф Plus (15 лимитов) — 250 ₽", callback_data="buy_tier_plus")], - [InlineKeyboardButton(text="💎 Тариф Pro (25 лимитов) — 350 ₽", callback_data="buy_tier_pro")], - [InlineKeyboardButton(text="🚀 Тариф Ultra (40 лимитов) — 450 ₽", callback_data="buy_tier_ultra")] - ] - ) - await bot.send_message(chat_id=callback.from_user.id, text=text, reply_markup=keyboard, parse_mode="HTML") - -@dp.callback_query(F.data.startswith("buy_tier_")) -async def process_invoice_real(callback: CallbackQuery): - await callback.answer() - tier = callback.data.split("_")[2].capitalize() - user_id = callback.from_user.id - prices = {"Plus": 250, "Pro": 350, "Ultra": 450} - amount = prices.get(tier, 250) - wait_msg = await bot.send_message(chat_id=user_id, text="⏳ Создаю безопасную ссылку...", parse_mode="HTML") - payment_url = await create_yookassa_payment(amount, f"Оплата тарифа {tier}", user_id, tier) - await bot.delete_message(chat_id=user_id, message_id=wait_msg.message_id) - if payment_url: - keyboard = InlineKeyboardMarkup(inline_keyboard=[[InlineKeyboardButton(text=f"💳 Оплатить {amount} ₽", url=payment_url)]]) - await bot.send_message(chat_id=user_id, text=f"🧾 Оформление заказа\n\n👇 Нажмите кнопку ниже:", reply_markup=keyboard, parse_mode="HTML") - else: - await bot.send_message(chat_id=user_id, text="❌ Ошибка связи с платежной системой.") - -# <<< END OF BLOCK 6 >>> -# >>> START OF BLOCK 7: HYBRID SUPPORT SYSTEM <<< - -# ═══════════════════════════════════════════════════════════════════ -# FORM-BASED SUPPORT (Тикеты с ИИ-анализом и B2 контекстом) -# ═══════════════════════════════════════════════════════════════════ -@dp.callback_query(F.data == "open_support") -async def process_open_support(callback: CallbackQuery, state: FSMContext): - await callback.answer() - # Защита: если пользователь в режиме чат-поддержки — сбрасываем - current_state = await state.get_state() - if current_state == SupportStates.in_support.state: - await state.clear() - await state.set_state(SupportForm.subject) - await bot.send_message(chat_id=callback.from_user.id, text="🆘 Напишите Тему вашей проблемы (макс. 100 символов).", parse_mode="HTML") - -@dp.message(SupportForm.subject) -async def support_subject(message: types.Message, state: FSMContext): - if len(message.text) > 100: - await message.answer("❌ Тема слишком длинная.") - return - await state.update_data(subject=message.text) - await state.set_state(SupportForm.description) - await message.answer("Опишите вашу проблему (макс. 1000 символов):", parse_mode="HTML") - -@dp.message(SupportForm.description) -async def support_description(message: types.Message, state: FSMContext): - if len(message.text) > 1000: - await message.answer("❌ Описание слишком длинное.") - return - user = message.from_user - data = await state.get_data() - subject = data['subject'] - description = message.text - await state.clear() - ticket_id = random.randint(10000, 99999) - await save_ticket(ticket_id, user.id) - wait_msg = await message.answer("⏳ Анализирую запрос...") - - model_name = "gpt-4o-mini" - role = "support_analyzer" - max_retries = 3 - backoff_delays = [1.0, 2.0, 4.0] - ai_answer = "Заявка зарегистрирована. Оператор свяжется с вами." - - # B2: Формируем текущее сообщение пользователя - current_user_message = f"Тема: {subject}\nПроблема: {description}" - - # B2.1: Сохраняем в историю ПЕРЕД запросом - await save_conversation_history(user.id, "user", current_user_message) - - # B4: Извлекаем знания из сообщения пользователя (фоново) - safe_background_task( - save_extracted_knowledge(user.id, current_user_message), - name=f"support_knowledge_{user.id}" - ) - - # B2.1: Собираем контекст с памятью и историей - user_context = await build_user_context(user.id, current_user_message) - - # B2.1: Формируем messages с контекстом в system prompt - context_enhanced_prompt = SUPPORT_SYSTEM_PROMPT + "\n\n" + user_context - messages = [ - {"role": "system", "content": context_enhanced_prompt}, - {"role": "user", "content": current_user_message} - ] - - # СТАБИЛИЗАЦИЯ: Внедряем отказоустойчивость и логирование для ИИ-саппорта - for attempt in range(max_retries): - start_time = time.perf_counter() - logging.info(f"📡 [GPT Саппо��т Запрос] Юзер: {user.id} | Попытка: {attempt + 1}/{max_retries}") - try: - async with openai_semaphore: - response = await client.chat.completions.create( - model=model_name, - messages=messages, - temperature=0.4 - ) - latency_ms = int((time.perf_counter() - start_time) * 1000) - prompt_tokens = response.usage.prompt_tokens if response.usage else 0 - completion_tokens = response.usage.completion_tokens if response.usage else 0 - - logging.info(f"✅ [GPT Саппорт Успех] Юзер: {user.id} | Latency: {latency_ms}ms") - ai_answer = response.choices[0].message.content - - # B2.1: Сохраняем ответ ассистента в историю - await save_conversation_history(user.id, "assistant", ai_answer) - - await save_ai_log( - user_id=user.id, model=model_name, role=role, - prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, - latency_ms=latency_ms, status_code=200, success=True, error=None - ) - break - except Exception as e: - latency_ms = int((time.perf_counter() - start_time) * 1000) - status_code = getattr(e, "status_code", 500) - error_msg = str(e) - - logging.error(f"❌ [GPT Саппорт Ошибка] Юзер: {user.id} | Попытка {attempt + 1} провалена: {error_msg}") - - await save_ai_log( - user_id=user.id, model=model_name, role=role, - prompt_tokens=0, completion_tokens=0, - latency_ms=latency_ms, status_code=status_code, success=False, error=error_msg - ) - - if attempt < max_retries - 1: - await asyncio.sleep(backoff_delays[attempt]) - else: - logging.critical(f"🚨 [GPT Саппорт Фатал] Ошибка ИИ поддержки: {e}") - - await wait_msg.edit_text(f"🤖 Ответ ИИ поддержки (#{ticket_id}):\n\n{html.escape(ai_answer)}", parse_mode="HTML") - - # Извлечение памяти из тикета поддержки - combined_support_text = f"Тема: {subject}\nПроблема: {description}" - safe_background_task(extract_and_save_memory(user.id, combined_support_text, ai_answer), name="support_memory") - - safe_name = html.escape(user.full_name) - safe_username = f"@{html.escape(user.username)}" if user.username else "Нет юзернейма" - admin_text = ( - f"🚨 НОВЫЙ ТИКЕТ #{ticket_id}\n\n👤 От: {safe_name} ({safe_username})\n🆔 ID: {user.id}\n📝 Тема: {html.escape(subject)}\n📄 Описание: {html.escape(description)}\n\n🤖 Ответ ИИ: {html.escape(ai_answer)}\n\n✍️ Сделайте REPLY...\n[TICKET_ID:{ticket_id}]" - ) - if ADMIN_ID != 0: - try: - await bot.send_message(chat_id=ADMIN_ID, text=admin_text, parse_mode="HTML") - except Exception: - pass - -# <<< END OF BLOCK 7 >>> -# >>> START OF BLOCK 8: INTERACTIVE MARKETING BRIEF SURVEY <<< -# --- SUB-BLOCK 8.1: SPRINT ENTRY & NICHE QUESTION --- -@dp.callback_query(F.data == "start_sprint") -async def process_start_sprint(callback: CallbackQuery, state: FSMContext): - await callback.answer() - - # Строго берем ID как число - user_id = int(callback.from_user.id) - username = callback.from_user.username or "unknown" - - count, max_limit, tier = await get_user_data(user_id, username) - if count >= max_limit: - await bot.send_message( - chat_id=user_id, - text=f"❌ Ваш лимит исчерпан ({count}/{max_limit} на тарифе {tier}).\nПожалуйста, докупите подписку.", - parse_mode="HTML" - ) - return - - await state.set_state(SprintForm.niche) - q1_text = ( - "🚀 Вопрос 1 из 4: Что за продукт\n\n" - "• Опиши продукт одной фразой, как если бы объяснял коллеге в лифте: что это и для кого.\n" - "• Кто эксперт: бэкграунд, манера подачи, суперсила, почему ему должны поверить?\n" - "• Какую одну главную трансформацию получает клиент? Опиши «до» и «после» бытовым языком.\n" - "• Формат и длительность: вебинары, записи, наставничество, разборы, клуб? Уровень сопровождения?\n" - "• Это новый продукт или переупаковка? Если переупаковка — что изменилось по сравнению с прошлым разом?" - ) - await bot.send_message(chat_id=user_id, text=q1_text, parse_mode="HTML") - - -# --- SUB-BLOCK 8.2: EXPERT & AUDIENCE QUESTIONS --- -@dp.message(SprintForm.niche) -async def process_niche(message: types.Message, state: FSMContext): - await state.update_data(niche=message.text) - await state.set_state(SprintForm.expert) - q2_text = ( - "🎯 Вопрос 2 из 4: Кому это нужно (платёжеспособная ЦА)\n\n" - "• Опиши 2–3 типа людей, которые уже покупали похожее у этого эксперта или конкурентов: род занятий, доход, образ жизни.\n" - "• Приведи 2–3 дословные фразы, которыми они жалуются на свою ситуацию.\n" - "• Что они уже пробовали для решения проблемы и почему это не сработало?\n" - "• Есть ли у них «момент истины» — конкретная ситуация, после которой они начинают искать решение?\n" - "• Сформулируй главную боль без цензуры: что реально не даёт спать, плюс какое глубинное желание за ней скрыто." - ) - await message.answer(q2_text, parse_mode="HTML") - -@dp.message(SprintForm.expert) -async def process_expert(message: types.Message, state: FSMContext): - await state.update_data(expert=message.text) - await state.set_state(SprintForm.audience) - q3_text = ( - "🧠 Вопрос 3 из 4: Возражения и страхи ЦА\n\n" - "• Опиши 2–3 типа людей, которые уже покупали похожее у этого эксперта или конкурентов: род занятий, доход, образ жизни.\n" - "• Приведи 2–3 дословные фразы, которыми они жалуются на свою ситуацию.\n" - "• Что они уже пробовали для решения проблемы и почему это не сработало?\n" - "• Сформулируй главную боль без цензуры: что реально не даёт спать, плюс какое глубинное желание за ней скрыто.\n" - "• Какое ключевое возражение мешает купить именно у этого эксперта и чего клиент боится, если не решит проблему сейчас?" - ) - await message.answer(q3_text, parse_mode="HTML") - -# --- SUB-BLOCK 8.3: FUNNEL QUESTION & AI GENERATION LAUNCH --- - -@dp.message(SprintForm.audience) -async def process_audience(message: types.Message, state: FSMContext): - await state.update_data(audience=message.text) - await state.set_state(SprintForm.funnel) - q4_text = ( - "💸 Вопрос 4 из 4: Модель продаж и точки касания\n\n" - "• Основной источник трафика прямо сейчас и точка входа для клиента: что именно и откуда?\n" - "• Как устроен прогрев после касания: канал, бот, сторис, автовебинар? Опиши путь коротко.\n" - "• Где и как клиент получает предложение купить и через что платит: бот, директ, созвон, касса?\n" - "• Кто закрывает сделки и где клиент чаще всего отваливается до покупки?\n" - "• Какие инструменты уже работают, что хромает и какой бюджет на тест готов выделить прямо сейчас?" - ) - await message.answer(q4_text, parse_mode="HTML") - - -@dp.message(SprintForm.funnel) -async def process_funnel(message: types.Message, state: FSMContext): - await state.update_data(funnel=message.text) - user_data = await state.get_data() - user_id = message.from_user.id - username = message.from_user.username or str(user_id) - await state.clear() - - lock = await user_lock_manager.get(user_id) - # Контекстный менеджер САМ гарантированно снимет лок при любом исходе - async with lock: - processing_msg = None # ФИКС NameError в finally - try: - count, max_limit, tier = await get_user_data(user_id, username) - if count >= max_limit: - await message.answer("❌ Превышен лимит генераций.") - return - - if tier == "Free" and not free_tier_protector.can_generate(): - await message.answer("🛑 Глобальный лимит бесплатных запросов исчерпан!", parse_mode="HTML") - return - - processing_msg = await message.answer("⏳ Модель собирает смыслы и пишет документ...") - await bot.send_chat_action(chat_id=user_id, action="typing") - - ai_text = await generate_ai_response(user_data, user_id) - docx_file = await asyncio.to_thread(create_docx_sync, ai_text, username) - - await bot.send_document( - chat_id=user_id, - document=docx_file, - caption=f"✅ Стратегия готова! Подписка: {tier} (Использовано: {count + 1}/{max_limit})", - parse_mode="HTML" - ) - await increment_limit(user_id) - - # Извлечение памяти из успешной генерации стратегии - combined_sprint_text = f"Ниша: {user_data.get('niche', '')}\nЭксперт: {user_data.get('expert', '')}\nЦА: {user_data.get('audience', '')}\nВоронка: {user_data.get('funnel', '')}" - safe_background_task(extract_and_save_memory(user_id, combined_sprint_text, ai_text), name="sprint_memory") - except Exception as e: - logging.error(f"Ошибка генерации: {e}", exc_info=True) - await message.answer("❌ Ошибка генерации на стороне AI. Попробуйте чуть позже.") - finally: - # Безопасное удаление сообщения - if processing_msg: - try: - await bot.delete_message(chat_id=user_id, message_id=processing_msg.message_id) - except Exception: - pass - # Очистка памяти ПОСЛЕ выхода из async with — лок гарантированно снят - await user_lock_manager.release(user_id) - -# <<< END OF BLOCK 8 >>> -# >>> START OF BLOCK 8.5: YOOKASSA INTEGRATION AND WEBHOOK <<< - -async def create_yookassa_payment(amount: int, description: str, user_id: int, tier: str) -> str: - url = "https://api.yookassa.ru/v3/payments" - idempotence_key = str(uuid.uuid4()) - shop_id = os.getenv("YOOKASSA_SHOP_ID") - secret_key = os.getenv("YOOKASSA_SECRET_KEY") - if not shop_id or not secret_key: return None - auth = aiohttp.BasicAuth(login=shop_id, password=secret_key) - bot_info = await bot.get_me() - payload = { - "amount": {"value": f"{amount}.00", "currency": "RUB"}, - "capture": True, - "confirmation": {"type": "redirect", "return_url": f"https://t.me/{bot_info.username}"}, - "description": description, - "metadata": {"user_id": str(user_id), "tier": tier} - } - try: - async with aiohttp.ClientSession(auth=auth) as yookassa_session: - headers = {"Idempotence-Key": idempotence_key, "Content-Type": "application/json"} - async with yookassa_session.post(url, json=payload, headers=headers) as resp: - if resp.status == 200: - data = await resp.json() - return data["confirmation"]["confirmation_url"] - return None - except Exception: - return None - -YOOKASSA_NETWORKS = [ - ipaddress.ip_network('185.71.76.0/27'), ipaddress.ip_network('185.71.77.0/27'), - ipaddress.ip_network('77.75.153.0/25'), ipaddress.ip_network('77.75.154.128/25'), - ipaddress.ip_network('2a02:5180::/32') -] - -def is_yookassa_ip(ip_str: str) -> bool: - try: - ip = ipaddress.ip_address(ip_str) - return any(ip in net for net in YOOKASSA_NETWORKS) - except ValueError: - return False - -async def yookassa_webhook_handler(request: web.Request): - client_ip = request.headers.get('X-Forwarded-For', request.remote) - if client_ip: client_ip = client_ip.split(',')[0].strip() - if not client_ip or not is_yookassa_ip(client_ip): - return web.Response(status=403, text="Forbidden") - - try: - data = await request.json() - if data.get("event") == "payment.succeeded": - metadata = data.get("object", {}).get("metadata", {}) - user_id_str, tier = metadata.get("user_id"), metadata.get("tier") - if user_id_str and tier: - user_id, limit = int(user_id_str), {"Plus": 15, "Pro": 25, "Ultra": 40}.get(tier, 15) - await update_user_subscription(user_id, tier, limit) - await bot.send_message(chat_id=user_id, text=f"✅ Оплата прошла успешно!\n\nТариф {tier} активен.", parse_mode="HTML") - return web.Response(status=200) - except Exception: - return web.Response(status=500) - -# <<< END OF BLOCK 8.5 >>> -# >>> START OF NEW BLOCK: AI TECH SUPPORT ASSISTANT (TICKETS) <<< - -def get_support_keyboard(): - return InlineKeyboardMarkup(inline_keyboard=[[InlineKeyboardButton(text="❌ Завершить диалог с поддержкой", callback_data="exit_support")]]) - -@dp.message(F.text == "Связаться с поддержкой") -@dp.message(Command("support")) -async def start_support_mode(message: Message, state: FSMContext): - # Защита: если пользователь в форме тикета — сбрасываем - current_state = await state.get_state() - if current_state in (SupportForm.subject.state, SupportForm.description.state): - await state.clear() - await state.set_state(SupportStates.in_support) - await message.answer("🤖 **Вы переключены на ИИ-ассистента техподдержки!**", reply_markup=get_support_keyboard(), parse_mode="Markdown") - -@dp.callback_query(F.data == "exit_support") -async def exit_support_callback(callback: CallbackQuery, state: FSMContext): - if await state.get_state() == SupportStates.in_support.state: - await state.clear() - await callback.message.answer("🔄 **Вы успешно вернулись в основное меню бота.**") - else: - await callback.answer("Вы уже вышли из режима поддержки.", show_alert=True) - -@dp.message(SupportStates.in_support) -async def handle_support_ticket(message: Message, state: FSMContext): - user_text = message.text - if not user_text: - await message.answer("⚠️ Я принимаю только текстовые обращения.") - return - - user_id = message.from_user.id - lock = await user_lock_manager.get(user_id) - - # async with сам гарантированно снимет лок при любом исходе - async with lock: - try: - await message.bot.send_chat_action(chat_id=message.chat.id, action="typing") - system_prompt = "Ты — лаконичный ИИ-ассистент техподдержки. Отвечай кратко, по существу, на русском языке. Если не знаешь ответ — скажи, что переадресуешь оператору." - model_name = "gpt-4o-mini" - role = "support_assistant_chat" - max_retries = 3 - backoff_delays = [1.0, 2.0, 4.0] - - # B2.1: Сохраняем текущее сообщение в историю ПЕРЕД запросом - await save_conversation_history(user_id, "user", user_text) - - # B4: Извлекаем знания из сообщения пользователя (фоново) - safe_background_task( - save_extracted_knowledge(user_id, user_text), - name=f"support_chat_knowledge_{user_id}" - ) - - # B2.1: Собираем контекст с памятью и историей - user_context = await build_user_context(user_id, user_text) - - # B2.1: Формируем messages с контекстом в system prompt - context_enhanced_prompt = system_prompt + "\n\n" + user_context - messages = [ - {"role": "system", "content": context_enhanced_prompt}, - {"role": "user", "content": user_text} - ] - - for attempt in range(max_retries): - start_time = time.perf_counter() - logging.info(f"📡 [GPT Чат Саппорт] Юзер: {user_id} | Попытка: {attempt + 1}/{max_retries}") - try: - async with openai_semaphore: - response = await client_support.chat.completions.create( - model=model_name, - messages=messages, - temperature=0.1, max_tokens=300 - ) - latency_ms = int((time.perf_counter() - start_time) * 1000) - - prompt_tokens = response.usage.prompt_tokens if response.usage else 0 - completion_tokens = response.usage.completion_tokens if response.usage else 0 - - # ИСПРАВЛЕНИЕ БАГА: ai_response определяется ПЕРЕД использованием - ai_response = response.choices[0].message.content - - logging.info(f"✅ [GPT Чат Саппорт Успех] Юзер: {user_id} | Latency: {latency_ms}ms") - - # B2.1: Сохраняем ответ ассистента в историю - await save_conversation_history(user_id, "assistant", ai_response) - - await save_ai_log( - user_id=user_id, model=model_name, role=role, - prompt_tokens=prompt_tokens, completion_tokens=completion_tokens, - latency_ms=latency_ms, status_code=200, success=True, error=None - ) - await message.answer(ai_response, reply_markup=get_support_keyboard()) - - # Извлечение памяти из чата поддержки - safe_background_task(extract_and_save_memory(user_id, user_text, ai_response), name="support_chat_memory") - break - except Exception as e: - latency_ms = int((time.perf_counter() - start_time) * 1000) - status_code = getattr(e, "status_code", 500) - error_msg = str(e) - - logging.error(f"❌ [GPT Чат Саппорт Ошибка] Юзер: {user_id} | Попытка {attempt + 1} провалена: {error_msg}") - - await save_ai_log( - user_id=user_id, model=model_name, role=role, - prompt_tokens=0, completion_tokens=0, - latency_ms=latency_ms, status_code=status_code, success=False, error=error_msg - ) - - if attempt < max_retries - 1: - await asyncio.sleep(backoff_delays[attempt]) - else: - logging.critical(f"🚨 [GPT Чат Саппорт Фатал] Ошибка: {e}") - await message.answer("⚠️ Возникла техническая заминка.", reply_markup=get_support_keyboard()) - except Exception as top_e: - logging.error(f"❌ Исключение в ИИ-саппорте: {top_e}") - await message.answer("⚠️ Возникла техническая заминка.", reply_markup=get_support_keyboard()) - # Очистка памяти ПОСЛЕ выхода из async with — лок гарантированно снят - await user_lock_manager.release(user_id) - - -# <<< END OF NEW BLOCK >>> -# >>> START OF BLOCK 9: AIOHTTP WEBSERVER AND WEBHOOK MAIN LOOP <<< - -@web.middleware -async def debug_logging_middleware(request, handler): - """Debug middleware: логирует ВСЕ входящие запросы.""" - logging.info(f"🔍 [DEBUG] {request.method} {request.path} from {request.remote}") - try: - response = await handler(request) - logging.info(f"🔍 [DEBUG] Response status: {response.status}") - return response - except web.HTTPException as e: - logging.info(f"🔍 [DEBUG] HTTP Exception: {e.status} {e.reason}") - raise - except Exception as e: - logging.exception(f"🔥 [CRITICAL] Unhandled exception in handler: {e}") - raise - -async def health_check(request): - """HF Spaces Health Check.""" - return web.Response(text="🚀 AI Sprint Bot РАБОТАЕТ!") - -# ============================================================ -# НАШ СОБСТВЕННЫЙ ОБРАБОТЧИК WEBHOOK (Без setup_application!) -# ============================================================ -async def custom_webhook_handler(request: web.Request): - """Принимает POST от Telegram через Worker и передает в aiogram.""" - logging.critical("📩 [WEBHOOK] Получен POST запрос на /webhook!") - - try: - # 1. Читаем тело запроса ОДИН раз - body = await request.read() - logging.critical(f"📩 [WEBHOOK] Тело прочитано ({len(body)} байт)") - - # 2. Вызываем диспетчер aiogram напрямую - update = await dp.feed_raw_update(bot, body) - logging.critical(f"📩 [WEBHOOK] Update обработан! Result: {update}") - return web.Response(text="OK") - - except Exception as e: - logging.exception(f"🔥 [FATAL] Ошибка в custom_webhook_handler: {e}") - return web.Response(status=500, text="Internal Server Error") - -# ============================================================ -# СТАРТУП / ШУТДАУН -# ============================================================ -async def on_startup(app: web.Application): - await init_db() - await start_cleanup_scheduler() - - proxy_server = os.getenv("TELEGRAM_API_SERVER") - if proxy_server: - worker_base = proxy_server.strip().rstrip('/') - full_webhook_link = f"{worker_base}{WEBHOOK_PATH}" - - for attempt in range(5): - try: - await bot.set_webhook(url=full_webhook_link, drop_pending_updates=True, request_timeout=30) - logging.critical(f"✅ Webhook установлен ЧЕРЕЗ WORKER: {full_webhook_link}") - break - except Exception as e: - logging.warning(f"⚠️ Попытка {attempt + 1}/5 установки webhook: {e}") - await asyncio.sleep(5) - else: - logging.warning("⚠️ TELEGRAM_API_SERVER не найден — webhook не установлен!") - -async def on_shutdown(app: web.Application): - global _cleanup_task - if _cleanup_task is not None and not _cleanup_task.done(): - logging.info("🛑 [Cleanup Scheduler] Отмена фоновой задачи очистки...") - _cleanup_task.cancel() - try: - await _cleanup_task - except asyncio.CancelledError: - logging.info("✅ [Cleanup Scheduler] Фоновая задача успешно отменена.") - except Exception as e: - logging.warning(f"⚠️ [Cleanup Scheduler] Ошибка при отмене задачи: {e}") - _cleanup_task = None - - try: - await dp.storage.close() - await bot.session.close() - await client.close() - await client_support.close() - if db_pool: - await db_pool.close() - except Exception as e: - logging.warning(f"Ошибка при shutdown: {e}") - -def main(): - logging.info("="*50) - logging.info("🟢 PROCESS STARTED. Python is running.") - logging.info("="*50) - - from aiohttp import web - - port = int(os.environ.get("PORT", 7860)) - app = web.Application(middlewares=[debug_logging_middleware]) - # ============================================================ - # MIDDLEWARE: Логирование всех входящих HTTP-запросов - # ============================================================ - @aiohttp.web.middleware - async def logging_middleware(request, handler): - logger.warning( - f"🌐 [HTTP] {request.method} {request.path} " - f"от {request.remote} | Headers: {dict(request.headers)}" - ) - try: - response = await handler(request) - logger.warning( - f"🌐 [HTTP] {request.method} {request.path} → {response.status}" - ) - return response - except Exception as e: - logger.error(f"🌐 [HTTP] {request.method} {request.path} → ERROR: {e}") - raise - - app.middlewares.append(logging_middleware) - - # 1. Health Check (Корень для HuggingFace Spaces) - app.router.add_get('/', health_check, allow_head=False) - - # 2. Стандартный и надежный обработчик aiogram - webhook_requests_handler = SimpleRequestHandler(dispatcher=dp, bot=bot) - webhook_requests_handler.register(app, path=WEBHOOK_PATH) - setup_application(app, dp, bot=bot) - - # ============================================================ - # DIAGNOSTIC: /health — расширенная проверка доступности - # ============================================================ - async def health_check_handler(request: aiohttp.web.Request) -> aiohttp.web.Response: - import datetime as _dt - logger.info("🔍 [HEALTH] GET /health — внешний запрос проверки") - return aiohttp.web.json_response({ - "status": "alive", - "service": "freshpixels-ai-sprint-bot", - "webhook_set": True, - "timestamp": _dt.datetime.utcnow().isoformat(), - }) - - app.router.add_get('/health', health_check_handler, allow_head=False) - - # 3. Запускаем бизнес-логику - app.on_startup.append(on_startup) - app.on_shutdown.append(on_shutdown) - - logging.info(f"🚀 Запуск сервера на {port}...") - - try: - web.run_app(app, host='0.0.0.0', port=port) - except Exception as e: - logging.exception(f"🔥 [FATAL] Server crashed: {e}") - sys.exit(1) - -# <<< END OF BLOCK 9 - -if __name__ == "__main__": - main() - \ No newline at end of file