import asyncio import logging import time from aiogram import Router, types, F from aiogram.exceptions import TelegramAPIError from aiogram.types import InlineKeyboardMarkup, InlineKeyboardButton from aiogram.enums import ParseMode from config import config from database import db from services.glm import glm_service, LLMServiceError from utils.helpers import send_long_message, edit_or_send_message logger = logging.getLogger(__name__) router = Router() # In-memory cache for last user messages (for regeneration) _last_user_messages: dict[int, str] = {} SYSTEM_PROMPT = { "role": "system", "content": ( "Ты — AI ассистент на базе GLM-5.1, работающий через NVIDIA NIM API. " "Ты помогаешь пользователю как Senior Staff Python Engineer и Acting CTO.\n\n" "Правила работы:\n" "1. Отвечай на русском языке, если не указано иное.\n" "2. Код пиши production-ready, Python 3.11+, async-safe.\n" "3. Используй PATCH MODE для изменений >200 строк.\n" "4. Не ломай существующую архитектуру.\n" "5. Минимальные изменения — локальные патчи.\n" "6. Не используй TODO, FIXME, PLACEHOLDER, PASS.\n" "7. Проверяй безопасность: SQL Injection, Race Conditions, Async Safety.\n" "8. Если не хватает контекста — запроси дополнительный код.\n" "9. Не придумывай существующие методы/классы/API.\n" "10. Перед кодом: анализ, архитектура, точка интеграции, зависимости." ), } def get_chat_buttons() -> InlineKeyboardMarkup: """Кнопки под ответом бота.""" return InlineKeyboardMarkup(inline_keyboard=[ [ InlineKeyboardButton(text="🔄 Регенерировать", callback_data="regenerate"), InlineKeyboardButton(text="🗑 Очистить", callback_data="clear_memory"), ], ]) async def _keep_typing(message: types.Message, stop_event: asyncio.Event) -> None: """Периодически отправляет typing action, пока stop_event не установлен.""" while not stop_event.is_set(): try: await message.bot.send_chat_action(message.chat.id, "typing") except TelegramAPIError as e: logger.debug("Typing action failed: %s", e) try: await asyncio.wait_for(stop_event.wait(), timeout=config.TYPING_ACTION_INTERVAL) except asyncio.TimeoutError: continue async def _maybe_summarize(user_id: int) -> None: """Фоновая суммаризация старой истории.""" try: count = await db.count_unsummarized(user_id) threshold = config.SUMMARIZE_THRESHOLD keep_recent = config.KEEP_RECENT_MESSAGES if count <= threshold: return to_summarize = count - keep_recent if to_summarize < 5: return old_messages = await db.get_oldest_unsummarized(user_id, to_summarize) if len(old_messages) < 5: return dialog_text = "\n".join([ f"{m['role']}: {m['content']}" for m in old_messages ]) summary = await glm_service.summarize(dialog_text) if not summary: return cutoff_id = old_messages[-1]["id"] await db.save_summary(user_id, summary, len(old_messages)) await db.mark_summarized(user_id, cutoff_id) logger.info( "Summarized %d messages for user %s, cutoff_id=%s", len(old_messages), user_id, cutoff_id ) except Exception as e: logger.error("Summarization failed for user %s: %s", user_id, e, exc_info=True) async def _build_messages(user_id: int, user_text: str) -> list[dict]: """Строит список сообщений для LLM с учётом бюджета токенов.""" messages: list[dict] = [SYSTEM_PROMPT] # Добавляем summary как контекст summary = await db.get_summary(user_id) if summary: messages.append({ "role": "system", "content": f"[Контекст предыдущих диалогов: {summary}]" }) # Получаем историю с учётом токенового бюджета history = await db.get_messages_with_token_budget( user_id, max_tokens=config.MAX_CONTEXT_TOKENS ) # Ограничиваем количество сообщений if len(history) > config.MAX_HISTORY: history = history[-config.MAX_HISTORY:] for h in history: messages.append({"role": h["role"], "content": h["content"]}) # Добавляем текущее сообщение (если его ещё нет в истории) if not history or history[-1]["content"] != user_text: messages.append({"role": "user", "content": user_text}) return messages @router.message(F.text) async def handle_chat(message: types.Message) -> None: """Основной обработчик текстовых сообщений.""" logger.info("→ handle_chat START, user_id=%s, text_len=%d", message.from_user.id, len(message.text or "")) if not message.text: await message.answer("❌ Поддерживаются только текстовые сообщения.") return user_id = message.from_user.id user_text = message.text _last_user_messages[user_id] = user_text # Upsert user await db.upsert_user( user_id, message.from_user.username, message.from_user.first_name, message.from_user.last_name, ) # Save user message tokens_estimate = glm_service.estimate_tokens(user_text) await db.save_message(user_id, "user", user_text, tokens_estimate) # Build messages for LLM messages = await _build_messages(user_id, user_text) # Typing action stop_typing = asyncio.Event() typing_task = asyncio.create_task(_keep_typing(message, stop_typing)) start_time = time.perf_counter() bot_message = None try: logger.info("→ Calling GLM, messages_count=%d", len(messages)) if config.STREAMING_ENABLED: # Streaming mode — обновляем сообщение по мере поступления чанков response_text = "" chunk_buffer = "" last_update = time.perf_counter() message_sent = False async for chunk in glm_service.chat_stream(messages, user_id=user_id): response_text += chunk chunk_buffer += chunk # Обновляем сообщение не чаще чем раз в N секунд now = time.perf_counter() if now - last_update >= config.STREAMING_UPDATE_INTERVAL and len(response_text) > 10: if not message_sent: bot_message = await message.answer( response_text + "▌", parse_mode=ParseMode.MARKDOWN, ) message_sent = True else: try: await bot_message.edit_text( response_text + "▌", parse_mode=ParseMode.MARKDOWN, ) except TelegramAPIError: pass # Message not modified или другая ошибка last_update = now chunk_buffer = "" # Финальное обновление stop_typing.set() await typing_task if message_sent and bot_message: try: await bot_message.edit_text( response_text, parse_mode=ParseMode.MARKDOWN, reply_markup=get_chat_buttons(), ) except TelegramAPIError: # Если edit не удался, отправляем новое await send_long_message( message, response_text, reply_markup=get_chat_buttons(), ) else: await send_long_message( message, response_text, reply_markup=get_chat_buttons(), ) else: # Non-streaming mode response = await glm_service.chat(messages, user_id=user_id) stop_typing.set() await typing_task await send_long_message(message, response, reply_markup=get_chat_buttons()) response_text = response duration_ms = (time.perf_counter() - start_time) * 1000 logger.info("← Response sent, duration=%.1fms, length=%d", duration_ms, len(response_text)) # Save assistant message assistant_tokens = glm_service.estimate_tokens(response_text) await db.save_message(user_id, "assistant", response_text, assistant_tokens) # Trigger background summarization asyncio.create_task(_maybe_summarize(user_id)) except LLMServiceError as e: stop_typing.set() await typing_task logger.error("LLM Service Error: %s", e) await message.answer( "❌ Ошибка модели:\n" f"{str(e)[:200]}\n\n" "Попробуйте повторить запрос позже.", parse_mode="HTML", ) except Exception as e: stop_typing.set() await typing_task logger.error("❌ Unexpected error in chat handler: %s", e, exc_info=True) await message.answer( "❌ Неожиданная ошибка.\n" "Разработчик уведомлён. Попробуйте позже.", parse_mode="HTML", )