diff --git a/scripts/ops/keep_ingesting.py b/scripts/ops/keep_ingesting.py index 8570274..f47f44a 100755 --- a/scripts/ops/keep_ingesting.py +++ b/scripts/ops/keep_ingesting.py @@ -15,11 +15,36 @@ Останавливается сам, когда источник исчерпан: парсер перестаёт отдавать статьи, прогон закрывается с нулём добавленных, и скрипт больше его не трогает (--stop-when-empty, по умолчанию включено). + +Перед подсчётом «занято ли» снимает зависшие прогоны: если worker-indexer упал +или перезапустился, прогон навсегда остаётся в статусе running с замолчавшим +heartbeat — без этой чистки сторож видит «занято» бесконечно и не запускает +вообще ничего, для всех источников сразу (порог тот же, что и в панели +отладки — services/api/app/api/admin.py, stale_cutoff). + +Молчащий heartbeat сам по себе ещё не значит «мёртв»: пачка COPY на большой +базе идёт десятки минут без единого тика. Поэтому прогон снимается, только +если его таска нет среди выполняющихся на воркерах queue.index; не ответил +хоть один такой воркер — не снимается ничего. Иначе сторож снимал живой +медленный прогон и тут же запускал его дубль по тем же статьям. """ import argparse BULK_TYPES = ("wikipedia_ru", "pmc_bulk") +STALE_MINUTES = 10 + + +def live_task_ids(celery_app) -> set[str] | None: + """ID тасков, выполняющихся на воркерах queue.index; None — кто-то из них не ответил.""" + queues = celery_app.control.inspect(timeout=10).active_queues() or {} + workers = [w for w, qs in queues.items() if any(q["name"] == "queue.index" for q in qs)] + if not workers: + return None + active = celery_app.control.inspect(destination=workers, timeout=10).active() + if not active or set(active) != set(workers): + return None + return {t["id"] for tasks in active.values() for t in tasks} def main() -> None: @@ -38,7 +63,44 @@ def main() -> None: types = tuple(t.strip() for t in args.types.split(",") if t.strip()) + live = live_task_ids(celery_app) + with db_session() as s: + if live is None: + print("воркеры queue.index не ответили — живость прогонов не проверить, не снимаю") + reaped = [] + else: + reaped = s.execute(text(f""" + UPDATE parse_runs SET status='error', stage='finished', + error='зависший прогон снят автоматически: таск не выполняется, heartbeat молчал дольше {STALE_MINUTES} минут', + finished_at=now() + WHERE status='running' + AND (heartbeat_at IS NULL OR heartbeat_at < now() - interval '{STALE_MINUTES} minutes') + AND (celery_task_id IS NULL OR NOT (celery_task_id = ANY(CAST(:live AS text[])))) + RETURNING id, source_id + """), {"live": sorted(live)}).all() + # queued без heartbeat — тот же зомби другого вида: send_task не дошёл + # или упал между двумя commit при постановке в очередь (celery_task_id + # так и остался пустым), воркер о такой задаче никогда не узнает + reaped += s.execute(text(f""" + UPDATE parse_runs SET status='error', stage='finished', + error='зависший прогон снят автоматически: висел в очереди дольше {STALE_MINUTES} минут', + finished_at=now() + WHERE status='queued' + AND celery_task_id IS NULL + AND started_at < now() - interval '{STALE_MINUTES} minutes' + RETURNING id, source_id + """)).all() + if reaped: + for source_id in {row.source_id for row in reaped}: + s.execute(text(""" + UPDATE parse_sources SET last_status='error', + last_error='зависший прогон снят автоматически' + WHERE id=:sid AND last_status='running' + """), {"sid": source_id}) + s.commit() + print(f"снято зависших прогонов: {len(reaped)} ({', '.join(str(r.id) for r in reaped)})") + active = s.execute(text( "SELECT count(*) FROM parse_runs WHERE status IN ('queued','running')" )).scalar_one()