gendesign/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py
bot-backend d9b4f124a6
All checks were successful
CI / changes (pull_request) Successful in 10s
CI Trade-In / changes (pull_request) Successful in 11s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
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 4m24s
CI / backend-tests (pull_request) Has been skipped
feat(tradein/scheduler): планировщик подхватывает чекпоинт оборванного прогона (#930)
#930 сделал обе половины механизма — запись точки (counters.done_buckets) и её
чтение (run_*_full_load(resume_run_id=...)) — но единственным входом оставил
админку. У avito full-load её нет вовсе, а планировщик передавал resume_run_id
литеральным None. То есть боевого пути возобновления не существовало ни дня:
433 корзины в 30 оборванных прогонах за 90 суток (avito 242, cian 134,
exhaustive 57) перебирались заново, включая 35 корзин прогона 3547, убитого
деплоем на третьем часу и лежащего с 09.08.

Одной строки resume_run_id=prev было бы мало и опасно — правка из четырёх частей.

1. Полнота корзины (иначе видимая потеря стала бы невидимой). Ключ
   done_buckets одинаков у целиком и частично собранной корзины, а частичность
   возникает тремя путями: страница молча выпала (page_html=None), исключение
   страницы проглочено gather'ом, признанный tail-loss/hard-cap. SERP-слой
   теперь отдаёт признак полноты третьим аргументом on_bucket, пайплайн пишет
   в чекпоинт ТОЛЬКО полную корзину, частичные считает (partial_buckets).

2. Точка стала монотонной. update_heartbeat/mark_* мержат counters-jsonb
   вместо замены: раньше _on_progress и фоновый heartbeat cian'а (каждые 60 с)
   слали counters без done_buckets и СТИРАЛИ точку, а mark_failed уничтожал её
   насовсем — 15 оборванных прогонов с 45k собранных лотов и без ключа вовсе.

3. Тождество задания. Кандидат — ПОСЛЕДНИЙ прогон того же source (после 'done'
   возобновлять нечего), params сверяются побайтово: incremental_days меняет
   СМЫСЛ ключа, price_cap_per_bucket — само дерево бисекции. Статусы banned/
   cancelled/failed; zombie исключён намеренно (reap не убивает процесс,
   точка может двигаться). Срок годности — такт источника плюс сутки сетки
   запуска. Цепочка возобновлений ограничена STALE_DIGEST_INTERVAL_FACTOR-1.

4. Наблюдаемость. Вердикт (resume_from/resume_reason/resume_buckets/
   resume_chain) пишется в counters нового прогона: no_prev_run, status_*,
   params_changed, no_checkpoint, checkpoint_stale, chain_limit, ok.

Плюс skip_buckets доехал до инкрементальной ветки avito — без этого resume в
боевом режиме avito_full_load (incremental_days=7) был бы чистым no-op.
2026-08-12 21:40:41 +05:00

295 lines
14 KiB
Python
Raw Permalink 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.

"""#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), ("<page2/>", 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"<page{page}/>"
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="<page1/>",
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, "лоты частичного бакета обязаны быть сохранены"