Files
trueconf_bot/search_bot/handlers.py

685 lines
36 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# /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")