"""Пропуск расписания оставляет след — строку scrape_runs(status='skipped') (#2658). До #2658 планировщик пропускал наступившее окно НЕМО: logger + сдвиг next_run_at, ни строки прогона, ни изменения last_run_at. cian_history_backfill так простоял 37 дней на протухших куках Циана, и снаружи это выглядело как «всё по расписанию». Покрываем: 1. mark_skipped — INSERT новой строки и схлопывание подряд идущих одинаковых причин. 2. Все ПЯТЬ мест, которые раньше пропускали молча: kit `_claim_run` ×3 (already_running / concurrent_claim / running_appeared_under_lock), kit `scheduler_loop` (unknown_source) и продуктовый cian `pre_claim`. 3. Достижимость алерта: протухшие куки (load_session → None) дают ERROR-запись, которая в scraper-контейнере становится событием GlitchTip (event_level=ERROR), — раньше алерт стоял во второй ветке, куда на протухших куках исполнение не доходит. 4. Неломание монитора нулевых прогонов (#2625): 'skipped' — не «прогон вернул ноль лотов», обе alert-выборки его не видят и стрик им не прерывается. Без сети, без БД. """ from __future__ import annotations import logging import os from datetime import UTC, datetime, timedelta from typing import Any from unittest.mock import AsyncMock, MagicMock, patch import pytest os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") from scraper_kit.orchestration import runs as kit_runs from scraper_kit.orchestration import scheduler as kit_sched from scraper_kit.orchestration.scheduler import ( SKIP_ALREADY_RUNNING, SKIP_CONCURRENT_CLAIM, SKIP_RUNNING_UNDER_LOCK, SKIP_UNKNOWN_SOURCE, SchedulerContext, _claim_run, ) def _make_sched(source: str) -> dict[str, Any]: return { "id": 1, "source": source, "enabled": True, "window_start_hour": 2, "window_end_hour": 5, "default_params": {}, "last_run_id": None, "last_run_at": None, "next_run_at": None, } # ── 1. mark_skipped: INSERT + схлопывание ──────────────────────────────────── class _FakeSkipDB: """Session-мок для mark_skipped: UPDATE-ветка (схлопывание) vs INSERT-ветка.""" def __init__(self, *, collapse: bool) -> None: self._collapse = collapse self.statements: list[tuple[str, dict[str, Any]]] = [] self.committed = False def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: sql = str(stmt) self.statements.append((sql, params or {})) result = MagicMock() if "UPDATE scrape_runs r" in sql: result.fetchone.return_value = MagicMock(id=7) if self._collapse else None else: result.fetchone.return_value = MagicMock(id=8) return result def commit(self) -> None: self.committed = True def sql_of(self, needle: str) -> tuple[str, dict[str, Any]] | None: for sql, params in self.statements: if needle in sql: return sql, params return None def test_mark_skipped_inserts_row_with_machine_readable_reason() -> None: db = _FakeSkipDB(collapse=False) run_id = kit_runs.mark_skipped( db, source="cian_history_backfill", reason="cian_cookies_expired", details="протухли" ) assert run_id == 8 insert = db.sql_of("INSERT INTO scrape_runs") assert insert is not None, "новая строка прогона не создана" sql, params = insert assert "'skipped'" in sql assert params["source"] == "cian_history_backfill" assert params["reason"] == "cian_cookies_expired" # слаг, не человеческий текст assert "протухли" in params["counters"] assert db.committed is True def test_mark_skipped_collapses_consecutive_same_reason() -> None: """Подряд идущие одинаковые пропуски не плодят строки — иначе тик 60с = строка/мин.""" db = _FakeSkipDB(collapse=True) run_id = kit_runs.mark_skipped(db, source="avito_full_load", reason=SKIP_ALREADY_RUNNING) assert run_id == 7 assert db.sql_of("INSERT INTO scrape_runs") is None, "схлопывание не сработало" update = db.sql_of("UPDATE scrape_runs r") assert update is not None sql, _params = update assert "skips" in sql # счётчик повторов растёт вместо новой строки def test_mark_skipped_collapse_refreshes_detail_and_started_at() -> None: """Схлопывание освежает строку: detail, started_at (иначе живой стрик тонет в списке). Список прогонов сортирует ORDER BY started_at DESC с limit=20 — замороженный started_at утопил бы 37-дневный стрик под свежими прогонами других источников. Начало стрика сохраняется в counters.first_skip_at. """ db = _FakeSkipDB(collapse=True) kit_runs.mark_skipped( db, source="cian_history_backfill", reason="cian_cookies_expired", details="37 дн. назад" ) update = db.sql_of("UPDATE scrape_runs r") assert update is not None sql, params = update # clock_timestamp(), а не NOW(): отметки времени прогона пишутся настоящими # часами, иначе внутри долгой открытой транзакции они замерзают (#2702). assert "started_at = clock_timestamp()" in sql assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска" assert "first_skip_at" in sql, "начало стрика потеряно" assert params["details"] == "37 дн. назад" def test_mark_skipped_latest_lookup_uses_indexed_order() -> None: """Поиск последней строки идёт по (source, started_at DESC) — индекс из миграции 015. ORDER BY id DESC этот индекс не использует: для unknown_source (тик каждые 60 с бессрочно) это был бы отбор всех строк источника с сортировкой раз в минуту. """ db = _FakeSkipDB(collapse=True) kit_runs.mark_skipped(db, source="avito_full_load", reason=SKIP_ALREADY_RUNNING) update = db.sql_of("UPDATE scrape_runs r") assert update is not None sql, _params = update assert "ORDER BY started_at DESC, id DESC" in sql # ── 2. пять мест: kit _claim_run ×3 ────────────────────────────────────────── class _FakeResult: def __init__(self, *, scalar: Any = None, fetchone: Any = None) -> None: self._scalar = scalar self._fetchone = fetchone def scalar(self) -> Any: return self._scalar def fetchone(self) -> Any: return self._fetchone class _FakeClaimDB: """Мок Session для _claim_run (тот же сценарный контракт, что в parity-тестах).""" def __init__(self, *, running_states: list[bool], lock: bool = True) -> None: self._running = list(running_states) self._lock = lock self.rolled_back = False def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: sql = str(stmt) if "SELECT 1 FROM scrape_runs" in sql: return _FakeResult(fetchone=(1,) if self._running.pop(0) else None) if "pg_try_advisory_xact_lock" in sql: return _FakeResult(scalar=self._lock) return _FakeResult() def commit(self) -> None: pass def rollback(self) -> None: self.rolled_back = True def _ctx() -> SchedulerContext: runs = MagicMock() runs.create_run = MagicMock(return_value=42) return SchedulerContext( config=MagicMock(), matcher=MagicMock(), enrichment=MagicMock(), session_factory=MagicMock(), runs=runs, ) @pytest.mark.parametrize( ("running_states", "lock", "expected_reason"), [ ([True], True, SKIP_ALREADY_RUNNING), ([False], False, SKIP_CONCURRENT_CLAIM), ([False, True], True, SKIP_RUNNING_UNDER_LOCK), ], ) def test_claim_run_skip_writes_skipped_row( running_states: list[bool], lock: bool, expected_reason: str ) -> None: """Каждая из трёх skip-веток _claim_run пишет строку прогона с своей причиной.""" db = _FakeClaimDB(running_states=running_states, lock=lock) with patch.object(kit_runs, "mark_skipped") as mark: run_id = _claim_run(db, _make_sched("avito_city_sweep"), _ctx()) assert run_id is None mark.assert_called_once() assert mark.call_args.kwargs["reason"] == expected_reason assert mark.call_args.kwargs["source"] == "avito_city_sweep" def test_claim_run_under_lock_rolls_back_before_writing() -> None: """Строка пишется ПОСЛЕ rollback'а — иначе её INSERT улетел бы в откат лока.""" db = _FakeClaimDB(running_states=[False, True], lock=True) order: list[str] = [] real_rollback = db.rollback def _rollback() -> None: order.append("rollback") real_rollback() db.rollback = _rollback # type: ignore[method-assign] with patch.object(kit_runs, "mark_skipped", side_effect=lambda *a, **k: order.append("skip")): _claim_run(db, _make_sched("avito_city_sweep"), _ctx()) assert order == ["rollback", "skip"] def test_claim_run_happy_path_writes_no_skip_row() -> None: """Успешный claim не должен оставлять skip-строк (иначе счётчики мусорные).""" db = _FakeClaimDB(running_states=[False, False], lock=True) with patch.object(kit_runs, "mark_skipped") as mark: run_id = _claim_run(db, _make_sched("avito_city_sweep"), _ctx()) assert run_id == 42 mark.assert_not_called() # ── 2b. пятое место: scheduler_loop с неизвестным source ───────────────────── async def test_scheduler_loop_unknown_source_writes_skipped_row() -> None: """enabled-расписание без handler'а: раньше — warning каждый тик и ноль следов.""" ctx = SchedulerContext( config=MagicMock(), matcher=MagicMock(), enrichment=MagicMock(), session_factory=MagicMock(return_value=MagicMock()), runs=MagicMock(), # первый вызов (верх тика) — работаем, дальше — drain, чтобы выйти из while True shutdown_requested=MagicMock(side_effect=[False, True, True]), ) with ( patch.object(kit_sched.asyncio, "sleep", AsyncMock()), patch.object(kit_sched, "reap_zombies", MagicMock(return_value=0)), patch.object( kit_sched, "get_due_schedules", MagicMock(return_value=[_make_sched("source_from_mars")]), ), patch.object(kit_runs, "mark_skipped") as mark, ): await kit_sched.scheduler_loop(ctx, registry={}) mark.assert_called_once() assert mark.call_args.kwargs["reason"] == SKIP_UNKNOWN_SOURCE assert mark.call_args.kwargs["source"] == "source_from_mars" # ── 3. cian pre_claim: строка + достижимый алерт ───────────────────────────── def _patch_cian( *, cookies: dict[str, str] | None, expires_at: datetime | None, verify: Any = None, ) -> Any: from app.services import cian_session return ( patch.object(cian_session, "load_session", MagicMock(return_value=cookies)), patch.object(cian_session, "session_expires_at", MagicMock(return_value=expires_at)), patch.object(cian_session, "verify_session", AsyncMock(return_value=verify)), ) async def test_cian_pre_claim_expired_cookies_logs_error() -> None: """Фальсификация #2658-2: на протухших куках ДОЛЖНА быть ERROR-запись. Именно ERROR — в scraper-контейнере GlitchTip поднят с LoggingIntegration(event_level=ERROR) (scheduler_main.py), warning событием не станет. Старый код в этой ветке писал только logger.warning, а capture_message стоял во второй ветке (verify_session → None), недостижимой при протухании: load_session сам фильтрует expires_at_estimate > NOW(). """ from app.services import product_handlers expired = datetime.now(tz=UTC) - timedelta(days=37) p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=expired) records: list[logging.LogRecord] = [] class _Collector(logging.Handler): def emit(self, record: logging.LogRecord) -> None: records.append(record) handler = _Collector() product_handlers.logger.addHandler(handler) try: with ( p_load, p_exp, p_verify, patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()), patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()), ): proceed = await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) finally: product_handlers.logger.removeHandler(handler) assert proceed is False errors = [r for r in records if r.levelno >= logging.ERROR] assert errors, "протухание кук не породило ERROR-запись → события GlitchTip не будет" async def test_cian_pre_claim_expired_cookies_writes_skipped_row() -> None: """Немой `return False` заменён строкой прогона с причиной cian_cookies_expired.""" from app.services import product_handlers expired = datetime.now(tz=UTC) - timedelta(days=37) p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=expired) with ( p_load, p_exp, p_verify, patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()) as defer, patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark, ): proceed = await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) assert proceed is False mark.assert_called_once() assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_EXPIRED assert mark.call_args.kwargs["source"] == "cian_history_backfill" defer.assert_called_once() # next_run_at по-прежнему двигаем (не долбим Циан каждые 60с) async def test_cian_pre_claim_missing_cookies_reason_differs_from_expired() -> None: """«Кук нет вовсе» и «протухли» — разные слаги: причина машиночитаема.""" from app.services import product_handlers p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=None) with ( p_load, p_exp, p_verify, patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()), patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark, ): await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_MISSING assert product_handlers.SKIP_CIAN_COOKIES_MISSING != product_handlers.SKIP_CIAN_COOKIES_EXPIRED async def test_cian_pre_claim_rejected_cookies_writes_skipped_row() -> None: """Вторая ветка (Циан не принимает куки) — тоже строка, а не только алерт.""" from app.services import product_handlers future = datetime.now(tz=UTC) + timedelta(days=20) p_load, p_exp, p_verify = _patch_cian( cookies={"DMIR_AUTH": "x"}, expires_at=future, verify=None ) with ( p_load, p_exp, p_verify, patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()), patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark, ): proceed = await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) assert proceed is False assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_INVALID async def test_cian_pre_claim_warns_before_expiry_not_after() -> None: """Предупреждаем ЗАРАНЕЕ: куки ещё рабочие, но жить им меньше COOKIE_EXPIRY_WARN_DAYS. Обновление кук — ручная операция; алерт по факту протухания приходит, когда сбор уже встал. Гейт при этом пропускает прогон (proceed=True) — предупреждение, не блокировка. """ from app.services import cian_session, product_handlers soon = datetime.now(tz=UTC) + timedelta(days=cian_session.COOKIE_EXPIRY_WARN_DAYS - 1) p_load, p_exp, p_verify = _patch_cian( cookies={"DMIR_AUTH": "x"}, expires_at=soon, verify={"user": {"isAuthenticated": True}} ) records: list[logging.LogRecord] = [] class _Collector(logging.Handler): def emit(self, record: logging.LogRecord) -> None: records.append(record) handler = _Collector() product_handlers.logger.addHandler(handler) try: with p_load, p_exp, p_verify: proceed = await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) # Срок считаем по ЗАПИСИ, которую взял load_session: при нескольких аккаунтах # свежайшая-любая может быть чужой протухшей строкой (валидности она не знает). assert cian_session.session_expires_at.call_args.kwargs == {"valid_only": True} finally: product_handlers.logger.removeHandler(handler) assert proceed is True assert [r for r in records if r.levelno >= logging.ERROR], "не предупредили заранее" async def test_cian_pre_claim_fresh_cookies_are_silent() -> None: """Свежие куки — ни алерта, ни skip-строки (иначе алерт-усталость).""" from app.services import product_handlers far = datetime.now(tz=UTC) + timedelta(days=25) p_load, p_exp, p_verify = _patch_cian( cookies={"DMIR_AUTH": "x"}, expires_at=far, verify={"user": {"isAuthenticated": True}} ) records: list[logging.LogRecord] = [] class _Collector(logging.Handler): def emit(self, record: logging.LogRecord) -> None: records.append(record) handler = _Collector() product_handlers.logger.addHandler(handler) try: with ( p_load, p_exp, p_verify, patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark, ): proceed = await product_handlers._cian_pre_claim( MagicMock(), _make_sched("cian_history_backfill"), MagicMock() ) finally: product_handlers.logger.removeHandler(handler) assert proceed is True mark.assert_not_called() assert not [r for r in records if r.levelno >= logging.ERROR] # ── 4. монитор нулевых прогонов не смешивается с пропусками ────────────────── class _RecordingDB: """Возвращает заданные строки на SELECT и запоминает SQL (для проверки фильтров).""" def __init__(self, rows: list[Any]) -> None: self.rows = rows self.sql: list[str] = [] def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: self.sql.append(str(stmt)) result = MagicMock() result.fetchall.return_value = self.rows return result def test_zero_result_monitor_ignores_skipped_rows() -> None: """'skipped' НЕ участвует в стрике нулевых прогонов — это не «прогон вернул ноль». Смешать их в одном счётчике нельзя в обе стороны: пропуск не должен ни считаться нулевым прогоном, ни прерывать стрик реальных нулевых. Отсечка делается в SQL — проверяем, что 'skipped' не попал в список статусов выборки. """ db = _RecordingDB(rows=[]) kit_runs._alert_if_consecutive_zero_results(db, "cian_city_sweep") kit_runs._alert_if_consecutive_failures(db, "cian_city_sweep") assert db.sql, "alert-выборка не выполнилась" for sql in db.sql: assert "status IN ('failed', 'banned', 'done', 'cancelled')" in sql assert "skipped" not in sql def test_admin_runs_filter_accepts_skipped_status() -> None: """GET /admin/scrape/runs?status=skipped не должен отдавать 422. Единственная поверхность, где оператор спрашивает «что сейчас пропускается» — фильтр статуса в таблице прогонов. Без 'skipped' в Literal строки видны, а вопрос задать нельзя, то есть цель #2658 на UI достигнута наполовину. """ import typing from app.api.v1.admin import list_scrape_runs_unified annotation = typing.get_type_hints(list_scrape_runs_unified, include_extras=True)["status"] allowed: set[str] = set() for arg in typing.get_args(typing.get_args(annotation)[0]): allowed.update(typing.get_args(arg)) assert "skipped" in allowed def test_mark_skipped_status_is_not_a_failure_status() -> None: """Строка-пропуск не попадает и в failed/banned-стрик (алерт «3 подряд ошибки»).""" db = _FakeSkipDB(collapse=False) kit_runs.mark_skipped(db, source="cian_history_backfill", reason="cian_cookies_expired") insert = db.sql_of("INSERT INTO scrape_runs") assert insert is not None sql, _params = insert assert "'skipped'" in sql assert "'failed'" not in sql and "'banned'" not in sql