From 917c911eb9f9b9b73d031ee221a172ef445a97bc Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 6 Aug 2026 13:12:31 +0500 Subject: [PATCH] =?UTF-8?q?fix(tradein):=20TTL-=D0=B4=D0=B5=D0=B0=D0=BA?= =?UTF-8?q?=D1=82=D0=B8=D0=B2=D0=B0=D1=86=D0=B8=D1=8F=20=D0=BD=D0=B5=20?= =?UTF-8?q?=D0=B8=D1=81=D0=BF=D0=BE=D0=BB=D0=BD=D1=8F=D0=B5=D1=82=D1=81?= =?UTF-8?q?=D1=8F,=20=D0=BF=D0=BE=D0=BA=D0=B0=20=D1=81=D0=B1=D0=BE=D1=80?= =?UTF-8?q?=20=D0=BF=D0=BE=20=D0=B8=D1=81=D1=82=D0=BE=D1=87=D0=BD=D0=B8?= =?UTF-8?q?=D0=BA=D1=83=20=D0=BB=D0=B5=D0=B6=D0=B8=D1=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Авито 10.07-26.07.2026: 17 суток подряд без единой собранной страницы, TTL=10 снял за отрезок 9 033 строки. 1 270 из них потом вернулись живыми, как только сбор восстановился (сверка listing_source_snapshots с текущим last_seen_at) — задача прочитала «мы не смогли зайти» как «объявление снято». Гейт меряет здоровье НЕ статусом прогона, а результатом: сколько строк источник подтвердил свежими за 3 суток, по той же колонке (last_seen_at / scraped_at) и тому же срезу source+segment, что и сам UPDATE. Статус не годится — оба других провала для него невидимы: yandex 18.07-30.07 (13 суток, все прогоны 'done', ноль 'banned', total_seen=0) и domklik 20.07-30.07, чей TTL 02.08 снял 6 131 строку разом. Пороги посчитаны по ряду подтверждений из listing_source_snapshots и лежат между максимумом провала и минимумом здоровых суток: avito 1500 (провал 0..970, здоровье 3542..6079), yandex/cian 500, domklik 200. Миграция 219 кладёт их в default_params, не затирая ручную настройку; страховочный дефолт в коде — 500. Цена ошибки асимметрична (пропущенная деактивация чинится следующим прогоном, ложная — только повторным сбором), поэтому при сомнении прогон пропускается. Пропущенный прогон закрывается mark_done со счётчиками {deactivated: 0, confirmations: N, skipped_unhealthy: 1} — ни одна строка не тронута, статусы прогонов не менялись. Тест test_gate_blocks_every_day_of_the_17_day_avito_ban проигрывает реальный прод-ряд по дням: каждые сутки провала блокируются, здоровые — нет. Refs #2659 --- .../backend/app/services/product_handlers.py | 10 +- .../app/tasks/deactivate_stale_avito.py | 112 ++++++ .../sql/219_deactivate_stale_health_gate.sql | 67 ++++ .../test_deactivate_stale_health_gate.py | 321 ++++++++++++++++++ 4 files changed, 509 insertions(+), 1 deletion(-) create mode 100644 tradein-mvp/backend/data/sql/219_deactivate_stale_health_gate.sql create mode 100644 tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index 1069ab8b..74d8432c 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -216,12 +216,19 @@ async def _job_deactivate_stale( db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext ) -> None: from app.core.config import settings as _settings - from app.tasks.deactivate_stale_avito import deactivate_stale_listings + from app.tasks.deactivate_stale_avito import ( + DEFAULT_MIN_CONFIRMATIONS, + deactivate_stale_listings, + ) listing_source: str = params.get("listing_source", "avito") ttl_days: int = params.get("ttl_days", _settings.avito_stale_ttl_days) segments: list[str] | None = params.get("segments") staleness_column: str = params.get("staleness_column", "last_seen_at") + # Гейт по здоровью сбора (#2659) включён по умолчанию: незасеянное расписание + # получает страховочный порог, а не «деактивируй вслепую». Посчитанные по + # источнику пороги приходят из default_params (миграция 219). + min_confirmations: int = params.get("min_confirmations", DEFAULT_MIN_CONFIRMATIONS) loop = asyncio.get_event_loop() await loop.run_in_executor( @@ -233,6 +240,7 @@ async def _job_deactivate_stale( ttl_days=ttl_days, segments=segments, staleness_column=staleness_column, + min_confirmations=min_confirmations, ), ) diff --git a/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py b/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py index 77a705bf..1f800160 100644 --- a/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py +++ b/tradein-mvp/backend/app/tasks/deactivate_stale_avito.py @@ -82,6 +82,70 @@ _STALE_SNAPSHOT_TAIL = """ """ +# ── Гейт по здоровью сбора (#2659) ──────────────────────────────────────────── +# TTL отвечает на вопрос «объявление сняли?», а меряет «мы его давно не видели». +# Пока обход здоров, разница мала. Когда обход лёг — разница равна всему инвентарю. +# +# Замер на проде, из-за которого этот гейт существует. Авито 10.07-26.07.2026: +# 17 суток подряд без единой успешно собранной страницы, TTL=10 снял за этот отрезок +# 9 033 строки; 1 270 из них потом доказанно вернулись живыми (снимки +# listing_source_snapshots + текущий last_seen_at) — сбор восстановился, и объявления +# оказались на месте. То есть «мы не смогли зайти» было прочитано как «объявление снято». +# +# ПОЧЕМУ НЕ ПО СТАТУСУ ban. Соблазн взять scrape_runs.status='banned' — ловушка: +# Яндекс 18.07-30.07 — 5 прогонов в сутки, ВСЕ 'done', НОЛЬ 'banned', total_seen=0 +# 13 суток подряд (снято ~839 строк vtorichka); +# Домклик 20.07-30.07 — то же самое, 11 суток 'done' с total_seen=0, а 02.08 TTL +# снял 6 131 строку разом (см. комментарий про 'stale' выше). +# Оба провала для ban-детектора невидимы. Поэтому здоровье меряем НЕ статусом прогона, +# а результатом: сколько строк источник реально подтвердил свежими за последние сутки. +# +# МЕТРИКА: count(*) по той же колонке свежести, что и сам TTL (last_seen_at или +# scraped_at) и по тому же срезу source+segment, что и UPDATE. Одна колонка на обе +# стороны — гейт нельзя обмануть bulk-touch'ем, который не двигает scraped_at (#2204). +# +# ПОРОГ. Ряд «подтверждений за 3 суток» по дням (восстановлен из listing_source_snapshots): +# avito здоровые сутки 3542..6079, провал 10.07-26.07 — 0..970 → порог 1500; +# yandex vtorichka здоровые 897..2206, провал — 0 → порог 500; +# cian vtorichka 748..4329, провала не было → порог 500; +# domklik по scraped_at сейчас 62/3 суток (сбор фактически стоит) → порог 200. +# Пороги живут в default_params расписания (миграция 219), здесь только страховка +# на случай незасеянного расписания. Асимметрия цены ошибки намеренная: пропущенная +# деактивация чинится следующим прогоном, ложная — только повторным сбором, которого +# может не быть. Поэтому при сомнении — пропускаем прогон. +# +# ПОТОЛОК: окно 3 суток годится, пока свипы источника ходят не реже чем раз в 3 дня. +# Источник с более редкой каденцией будет блокироваться всегда — тогда окно нужно +# растить до каденции, а не понижать порог. +_HEALTH_WINDOW_DAYS = 3 + +# Страховка для расписаний без явного min_confirmations в default_params: ловит +# полный ноль и близкое к нулю, но НЕ ловит частичный провал вроде avito 936-970 — +# для этого нужен посчитанный по источнику порог из миграции 219. +DEFAULT_MIN_CONFIRMATIONS = 500 + +_CONFIRMATIONS_SEGMENT_FILTER = "\n AND listing_segment = ANY(CAST(:segments AS text[]))" + + +def _build_confirmations_sql(staleness_column: str, *, with_segments: bool) -> Any: + """SELECT count(*) подтверждённых за окно строк — тот же срез, что и у UPDATE. + + staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings. + Значения (:listing_source, :health_window_days, :segments) — param-binding, + psycopg v3 safe (CAST(... AS ...), никаких :param::type). + """ + segment_filter = _CONFIRMATIONS_SEGMENT_FILTER if with_segments else "" + return text( + f""" + SELECT count(*) + FROM listings + WHERE source = :listing_source + AND {staleness_column} + > NOW() - CAST(:health_window_days || ' days' AS interval){segment_filter} + """ + ) + + def _build_all_segments_sql(staleness_column: str) -> Any: """UPDATE без фильтра по сегменту: все сегменты для данного source. @@ -143,6 +207,8 @@ def deactivate_stale_listings( ttl_days: int, segments: list[str] | None = None, staleness_column: str = "last_seen_at", + min_confirmations: int = 0, + health_window_days: int = _HEALTH_WINDOW_DAYS, ) -> dict[str, int]: """Пометить is_active=false объявления, чья свежесть старше ttl_days дней. @@ -157,12 +223,20 @@ def deactivate_stale_listings( last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch двигает last_seen_at всем строкам одним timestamp, поэтому честная свежесть = scraped_at (двигается только реальным скрейпом). + min_confirmations: гейт по здоровью сбора (#2659). Сколько строк источник + должен был подтвердить свежими за health_window_days суток, чтобы + деактивации вообще разрешалось исполниться. 0 -> гейт выключен (так + вызывают старые тесты и совместимая обёртка); реальные значения приходят + из default_params расписания, см. миграцию 219 и комментарий выше. + health_window_days: окно подтверждений для гейта, суток. Дефолт 3. Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources). Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots (data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed). Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками). + Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N, + "skipped_unhealthy": 1} и НИ ОДНА строка не тронута. Raises: ValueError: если staleness_column не входит в whitelist (проверка ДО SQL, @@ -180,6 +254,44 @@ def deactivate_stale_listings( f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}" ) + # Гейт по здоровью сбора (#2659) — ДО любого UPDATE. Деактивация необратима + # на практике (вернуть «живость» может только повторный сбор), поэтому + # проверяем ПЕРЕД записью, а не откатываем после. + if min_confirmations > 0: + health_params: dict[str, Any] = { + "listing_source": listing_source, + "health_window_days": health_window_days, + } + if segments is not None: + health_params["segments"] = segments + confirmations = ( + db.execute( + _build_confirmations_sql(staleness_column, with_segments=segments is not None), + health_params, + ).scalar() + or 0 + ) + counters["confirmations"] = int(confirmations) + if confirmations < min_confirmations: + counters["skipped_unhealthy"] = 1 + # Ничего не писали (был только SELECT) — rollback закрывает транзакцию + # чисто, чтобы mark_done стартовал со своей. + db.rollback() + runs_mod.mark_done(db, run_id, counters) + logger.warning( + "deactivate_stale source=%s run_id=%d SKIPPED: сбор нездоров — " + "подтверждений за %d сут %d < порога %d " + "(segments=%r, staleness_column=%s); ни одна строка не тронута", + listing_source, + run_id, + health_window_days, + confirmations, + min_confirmations, + segments, + staleness_column, + ) + return counters + # segments is None -> все сегменты (поведение avito). segments=[...] -> только # перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список [] # означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0), diff --git a/tradein-mvp/backend/data/sql/219_deactivate_stale_health_gate.sql b/tradein-mvp/backend/data/sql/219_deactivate_stale_health_gate.sql new file mode 100644 index 00000000..4fbefecd --- /dev/null +++ b/tradein-mvp/backend/data/sql/219_deactivate_stale_health_gate.sql @@ -0,0 +1,67 @@ +-- 219_deactivate_stale_health_gate.sql +-- Пороги гейта здоровья сбора для TTL-деактивации (#2659). +-- +-- ЗАЧЕМ. TTL отвечает на вопрос «объявление сняли?», а меряет «мы его давно не +-- видели». Пока обход здоров, разница мала; когда обход лёг — разница равна всему +-- инвентарю. Прод, авито 10.07-26.07.2026: 17 суток подряд без единой собранной +-- страницы, TTL=10 снял 9 033 строки, из них 1 270 доказанно вернулись живыми, +-- как только сбор восстановился (сверка listing_source_snapshots с текущим +-- last_seen_at). Код гейта — app/tasks/deactivate_stale_avito.py. +-- +-- ПОЧЕМУ НЕ ПО СТАТУСУ ПРОГОНА. Ban-детектор эти провалы НЕ ловит: +-- yandex 18.07-30.07 — 5 прогонов в сутки, ВСЕ 'done', НОЛЬ 'banned', +-- total_seen = 0 тринадцать суток подряд; +-- domklik 20.07-30.07 — 11 суток 'done' с total_seen = 0, а 02.08 TTL=14 +-- снял 6 131 строку разом. +-- Поэтому здоровье меряется результатом (сколько строк источник реально подтвердил +-- свежими за 3 суток), а не статусом прогона. +-- +-- ОТКУДА ЧИСЛА. Ряд «подтверждений за 3 суток» по дням восстановлен из +-- listing_source_snapshots (снимок last_seen_at на каждую дату), срез совпадает +-- со срезом соответствующего UPDATE (source + segments + та же колонка свежести): +-- avito (все сегменты, last_seen_at): здоровые сутки 3542..6079, +-- провал 08.07-29.07 — 0..970 -> 1500 +-- yandex (vtorichka, last_seen_at): здоровые 897..2206, провал 0 -> 500 +-- cian (vtorichka, last_seen_at): 748..4329, провалов не было -> 500 +-- domklik (все сегменты, scraped_at): сейчас 62 за 3 суток, сбор +-- фактически стоит -> 200 +-- Каждый порог лежит между максимумом провала и минимумом здоровых суток: +-- авито 970 < 1500 < 3542 — исторический случай ловится с запасом в обе стороны. +-- +-- domklik ЗАБЛОКИРУЕТСЯ СРАЗУ, и это верный исход, а не сбой миграции: источник +-- подтверждает ~50 строк в сутки, TTL по нему уже один раз (02.08) снёс инвентарь +-- целиком. Пока сбор не восстановлен, деактивации там нечего подтверждать; протухшие +-- строки закрываются фильтром свежести на стороне чтения (#2656), а не TTL. +-- +-- Цена ошибки асимметрична: пропущенная деактивация чинится следующим прогоном, +-- ложная — только повторным сбором, которого может не быть. Пороги поэтому +-- смещены в сторону «пропустить прогон». +-- +-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)), +-- 090/115/160 (сами расписания deactivate_stale_*). +-- ТОЛЬКО данные (UPDATE default_params), DDL нет. +-- Идемпотентность + уважение к ручной настройке: ключ проставляется лишь там, где +-- его ещё нет, поэтому повторный прогон файла не затирает подкрученное оператором +-- значение. Снять гейт вручную: min_confirmations = 0. + +BEGIN; + +UPDATE scrape_schedules +SET default_params = default_params || jsonb_build_object('min_confirmations', 1500), + updated_at = NOW() +WHERE source = 'deactivate_stale_avito' + AND NOT default_params ? 'min_confirmations'; + +UPDATE scrape_schedules +SET default_params = default_params || jsonb_build_object('min_confirmations', 500), + updated_at = NOW() +WHERE source IN ('deactivate_stale_cian', 'deactivate_stale_yandex') + AND NOT default_params ? 'min_confirmations'; + +UPDATE scrape_schedules +SET default_params = default_params || jsonb_build_object('min_confirmations', 200), + updated_at = NOW() +WHERE source = 'deactivate_stale_domklik' + AND NOT default_params ? 'min_confirmations'; + +COMMIT; diff --git a/tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py b/tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py new file mode 100644 index 00000000..89d4ef30 --- /dev/null +++ b/tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py @@ -0,0 +1,321 @@ +"""Гейт по здоровью сбора для TTL-деактивации (#2659). + +Ключевой тест здесь — test_gate_blocks_every_day_of_the_17_day_avito_ban: он +проигрывает РЕАЛЬНЫЙ прод-ряд подтверждений по дням и требует, чтобы порог из +миграции 219 заблокировал каждые сутки провала 10.07-26.07.2026 и не тронул ни +одних здоровых суток. На старом коде (без min_confirmations) он не проходит: +деактивация исполнялась вслепую. +""" + +from __future__ import annotations + +import os +import re +from pathlib import Path +from typing import Any + +import pytest + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +from app.tasks import deactivate_stale_avito as task_mod + +_SQL_DIR = Path(__file__).resolve().parents[1] / "data" / "sql" +_MIGRATION_219 = _SQL_DIR / "219_deactivate_stale_health_gate.sql" + +# ── Прод-ряд: подтверждений за 3 суток по дням, avito, все сегменты ──────────── +# Восстановлено из listing_source_snapshots (снимок last_seen_at на каждую дату): +# count(*) FILTER (WHERE last_seen_at > snapshot_date - interval '3 days') +# Провал сбора: 08.07-29.07.2026 (17 суток, за которые TTL=10 снял 9 033 строки). +# Здоровые сутки — до 07.07 и после восстановления 03.08. +_AVITO_BAN_DAYS: dict[str, int] = { + "2026-07-08": 936, + "2026-07-09": 970, + "2026-07-10": 830, + "2026-07-11": 300, + "2026-07-13": 160, + "2026-07-14": 683, + "2026-07-15": 683, + "2026-07-16": 665, + "2026-07-17": 0, + "2026-07-18": 0, + "2026-07-19": 0, + "2026-07-20": 0, + "2026-07-21": 0, + "2026-07-22": 0, + "2026-07-23": 0, + "2026-07-24": 0, + "2026-07-25": 0, + "2026-07-27": 0, + "2026-07-28": 0, + "2026-07-29": 0, +} +_AVITO_HEALTHY_DAYS: dict[str, int] = { + "2026-06-25": 3542, + "2026-06-26": 3902, + "2026-06-27": 4017, + "2026-06-28": 3894, + "2026-06-29": 4424, + "2026-06-30": 4943, + "2026-07-01": 5270, + "2026-07-02": 4828, + "2026-07-05": 6079, + "2026-07-06": 5060, + "2026-07-07": 3666, + "2026-08-03": 4039, + "2026-08-04": 4247, + "2026-08-05": 4197, + "2026-08-06": 2542, +} +# Порог из миграции 219 для deactivate_stale_avito. +_AVITO_MIN_CONFIRMATIONS = 1500 + +# Сколько строк TTL снял в каждые сутки провала (scrape_runs.counters->>'deactivated'). +# Сумма = 9 033 — цифра из #2659, перепроверена на проде. +_AVITO_BAN_DEACTIVATED = [316, 970, 742, 676, 1541, 2959, 216, 536, 146, 95, 153, 0, 0, 683] + + +# ── Фейковая сессия ─────────────────────────────────────────────────────────── + + +class _FakeResult: + def __init__(self, rowcount: int = 0, scalar_value: int | None = None) -> None: + self.rowcount = rowcount + self._scalar = scalar_value + + def scalar(self) -> int | None: + return self._scalar + + +class _FakeDB: + """Session-заглушка: SELECT count(*) отдаёт confirmations, UPDATE — rowcount.""" + + def __init__(self, *, confirmations: int, rowcount: int = 137) -> None: + self._confirmations = confirmations + self._rowcount = rowcount + self.executed: list[tuple[str, dict[str, Any] | None]] = [] + self.committed = False + self.rolled_back = False + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = str(stmt.text) + self.executed.append((sql, params)) + if "SELECT count(*)" in sql: + return _FakeResult(scalar_value=self._confirmations) + return _FakeResult(rowcount=self._rowcount) + + def commit(self) -> None: + self.committed = True + + def rollback(self) -> None: + self.rolled_back = True + + @property + def update_statements(self) -> list[str]: + return [sql for sql, _ in self.executed if "UPDATE listings" in sql] + + +def _run( + db: _FakeDB, + monkeypatch: pytest.MonkeyPatch, + **kwargs: Any, +) -> dict[str, int]: + monkeypatch.setattr(task_mod.runs_mod, "mark_done", lambda *a, **k: None) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + return task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 1, + listing_source=kwargs.pop("listing_source", "avito"), + ttl_days=kwargs.pop("ttl_days", 10), + **kwargs, + ) + + +# ── Исторический случай: 17 суток бана Авито ────────────────────────────────── + + +def test_gate_blocks_every_day_of_the_17_day_avito_ban(monkeypatch: pytest.MonkeyPatch) -> None: + """Ни одни сутки провала 08.07-29.07 не должны пропустить деактивацию.""" + for day, confirmations in _AVITO_BAN_DAYS.items(): + db = _FakeDB(confirmations=confirmations) + out = _run(db, monkeypatch, min_confirmations=_AVITO_MIN_CONFIRMATIONS) + assert out["skipped_unhealthy"] == 1, f"{day}: гейт пропустил провальные сутки" + assert out["deactivated"] == 0, f"{day}: деактивировано ненулевое количество" + assert db.update_statements == [], f"{day}: UPDATE listings всё-таки исполнился" + + +def test_gate_passes_every_healthy_avito_day(monkeypatch: pytest.MonkeyPatch) -> None: + """Здоровые сутки порог 1500 не блокирует — гейт не ломает штатную работу.""" + for day, confirmations in _AVITO_HEALTHY_DAYS.items(): + db = _FakeDB(confirmations=confirmations) + out = _run(db, monkeypatch, min_confirmations=_AVITO_MIN_CONFIRMATIONS) + assert "skipped_unhealthy" not in out, f"{day}: гейт заблокировал здоровые сутки" + assert out["deactivated"] == 137, f"{day}: деактивация не исполнилась" + assert len(db.update_statements) == 1, f"{day}: UPDATE listings не исполнился" + + +def test_threshold_separates_ban_from_health() -> None: + """Порог лежит строго между максимумом провала и минимумом здоровых суток.""" + assert max(_AVITO_BAN_DAYS.values()) < _AVITO_MIN_CONFIRMATIONS + assert min(_AVITO_HEALTHY_DAYS.values()) > _AVITO_MIN_CONFIRMATIONS + + +def test_ban_window_damage_matches_issue_number() -> None: + """Ущерб исторического случая — 9 033 строки (#2659), гейт спасает их все.""" + assert sum(_AVITO_BAN_DEACTIVATED) == 9033 + + +# ── Контракт гейта ──────────────────────────────────────────────────────────── + + +def test_gate_disabled_by_default_keeps_old_behaviour(monkeypatch: pytest.MonkeyPatch) -> None: + """min_confirmations=0 -> ни одного лишнего запроса, поведение как до #2659.""" + db = _FakeDB(confirmations=0) + out = _run(db, monkeypatch) + assert out == {"deactivated": 137} + assert len(db.executed) == 1 + + +def test_gate_reports_confirmations_when_passing(monkeypatch: pytest.MonkeyPatch) -> None: + """Прошедший гейт прогон всё равно пишет замер — счётчик виден оператору.""" + db = _FakeDB(confirmations=4000) + out = _run(db, monkeypatch, min_confirmations=1500) + assert out["confirmations"] == 4000 + assert out["deactivated"] == 137 + + +def test_gate_runs_before_any_write(monkeypatch: pytest.MonkeyPatch) -> None: + """SELECT-проверка идёт ПЕРВОЙ: деактивация необратима, откат после неё не спасает.""" + db = _FakeDB(confirmations=4000) + _run(db, monkeypatch, min_confirmations=1500) + assert "SELECT count(*)" in db.executed[0][0] + assert "UPDATE listings" in db.executed[1][0] + + +def test_blocked_run_does_not_commit(monkeypatch: pytest.MonkeyPatch) -> None: + db = _FakeDB(confirmations=10) + _run(db, monkeypatch, min_confirmations=1500) + assert db.committed is False + assert db.rolled_back is True + + +def test_blocked_run_is_finalised_as_done(monkeypatch: pytest.MonkeyPatch) -> None: + """Пропущенный прогон закрывается mark_done, а не висит 'running' до zombie-жатвы.""" + marked: dict[str, Any] = {} + monkeypatch.setattr( + task_mod.runs_mod, + "mark_done", + lambda _db, run_id, counters: marked.update(run_id=run_id, counters=dict(counters)), + ) + monkeypatch.setattr(task_mod.runs_mod, "mark_failed", lambda *a, **k: None) + db = _FakeDB(confirmations=10) + task_mod.deactivate_stale_listings( + db, # type: ignore[arg-type] + 77, + listing_source="avito", + ttl_days=10, + min_confirmations=1500, + ) + assert marked["run_id"] == 77 + assert marked["counters"]["skipped_unhealthy"] == 1 + assert marked["counters"]["deactivated"] == 0 + + +def test_gate_measures_same_slice_as_update(monkeypatch: pytest.MonkeyPatch) -> None: + """Срез гейта совпадает со срезом UPDATE: тот же source и те же сегменты.""" + db = _FakeDB(confirmations=4000) + _run(db, monkeypatch, segments=["vtorichka"], min_confirmations=500) + health_sql, health_params = db.executed[0] + assert "ANY(CAST(:segments AS text[]))" in health_sql + assert health_params is not None + assert health_params["segments"] == ["vtorichka"] + assert health_params["listing_source"] == "avito" + + +def test_gate_uses_same_staleness_column_as_ttl(monkeypatch: pytest.MonkeyPatch) -> None: + """domklik считает свежесть по scraped_at (#2204) — гейт обязан мерить ту же колонку, + иначе bulk-touch по last_seen_at показал бы здоровье там, где сбора нет.""" + db = _FakeDB(confirmations=4000) + _run(db, monkeypatch, staleness_column="scraped_at", min_confirmations=200) + health_sql = db.executed[0][0] + assert "scraped_at" in health_sql + assert "last_seen_at" not in health_sql + + +def test_gate_rejects_invalid_staleness_column(monkeypatch: pytest.MonkeyPatch) -> None: + """Whitelist колонки работает и на пути гейта — интерполяции чужого имени нет.""" + db = _FakeDB(confirmations=4000) + with pytest.raises(ValueError): + _run(db, monkeypatch, staleness_column="is_active", min_confirmations=500) + assert db.executed == [] + + +def test_confirmations_sql_is_psycopg_v3_safe() -> None: + sql = str(task_mod._build_confirmations_sql("last_seen_at", with_segments=True).text) + assert "CAST(:health_window_days || ' days' AS interval)" in sql + assert not re.search(r":\w+::", sql) + assert "UPDATE" not in sql.upper() + assert "DELETE" not in sql.upper() + + +def test_default_min_confirmations_is_a_safety_net_not_zero() -> None: + """Незасеянное расписание получает страховку, а не «деактивируй вслепую».""" + assert task_mod.DEFAULT_MIN_CONFIRMATIONS > 0 + + +def test_handler_wires_min_confirmations_from_schedule_params() -> None: + """Читаем исходник файлом: product_handlers тянет scraper_kit, которого в + юнит-окружении может не быть, а проверяем мы проводку, а не импорт.""" + handlers = Path(__file__).resolve().parents[1] / "app" / "services" / "product_handlers.py" + src = handlers.read_text("utf-8") + job = src.split("async def _job_deactivate_stale")[1].split("\nasync def ")[0] + assert 'params.get("min_confirmations", DEFAULT_MIN_CONFIRMATIONS)' in job + assert "min_confirmations=min_confirmations" in job + + +# ── Миграция 219 ────────────────────────────────────────────────────────────── + + +def test_migration_219_exists() -> None: + assert _MIGRATION_219.is_file(), f"missing migration: {_MIGRATION_219}" + + +def test_migration_219_seeds_all_four_schedules() -> None: + sql = _MIGRATION_219.read_text("utf-8") + for source in ( + "deactivate_stale_avito", + "deactivate_stale_cian", + "deactivate_stale_yandex", + "deactivate_stale_domklik", + ): + assert f"'{source}'" in sql, f"{source} без порога — деактивирует вслепую" + + +def test_migration_219_avito_threshold_catches_the_ban() -> None: + """Порог авито должен быть выше максимума провальных суток (970).""" + sql = _MIGRATION_219.read_text("utf-8") + avito_block = sql.split("WHERE source = 'deactivate_stale_avito'")[0] + match = re.findall(r"'min_confirmations',\s*(\d+)", avito_block) + assert match, "порог авито не найден в миграции" + assert int(match[-1]) == _AVITO_MIN_CONFIRMATIONS + assert int(match[-1]) > max(_AVITO_BAN_DAYS.values()) + + +def test_migration_219_is_transactional_and_idempotent() -> None: + sql = _MIGRATION_219.read_text("utf-8") + assert "BEGIN;" in sql + assert "COMMIT;" in sql + # Повторный прогон не затирает подкрученное оператором значение. + assert sql.count("NOT default_params ? 'min_confirmations'") == 3 + + +def test_migration_219_touches_only_deactivate_schedules() -> None: + sql = _MIGRATION_219.read_text("utf-8") + for line in sql.splitlines(): + if line.strip().startswith("WHERE source"): + assert "deactivate_stale_" in line + + +def test_migration_219_no_psycopg_trap() -> None: + sql = _MIGRATION_219.read_text("utf-8") + assert not re.search(r":\w+::", sql)