Initial commit: TrueConf Chatbot КЛЕВЕР

This commit is contained in:
2026-07-15 12:00:27 +07:00
commit fc0288d13d
42 changed files with 7818 additions and 0 deletions
+685
View File
@@ -0,0 +1,685 @@
# /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")
+173
View File
@@ -0,0 +1,173 @@
import os
import sys
import string
import pickle # Библиотека для сохранения фрагментов текста на диск
# Загрузчики и разбивка текста
from langchain_community.document_loaders import DirectoryLoader, PyMuPDFLoader, Docx2txtLoader
from langchain_text_splitters import RecursiveCharacterTextSplitter
# Эмбеддинги и векторная база
from langchain_community.embeddings import HuggingFaceEmbeddings
from langchain_community.vectorstores import FAISS
# Ретриверы
from langchain_community.retrievers import BM25Retriever
from langchain_classic.retrievers import EnsembleRetriever
# Цепочки
from langchain_classic.chains.combine_documents import create_stuff_documents_chain
from langchain_classic.chains import create_retrieval_chain
# Базовые модули и LLM
from langchain_core.prompts import ChatPromptTemplate
from langchain_openai import ChatOpenAI
# ==========================================
# НАСТРОЙКА ПУТЕЙ
# ==========================================
BASE_DIR = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
BASE_REGULATIONS_PATH = os.path.join(BASE_DIR, "search_bot", "Регламенты", "АОХКСибцем", "07 Служба Вице-президента по экономике и финансам", "ДИТ")
MODEL_PATH = os.path.join(BASE_DIR, "models", "mpnet_model")
# Пути для сохранения базы данных
FAISS_DB_PATH = os.path.join(BASE_DIR, "search_bot", "faiss_index")
CHUNKS_PATH = os.path.join(BASE_DIR, "search_bot", "chunks.pkl")
# ==========================================
# 1. ИНИЦИАЛИЗАЦИЯ ЭМБЕДДИНГОВ
# ==========================================
print("Загрузка модели эмбеддингов (оффлайн)...")
embedding = HuggingFaceEmbeddings(
model_name=MODEL_PATH,
model_kwargs={'device': 'cpu'},
encode_kwargs={'normalize_embeddings': False}
)
# ==========================================
# 2. ЗАГРУЗКА ИЛИ СОЗДАНИЕ БАЗЫ ДАННЫХ
# ==========================================
# Проверяем, существует ли уже сохраненная база
if os.path.exists(FAISS_DB_PATH) and os.path.exists(CHUNKS_PATH):
print("✅ Найдена сохраненная база данных! Загружаю с диска (это быстро)...")
# Загружаем фрагменты текста для BM25
with open(CHUNKS_PATH, 'rb') as f:
split_docs = pickle.load(f)
# Загружаем векторы FAISS
vector_store = FAISS.load_local(
FAISS_DB_PATH,
embedding,
allow_dangerous_deserialization=True # Обязательно для локальных файлов
)
else:
print("⚠️ Сохраненная база не найдена. Начинаю чтение и индексацию документов...\n")
if not os.path.exists(BASE_REGULATIONS_PATH):
print("❌ ОШИБКА: Указанный путь к документам не существует.")
sys.exit(1)
# Загружаем PDF и DOCX
pdf_loader = DirectoryLoader(BASE_REGULATIONS_PATH, glob="**/*.pdf", loader_cls=PyMuPDFLoader, show_progress=True)
docx_loader = DirectoryLoader(BASE_REGULATIONS_PATH, glob="**/*.docx", loader_cls=Docx2txtLoader, show_progress=True)
docs = pdf_loader.load() + docx_loader.load()
if len(docs) == 0:
print("❌ В указанной папке нет файлов .pdf или .docx. Завершение.")
sys.exit(1)
print(f"✅ Загружено страниц/файлов: {len(docs)}. Разбивка на фрагменты...")
text_splitter = RecursiveCharacterTextSplitter(chunk_size=1200, chunk_overlap=200, separators=["\n\n", "\n", ".", " "])
split_docs = text_splitter.split_documents(docs)
print("Создаю векторную базу FAISS...")
vector_store = FAISS.from_documents(split_docs, embedding=embedding)
# СОХРАНЕНИЕ НА ДИСК ДЛЯ БУДУЩИХ ЗАПУСКОВ
print("💾 Сохраняю базу данных на диск...")
vector_store.save_local(FAISS_DB_PATH)
with open(CHUNKS_PATH, 'wb') as f:
pickle.dump(split_docs, f)
print("✅ База успешно сохранена!")
# ==========================================
# 3. НАСТРОЙКА РЕТРИВЕРОВ (ПОИСКА)
# ==========================================
# 3.1 Семантический поиск
embedding_retriever = vector_store.as_retriever(search_kwargs={"k": 4})
# 3.2 Лексический поиск
def tokenize(s):
return s.lower().translate(str.maketrans("", "", string.punctuation)).split(" ")
bm25_retriever = BM25Retriever.from_documents(
documents=split_docs,
preprocess_func=tokenize,
k=5
)
# 3.3 Гибридный поиск (Ансамбль)
ensemble_retriever = EnsembleRetriever(
retrievers=[embedding_retriever, bm25_retriever],
weights=[0.4, 0.6]
)
# ==========================================
# 4. ПОДКЛЮЧЕНИЕ К LLAMA.CPP
# ==========================================
print("Подключение к серверу llama.cpp...")
llm = ChatOpenAI(
base_url="http://127.0.0.1:8080/v1",
api_key="not-needed",
temperature=0.0,
max_tokens=1024
)
prompt = ChatPromptTemplate.from_template('''Ты — строгий корпоративный ИИ-помощник.
Твоя задача — отвечать на вопросы строго на основании текста предоставленных регламентов.
Внимательно изучи контекст. Если там есть ответ, сформулируй его четко и по делу. Обязательно указывай номер пункта или название документа, откуда взята информация.
Если в контексте НЕТ ответа на вопрос, не придумывай информацию, а выведи фразу: "В предоставленных регламентах нет информации по данному вопросу."
Контекст:
{context}
Вопрос пользователя: {input}
Ответ:'''
)
document_chain = create_stuff_documents_chain(llm=llm, prompt=prompt)
rag_chain = create_retrieval_chain(ensemble_retriever, document_chain)
# ==========================================
# 5. ТЕСТИРОВАНИЕ СИСТЕМЫ
# ==========================================
print("\n🚀 Система готова к работе! Выполняем запросы...\n")
# Тестовые вопросы по Положению № ПОЛ-177
questions = [
"Где работник обязан хранить электронные документы, связанные с производственной деятельностью?",
"Через какой срок удаляются электронные документы, если их не открывали, и можно ли их восстановить?",
"Кому разрешен доступ к сетевым ресурсам Soft/Distr?",
"Какие правила хранения и удаления установлены для медиафайлов в сетевой папке 'СВК'?",
"Разрешено ли самостоятельно предоставлять общий доступ к папкам на своем рабочем компьютере?",
"Что имеют право сделать сотрудники ДИТ или СИБ, если файл несет угрозу ресурсам?"
]
for q in questions:
print(f"❓ Вопрос: {q}")
try:
response = rag_chain.invoke({'input': q})
print(f"🤖 Ответ: {response['answer']}\n")
# Распечатка источников (полезно для отладки)
print("🔍 Найденные источники:")
for i, doc in enumerate(response['context']):
source_name = doc.metadata.get('source', 'Неизвестный источник').split('/')[-1] # Берем только имя файла
print(f" [{i+1}] Файл: {source_name}")
print("\n" + "-" * 50 + "\n")
except Exception as e:
print(f"❌ Ошибка при запросе к серверу llama.cpp: {e}")
print("-" * 50)
+70
View File
@@ -0,0 +1,70 @@
import os
import requests
import json
# Адрес вашего сервера с поднятым Docker-контейнером Unstructured
API_URL = "http://172.16.80.250:9000/general/v0/general"
def main():
# 1. Ищем все файлы .pdf и .docx в текущей директории
target_files = [f for f in os.listdir('.') if f.lower().endswith(('.pdf', '.docx'))]
target_files.sort() # Для порядка
if not target_files:
print("❌ В текущей папке не найдено ни одного PDF или DOCX файла.")
return
print(f"📂 Найдено файлов для обработки: {len(target_files)}")
# 2. Проходимся по всем найденным файлам
for target_file in target_files:
print(f"\n📄 Обрабатывается файл: {target_file}")
print(f"⏳ Отправляем на сервер {API_URL}... (это может занять время)")
# Настройки запроса (параметры PDF просто игнорируются сервером для DOCX)
payload_data = {
"strategy": "hi_res",
"languages": ["rus"],
"pdf_infer_table_structure": "true"
}
# Определяем правильный MIME-тип для файла
ext = target_file.lower().split('.')[-1]
if ext == 'pdf':
mime_type = "application/pdf"
else:
mime_type = "application/vnd.openxmlformats-officedocument.wordprocessingml.document"
# 3. Открываем файл и делаем POST-запрос
try:
with open(target_file, "rb") as f:
files = {
"files": (target_file, f, mime_type)
}
# Таймаут 300 секунд
response = requests.post(API_URL, files=files, data=payload_data, timeout=300)
response.raise_for_status()
# 4. Получаем JSON и сохраняем его в файл
parsed_data = response.json()
output_filename = f"{os.path.splitext(target_file)[0]}_parsed.json"
with open(output_filename, "w", encoding="utf-8") as out_f:
json.dump(parsed_data, out_f, ensure_ascii=False, indent=4)
print(f"✅ Готово! Результат успешно сохранен в файл: {output_filename}")
except requests.exceptions.ConnectionError:
print(f"❌ Ошибка соединения! Проверьте доступность {API_URL}")
except requests.exceptions.Timeout:
print("❌ Сервер слишком долго не отвечал.")
except requests.exceptions.HTTPError as e:
print(f"❌ Сервер вернул ошибку: {e}")
if response.text:
print(f"Детали ошибки от сервера: {response.text}")
except Exception as e:
print(f"❌ Произошла непредвиденная ошибка при обработке {target_file}: {e}")
if __name__ == "__main__":
main()
+650
View File
@@ -0,0 +1,650 @@
import os
import sys
import re
import time
import uuid
import json
import hashlib
import requests
from glob import glob
from datetime import datetime
# Импорт SentenceTransformer УДАЛЕН для разгрузки ОЗУ сервера
from qdrant_client import QdrantClient
from qdrant_client.models import (
Distance,
VectorParams,
PointStruct,
PayloadSchemaType,
Filter,
FieldCondition,
MatchValue
)
# =========================================================
# ИМПОРТ ЦЕНТРАЛЬНОГО КОНФИГА БОТА
# =========================================================
# Добавляем корневую папку бота в sys.path, чтобы скрипт увидел модуль config,
# даже если скрипт запускается из подпапки.
sys.path.append(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
import config.config as config
# =========================================================
# CONFIG (Берем данные из центрального config.py)
# =========================================================
DOCS_DIR = config.SEARCH_DOCS_PATH
QDRANT_HOST = config.QDRANT_HOST
QDRANT_PORT = config.QDRANT_PORT
COLLECTION_NAME = config.QDRANT_COLLECTION
# Формируем путь к реестру внутри папки логов, указанной в конфиге
REGISTRY_PATH = os.path.join(config.LOGS_PATH, "indexing_registry.json")
# =========================================================
# ЛОКАЛЬНЫЕ НАСТРОЙКИ ИНДЕКСАТОРА
# =========================================================
UNSTRUCTURED_API_URL = "http://172.16.80.250:9000/general/v0/general"
BATCH_SIZE = 32
RECREATE_COLLECTION = False
MIN_TEXT_LENGTH = 25
# =========================================================
# INIT
# =========================================================
# Локальная инициализация нейросети убрана. Работаем через API.
print("🔄 Подключение к Qdrant...")
client = QdrantClient(
host=QDRANT_HOST,
port=QDRANT_PORT
)
# =========================================================
# REGISTRY & HASHING
# =========================================================
def get_file_hash(filepath):
"""Вычисляет MD5 хэш физического файла"""
hasher = hashlib.md5()
with open(filepath, 'rb') as f:
buf = f.read()
hasher.update(buf)
return hasher.hexdigest()
def load_registry():
"""Загружает текущее состояние проиндексированных файлов"""
if os.path.exists(REGISTRY_PATH):
with open(REGISTRY_PATH, 'r', encoding='utf-8') as f:
return json.load(f)
return {}
def save_registry(registry):
"""Сохраняет реестр на диск"""
os.makedirs(os.path.dirname(REGISTRY_PATH), exist_ok=True)
with open(REGISTRY_PATH, 'w', encoding='utf-8') as f:
json.dump(registry, f, ensure_ascii=False, indent=4)
# =========================================================
# HELPERS
# =========================================================
def sha256_text(text):
return hashlib.sha256(
text.encode("utf-8")
).hexdigest()
def clean_text(text):
if not text:
return ""
text = re.sub(
r"\s+",
" ",
text
)
return text.strip()
def is_page_number(text):
return bool(
re.match(
r"^Стр\.\s*\d+\s*из\s*\d+$",
text
)
)
def is_header_noise(text):
patterns = [
"АО «ХК «Сибцем»",
"Тип документа:",
"Ведущее подразделение:",
"Дата утверждения:",
"Редакция 1",
"Оглавление"
]
return any(
p in text
for p in patterns
)
def is_junk(item, text):
if not text:
return True
if len(text) < MIN_TEXT_LENGTH:
return True
if item["type"] in [
"Header",
"Footer",
"Image"
]:
return True
if is_page_number(text):
return True
if is_header_noise(text):
return True
return False
# =========================================================
# TABLE PARSER
# =========================================================
def html_table_to_text(html):
html = re.sub(
r"\n+",
" ",
html
)
rows = re.findall(
r"<tr.*?>(.*?)</tr>",
html,
re.I | re.S
)
parsed_rows = []
for row in rows:
cells = re.findall(
r"<t[hd].*?>(.*?)</t[hd]>",
row,
re.I | re.S
)
clean_cells = []
for cell in cells:
box = re.sub(
r"<[^>]+>",
" ",
cell
)
box = clean_text(box)
clean_cells.append(box)
if clean_cells:
parsed_rows.append(
clean_cells
)
if not parsed_rows:
return ""
headers = parsed_rows[0]
result_lines = []
for row_idx, row in enumerate(parsed_rows[1:], 1):
result_lines.append(
f"Запись {row_idx}:"
)
for col_idx, cell in enumerate(row):
if not cell:
continue
if col_idx < len(headers):
header = headers[col_idx]
else:
header = f"Колонка {col_idx+1}"
result_lines.append(
f"{header}: {cell}"
)
result_lines.append("")
return "\n".join(result_lines).strip()
# =========================================================
# MERGE FRAGMENTS
# =========================================================
def merge_fragments(elements):
merged = []
buffer = []
buffer_page = None
buffer_section = None
def flush_buffer():
nonlocal buffer
if not buffer:
return
text = clean_text(
" ".join(buffer)
)
if text:
merged.append({
"type": "MergedText",
"text": text,
"page_number": buffer_page,
"section_title": buffer_section
})
buffer = []
for el in elements:
text = clean_text(
el["text"]
)
if not text:
continue
el_type = el["type"]
page = el["metadata"].get("page_number", 0)
section = el.get("section_title", "")
# =========================================
# TABLE
# =========================================
if el_type == "Table":
flush_buffer()
merged.append({
"type": "Table",
"text": text,
"html": el["metadata"].get("text_as_html", ""),
"page_number": page,
"section_title": section
})
continue
# =========================================
# SHORT TEXT
# =========================================
if el_type in [
"NarrativeText",
"ListItem",
"UncategorizedText"
]:
if len(text) < 80:
buffer.append(text)
buffer_page = page
buffer_section = section
else:
flush_buffer()
merged.append({
"type": "Text",
"text": text,
"page_number": page,
"section_title": section
})
else:
flush_buffer()
flush_buffer()
return merged
# =========================================================
# QDRANT INIT & DELETE
# =========================================================
def init_collection():
collections = [
c.name
for c in client.get_collections().collections
]
if (
RECREATE_COLLECTION
and COLLECTION_NAME in collections
):
print("🗑 Удаление коллекции...")
client.delete_collection(COLLECTION_NAME)
collections = [
c.name
for c in client.get_collections().collections
]
if COLLECTION_NAME not in collections:
print("⚙️ Создание коллекции...")
client.create_collection(
collection_name=COLLECTION_NAME,
vectors_config=VectorParams(
size=1024,
distance=Distance.COSINE
)
)
indexes = [
("metadata.filename", PayloadSchemaType.KEYWORD),
("metadata.section_title", PayloadSchemaType.KEYWORD),
("metadata.page_number", PayloadSchemaType.INTEGER),
("metadata.content_type", PayloadSchemaType.KEYWORD)
]
for field_name, field_schema in indexes:
client.create_payload_index(
collection_name=COLLECTION_NAME,
field_name=field_name,
field_schema=field_schema
)
def delete_file_from_qdrant(filename):
"""Удаляет все векторы, привязанные к конкретному файлу"""
client.delete(
collection_name=COLLECTION_NAME,
points_selector=Filter(
must=[
FieldCondition(
key="metadata.filename",
match=MatchValue(value=filename)
)
]
)
)
print(f" 🗑 Удалены старые данные файла {filename} из базы.")
# =========================================================
# PARSE DOCUMENT
# =========================================================
def parse_document(file_path):
filename = os.path.basename(
file_path
)
print(f"\n📄 Парсинг: {filename}")
payload_data = {
"strategy": "hi_res",
"languages": ["rus"],
"pdf_infer_table_structure": "true"
}
ext = filename.lower().split(".")[-1]
if ext == "pdf":
mime_type = "application/pdf"
else:
mime_type = (
"application/vnd.openxmlformats-officedocument."
"wordprocessingml.document"
)
with open(file_path, "rb") as f:
files = {
"files": (
filename,
f,
mime_type
)
}
response = requests.post(
UNSTRUCTURED_API_URL,
files=files,
data=payload_data,
timeout=900
)
response.raise_for_status()
return response.json()
# =========================================================
# INDEX DOCUMENT
# =========================================================
def index_document(file_path, data):
filename = os.path.basename(file_path)
cleaned = []
current_title = "Общая информация"
# =====================================================
# CLEAN
# =====================================================
for item in data:
raw_text = item.get("text", "")
text = clean_text(raw_text)
if is_junk(item, text):
continue
if item["type"] == "Title":
current_title = text
continue
item["section_title"] = current_title
cleaned.append(item)
# =====================================================
# MERGE
# =====================================================
merged = merge_fragments(cleaned)
points = []
# =====================================================
# BUILD CHUNKS
# =====================================================
for idx, chunk in enumerate(merged):
chunk_type = chunk["type"]
section_title = chunk["section_title"]
page_number = chunk["page_number"]
# =================================================
# TABLE / TEXT
# =================================================
if chunk_type == "Table":
table_text = html_table_to_text(chunk["html"])
content = (
f"Раздел: {section_title}\n"
f"Тип: Таблица\n\n"
f"{table_text}"
)
else:
content = (
f"Раздел: {section_title}\n\n"
f"{chunk['text']}"
)
content = clean_text(content)
if len(content) < 40:
continue
# =================================================
# VECTOR (Переведено на центральный Embedding API)
# =================================================
try:
payload_api = {
"model": config.EMBEDDING_MODEL_NAME, # Отправляем "bge-m3"
"input_type": "passage", # Для сохранения документов используем тип passage
"texts": [content]
}
headers_api = {"Content-Type": "application/json"}
if hasattr(config, 'EMBEDDING_API_TOKEN') and config.EMBEDDING_API_TOKEN:
headers_api["X-Embedding-Token"] = config.EMBEDDING_API_TOKEN
api_resp = requests.post(config.EMBEDDING_API_URL, json=payload_api, headers=headers_api)
api_resp.raise_for_status()
vector = api_resp.json()["embeddings"][0]
except Exception as e:
print(f"❌ Ошибка векторизации чанка {idx} через API: {e}")
continue # Пропускаем чанк, если центральный сервис недоступен
chunk_hash = sha256_text(
filename + str(idx) + content
)
point_id = str(
uuid.uuid5(
uuid.NAMESPACE_DNS,
chunk_hash
)
)
payload = {
"page_content": content,
"metadata": {
"filename": filename,
"source": file_path,
"chunk_index": idx,
"section_title": section_title,
"page_number": page_number,
"content_type": chunk_type,
"chunk_hash": chunk_hash
}
}
points.append(
PointStruct(
id=point_id,
vector=vector,
payload=payload
)
)
# =================================================
# UPSERT BATCH
# =================================================
if len(points) >= BATCH_SIZE:
client.upsert(
collection_name=COLLECTION_NAME,
points=points
)
print(f" ⬆️ batch={len(points)}")
points = []
# =====================================================
# FINAL UPSERT
# =====================================================
if points:
client.upsert(
collection_name=COLLECTION_NAME,
points=points
)
print(f" ✅ Загружено chunks: {len(merged)}")
return len(merged)
# =========================================================
# MAIN
# =========================================================
def main():
init_collection()
target_files = []
for ext in ["*.pdf", "*.docx"]:
target_files.extend(
glob(
os.path.join(DOCS_DIR, "**", ext),
recursive=True
)
)
# ФИЛЬТРАЦИЯ СКРЫТЫХ ВРЕМЕННЫХ ФАЙЛОВ WORD (~$)
target_files = [f for f in target_files if not os.path.basename(f).startswith("~$")]
target_files.sort()
if not target_files:
print("❌ PDF/DOCX файлы не найдены")
return
print(f"📂 Найдено файлов для индексации: {len(target_files)}")
print(f"📡 Использование удаленной векторизации: {config.EMBEDDING_API_URL} ({config.EMBEDDING_MODEL_NAME})")
# Загружаем реестр
registry = load_registry()
# =====================================================
# PROCESS
# =====================================================
for file_path in target_files:
filename = os.path.basename(file_path)
current_hash = get_file_hash(file_path)
file_info = registry.get(filename)
# Проверка: Файл уже есть, хэш совпадает, статус был успешным
if file_info and file_info.get("hash") == current_hash and file_info.get("status") == "SUCCESS":
print(f"\n⏭️ Пропуск (не изменился): {filename}")
continue
print(f"\n🔄 Обработка: {filename}")
# Если файл есть в реестре со статусом SUCCESS, но хэш другой (файл обновили)
if file_info and file_info.get("status") == "SUCCESS":
print(f" ⚠️ Обнаружена новая версия файла! Обновляем индекс...")
delete_file_from_qdrant(filename)
try:
data = parse_document(file_path)
chunks_count = index_document(file_path, data)
# Запись успешного результата
registry[filename] = {
"hash": current_hash,
"status": "SUCCESS",
"chunks": chunks_count,
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
# ПАУЗА ПРИ НЕДОСТУПНОСТИ СЕРВЕРА ПАРСИНГА
except requests.exceptions.ConnectionError:
print(f" ⚠️ Сервер парсинга документов недоступен. Ждем 60 секунд...")
time.sleep(60)
registry[filename] = {
"hash": current_hash,
"status": "ERROR",
"error_details": "Connection refused",
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
except Exception as e:
print(f" ❌ Ошибка обработки файла {filename}: {e}")
# Запись ошибки в реестр
registry[filename] = {
"hash": current_hash,
"status": "ERROR",
"error_details": str(e),
"timestamp": datetime.now().strftime("%Y-%m-%d %H:%M:%S")
}
save_registry(registry)
print("\n🎉 ГОТОВО. База Qdrant полностью синхронизирована через API.")
# =========================================================
if __name__ == "__main__":
main()