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"(.*?)", html, re.I | re.S ) parsed_rows = [] for row in rows: cells = re.findall( r"(.*?)", 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()