"""Гейт по здоровью сбора для 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)