"""Celery задачи индексации документов и проверки плагиата (уровни 1-2).""" import contextlib import io from datetime import UTC, datetime from pathlib import Path from typing import Any from celery.utils.log import get_task_logger from sqlalchemy import func, select, update from sqlalchemy.exc import IntegrityError from app.algorithms.minhash import add_to_lsh, find_similar from app.algorithms.winnowing import sample_evenly, winnow, winnow_ordered from app.celery_app import celery_app from app.config import settings from app.db import db_session, get_minio, refund_plagiarism_quota, 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 from app.fragments import split_into_fragments from app.staging import staged_work_to_doc_data logger = get_task_logger(__name__) @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 (пофрагментно, с локализацией) ─── # Winnow'им КАЖДЫЙ фрагмент отдельно и ищем, с каким источником и на # какой позиции он совпадает. Так в отчёте видно не «документ похож на X», # а «фрагмент на символах A–B скопирован из источника X». level1_matches: list[dict] = [] from app.models import Document, Fingerprint with db_session() as session: doc_cache: dict[int, Document] = {} for frag in fragments: frag_fp = winnow(frag["text"]) if not frag_fp: continue frag_hashes = list(frag_fp) # Источник, разделяющий больше всего отпечатков с этим фрагментом. # user_submission исключены: это чужие непубличные загрузки (и # свои же прошлые прогоны того же файла) — сравнение с ними даёт # ложные 100%-совпадения, а не реальный плагиат из источника. row = session.execute( select( Fingerprint.doc_id, func.count(Fingerprint.id).label("cnt"), ) .join(Document, Document.id == Fingerprint.doc_id) .where( Fingerprint.hash_value.in_(frag_hashes), Document.source != "user_submission", ) .group_by(Fingerprint.doc_id) .order_by(func.count(Fingerprint.id).desc()) .limit(1) ).first() if row is None: continue doc_id, cnt = row # Доля отпечатков фрагмента, найденных в источнике similarity = cnt / len(frag_hashes) * 100 if similarity < settings.EXACT_FRAGMENT_THRESHOLD: continue doc = doc_cache.get(doc_id) if doc is None: doc = session.get(Document, doc_id) if doc is None: continue doc_cache[doc_id] = doc level1_matches.append({ "fragment": frag["text"][:300], "position_start": frag["start"], "position_end": frag["end"], "similarity": round(min(similarity, 100.0), 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 or doc.source == "user_submission": 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 ValueError as exc: # Детерминированная ошибка (битый файл, не тот формат, пустой текст) — # повтор не поможет: результат будет тем же на 2-й и 3-й попытке. Не # ретраим (быстрее фидбек юзеру) и возвращаем месячную квоту — реальная # проверка так и не началась, списывать не за что. logger.warning(f"Задача {task_id!r}: не удалось обработать файл: {exc}") update_task_status(task_id, "failed", str(exc)) refund_plagiarism_quota(task_id) return {"task_id": task_id, "status": "failed", "error": str(exc)} 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) from exc @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) try: session.flush() except IntegrityError: # Гонка: другой воркер параллельно вставил этот ext_id между нашей # проверкой existing и flush. Уникальный индекс отработал — это дубль. session.rollback() existing = session.execute( select(Document).where(Document.ext_id == ext_id) ).scalar_one_or_none() return {"status": "duplicate", "doc_id": existing.id if existing else None} doc_id = doc.id # Вычислить fingerprints text = doc_data.get("full_text") or doc_data.get("abstract", "") or "" if text: fingerprints_to_add = sample_evenly( winnow_ordered(text), 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. # Если full_text уже есть (например, OCR-фрагмент из ответа поиска # КиберЛенинки) — не дублируем запрос: страница статьи там HTML, а не PDF, # enrich_full_text гарантированно ничего не найдёт, только зря нагрузит # источник (и получит 503 при массовой заливке). if settings.FETCH_FULL_TEXT and doc_data.get("url") and not doc_data.get("full_text"): celery_app.send_task( "index.enrich_full_text", args=[doc_id, doc_data["url"]], queue="queue.index", ) return {"status": "indexed", "doc_id": doc_id} def store_full_text(doc_id: int, text: str) -> dict[str, Any]: """Сохранить полный текст документа и переиндексировать его по нему. Общая часть для всех путей получения полного текста: скачанный PDF (enrich_full_text) и текст, пришедший прямо из API источника (бэкфилл PMC — scripts/ops/). Провизорные fingerprints, посчитанные по аннотации, заменяются на посчитанные по полному тексту — ради этого всё и делается: L1 начинает видеть тело статьи, а не только её краткое описание. Args: doc_id: документ в PostgreSQL text: полный текст статьи Returns: dict со статусом, объёмом текста и числом отпечатков """ from sqlalchemy import delete from app.models import Document, Fingerprint 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", ) hashes = sample_evenly(winnow_ordered(text), settings.MAX_FINGERPRINTS_PER_DOC) 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 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 (L2) тоже должен считаться по полному тексту add_to_lsh(f"doc:{doc_id}", text) logger.info(f"store_full_text: doc={doc_id} {len(text)} симв., fingerprints={len(hashes)}") return {"status": "ok", "doc_id": doc_id, "chars": len(text), "fingerprints": len(hashes)} @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 text = fetch_full_text(url) if not text: return {"status": "no_fulltext", "doc_id": doc_id} try: return store_full_text(doc_id, text) except Exception as exc: logger.error(f"enrich_full_text: не удалось сохранить текст для {doc_id}: {exc}") raise self.retry(exc=exc, countdown=120) from exc @celery_app.task( name="index.auto_approve_submission", bind=True, max_retries=2, default_retry_delay=60, ) def auto_approve_submission(self, task_id: str) -> dict[str, Any]: """Автоматически добавить проверенную работу студента в базу для сравнения. Вызывается из gpu.check_plagiarism ПОСЛЕ сохранения результата проверки — так работа не сматчится сама с собой. Идемпотентно: пропускает, если StagedWork не 'pending' (уже одобрена/отклонена вручную или повторный вызов) или если AUTO_APPROVE_SUBMISSIONS выключен настройкой — тогда работа так и остаётся ждать ручного решения в админке. """ from app.models import StagedWork if not settings.AUTO_APPROVE_SUBMISSIONS: return {"status": "skipped", "reason": "auto-approve выключен настройкой"} with db_session() as session: sw = session.execute( select(StagedWork).where(StagedWork.task_id == task_id) ).scalar_one_or_none() if sw is None: return {"status": "skipped", "reason": "StagedWork не найден"} if sw.status != "pending": return {"status": "skipped", "reason": f"уже {sw.status}"} full_text = "" if sw.text_key: try: minio = get_minio() obj = minio.get_object(settings.MINIO_BUCKET_STAGING, sw.text_key) full_text = obj.read().decode("utf-8", errors="replace") except Exception as e: logger.warning(f"auto_approve_submission: не прочитан текст {task_id!r}: {e}") return {"status": "error", "reason": str(e)} doc_data = staged_work_to_doc_data(sw, full_text) sw.status = "approved" sw.reviewed_by = None # None = одобрено автоматически, не человеком sw.reviewed_at = datetime.now(UTC) session.commit() result = add_document(doc_data, dispatch_embed=True) logger.info(f"auto_approve_submission: задача {task_id!r} → {result}") return result 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}") def _parser_for(source_type: str, cfg: dict[str, Any]) -> tuple[Any, dict[str, Any]]: """Инстанс парсера и совместимые с его fetch() аргументы по типу источника.""" limit = cfg["limit"] query = cfg.get("query") or "" if source_type == "openalex": from openalex import OpenAlexParser as P kwargs = { "query": query, "limit": limit, "lang": cfg.get("lang"), "year_from": cfg.get("year_from"), "year_to": cfg.get("year_to"), "open_access_only": settings.INGEST_OPEN_ACCESS_ONLY, } elif source_type == "cyberleninka": from cyberleninka import CyberLeninkaParser as P kwargs = {"query": query, "limit": limit} elif source_type == "arxiv": from arxiv import ArxivParser as P kwargs = {"query": query, "limit": limit, "year_from": cfg.get("year_from")} elif source_type == "pmc": from pmc import PMCParser as P kwargs = { "query": query, "limit": limit, "year_from": cfg.get("year_from"), "year_to": cfg.get("year_to"), } elif source_type == "wikipedia_ru": # Массовый источник: дамп читается с позиции прошлого прогона from wikipedia_ru import WikipediaRuParser as P kwargs = { "limit": limit, "dump_path": query or None, # query = путь к дампу, если задан "skip": int(cfg.get("resume_token") or 0), } elif source_type == "pmc_bulk": from pmc_bulk import PMCBulkParser as P kwargs = {"limit": limit, "resume_token": cfg.get("resume_token")} else: raise ValueError(f"неизвестный тип источника: {source_type}") return P(), kwargs def _ingest_one_by_one(parser: Any, raw_docs: Any, prog: Any, run_id: int) -> tuple[bool, bool]: """Обычный путь: документ за документом через add_document. Подходит источникам, отдающим сотни статей: работает ORM-логика с дедупликацией, индексацией в Elasticsearch и обогащением полным текстом. Returns: (отменён, исчерпан бюджет времени) """ cancelled = exhausted = False prog.stage = "index" prog.count_fetched(len(raw_docs)) prog.log("info", f"получено {len(raw_docs)} документов, индексация") _write_run(run_id, **prog.snapshot()) # Эмбеддинги диспатчим пачками, а не по одному документу: так 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: if not cancelled and prog.over_budget: exhausted = True prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на индексации") if cancelled or exhausted: break try: doc = parser.transform(raw) if not (doc and doc.get("title") and doc.get("ext_id")): prog.count_result("skipped") continue result = add_document(doc, dispatch_embed=False) prog.count_result(result.get("status", "error")) if result.get("status") == "indexed": embed_batch.append(result["doc_id"]) if len(embed_batch) >= batch_size: _flush_embed() except Exception as e: prog.count_result("error") prog.log("warning", f"документ пропущен: {e}") logger.warning(f"run_parser: ошибка документа: {e}") if prog.should_flush(): cancelled = _write_run(run_id, **prog.snapshot()) prog.mark_flushed() if cancelled: prog.log("warning", "отмена по запросу из админки") _flush_embed() return cancelled, exhausted def _ingest_bulk( parser: Any, raw_docs: Any, prog: Any, run_id: int, source_id: int ) -> tuple[bool, bool]: """Массовый путь: статьи приходят потоком и пишутся пачками через COPY. Отличия от обычного пути и почему они нужны: - парсер отдаёт генератор, поэтому счётчик «получено» растёт по ходу дела, а не известен заранее; - запись идёт пакетно (bulk_writer), иначе миллионы отпечатков занимают недели вместо суток; - позиция продолжения сохраняется в источнике, чтобы следующий прогон начинал с места остановки, а не перечитывал дамп с начала. Эмбеддинги здесь не диспатчатся: они считаются заметно медленнее заливки и делали бы её узким местом. Векторы досчитываются отдельно — `scripts/ops/reembed_missing.py`. Returns: (отменён, исчерпан бюджет времени) """ from app.bulk_writer import connect, write_batch from app.models import ParseSource cancelled = exhausted = False prog.stage = "index" prog.log("info", "массовый источник: запись пачками") _write_run(run_id, **prog.snapshot()) conn = connect() batch: list[dict[str, Any]] = [] resume_token: str | None = None batch_size = settings.BULK_WRITE_BATCH def save_resume() -> None: """Запомнить позицию в источнике — с неё продолжит следующий прогон.""" if resume_token is None: return with db_session() as session: src = session.get(ParseSource, source_id) if src: src.resume_token = str(resume_token) session.commit() def flush() -> bool: """Записать накопленное; True — попросили остановиться.""" nonlocal conn, batch if not batch: return False conn, counts = write_batch(conn, batch, settings.MAX_FINGERPRINTS_PER_DOC) for _ in range(counts["added"]): prog.count_result("indexed") for _ in range(counts["duplicates"]): prog.count_result("duplicate") batch = [] save_resume() stop = _write_run(run_id, **prog.snapshot()) prog.mark_flushed() return stop try: for raw in raw_docs: if prog.over_budget: exhausted = True prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан") break try: doc = parser.transform(raw) except Exception as e: prog.count_result("error") logger.warning(f"run_parser bulk: ошибка документа: {e}") continue if not (doc and doc.get("ext_id") and doc.get("text")): prog.count_result("skipped") continue # Позиция продолжения: у Википедии — номер статьи в дампе, # у PMC — токен страницы бакета resume_token = doc.get("dump_position") or doc.get("resume_token") or resume_token batch.append(doc) prog.count_fetched(prog.fetched + 1) if len(batch) >= batch_size and flush(): cancelled = True prog.log("warning", "отмена по запросу из админки") break if not cancelled and flush(): cancelled = True finally: with contextlib.suppress(Exception): conn.close() return cancelled, exhausted def _write_run(run_id: int, **fields: Any) -> bool: """Записать прогресс прогона и вернуть True, если запрошена отмена. Отмена кооперативная: админка ставит cancel_requested, воркер узнаёт о ней на ближайшем тике прогресса и останавливается сам. Так прогон завершается с осмысленным статусом и не оставляет источник висеть в running (что было бы при жёстком revoke уже начатого таска). """ from app.models import ParseRun with db_session() as session: row = session.execute( update(ParseRun) .where(ParseRun.id == run_id) .values(heartbeat_at=datetime.now(UTC), **fields) .returning(ParseRun.cancel_requested) ).first() return bool(row and row[0]) @celery_app.task(name="index.run_parser", bind=True) def run_parser(self, source_id: int, run_id: int | None = None) -> dict[str, Any]: """Запустить парсинг источника и наполнить базу документов. Переиспользует парсеры из scripts/parsers и задачу add_document для каждого полученного документа. Ход заливки пишется в parse_runs (шкала загрузки и журнал в админке); прогон можно отменить из админки и он сам закругляется по бюджету времени, не доводя воркер до падения по consumer_timeout RabbitMQ (см. app/progress.py). Args: source_id: ID источника в parse_sources run_id: ID заранее созданной строки parse_runs (её создаёт админка, чтобы прогон был виден в очереди ещё до старта). Без него строка создаётся здесь — для запусков из скриптов. """ import sys from app.models import ParseRun, ParseSource from app.progress import RunProgress # Парсеры лежат в /parsers (bind-mount scripts/parsers, см. compose) 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, # Массовые источники продолжают с сохранённой позиции "resume_token": src.resume_token, } src.last_status = "running" src.last_error = None src.last_run_at = datetime.now(UTC) if run_id is None: run = ParseRun(source_id=source_id, status="running", stage="fetch") session.add(run) session.flush() run_id = run.id src.last_run_id = run_id session.commit() prog = RunProgress(target=cfg["limit"], budget_s=settings.PARSER_TIME_BUDGET_S) prog.stage = "fetch" prog.log("info", f"старт: {cfg['source_type']} q={cfg.get('query') or '—'} limit={cfg['limit']}") _write_run( run_id, status="running", celery_task_id=self.request.id, error=None, # Отдельно от started_at (постановка в очередь): при массовом запуске # между ними часы ожидания, и без этой отметки «длительность прогона» # в отладке показывала очередь, а не работу run_started_at=datetime.now(UTC), **prog.snapshot(), ) cancelled = False exhausted = False # остановлены бюджетом времени, а не концом выдачи error_msg = None try: parser, fetch_kwargs = _parser_for(cfg["source_type"], cfg) def on_fetch_progress(fetched: int) -> bool: """Тик прогресса выборки; False — парсеру пора остановиться.""" nonlocal cancelled, exhausted prog.count_fetched(fetched) if prog.over_budget: exhausted = True prog.log("warning", f"бюджет времени {prog.budget_s:.0f}с исчерпан на выборке") return False if prog.should_flush(): if _write_run(run_id, **prog.snapshot()): cancelled = True prog.log("warning", "отмена по запросу из админки") return False prog.mark_flushed() return True # fetch+transform без записи в JSONL — работаем in-memory raw_docs = parser.fetch(**fetch_kwargs, progress_cb=on_fetch_progress) if getattr(parser, "bulk", False): # Массовые источники (дамп Википедии, бакет PMC) отдают генератор: # миллионы статей нельзя ни держать в памяти, ни писать по одной # через ORM — для них отдельный путь с пакетной записью cancelled, exhausted = _ingest_bulk(parser, raw_docs, prog, run_id, source_id) else: cancelled, exhausted = _ingest_one_by_one(parser, raw_docs, prog, run_id) except Exception as e: error_msg = str(e)[:500] prog.log("error", f"прогон упал: {error_msg}") logger.error(f"run_parser source={source_id} ошибка: {e}", exc_info=True) if error_msg: status = "error" elif cancelled: status = "cancelled" elif exhausted: status = "partial" else: status = "done" prog.stage = "finished" prog.log( "info" if status in ("done", "partial") else "warning", f"итог: {status}, добавлено {prog.added}, дублей {prog.duplicates}, " f"ошибок {prog.failed}, за {prog.elapsed:.0f}с", ) _write_run( run_id, status=status, error=error_msg, finished_at=datetime.now(UTC), **prog.snapshot(), ) with db_session() as session: src = session.get(ParseSource, source_id) if src: src.last_status = status src.last_error = error_msg src.docs_added = (src.docs_added or 0) + prog.added src.last_run_at = datetime.now(UTC) session.commit() return { "status": status, "run_id": run_id, "added": prog.added, "duplicates": prog.duplicates, "failed": prog.failed, "error": error_msg, } @celery_app.task(name="index.ingest_upload", bind=True, max_retries=2, default_retry_delay=60) def ingest_upload( self, minio_key: str, filename: str, meta: dict[str, Any] | None = None, ) -> dict[str, Any]: """Добавить загруженный админом файл в базу документов (корпус для сравнения). Это не проверка плагиата: файл сразу становится источником, с которым сравниваются работы студентов. Текст извлекается тем же кодом, что и в extract_and_check, дальше — обычный add_document (дедуп, fingerprints, LSH, Elasticsearch, эмбеддинги). """ meta = meta or {} try: minio = get_minio() response = minio.get_object(settings.MINIO_BUCKET_DOCS, minio_key) file_data = response.read() response.close() response.release_conn() 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("не удалось извлечь текст") except ValueError as exc: logger.warning(f"ingest_upload {minio_key}: {exc}") return {"status": "failed", "minio_key": minio_key, "error": str(exc)} except Exception as exc: logger.error(f"ingest_upload {minio_key}: {exc}", exc_info=True) raise self.retry(exc=exc, countdown=60) from exc doc_data = { "source": meta.get("source") or "manual_upload", # Ключ MinIO уникален (UUID в имени) — годится как ext_id для дедупа "ext_id": f"upload:{minio_key}", "title": meta.get("title") or Path(filename).stem, "authors": meta.get("authors") or [], "year": meta.get("year"), "lang": meta.get("lang"), "abstract": text[:2000], "full_text": text, "minio_key": minio_key, } result = add_document(doc_data) logger.info( f"ingest_upload: {filename!r} → {result.get('status')} " f"(doc_id={result.get('doc_id')}, {len(text)} симв.)" ) return {"status": result.get("status"), "doc_id": result.get("doc_id"), "filename": filename}