Files
anti-plagiarism/services/worker-indexer/app/tasks/index.py
jze9 b471ec767a
Some checks failed
Deploy / test (push) Failing after 4m15s
Deploy / deploy (push) Has been skipped
feat: две доработки в духе коммерческих систем — цитаты и авто-корпус
Ответ на "неужели больше нет" / "мы ведь можем это исправить": закрывает
две дыры относительно Антиплагиат.ру/Turnitin, о которых договорились.

1. Различение цитаты и голого плагиата (app.scoring.is_cited, worker-gpu):
   эвристика — фрагмент считается процитированным, если обрамлён кавычками
   («…», "…") либо сразу за ним (в пределах ~150 симв.) идёт скобочная ссылка
   с годом: [Иванов, 2023], (Smith, 2020) — совпадает и с нашим же форматом
   ГОСТ 7.0.5. aggregate_results теперь принимает full_text, размечает
   match["cited"] и считает uncited_similarity (доля БЕЗ похожих на цитаты —
   ближе к тому, что коммерческие системы называют "% некорректных
   заимствований") отдельно от overall_similarity (как было, для совместимости).
   Фронтенд: бейдж "Цитата" на совпадении + строка с разбивкой в отчёте
   (аддитивные опциональные поля в типах — старые задачи не ломаются).

2. Автопополнение корпуса проверенными работами (как у коммерческих систем —
   так ловится списывание у предыдущих потоков). Раньше загруженная на проверку
   работа складывалась в StagedWork и ждала РУЧНОГО одобрения админом — де-факто
   не пополняла базу для сравнения. Теперь index.auto_approve_submission
   (диспатчится из gpu.check_plagiarism ПОСЛЕ сохранения результата — чтобы
   работа не сматчилась сама с собой) добавляет её в documents автоматически,
   под настройкой AUTO_APPROVE_SUBMISSIONS (default True). Ручное
   approve/reject в админке остаётся рабочим (идемпотентно — auto-approve
   пропускает уже не-pending записи), пригодится при AUTO_APPROVE=False.
   Конвертация StagedWork→doc_data вынесена в чистый app/staging.py (без
   Celery/SQLAlchemy/MinIO) — тестируется изолированно, идентична ручному
   пути в admin.py (POST /admin/staging/{id}/approve).

Тестов добавлено 16 (scoring 8→18, новый staging.py — 6). Оба mypy-гейта
расширены (scoring.py, staging.py). Тестов всего: 112 (было 82 в последнем
подсчёте README — таблица давно отставала, заодно поправил на актуальные цифры
по всем сервисам, включая забытый в прошлый раз Qdrant).

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-08-24 15:09:43 +05:00

612 lines
25 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Celery задачи индексации документов и проверки плагиата (уровни 1-2)."""
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
from sqlalchemy.exc import IntegrityError
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
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)
# Источник, разделяющий больше всего отпечатков с этим фрагментом
row = session.execute(
select(
Fingerprint.doc_id,
func.count(Fingerprint.id).label("cnt"),
)
.where(Fingerprint.hash_value.in_(frag_hashes))
.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:
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) 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:
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) from exc
# Пересчитать 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)}
@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}")
@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
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"),
"open_access_only": settings.INGEST_OPEN_ACCESS_ONLY,
}
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"),
}
elif stype == "pmc":
from pmc import PMCParser as P
fetch_kwargs = {
"query": cfg.get("query") or "",
"limit": cfg["limit"],
"year_from": cfg.get("year_from"),
"year_to": cfg.get("year_to"),
}
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(UTC)
session.commit()
return {"status": "error" if error_msg else "done", "added": added, "error": error_msg}