feat(admin): шкала загрузки источников, отладка и загрузка работ в корпус
All checks were successful
Deploy / test (push) Successful in 3m56s
Deploy / deploy (push) Successful in 31s

Заливка корпуса была чёрным ящиком: у источника только last_status
(idle/running/done/error), без «сколько из скольки», без причины падения и без
способа остановить начатое. Теперь каждый запуск создаёт строку parse_runs,
куда воркер раз в ~2с пишет стадию, счётчики и журнал событий.

Админка:
- шкала загрузки у каждого источника (0→50% выборка, 50→100% индексация),
  раскрытая строка — журнал прогона по шагам с таймингами;
- «Запустить всё» / «Остановить всё» и остановка по одному источнику
  (кооперативная отмена: воркер останавливается сам, не рвя запись в базу);
- пакетное добавление источников (тип + список тем), тип pmc в форме;
- загрузка PDF/DOCX/TXT прямо в базу сравнения (index.ingest_upload);
- страница «Отладка»: воркеры Celery и их текущие таски, очереди RabbitMQ,
  покрытие корпуса эмбеддингами, зависшие и упавшие прогоны, конфиг бэкендов.

Защита от краш-лупа по consumer_timeout RabbitMQ (docs/DR-HA.md §6), без неё
массовый запуск 170+ источников гарантированно ронял воркер:
- PARSER_TIME_BUDGET_S (1500с) — прогон закругляется сам и помечается partial;
- worker_prefetch_multiplier=1 — таймаут считается от ДОСТАВКИ сообщения, и с
  дефолтным префетчем очередь долгих run_parser убивала канал на задачах,
  которые ещё не начинались.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
jze9
2026-08-27 17:40:23 +05:00
parent 95d2903766
commit 78e27806b9
29 changed files with 2118 additions and 168 deletions

View File

@@ -26,5 +26,12 @@ celery_app.conf.update(
},
task_acks_late=True,
task_reject_on_worker_lost=True,
# Один неподтверждённый месседж на процесс пула. При acks_late=True всё
# предвыбранное висит unacked, а RabbitMQ рвёт канал по consumer_timeout
# (1800с) с момента ДОСТАВКИ, а не начала выполнения. С дефолтным префетчем
# (4×concurrency) массовая заливка — сотни долгих run_parser в очереди —
# гарантированно роняет воркер на задачах, которые ещё даже не начинались,
# и он уходит в краш-луп на передоставленных сообщениях (docs/DR-HA.md §6).
worker_prefetch_multiplier=1,
result_expires=86400,
)

View File

@@ -59,6 +59,13 @@ class Settings(BaseSettings):
FULL_TEXT_MIN_CHARS: int = 500 # Минимум символов, иначе считаем извлечение неудачным
EMBED_BATCH_SIZE: int = 64 # Размер пачки документов для диспатча эмбеддингов
# Бюджет времени одного прогона index.run_parser. RabbitMQ рвёт канал
# consumer'а, не сделавшего ack за consumer_timeout (по умолчанию 1800с):
# воркер падает, сообщение передоставляется, таск начинается заново —
# бесконечный краш-луп (docs/DR-HA.md §6). Прогон закругляется раньше и
# честно помечается partial: остаток дозаливается повторным запуском.
PARSER_TIME_BUDGET_S: float = 1500.0
# Автоматически добавлять проверенные работы студентов в базу для сравнения
# (как в коммерческих системах — Антиплагиат.ру/Turnitin ловят списывание у
# предыдущих потоков именно так). Выключено по умолчанию: без ручной

View File

@@ -2,7 +2,7 @@
from datetime import datetime
from sqlalchemy import JSON, BigInteger, ForeignKey, Integer, String, Text, func
from sqlalchemy import JSON, BigInteger, Boolean, ForeignKey, Integer, String, Text, func
from sqlalchemy.orm import DeclarativeBase, Mapped, mapped_column
@@ -77,9 +77,34 @@ class ParseSource(Base):
last_error: Mapped[str | None] = mapped_column(Text, nullable=True)
last_run_at: Mapped[datetime | None] = mapped_column(nullable=True)
docs_added: Mapped[int] = mapped_column(default=0)
last_run_id: Mapped[int | None] = mapped_column(nullable=True)
created_at: Mapped[datetime] = mapped_column(server_default=func.now())
class ParseRun(Base):
"""Прогон парсинга: счётчики прогресса и журнал (пишет run_parser)."""
__tablename__ = "parse_runs"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
source_id: Mapped[int] = mapped_column(ForeignKey("parse_sources.id"))
celery_task_id: Mapped[str | None] = mapped_column(String(64), nullable=True)
status: Mapped[str] = mapped_column(String(20), default="queued")
stage: Mapped[str] = mapped_column(String(20), default="queued")
target: Mapped[int] = mapped_column(default=0)
fetched: Mapped[int] = mapped_column(default=0)
processed: Mapped[int] = mapped_column(default=0)
added: Mapped[int] = mapped_column(default=0)
duplicates: Mapped[int] = mapped_column(default=0)
skipped: Mapped[int] = mapped_column(default=0)
failed: Mapped[int] = mapped_column(default=0)
cancel_requested: Mapped[bool] = mapped_column(Boolean, default=False)
error: Mapped[str | None] = mapped_column(Text, nullable=True)
log: Mapped[list | None] = mapped_column(JSON, default=list)
started_at: Mapped[datetime] = mapped_column(server_default=func.now())
heartbeat_at: Mapped[datetime | None] = mapped_column(nullable=True)
finished_at: Mapped[datetime | None] = mapped_column(nullable=True)
class StagedWork(Base):
__tablename__ = "staged_works"
id: Mapped[int] = mapped_column(Integer, primary_key=True)
@@ -100,4 +125,13 @@ class StagedWork(Base):
created_at: Mapped[datetime] = mapped_column(server_default=func.now())
__all__ = ["Base", "User", "Task", "Document", "Fingerprint", "ParseSource", "StagedWork"]
__all__ = [
"Base",
"User",
"Task",
"Document",
"Fingerprint",
"ParseSource",
"ParseRun",
"StagedWork",
]

