All checks were successful
Deploy Trade-In / changes (push) Successful in 10s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Successful in 2m56s
Deploy Trade-In / build-backend (push) Successful in 1m1s
Deploy Trade-In / deploy (push) Successful in 1m16s
321 lines
14 KiB
Python
321 lines
14 KiB
Python
"""Гейт по здоровью сбора для 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)
|