gendesign/backend/app/workers/lifecycle.py
bot-backend 392cc2ffb8
All checks were successful
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 9s
CI Trade-In / backend-tests (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 2m5s
CI / backend-tests (pull_request) Successful in 17m2s
fix(ptica): осиротевшие прогоны Объектива закрываются, а не висят вечно (#2464)
`objective_scrape_runs` не подметалась ничем: `worker_ready` знал только
про `kn_scrape_runs` и `nspd_geo_jobs`. На проде 20.08.2026 в ней висело
6 строк `status='running'` с 17.05 — 94 суток, при 71 `done` и НИ ОДНОМ
`failed`. Отсутствие `failed` — след отравления сессии, из-за которого
`_finish_run(status='failed')` не мог записаться (причина починена
#2972). Причина устранена, но жёсткое убийство воркера (редеплой, OOM)
по-прежнему оставляет `running` навсегда: у Объектива нет ни своего
cleanup_zombies, ни снапшота для resume.

Тот же инвариант, что у kn: на worker_ready активных воркеров нет,
значит любая строка `running` осиротела. Resume не ставим —
возобновлять нечего.

`finished_at` ставится НЕ NOW(), а `COALESCE(heartbeat_at, started_at)`:
прогон, умерший 94 дня назад, не должен читаться как «завершён только
что». Монитору свежести это безразлично в обе стороны — `last_success_at`
и `recent_output` считаются только по `status='done'`, а
`last_attempt_at`/`last_status` — по `started_at`, так что зомби-строки
не попадают в него ни одним столбцом (проверено по коду
_FRESHNESS_SOURCES, а не предположено).

Двусторонне: против origin/main три теста красные по существу («не
трогает objective_scrape_runs», функция при этом отрабатывает 4 запроса
— то есть краснота не от отсутствующего символа). Мутационно проверены
оба контроля: снятие `WHERE status='running'` роняет
test_only_running_rows_are_touched, замена на `finished_at = NOW()`
роняет test_finished_at_is_last_sign_of_life_not_now.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-20 17:36:23 +05:00

318 lines
15 KiB
Python
Raw Permalink 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.

"""Worker lifecycle hooks — SQLAlchemy engine dispose on fork, zombie-resume on ready.
Handlers регистрируются через Celery signals декораторами при импорте модуля.
Импортируй этот модуль как side-effect из celery_app.py:
from app.workers import lifecycle # noqa: F401
worker_process_init — dispose SQLAlchemy engine в каждом prefork child-процессе,
чтобы избежать shared TCP-сокетов к PostgreSQL после fork().
worker_ready — при рестарте воркера находит 'running'/'paused' записи
kn_scrape_runs и nspd_geo_jobs и re-enqueue'ит их как zombie-resume.
Для nspd_geo 'paused' (WAF-пауза) применяется 30-минутный cooldown
(_ZOMBIE_PAUSED_THRESHOLD), чтобы редеплой не аннулировал WAF-защиту.
"""
import logging
from celery.signals import worker_process_init, worker_ready
logger = logging.getLogger(__name__)
@worker_process_init.connect
def _reset_db_connections(sender=None, **_kwargs) -> None:
"""Dispose SQLAlchemy engine in each prefork child process.
Background: Celery prefork pool fork()'s child processes from the master
that already imported SQLAlchemy + psycopg. The child inherits open TCP
sockets to PostgreSQL, but psycopg connections cannot be safely shared
across processes after fork — first use raises:
psycopg.ProgrammingError: can't change 'autocommit' now:
connection in transaction status INTRANS
Fix: dispose the engine in each child. SQLAlchemy will lazily open fresh
connections per-child. Standard Celery + SQLA recipe.
"""
try:
from app.core.db import engine
engine.dispose()
logger.info("worker_process_init: SQLAlchemy engine disposed for fork")
except Exception as e:
logger.warning("worker_process_init: engine dispose failed: %s", e)
@worker_ready.connect
def _resume_zombie_runs(sender=None, **_kwargs) -> None:
"""When a worker finishes booting (after redeploy/restart), find any sweep
that was 'running' or 'paused' and re-enqueue a resume task. Each resume
creates a fresh run_id linked via resumed_from_run_id; the original row is
marked 'zombie' so the audit trail is preserved.
No time threshold: by definition, on worker_ready ANY 'running' row is a
zombie because there is no active worker. Previously we required heartbeat
>5 min, which skipped jobs interrupted seconds before redeploy — they
stayed in 'running' status forever and required manual cancel/resume.
"""
logger.info("worker_ready: resume scan starting")
from sqlalchemy import text
# Persistent breadcrumb: write to nspd_geo_log so we can confirm via DB
# query whether worker_ready actually fired (independent of container logs).
# NULL job_id to avoid FK violation (FK is to nspd_geo_jobs.job_id which
# has no 0 row).
try:
from app.core.db import SessionLocal
_db = SessionLocal()
try:
_db.execute(
text(
"INSERT INTO nspd_geo_log (job_id, level, stage, message) "
"VALUES (NULL, 'info', 'worker_ready', 'worker_ready signal fired')"
)
)
_db.commit()
except Exception:
_db.rollback()
raise
finally:
_db.close()
except Exception as _e:
logger.warning("worker_ready: breadcrumb insert failed: %s", _e)
from app.core.db import SessionLocal
# ── kn_scrape_runs resume ──
db = SessionLocal()
ids: list[int] = []
try:
# #2464: берём ЛЮБУЮ строку в 'running', без фильтра по objects_snapshot.
# Докстринг этой функции формулирует инвариант прямо: «by definition, on
# worker_ready ANY 'running' row is a zombie because there is no active
# worker». Фильтр ему противоречил: строка без снапшота не попадала в
# выборку и оставалась 'running' НАВСЕГДА — ровно то состояние, ради
# устранения которого функция и заводилась.
#
# Снапшот всё равно нужен — но не для пометки, а для ВОЗОБНОВЛЕНИЯ:
# resume_kn_run восстанавливает обход «using objects_snapshot». Поэтому
# помечаем зомби всех, а resume ставим только тем, кого есть чем
# возобновить. Остальные получают честную причину вместо тишины.
rows = (
db.execute(
text(
"""
SELECT run_id, (objects_snapshot IS NOT NULL) AS resumable
FROM kn_scrape_runs
WHERE status = 'running'
ORDER BY started_at ASC
LIMIT 20
"""
)
)
.mappings()
.all()
)
if rows:
ids = [int(r["run_id"]) for r in rows if r["resumable"]]
orphan_ids = [int(r["run_id"]) for r in rows if not r["resumable"]]
if ids:
# Помечаем как 'zombie' одним апдейтом — resume создаст новые
# run_id со ссылкой resumed_from_run_id.
db.execute(
text(
"""
UPDATE kn_scrape_runs
SET status = 'zombie',
finished_at = NOW(),
error = COALESCE(error,
'auto-zombie at worker_ready, resume scheduled')
WHERE run_id = ANY(:ids)
"""
),
{"ids": ids},
)
if orphan_ids:
db.execute(
text(
"""
UPDATE kn_scrape_runs
SET status = 'zombie',
finished_at = NOW(),
error = COALESCE(error,
'auto-zombie at worker_ready, resume невозможен: '
'нет objects_snapshot')
WHERE run_id = ANY(:orphans)
"""
),
{"orphans": orphan_ids},
)
logger.warning(
"worker_ready: %d kn-прогонов без objects_snapshot помечены zombie"
" без resume — возобновлять нечем: %s",
len(orphan_ids),
orphan_ids,
)
db.commit()
else:
logger.info("worker_ready: нет stale kn runs для resume")
except Exception as e:
logger.exception("worker_ready kn resume scan failed: %s", e)
try:
db.rollback()
except Exception:
pass
finally:
db.close()
# Enqueue resume tasks. Late import to avoid circular at module load.
from app.workers.tasks.scrape_kn import resume_kn_run
for rid in ids:
try:
resume_kn_run.apply_async(args=[rid])
logger.info("worker_ready: resume_kn_run enqueued for run=%s", rid)
except Exception as e:
logger.warning("worker_ready: failed to enqueue resume for run=%s: %s", rid, e)
# NSPD-runs (старый Playwright-scraper) resume УДАЛЁН 2026-05-11. Скрапер
# снят с эксплуатации в пользу bulk geo-fetcher (nspd_geo). История
# старых runs сохранена в nspd_scrape_runs — никаких side-effects.
# NSPD geo-jobs: bulk-fetcher с собственной resume-логикой через
# nspd_geo_jobs / nspd_geo_targets. Resume любых 'running' jobs — на
# worker_ready по определению нет активных воркеров, всё running ==
# zombie. Раньше требовали heartbeat >10мин, что пропускало jobs убитых
# за минуту до редеплоя и оставляло их вечно висеть.
#
# 'paused' (consecutive_waf>=8 — NSPD-WAF забанил IP VPS) НЕ ре-enqueue'им
# безусловно: иначе каждый рестарт/редеплой воркера мгновенно аннулировал бы
# WAF-cooldown и worker снова долбил бы забаненный сервис. Применяем тот же
# 30-минутный порог, что и периодический cleanup_zombies
# (_ZOMBIE_PAUSED_THRESHOLD), переиспользуя константу чтобы избежать дрейфа.
from app.workers.tasks.nspd_geo import _ZOMBIE_PAUSED_THRESHOLD
db = SessionLocal()
geo_resume_jobs: list[int] = []
try:
rows = (
db.execute(
text(
"""
UPDATE nspd_geo_jobs
SET status = 'queued',
error = COALESCE(error, 'auto-resume at worker_ready')
WHERE status = 'running'
OR (
status = 'paused'
AND heartbeat_at
< NOW() - CAST(:paused_threshold AS interval)
)
RETURNING job_id
"""
),
{"paused_threshold": _ZOMBIE_PAUSED_THRESHOLD},
)
.mappings()
.all()
)
db.commit()
geo_resume_jobs = [int(r["job_id"]) for r in rows]
for jid in geo_resume_jobs:
logger.info("worker_ready: NSPD geo job=%s — resume scheduled", jid)
except Exception as e:
logger.warning("worker_ready nspd_geo resume scan failed: %s", e)
try:
db.rollback()
except Exception:
pass
finally:
db.close()
if geo_resume_jobs:
from app.workers.tasks.nspd_geo import process_nspd_geo_job
for jid in geo_resume_jobs:
try:
process_nspd_geo_job.apply_async(args=[jid], queue="geo")
logger.info("worker_ready: process_nspd_geo_job enqueued job=%s", jid)
except Exception as e:
logger.warning("worker_ready: failed to enqueue geo resume job=%s: %s", jid, e)
logger.info("worker_ready: resume scan finished (geo_jobs=%d)", len(geo_resume_jobs))
# objective_scrape_runs: тот же инвариант, что у kn — на worker_ready активных
# воркеров нет, значит любая строка в 'running' осиротела. Подметальщика у этой
# таблицы не было вовсе, и на проде 2026-08-20 висело 6 строк со статусом
# 'running' с 17.05 (94 суток), при 71 'done' и НИ ОДНОГО 'failed' — след
# отравления сессии, из-за которого _finish_run(status='failed') не мог
# записаться (починено #2972). Причина устранена, но жёсткое убийство воркера
# (редеплой, OOM) по-прежнему оставляет 'running' навсегда: у Объектива нет
# ни своего cleanup_zombies, ни snapshot'а для resume.
#
# Resume не делаем — возобновлять нечего (снапшота обхода нет), только честно
# закрываем. finished_at ставим НЕ NOW(), а по последнему признаку жизни:
# прогон, умерший 94 дня назад, не должен читаться как «завершён только что».
# Монитору свежести это безразлично в обе стороны — он считает last_success_at
# и recent_output только по status='done', а last_attempt_at/last_status — по
# started_at (см. _FRESHNESS_SOURCES в admin_scrape.py), так что зомби-строки
# в него не попадают ни одним столбцом.
db = SessionLocal()
try:
rows = (
db.execute(
text(
"""
UPDATE objective_scrape_runs
SET status = 'zombie',
finished_at = COALESCE(heartbeat_at, started_at),
error = COALESCE(error,
'auto-zombie at worker_ready: воркер перезапущен '
'во время прогона, возобновление невозможно')
WHERE status = 'running'
RETURNING run_id
"""
)
)
.mappings()
.all()
)
db.commit()
if rows:
logger.info(
"worker_ready: objective_scrape_runs — помечено зомби: %s",
[int(r["run_id"]) for r in rows],
)
else:
logger.info("worker_ready: нет осиротевших objective-прогонов")
except Exception as e:
logger.warning("worker_ready objective zombie sweep failed: %s", e)
try:
db.rollback()
except Exception:
pass
finally:
db.close()
# Sanity check: nspd_quarter_dumps table must exist (migration 88).
# Logs critical error but does NOT crash the worker — table may be absent
# in dev/staging before migration is applied.
try:
from app.core.db import SessionLocal as _SessionLocal
_db = _SessionLocal()
try:
_db.execute(text("SELECT 1 FROM nspd_quarter_dumps LIMIT 0"))
logger.info("worker_ready: nspd_quarter_dumps table OK")
except Exception as _te:
logger.critical(
"worker_ready: nspd_quarter_dumps table missing or inaccessible — "
"apply migration 88_nspd_quarter_dumps.sql before using harvest_quarter. "
"Error: %s",
_te,
)
finally:
_db.close()
except Exception as _e:
logger.warning("worker_ready: nspd_quarter_dumps sanity check failed: %s", _e)