gendesign/tradein-mvp/backend/tests/test_scraper_kit_scheduler_parity.py
bot-backend 85059aeb1b refactor(tradein/scheduler): удалить legacy scheduler_loop + scraper-scheduling, kit единственный путь (#2397 Part C)
Топология подтверждена перед удалением (docker-compose.prod.yml): tradein-backend
(uvicorn app.main:app) — SCHEDULER_ENABLE=false; tradein-scraper (python -m
app.scheduler_main) — SCHEDULER_ENABLE=true + USE_KIT_SCHEDULER=true. Kit-путь
(_run_kit_scheduler → scraper_kit.orchestration.scheduler + product_handlers)
самодостаточен: не импортирует ничего из app.services.scheduler.scheduler_loop
или app.services.scrape_pipeline. Все НЕ-sweep джобы, которые kit-scheduler
диспетчерит через build_product_handlers, идут напрямую в app.tasks.*/
app.services.* (либо lazy-импортят import_rosreestr_dkp/_execute_cian_backfill
из scheduler.py) — мимо удаляемой legacy-машинерии.

app/services/scheduler.py: 2098 → 418 строк. Удалено: scheduler_loop,
get_due_schedules, reap_zombies, _claim_run, _defer_next_run_at, _spawn_tracked/
_drain_inflight/_inflight_tasks, все 27 trigger_*_run-функций, импорт
app.services.scrape_pipeline, константы SCHEDULER_TICK_SEC/ZOMBIE_THRESHOLD_HOURS
(достижимы были только через удалённый scheduler_loop-путь). Оставлено (живые
импортёры вне удалённого): compute_next_run_at + has_running_run (admin.py),
import_rosreestr_dkp + _execute_cian_backfill (lazy-импорты в
product_handlers.py — job-тела kit-handler'ов).

main.py: убран `from app.services.scheduler import scheduler_loop` + lifespan-блок
запуска (`if settings.scheduler_enable: asyncio.create_task(scheduler_loop())`);
прод-backend всегда шёл с SCHEDULER_ENABLE=false, так что это был мёртвый код.

scheduler_main.py: убрана ship-dark развилка #2192 (USE_KIT_SCHEDULER=false →
legacy scheduler_loop fallback) — _run_kit_scheduler() теперь безусловный путь.
Поле settings.use_kit_scheduler оставлено в конфиге (Settings extra="ignore"
защищает от startup-краха на leftover env var), но на ветвление не влияет.

app.services.scrape_pipeline: 0 runtime-импортёров в app/+scripts/+packages/
после этого PR (только тесты, которые Part E удалит вместе с самим файлом) —
подтверждено grep. scrape_pipeline.py не тронут (Part E).

Тесты: удалены test_house_imv_backfill_scheduler.py (100% legacy-триггер,
backfill_house_imv сервис покрыт в test_house_imv_backfill_browser_flag.py /
test_backfill_wave2.py) и test_kit_registry_completeness.py (parity-инвариант
против удалённого dispatch, дублирует test_scraper_kit_scheduler_parity.py).
Точечно вырезаны "Scheduler wiring" секции (trigger_fn_exists/dispatch_branch_
wired/runs_in_executor) из ~10 файлов, тестирующих сами task-функции — сами
task-тесты (SQL-shape, миграции, fake-db поведение) оставлены нетронутыми.
test_scheduler.py: 825 → ~90 строк (остались только compute_next_run_at-тесты).
test_scraper_kit_scheduler_parity.py: убрана golden-parity секция против
удалённого scheduler_loop (SOURCE_TO_OLD_TRIGGER/_drive_old_one_tick/
test_routing_parity_per_source), остальное (claim/reap_zombies/dispatch/
registry-shape тесты kit-модуля) сохранено — источник этих инвариантов не
app.services.scheduler, а сам scraper_kit.orchestration.scheduler.
test_scheduler_main.py: 2 теста, патчившие app.services.scheduler.scheduler_loop,
переведены на монкипатч sm._run_kit_scheduler (единственный путь после этого PR).
test_sweep_imv_phase.py:171-371 (6 прямых импортов run_avito_city_sweep из
scrape_pipeline) намеренно НЕ тронуты — Part E.

Verify: полный pytest 3179 passed / 6 skipped / 1 known-unrelated fail
(test_search_cache_hit, #2208, не связан с этим PR); ruff 0.7.4 чист на всех
изменённых файлах; `python -c "import app.main; import app.scheduler_main"` OK.
2026-07-04 13:00:30 +03:00

411 lines
15 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Smoke/shape tests for `scraper_kit.orchestration.scheduler` (#2136, post-strangler).
Kit — единственный scheduler-путь (#2397 Part C убрал legacy
`app.services.scheduler.scheduler_loop` + trigger_*/get_due_schedules/_claim_run/
reap_zombies machinery; сравнивать «golden-parity» больше не с чем). Этот файл проверяет
форму и инварианты kit-реестра самого по себе:
1. REGISTRY: kit-native sweep-handler'ы зарегистрированы (`build_registry`), продуктовые
source'ы резолвятся через `resolve_handler` (включая deactivate_stale_* wildcard),
неизвестный source → None.
2. CLAIM: `_claim_run` — happy-path, already-running skip, lock-busy skip, appeared-under-
lock rollback. Мокаем БД + ctx.runs.
3. ZOMBIE: `reap_zombies` возвращает число + commit.
4. SUB-HOURLY: `reschedule_after_minutes` post_claim UPDATE + commit (proxy_healthcheck).
5. DISPATCH e2e: `_dispatch` для kit-native sweep (мок pipeline) и продуктового handler —
claim → свежая сессия → job вызван → close.
Без сети, без БД.
"""
from __future__ import annotations
import asyncio
import os
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.orchestration import scheduler as kit_sched
from scraper_kit.orchestration.scheduler import (
Handler,
SchedulerContext,
_claim_run,
_dispatch,
build_registry,
reap_zombies,
reschedule_after_minutes,
resolve_handler,
)
# ── продуктовые source'ы (НЕ kit-native), которые build_product_handlers регистрирует ──
# (было ключами таблицы source → старый trigger_*; сам старый dispatch удалён #2397 Part C,
# но список source-имён остаётся полезным для recording-registry ниже).
_PRODUCT_SOURCES: set[str] = {
"rosreestr_dkp_import",
"listing_source_snapshot",
"asking_to_sold_ratio_refresh",
"refresh_search_matview",
"yandex_address_backfill",
"deactivate_stale_avito",
"deactivate_stale_yandex",
"deactivate_stale_cian",
"sber_index_pull",
"rosreestr_quarter_poll",
"newbuilding_enrich",
"yandex_newbuilding_sweep",
"geoportal_coords_backfill",
"geocode_missing_listings",
"avito_detail_backfill",
"yandex_detail_backfill",
"cadastral_geo_match",
"house_imv_backfill",
"house_dedup_merge",
"proxy_healthcheck",
}
# kit-native (тело в scraper_kit.orchestration.pipeline) — регистрируются встроенно
_KIT_NATIVE_SOURCES = {
"avito_city_sweep",
"avito_full_load",
"avito_full_load_exhaustive",
"avito_newbuilding_sweep",
"yandex_city_sweep",
"cian_city_sweep",
"cian_full_load",
"domclick_city_sweep",
}
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,
}
# ── kit-реестр с продуктовыми handler'ами (recording stubs) ──────────────────
def _build_recording_registry() -> tuple[dict[str, Handler], dict[str, MagicMock]]:
"""Kit-native sweeps + продуктовые recording-handler'ы. Возвращает (registry, fired-map)."""
fired: dict[str, MagicMock] = {}
def _make_job(name: str) -> Any:
rec = MagicMock(name=name)
fired[name] = rec
async def _job(db: Any, run_id: int, params: dict[str, Any], ctx: Any) -> None:
rec(db, run_id, params)
return _job
product: dict[str, Handler] = {}
# продуктовые exact-match источники
for src in _PRODUCT_SOURCES:
if src in _KIT_NATIVE_SOURCES or src.startswith("deactivate_stale_"):
continue
product[src] = Handler(_make_job(src), src)
# deactivate_stale_* — одна wildcard-запись (как боевой startswith)
product["deactivate_stale_*"] = Handler(_make_job("deactivate_stale"), "deactivate_stale")
return build_registry(product), fired
# ── 1. REGISTRY shape ────────────────────────────────────────────────────────
def test_routing_coverage_sets_match() -> None:
"""Множество source'ов kit-реестра ⊇ множество продуктовых source'ов."""
registry, _ = _build_recording_registry()
for source in _PRODUCT_SOURCES:
assert resolve_handler(source, registry) is not None, f"kit misses source={source}"
def test_kit_native_handler_set() -> None:
"""Ровно 8 kit-native sweep-обработчиков зарегистрированы встроенно."""
registry = build_registry()
assert set(registry) == _KIT_NATIVE_SOURCES
for src in _KIT_NATIVE_SOURCES:
assert resolve_handler(src, registry).log_name == src
def test_unknown_source_resolves_none() -> None:
registry, _ = _build_recording_registry()
assert resolve_handler("totally_unknown_source", registry) is None
def test_kit_scheduler_has_no_app_imports() -> None:
"""kit scheduler развязан от app.* — ни одного `import app` / `from app` (AST, не substring)."""
import ast
import inspect
source = inspect.getsource(kit_sched)
tree = ast.parse(source)
for node in ast.walk(tree):
if isinstance(node, ast.Import):
for alias in node.names:
assert not alias.name.startswith("app"), f"illegal import app: {alias.name}"
elif isinstance(node, ast.ImportFrom):
assert not (node.module or "").startswith("app"), f"illegal from app: {node.module}"
def test_deactivate_stale_wildcard_prefix() -> None:
"""deactivate_stale_{avito,yandex,cian} → одна wildcard-запись (как боевой startswith)."""
registry, _ = _build_recording_registry()
h_avito = resolve_handler("deactivate_stale_avito", registry)
h_yandex = resolve_handler("deactivate_stale_yandex", registry)
h_cian = resolve_handler("deactivate_stale_cian", registry)
assert h_avito is not None
assert h_avito is h_yandex is h_cian # один и тот же handler на всё семейство
# ── 2. _claim_run advisory-lock parity ───────────────────────────────────────
class _FakeResult:
def __init__(self, *, scalar: Any = None, fetchone: Any = None, fetchall: Any = None) -> None:
self._scalar = scalar
self._fetchone = fetchone
self._fetchall = fetchall or []
def scalar(self) -> Any:
return self._scalar
def fetchone(self) -> Any:
return self._fetchone
def fetchall(self) -> Any:
return self._fetchall
class _FakeClaimDB:
"""Мок Session для _claim_run: сценарий has_running_run + advisory-lock."""
def __init__(self, *, running_states: list[bool], lock: bool = True) -> None:
# running_states: очередь ответов has_running_run (по вызовам SELECT 1 FROM scrape_runs)
self._running = list(running_states)
self._lock = lock
self.committed = False
self.rolled_back = False
self.update_calls = 0
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
if "SELECT 1 FROM scrape_runs" in sql:
is_running = self._running.pop(0)
return _FakeResult(fetchone=(1,) if is_running else None)
if "pg_try_advisory_xact_lock" in sql:
return _FakeResult(scalar=self._lock)
if "UPDATE scrape_schedules" in sql:
self.update_calls += 1
return _FakeResult()
return _FakeResult()
def commit(self) -> None:
self.committed = True
def rollback(self) -> None:
self.rolled_back = True
def _ctx_with_runs(create_run_ret: int = 42) -> SchedulerContext:
runs = MagicMock()
runs.create_run = MagicMock(return_value=create_run_ret)
return SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(),
runs=runs,
)
def test_claim_run_happy_path() -> None:
db = _FakeClaimDB(running_states=[False, False], lock=True)
ctx = _ctx_with_runs(create_run_ret=42)
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id == 42
ctx.runs.create_run.assert_called_once()
assert db.update_calls == 1
assert db.committed is True
def test_claim_run_skip_already_running() -> None:
db = _FakeClaimDB(running_states=[True], lock=True)
ctx = _ctx_with_runs()
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id is None
ctx.runs.create_run.assert_not_called()
def test_claim_run_skip_lock_busy() -> None:
# pre-check clean, но advisory-lock занят конкурентным тиком → skip
db = _FakeClaimDB(running_states=[False], lock=False)
ctx = _ctx_with_runs()
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id is None
ctx.runs.create_run.assert_not_called()
def test_claim_run_running_appeared_under_lock() -> None:
# pre-check clean, lock взят, но под локом появился running → rollback + skip (double-check)
db = _FakeClaimDB(running_states=[False, True], lock=True)
ctx = _ctx_with_runs()
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id is None
assert db.rolled_back is True
ctx.runs.create_run.assert_not_called()
# ── 3. reap_zombies ──────────────────────────────────────────────────────────
def test_reap_zombies_counts_and_commits() -> None:
db = MagicMock()
db.execute.return_value = _FakeResult(fetchall=[MagicMock(id=1), MagicMock(id=2)])
n = reap_zombies(db)
assert n == 2
db.commit.assert_called_once()
# порог 6 часов передан в SQL-параметр
_stmt, params = db.execute.call_args[0]
assert params["interval"] == "6 hours"
def test_reap_zombies_threshold_constant() -> None:
assert kit_sched.ZOMBIE_THRESHOLD_HOURS == 6
# ── 4. proxy_healthcheck sub-hourly post_claim ───────────────────────────────
def test_reschedule_after_minutes_updates_next_run() -> None:
db = MagicMock()
hook = reschedule_after_minutes(param="interval_minutes", default=30)
ctx = _ctx_with_runs()
hook(db, run_id=7, params={"interval_minutes": 15}, ctx=ctx)
db.commit.assert_called_once()
_stmt, params = db.execute.call_args[0]
assert params["mins"] == 15
assert params["rid"] == 7
def test_reschedule_after_minutes_default() -> None:
db = MagicMock()
hook = reschedule_after_minutes()
hook(db, run_id=9, params={}, ctx=_ctx_with_runs())
_stmt, params = db.execute.call_args[0]
assert params["mins"] == 30
# ── 5. _dispatch end-to-end (claim → fresh session → job → close) ────────────
async def test_dispatch_kit_native_sweep_invokes_pipeline() -> None:
"""kit-native handler: _dispatch клеймит, открывает сессию и зовёт pipeline.run_*."""
run_db = MagicMock()
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(return_value=run_db),
runs=MagicMock(),
)
ctx.runs.create_run = MagicMock(return_value=100)
db = _FakeClaimDB(running_states=[False, False], lock=True)
registry = build_registry()
handler = registry["avito_city_sweep"]
with patch.object(kit_sched, "run_avito_city_sweep", AsyncMock()) as mock_run:
run_id = await _dispatch(handler, db, _make_sched("avito_city_sweep"), ctx)
assert run_id == 100
# дождаться detached run-задачи
await asyncio.gather(*list(ctx._inflight_tasks))
mock_run.assert_awaited_once()
_args, kwargs = mock_run.call_args
assert kwargs["run_id"] == 100
assert kwargs["config"] is ctx.config
assert kwargs["matcher"] is ctx.matcher
run_db.close.assert_called_once()
async def test_dispatch_product_handler_invokes_job() -> None:
"""Продуктовый handler: _dispatch клеймит и зовёт инжектированный job на свежей сессии."""
run_db = MagicMock()
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(return_value=run_db),
runs=MagicMock(),
)
ctx.runs.create_run = MagicMock(return_value=200)
db = _FakeClaimDB(running_states=[False, False], lock=True)
registry, fired = _build_recording_registry()
handler = resolve_handler("sber_index_pull", registry)
run_id = await _dispatch(handler, db, _make_sched("sber_index_pull"), ctx)
assert run_id == 200
await asyncio.gather(*list(ctx._inflight_tasks))
fired["sber_index_pull"].assert_called_once()
called_args = fired["sber_index_pull"].call_args[0]
assert called_args[0] is run_db # свежая сессия
assert called_args[1] == 200 # run_id
run_db.close.assert_called_once()
async def test_dispatch_pre_claim_gate_skips() -> None:
"""pre_claim → False: _dispatch пропускает claim (cian cookie-gate семантика)."""
ctx = _ctx_with_runs()
async def _deny(db: Any, sch: dict[str, Any], c: Any) -> bool:
return False
job = AsyncMock()
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
job()
handler = Handler(_job, "cian_history_backfill", pre_claim=_deny)
db = _FakeClaimDB(running_states=[], lock=True)
run_id = await _dispatch(handler, db, _make_sched("cian_history_backfill"), ctx)
assert run_id is None
ctx.runs.create_run.assert_not_called()
job.assert_not_called()
async def test_dispatch_post_claim_hook_runs() -> None:
"""post_claim вызывается сразу после claim (proxy_healthcheck sub-hourly reschedule)."""
run_db = MagicMock()
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(return_value=run_db),
runs=MagicMock(),
)
ctx.runs.create_run = MagicMock(return_value=300)
post = MagicMock()
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
pass
handler = Handler(_job, "proxy_healthcheck", post_claim=post)
db = _FakeClaimDB(running_states=[False, False], lock=True)
run_id = await _dispatch(handler, db, _make_sched("proxy_healthcheck"), ctx)
assert run_id == 300
await asyncio.gather(*list(ctx._inflight_tasks))
post.assert_called_once()
assert post.call_args[0][1] == 300 # run_id