Разбор итогов массовой заливки показал в отладке «среднюю длительность прогона» в 2.8 часа там, где парсинг занимал 50 секунд: started_at пишется в момент постановки в очередь, а очередь из 173 источников разбирается часами. Теперь момент реального старта пишется отдельно (run_started_at, миграция 006), а схема отдаёт обе величины — сколько ждал очереди и сколько работал. Плюс scripts/ops/reembed_missing.py: документ попадает в корпус сразу, а вектор для L3 считает отдельная задача; когда worker-gpu или Ollama недоступны, эти задачи теряются и документ остаётся невидимым для семантического поиска. Скрипт находит faiss_id IS NULL и переотправляет пачками (dry-run по умолчанию) — сейчас таких 23 101 из 177 147. Документация: актуальные цифры корпуса, дубли при повторном прогоне, лимит OpenAlex, и главное — гипервизор .254 зафиксирован в DR-HA как самая широкая единая точка отказа (брокер, эмбеддинги, секреты и прокси на одном железе; подтверждено аварией 28.08). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
816 lines
35 KiB
Python
816 lines
35 KiB
Python
"""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, update
|
||
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, 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:
|
||
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.
|
||
# Если 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}
|
||
|
||
|
||
@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}")
|
||
|
||
|
||
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"),
|
||
}
|
||
else:
|
||
raise ValueError(f"неизвестный тип источника: {source_type}")
|
||
|
||
return P(), kwargs
|
||
|
||
|
||
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,
|
||
}
|
||
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)
|
||
|
||
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()
|
||
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}
|