Files
anti-plagiarism/services/worker-indexer/app/progress.py
jze9 78e27806b9
All checks were successful
Deploy / test (push) Successful in 3m56s
Deploy / deploy (push) Successful in 31s
feat(admin): шкала загрузки источников, отладка и загрузка работ в корпус
Заливка корпуса была чёрным ящиком: у источника только 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>
2026-08-27 17:40:23 +05:00

125 lines
6.0 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""Счётчики прогресса прогона парсинга — чистая логика, без БД и Celery.
Отделено от 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),
}