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
`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>
318 lines
15 KiB
Python
318 lines
15 KiB
Python
"""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)
|