View File

@@ -0,0 +1,124 @@
"""Счётчики прогресса прогона парсинга — чистая логика, без БД и Celery.
Отделено от tasks/index.py по той же причине, что и staging.py: здесь живут
решения, которые легко ломаются молча (как часто писать в БД, когда пора
закругляться по бюджету времени, сколько строк журнала хранить) — их надо
тестировать изолированно, без RabbitMQ/PostgreSQL.
Бюджет времени — не оптимизация, а защита: RabbitMQ рвёт канал consumer'а,
не сделавшего ack за `consumer_timeout` (по умолчанию 1800с), Celery-процесс
падает, сообщение передоставляется и таск начинается заново — бесконечный
краш-луп, который уже дважды съедал весь пул воркера (см. docs/DR-HA.md §6).
Поэтому прогон сам останавливается раньше дедлайна и честно помечается
`partial`, а не доводит воркер до падения.
"""
import time
from collections.abc import Callable
from datetime import UTC, datetime
from typing import Any
# Значения по умолчанию: бюджет заметно меньше consumer_timeout (1800с),
# чтобы после остановки успеть дописать статус в БД и сдаться брокеру.
DEFAULT_BUDGET_S = 1500.0
DEFAULT_FLUSH_INTERVAL_S = 2.0
LOG_TAIL_LIMIT = 200
class RunProgress:
"""Состояние одного прогона: счётчики, журнал, тайминги.
Воркер дёргает `count_*`/`log`, а по `should_flush()` сбрасывает
`snapshot()` в таблицу parse_runs.
"""
def __init__(
self,
target: int,
budget_s: float = DEFAULT_BUDGET_S,
flush_interval_s: float = DEFAULT_FLUSH_INTERVAL_S,
log_limit: int = LOG_TAIL_LIMIT,
clock: Callable[[], float] = time.monotonic,
) -> None:
self.target = max(0, target)
self.budget_s = budget_s
self.flush_interval_s = flush_interval_s
self.log_limit = log_limit
self._clock = clock
self._started = clock()
self._last_flush = self._started
self.stage = "queued"
self.fetched = 0
self.processed = 0
self.added = 0
self.duplicates = 0
self.skipped = 0
self.failed = 0
self.entries: list[dict[str, Any]] = []
# ── тайминги ──────────────────────────────────────────────────────────────
@property
def elapsed(self) -> float:
return self._clock() - self._started
@property
def over_budget(self) -> bool:
"""Пора закругляться, пока брокер не оборвал канал."""
return self.elapsed >= self.budget_s
def should_flush(self) -> bool:
"""Прошло ли достаточно времени с прошлой записи в БД."""
return (self._clock() - self._last_flush) >= self.flush_interval_s
def mark_flushed(self) -> None:
self._last_flush = self._clock()
# ── счётчики ──────────────────────────────────────────────────────────────
def count_fetched(self, total: int) -> None:
"""Сколько сырых документов получено от источника (нарастающим итогом)."""
self.fetched = total
def count_result(self, status: str) -> None:
"""Учесть результат обработки одного документа.
indexed / duplicate / skipped (мусор от источника: нет title или
ext_id) считаются отдельно от failed — иначе панель отладки показывала
бы тысячи «ошибок» там, где источник просто отдал неполные записи.
"""
self.processed += 1
if status == "indexed":
self.added += 1
elif status == "duplicate":
self.duplicates += 1
elif status == "skipped":
self.skipped += 1
else:
self.failed += 1
# ── журнал ────────────────────────────────────────────────────────────────
def log(self, level: str, msg: str) -> None:
"""Добавить запись журнала, храня только хвост последних log_limit строк."""
self.entries.append({
"ts": datetime.now(UTC).isoformat(timespec="seconds"),
"elapsed": round(self.elapsed, 1),
"level": level,
"msg": msg[:500],
})
if len(self.entries) > self.log_limit:
del self.entries[: len(self.entries) - self.log_limit]
# ── выгрузка ──────────────────────────────────────────────────────────────
def snapshot(self) -> dict[str, Any]:
"""Поля для UPDATE parse_runs."""
return {
"stage": self.stage,
"target": self.target,
"fetched": self.fetched,
"processed": self.processed,
"added": self.added,
"duplicates": self.duplicates,
"skipped": self.skipped,
"failed": self.failed,
"log": list(self.entries),
}

