"""Celery задачи индексации документов и проверки плагиата (уровни 1-2).""" import io import logging from pathlib import Path from typing import Any from celery.utils.log import get_task_logger from sqlalchemy import func, select from app.algorithms.minhash import add_to_lsh, find_similar from app.algorithms.winnowing import winnow from app.celery_app import celery_app from app.config import settings from app.db import db_session, get_minio, update_task_status from app.extractors.docx import extract_text_from_docx, extract_text_from_txt from app.extractors.pdf import extract_text_from_pdf logger = get_task_logger(__name__) def _split_into_fragments( text: str, window: int = 200, overlap: int = 50, ) -> list[dict[str, Any]]: """ Разбить текст на фрагменты для проверки плагиата. Использует скользящее окно с перекрытием. Args: text: Исходный текст window: Размер окна в словах overlap: Перекрытие между фрагментами в словах Returns: Список словарей {"text": str, "start": int, "end": int} """ words = text.split() if not words: return [] fragments = [] step = window - overlap char_positions = [] # Вычислить позиции символов для каждого слова pos = 0 for word in words: char_positions.append(pos) pos += len(word) + 1 # +1 для пробела for i in range(0, max(1, len(words) - window + 1), step): chunk_words = words[i : i + window] if len(chunk_words) < 20: # Пропустить слишком короткие фрагменты continue start_char = char_positions[i] end_idx = min(i + window - 1, len(words) - 1) end_char = char_positions[end_idx] + len(words[end_idx]) fragments.append({ "text": " ".join(chunk_words), "start": start_char, "end": end_char, }) return fragments @celery_app.task( name="index.extract_and_check", bind=True, max_retries=3, default_retry_delay=60, ) def extract_and_check( self, task_id: str, minio_key: str, filename: str, ) -> dict[str, Any]: """ Извлечь текст из документа и проверить плагиат (уровни 1-2). Алгоритм: 1. Скачать файл из MinIO 2. Извлечь текст (PDF/DOCX/TXT) 3. Разбить на фрагменты по 200 слов с перекрытием 50 слов 4. Уровень 1: Winnowing против базы fingerprints в PostgreSQL 5. Уровень 2: MinHash LSH нечёткий поиск 6. Диспатч gpu.check_plagiarism для уровней 3 и 4 Args: task_id: ID задачи в PostgreSQL minio_key: Ключ объекта в MinIO filename: Оригинальное имя файла """ logger.info(f"Извлечение текста для задачи {task_id!r}, файл: {filename!r}") update_task_status(task_id, "processing") try: # Скачать из MinIO minio = get_minio() response = minio.get_object(settings.MINIO_BUCKET_DOCS, minio_key) file_data = response.read() response.close() response.release_conn() logger.info(f"Файл скачан из MinIO: {minio_key} ({len(file_data)} байт)") # Извлечь текст в зависимости от формата ext = Path(filename).suffix.lower() if ext == ".pdf": text = extract_text_from_pdf(file_data) elif ext == ".docx": text = extract_text_from_docx(file_data) else: text = extract_text_from_txt(file_data) if not text.strip(): raise ValueError("Не удалось извлечь текст из документа") word_count = len(text.split()) logger.info(f"Текст извлечён: {len(text)} символов, {word_count} слов") # ──── Автосбор в отстойник (буфер для пополнения базы) ──────────────── # Сохраняем извлечённый текст и метаданные для последующего ручного # одобрения админом. Не дублируем в базу источников до решения. _stage_work(task_id, minio_key, filename, text, word_count) # Разбить на фрагменты fragments = _split_into_fragments( text, window=settings.FRAGMENT_WINDOW_WORDS, overlap=settings.FRAGMENT_OVERLAP_WORDS, ) logger.info(f"Фрагментов создано: {len(fragments)}") # ──── Уровень 1: Winnowing fingerprints ──────────────────────────────── level1_matches: list[dict] = [] doc_fingerprint = winnow(text) if doc_fingerprint: from app.models import Document, Fingerprint with db_session() as session: hashes = list(doc_fingerprint)[: settings.MAX_FINGERPRINTS_PER_DOC] # Найти совпадения в базе fingerprints matching_docs = session.execute( select( Fingerprint.doc_id, func.count(Fingerprint.id).label("match_count"), ) .where(Fingerprint.hash_value.in_(hashes)) .group_by(Fingerprint.doc_id) .having(func.count(Fingerprint.id) > len(hashes) * 0.1) .order_by(func.count(Fingerprint.id).desc()) .limit(20) ).all() for doc_id, match_count in matching_docs: similarity = match_count / len(hashes) * 100 if similarity < 20: continue doc = session.get(Document, doc_id) if not doc: continue level1_matches.append({ "fragment": text[:200], "position_start": 0, "position_end": len(text), "similarity": round(similarity, 1), "method": "exact", "source_title": doc.title, "source_url": doc.url, "source_db": doc.source, }) logger.info(f"Уровень 1 (Winnowing): {len(level1_matches)} совпадений") # ──── Уровень 2: MinHash LSH ──────────────────────────────────────────── level2_matches: list[dict] = [] similar_keys = find_similar(text) if similar_keys: from app.models import Document with db_session() as session: for key in similar_keys[:10]: # Ключ формата "doc:{id}" try: doc_id = int(key.split(":")[-1]) doc = session.get(Document, doc_id) if not doc: continue level2_matches.append({ "fragment": text[:200], "position_start": 0, "position_end": len(text), "similarity": 60.0, # MinHash даёт только факт похожести "method": "fuzzy", "source_title": doc.title, "source_url": doc.url, "source_db": doc.source, }) except (ValueError, IndexError): continue logger.info(f"Уровень 2 (MinHash): {len(level2_matches)} совпадений") # ──── Диспатч GPU задачи (уровни 3-4) ───────────────────────────────── # Передаём только текстовые данные (JSON-сериализуемые) celery_app.send_task( "gpu.check_plagiarism", args=[task_id, text, fragments], kwargs={ "level1_matches": level1_matches, "level2_matches": level2_matches, }, queue="queue.gpu", ) logger.info(f"GPU задача отправлена для задачи {task_id!r}") return { "task_id": task_id, "fragments": len(fragments), "level1": len(level1_matches), "level2": len(level2_matches), } except Exception as exc: logger.error(f"Ошибка при обработке задачи {task_id!r}: {exc}", exc_info=True) update_task_status(task_id, "failed", str(exc)) raise self.retry(exc=exc, countdown=60) @celery_app.task(name="index.add_document") def add_document(doc_data: dict[str, Any], dispatch_embed: bool = True) -> dict[str, Any]: """ Добавить документ из внешнего источника в систему. Алгоритм: 1. Дедупликация по ext_id 2. Сохранить метаданные в PostgreSQL 3. Вычислить провизорные Winnowing fingerprints (из аннотации) 4. Добавить в MinHash LSH 5. Индексировать в Elasticsearch 6. Диспатч gpu.embed_documents для FAISS эмбеддингов (если dispatch_embed) 7. Если включён FETCH_FULL_TEXT и есть url — диспатч index.enrich_full_text, который скачает полный текст и пересчитает fingerprints по нему. Args: doc_data: Словарь с метаданными документа dispatch_embed: Диспатчить ли эмбеддинг по одному документу. При массовой заливке run_parser выключает это и батчит эмбеддинги сам. Returns: dict со статусом операции и doc_id """ from app.models import Document, Fingerprint ext_id = doc_data.get("ext_id") if not ext_id: return {"status": "error", "reason": "ext_id обязателен"} with db_session() as session: # Проверить дублирование existing = session.execute( select(Document).where(Document.ext_id == ext_id) ).scalar_one_or_none() if existing: return {"status": "duplicate", "doc_id": existing.id} # Создать документ allowed_fields = {c.key for c in Document.__table__.columns} doc_kwargs = {k: v for k, v in doc_data.items() if k in allowed_fields} doc = Document(**doc_kwargs) session.add(doc) session.flush() doc_id = doc.id # Вычислить fingerprints text = doc_data.get("full_text") or doc_data.get("abstract", "") or "" if text: fp = winnow(text) fingerprints_to_add = list(fp)[: settings.MAX_FINGERPRINTS_PER_DOC] for i, hash_val in enumerate(fingerprints_to_add): session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i)) # Добавить в MinHash LSH (in-memory) add_to_lsh(f"doc:{doc_id}", text) session.commit() logger.info(f"Документ {doc_id} добавлен в PostgreSQL: {doc_data.get('title', '')[:50]!r}") # Индексация в Elasticsearch try: from elasticsearch import Elasticsearch from app.config import settings as cfg es = Elasticsearch(cfg.ELASTICSEARCH_URL) es_doc = { "doc_id": doc_id, "source": doc_data.get("source"), "title": doc_data.get("title", ""), "abstract": doc_data.get("abstract", ""), "authors": " ".join( f"{a.get('last_name', '')} {a.get('first_name', '')}" for a in doc_data.get("authors", []) ), "year": doc_data.get("year"), "lang": doc_data.get("lang"), "journal": doc_data.get("journal"), "doi": doc_data.get("doi"), "url": doc_data.get("url"), } es.index(index="documents", id=str(doc_id), document=es_doc) logger.info(f"Документ {doc_id} проиндексирован в Elasticsearch") except Exception as e: logger.warning(f"Ошибка индексации в ES для документа {doc_id}: {e}") # Диспатч FAISS эмбеддингов (по одному документу; при массовой заливке # run_parser выключает это и батчит сам) if dispatch_embed: celery_app.send_task( "gpu.embed_documents", args=[[doc_id]], queue="queue.gpu", ) # Обогащение полным текстом: скачать PDF по url и пересчитать fingerprints if settings.FETCH_FULL_TEXT and doc_data.get("url"): celery_app.send_task( "index.enrich_full_text", args=[doc_id, doc_data["url"]], queue="queue.index", ) return {"status": "indexed", "doc_id": doc_id} @celery_app.task( name="index.enrich_full_text", bind=True, max_retries=2, default_retry_delay=120, ) def enrich_full_text(self, doc_id: int, url: str) -> dict[str, Any]: """Скачать полный текст статьи и пересчитать по нему fingerprints. Метаданные и провизорные fingerprints (из аннотации) уже сохранены в add_document. Здесь мы: 1. Скачиваем PDF по url и извлекаем текст (best-effort). 2. Сохраняем полный текст в MinIO (bucket documents, префикс corpus/). 3. Заменяем fingerprints документа на посчитанные по полному тексту. 4. Обновляем MinHash LSH. Недоступный/не-PDF источник — не ошибка: возвращаем no_fulltext. """ from app.fulltext import fetch_full_text from app.models import Document, Fingerprint text = fetch_full_text(url) if not text: return {"status": "no_fulltext", "doc_id": doc_id} # Сохранить полный текст в MinIO try: minio = get_minio() key = f"corpus/{doc_id}.txt" data = text.encode("utf-8") minio.put_object( settings.MINIO_BUCKET_DOCS, key, io.BytesIO(data), length=len(data), content_type="text/plain; charset=utf-8", ) except Exception as exc: logger.error(f"enrich_full_text: не удалось сохранить текст в MinIO для {doc_id}: {exc}") raise self.retry(exc=exc, countdown=120) # Пересчитать fingerprints по полному тексту doc_fp = winnow(text) hashes = list(doc_fp)[: settings.MAX_FINGERPRINTS_PER_DOC] from sqlalchemy import delete with db_session() as session: doc = session.get(Document, doc_id) if doc is None: return {"status": "doc_gone", "doc_id": doc_id} doc.minio_key = key # Удалить провизорные fingerprints и записать новые session.execute(delete(Fingerprint).where(Fingerprint.doc_id == doc_id)) for i, hash_val in enumerate(hashes): session.add(Fingerprint(doc_id=doc_id, hash_value=hash_val, position=i)) session.commit() # Обновить MinHash LSH по полному тексту add_to_lsh(f"doc:{doc_id}", text) logger.info( f"enrich_full_text: doc={doc_id} полный текст {len(text)} симв., " f"fingerprints={len(hashes)}" ) return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)} def _stage_work( task_id: str, minio_key: str, filename: str, text: str, word_count: int, ) -> None: """Сохранить проверенную работу в отстойник (StagedWork) + текст в MinIO. Идемпотентно по task_id: повторный вызов (например при requeue) не дублирует. Ошибки логируются, но не валят основную проверку плагиата. """ try: from app.models import StagedWork # Сохранить извлечённый текст в bucket staging text_key = f"text/{task_id}.txt" minio = get_minio() bucket = settings.MINIO_BUCKET_STAGING if not minio.bucket_exists(bucket): minio.make_bucket(bucket) data = text.encode("utf-8") minio.put_object( bucket, text_key, io.BytesIO(data), length=len(data), content_type="text/plain; charset=utf-8", ) with db_session() as session: existing = session.execute( select(StagedWork).where(StagedWork.task_id == task_id) ).scalar_one_or_none() if existing: return # Достать user_id из задачи from app.models import Task task = session.get(Task, task_id) user_id = task.user_id if task else None session.add(StagedWork( user_id=user_id, task_id=task_id, filename=filename, minio_key=minio_key, text_key=text_key, title=Path(filename).stem if filename else None, word_count=word_count, status="pending", )) session.commit() logger.info(f"Работа задачи {task_id!r} добавлена в отстойник") except Exception as e: logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}") @celery_app.task(name="index.run_parser") def run_parser(source_id: int) -> dict[str, Any]: """Запустить парсинг источника и наполнить базу документов. Переиспользует парсеры из scripts/parsers (BaseParser.run) и задачу add_document для каждого полученного документа. """ import sys from datetime import datetime, timezone from app.models import ParseSource # Парсеры лежат в /parsers (скопированы в образ) if "/parsers" not in sys.path: sys.path.insert(0, "/parsers") with db_session() as session: src = session.get(ParseSource, source_id) if src is None: return {"status": "error", "reason": "источник не найден"} cfg = { "source_type": src.source_type, "query": src.query, "lang": src.lang, "year_from": src.year_from, "year_to": src.year_to, "limit": src.limit, } src.last_status = "running" session.commit() added = 0 error_msg = None try: # Выбрать парсер по типу и собрать совместимые с его fetch() аргументы stype = cfg["source_type"] if stype == "openalex": from openalex import OpenAlexParser as P fetch_kwargs = { "query": cfg.get("query") or "", "limit": cfg["limit"], "lang": cfg.get("lang"), "year_from": cfg.get("year_from"), "year_to": cfg.get("year_to"), } elif stype == "cyberleninka": from cyberleninka import CyberLeninkaParser as P fetch_kwargs = {"query": cfg.get("query") or "", "limit": cfg["limit"]} elif stype == "arxiv": from arxiv import ArxivParser as P fetch_kwargs = { "query": cfg.get("query") or "", "limit": cfg["limit"], "year_from": cfg.get("year_from"), } else: raise ValueError(f"неизвестный тип источника: {stype}") parser = P() # fetch+transform без записи в JSONL — работаем in-memory raw_docs = parser.fetch(**fetch_kwargs) # Эмбеддинги диспатчим пачками, а не по одному документу: так worker-gpu # кодирует батч разом и переписывает FAISS-индекс на диск раз в N добавлений, # а не на каждый документ. embed_batch: list[int] = [] batch_size = settings.EMBED_BATCH_SIZE def _flush_embed() -> None: if embed_batch: celery_app.send_task( "gpu.embed_documents", args=[list(embed_batch)], queue="queue.gpu" ) embed_batch.clear() for raw in raw_docs: try: doc = parser.transform(raw) if not (doc and doc.get("title") and doc.get("ext_id")): continue result = add_document(doc, dispatch_embed=False) if result.get("status") == "indexed": added += 1 embed_batch.append(result["doc_id"]) if len(embed_batch) >= batch_size: _flush_embed() except Exception as e: logger.warning(f"run_parser: ошибка документа: {e}") _flush_embed() except Exception as e: error_msg = str(e)[:500] logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True) with db_session() as session: src = session.get(ParseSource, source_id) if src: src.last_status = "error" if error_msg else "done" src.last_error = error_msg src.docs_added = (src.docs_added or 0) + added src.last_run_at = datetime.now(timezone.utc) session.commit() return {"status": "error" if error_msg else "done", "added": added, "error": error_msg}