Spaces:
Runtime error
Runtime error
| 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 | |
| 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( | |
| "❌ <b>Ошибка модели:</b>\n" | |
| f"<code>{str(e)[:200]}</code>\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( | |
| "❌ <b>Неожиданная ошибка.</b>\n" | |
| "Разработчик уведомлён. Попробуйте позже.", | |
| parse_mode="HTML", | |
| ) | |