"""#930 добивка: планировщик не подхватывал чекпоинт оборванного прогона.
#930 сделал обе половины механизма — запись точки (`counters.done_buckets`, per-bucket
heartbeat) и её чтение (`run_*_full_load(resume_run_id=...)`, skip-set в SERP-слое), —
но единственным входом оставил админку. У avito full-load админского эндпоинта нет
вовсе, а планировщик передавал `resume_run_id=None` ЛИТЕРАЛОМ (scheduler.py 708/728/799
на origin/main). То есть боевой путь возобновления не существовал ни одного дня.
Цена на проде (замер 2026-08-12, 90 суток, read-only): 433 корзины в 30 оборванных
прогонах с ЖИВОЙ незабранной точкой — avito_full_load 242, cian_full_load 134,
avito_full_load_exhaustive 57. Прогон 3547 (09.08, убит деплоем на третьем часу, 35 из
84 корзин дерева) лежит до сих пор и будет подхвачен расписанием 139 16.08.
Красный прогон на origin/main:
1. `test_scheduler_hands_checkpoint_to_pipeline` — планировщик отдаёт в пайплайн
resume_run_id=None вместо id прошлого прогона (AssertionError на 3 источниках);
2. `test_partial_bucket_is_not_complete` — бакет с выпавшей страницей приезжает в
on_bucket неотличимым от целого (у колбэка нет аргумента полноты вообще);
3. `test_pipeline_keeps_partial_bucket_out_of_checkpoint` — TypeError: `_on_bucket`
на main принимает два аргумента, признаку полноты некуда приехать.
Тесты ладдера (`_resume_decision`) на main падают с AttributeError — функции нет.
"""
from __future__ import annotations
import os
from types import SimpleNamespace
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 scheduler as sched
from scraper_kit.orchestration.pipeline import run_avito_full_load
from scraper_kit.providers.avito.serp import AvitoScraper
PFX = "scraper_kit.orchestration.pipeline"
# Прод-слепок расписания 139 (avito_full_load_exhaustive) на 2026-08-12: params прогона
# 3547 совпадают с default_params расписания байт-в-байт — это и есть «то же задание».
_PARAMS = {
"concurrency": 1,
"interval_days": 7,
"secondary_only": True,
"request_delay_sec": 7.0,
"price_cap_per_bucket": 1400,
}
def _candidate(**over: Any) -> SimpleNamespace:
"""Строка-кандидат из _RESUME_CANDIDATE_SQL: прогон 3547 как он лежит на проде."""
base = {
"prev_id": 3547,
"prev_status": "cancelled",
"prev_counters": {
"unique_fetched": 5496,
"done_buckets": [f"room_1_komn:{i}:0" for i in range(35)],
},
"same_params": True,
"age_h": 164.6, # 6.86 суток — столько будет точке 3547 к подхвату 16.08
"interval_days": "7",
}
base.update(over)
return SimpleNamespace(**base)
class _FakeDb:
"""Двойник сессии: отдаёт ОДНУ строку-кандидата на любой SELECT, глотает UPDATE."""
def __init__(self, row: Any) -> None:
self.row = row
self.written: list[dict[str, Any]] = []
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params: # update_heartbeat пишет вердикт
self.written.append(params)
return MagicMock()
return MagicMock(fetchone=lambda: self.row)
def commit(self) -> None:
pass
# ── 1. Главное: планировщик обязан отдать точку в пайплайн ───────────────────
@pytest.mark.parametrize(
("job", "pipeline_fn"),
[
(sched._job_avito_full_load, "run_avito_full_load"),
(sched._job_avito_full_load_exhaustive, "run_avito_full_load"),
(sched._job_cian_full_load, "run_cian_full_load"),
],
)
async def test_scheduler_hands_checkpoint_to_pipeline(job: Any, pipeline_fn: str) -> None:
"""Оборванный прогон с валидной точкой → новый прогон продолжает его, а не с нуля.
Падает на origin/main: планировщик передаёт литеральный None — 433 корзины за 90
суток перебирались заново, включая 35 корзин прогона 3547.
"""
db = _FakeDb(_candidate())
captured: dict[str, Any] = {}
async def _spy(*_a: Any, **kw: Any) -> None:
captured.update(kw)
with patch.object(sched, pipeline_fn, _spy):
await job(db, 4000, dict(_PARAMS), MagicMock())
assert captured["resume_run_id"] == 3547
async def test_verdict_lands_in_counters_of_new_run() -> None:
"""Подхватили или нет — видно В СЧЁТЧИКАХ прогона, а не только в docker-логах.
Логи теряются при редеплое (контейнер tradein-scraper пересоздаётся), поэтому
молчаливый отказ подхватить неотличим от отсутствия правки.
"""
db = _FakeDb(_candidate(prev_status="zombie"))
with patch.object(sched, "run_avito_full_load", AsyncMock()):
await sched._job_avito_full_load(db, 4000, dict(_PARAMS), MagicMock())
assert db.written, "вердикт о подхвате не записан в counters нового прогона"
written = db.written[-1]["counters"]
assert '"resume_reason": "status_zombie"' in written
assert '"resume_candidate": 3547' in written
# ── 2. Ладдер отказов: у каждого нуля своя причина ───────────────────────────
@pytest.mark.parametrize(
("row", "reason"),
[
(None, "no_prev_run"),
(_candidate(prev_status="done"), "status_done"),
(_candidate(prev_status="zombie"), "status_zombie"),
(_candidate(same_params=False), "params_changed"),
(_candidate(prev_counters={"unique_fetched": 2977}), "no_checkpoint"),
(_candidate(age_h=200.0), "checkpoint_stale"),
(_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2}), "chain_limit"),
],
)
def test_resume_refusals_are_named(row: Any, reason: str) -> None:
"""«Не подхватили» — это семь РАЗНЫХ фактов, и в counters они различимы."""
resume_id, verdict = sched._resume_decision(row)
assert resume_id is None
assert verdict["resume_reason"] == reason
assert verdict["resume_from"] is None
def test_resume_chain_is_bounded() -> None:
"""Цепочка возобновлений считается и упирается в потолок, а не тянется вечно.
Потолок выведен из STALE_DIGEST_INTERVAL_FACTOR (см. scheduler.py): полный обход
обязан начаться раньше, чем сводка объявит источник просроченным.
"""
assert sched._MAX_RESUME_CHAIN == sched.STALE_DIGEST_INTERVAL_FACTOR - 1
_id, first = sched._resume_decision(_candidate())
assert first["resume_chain"] == 1
_id2, second = sched._resume_decision(
_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 1})
)
assert second["resume_chain"] == sched._MAX_RESUME_CHAIN
third_id, third = sched._resume_decision(
_candidate(prev_counters={"done_buckets": ["a"], "resume_chain": 2})
)
assert third_id is None and third["resume_reason"] == "chain_limit"
def test_stale_threshold_follows_the_source_tick() -> None:
"""Срок годности точки считается от такта ИСТОЧНИКА, а не общей константой.
cian ходит раз в 3 суток, avito — раз в 7; одна и та же точка возрастом 100 ч для
первого просрочена, для второго свежая. Плюс сутки — сетка запуска (см.
_resume_decision): 164.6 ч прогона 3547 при такте 7 суток обязаны пройти, иначе
точку отвергал бы jitter расписания, а пропущенный цикл (13 суток) — нет.
"""
assert sched._resume_decision(_candidate(age_h=100.0, interval_days="3"))[0] is None
assert sched._resume_decision(_candidate(age_h=100.0, interval_days="7"))[0] == 3547
assert sched._resume_decision(_candidate(age_h=164.6, interval_days="7"))[0] == 3547
assert sched._resume_decision(_candidate(age_h=13 * 24.0, interval_days="7"))[0] is None
# ── 3. Недособранный бакет не имеет права попасть в чекпоинт ─────────────────
def _serp_config() -> SimpleNamespace:
return SimpleNamespace(
scraper_fetch_mode="curl_cffi",
browser_http_endpoint="http://browser.test/fetch",
scraper_proxy_url=None,
avito_proxy_max_rotations=0,
avito_serp_ok_not_banned=True,
avito_proxy_rotate_settle_s=0.0,
proxy_rotate_attempts=1,
proxy_rotate_attempt_timeout_s=1.0,
scraper_skip_seen_today=False,
)
@pytest.mark.parametrize(
("page2_html", "expected_complete"),
[(None, False), ("", True)],
)
async def test_partial_bucket_is_not_complete(
page2_html: str | None, expected_complete: bool
) -> None:
"""Страница 2 из 3 выпала → бакет НЕ «сделан»; все три пришли → «сделан».
Контрольная половина обязательна: реализация «всегда False» тоже прошла бы
одностороннюю проверку, но убила бы возобновление целиком.
Падает на origin/main: `on_bucket` вызывается двумя аргументами, признака полноты
в протоколе нет — частичный бакет неотличим от целого и попадает в done_buckets.
"""
scraper = AvitoScraper(_serp_config())
scraper.request_delay_sec = 0.0
calls: list[tuple[str, bool]] = []
def _on_bucket(key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg]
calls.append((key, complete))
async def _fetch_page(_self: Any, _slug: str, page: int, *_a: Any, **_k: Any) -> str | None:
return page2_html if page == 2 else f""
with (
patch.object(AvitoScraper, "_fetch_rooms_page_html", _fetch_page),
patch.object(
AvitoScraper,
"_parse_html",
lambda _self, html, **_k: [MagicMock(source_id=html, listing_segment="secondary")],
),
):
await scraper._paginate_leaf_bucket(
room_slug="kvartiry_1_komnatnye",
room_label="room_1_komn",
lo=0,
hi=3999999,
html="",
max_pages=3,
seen={},
price_cap_per_bucket=1400,
max_pages_per_bucket=100,
concurrency=2,
secondary_only=False,
on_bucket=_on_bucket,
skip_buckets=None,
expected_total=3 * 50,
)
assert [c[1] for c in calls] == [expected_complete]
async def test_pipeline_keeps_partial_bucket_out_of_checkpoint() -> None:
"""Пайплайн: лоты частичного бакета СОХРАНЕНЫ, но в чекпоинт он не попал.
Именно здесь «видимая потеря» (перескрап) не превращается в «невидимую»: пропустить
частичный бакет на следующем прогоне значит не перечитать его страницы уже никогда.
"""
finals: list[dict[str, Any]] = []
class _Recorder:
def is_cancelled(self, *_a: Any, **_k: Any) -> bool:
return False
def update_heartbeat(self, *_a: Any, **_k: Any) -> None:
pass
def mark_done(self, _db: Any, _rid: int, counters: dict[str, Any]) -> None:
finals.append(dict(counters))
async def _fetch(*_a: Any, on_bucket: Any = None, **_k: Any) -> None:
on_bucket("room_1_komn:0:3999999", [MagicMock(source_id="a1")], True)
on_bucket("room_1_komn:4000000:4999999", [MagicMock(source_id="a2")], False)
scraper = MagicMock()
scraper.__aenter__ = AsyncMock(return_value=scraper)
scraper.__aexit__ = AsyncMock(return_value=None)
scraper.fetch_all_secondary = _fetch
with (
patch(f"{PFX}.AvitoScraper", return_value=scraper),
patch(f"{PFX}.save_listings", MagicMock(return_value=(1, 0))),
patch(f"{PFX}.runs", _Recorder()),
):
counters = await run_avito_full_load(
MagicMock(), run_id=1, config=_serp_config(), matcher=MagicMock()
)
assert finals[0]["done_buckets"] == ["room_1_komn:0:3999999"]
assert finals[0]["partial_buckets"] == 1
assert counters.unique_fetched == 2, "лоты частичного бакета обязаны быть сохранены"