Две связанные проблемы, ломавшие качество детекции: 1. Уровень 1 (Winnowing) репортил весь документ одним совпадением на позиции 0..длина — в отчёте нельзя было понять, ГДЕ плагиат. Теперь winnow'им каждый фрагмент и находим источник + позицию для каждого, так отчёт показывает "символы A-B скопированы из источника X". 2. MAX_FINGERPRINTS_PER_DOC=500 обрезал отпечатки статьи (у типичной статьи их ~3600), причём произвольную выборку. Из-за этого скопированный фрагмент почти не разделял отпечатки с источником (проверено: хранилось ~14%, фрагмент находил 17% своих хэшей → не срабатывало). Winnowing рассчитан на ПОЛНОЕ хранение; поднял лимит до 20000. Также overall_similarity считается по числу уникальных помеченных позиций, а не совпадений — фрагмент, совпавший с несколькими источниками, больше не раздувает процент выше 100. Проверено end-to-end: документ с дословной вставкой из статьи корпуса → вставка локализована (chars 1027-2385 = 100%, источник атрибутирован); чисто оригинальный текст → 0% (нет ложных срабатываний). Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
590 lines
23 KiB
Python
590 lines
23 KiB
Python
"""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 (пофрагментно, с локализацией) ───
|
||
# 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)
|
||
|
||
|
||
@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}
|