All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
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 / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Successful in 1m54s
CI / backend-tests (pull_request) Successful in 17m0s
11 тасок объявляли `max_retries=2`, но ретраи не реализовывали: ни `autoretry_for` в декораторе, ни вызова `self.retry()` в теле. Celery в таком виде параметр не применяет — при исключении таска падает с первой попытки. Читающий код видит «до 3 попыток», а их одна. Убран `max_retries` у: cbr_macro_sync, rosstat_macro_sync, developer_registry_refresh, location_refresh, mv_sales_tracker_refresh, refresh_analytics, refresh_layout_velocity, refresh_quarter_price_index, scrape_objective.sync_objective_group, supply_layers_refresh, scrape_kn.scrape_kn_region. Заодно убран `bind=True` там, где `self` не использовался вовсе; в `scrape_kn_region` он оставлен — `self.request.id` пишется в kn_scrape_log. Не тронуты и не должны быть: `resume_kn_run` (max_retries=12 + настоящий self.retry()), `nspd_sync`/`scrape_cadastre` (autoretry_for), `nspd_geo`/`objective_etl` (max_retries=0 — честное «ретраев нет»). Гейт `test_2464_retry_config_is_real.py` разбирает AST всех модулей `app/workers/tasks/` и требует: если декоратор объявляет ненулевой max_retries, в нём есть autoretry_for либо в теле функции есть self.retry(). Три таски из одиннадцати гейт нашёл сверх списка эпика. Проверка гейта: с фиксом зелено, при возврате `max_retries=2` в supply_layers_refresh — красно с указанием на эту таску. Плюс два контроля: гейт видит ≥20 тасок (не молчит из-за пустой выборки) и признаёт обе законные формы ретраев. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
321 lines
14 KiB
Python
321 lines
14 KiB
Python
"""Celery task wrapper for naш.дом.рф kn-API sweeps."""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import logging
|
||
import random
|
||
import time
|
||
import uuid
|
||
from collections.abc import Iterator
|
||
from contextlib import contextmanager
|
||
from pathlib import Path
|
||
from typing import Any
|
||
|
||
import redis
|
||
|
||
from app.core.config import settings
|
||
from app.services.scrapers.domrf_kn import run_region_sweep
|
||
from app.workers.celery_app import celery_app
|
||
|
||
logger = logging.getLogger(__name__)
|
||
|
||
# Lock TTL — full Свердл sweep с extras занимает ~25-30 мин. Держим 45 мин —
|
||
# полтора max sweep duration: запас на хвостовые retry внутри sweep'а и не
|
||
# даём lock'у истечь под живым worker'ом (фикс P2 issue #1216 — раньше TTL
|
||
# равнялся sweep duration → race на границе).
|
||
# Если worker умер с lock'ом — он истечёт за 45 мин, а не за 2 часа (раньше),
|
||
# heartbeat-based autoclean в /queue endpoint всё равно подхватит зомби-runs.
|
||
_LOCK_TTL_SECONDS = 45 * 60
|
||
|
||
# Lua check-and-delete: освобождаем lock ТОЛЬКО если value совпадает с нашим
|
||
# owner-токеном. Защита от P2 race: TTL может истечь под живым sweep'ом, второй
|
||
# sweep встанет под тем же key с НОВЫМ токеном; финализация первого без guard
|
||
# снесла бы чужой lock → возможен третий параллельный sweep (issue #1216).
|
||
_RELEASE_LOCK_LUA = (
|
||
'if redis.call("GET", KEYS[1]) == ARGV[1] '
|
||
'then return redis.call("DEL", KEYS[1]) '
|
||
"else return 0 end"
|
||
)
|
||
|
||
|
||
def _lock_key(region_code: int, developers: list[str] | None) -> str:
|
||
"""Ключ синглтон-лока. Зависит от МНОЖЕСТВА разработчиков, не от их порядка.
|
||
|
||
#2464: раньше список джойнился как пришёл, а приходит он прямо из тела запроса
|
||
(`developers` в admin_scrape). Один и тот же набор, поданный в другом порядке,
|
||
давал ДРУГОЙ ключ — и синглтон молча переставал быть синглтоном: два свипа шли
|
||
параллельно по одним и тем же разработчикам.
|
||
|
||
Второе следствие того же: `force_release_lock` строит ключ этой же функцией.
|
||
Оператор, снимающий залипший лок и перечисливший разработчиков в ином порядке,
|
||
молча не снимал ничего.
|
||
|
||
sorted(set(...)) закрывает оба: порядок и повторы ('A','A' ≡ 'A').
|
||
"""
|
||
devs_key = ",".join(sorted(set(developers))) if developers else "*"
|
||
return f"scrape:kn:lock:{region_code}:{devs_key}"
|
||
|
||
|
||
def _log_task_received(region_code: int, devs: list[str] | None, task_id: str | None) -> None:
|
||
"""Записать в kn_scrape_log что worker получил задачу. Без run_id —
|
||
он появится позже в run_region_sweep. Best-effort: SQL ошибки гасим."""
|
||
from sqlalchemy import text
|
||
|
||
from app.core.db import SessionLocal
|
||
|
||
db = SessionLocal()
|
||
try:
|
||
db.execute(
|
||
text(
|
||
"""
|
||
INSERT INTO kn_scrape_log (run_id, level, stage, message)
|
||
VALUES (NULL, 'info', 'task_received', :msg)
|
||
"""
|
||
),
|
||
{
|
||
"msg": (
|
||
f"Worker подхватил task_id={task_id or '?'}, region={region_code},"
|
||
f" devs={devs or '*'}"
|
||
)[:500],
|
||
},
|
||
)
|
||
db.commit()
|
||
except Exception as e:
|
||
logger.warning("log_task_received failed: %s", e)
|
||
finally:
|
||
db.close()
|
||
|
||
|
||
def _log_task_skipped(
|
||
region_code: int, devs: list[str] | None, task_id: str | None, reason: str
|
||
) -> None:
|
||
from sqlalchemy import text
|
||
|
||
from app.core.db import SessionLocal
|
||
|
||
db = SessionLocal()
|
||
try:
|
||
db.execute(
|
||
text(
|
||
"""
|
||
INSERT INTO kn_scrape_log (run_id, level, stage, message)
|
||
VALUES (NULL, 'warn', 'task_skipped', :msg)
|
||
"""
|
||
),
|
||
{
|
||
"msg": (
|
||
f"task_id={task_id or '?'} skipped: {reason}, region={region_code},"
|
||
f" devs={devs or '*'}"
|
||
)[:500],
|
||
},
|
||
)
|
||
db.commit()
|
||
except Exception as e:
|
||
logger.warning("log_task_skipped failed: %s", e)
|
||
finally:
|
||
db.close()
|
||
|
||
|
||
def force_release_lock(region_code: int, developers: list[str] | None) -> bool:
|
||
"""Принудительно удалить Redis-lock (для emergency manual trigger,
|
||
когда зомби-worker оставил key). Возвращает True если ключ был."""
|
||
key = _lock_key(region_code, developers)
|
||
r = redis.Redis.from_url(settings.redis_url)
|
||
try:
|
||
return bool(r.delete(key))
|
||
except Exception as e:
|
||
logger.warning("force_release_lock %s failed: %s", key, e)
|
||
return False
|
||
|
||
|
||
@contextmanager
|
||
def _region_lock(region_code: int, developers: list[str] | None) -> Iterator[bool]:
|
||
"""Redis singleton lock so beat double-fire or two workers can't run the
|
||
same (region, devs) sweep concurrently. SET NX EX is atomic.
|
||
|
||
Owner-token (uuid4 hex) кладём в value и снимаем lock через Lua
|
||
check-and-delete — иначе истекший по TTL lock, перехваченный другим
|
||
sweep'ом, был бы снесён первым worker'ом при finally → возможен третий
|
||
параллельный sweep (P2 race, issue #1216).
|
||
"""
|
||
key = _lock_key(region_code, developers)
|
||
token = uuid.uuid4().hex
|
||
r = redis.Redis.from_url(settings.redis_url)
|
||
acquired = r.set(key, token, nx=True, ex=_LOCK_TTL_SECONDS)
|
||
try:
|
||
yield bool(acquired)
|
||
finally:
|
||
if acquired:
|
||
try:
|
||
# Lua check-and-delete: атомарно, мимо python GIL → нельзя
|
||
# снести чужой lock даже если наш TTL уже истёк.
|
||
r.eval(_RELEASE_LOCK_LUA, 1, key, token)
|
||
except Exception as e:
|
||
# Глотаем намеренно: lock истечёт сам по TTL. Логируем,
|
||
# чтобы не молча терять Redis-фейлы (не bare except).
|
||
logger.warning("release lock %s failed: %s", key, e)
|
||
|
||
|
||
# bind=True здесь настоящий: self.request.id пишется в kn_scrape_log. А вот
|
||
# max_retries=2 был инертен — self.retry() в этой таске не вызывается и
|
||
# autoretry_for не задан (#2464). self.retry() ниже принадлежит resume_kn_run.
|
||
@celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region")
|
||
def scrape_kn_region(
|
||
self: Any,
|
||
region_code: int,
|
||
developers: list[str] | None = None,
|
||
skip_jitter: bool = False,
|
||
extras: bool = True,
|
||
download_photos: bool = False,
|
||
fetch_flats: bool = True,
|
||
) -> dict[str, Any]:
|
||
"""Run a full kn-API sweep for one region. Returns summary dict.
|
||
|
||
For scheduled runs (skip_jitter=False), sleeps a random 0..SCRAPE_KN_JITTER_SECONDS
|
||
interval before starting so multiple workers / repeated weekly cycles don't hit
|
||
DOM.РФ at exactly the same minute. Manual triggers pass skip_jitter=True for
|
||
immediate execution.
|
||
"""
|
||
# Сразу пишем «task received» в kn_scrape_log БЕЗ run_id (run_id ещё не
|
||
# создан — он появится в run_region_sweep). Это позволяет UI видеть
|
||
# «worker подхватил задачу» даже если потом сработает lock_held или
|
||
# упадём ещё до Phase A.
|
||
_log_task_received(region_code, developers, self.request.id)
|
||
|
||
if not skip_jitter and settings.scrape_kn_jitter_seconds > 0:
|
||
delay = random.uniform(0, settings.scrape_kn_jitter_seconds)
|
||
logger.info("jitter delay %.0fs before scrape", delay)
|
||
time.sleep(delay)
|
||
|
||
state_path: str | None = settings.scrape_kn_state_path or None
|
||
if state_path and not Path(state_path).exists():
|
||
logger.warning("storage_state %s not found — bootstrap will run cold", state_path)
|
||
state_path = None
|
||
|
||
with _region_lock(region_code, developers) as got_lock:
|
||
if not got_lock:
|
||
logger.warning(
|
||
"scrape_kn_region SKIPPED region=%s devs=%s — another sweep is"
|
||
" already running (Redis lock held)",
|
||
region_code,
|
||
developers,
|
||
)
|
||
_log_task_skipped(region_code, developers, self.request.id, "lock_held")
|
||
return {
|
||
"skipped": True,
|
||
"reason": "lock_held",
|
||
"region_code": region_code,
|
||
"developers": developers,
|
||
"lock_key": _lock_key(region_code, developers),
|
||
}
|
||
|
||
logger.info(
|
||
"scrape_kn_region start region=%s devs=%s extras=%s photos=%s state=%s",
|
||
region_code,
|
||
developers,
|
||
extras,
|
||
download_photos,
|
||
state_path,
|
||
)
|
||
return asyncio.run(
|
||
run_region_sweep(
|
||
region_code=region_code,
|
||
developers=developers,
|
||
load_state=state_path,
|
||
fetch_flats=fetch_flats,
|
||
extras=extras,
|
||
download_photos_binary=download_photos,
|
||
# #1945 anti-ban: throttle + optional proxy + flats/extras isolation.
|
||
browser_concurrency=settings.scrape_kn_browser_concurrency,
|
||
request_jitter_min_ms=settings.scrape_kn_request_jitter_min_ms,
|
||
request_jitter_max_ms=settings.scrape_kn_request_jitter_max_ms,
|
||
proxy_url=settings.scrape_kn_proxy_url,
|
||
extras_isolated=settings.scrape_kn_extras_isolated,
|
||
)
|
||
)
|
||
|
||
|
||
@celery_app.task(bind=True, name="tasks.scrape_kn.resume_kn_run", max_retries=12)
|
||
def resume_kn_run(self: Any, run_id: int) -> dict[str, Any]:
|
||
"""Resume a previously interrupted sweep using objects_snapshot +
|
||
progress_obj_index from kn_scrape_runs. Triggered by worker_ready hook
|
||
when a stale 'running' row is found (heartbeat >5 min).
|
||
|
||
На lock_held делаем self.retry(countdown=300, max_retries=12) — окно
|
||
в ~час суммарно. Без retry resume пропускался до недельного beat
|
||
(P2 issue #1216): worker_ready помечал run 'zombie' и enqueue'ил resume
|
||
без countdown; live sweep держал lock → resume получал skipped и run
|
||
терялся до следующего планового запуска.
|
||
"""
|
||
from sqlalchemy import text
|
||
|
||
from app.core.db import SessionLocal
|
||
|
||
state_path: str | None = settings.scrape_kn_state_path or None
|
||
if state_path and not Path(state_path).exists():
|
||
state_path = None
|
||
|
||
# Reconstruct region/devs from the stale row to scope the lock correctly.
|
||
db = SessionLocal()
|
||
try:
|
||
row = (
|
||
db.execute(
|
||
text(
|
||
"""
|
||
SELECT region_codes, developer_ids
|
||
FROM kn_scrape_runs
|
||
WHERE run_id = :rid
|
||
"""
|
||
),
|
||
{"rid": run_id},
|
||
)
|
||
.mappings()
|
||
.first()
|
||
)
|
||
finally:
|
||
db.close()
|
||
if not row:
|
||
return {"skipped": True, "reason": "run_not_found", "run_id": run_id}
|
||
|
||
region_codes = list(row["region_codes"] or [])
|
||
developers = list(row["developer_ids"]) if row["developer_ids"] else None
|
||
if not region_codes:
|
||
return {"skipped": True, "reason": "no_region", "run_id": run_id}
|
||
region_code = region_codes[0]
|
||
|
||
_log_task_received(region_code, developers, self.request.id)
|
||
|
||
with _region_lock(region_code, developers) as got_lock:
|
||
if not got_lock:
|
||
logger.warning(
|
||
"resume_kn_run lock_held run=%s region=%s — retry in 5min"
|
||
" (live sweep still holds the lock)",
|
||
run_id,
|
||
region_code,
|
||
)
|
||
_log_task_skipped(region_code, developers, self.request.id, "lock_held_retry")
|
||
# 5 мин × 12 retries = 60 мин окно. Покрывает full sweep
|
||
# duration (~30 мин) + запас. Celery поднимет Retry exception,
|
||
# task поставится обратно в очередь с countdown=300.
|
||
raise self.retry(
|
||
exc=RuntimeError(f"lock_held for run={run_id} region={region_code}"),
|
||
countdown=300,
|
||
max_retries=12,
|
||
)
|
||
|
||
logger.info("resume_kn_run start run=%s region=%s devs=%s", run_id, region_code, developers)
|
||
return asyncio.run(
|
||
run_region_sweep(
|
||
region_code=region_code,
|
||
developers=developers,
|
||
load_state=state_path,
|
||
resume_from_run_id=run_id,
|
||
# #1945 anti-ban: throttle + optional proxy + flats/extras isolation.
|
||
browser_concurrency=settings.scrape_kn_browser_concurrency,
|
||
request_jitter_min_ms=settings.scrape_kn_request_jitter_min_ms,
|
||
request_jitter_max_ms=settings.scrape_kn_request_jitter_max_ms,
|
||
proxy_url=settings.scrape_kn_proxy_url,
|
||
extras_isolated=settings.scrape_kn_extras_isolated,
|
||
)
|
||
)
|