Заливка корпуса была чёрным ящиком: у источника только 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>
69 lines
3.5 KiB
Python
69 lines
3.5 KiB
Python
"""Журнал прогонов парсинга (parse_runs) + ссылка на последний прогон источника.
|
||
|
||
Revision ID: 005
|
||
Revises: 004
|
||
Create Date: 2026-08-27
|
||
|
||
Прогресс заливки раньше был не виден: у источника был только last_status
|
||
(idle/running/done/error) без «сколько из скольки». Таблица parse_runs хранит
|
||
счётчики и структурный журнал каждого прогона — по ним админка рисует шкалу
|
||
загрузки и показывает, на чём именно споткнулся источник.
|
||
"""
|
||
|
||
import sqlalchemy as sa
|
||
|
||
from alembic import op
|
||
|
||
revision = "005"
|
||
down_revision = "004"
|
||
branch_labels = None
|
||
depends_on = None
|
||
|
||
|
||
def upgrade() -> None:
|
||
op.create_table(
|
||
"parse_runs",
|
||
sa.Column("id", sa.Integer(), primary_key=True, autoincrement=True),
|
||
sa.Column(
|
||
"source_id",
|
||
sa.Integer(),
|
||
sa.ForeignKey("parse_sources.id", ondelete="CASCADE"),
|
||
nullable=False,
|
||
),
|
||
sa.Column("celery_task_id", sa.String(64), nullable=True),
|
||
# queued / running / done / partial / error / cancelled
|
||
sa.Column("status", sa.String(20), nullable=False, server_default="queued"),
|
||
# queued / fetch / index / finished
|
||
sa.Column("stage", sa.String(20), nullable=False, server_default="queued"),
|
||
sa.Column("target", sa.Integer(), nullable=False, server_default="0"),
|
||
sa.Column("fetched", sa.Integer(), nullable=False, server_default="0"),
|
||
sa.Column("processed", sa.Integer(), nullable=False, server_default="0"),
|
||
sa.Column("added", sa.Integer(), nullable=False, server_default="0"),
|
||
sa.Column("duplicates", sa.Integer(), nullable=False, server_default="0"),
|
||
# Мусор от источника (нет title/ext_id) — считаем отдельно от ошибок
|
||
sa.Column("skipped", sa.Integer(), nullable=False, server_default="0"),
|
||
sa.Column("failed", sa.Integer(), nullable=False, server_default="0"),
|
||
# Флаг кооперативной отмены: воркер проверяет его на каждом тике прогресса
|
||
sa.Column("cancel_requested", sa.Boolean(), nullable=False, server_default=sa.false()),
|
||
sa.Column("error", sa.Text(), nullable=True),
|
||
sa.Column("log", sa.JSON(), nullable=True),
|
||
sa.Column("started_at", sa.DateTime(), server_default=sa.func.now(), nullable=False),
|
||
sa.Column("heartbeat_at", sa.DateTime(), nullable=True),
|
||
sa.Column("finished_at", sa.DateTime(), nullable=True),
|
||
)
|
||
op.create_index("ix_parse_runs_source_id", "parse_runs", ["source_id"])
|
||
op.create_index("ix_parse_runs_status", "parse_runs", ["status"])
|
||
op.create_index("ix_parse_runs_started_at", "parse_runs", ["started_at"])
|
||
|
||
# Последний (он же текущий, пока идёт) прогон источника — чтобы список
|
||
# источников отдавался одним джойном, без подзапроса «максимальный id».
|
||
op.add_column("parse_sources", sa.Column("last_run_id", sa.Integer(), nullable=True))
|
||
|
||
|
||
def downgrade() -> None:
|
||
op.drop_column("parse_sources", "last_run_id")
|
||
op.drop_index("ix_parse_runs_started_at", table_name="parse_runs")
|
||
op.drop_index("ix_parse_runs_status", table_name="parse_runs")
|
||
op.drop_index("ix_parse_runs_source_id", table_name="parse_runs")
|
||
op.drop_table("parse_runs")
|