685 lines
36 KiB
Python
685 lines
36 KiB
Python
# /opt/trueconf_bot/search_bot/handlers.py
|
||
import re
|
||
import os
|
||
import sys
|
||
import csv
|
||
import time
|
||
import math
|
||
import httpx
|
||
import smtplib
|
||
import asyncio
|
||
import torch
|
||
from email.mime.text import MIMEText
|
||
from email.mime.multipart import MIMEMultipart
|
||
from datetime import datetime
|
||
|
||
from trueconf import Router, Message
|
||
from sentence_transformers import CrossEncoder
|
||
from qdrant_client import QdrantClient
|
||
from qdrant_client.models import Filter, FieldCondition, MatchValue, Range
|
||
|
||
from config.config import *
|
||
from utils import states
|
||
from utils.menu import MENU_TEXT
|
||
from utils.texts import SEARCH_MAIN_MENU_TEXT, SEARCH_UNKNOWN_CMD_TEXT, search_searching_text, search_not_found_text, search_feedback_thanks_text, search_feedback_info_text, search_unavailable_action
|
||
|
||
# 🔌 Импортируем централизованную функцию сбора статистики из main
|
||
from utils.stats_logger import log_menu_stats
|
||
|
||
# =========================================================
|
||
# ГЛОБАЛЬНЫЙ МАРШРУТИЗАТОР ЦЕЛЬНЫХ VK-ЭМОДЗИ
|
||
# =========================================================
|
||
EMOJI_DIGITS = {
|
||
"0": "0⃣",
|
||
"1": "1⃣",
|
||
"2": "2⃣",
|
||
"3": "3⃣",
|
||
"4": "4⃣",
|
||
"5": "5⃣",
|
||
"6": "6⃣",
|
||
"7": "7⃣",
|
||
"8": "8⃣",
|
||
"9": "9⃣",
|
||
"10": "🔟"
|
||
}
|
||
|
||
# Глобальная константа для меню Поискового Бота
|
||
|
||
# Универсальный стандарт сообщения об ошибке ввода
|
||
|
||
# Проверка флага глубокого профилирования
|
||
ENABLE_DEEP_PROFILING = getattr(sys.modules['config.config'], 'ENABLE_DEEP_PROFILING', False) if 'config.config' in sys.modules else True
|
||
|
||
# =========================================================
|
||
# INIT
|
||
# =========================================================
|
||
router = Router()
|
||
|
||
# Подключаем Qdrant динамически из живых настроек config.py
|
||
live_config = sys.modules.get('config.config')
|
||
client = QdrantClient(
|
||
host=getattr(live_config, 'QDRANT_HOST', 'localhost'),
|
||
port=getattr(live_config, 'QDRANT_PORT', 6333)
|
||
)
|
||
|
||
# Универсальная гибридная инициализация локального реранкера (если выбран LOCAL)
|
||
reranker = None
|
||
if getattr(live_config, 'ENABLE_RERANKER', False) and getattr(live_config, 'RERANKER_MODE', '') == "LOCAL":
|
||
try:
|
||
device = "cuda" if torch.cuda.is_available() else "cpu"
|
||
model_path = getattr(live_config, 'RERANK_MODEL_PATH', '/opt/models/bge-reranker')
|
||
reranker = CrossEncoder(model_path, max_length=512, device=device)
|
||
print(f"📦 Реранкинг: ЛОКАЛЬНЫЙ режим на {device.upper()}")
|
||
except Exception as e:
|
||
print(f"⚠️ Ошибка загрузки локальной модели реранкера: {e}")
|
||
|
||
# =========================================================
|
||
# EMAIL ALERT FUNCTION
|
||
# =========================================================
|
||
last_email_time = 0.0
|
||
|
||
def send_alert_email(subject: str, body: str, is_critical_error: bool = False):
|
||
global last_email_time
|
||
current_config = sys.modules.get('config.config')
|
||
if not getattr(current_config, 'ALERTS_SUPPORT_EMAILS', []):
|
||
return
|
||
|
||
current_time = time.time()
|
||
if is_critical_error and (current_time - last_email_time < EMAIL_COOLDOWN_SEC):
|
||
return
|
||
|
||
if is_critical_error:
|
||
last_email_time = current_time
|
||
|
||
msg = MIMEMultipart()
|
||
msg["From"] = getattr(current_config, 'DEFAULT_EMAIL_FROM', 'ai@sibcem.ru')
|
||
msg["To"] = ", ".join(getattr(current_config, 'ALERTS_SUPPORT_EMAILS', []))
|
||
msg["Subject"] = subject
|
||
msg.attach(MIMEText(body, "plain", "utf-8"))
|
||
|
||
try:
|
||
with smtplib.SMTP(SMTP_SERVER, SMTP_PORT) as server:
|
||
server.send_message(msg)
|
||
except Exception as e:
|
||
print(f"⚠️ Ошибка отправки email-алерта: {e}")
|
||
|
||
# =========================================================
|
||
# SYSTEM ERRORS & LOGGING
|
||
# =========================================================
|
||
def log_system_error(action: str, error_details: str):
|
||
os.makedirs(LOGS_PATH, exist_ok=True)
|
||
log_file = os.path.join(LOGS_PATH, "search_errors.log")
|
||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
with open(log_file, "a", encoding="utf-8") as f:
|
||
f.write(f"[{timestamp}] Действие: {action} | Ошибка: {error_details}\n")
|
||
|
||
async def handle_search_system_error(msg: Message, user_id: str, action: str, error_details: str):
|
||
await asyncio.to_thread(log_system_error, action, error_details)
|
||
email_body = f"Пользователь: {user_id}\nДействие: {action}\nДетали проблемы:\n{error_details}"
|
||
await asyncio.to_thread(send_alert_email, f"🚨 Критический сбой: Поиск Бот - {action}", email_body, True)
|
||
|
||
await msg.answer(search_unavailable_action(EMOJI_DIGITS), parse_mode="html")
|
||
|
||
def log_analytics(login: str, query: str, answer: str, is_success: bool, response_time: float, qdrant_debug: str, sources: set):
|
||
os.makedirs(LOGS_PATH, exist_ok=True)
|
||
csv_file = os.path.join(LOGS_PATH, "chat_history.csv")
|
||
text_log = os.path.join(LOGS_PATH, "chat_debug.log")
|
||
timestamp = datetime.now().strftime("%Y-%m-%d %H:%M:%S")
|
||
sources_str = ", ".join(sorted(sources)) if sources else "Нет"
|
||
file_exists = os.path.isfile(csv_file)
|
||
|
||
time_str = f"{response_time:.2f}".replace(".", ",")
|
||
|
||
with open(csv_file, mode="a", encoding="utf-8-sig", newline="") as f:
|
||
writer = csv.writer(f, delimiter=";")
|
||
if not file_exists:
|
||
writer.writerow(["Дата и Время", "Логин", "Вопрос", "Ответ", "Успешно?", "Время ответа (сек)", "Документы", "Ответ устроил?", "Текст отзыва"])
|
||
|
||
writer.writerow([
|
||
timestamp, login, query.replace("\n", " ").strip(),
|
||
answer.replace("\n", " ").strip(), "Да" if is_success else "Нет",
|
||
time_str, sources_str, "Да", ""
|
||
])
|
||
|
||
if not getattr(sys.modules['config.config'], 'ENABLE_QDRANT_DEBUG', False):
|
||
return
|
||
|
||
with open(text_log, mode="a", encoding="utf-8") as f:
|
||
f.write(f"╔{'═'*78}╗\n")
|
||
f.write(f"║ ВРЕМЯ: {timestamp}\n")
|
||
f.write(f"║ ПОЛЬЗОВАТЕЛЬ: {login}\n")
|
||
f.write(f"║ ВОПРОС: {query}\n")
|
||
f.write(f"╔{'═'*78}╗\n")
|
||
f.write("║ ВЫВОД QDRANT / RERANKER:\n")
|
||
f.write(f"{qdrant_debug}\n")
|
||
f.write(f"╠{'═'*78}╣\n")
|
||
f.write("║ ПОЛНЫЙ ОТВЕТ ИИ / ОШИБКА:\n")
|
||
f.write(f"║ {answer.replace(chr(10), chr(10) + '║ ')}\n")
|
||
f.write(f"╠{'═'*78}╣\n")
|
||
f.write(f"║ ИСТОЧНИКИ: {sources_str}\n")
|
||
f.write(f"║ ВРЕМЯ ОБРАБОТКИ: {response_time:.2f} сек.\n")
|
||
f.write(f"╚{'═'*78}╝\n\n")
|
||
|
||
def update_feedback_log(login: str, rating: str = None, feedback_text: str = None):
|
||
csv_file = os.path.join(LOGS_PATH, "chat_history.csv")
|
||
if not os.path.isfile(csv_file):
|
||
return
|
||
try:
|
||
with open(csv_file, mode="r", encoding="utf-8-sig") as f:
|
||
reader = list(csv.reader(f, delimiter=";"))
|
||
target_idx = -1
|
||
for i in range(len(reader) - 1, 0, -1):
|
||
if reader[i][1] == login:
|
||
target_idx = i
|
||
break
|
||
if target_idx != -1:
|
||
while len(reader[target_idx]) < 9:
|
||
reader[target_idx].append("")
|
||
if rating is not None:
|
||
reader[target_idx][7] = rating
|
||
if feedback_text is not None:
|
||
reader[target_idx][8] = feedback_text.replace("\n", " ").strip()
|
||
with open(csv_file, mode="w", encoding="utf-8-sig", newline="") as f:
|
||
writer = csv.writer(f, delimiter=";")
|
||
writer.writerows(reader)
|
||
except Exception as e:
|
||
print(f"⚠️ Ошибка обновления отзыва: {e}")
|
||
|
||
def get_last_query_and_answer(login: str):
|
||
csv_file = os.path.join(LOGS_PATH, "chat_history.csv")
|
||
if not os.path.isfile(csv_file):
|
||
return "Неизвестно", "Неизвестно"
|
||
try:
|
||
with open(csv_file, mode="r", encoding="utf-8-sig") as f:
|
||
reader = list(csv.reader(f, delimiter=";"))
|
||
for i in range(len(reader) - 1, 0, -1):
|
||
if reader[i][1] == login:
|
||
query = reader[i][2] if len(reader[i]) > 2 else "Неизвестно"
|
||
answer = reader[i][3] if len(reader[i]) > 3 else "Неизвестно"
|
||
return query, answer
|
||
except Exception as e:
|
||
pass
|
||
return "Неизвестно", "Неизвестно"
|
||
|
||
# =========================================================
|
||
# HELPERS
|
||
# =========================================================
|
||
def truncate_context(text: str) -> str:
|
||
if len(text) <= MAX_SINGLE_CHUNK_LENGTH:
|
||
return text
|
||
return text[:MAX_SINGLE_CHUNK_LENGTH] + "\n...[обрезано]"
|
||
|
||
def normalize_text(text: str) -> str:
|
||
return re.sub(r"\s+", " ", text.lower()).strip()
|
||
|
||
def is_table_chunk(text: str) -> bool:
|
||
table_patterns = ["Контекст:", "Параметр:", "Значение:", "Col", "|"]
|
||
hits = sum(1 for p in table_patterns if p in text)
|
||
return hits >= 2
|
||
|
||
# =========================================================
|
||
# AI REQUEST
|
||
# =========================================================
|
||
async def ask_ai(query: str, context_texts: list) -> str:
|
||
刻_texts = "\n\n".join(context_texts)
|
||
|
||
current_config = sys.modules.get('config.config')
|
||
live_ollama_url = getattr(current_config, 'OLLAMA_URL', 'http://172.16.80.250:8080/v1/chat/completions')
|
||
live_ollama_model = getattr(current_config, 'OLLAMA_MODEL', 'gemma-4-E4B-it-Q4_K_M')
|
||
live_ollama_temp = getattr(current_config, 'OLLAMA_TEMPERATURE', 0.0)
|
||
live_ollama_predict = getattr(current_config, 'OLLAMA_NUM_PREDICT', 1500)
|
||
live_ai_timeout = getattr(current_config, 'AI_HTTP_TIMEOUT', 180.0)
|
||
|
||
full_prompt = f"""Ты — строгий и точный корпоративный аналитик документов АО «ХК «Сибцем».
|
||
Твоя единственная задача — отвечать на вопросы пользователя, используя ИСКЛЮЧИТЕЛЬНО информацию из предоставленного контекста.
|
||
|
||
<context>
|
||
{刻_texts}
|
||
</context>
|
||
|
||
ИНСТРУКЦИИ И ПРАВИЛА:
|
||
1. Опирайся строго на текст внутри тегов <context>. Категорически запрещено использовать внешние знания или додумывать факты.
|
||
2. Приоритет данных: Если в тексте есть таблицы, списки или структурированные блоки, в первую очередь бери сроки и параметры оттуда.
|
||
3. Внимательность: Всегда дочитывай абзац до конца, чтобы не пропустить исключения.
|
||
4. Косвенные ответы: Если в тексте прямо не описан процесс (например, "Как создать?"), но указано, на основании чего он выполняется (например, "создается по заявке"), ОБЯЗАТЕЛЬНО сообщи об этом пользователю.
|
||
5. Формат ответа: Пиши кратко, ясно и только по делу на русском языке.
|
||
6. ЕСЛИ ОТВЕТ НАЙДЕН: В конце ответа всегда с новой строки указывай источник в формате: "Источник: [название файла]".
|
||
7. ЕСЛИ ОТВЕТА НЕТ: Если информации в тегах <context> недостаточно даже для косвенного ответа, выведи строго одну фразу: «В данных документах этот вопрос не регламентирован.». Источник указывать не нужно.
|
||
|
||
ВОПРОС ПОЛЬЗОВАТЕЛЯ:
|
||
{query}"""
|
||
|
||
safe_temp = 0.1 if live_ollama_temp <= 0.0 else live_ollama_temp
|
||
|
||
payload = {
|
||
"model": live_ollama_model,
|
||
"messages": [
|
||
{"role": "user", "content": full_prompt}
|
||
],
|
||
"temperature": safe_temp,
|
||
"max_tokens": live_ollama_predict,
|
||
"stream": False
|
||
}
|
||
|
||
async with httpx.AsyncClient(timeout=httpx.Timeout(live_ai_timeout)) as http_client:
|
||
response = await http_client.post(live_ollama_url, json=payload)
|
||
response.raise_for_status()
|
||
data = response.json()
|
||
return data["choices"][0]["message"]["content"].strip()
|
||
|
||
# =========================================================
|
||
# MAIN MESSAGE HANDLER
|
||
# =========================================================
|
||
@router.message()
|
||
async def main_message_handler(msg: Message):
|
||
user_id = msg.from_user.id
|
||
current_state = states.get_state(user_id)
|
||
|
||
current_config = sys.modules.get('config.config')
|
||
live_use_reranker = getattr(current_config, 'ENABLE_RERANKER', False)
|
||
live_reranker_mode = getattr(current_config, 'RERANKER_MODE', 'REMOTE_LLAMACPP')
|
||
live_reranker_url = getattr(current_config, 'RERANKER_API_URL', 'http://172.16.80.250:8080/v1/rerank')
|
||
live_reranker_model = getattr(current_config, 'RERANKER_MODEL_NAME', 'bge-reranker-v2-m3-Q8_0')
|
||
|
||
msg_text = msg.text.strip() if msg.text else ""
|
||
cmd = msg_text.lower()
|
||
|
||
# 🔄 НОРМАЛИЗАЦИЯ КНОПОК ВК-ЭМОДЗИ
|
||
for raw_num, emoji_num in EMOJI_DIGITS.items():
|
||
if cmd == emoji_num:
|
||
cmd = raw_num
|
||
break
|
||
|
||
# -----------------------------------------------------
|
||
# ЛОГИКА ФИДБЕКА
|
||
# -----------------------------------------------------
|
||
if current_state in ["AWAITING_FEEDBACK", "FEEDBACK_DONE"]:
|
||
msg.handled = True
|
||
if cmd in ["0", "9", "/start", "меню"]:
|
||
log_menu_stats(user_id, "Search Bot", "Выход в главное меню")
|
||
states.clear_state(user_id)
|
||
await msg.answer(MENU_TEXT, parse_mode="html")
|
||
return
|
||
elif cmd == "1":
|
||
log_menu_stats(user_id, "Search Bot", "Повторный ввод вопроса")
|
||
states.set_state(user_id, "SEARCH_MODE")
|
||
await msg.answer(f"🔍 <b>Жду ваш следующий вопрос по регламентам!</b>\n\n<i>Для выхода отправьте {EMOJI_DIGITS['0']}.</i>", parse_mode="html")
|
||
return
|
||
elif current_state == "AWAITING_FEEDBACK" and msg_text:
|
||
log_menu_stats(user_id, "Search Bot", "Отправлен развернутый отзыв")
|
||
update_feedback_log(str(user_id), feedback_text=msg_text)
|
||
states.set_state(user_id, "SEARCH_MODE")
|
||
last_query, last_answer = get_last_query_and_answer(str(user_id))
|
||
send_alert_email(
|
||
subject="💬 Новый текстовый отзыв к боту-регламентов",
|
||
body=(f"Пользователь: {user_id}\n\nВОПРОС ПОЛЬЗОВАТЕЛЯ:\n{last_query}\n\nОТВЕТ БОТА:\n{last_answer}\n\nТЕКСТ ОТЗЫВА:\n{msg_text}")
|
||
)
|
||
await msg.answer(search_feedback_thanks_text(EMOJI_DIGITS), parse_mode="html")
|
||
return
|
||
else:
|
||
await msg.answer(UNKNOWN_CMD_TEXT, parse_mode="html")
|
||
return
|
||
|
||
# -----------------------------------------------------
|
||
# ЛОГИКА ПОИСКА (SEARCH_MODE)
|
||
# -----------------------------------------------------
|
||
if current_state == "SEARCH_MODE":
|
||
msg.handled = True
|
||
|
||
if cmd in ["0", "9", "/start", "меню"]:
|
||
log_menu_stats(user_id, "Search Bot", "Выход в главное меню")
|
||
states.clear_state(user_id)
|
||
await msg.answer(MENU_TEXT, parse_mode="html")
|
||
return
|
||
|
||
if cmd == "1":
|
||
log_menu_stats(user_id, "Search Bot", "Клик: Ответ не удовлетворил (1)")
|
||
states.set_state(user_id, "AWAITING_FEEDBACK")
|
||
update_feedback_log(str(user_id), rating="Нет")
|
||
last_query, last_answer = get_last_query_and_answer(str(user_id))
|
||
send_alert_email(
|
||
subject="⚠️ Пользователя не устроил answer бота",
|
||
body=(f"Пользователь: {user_id}\n\nВОПРОС ПОЛЬЗОВАТЕЛЯ:\n{last_query}\n\nОТВЕТ БОТА:\n{last_answer}\n\nНажал '1' (Ответ не удовлетворил).")
|
||
)
|
||
await msg.answer(search_feedback_info_text(EMOJI_DIGITS), parse_mode="html")
|
||
return
|
||
|
||
if not msg_text:
|
||
await msg.answer(UNKNOWN_CMD_TEXT, parse_mode="html")
|
||
return
|
||
|
||
# ФИКСАЦИЯ СТАТИСТИКИ: Начало семантического поиска
|
||
log_menu_stats(user_id, "Search Bot", "Выполнение поиска по регламентам")
|
||
|
||
# =====================================================
|
||
# ГЛУБОКОЕ ПРОФИЛИРОВАНИЕ (DEEP PROFILING)
|
||
# =====================================================
|
||
start_time = time.time()
|
||
prof_log = []
|
||
|
||
def add_prof(msg_text):
|
||
if ENABLE_DEEP_PROFILING:
|
||
elapsed = time.time() - start_time
|
||
prof_log.append(f"[{elapsed:>6.2f}s] {msg_text}")
|
||
|
||
add_prof(f"=== СТАРТ ЗАПРОСА: '{msg_text}' ===")
|
||
add_prof(f"[НАСТРОЙКИ] QDRANT_LIMIT={QDRANT_LIMIT}, RERANK_TOP_K={RERANK_TOP_K}, USE_RERANKER={live_use_reranker}, CHUNK_WINDOW={QDRANT_CHUNK_WINDOW}")
|
||
|
||
await msg.answer(search_searching_text(), parse_mode="html")
|
||
|
||
# =====================================================
|
||
# ВЕКТОРИЗАЦИЯ И ПОИСК (API ЭМБЕДДИНГОВ)
|
||
# =====================================================
|
||
try:
|
||
t_vec = time.time()
|
||
|
||
payload_api = {
|
||
"model": getattr(current_config, 'EMBEDDING_MODEL_NAME', 'bge-m3'),
|
||
"input_type": "query",
|
||
"texts": [msg_text]
|
||
}
|
||
headers_api = {"Content-Type": "application/json"}
|
||
api_token = getattr(current_config, 'EMBEDDING_API_TOKEN', None)
|
||
if api_token:
|
||
headers_api["X-Embedding-Token"] = api_token
|
||
|
||
async with httpx.AsyncClient(timeout=httpx.Timeout(API_HTTP_TIMEOUT)) as http_client:
|
||
api_resp = await http_client.post(EMBEDDING_API_URL, json=payload_api, headers=headers_api)
|
||
api_resp.raise_for_status()
|
||
query_vector = api_resp.json()["embeddings"][0]
|
||
|
||
add_prof(f"[ВЕКТОРИЗАЦИЯ API] Текст преобразован в вектор через API за {time.time() - t_vec:.2f}s.")
|
||
|
||
fetch_limit = QDRANT_LIMIT * 3 if live_use_reranker else QDRANT_LIMIT
|
||
|
||
t_qdr = time.time()
|
||
search_results = client.query_points(
|
||
collection_name=QDRANT_COLLECTION,
|
||
query=query_vector,
|
||
limit=fetch_limit,
|
||
with_payload=True
|
||
).points
|
||
add_prof(f"[QDRANT ПОИСК] Извлечено чанков: {len(search_results)} (запрашивали {fetch_limit}). Заняло {time.time() - t_qdr:.2f}s.")
|
||
|
||
# --- УНИВЕРСАЛЬНЫЙ ГИБРИДНЫЙ БЛОК РЕРАНКИНГА ---
|
||
if live_use_reranker and search_results:
|
||
t_rerank = time.time()
|
||
logits = None
|
||
rerank_success = False
|
||
|
||
# ---------------------------------------------------------
|
||
# ВАРИАНТ 1: Новый API Llama.cpp (Standard Jina API format) с БАТЧИНГОМ
|
||
# ---------------------------------------------------------
|
||
if live_reranker_mode == "REMOTE_LLAMACPP":
|
||
add_prof(f"[RERANKER] Вызов Llama.cpp API для {len(search_results)} чанков (Батчинг).")
|
||
|
||
documents = [res.payload.get("page_content", "")[:1500] for res in search_results]
|
||
batch_size = 4
|
||
logits = [0.0] * len(documents)
|
||
rerank_success = True
|
||
|
||
try:
|
||
async with httpx.AsyncClient(timeout=httpx.Timeout(API_HTTP_TIMEOUT)) as http_client:
|
||
for b_start in range(0, len(documents), batch_size):
|
||
b_docs = documents[b_start:b_start + batch_size]
|
||
payload = {
|
||
"model": live_reranker_model,
|
||
"query": msg_text,
|
||
"documents": b_docs,
|
||
"top_n": len(b_docs)
|
||
}
|
||
|
||
response = await http_client.post(live_reranker_url, json=payload)
|
||
response.raise_for_status()
|
||
rerank_data = response.json()
|
||
|
||
if "results" in rerank_data:
|
||
for item in rerank_data["results"]:
|
||
local_index = item.get("index")
|
||
score = item.get("relevance_score", 0.0)
|
||
if local_index is not None:
|
||
global_index = b_start + local_index
|
||
if 0 <= global_index < len(logits):
|
||
logits[global_index] = score
|
||
else:
|
||
add_prof(f"[RERANKER ERROR] Неверный формат ответа в батче {b_start}: {rerank_data}")
|
||
rerank_success = False
|
||
break
|
||
except Exception as e:
|
||
add_prof(f"[RERANKER CRASH] Ошибка Llama.cpp API (батч {b_start}): {e}")
|
||
rerank_success = False
|
||
|
||
# ---------------------------------------------------------
|
||
# ВАРИАНТ 2: Старый самописный удаленный API (Legacy)
|
||
# ---------------------------------------------------------
|
||
elif live_reranker_mode == "REMOTE_CUSTOM":
|
||
add_prof(f"[RERANKER] Вызов старого кастомного API для {len(search_results)} чанков.")
|
||
passages = [res.payload.get("page_content", "") for res in search_results]
|
||
payload = {
|
||
"query": msg_text,
|
||
"passages": passages
|
||
}
|
||
try:
|
||
async with httpx.AsyncClient(timeout=httpx.Timeout(API_HTTP_TIMEOUT)) as http_client:
|
||
response = await http_client.post(live_reranker_url, json=payload)
|
||
response.raise_for_status()
|
||
rerank_data = response.json()
|
||
|
||
if rerank_data.get("status") == "success":
|
||
logits = rerank_data.get("scores", [])
|
||
rerank_success = True
|
||
else:
|
||
add_prof(f"[RERANKER ERROR] Сервер вернул ошибку: {rerank_data.get('detail')}")
|
||
except Exception as e:
|
||
add_prof(f"[RERANKER CRASH] Ошибка кастомного API: {e}")
|
||
|
||
# ---------------------------------------------------------
|
||
# ВАРИАНТ 3: Локальный инференс (CrossEncoder)
|
||
# ---------------------------------------------------------
|
||
elif live_reranker_mode == "LOCAL" and reranker is not None:
|
||
add_prof(f"[RERANKER] Локальный инференс для {len(search_results)} пар.")
|
||
pairs = [[msg_text, res.payload.get("page_content", "")] for res in search_results]
|
||
try:
|
||
logits = reranker.predict(pairs)
|
||
rerank_success = True
|
||
except Exception as e:
|
||
add_prof(f"[RERANKER CRASH] Ошибка локального инференса: {e}")
|
||
|
||
# ---------------------------------------------------------
|
||
# ПРИМЕНЕНИЕ РЕЗУЛЬТАТОВ
|
||
# ---------------------------------------------------------
|
||
if rerank_success and logits is not None:
|
||
for i, res in enumerate(search_results):
|
||
res.payload["original_qdrant_score"] = res.score
|
||
res.score = float(logits[i])
|
||
|
||
search_results = sorted(search_results, key=lambda r: r.score, reverse=True)
|
||
search_results = search_results[:RERANK_TOP_K]
|
||
add_prof(f"[RERANKER END] Успешно отобрано {len(search_results)} чанков за {time.time() - t_rerank:.2f}s.")
|
||
else:
|
||
add_prof("[RERANKER FALLBACK] Ошибка реранкинга. Падаем на базовые скоры Qdrant.")
|
||
search_results = sorted(search_results, key=lambda r: r.score, reverse=True)
|
||
search_results = search_results[:RERANK_TOP_K]
|
||
else:
|
||
search_results = sorted(search_results, key=lambda r: r.score, reverse=True)
|
||
|
||
except Exception as e:
|
||
error_str = f"[ОШИБКА ПОИСКА]: {type(e).__name__} - {str(e)}"
|
||
add_prof(f"[КРИТИЧЕСКАЯ ОШИБКА] {error_str}")
|
||
log_analytics(str(user_id), msg_text, error_str, False, time.time() - start_time, error_str, set())
|
||
await handle_search_system_error(msg, str(user_id), "Семантический поиск", error_str)
|
||
return
|
||
|
||
if not search_results:
|
||
error_str = "Поиск вернул 0 результатов."
|
||
add_prof(error_str)
|
||
log_analytics(str(user_id), msg_text, error_str, False, time.time() - start_time, error_str, set())
|
||
await handle_search_system_error(msg, str(user_id), "Проверка результатов Qdrant", error_str)
|
||
return
|
||
|
||
# =====================================================
|
||
# ФИЛЬТРАЦИЯ И СБОРКА КОНТЕКСТА
|
||
# =====================================================
|
||
add_prof("[СБОРКА КОНТЕКСТА] Начат анализ найденных чанков...")
|
||
raw_contexts, sources, processed_windows, seen_chunk_hashes, qdrant_debug_lines = [], set(), set(), set(), []
|
||
contexts_added, total_context_size = 0, 0
|
||
query_lower = msg_text.lower()
|
||
query_words = set(w for w in re.findall(r"\w+", query_lower) if len(w) > 3)
|
||
|
||
for i, res in enumerate(search_results, 1):
|
||
if contexts_added >= MAX_AI_CONTEXTS:
|
||
add_prof(f" -> ПРЕРЫВАНИЕ: Достигнут лимит MAX_AI_CONTEXTS ({MAX_AI_CONTEXTS}).")
|
||
break
|
||
|
||
score = res.score
|
||
metadata = res.payload.get("metadata", {})
|
||
filename = metadata.get("filename", "unknown")
|
||
source_path = metadata.get("source", "")
|
||
chunk_idx = metadata.get("chunk_index", -1)
|
||
chunk_text = res.payload.get("page_content", "")
|
||
chunk_text_lower = chunk_text.lower()
|
||
|
||
chunk_words = set(w for w in re.findall(r"\w+", chunk_text_lower) if len(w) > 3)
|
||
keyword_overlap = len(query_words & chunk_words)
|
||
|
||
add_prof(f" -> [{i}] Оценка чанка: '{filename}' (idx={chunk_idx}). Score: {score:.4f}. Пересечение: {keyword_overlap} слов.")
|
||
|
||
if not live_use_reranker:
|
||
if keyword_overlap == 0: score -= 0.05
|
||
elif keyword_overlap == 1: score += 0.03
|
||
elif keyword_overlap >= 2: score += 0.10
|
||
|
||
if "хран" in query_lower and "хран" in filename.lower(): score += 0.10
|
||
if "долгосроч" in query_lower and ("долгосроч" in chunk_text_lower or "срок хран" in chunk_text_lower or "архив" in chunk_text_lower): score += 0.15
|
||
if "файл" in query_lower and ("файл" in chunk_text_lower or "электрон" in chunk_text_lower): score += 0.05
|
||
|
||
is_table = is_table_chunk(chunk_text)
|
||
|
||
# --- УНИВЕРСАЛЬНАЯ ДИНАМИЧЕСКАЯ ФИЛЬТРАЦИЯ ОТ МУСОРА ---
|
||
if live_use_reranker:
|
||
effective_threshold = getattr(current_config, 'RERANK_SCORE_THRESHOLD', -3.5)
|
||
else:
|
||
effective_threshold = TABLE_SCORE_THRESHOLD if is_table else QDRANT_SCORE_THRESHOLD
|
||
if is_table:
|
||
add_prof(f" - Обнаружена таблица. Порог снижен до {effective_threshold}.")
|
||
|
||
if score < effective_threshold:
|
||
add_prof(f" - ПРОПУСК: Итоговый score {score:.4f} ниже порога {effective_threshold}.")
|
||
qdrant_debug_lines.append(f"[SKIP] {score:.4f} < {effective_threshold} | {filename}")
|
||
continue
|
||
|
||
min_idx, max_idx = max(0, chunk_idx - QDRANT_CHUNK_WINDOW), chunk_idx + QDRANT_CHUNK_WINDOW
|
||
|
||
try:
|
||
t_scroll = time.time()
|
||
neighbor_records, _ = client.scroll(
|
||
collection_name=QDRANT_COLLECTION,
|
||
scroll_filter=Filter(
|
||
must=[
|
||
FieldCondition(key="metadata.source", match=MatchValue(value=source_path)),
|
||
FieldCondition(key="metadata.chunk_index", range=Range(gte=min_idx, lte=max_idx))
|
||
]
|
||
),
|
||
limit=QDRANT_SCROLL_LIMIT,
|
||
with_payload=True
|
||
)
|
||
add_prof(f" - SCROLL соседей ({min_idx}-{max_idx}) занял {time.time() - t_scroll:.2f}s. Найдено соседей: {len(neighbor_records)}.")
|
||
except Exception as e:
|
||
add_prof(f" - SCROLL ERROR: {e}")
|
||
qdrant_debug_lines.append(f"[SCROLL ERROR] {e}")
|
||
continue
|
||
|
||
if not neighbor_records:
|
||
continue
|
||
|
||
neighbor_records.sort(key=lambda x: x.payload["metadata"]["chunk_index"])
|
||
|
||
merged_parts = []
|
||
for p in neighbor_records:
|
||
idx = p.payload["metadata"].get("chunk_index", -1)
|
||
txt = p.payload.get("page_content", "")
|
||
if txt:
|
||
normalized_sub = normalize_text(txt)[:1200]
|
||
sub_hash = hash(normalized_sub)
|
||
if sub_hash in seen_chunk_hashes:
|
||
continue
|
||
|
||
seen_chunk_hashes.add(sub_hash)
|
||
merged_parts.append(f"[Чанк {idx}]\n{truncate_context(txt)}")
|
||
|
||
if not merged_parts:
|
||
continue
|
||
|
||
merged_text = "\n\n".join(merged_parts)
|
||
final_context = f"ФАЙЛ: {filename}\n\n{merged_text}"
|
||
|
||
if total_context_size + len(final_context) > MAX_CONTEXT_LENGTH:
|
||
add_prof(f" - ПРЕРЫВАНИЕ: Лимит контекста ({total_context_size + len(final_context)} > {MAX_CONTEXT_LENGTH} симв.)")
|
||
qdrant_debug_lines.append("[STOP] Достигнут лимит контекста")
|
||
break
|
||
|
||
total_context_size += len(final_context)
|
||
raw_contexts.append(final_context)
|
||
sources.add(filename)
|
||
contexts_added += 1
|
||
|
||
add_prof(f" - ДОБАВЛЕНО В ИТОГ. Текущий размер промпта: {total_context_size} симв.")
|
||
|
||
debug_preview = merged_text[:DEBUG_TEXT_LIMIT].replace("\n", " ")
|
||
orig_score = res.payload.get("original_qdrant_score", res.score)
|
||
qdrant_debug_lines.append(f"[{i}] orig={orig_score:.4f} | final={score:.4f} | overlap={keyword_overlap} | table={is_table} | file={filename} | chunk={chunk_idx}\npreview={debug_preview}...")
|
||
|
||
qdrant_debug_full = "\n".join(qdrant_debug_lines)
|
||
add_prof(f"[СБОРКА КОНТЕКСТА ЗАВЕРШЕНА] Итого контекстов: {contexts_added}. Источников: {len(sources)}.")
|
||
|
||
if not raw_contexts:
|
||
error_str = "Документы найдены Qdrant, но отсеяны фильтрами (score threshold)."
|
||
add_prof(f"[СТОП] {error_str}")
|
||
log_analytics(str(user_id), msg_text, "Не найдено", False, time.time() - start_time, qdrant_debug_full, set())
|
||
send_alert_email(subject="📭 Бот не нашел релевантных документов", body=f"Пользователь: {user_id}\nВопрос: {msg_text}\nПричина: {error_str}")
|
||
await msg.answer(search_not_found_text(EMOJI_DIGITS), parse_mode="html")
|
||
|
||
if ENABLE_DEEP_PROFILING:
|
||
with open(os.path.join(LOGS_PATH, "deep_profiling.log"), "a", encoding="utf-8") as f:
|
||
f.write(f"\n[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] USER: {user_id}\n" + "\n".join(prof_log) + "\n" + "="*80 + "\n")
|
||
return
|
||
|
||
# =====================================================
|
||
# AI REQUEST (LLM)
|
||
# =====================================================
|
||
try:
|
||
add_prof(f"[LLM START] Подготовка запроса к {getattr(current_config, 'OLLAMA_MODEL', 'gemma')}.")
|
||
t_llm = time.time()
|
||
|
||
ai_answer = await ask_ai(msg_text, raw_contexts)
|
||
|
||
add_prof(f"[LLM END] Ответ получен за {time.time() - t_llm:.2f}s.")
|
||
|
||
if not ai_answer or ai_answer.strip() == "":
|
||
add_prof(f"[LLM WARNING] Модель вернула пустую строку. Применяется фоллбэк.")
|
||
ai_answer = "В данных документах этот вопрос не регламентирован."
|
||
|
||
is_success = "не регламентирован" not in ai_answer.lower()
|
||
|
||
if not is_success:
|
||
send_alert_email(
|
||
subject="👁 Бот не смог найти answer в тексте (Не регламентирован)",
|
||
body=f"Пользователь: {user_id}\nВопрос: {msg_text}\n\nИИ ответил, что вопрос не регламентирован."
|
||
)
|
||
|
||
final_response = f"🤖 <b>ОТВЕТ ИИ:</b>\n\n{ai_answer}\n\n<i>{EMOJI_DIGITS['1']} — Ответ не удовлетворил\n{EMOJI_DIGITS['0']} — Выход в меню</i>"
|
||
|
||
add_prof(f"[ИТОГ] Запрос полностью выполнен за {time.time() - start_time:.2f}s.")
|
||
log_analytics(str(user_id), msg_text, ai_answer, is_success, time.time() - start_time, qdrant_debug_full, sources)
|
||
await msg.answer(final_response, parse_mode="html")
|
||
|
||
except Exception as e:
|
||
error_str = f"[ОШИБКА LLM]: {type(e).__name__} - {str(e)}"
|
||
add_prof(f"[FATAL ERROR] {error_str}")
|
||
qdrant_debug_full += f"\n\n[FATAL ERROR]\n{error_str}"
|
||
log_analytics(str(user_id), msg_text, error_str, False, time.time() - start_time, qdrant_debug_full, sources)
|
||
await handle_search_system_error(msg, str(user_id), "Запрос к серверу LLM", error_str)
|
||
|
||
if ENABLE_DEEP_PROFILING:
|
||
with open(os.path.join(LOGS_PATH, "deep_profiling.log"), "a", encoding="utf-8") as f:
|
||
f.write(f"\n[{datetime.now().strftime('%Y-%m-%d %H:%M:%S')}] USER: {user_id}\n" + "\n".join(prof_log) + "\n" + "="*80 + "\n") |