"""Потолок объёма снятия за один прогон -- аварийный, наблюдательный (PR-B, #2659 продолжение). У джобы деактивации не было НИ ОДНОГО ограничителя ОБЪЁМА: UPDATE идёт одним statement'ом без LIMIT, counters["deactivated"] = result.rowcount. Долевой порог ("не больше N% активного пула") откалибровать по историческим данным нельзя -- исторические снятия дают доли 40%/97%/317%/1652% от восстановленного пула, числа не образуют осмысленного ряда (гранулярность listings vs listing_sources разная). Поэтому: preflight count(*) ПО ТОМУ ЖЕ предикату, что и UPDATE, пишет deactivation_candidates/active_pool/deactivated_pct В КАЖДЫЙ прогон (наблюдение, задел под будущую калибровку долевого порога), а блокирует ТОЛЬКО абсолютный аварийный порог max_deactivated (дефолт 15000 -- запас 1.6x над историческим максимумом легитимного снятия 9300, avito 2026-06-06). """ from __future__ import annotations import os 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 # ── Фейковая сессия (тот же контракт, что test_deactivate_stale_floor_degradation.py) ── 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_cap_trips_when_candidates_exceed_max_deactivated(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(deactivation_candidates=20_000, active_pool=100_000) out = _run(db, monkeypatch, max_deactivated=15000) assert out["skipped_cap_exceeded"] == 1 assert out["deactivation_candidates"] == 20000 assert out["active_pool"] == 100000 assert out["deactivated"] == 0 assert db.update_query[0] == "" assert db.committed is False assert db.rolled_back is True def test_cap_does_not_trip_below_threshold(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(deactivation_candidates=9300, active_pool=100_000, rowcount=9300) out = _run(db, monkeypatch, max_deactivated=15000) assert "skipped_cap_exceeded" not in out assert out["deactivation_candidates"] == 9300 assert out["deactivated"] == 9300 assert db.update_query[0] != "" assert db.committed is True def test_cap_trips_exactly_above_threshold_not_at_it(monkeypatch: pytest.MonkeyPatch) -> None: """== max_deactivated не триггерит, только > (строго больше).""" db = _FakeDB(deactivation_candidates=15000, active_pool=100_000, rowcount=15000) out = _run(db, monkeypatch, max_deactivated=15000) assert "skipped_cap_exceeded" not in out db2 = _FakeDB(deactivation_candidates=15001, active_pool=100_000) out2 = _run(db2, monkeypatch, max_deactivated=15000) assert out2["skipped_cap_exceeded"] == 1 def test_cap_check_runs_before_update(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(deactivation_candidates=20_000) _run(db, monkeypatch, max_deactivated=15000) assert "update" not in db.executed_kinds def test_blocked_run_is_finalised_as_done(monkeypatch: pytest.MonkeyPatch) -> None: 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(deactivation_candidates=20_000) task_mod.deactivate_stale_listings( db, # type: ignore[arg-type] 55, listing_source="avito", ttl_days=10, revisit_floor_quantile=task_mod.DEFAULT_REVISIT_FLOOR_QUANTILE, max_deactivated=15000, ) assert marked["run_id"] == 55 assert marked["counters"]["skipped_cap_exceeded"] == 1 assert marked["counters"]["deactivated"] == 0 # ── Наблюдательные counters (пишутся ВСЕГДА, не только при abort) ──────────────── def test_deactivated_pct_is_computed(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(deactivation_candidates=5000, active_pool=10_000, rowcount=5000) out = _run(db, monkeypatch) assert out["deactivated_pct"] == 50 def test_deactivated_pct_zero_when_active_pool_empty(monkeypatch: pytest.MonkeyPatch) -> None: """active_pool=0 -> деление защищено, не ZeroDivisionError.""" db = _FakeDB(deactivation_candidates=0, active_pool=0, rowcount=0) out = _run(db, monkeypatch) assert out["deactivated_pct"] == 0 assert out["active_pool"] == 0 def test_candidates_and_pool_written_on_every_healthy_run(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB(deactivation_candidates=137, active_pool=5000, rowcount=137) out = _run(db, monkeypatch) assert out["deactivation_candidates"] == 137 assert out["active_pool"] == 5000 assert db.update_query[0] != "" # ── Предикат preflight COUNT совпадает с предикатом UPDATE ─────────────────────── def test_candidates_predicate_matches_update_predicate_all_segments() -> None: update_sql = str(task_mod._build_all_segments_sql("last_seen_at").text) count_sql = str(task_mod._build_all_segments_candidates_count_sql("last_seen_at").text) for fragment in ( "source = :listing_source", "is_active = true", "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", ): assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" assert "listing_segment" not in update_sql assert "listing_segment" not in count_sql def test_candidates_predicate_matches_update_predicate_segments() -> None: update_sql = str(task_mod._build_segments_sql("last_seen_at").text) count_sql = str(task_mod._build_segments_candidates_count_sql("last_seen_at").text) for fragment in ( "source = :listing_source", "is_active = true", "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", "listing_segment = ANY(CAST(:segments AS text[]))", ): assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" def test_candidates_predicate_matches_update_predicate_null_segment() -> None: update_sql = str(task_mod._build_null_segment_sql("last_seen_at").text) count_sql = str(task_mod._build_null_segment_candidates_count_sql("last_seen_at").text) for fragment in ( "source = :listing_source", "is_active = true", "last_seen_at < NOW() - CAST(:ttl_days || ' days' AS interval)", "listing_segment IS NULL", ): assert fragment in update_sql, f"{fragment!r} missing from UPDATE sql" assert fragment in count_sql, f"{fragment!r} missing from candidates-count sql" def test_candidates_count_sql_uses_the_effective_ttl_days_param_name() -> None: """Preflight использует :ttl_days -- caller обязан передать effective_ttl_days под этим именем (см. deactivate_stale_listings), а не сырой ttl_days.""" sql = str(task_mod._build_all_segments_candidates_count_sql("last_seen_at").text) assert ":ttl_days" in sql def test_active_pool_sql_has_no_segment_or_ttl_filter() -> None: """active_pool -- весь активный пул source, без сегмента и без TTL (см. докстринг deactivate_stale_listings, Returns): денормализатор пула, а не срез UPDATE.""" sql = str(task_mod._ACTIVE_POOL_SQL.text) assert "is_active = true" in sql assert "listing_segment" not in sql assert ":ttl_days" not in sql # ── Параметры / контракт ────────────────────────────────────────────────────────── def test_default_max_deactivated_is_15000() -> None: assert task_mod.DEFAULT_MAX_DEACTIVATED == 15000 def test_default_does_not_reject_any_known_legit_historical_run() -> None: """Исторический максимум легитимного снятия (avito, 2026-06-06) и следующие по убыванию -- ни один не должен упереться в дефолтный потолок.""" known_legit = [9300, 6531, 6131, 4959, 3909] assert max(known_legit) < task_mod.DEFAULT_MAX_DEACTIVATED def test_max_deactivated_rejects_non_positive(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="max_deactivated"): _run(db, monkeypatch, max_deactivated=0) assert db.executed == [] def test_max_deactivated_rejects_bool(monkeypatch: pytest.MonkeyPatch) -> None: db = _FakeDB() with pytest.raises(ValueError, match="max_deactivated"): _run(db, monkeypatch, max_deactivated=True) assert db.executed == [] def test_handler_wires_max_deactivated_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("max_deactivated", DEFAULT_MAX_DEACTIVATED)' in job assert "max_deactivated=max_deactivated" in job