All checks were successful
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 7s
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m35s
Из таблицы убитых деплоем: yandex_city_sweep 15.08 прожил 65 мин, 12.08 — 2ч33м; оба потеряны целиком — у свипа не было чекпоинтов вовсе. Единица чекпоинта — combo (сегмент × комнатность × ценовой диапазон), ровно то, чем цикл обхода уже итерируется. Три слоя: - провайдер: skip_combos (ни одного HTTP по собранным) + on_combo для каждого ПРОЙДЕННОГО combo, включая пустые — иначе пустой combo не попадал бы в чекпоинт и перечитывался бы вечно; оборванный отказом combo (gate failure) on_combo по-прежнему не вызывает; - пайплайн: done_combos → heartbeat с done_buckets (мерж jsonb, финализаторы не затирают); подхват гейтится единственным якорем — combo_label не содержит якоря, multi-anchor подхват пропускал бы чужие якоря; - планировщик: resume_run_id=_pick_resume(...) в диспатче (generic-механизм #2845 — params-идентичность, свежесть точки, потолок цепочки — бесплатно). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
241 lines
9.6 KiB
Python
241 lines
9.6 KiB
Python
"""#3074: чекпоинты для yandex_city_sweep — combo как единица возобновления.
|
||
|
||
Из таблицы убитых деплоем прогонов (#3074): yandex_city_sweep 15.08 прожил 65 мин,
|
||
12.08 — 2 ч 33 мин; оба потеряны целиком, потому что у свипа не было чекпоинтов
|
||
вовсе (avito/cian full-load обзавелись ими в #930/#2845).
|
||
|
||
Единица чекпоинта — combo (сегмент × комнатность × ценовой диапазон): ровно то,
|
||
чем цикл обхода уже итерируется, метка та же, что у бакетов avito
|
||
("vtorichka/room_1:0-3000000"). Три слоя:
|
||
провайдер — skip_combos (ни одного HTTP по собранным) + on_combo для каждого
|
||
ПРОЙДЕННОГО combo, включая пустые (иначе пустой combo не попадал
|
||
бы в чекпоинт и перечитывался бы вечно);
|
||
пайплайн — done_combos → heartbeat c done_buckets (мерж jsonb);
|
||
планировщик — resume_run_id=_pick_resume(...) в диспатче.
|
||
|
||
Красные на origin/main: планировщик передаёт None литералом (по значению),
|
||
on_combo не вызывается для пустых combo (по значению).
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import os
|
||
|
||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
|
||
|
||
import json
|
||
import types
|
||
from typing import Any
|
||
from unittest.mock import MagicMock, patch
|
||
|
||
import pytest
|
||
from scraper_kit.orchestration import scheduler as sched
|
||
from scraper_kit.providers.yandex.serp import YandexRealtyScraper
|
||
|
||
# ── провайдер: заглушка gate-API ─────────────────────────────────────────────
|
||
|
||
|
||
def _gate_payload(entities: list[dict[str, Any]], total_pages: int = 1) -> str:
|
||
offers = {"entities": entities, "pager": {"totalPages": total_pages}}
|
||
return json.dumps({"response": {"search": {"offers": offers}}})
|
||
|
||
|
||
class _RecordingBrowser:
|
||
"""BrowserFetcher-заглушка: отвечает одним и тем же телом, считает вызовы."""
|
||
|
||
def __init__(self, body: str) -> None:
|
||
self.body = body
|
||
self.urls: list[str] = []
|
||
|
||
async def fetch(self, url: str) -> str:
|
||
self.urls.append(url)
|
||
return self.body
|
||
|
||
|
||
def _scraper(body: str) -> YandexRealtyScraper:
|
||
s = YandexRealtyScraper(types.SimpleNamespace())
|
||
s._browser = _RecordingBrowser(body) # type: ignore[assignment]
|
||
s.sleep_between_requests = _no_sleep # type: ignore[method-assign]
|
||
return s
|
||
|
||
|
||
async def _no_sleep() -> None:
|
||
pass
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_on_combo_fires_for_complete_empty_combo() -> None:
|
||
"""Пройденный до конца combo с ПУСТОЙ выдачей обязан дойти до on_combo —
|
||
иначе он не попадёт в чекпоинт и будет перечитываться каждым продолжением.
|
||
На origin/main on_combo вызывался только при непустых new_lots."""
|
||
scraper = _scraper(_gate_payload([]))
|
||
seen_labels: list[str] = []
|
||
|
||
await scraper.fetch_around_multi_room(
|
||
56.84,
|
||
60.60,
|
||
1000,
|
||
max_pages=1,
|
||
rooms_list=["room_1"],
|
||
price_ranges=[(None, 3_000_000)],
|
||
segments=["NO"],
|
||
on_combo=lambda label, lots: seen_labels.append(label),
|
||
)
|
||
|
||
assert seen_labels == ["vtorichka/room_1:None-3000000"], (
|
||
"пустой, но полностью пройденный combo не дошёл до on_combo — "
|
||
"в чекпоинт он не попадёт никогда"
|
||
)
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_skip_combos_makes_zero_http_requests() -> None:
|
||
"""Combo из чекпоинта не порождает ни одного HTTP-запроса и не зовёт on_combo."""
|
||
scraper = _scraper(_gate_payload([]))
|
||
seen_labels: list[str] = []
|
||
|
||
await scraper.fetch_around_multi_room(
|
||
56.84,
|
||
60.60,
|
||
1000,
|
||
max_pages=1,
|
||
rooms_list=["room_1", "room_2"],
|
||
price_ranges=[(None, 3_000_000)],
|
||
segments=["NO"],
|
||
on_combo=lambda label, lots: seen_labels.append(label),
|
||
skip_combos={"vtorichka/room_1:None-3000000"},
|
||
)
|
||
|
||
browser: _RecordingBrowser = scraper._browser # type: ignore[assignment]
|
||
assert seen_labels == ["vtorichka/room_2:None-3000000"]
|
||
assert len(browser.urls) == 1, f"скипнутый combo всё равно ходил в сеть: {browser.urls}"
|
||
|
||
|
||
# ── планировщик: диспатч отдаёт точку в пайплайн (зеркало test_930) ─────────
|
||
|
||
|
||
def _candidate(**over: Any) -> types.SimpleNamespace:
|
||
base = {
|
||
"prev_id": 4117,
|
||
"prev_status": "banned",
|
||
"prev_counters": {"done_buckets": ["vtorichka/room_1:None-3000000"]},
|
||
"same_params": True,
|
||
"age_h": 20.0,
|
||
"interval_days": "1",
|
||
}
|
||
base.update(over)
|
||
return types.SimpleNamespace(**base)
|
||
|
||
|
||
class _FakeDb:
|
||
def __init__(self, row: Any) -> None:
|
||
self.row = row
|
||
|
||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||
if params and "counters" in params:
|
||
return MagicMock()
|
||
return MagicMock(fetchone=lambda: self.row)
|
||
|
||
def commit(self) -> None:
|
||
pass
|
||
|
||
|
||
async def test_scheduler_hands_checkpoint_to_yandex_sweep() -> None:
|
||
"""_job_yandex_city_sweep передаёт resume_run_id, а не None литералом.
|
||
Красный на origin/main по значению (ключа в kwargs нет → None != 4117)."""
|
||
db = _FakeDb(_candidate())
|
||
captured: dict[str, Any] = {}
|
||
|
||
async def _spy(*_a: Any, **kw: Any) -> None:
|
||
captured.update(kw)
|
||
|
||
with patch.object(sched, "run_yandex_city_sweep", _spy):
|
||
await sched._job_yandex_city_sweep(db, 5000, {}, MagicMock())
|
||
|
||
assert captured.get("resume_run_id") == 4117, (
|
||
"планировщик не отдал чекпоинт яндекс-свипу — прогон пойдёт с нуля"
|
||
)
|
||
|
||
|
||
# ── пайплайн: проводка skip→scraper и done_buckets→heartbeat ────────────────
|
||
|
||
|
||
class _SweepFakeDb:
|
||
"""Резюм-SELECT отдаёт counters предшественника; UPDATE'ы записываются."""
|
||
|
||
def __init__(self, prev_counters: dict[str, Any]) -> None:
|
||
self.prev_counters = prev_counters
|
||
self.heartbeats: list[dict[str, Any]] = []
|
||
|
||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||
if params and "counters" in params:
|
||
self.heartbeats.append(json.loads(params["counters"]))
|
||
return MagicMock()
|
||
if params and "rid" in params:
|
||
return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters))
|
||
return MagicMock() # is_cancelled/сторожа — best-effort, глотаем
|
||
|
||
def commit(self) -> None:
|
||
pass
|
||
|
||
def rollback(self) -> None:
|
||
pass
|
||
|
||
|
||
class _FakeSweepScraper:
|
||
"""Двойник YandexRealtyScraper для пайплайна: фиксирует kwargs fetch'а,
|
||
отдаёт один пройденный пустой combo через on_combo."""
|
||
|
||
captured: dict[str, Any] = {} # noqa: RUF012 — тестовый сборник kwargs
|
||
gate_fetch_attempts = 1
|
||
gate_fetch_failures = 0
|
||
|
||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||
pass
|
||
|
||
async def __aenter__(self) -> _FakeSweepScraper:
|
||
return self
|
||
|
||
async def __aexit__(self, *_exc: Any) -> None:
|
||
return None
|
||
|
||
async def fetch_around_multi_room(self, *_a: Any, **kw: Any) -> list:
|
||
_FakeSweepScraper.captured = dict(kw)
|
||
on_combo = kw.get("on_combo")
|
||
if on_combo is not None:
|
||
on_combo("vtorichka/room_2:None-3000000", [])
|
||
return []
|
||
|
||
|
||
async def test_pipeline_resumes_and_checkpoints() -> None:
|
||
"""run_yandex_city_sweep: чекпоинт предшественника уезжает в scraper как
|
||
skip_combos, пройденный combo дописывается в done_buckets heartbeat'а."""
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
prev = {"done_buckets": ["vtorichka/room_1:None-3000000"]}
|
||
db = _SweepFakeDb(prev)
|
||
|
||
with (
|
||
patch.object(pl, "YandexRealtyScraper", _FakeSweepScraper),
|
||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||
):
|
||
await pl.run_yandex_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
run_id=5001,
|
||
config=types.SimpleNamespace(scraper_proxy_url=None),
|
||
matcher=MagicMock(),
|
||
enrichment=MagicMock(),
|
||
enrich_address=False,
|
||
resume_run_id=4999,
|
||
)
|
||
|
||
assert _FakeSweepScraper.captured.get("skip_combos") == {"vtorichka/room_1:None-3000000"}, (
|
||
"чекпоинт предшественника не доехал до scraper'а"
|
||
)
|
||
|
||
with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb]
|
||
assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится"
|
||
assert with_ckpt[-1]["done_buckets"] == [
|
||
"vtorichka/room_1:None-3000000",
|
||
"vtorichka/room_2:None-3000000",
|
||
], "done_buckets не аккумулирует пройденный combo поверх унаследованных"
|