From 1ff6699b958aab09e9e2d4a3ee8e1218a29e69f0 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Wed, 12 Aug 2026 18:51:02 +0000 Subject: [PATCH] =?UTF-8?q?feat(tradein/scheduler):=20=D0=BF=D0=BB=D0=B0?= =?UTF-8?q?=D0=BD=D0=B8=D1=80=D0=BE=D0=B2=D1=89=D0=B8=D0=BA=20=D0=BF=D0=BE?= =?UTF-8?q?=D0=B4=D1=85=D0=B2=D0=B0=D1=82=D1=8B=D0=B2=D0=B0=D0=B5=D1=82=20?= =?UTF-8?q?=D1=87=D0=B5=D0=BA=D0=BF=D0=BE=D0=B8=D0=BD=D1=82=20=D0=BE=D0=B1?= =?UTF-8?q?=D0=BE=D1=80=D0=B2=D0=B0=D0=BD=D0=BD=D0=BE=D0=B3=D0=BE=20=D0=BF?= =?UTF-8?q?=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=D0=B0=20(#930)=20(#2845)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../test_930_scheduler_resume_checkpoint.py | 295 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 82 ++++- .../src/scraper_kit/orchestration/runs.py | 32 +- .../scraper_kit/orchestration/scheduler.py | 132 +++++++- .../src/scraper_kit/providers/avito/serp.py | 61 +++- .../src/scraper_kit/providers/cian/serp.py | 18 +- 6 files changed, 588 insertions(+), 32 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py diff --git a/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py b/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py new file mode 100644 index 00000000..cc0ba37a --- /dev/null +++ b/tradein-mvp/backend/tests/test_930_scheduler_resume_checkpoint.py @@ -0,0 +1,295 @@ +"""#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, "лоты частичного бакета обязаны быть сохранены" diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 5c090d75..1f3bdd1a 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -2953,6 +2953,11 @@ class CianFullLoadCounters: detail_enriched: int = 0 detail_failed: int = 0 errors_count: int = 0 + # Бакеты, отданные SERP-слоем как НЕполные (страница выпала / исключение + # проглочено / hard-cap): лоты сохранены, но в done_buckets бакет не попал. + # Без этого счётчика «сколько бакетов чекпоинт не покрывает» видно только грепом + # логов, которые теряются при редеплое. + partial_buckets: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3007,8 +3012,32 @@ async def run_cian_full_load( done: set[str] = set(skip_set) # накапливаем завершённые бакеты этого прогона - def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] - """Инкрементальный save после каждого leaf-бакета. Дописывает bucket_key в done.""" + def _mark_bucket(bucket_key: str, complete: bool) -> None: + """В чекпоинт — только ПОЛНОСТЬЮ собранный бакет; частичный лишь считаем. + + Ключ done_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета, + поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск + частичного бакета на следующем прогоне означал бы, что его непрочитанные + страницы не перечитает уже никто, а счётчики покажут успех. + """ + if complete: + done.add(bucket_key) + return + counters.partial_buckets += 1 + logger.warning( + "cian-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, " + "следующий прогон перечитает его целиком (partial_buckets=%d)", + run_id, + bucket_key, + counters.partial_buckets, + ) + + def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg] + """Инкрементальный save после каждого leaf-бакета. + + complete=False → лоты сохраняем, бакет в чекпоинт не пишем + (см. _mark_bucket). Дефолт True — для вызывающих без пагинации. + """ nonlocal done if runs.is_cancelled(db, run_id): logger.info("cian-full-load run_id=%d: cancel detected in on_bucket", run_id) @@ -3023,7 +3052,7 @@ async def run_cian_full_load( ) raise RuntimeError("shutdown") if not lots: - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)} ) @@ -3042,7 +3071,7 @@ async def run_cian_full_load( counters.saved_inserted += inserted counters.saved_updated += updated counters.unique_fetched += len(lots) - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "cian-full-load run_id=%d: bucket %s saved ins=%d upd=%d total_unique=%d", @@ -3225,7 +3254,8 @@ async def run_cian_full_load( except NoProxyAvailableError as exc: # #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан - # сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет. + # сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который + # его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930). logger.error("cian-full-load run_id=%d: no proxy available — %s", run_id, exc) counters.errors_count += 1 runs.mark_banned( @@ -3467,7 +3497,8 @@ async def run_yandex_full_load( except NoProxyAvailableError as exc: # #2687, см. ту же ветку в run_avito_full_load: наш отказ (пул пуст) обязан - # сохранять done_buckets, а generic-RuntimeError → mark_failed его теряет. + # сохранять done_buckets, а generic-RuntimeError уходил в mark_failed, который + # его терял (теперь counters мержатся — runs.update_heartbeat/mark_*, #930). logger.error("yandex-full-load run_id=%d: no proxy available — %s", run_id, exc) counters.errors_count += 1 runs.mark_banned( @@ -3527,6 +3558,8 @@ class AvitoFullLoadCounters: saved_inserted: int = 0 saved_updated: int = 0 errors_count: int = 0 + # Бакеты, отданные SERP-слоем как НЕполные — см. CianFullLoadCounters. + partial_buckets: int = 0 def to_dict(self) -> dict[str, int]: return {f.name: getattr(self, f.name) for f in fields(self)} @@ -3599,8 +3632,32 @@ async def run_avito_full_load( ) done: set[str] = set(skip_set) - def _on_bucket(bucket_key: str, lots: list) -> None: # type: ignore[type-arg] - """Инкрементальный save после каждого leaf-бакета.""" + def _mark_bucket(bucket_key: str, complete: bool) -> None: + """В чекпоинт — только ПОЛНОСТЬЮ собранный бакет; частичный лишь считаем. + + Ключ done_buckets ("room:lo:hi") одинаков у целого и у недособранного бакета, + поэтому признак полноты приезжает отдельным аргументом из SERP-слоя. Пропуск + частичного бакета на следующем прогоне означал бы, что его непрочитанные + страницы не перечитает уже никто, а счётчики покажут успех. + """ + if complete: + done.add(bucket_key) + return + counters.partial_buckets += 1 + logger.warning( + "avito-full-load run_id=%d: bucket=%s собран ЧАСТИЧНО — в чекпоинт НЕ пишем, " + "следующий прогон перечитает его целиком (partial_buckets=%d)", + run_id, + bucket_key, + counters.partial_buckets, + ) + + def _on_bucket(bucket_key: str, lots: list, complete: bool = True) -> None: # type: ignore[type-arg] + """Инкрементальный save после каждого leaf-бакета. + + complete=False → лоты сохраняем, бакет в чекпоинт не пишем + (см. _mark_bucket). Дефолт True — для вызывающих без пагинации. + """ nonlocal done if runs.is_cancelled(db, run_id): logger.info("avito-full-load run_id=%d: cancel detected in on_bucket", run_id) @@ -3613,7 +3670,7 @@ async def run_avito_full_load( ) raise RuntimeError("shutdown") if not lots: - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)} ) @@ -3632,7 +3689,7 @@ async def run_avito_full_load( counters.saved_inserted += inserted counters.saved_updated += updated counters.unique_fetched += len(lots) - done.add(bucket_key) + _mark_bucket(bucket_key, complete) runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)}) logger.info( "avito-full-load run_id=%d: bucket=%s saved ins=%d upd=%d total_unique=%d", @@ -3699,9 +3756,10 @@ async def run_avito_full_load( except NoProxyAvailableError as exc: # #2687: пул опустел mid-run. Ветка стоит ДО generic-RuntimeError намеренно — # NoProxyAvailableError его подкласс, и без неё отказ уходил в mark_failed, - # который (в отличие от mark_banned) НЕ пишет done_buckets. То есть чекпоинт + # который (в отличие от mark_banned) НЕ передавал done_buckets. То есть чекпоинт # терялся ровно на НАШЕМ отказе — том исходе, для которого #2686 требовал его - # сохранять наравне с блокировкой площадкой. + # сохранять наравне с блокировкой площадкой. Диагноз ban_kind='infra' ветка + # даёт по-прежнему; сам чекпоинт с #930-мержем counters переживает и mark_failed. logger.error("avito-full-load run_id=%d: no proxy available — %s", run_id, exc) counters.errors_count += 1 runs.mark_banned( diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index c0996e1b..86729701 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -100,8 +100,8 @@ def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: # 1. Побочная функция 'banned' — сохранение done_buckets-чекпоинта (mark_failed -# его теряет) — нужна ОБОИМ исходам. Оставив статус, получаем её даром; расщепив -# статус, пришлось бы дублировать её в каждом потребителе. +# его тогда терял) — нужна ОБОИМ исходам. Оставив статус, получаем её даром; +# расщепив статус, пришлось бы дублировать её в каждом потребителе. # 2. Новое значение статуса пришлось бы доучить пяти местам, каждое из которых # молча даёт неверный ответ, если про него забыть: CHECK-констрейнт схемы, # IN-списки обоих сторожей (_alert_if_consecutive_failures / _zero_results), @@ -570,12 +570,21 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = return int(row.id) -def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None: - """UPDATE heartbeat_at + counters=:counters + total_seen/new_count колонки. +def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None: + """UPDATE heartbeat_at + counters (МЕРЖ, не замена) + total_seen/new_count колонки. total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926). COALESCE: если ключа нет в counters — старое значение колонки сохраняется. + + `counters || :counters` вместо замены (#930 добивка): чекпоинт `done_buckets` + писали ТОЛЬКО сайты, знающие о нём (`_on_bucket`), а heartbeat'ы, которые о нём не + знают, целиком перезаписывали объект и СТИРАЛИ точку. На проде это давало + немонотонную точку: `_on_progress` (после каждой комнатности) и фоновый heartbeat + cian'а (каждые 60 с) отправляли `counters.to_dict()` без ключа — то есть у cian + точка в БД жила лишь от сохранения бакета до ближайшего тика. Прод-след: 15 + оборванных прогонов с доказанной работой (35 706 + 9 222 fetched) и БЕЗ ключа + вообще. Мерж делает точку монотонной для любого писателя, а не только для знающих. """ total_seen, new_count = _column_counts(counters) db.execute( @@ -583,7 +592,7 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None """ UPDATE scrape_runs SET heartbeat_at = clock_timestamp(), - counters = CAST(:counters AS jsonb), + counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) WHERE id = :run_id @@ -641,7 +650,7 @@ def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None: UPDATE scrape_runs SET status = 'done', finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), - counters = CAST(:counters AS jsonb), + counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) WHERE id = :run_id AND status = 'running' @@ -681,7 +690,8 @@ def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) UPDATE scrape_runs SET status = 'failed', finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), - error = :error, counters = CAST(:counters AS jsonb), + error = :error, + counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb), total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) WHERE id = :run_id AND status = 'running' @@ -715,8 +725,9 @@ def mark_banned( Per migration 015 — 'banned' задокументирован как 'Avito вернул 403/captcha'. Отличается от 'failed': прогон оборван внешним/блокирующим условием, а не нашим - багом, и — важно — СОХРАНЯЕТ done_buckets-чекпоинт в counters (mark_failed его - теряет). Cooldown 2-4 часа. + багом. Чекпоинт done_buckets раньше сохранял только этот финализатор — теперь + counters мержатся во всех (см. update_heartbeat), и точка переживает любой из них. + Cooldown 2-4 часа. `ban_kind` разводит два исхода, которые раньше схлопывались в один статус: - BAN_KIND_PLATFORM — площадка нас заблокировала (firewall/403/captcha); @@ -742,7 +753,8 @@ def mark_banned( UPDATE scrape_runs SET status = 'banned', finished_at = clock_timestamp(), heartbeat_at = clock_timestamp(), - error = :error, counters = CAST(:counters AS jsonb), + error = :error, + counters = COALESCE(counters, CAST('{}' AS jsonb)) || CAST(:counters AS jsonb), ban_kind = :ban_kind, total_seen = COALESCE(CAST(:total_seen AS int), total_seen), new_count = COALESCE(CAST(:new_count AS int), new_count) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py index 06d778a8..8b0b5009 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/scheduler.py @@ -497,6 +497,132 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) return run_id +# ── подхват чекпоинта оборванного прогона (#930, вторая половина) ──────────── +# #930 сделал чекпоинт (`counters.done_buckets`) и приёмную сторону +# (`run_*_full_load(resume_run_id=...)`), но единственным входом оставил админку. У +# avito full-load её нет вовсе, поэтому 299 из 433 «впустую перебранных» корзин за 90 +# суток не имели НИКАКОГО пути возобновления, даже ручного. Планировщик передавал +# resume_run_id=None литералом. +# +# Точка берётся только когда выполнены ВСЕ условия ниже; иначе прогон честно начинает с +# нуля, а ПРИЧИНА пишется в его counters (молчаливый отказ неотличим от отсутствия +# правки — см. _resume_decision). + +# 'zombie' НЕ в списке НАМЕРЕННО. reap_zombies снимает пометку 'running', но НЕ убивает +# процесс (прямо задокументировано в app/tasks/listing_source_snapshot.py) — а +# has_running_run гейтит claim именно по статусу. То есть после reap'а старый сборщик +# может продолжать писать в ту же строку: подхват читал бы ДВИЖУЩУЮСЯ точку и запускал +# второй сборщик на ту же площадку. 147 корзин в 10 zombie-прогонах за 90 суток +# остаются несобранными сознательно — это цена, а не недосмотр. +# 'failed' в списке: его чекпоинт больше не стирается финализатором (runs.py, мерж +# counters), а причина отказа («наш баг») ничего не говорит о полноте УЖЕ записанных +# корзин — они записаны тем же heartbeat'ом, что и у banned. +_RESUME_STATUSES = frozenset({"banned", "cancelled", "failed"}) + +# Длина цепочки возобновлений. Не круглое число: STALE_DIGEST_INTERVAL_FACTOR (=3) — +# уже существующий в этом файле порог «источник не собирал дольше 3× своего такта = +# сломан». Цепочка не имеет права отодвинуть полный обход дальше этой же черты, +# поэтому подряд идущих подхватов допускается на один меньше: прогоны 1 и 2 могут +# продолжать предшественника, третий обязан пойти с нуля. При такте avito 7 суток это +# гарантирует попытку полного обхода не реже, чем раз в 21 сутки — ровно в тот момент, +# когда сводка объявляет источник просроченным. +_MAX_RESUME_CHAIN = STALE_DIGEST_INTERVAL_FACTOR - 1 + +# Кандидат — ПОСЛЕДНИЙ прогон источника, а не последний подходящий: если после обрыва +# уже прошёл полный ('done') прогон, дерево обойдено и возобновлять нечего. Строки +# 'skipped' — бухгалтерия планировщика, а не прогоны, поэтому не в счёт. +_RESUME_CANDIDATE_SQL = text(""" + WITH cur AS ( + SELECT source, params FROM scrape_runs WHERE id = CAST(:rid AS bigint) + ) + SELECT r.id AS prev_id, + r.status AS prev_status, + r.counters AS prev_counters, + (r.params IS NOT DISTINCT FROM cur.params) AS same_params, + EXTRACT(EPOCH FROM (clock_timestamp() - r.heartbeat_at)) / 3600.0 AS age_h, + cur.params ->> 'interval_days' AS interval_days + FROM scrape_runs r, cur + WHERE r.source = cur.source + AND r.id <> CAST(:rid AS bigint) + AND r.status <> 'skipped' + ORDER BY r.started_at DESC + LIMIT 1 +""") + + +def _resume_decision(row: Any) -> tuple[int | None, dict[str, Any]]: + """Решение «подхватывать ли точку» + счётчики-объяснение. Чистая функция. + + Возвращает (resume_run_id | None, counters-заготовка нового прогона). Причина + отказа — машиночитаемый слаг в `resume_reason`, по нему «предыдущего прогона не + было» отличается от «параметры разъехались» ЗАПРОСОМ, а не грепом логов. + + Срок годности точки — такт источника ПЛЮС сутки (`interval_days` из его же params). + Ни одно из слагаемых не выбрано произвольно. Такт — объявленный самим расписанием + срок, в течение которого собранное считается свежим; точка старше него пережила + цикл, в котором источник обязан был обойти дерево целиком. Сутки — сетка запуска: + `compute_next_run_at` выбирает ДЕНЬ (сегодня+interval_days) и случайное время внутри + окна, поэтому два соседних запуска отстоят друг от друга на interval_days ± меньше + суток. Без этого слагаемого точку отвергал бы jitter расписания, а не устаревание: + прогон 3547 убит деплоем 09.08 16:53, расписание 139 подхватит его 16.08 13:37 — + 164.6 ч при такте 168 ч, запас 3.4 ч при ширине окна 2 ч. Пропущенный цикл в окно + всё равно не влезает: для avito это 13 суток против порога 8. + """ + if row is None: + return None, {"resume_from": None, "resume_reason": "no_prev_run", "resume_chain": 0} + + prev_counters = row.prev_counters if isinstance(row.prev_counters, dict) else {} + done_buckets = prev_counters.get("done_buckets") + done_n = len(done_buckets) if isinstance(done_buckets, list) else 0 + chain_raw = prev_counters.get("resume_chain") + prev_chain = chain_raw if isinstance(chain_raw, int) else 0 + verdict: dict[str, Any] = { + "resume_from": None, + "resume_candidate": int(row.prev_id), + "resume_buckets": done_n, + "resume_chain": 0, + } + + if row.prev_status not in _RESUME_STATUSES: + verdict["resume_reason"] = f"status_{row.prev_status}" + elif not row.same_params: + verdict["resume_reason"] = "params_changed" + elif done_n == 0: + verdict["resume_reason"] = "no_checkpoint" + elif row.age_h is None or float(row.age_h) > 24.0 * ( + _schedule_interval_days(row.interval_days) + 1 + ): + verdict["resume_reason"] = "checkpoint_stale" + elif prev_chain >= _MAX_RESUME_CHAIN: + verdict["resume_reason"] = "chain_limit" + else: + verdict["resume_from"] = int(row.prev_id) + verdict["resume_reason"] = "ok" + verdict["resume_chain"] = prev_chain + 1 + return int(row.prev_id), verdict + return None, verdict + + +def _pick_resume(db: Session, run_id: int) -> int | None: + """Чекпоинт какого прогона наследует `run_id` (или None) + запись вердикта. + + Тождество задания сверяется РОВНО по тем полям, которыми задание задаётся: source + (кандидат ищется в пределах одного source) и params целиком, побайтово. Кандидат + сравнивается с ТЕКУЩИМ прогоном, а не с расписанием, потому что именно params + прогона поехали в pipeline. На проде за 90 суток 89 корзин из 433 (21%) лежат в + прогонах, чьи params отличаются от следующего — по ним пропуск был бы неверным: + `incremental_days` меняет СМЫСЛ ключа (дочитано до watermark ≠ бакет перебран), а + `price_cap_per_bucket` меняет само дерево бисекции, то есть какие ключи существуют. + """ + row = db.execute(_RESUME_CANDIDATE_SQL, {"rid": run_id}).fetchone() + resume_run_id, verdict = _resume_decision(row) + # Вердикт кладём в counters НОВОГО прогона: `update_heartbeat` мержит jsonb, поэтому + # последующие heartbeat'ы пайплайна его не затрут и он доживёт до финализатора. + _kit_runs.update_heartbeat(db, run_id, verdict) + logger.info("scheduler: resume run_id=%d — %s", run_id, verdict) + return resume_run_id + + def _defer_next_run_at(db: Session, schedule_row: dict[str, Any]) -> None: """Сдвинуть next_run_at на следующее окно БЕЗ создания run (#1522). @@ -705,7 +831,7 @@ async def _job_avito_full_load( concurrency=int(params.get("concurrency", 5)), request_delay_sec=float(params.get("request_delay_sec", 7.0)), secondary_only=bool(params.get("secondary_only", True)), - resume_run_id=None, + resume_run_id=_pick_resume(db, run_id), incremental_days=incremental_days, ) @@ -725,7 +851,7 @@ async def _job_avito_full_load_exhaustive( concurrency=int(params.get("concurrency", 5)), request_delay_sec=float(params.get("request_delay_sec", 7.0)), secondary_only=bool(params.get("secondary_only", True)), - resume_run_id=None, + resume_run_id=_pick_resume(db, run_id), incremental_days=None, ) @@ -796,7 +922,7 @@ async def _job_cian_full_load( request_delay_sec=float(params.get("request_delay_sec", 4.0)), enrich_detail=bool(params.get("enrich_detail", False)), detail_top_n=int(params.get("detail_top_n", 0)), - resume_run_id=None, + resume_run_id=_pick_resume(db, run_id), ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py index 3c3d717a..95e2618a 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/avito/serp.py @@ -1095,7 +1095,9 @@ class AvitoScraper(BaseScraper): leaf-бакета. Может быть async или sync. Исключение прерывает прогон. on_progress: опциональный callback(unique_count) для heartbeat (per room-bucket). skip_buckets: множество ключей «room_label:lo:hi» уже завершённых бакетов — - пагинация и on_bucket для них пропускаются. Probe-запросы выполняются. + пагинация и on_bucket для них пропускаются. В exhaustive-режиме probe-запросы + всё равно выполняются (skip проверяется уже в листе, после probe), в + инкрементальном probe'а нет — там пропускается весь бакет целиком. since: если None (default) — EXHAUSTIVE bisection-обход (поведение без изменений). Если задана date — INCREMENTAL: на каждый (комнатность × seed-брекет) последовательная пагинация newest-first с ранней остановкой, @@ -1141,6 +1143,7 @@ class AvitoScraper(BaseScraper): max_pages_per_bucket=max_pages_per_bucket, secondary_only=secondary_only, on_bucket=on_bucket, + skip_buckets=skip_buckets, ) else: await self._walk_price_range( @@ -1390,14 +1393,17 @@ class AvitoScraper(BaseScraper): # иначе фетчим её как обычную страницу (открытый брекет без probe-html). sem = asyncio.Semaphore(concurrency) first_url = self._build_rooms_url(room_slug, 1, _lo_param, _hi_param) + dropped_pages = 0 # страницы, не отдавшие карточки по отказу (не по пустоте) async def _one_page(p: int) -> list[ScrapedLot]: + nonlocal dropped_pages if p == 1 and html is not None: return self._parse_html(html, source_url_base=first_url) async with sem: page_html = await self._fetch_rooms_page_html(room_slug, p, _lo_param, _hi_param) await asyncio.sleep(self.request_delay_sec) if page_html is None: + dropped_pages += 1 logger.warning( "avito: page_html=None %s [%d, %s] page=%d — skipping page", room_label, @@ -1420,6 +1426,7 @@ class AvitoScraper(BaseScraper): # Блокировки пробрасываем наверх (mark_banned в pipeline-обёртке). if isinstance(res, AvitoBlockedError | AvitoRateLimitedError): raise res + dropped_pages += 1 logger.warning( "avito: page exception %s [%d, %s] page=%d — %r", room_label, @@ -1464,8 +1471,25 @@ class AvitoScraper(BaseScraper): if key: seen[key] = lot + # ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ───── + # `bucket_key` — это "room:lo:hi" и больше ничего, поэтому «сделано» едет + # отдельным аргументом. Три пути частичности сходятся здесь: + # 1. страница молча выпала (page_html=None выше); + # 2. исключение страницы проглочено gather'ом (return_exceptions=True); + # 3. признанный tail-loss — probe провалился (expected_total=None, открытый + # брекет) или страниц нужно больше, чем max_pages. + # Ни один из них не оставлял следа в чекпоинте: следующий прогон видел ключ и + # пропускал бакет. Пока признак не доехал до done_buckets, включать подхват + # нельзя — видимая потеря (перескрап) стала бы невидимой (пропуск страниц). + pages_needed = ( + math.ceil(expected_total / _AVITO_OFFERS_PER_PAGE) + if expected_total is not None + else None + ) + complete = dropped_pages == 0 and pages_needed is not None and pages_needed <= max_pages logger.info( - "avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d", + "avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d " + "complete=%s dropped_pages=%d", room_label, lo, _hi_repr, @@ -1473,11 +1497,13 @@ class AvitoScraper(BaseScraper): collected_this_bucket, dropped_nb, len(seen), + complete, + dropped_pages, ) # ── on_bucket callback: инкрементальный save ────────────────────────── if on_bucket is not None and bucket_lots: - res_cb = on_bucket(bucket_key, bucket_lots) + res_cb = on_bucket(bucket_key, bucket_lots, complete) if inspect.isawaitable(res_cb): await res_cb @@ -1493,6 +1519,7 @@ class AvitoScraper(BaseScraper): max_pages_per_bucket: int, secondary_only: bool, on_bucket: Callable[..., Any] | None, + skip_buckets: set[str] | None = None, ) -> None: """INCREMENTAL пагинация одного (комнатность × seed-брекет) с ранней остановкой. @@ -1514,16 +1541,37 @@ class AvitoScraper(BaseScraper): bucket_key, secondary_only-фильтр и дедуп в seen — идентичны _paginate_leaf_bucket. on_bucket вызывается один раз для собранного бакета (async/sync-aware). AvitoBlockedError/AvitoRateLimitedError из page-фетчей пробрасываются наверх. + + skip_buckets: ключи, дочитанные ПРЕДЫДУЩИМ прогоном с ТЕМИ ЖЕ params (тождество + задания проверяет планировщик, `_pick_resume`). В инкрементальном режиме + «сделано» значит «дочитал до watermark `since`», а не «перебрал бакет целиком», + поэтому смешивать такой ключ с exhaustive-ключом нельзя — они дословно совпадают + (6 из 11 seed-ключей), но означают разное. Пропуск здесь безопасен по покрытию + ровно потому, что окно ретроспективы не уже такта (#2674, гарантируется + _job_avito_full_load): бакет, дочитанный до watermark N суток назад, следующий + плановый прогон перечитает со своим since = сегодня−N и увидит всё, что успело + появиться. Раньше аргумент сюда не передавался вовсе — resume в боевом + (инкрементальном) режиме avito_full_load был чистым no-op: 134 из 242 корзин. + + complete=False (см. `_paginate_leaf_bucket`) отдаётся, когда бакет НЕ дочитан до + watermark: страница выпала, или страниц не хватило (max_pages), или остановка + произошла по grace-эвристике «2 подряд недатированные страницы» — там watermark + не доказан, а не достигнут. """ _lo_param = lo if lo > 0 else None _hi_param = hi # None → _build_rooms_url не ставит pmax _hi_repr = "open" if hi is None else str(hi) bucket_key = f"{room_label}:{lo}:{_hi_repr}" + if skip_buckets and bucket_key in skip_buckets: + logger.info("avito: skip bucket %s — already read to watermark (resume)", bucket_key) + return + bucket_lots: list[ScrapedLot] = [] pages_fetched = 0 not_fresh_streak = 0 # подряд идущие не-свежие (all None-или-старые) страницы stop_reason = "max-pages" # перетирается ниже на реальную причину + complete = False # дочитан ли бакет до watermark; True только на честных стопах for p in range(1, max_pages_per_bucket + 1): # Последовательный фетч (НЕ asyncio.gather): early-stop требует читать @@ -1540,6 +1588,7 @@ class AvitoScraper(BaseScraper): page_lots = self._parse_html(page_html, source_url_base=page_url) if not page_lots: stop_reason = "end-of-pages (empty parse)" + complete = True # выдача кончилась — читать в этом брекете больше нечего break bucket_lots.extend(page_lots) @@ -1556,6 +1605,7 @@ class AvitoScraper(BaseScraper): # Есть даты, но все < since → newest-first гарантирует, что дальше # только старее → стоп немедленно. stop_reason = "early-stop (page all older than since)" + complete = True # watermark достигнут — ровно то, что значит «сделано» break # Все карточки undated (None) → не стопим сразу (могут быть свежие без # даты), но копим streak; 2 подряд недатированные страницы → grace-стоп. @@ -1581,7 +1631,7 @@ class AvitoScraper(BaseScraper): logger.info( "avito: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d " - "unique_total=%d incremental since=%s stop=%s", + "unique_total=%d incremental since=%s stop=%s complete=%s", room_label, lo, _hi_repr, @@ -1591,11 +1641,12 @@ class AvitoScraper(BaseScraper): len(seen), since.isoformat(), stop_reason, + complete, ) # ── on_bucket callback: инкрементальный save ────────────────────────── if on_bucket is not None and bucket_lots: - res_cb = on_bucket(bucket_key, bucket_lots) + res_cb = on_bucket(bucket_key, bucket_lots, complete) if inspect.isawaitable(res_cb): await res_cb diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py index e046ef67..67b5c302 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/cian/serp.py @@ -621,8 +621,10 @@ class CianScraper(BaseScraper): # Страница 1 уже есть (html из probe выше); остальные — параллельно. sem = asyncio.Semaphore(concurrency) + dropped_pages = 0 if html else 1 # пустой probe-html = страница 1 не собрана async def _one_page(p: int) -> list[ScrapedLot]: + nonlocal dropped_pages if p == 1: # Используем уже полученный HTML от probe return self._parse_serp_html(html) if html else [] @@ -630,6 +632,7 @@ class CianScraper(BaseScraper): page_html = await self._fetch_page_html(rooms, p, _lo_param, hi) await asyncio.sleep(self.request_delay_sec) if page_html is None: + dropped_pages += 1 logger.warning( "cian: page_html=None %s [%d, %s] page=%d — skipping page", room_label, @@ -648,6 +651,7 @@ class CianScraper(BaseScraper): bucket_lots: list[ScrapedLot] = [] for p_idx, res in enumerate(page_results, start=1): if isinstance(res, BaseException): + dropped_pages += 1 logger.warning( "cian: page exception %s [%d, %s] page=%d — %r", room_label, @@ -674,8 +678,16 @@ class CianScraper(BaseScraper): if key: seen[key] = lot + # ── Полнота бакета: ключ у целиком и частично собранного ОДИНАКОВ ───── + # См. avito/serp.py — те же два пути частичности (выпавшая страница, + # проглоченное gather'ом исключение) плюс признанный hard-cap выше. У cian это + # тяжелее: в SERP-слое НЕТ класса блок-исключения вообще (ban-детект #2625 + # агрегатный, на выходе из скрапера), поэтому капча посреди бакета приходит сюда + # как page_html=None и раньше давала «сделанный» бакет из уцелевших страниц. + complete = dropped_pages == 0 and pages_needed <= max_pages logger.info( - "cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d", + "cian: %s [%d, %s] paginated=%d pages collected=%d dropped_nb=%d unique_total=%d " + "complete=%s dropped_pages=%d", room_label, lo, _hi_repr, @@ -683,11 +695,13 @@ class CianScraper(BaseScraper): collected_this_bucket, dropped_nb, len(seen), + complete, + dropped_pages, ) # ── on_bucket callback: инкрементальный save ────────────────────────── if on_bucket is not None and bucket_lots: - res_cb = on_bucket(bucket_key, bucket_lots) + res_cb = on_bucket(bucket_key, bucket_lots, complete) if inspect.isawaitable(res_cb): await res_cb