View File

@@ -6,7 +6,7 @@ from pathlib import Path
from typing import Any
from celery.utils.log import get_task_logger
from sqlalchemy import func, select
from sqlalchemy import func, select, update
from sqlalchemy.exc import IntegrityError
from app.algorithms.minhash import add_to_lsh, find_similar
@@ -517,19 +517,83 @@ def _stage_work(
logger.warning(f"Не удалось добавить работу {task_id!r} в отстойник: {e}")
@celery_app.task(name="index.run_parser")
def run_parser(source_id: int) -> dict[str, Any]:
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 (BaseParser.run) и
задачу add_document для каждого полученного документа.
Переиспользует парсеры из 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 datetime import datetime
from app.models import ParseSource
from app.models import ParseRun, ParseSource
from app.progress import RunProgress
# Парсеры лежат в /parsers (скопированы в образ)
# Парсеры лежат в /parsers (bind-mount scripts/parsers, см. compose)
if "/parsers" not in sys.path:
sys.path.insert(0, "/parsers")
@@ -546,47 +610,54 @@ def run_parser(source_id: int) -> dict[str, Any]:
"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()
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}")
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, **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
parser = P()
# fetch+transform без записи в JSONL — работаем in-memory
raw_docs = parser.fetch(**fetch_kwargs)
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 добавлений,
@@ -602,31 +673,135 @@ def run_parser(source_id: int) -> dict[str, Any]:
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":
added += 1
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 = "error" if error_msg else "done"
src.last_status = status
src.last_error = error_msg
src.docs_added = (src.docs_added or 0) + added
src.docs_added = (src.docs_added or 0) + prog.added
src.last_run_at = datetime.now(UTC)
session.commit()
return {"status": "error" if error_msg else "done", "added": added, "error": error_msg}
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}

View File

@@ -0,0 +1,104 @@
"""Юнит-тесты счётчиков прогона парсинга (app.progress) — без БД и Celery."""
from app.progress import RunProgress
class FakeClock:
"""Управляемое время: тесты бюджета не должны ничего ждать по-настоящему."""
def __init__(self) -> None:
self.now = 0.0
def __call__(self) -> float:
return self.now
def advance(self, seconds: float) -> None:
self.now += seconds
def test_counts_split_by_result_status():
prog = RunProgress(target=10, clock=FakeClock())
for status in ("indexed", "indexed", "duplicate", "skipped", "boom"):
prog.count_result(status)
assert prog.processed == 5
assert prog.added == 2
assert prog.duplicates == 1
# Мусор от источника не должен смешиваться с реальными ошибками:
# иначе панель отладки показывает ошибки там, где их нет
assert prog.skipped == 1
assert prog.failed == 1
def test_over_budget_only_after_deadline():
clock = FakeClock()
prog = RunProgress(target=100, budget_s=1500, clock=clock)
assert not prog.over_budget
clock.advance(1499)
assert not prog.over_budget
clock.advance(2)
assert prog.over_budget
def test_flush_is_throttled_by_interval():
clock = FakeClock()
prog = RunProgress(target=100, flush_interval_s=2.0, clock=clock)
assert not prog.should_flush()
clock.advance(2.5)
assert prog.should_flush()
prog.mark_flushed()
assert not prog.should_flush()
clock.advance(2.5)
assert prog.should_flush()
def test_log_keeps_only_tail():
prog = RunProgress(target=1, log_limit=3, clock=FakeClock())
for i in range(10):
prog.log("info", f"строка {i}")
assert len(prog.entries) == 3
assert [e["msg"] for e in prog.entries] == ["строка 7", "строка 8", "строка 9"]
def test_log_entry_shape():
clock = FakeClock()
prog = RunProgress(target=1, clock=clock)
clock.advance(12.34)
prog.log("warning", "x" * 900)
entry = prog.entries[0]
assert entry["level"] == "warning"
assert entry["elapsed"] == 12.3
# Длинные сообщения режем: журнал целиком лежит в одной JSON-колонке
assert len(entry["msg"]) == 500
assert entry["ts"]
def test_snapshot_carries_all_counters():
prog = RunProgress(target=42, clock=FakeClock())
prog.stage = "index"
prog.count_fetched(30)
prog.count_result("indexed")
prog.log("info", "поехали")
snap = prog.snapshot()
assert snap["stage"] == "index"
assert snap["target"] == 42
assert snap["fetched"] == 30
assert snap["processed"] == 1
assert snap["added"] == 1
assert snap["log"][0]["msg"] == "поехали"
def test_snapshot_log_is_detached_copy():
"""Снимок уходит в БД как есть — последующие записи не должны его менять."""
prog = RunProgress(target=1, clock=FakeClock())
prog.log("info", "первая")
snap = prog.snapshot()
prog.log("info", "вторая")
assert len(snap["log"]) == 1