"""Гейт деградации пола переобхода (PR-B, #2659 продолжение). PR-A (fb38d657) перевёл поиск предшественника в поле переобхода на LATERAL и добавил counters["floor_n_pairs"] -- ЧИСТО наблюдательный счётчик размера выборки, из которой percentile_disc берёт квантиль. Этот файл проверяет PR-B: тот же floor_n_pairs теперь ГЕЙТИТ прогон, если выборка вырождена -- либо ниже абсолютного порога (min_floor_pairs), либо упала более чем в floor_drop_ratio раз относительно предыдущего УСПЕШНОГО прогона того же расписания. """ from __future__ import annotations import os import re 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 # ── Фейковая сессия ───────────────────────────────────────────────────────────── # Различает ВСЕ SELECT-варианты, которые видит deactivate_stale_listings при # revisit_floor_quantile > 0: percentile_disc (пол), previous floor_n_pairs # (scrape_runs, PR-B), floor_n_pairs (LATERAL count, PR-A), confirmations (гейт # здоровья), deactivation_candidates / active_pool (потолок объёма, PR-B). # # Порядок веток важен: percentile_disc и count(*) LATERAL-запрос ОБА содержат # "JOIN LATERAL" -- percentile_disc проверяется первой веткой. LATERAL-запрос # также содержит "health_window_days" (дважды) -- "JOIN LATERAL" проверяется # раньше этой ветки. UPDATE содержит "ttl_days", но не "SELECT count(*)" -- # ветка candidates требует ОБА маркера, поэтому UPDATE в неё не попадает. class _FakeResult: def __init__(self, rowcount: int = 0, scalar_value: Any = None) -> None: self.rowcount = rowcount self._scalar = scalar_value def scalar(self) -> Any: return self._scalar class _FakeDB: def __init__( self, *, floor_days: float | None = 50.0, floor_n_pairs: int = 5000, prev_floor_n_pairs: int | None = 5000, confirmations: int = 10_000, deactivation_candidates: int | None = None, active_pool: int = 100_000, rowcount: int = 137, ) -> None: self._floor = floor_days self._floor_n_pairs = floor_n_pairs self._prev_floor_n_pairs = prev_floor_n_pairs self._confirmations = confirmations self._deactivation_candidates = ( rowcount if deactivation_candidates is None else deactivation_candidates ) self._active_pool = active_pool 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 "percentile_disc" in sql: return _FakeResult(scalar_value=self._floor) if "FROM scrape_runs prev" in sql: return _FakeResult(scalar_value=self._prev_floor_n_pairs) if "JOIN LATERAL" in sql: return _FakeResult(scalar_value=self._floor_n_pairs) if "health_window_days" in sql: return _FakeResult(scalar_value=self._confirmations) if "SELECT count(*)" in sql and "ttl_days" in sql: return _FakeResult(scalar_value=self._deactivation_candidates) if "SELECT count(*)" in sql: return _FakeResult(scalar_value=self._active_pool) return _FakeResult(rowcount=self._rowcount) def commit(self) -> None: self.committed = True def rollback(self) -> None: self.rolled_back = True @property def update_query(self) -> tuple[str, dict[str, Any] | None]: return next((e for e in self.executed if "UPDATE listings" in e[0]), ("", None)) @property def executed_kinds(self) -> list[str]: kinds = [] for sql, _ in self.executed: if "percentile_disc" in sql: kinds.append("floor") elif "FROM scrape_runs prev" in sql: kinds.append("prev_floor_n_pairs") elif "JOIN LATERAL" in sql: kinds.append("floor_n_pairs") elif "health_window_days" in sql: kinds.append("confirmations") elif "SELECT count(*)" in sql and "ttl_days" in sql: kinds.append("candidates") elif "SELECT count(*)" in sql: kinds.append("active_pool") elif "UPDATE listings" in sql: kinds.append("update") else: kinds.append("unknown") return kinds 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) kwargs.setdefault("revisit_floor_quantile", task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE) 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, ) # ── Абсолютный порог ──────────────────────────────────────────────────────────── def test_gate_trips_on_absolute_floor(monkeypatch: pytest.MonkeyPatch) -> None: """floor_n_pairs ниже min_floor_pairs -> skipped_floor_degraded, ни одна строка не тронута.""" db = _FakeDB(floor_n_pairs=10, prev_floor_n_pairs=None) out = _run(db, monkeypatch, min_floor_pairs=30) assert out["skipped_floor_degraded"] == 1 assert out["floor_n_pairs"] == 10 assert out["deactivated"] == 0 assert db.update_query[0] == "" assert db.committed is False assert db.rolled_back is True def test_absolute_floor_disabled_at_zero(monkeypatch: pytest.MonkeyPatch) -> None: """min_floor_pairs=0 -> абсолютная проверка отключена (тот же идиом, что у гейта здоровья).""" db = _FakeDB(floor_n_pairs=0, prev_floor_n_pairs=None) out = _run(db, monkeypatch, min_floor_pairs=0) assert "skipped_floor_degraded" not in out assert db.update_query[0] != "" # ── Относительный порог (падение в N раз) ─────────────────────────────────────── def test_gate_trips_on_ratio_drop(monkeypatch: pytest.MonkeyPatch) -> None: """floor_n_pairs упал в 10 раз против предыдущего успешного прогона (порог 5x).""" db = _FakeDB(floor_n_pairs=100, prev_floor_n_pairs=1000) out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) assert out["skipped_floor_degraded"] == 1 assert out["floor_n_pairs_previous"] == 1000 assert out["deactivated"] == 0 assert db.update_query[0] == "" assert db.committed is False assert db.rolled_back is True def test_gate_does_not_trip_at_exactly_the_ratio(monkeypatch: pytest.MonkeyPatch) -> None: """Падение РОВНО в floor_drop_ratio раз не триггерит -- задача требует «более чем».""" db = _FakeDB(floor_n_pairs=200, prev_floor_n_pairs=1000) # ровно 5x out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) assert "skipped_floor_degraded" not in out assert db.update_query[0] != "" def test_gate_ignores_ratio_when_no_previous_run(monkeypatch: pytest.MonkeyPatch) -> None: """Нет предыдущего успешного прогона (первый прогон после деплоя) -> относительная часть молчит, работает только абсолютный порог.""" db = _FakeDB(floor_n_pairs=50, prev_floor_n_pairs=None) out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) assert "skipped_floor_degraded" not in out assert "floor_n_pairs_previous" not in out assert db.update_query[0] != "" def test_gate_ignores_zero_previous(monkeypatch: pytest.MonkeyPatch) -> None: """Предыдущий прогон существует, но floor_n_pairs=0 в нём -- деление защищено guard'ом prev_floor_n_pairs > 0, ratio-часть молчит вместо ZeroDivisionError.""" db = _FakeDB(floor_n_pairs=0, prev_floor_n_pairs=0) out = _run(db, monkeypatch, min_floor_pairs=0, floor_drop_ratio=5.0) assert "skipped_floor_degraded" not in out assert out["floor_n_pairs_previous"] == 0 # ── Нормальный (не деградировавший) floor_n_pairs ─────────────────────────────── def test_gate_passes_with_healthy_floor_n_pairs(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(floor_n_pairs=5000, prev_floor_n_pairs=5200) out = _run(db, monkeypatch, min_floor_pairs=30, floor_drop_ratio=5.0) assert "skipped_floor_degraded" not in out assert out["floor_n_pairs"] == 5000 assert out["floor_n_pairs_previous"] == 5200 assert db.update_query[0] != "" assert db.committed is True def test_gate_disabled_when_revisit_floor_is_off(monkeypatch: pytest.MonkeyPatch) -> None: """revisit_floor_quantile=0 -> весь блок (пол + оба гейта PR-B) не исполняется вовсе.""" db = _FakeDB(floor_n_pairs=1, prev_floor_n_pairs=None) out = _run(db, monkeypatch, revisit_floor_quantile=0, min_floor_pairs=30) assert out == {"deactivated": 137} assert "skipped_floor_degraded" not in out assert db.executed_kinds == ["update"] # ── Abort -- ни одна строка не тронута ─────────────────────────────────────────── def test_gate_runs_before_update_and_before_the_volume_cap(monkeypatch: pytest.MonkeyPatch) -> None: """Деградация останавливает прогон ДО preflight-потолка (PR-B п.2) и ДО UPDATE.""" db = _FakeDB(floor_n_pairs=10, prev_floor_n_pairs=None) _run(db, monkeypatch, min_floor_pairs=30) kinds = db.executed_kinds assert "update" not in kinds assert "candidates" not in kinds assert "active_pool" not in kinds 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(floor_n_pairs=10, prev_floor_n_pairs=None) task_mod.deactivate_stale_listings( db, # type: ignore[arg-type] 99, listing_source="avito", ttl_days=10, revisit_floor_quantile=task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE, min_floor_pairs=30, ) assert marked["run_id"] == 99 assert marked["counters"]["skipped_floor_degraded"] == 1 assert marked["counters"]["deactivated"] == 0 # ── Параметры / контракт ───────────────────────────────────────────────────────── def test_default_min_floor_pairs_is_a_safety_net_not_zero() -> None: assert task_mod.DEFAULT_MIN_FLOOR_PAIRS > 0 def test_default_floor_drop_ratio_is_above_one() -> None: assert task_mod.DEFAULT_FLOOR_DROP_RATIO > 1 def test_min_floor_pairs_rejects_negative(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="min_floor_pairs"): _run(db, monkeypatch, min_floor_pairs=-1) assert db.executed == [] def test_min_floor_pairs_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="min_floor_pairs"): _run(db, monkeypatch, min_floor_pairs=True) assert db.executed == [] def test_floor_drop_ratio_rejects_below_one(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="floor_drop_ratio"): _run(db, monkeypatch, floor_drop_ratio=0.5) assert db.executed == [] def test_floor_drop_ratio_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="floor_drop_ratio"): _run(db, monkeypatch, floor_drop_ratio=True) assert db.executed == [] def test_previous_floor_n_pairs_sql_compares_by_scrape_runs_source() -> None: """Сравнение идёт по scrape_runs.source ТЕКУЩЕГО прогона (имя расписания), а НЕ по listing_source -- см. модульный докстринг у DEFAULT_MIN_FLOOR_PAIRS.""" sql = str(task_mod._PREVIOUS_FLOOR_N_PAIRS_SQL.text) assert "SELECT source FROM scrape_runs WHERE id = :run_id" in sql assert ":listing_source" not in sql assert "status = 'done'" in sql assert "id != :run_id" in sql assert not re.search(r":\w+::", sql) def test_handler_wires_floor_degradation_params_from_schedule() -> None: """Читаем исходник файлом (как соседние test_handler_wires_* в этом сьюте): product_handlers тянет scraper_kit, которого в юнит-окружении может не быть.""" from pathlib import Path 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_floor_pairs", DEFAULT_MIN_FLOOR_PAIRS)' in job assert "min_floor_pairs=min_floor_pairs" in job assert 'params.get("floor_drop_ratio", DEFAULT_FLOOR_DROP_RATIO)' in job assert "floor_drop_ratio=floor_drop_ratio" in job