Заливка корпуса была чёрным ящиком: у источника только 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>
125 lines
6.0 KiB
Python
125 lines
6.0 KiB
Python
"""Счётчики прогресса прогона парсинга — чистая логика, без БД и 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),
|
||
}
|