Первый прогон нового конвейера показал слабое место: листинг бакета шёл с начала, и заливка часами перемалывала уже существующие статьи как дубли — из 300 полученных 300 оказались дублями. S3 умеет start-after, и стартовый ключ выводится прямо из базы: ext_id вида `pmc:PMC10000000` соответствует ключу `PMC10000000.1`, бакет отдаётся лексикографически. Теперь при отсутствии сохранённого токена (первый прогон после ручных заливок) листинг начинается после максимального залитого. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
986 lines
43 KiB
Python
986 lines
43 KiB
Python
"""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 _last_pmc_key() -> str | None:
|
||
"""Ключ бакета после последней залитой статьи PMC — старт листинга.
|
||
|
||
Бакет отдаётся в лексикографическом порядке ключей, и ext_id вида
|
||
`pmc:PMC10000000` соответствует ключу `PMC10000000.1`. Берём максимальный
|
||
залитый и продолжаем после него.
|
||
"""
|
||
from sqlalchemy import text as sql_text
|
||
|
||
with db_session() as session:
|
||
row = session.execute(sql_text(
|
||
"SELECT max(ext_id) FROM documents WHERE ext_id LIKE 'pmc:PMC%'"
|
||
)).first()
|
||
if not row or not row[0]:
|
||
return None
|
||
return row[0].split(":", 1)[1] + ".1"
|
||
|
||
|
||
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"),
|
||
# Токена ещё нет (первый прогон после ручных заливок) — начинаем
|
||
# после последней уже залитой статьи, иначе листинг часами
|
||
# перемалывает существующие как дубли
|
||
"start_after": None if cfg.get("resume_token") else _last_pmc_key(),
|
||
}
|
||
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}
|