All checks were successful
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / test (push) Successful in 3m47s
Deploy Trade-In / build-backend (push) Successful in 1m43s
Deploy Trade-In / deploy (push) Successful in 2m1s
Deploy Trade-In / deploy-status (push) Successful in 1s
Deploy Trade-In / perimeter-smoke (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
311 lines
14 KiB
Python
311 lines
14 KiB
Python
"""Потолок объёма снятия за один прогон -- аварийный, наблюдательный (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
|