From 6982255fb3c516f64b2a2d5fd0d757202c2af531 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 20 Aug 2026 17:17:38 +0500 Subject: [PATCH] =?UTF-8?q?fix(workers):=20=D1=83=D0=B1=D1=80=D0=B0=D1=82?= =?UTF-8?q?=D1=8C=20max=5Fretries,=20=D0=BA=D0=BE=D1=82=D0=BE=D1=80=D1=8B?= =?UTF-8?q?=D0=B9=20=D0=BD=D0=B8=D1=87=D0=B5=D0=B3=D0=BE=20=D0=BD=D0=B5=20?= =?UTF-8?q?=D0=B4=D0=B5=D0=BB=D0=B0=D0=B5=D1=82=20(#2464)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- backend/app/workers/tasks/cbr_macro_sync.py | 15 ++- .../tasks/developer_registry_refresh.py | 9 +- backend/app/workers/tasks/location_refresh.py | 9 +- .../workers/tasks/mv_sales_tracker_refresh.py | 9 +- .../app/workers/tasks/refresh_analytics.py | 9 +- .../workers/tasks/refresh_layout_velocity.py | 10 +- .../tasks/refresh_quarter_price_index.py | 10 +- .../app/workers/tasks/rosstat_macro_sync.py | 16 ++- backend/app/workers/tasks/scrape_kn.py | 5 +- backend/app/workers/tasks/scrape_objective.py | 9 +- .../workers/tasks/supply_layers_refresh.py | 10 +- .../workers/test_2464_retry_config_is_real.py | 121 ++++++++++++++++++ 12 files changed, 201 insertions(+), 31 deletions(-) create mode 100644 backend/tests/workers/test_2464_retry_config_is_real.py diff --git a/backend/app/workers/tasks/cbr_macro_sync.py b/backend/app/workers/tasks/cbr_macro_sync.py index 92c27a5d..0c4be0c7 100644 --- a/backend/app/workers/tasks/cbr_macro_sync.py +++ b/backend/app/workers/tasks/cbr_macro_sync.py @@ -110,13 +110,22 @@ def _upsert_inflation(db: Session, rows: list[tuple[date, Decimal]]) -> int: return upserted +# Ретраев здесь НЕТ намеренно, и параметров, обещающих их, тоже быть не должно +# (#2464). Раньше стояло `bind=True, max_retries=2` — но self не использовался, +# self.retry() не вызывался и autoretry_for задан не был, поэтому конфигурация +# ретраев не имела ни малейшего эффекта: таска падала окончательно с первой ошибки, +# а параметр обещал до двух повторов. Соседи, где ретраи действительно нужны, +# задают их явно: autoretry_for в nspd_sync и scrape_cadastre, self.retry() в +# scrape_kn. +# +# Отсутствие ретраев — это и есть задуманное поведение, оно описано в докстринге +# ниже: «первая возникшая ошибка пробрасывается в конце (surfaces в +# Celery/GlitchTip), не глотается». Ряды тянутся по расписанию, следующий тик +# повторит попытку; молча ретраить внутри тика значило бы прятать отказ источника. @celery_app.task( - bind=True, name="tasks.cbr_macro_sync.cbr_macro_sync", - max_retries=2, ) def cbr_macro_sync( - self: Any, from_date: str | None = None, to_date: str | None = None, ) -> dict[str, Any]: diff --git a/backend/app/workers/tasks/developer_registry_refresh.py b/backend/app/workers/tasks/developer_registry_refresh.py index f20f7094..ea72d078 100644 --- a/backend/app/workers/tasks/developer_registry_refresh.py +++ b/backend/app/workers/tasks/developer_registry_refresh.py @@ -28,12 +28,15 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи +# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn). @celery_app.task( - bind=True, name="tasks.developer_registry_refresh.refresh_developer_registry", - max_retries=2, ) -def refresh_developer_registry(self: Any) -> dict[str, Any]: +def refresh_developer_registry() -> dict[str, Any]: """REFRESH MATERIALIZED VIEW CONCURRENTLY developer_registry. Лёгкая задача (реестр ~1024 застройщика). CONCURRENTLY — non-blocking для diff --git a/backend/app/workers/tasks/location_refresh.py b/backend/app/workers/tasks/location_refresh.py index 24c84b1b..903b9266 100644 --- a/backend/app/workers/tasks/location_refresh.py +++ b/backend/app/workers/tasks/location_refresh.py @@ -30,12 +30,15 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи +# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn). @celery_app.task( - bind=True, name="tasks.location_refresh.location_refresh", - max_retries=2, ) -def location_refresh(self: Any, region: str | None = None) -> dict[str, Any]: +def location_refresh(region: str | None = None) -> dict[str, Any]: """Пересчитать + upsert-нуть district-level индексы по всем районам в `location`. Идемпотентно (ON CONFLICT по district_name). Graceful: сбойный район diff --git a/backend/app/workers/tasks/mv_sales_tracker_refresh.py b/backend/app/workers/tasks/mv_sales_tracker_refresh.py index badee03f..2e87e56c 100644 --- a/backend/app/workers/tasks/mv_sales_tracker_refresh.py +++ b/backend/app/workers/tasks/mv_sales_tracker_refresh.py @@ -26,12 +26,15 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Соседи, где ретраи +# нужны, задают их явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn). @celery_app.task( - bind=True, name="tasks.mv_sales_tracker_refresh.refresh_sales_tracker_mvs", - max_retries=2, ) -def refresh_sales_tracker_mvs_task(self: Any) -> dict[str, Any]: +def refresh_sales_tracker_mvs_task() -> dict[str, Any]: """REFRESH both sales-tracker MVs (#61). Both MVs are refreshed CONCURRENTLY (non-blocking, require their UNIQUE diff --git a/backend/app/workers/tasks/refresh_analytics.py b/backend/app/workers/tasks/refresh_analytics.py index 5864fb37..d2a5ba77 100644 --- a/backend/app/workers/tasks/refresh_analytics.py +++ b/backend/app/workers/tasks/refresh_analytics.py @@ -17,13 +17,16 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они +# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с +# пояснением (nspd_geo, objective_etl.import_anton_objective). @celery_app.task( - bind=True, name="tasks.refresh_analytics.refresh_ekb_districts_medians", - max_retries=2, ) def refresh_ekb_districts_medians( - self: Any, window_months: int = 24, min_deals: int = 50, ) -> dict[str, Any]: diff --git a/backend/app/workers/tasks/refresh_layout_velocity.py b/backend/app/workers/tasks/refresh_layout_velocity.py index 7f67c4aa..b5706a78 100644 --- a/backend/app/workers/tasks/refresh_layout_velocity.py +++ b/backend/app/workers/tasks/refresh_layout_velocity.py @@ -22,12 +22,16 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они +# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с +# пояснением (nspd_geo, objective_etl.import_anton_objective). @celery_app.task( - bind=True, name="tasks.refresh_layout_velocity.refresh_layout_velocity", - max_retries=2, ) -def refresh_layout_velocity_task(self: Any) -> dict[str, Any]: +def refresh_layout_velocity_task() -> dict[str, Any]: """REFRESH MATERIALIZED VIEW mv_layout_velocity (best_layouts, #113 / #1666). MV рефрешится CONCURRENTLY (non-blocking, требует unique-индекс diff --git a/backend/app/workers/tasks/refresh_quarter_price_index.py b/backend/app/workers/tasks/refresh_quarter_price_index.py index 803e3457..d8f0d3b9 100644 --- a/backend/app/workers/tasks/refresh_quarter_price_index.py +++ b/backend/app/workers/tasks/refresh_quarter_price_index.py @@ -20,12 +20,16 @@ from app.workers.celery_app import celery_app logger = logging.getLogger(__name__) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они +# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с +# пояснением (nspd_geo, objective_etl.import_anton_objective). @celery_app.task( - bind=True, name="tasks.refresh_quarter_price_index.refresh_quarter_price_index_chain", - max_retries=2, ) -def refresh_quarter_price_index_chain(self: Any) -> dict[str, Any]: +def refresh_quarter_price_index_chain() -> dict[str, Any]: """Refresh mv_quarter_price_per_m2 then mv_quarter_price_index in sequence. Both MVs are refreshed CONCURRENTLY (non-blocking). Falls back to diff --git a/backend/app/workers/tasks/rosstat_macro_sync.py b/backend/app/workers/tasks/rosstat_macro_sync.py index cf84c5a3..c24ebe0c 100644 --- a/backend/app/workers/tasks/rosstat_macro_sync.py +++ b/backend/app/workers/tasks/rosstat_macro_sync.py @@ -231,12 +231,22 @@ def _upsert_emiss_rows(db: Session, rows: list[EmissRow]) -> int: return upserted +# Ретраев здесь НЕТ намеренно, и параметров, обещающих их, тоже быть не должно +# (#2464). Раньше стояло `bind=True, max_retries=2` — но self не использовался, +# self.retry() не вызывался и autoretry_for задан не был, поэтому конфигурация +# ретраев не имела ни малейшего эффекта: таска падала окончательно с первой ошибки, +# а параметр обещал до двух повторов. Соседи, где ретраи действительно нужны, +# задают их явно: autoretry_for в nspd_sync и scrape_cadastre, self.retry() в +# scrape_kn. +# +# Отсутствие ретраев — это и есть задуманное поведение, оно описано в докстринге +# ниже: «первая возникшая ошибка пробрасывается в конце (surfaces в +# Celery/GlitchTip), не глотается». Ряды тянутся по расписанию, следующий тик +# повторит попытку; молча ретраить внутри тика значило бы прятать отказ источника. @celery_app.task( - bind=True, name="tasks.rosstat_macro_sync.rosstat_macro_sync", - max_retries=2, ) -def rosstat_macro_sync(self: Any) -> dict[str, Any]: +def rosstat_macro_sync() -> dict[str, Any]: """Загрузить ряды Росстата (open-data + ЕМИСС + xlsx-СМР) и апсертить в macro_indicator. Источники выполняются НЕЗАВИСИМО (per-source try/except): сбой одного источника diff --git a/backend/app/workers/tasks/scrape_kn.py b/backend/app/workers/tasks/scrape_kn.py index 0e7af756..7ba512a1 100644 --- a/backend/app/workers/tasks/scrape_kn.py +++ b/backend/app/workers/tasks/scrape_kn.py @@ -157,7 +157,10 @@ def _region_lock(region_code: int, developers: list[str] | None) -> Iterator[boo logger.warning("release lock %s failed: %s", key, e) -@celery_app.task(bind=True, name="tasks.scrape_kn.scrape_kn_region", max_retries=2) +# 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, diff --git a/backend/app/workers/tasks/scrape_objective.py b/backend/app/workers/tasks/scrape_objective.py index cfbb0ca8..7ab5506d 100644 --- a/backend/app/workers/tasks/scrape_objective.py +++ b/backend/app/workers/tasks/scrape_objective.py @@ -146,13 +146,16 @@ def _save_raw( return int(row) +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они +# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с +# пояснением (nspd_geo, objective_etl.import_anton_objective). @celery_app.task( - bind=True, name="tasks.scrape_objective.sync_objective_group", - max_retries=2, ) def sync_objective_group( - self: Any, group_name: str | None = None, triggered_by: str = "beat", use_ddu: bool = True, diff --git a/backend/app/workers/tasks/supply_layers_refresh.py b/backend/app/workers/tasks/supply_layers_refresh.py index eddcd606..667450ed 100644 --- a/backend/app/workers/tasks/supply_layers_refresh.py +++ b/backend/app/workers/tasks/supply_layers_refresh.py @@ -155,12 +155,16 @@ def _upsert_rows(db: Session, rows: list[SupplyLayerRow]) -> tuple[int, int]: return upserted, skipped +# Ретраев здесь нет, и параметров, обещающих их, быть не должно (#2464). Стояло +# `bind=True, max_retries=2`, но self не использовался, self.retry() не вызывался и +# autoretry_for задан не был — конфигурация не имела эффекта. Где ретраи нужны, они +# задаются явно: autoretry_for (nspd_sync, scrape_cadastre) или self.retry() +# (scrape_kn.resume_kn_run). Где их сознательно нет — пишется max_retries=0 с +# пояснением (nspd_geo, objective_etl.import_anton_objective). @celery_app.task( - bind=True, name="tasks.supply_layers_refresh.supply_layers_refresh", - max_retries=2, ) -def supply_layers_refresh(self: Any) -> dict[str, Any]: +def supply_layers_refresh() -> dict[str, Any]: """Пересчитать 3-слойный склад предложения по всем районам и UPSERT в supply_layers. Coherent snapshot: ``run_date = date.today()`` фиксируется ОДИН раз и передаётся diff --git a/backend/tests/workers/test_2464_retry_config_is_real.py b/backend/tests/workers/test_2464_retry_config_is_real.py new file mode 100644 index 00000000..7c517638 --- /dev/null +++ b/backend/tests/workers/test_2464_retry_config_is_real.py @@ -0,0 +1,121 @@ +"""Таска, объявившая max_retries, обязана ретраи РЕАЛИЗОВАТЬ (#2464). + +`cbr_macro_sync` и `rosstat_macro_sync` были объявлены как `bind=True, max_retries=2`, но +`self` не использовался, `self.retry()` не вызывался и `autoretry_for` задан не был. То +есть конфигурация ретраев не имела ни малейшего эффекта: таска падала окончательно с первой +ошибки, а параметр обещал до двух повторов. Ни красного, ни ошибки — параметр просто +декорация, и отличить её от работающей настройки можно только чтением тела. + +Гейт разбирает КАЖДЫЙ декоратор `@celery_app.task(...)` через ast: если в нём есть +`max_retries`, то либо в том же декораторе должен стоять `autoretry_for`, либо в теле +функции — вызов `self.retry(...)`. + +Оговорка про силу проверки: гейт СТРУКТУРНЫЙ, он читает исходники. Поведенческим его +сделать нельзя — инертная конфигурация по определению ничего не меняет в поведении, и +поймать её можно только по несоответствию объявления и кода. +""" + +from __future__ import annotations + +import ast +import os +from pathlib import Path + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +_TASKS_DIR = Path(__file__).resolve().parents[2] / "app" / "workers" + + +def _task_decorators(tree: ast.AST): + """(имя функции, декоратор celery-таски, узел функции) для каждой таски модуля.""" + for node in ast.walk(tree): + if not isinstance(node, ast.FunctionDef | ast.AsyncFunctionDef): + continue + for dec in node.decorator_list: + if not isinstance(dec, ast.Call): + continue + target = dec.func + name = getattr(target, "attr", None) or getattr(target, "id", None) + if name == "task": + yield node.name, dec, node + + +def _calls_self_retry(fn: ast.AST) -> bool: + for node in ast.walk(fn): + if isinstance(node, ast.Call): + f = node.func + if isinstance(f, ast.Attribute) and f.attr == "retry": + owner = getattr(f.value, "id", None) + if owner == "self": + return True + return False + + +def _offenders() -> list[str]: + bad: list[str] = [] + for path in sorted(_TASKS_DIR.rglob("*.py")): + tree = ast.parse(path.read_text()) + for fn_name, dec, fn_node in _task_decorators(tree): + kwargs = {kw.arg: kw.value for kw in dec.keywords if kw.arg} + if "max_retries" not in kwargs: + continue + # max_retries=0 — ЧЕСТНОЕ объявление «ретраев нет», а не пустое обещание. + # Такие места в репозитории есть и снабжены пояснением (nspd_geo: + # «resume через worker_ready, не через retry»; objective_etl: «при сбое + # лучше человек посмотрит»). Флагуем только ненулевое обещание. + v = kwargs["max_retries"] + if isinstance(v, ast.Constant) and v.value == 0: + continue + if "autoretry_for" in kwargs or _calls_self_retry(fn_node): + continue + bad.append(f"{path.relative_to(_TASKS_DIR.parent.parent)}::{fn_name}") + return bad + + +def test_declared_retries_are_actually_implemented() -> None: + """max_retries без autoretry_for и без self.retry() — обещание без исполнения.""" + bad = _offenders() + assert not bad, ( + "таски объявляют max_retries, но ретраи не реализуют — параметр не имеет эффекта:\n " + + "\n ".join(bad) + + "\nЛибо задай autoretry_for / вызови self.retry(), либо убери max_retries." + ) + + +def test_gate_sees_the_tasks_at_all() -> None: + """Контроль на сам гейт: он должен что-то находить. + + Если разбор перестанет узнавать декораторы (переименуют celery_app, сменят + обёртку), проверка выше станет тавтологически зелёной. + """ + found = 0 + for path in sorted(_TASKS_DIR.rglob("*.py")): + found += sum(1 for _ in _task_decorators(ast.parse(path.read_text()))) + assert found >= 20, f"гейт нашёл всего {found} celery-тасок — разбор сломался" + + +def test_gate_recognises_both_valid_forms() -> None: + """Контроль: обе законные формы ретраев распознаются, а не только одна. + + Иначе гейт краснел бы на исправных тасках и его бы отключили. + """ + src_autoretry = ( + "@celery_app.task(bind=True, max_retries=2, autoretry_for=(ValueError,))\n" + "def t(self):\n return 1\n" + ) + src_self_retry = ( + "@celery_app.task(bind=True, max_retries=2)\n" + "def t(self):\n raise self.retry(countdown=1)\n" + ) + src_inert = "@celery_app.task(bind=True, max_retries=2)\ndef t(self):\n return 1\n" + + def _is_ok(src: str) -> bool: + tree = ast.parse(src) + for _n, dec, fn in _task_decorators(tree): + kwargs = {kw.arg for kw in dec.keywords if kw.arg} + return "autoretry_for" in kwargs or _calls_self_retry(fn) + raise AssertionError("декоратор не распознан") + + assert _is_ok(src_autoretry), "форма autoretry_for не распознана" + assert _is_ok(src_self_retry), "форма self.retry() не распознана" + assert not _is_ok(src_inert), "инертная форма ошибочно признана исправной" -- 2.45.3