SkyAlone / handlers /chat.py
FreshPixels's picture
Rename handlers chat.py to handlers/chat.py
dd596d8 verified
Raw
History Blame Contribute Delete
10.4 kB
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(
"❌ <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",
)