fix(tradein): TTL-деактивация не исполняется, пока сбор по источнику лежит (#2659) #2710
4 changed files with 509 additions and 1 deletions
|
|
@ -216,12 +216,19 @@ async def _job_deactivate_stale(
|
|||
db: Session, run_id: int, params: dict[str, Any], ctx: SchedulerContext
|
||||
) -> None:
|
||||
from app.core.config import settings as _settings
|
||||
from app.tasks.deactivate_stale_avito import deactivate_stale_listings
|
||||
from app.tasks.deactivate_stale_avito import (
|
||||
DEFAULT_MIN_CONFIRMATIONS,
|
||||
deactivate_stale_listings,
|
||||
)
|
||||
|
||||
listing_source: str = params.get("listing_source", "avito")
|
||||
ttl_days: int = params.get("ttl_days", _settings.avito_stale_ttl_days)
|
||||
segments: list[str] | None = params.get("segments")
|
||||
staleness_column: str = params.get("staleness_column", "last_seen_at")
|
||||
# Гейт по здоровью сбора (#2659) включён по умолчанию: незасеянное расписание
|
||||
# получает страховочный порог, а не «деактивируй вслепую». Посчитанные по
|
||||
# источнику пороги приходят из default_params (миграция 219).
|
||||
min_confirmations: int = params.get("min_confirmations", DEFAULT_MIN_CONFIRMATIONS)
|
||||
|
||||
loop = asyncio.get_event_loop()
|
||||
await loop.run_in_executor(
|
||||
|
|
@ -233,6 +240,7 @@ async def _job_deactivate_stale(
|
|||
ttl_days=ttl_days,
|
||||
segments=segments,
|
||||
staleness_column=staleness_column,
|
||||
min_confirmations=min_confirmations,
|
||||
),
|
||||
)
|
||||
|
||||
|
|
|
|||
|
|
@ -82,6 +82,70 @@ _STALE_SNAPSHOT_TAIL = """
|
|||
"""
|
||||
|
||||
|
||||
# ── Гейт по здоровью сбора (#2659) ────────────────────────────────────────────
|
||||
# TTL отвечает на вопрос «объявление сняли?», а меряет «мы его давно не видели».
|
||||
# Пока обход здоров, разница мала. Когда обход лёг — разница равна всему инвентарю.
|
||||
#
|
||||
# Замер на проде, из-за которого этот гейт существует. Авито 10.07-26.07.2026:
|
||||
# 17 суток подряд без единой успешно собранной страницы, TTL=10 снял за этот отрезок
|
||||
# 9 033 строки; 1 270 из них потом доказанно вернулись живыми (снимки
|
||||
# listing_source_snapshots + текущий last_seen_at) — сбор восстановился, и объявления
|
||||
# оказались на месте. То есть «мы не смогли зайти» было прочитано как «объявление снято».
|
||||
#
|
||||
# ПОЧЕМУ НЕ ПО СТАТУСУ ban. Соблазн взять scrape_runs.status='banned' — ловушка:
|
||||
# Яндекс 18.07-30.07 — 5 прогонов в сутки, ВСЕ 'done', НОЛЬ 'banned', total_seen=0
|
||||
# 13 суток подряд (снято ~839 строк vtorichka);
|
||||
# Домклик 20.07-30.07 — то же самое, 11 суток 'done' с total_seen=0, а 02.08 TTL
|
||||
# снял 6 131 строку разом (см. комментарий про 'stale' выше).
|
||||
# Оба провала для ban-детектора невидимы. Поэтому здоровье меряем НЕ статусом прогона,
|
||||
# а результатом: сколько строк источник реально подтвердил свежими за последние сутки.
|
||||
#
|
||||
# МЕТРИКА: count(*) по той же колонке свежести, что и сам TTL (last_seen_at или
|
||||
# scraped_at) и по тому же срезу source+segment, что и UPDATE. Одна колонка на обе
|
||||
# стороны — гейт нельзя обмануть bulk-touch'ем, который не двигает scraped_at (#2204).
|
||||
#
|
||||
# ПОРОГ. Ряд «подтверждений за 3 суток» по дням (восстановлен из listing_source_snapshots):
|
||||
# avito здоровые сутки 3542..6079, провал 10.07-26.07 — 0..970 → порог 1500;
|
||||
# yandex vtorichka здоровые 897..2206, провал — 0 → порог 500;
|
||||
# cian vtorichka 748..4329, провала не было → порог 500;
|
||||
# domklik по scraped_at сейчас 62/3 суток (сбор фактически стоит) → порог 200.
|
||||
# Пороги живут в default_params расписания (миграция 219), здесь только страховка
|
||||
# на случай незасеянного расписания. Асимметрия цены ошибки намеренная: пропущенная
|
||||
# деактивация чинится следующим прогоном, ложная — только повторным сбором, которого
|
||||
# может не быть. Поэтому при сомнении — пропускаем прогон.
|
||||
#
|
||||
# ПОТОЛОК: окно 3 суток годится, пока свипы источника ходят не реже чем раз в 3 дня.
|
||||
# Источник с более редкой каденцией будет блокироваться всегда — тогда окно нужно
|
||||
# растить до каденции, а не понижать порог.
|
||||
_HEALTH_WINDOW_DAYS = 3
|
||||
|
||||
# Страховка для расписаний без явного min_confirmations в default_params: ловит
|
||||
# полный ноль и близкое к нулю, но НЕ ловит частичный провал вроде avito 936-970 —
|
||||
# для этого нужен посчитанный по источнику порог из миграции 219.
|
||||
DEFAULT_MIN_CONFIRMATIONS = 500
|
||||
|
||||
_CONFIRMATIONS_SEGMENT_FILTER = "\n AND listing_segment = ANY(CAST(:segments AS text[]))"
|
||||
|
||||
|
||||
def _build_confirmations_sql(staleness_column: str, *, with_segments: bool) -> Any:
|
||||
"""SELECT count(*) подтверждённых за окно строк — тот же срез, что и у UPDATE.
|
||||
|
||||
staleness_column уже прошёл whitelist-проверку в deactivate_stale_listings.
|
||||
Значения (:listing_source, :health_window_days, :segments) — param-binding,
|
||||
psycopg v3 safe (CAST(... AS ...), никаких :param::type).
|
||||
"""
|
||||
segment_filter = _CONFIRMATIONS_SEGMENT_FILTER if with_segments else ""
|
||||
return text(
|
||||
f"""
|
||||
SELECT count(*)
|
||||
FROM listings
|
||||
WHERE source = :listing_source
|
||||
AND {staleness_column}
|
||||
> NOW() - CAST(:health_window_days || ' days' AS interval){segment_filter}
|
||||
"""
|
||||
)
|
||||
|
||||
|
||||
def _build_all_segments_sql(staleness_column: str) -> Any:
|
||||
"""UPDATE без фильтра по сегменту: все сегменты для данного source.
|
||||
|
||||
|
|
@ -143,6 +207,8 @@ def deactivate_stale_listings(
|
|||
ttl_days: int,
|
||||
segments: list[str] | None = None,
|
||||
staleness_column: str = "last_seen_at",
|
||||
min_confirmations: int = 0,
|
||||
health_window_days: int = _HEALTH_WINDOW_DAYS,
|
||||
) -> dict[str, int]:
|
||||
"""Пометить is_active=false объявления, чья свежесть старше ttl_days дней.
|
||||
|
||||
|
|
@ -157,12 +223,20 @@ def deactivate_stale_listings(
|
|||
last_seen_at. Для domklik (#2204) — scraped_at: нетрекаемый bulk-touch
|
||||
двигает last_seen_at всем строкам одним timestamp, поэтому честная
|
||||
свежесть = scraped_at (двигается только реальным скрейпом).
|
||||
min_confirmations: гейт по здоровью сбора (#2659). Сколько строк источник
|
||||
должен был подтвердить свежими за health_window_days суток, чтобы
|
||||
деактивации вообще разрешалось исполниться. 0 -> гейт выключен (так
|
||||
вызывают старые тесты и совместимая обёртка); реальные значения приходят
|
||||
из default_params расписания, см. миграцию 219 и комментарий выше.
|
||||
health_window_days: окно подтверждений для гейта, суток. Дефолт 3.
|
||||
|
||||
Sync (вызывается scheduler-триггером в executor, как snapshot_listing_sources).
|
||||
Один statement в транзакции: UPDATE флага + снимок 'stale' в listings_snapshots
|
||||
(data-modifying CTE, #2674). Финализирует scrape_runs (mark_done / mark_failed).
|
||||
|
||||
Returns {"deactivated": N} -- количество обновлённых строк (1:1 со снимками).
|
||||
Если гейт не пропустил прогон: {"deactivated": 0, "confirmations": N,
|
||||
"skipped_unhealthy": 1} и НИ ОДНА строка не тронута.
|
||||
|
||||
Raises:
|
||||
ValueError: если staleness_column не входит в whitelist (проверка ДО SQL,
|
||||
|
|
@ -180,6 +254,44 @@ def deactivate_stale_listings(
|
|||
f"allowed: {sorted(_ALLOWED_STALENESS_COLUMNS)}"
|
||||
)
|
||||
|
||||
# Гейт по здоровью сбора (#2659) — ДО любого UPDATE. Деактивация необратима
|
||||
# на практике (вернуть «живость» может только повторный сбор), поэтому
|
||||
# проверяем ПЕРЕД записью, а не откатываем после.
|
||||
if min_confirmations > 0:
|
||||
health_params: dict[str, Any] = {
|
||||
"listing_source": listing_source,
|
||||
"health_window_days": health_window_days,
|
||||
}
|
||||
if segments is not None:
|
||||
health_params["segments"] = segments
|
||||
confirmations = (
|
||||
db.execute(
|
||||
_build_confirmations_sql(staleness_column, with_segments=segments is not None),
|
||||
health_params,
|
||||
).scalar()
|
||||
or 0
|
||||
)
|
||||
counters["confirmations"] = int(confirmations)
|
||||
if confirmations < min_confirmations:
|
||||
counters["skipped_unhealthy"] = 1
|
||||
# Ничего не писали (был только SELECT) — rollback закрывает транзакцию
|
||||
# чисто, чтобы mark_done стартовал со своей.
|
||||
db.rollback()
|
||||
runs_mod.mark_done(db, run_id, counters)
|
||||
logger.warning(
|
||||
"deactivate_stale source=%s run_id=%d SKIPPED: сбор нездоров — "
|
||||
"подтверждений за %d сут %d < порога %d "
|
||||
"(segments=%r, staleness_column=%s); ни одна строка не тронута",
|
||||
listing_source,
|
||||
run_id,
|
||||
health_window_days,
|
||||
confirmations,
|
||||
min_confirmations,
|
||||
segments,
|
||||
staleness_column,
|
||||
)
|
||||
return counters
|
||||
|
||||
# segments is None -> все сегменты (поведение avito). segments=[...] -> только
|
||||
# перечисленные сегменты. Используем `is not None` (НЕ truthy): пустой список []
|
||||
# означает "ни один сегмент" (= ANY(ARRAY[]) ничего не матчит, деактивирует 0),
|
||||
|
|
|
|||
|
|
@ -0,0 +1,67 @@
|
|||
-- 219_deactivate_stale_health_gate.sql
|
||||
-- Пороги гейта здоровья сбора для TTL-деактивации (#2659).
|
||||
--
|
||||
-- ЗАЧЕМ. TTL отвечает на вопрос «объявление сняли?», а меряет «мы его давно не
|
||||
-- видели». Пока обход здоров, разница мала; когда обход лёг — разница равна всему
|
||||
-- инвентарю. Прод, авито 10.07-26.07.2026: 17 суток подряд без единой собранной
|
||||
-- страницы, TTL=10 снял 9 033 строки, из них 1 270 доказанно вернулись живыми,
|
||||
-- как только сбор восстановился (сверка listing_source_snapshots с текущим
|
||||
-- last_seen_at). Код гейта — app/tasks/deactivate_stale_avito.py.
|
||||
--
|
||||
-- ПОЧЕМУ НЕ ПО СТАТУСУ ПРОГОНА. Ban-детектор эти провалы НЕ ловит:
|
||||
-- yandex 18.07-30.07 — 5 прогонов в сутки, ВСЕ 'done', НОЛЬ 'banned',
|
||||
-- total_seen = 0 тринадцать суток подряд;
|
||||
-- domklik 20.07-30.07 — 11 суток 'done' с total_seen = 0, а 02.08 TTL=14
|
||||
-- снял 6 131 строку разом.
|
||||
-- Поэтому здоровье меряется результатом (сколько строк источник реально подтвердил
|
||||
-- свежими за 3 суток), а не статусом прогона.
|
||||
--
|
||||
-- ОТКУДА ЧИСЛА. Ряд «подтверждений за 3 суток» по дням восстановлен из
|
||||
-- listing_source_snapshots (снимок last_seen_at на каждую дату), срез совпадает
|
||||
-- со срезом соответствующего UPDATE (source + segments + та же колонка свежести):
|
||||
-- avito (все сегменты, last_seen_at): здоровые сутки 3542..6079,
|
||||
-- провал 08.07-29.07 — 0..970 -> 1500
|
||||
-- yandex (vtorichka, last_seen_at): здоровые 897..2206, провал 0 -> 500
|
||||
-- cian (vtorichka, last_seen_at): 748..4329, провалов не было -> 500
|
||||
-- domklik (все сегменты, scraped_at): сейчас 62 за 3 суток, сбор
|
||||
-- фактически стоит -> 200
|
||||
-- Каждый порог лежит между максимумом провала и минимумом здоровых суток:
|
||||
-- авито 970 < 1500 < 3542 — исторический случай ловится с запасом в обе стороны.
|
||||
--
|
||||
-- domklik ЗАБЛОКИРУЕТСЯ СРАЗУ, и это верный исход, а не сбой миграции: источник
|
||||
-- подтверждает ~50 строк в сутки, TTL по нему уже один раз (02.08) снёс инвентарь
|
||||
-- целиком. Пока сбор не восстановлен, деактивации там нечего подтверждать; протухшие
|
||||
-- строки закрываются фильтром свежести на стороне чтения (#2656), а не TTL.
|
||||
--
|
||||
-- Цена ошибки асимметрична: пропущенная деактивация чинится следующим прогоном,
|
||||
-- ложная — только повторным сбором, которого может не быть. Пороги поэтому
|
||||
-- смещены в сторону «пропустить прогон».
|
||||
--
|
||||
-- ЗАВИСИМОСТИ: 052_scrape_schedules.sql (таблица + UNIQUE(source)),
|
||||
-- 090/115/160 (сами расписания deactivate_stale_*).
|
||||
-- ТОЛЬКО данные (UPDATE default_params), DDL нет.
|
||||
-- Идемпотентность + уважение к ручной настройке: ключ проставляется лишь там, где
|
||||
-- его ещё нет, поэтому повторный прогон файла не затирает подкрученное оператором
|
||||
-- значение. Снять гейт вручную: min_confirmations = 0.
|
||||
|
||||
BEGIN;
|
||||
|
||||
UPDATE scrape_schedules
|
||||
SET default_params = default_params || jsonb_build_object('min_confirmations', 1500),
|
||||
updated_at = NOW()
|
||||
WHERE source = 'deactivate_stale_avito'
|
||||
AND NOT default_params ? 'min_confirmations';
|
||||
|
||||
UPDATE scrape_schedules
|
||||
SET default_params = default_params || jsonb_build_object('min_confirmations', 500),
|
||||
updated_at = NOW()
|
||||
WHERE source IN ('deactivate_stale_cian', 'deactivate_stale_yandex')
|
||||
AND NOT default_params ? 'min_confirmations';
|
||||
|
||||
UPDATE scrape_schedules
|
||||
SET default_params = default_params || jsonb_build_object('min_confirmations', 200),
|
||||
updated_at = NOW()
|
||||
WHERE source = 'deactivate_stale_domklik'
|
||||
AND NOT default_params ? 'min_confirmations';
|
||||
|
||||
COMMIT;
|
||||
321
tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py
Normal file
321
tradein-mvp/backend/tests/test_deactivate_stale_health_gate.py
Normal file
|
|
@ -0,0 +1,321 @@
|
|||
"""Гейт по здоровью сбора для 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)
|
||||
Loading…
Add table
Reference in a new issue