From b4e0618025ba39fa4b2ecf1e7199aca5e2033bd8 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 16:23:22 +0300 Subject: [PATCH] =?UTF-8?q?fix(trade-in):=20=D1=81=D0=B2=D0=B8=D0=BF=20?= =?UTF-8?q?=D0=94=D0=BE=D0=BC=D0=9A=D0=BB=D0=B8=D0=BA=D0=B0=20=D1=81=D0=BE?= =?UTF-8?q?=D1=85=D1=80=D0=B0=D0=BD=D1=8F=D0=B5=D1=82=20=D0=BB=D0=BE=D1=82?= =?UTF-8?q?=D1=8B=20=D0=BF=D0=BE=20=D0=BA=D0=BE=D1=80=D0=B7=D0=B8=D0=BD?= =?UTF-8?q?=D0=B0=D0=BC,=20=D0=B0=20=D0=BD=D0=B5=20=D0=BE=D0=B4=D0=BD?= =?UTF-8?q?=D0=B8=D0=BC=20=D0=BC=D0=B0=D1=85=D0=BE=D0=BC=20=D0=B2=20=D0=BA?= =?UTF-8?q?=D0=BE=D0=BD=D1=86=D0=B5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit На большой выдаче свип не отдавал НИЧЕГО. Замер: прогон 7344 (Москва, одиночный, прокси 13) за два часа прошёл 2 корзины из 6 — 'st' 34 минуты, '1' 86+ минут. Причина структурная: ДомКлик режет offset на 2000, поэтому диапазон делится бисекцией по цене и каждый лист пагинируется отдельно. У Москвы ≈23 690 лотов вторички против ≈6 300 у ЕКБ. Формула watchdog'а считает фетчи как 6 * pages и объём выдачи не учитывает вовсе, так что снятие по таймауту было гарантировано, а вместе с ним терялось всё собранное: сохранение было ОДНО, после всех корзин. Теперь save_listings зовётся из колбэка on_bucket сразу после каждой успешной корзины, и туда же переехал чекпоинт: done_buckets означает «собрано И сохранено». Раньше он писался из scraper.completed_buckets уже после except TimeoutError, то есть помечал пройденными корзины, у которых в БД ноль строк, — следующий прогон пропускал их через skip_buckets, и за два-три цикла чекпоинт закрывался целиком. Гард _saved (#2406) закрывал это «всё или ничего»; с инкрементальным сохранением он не нужен и снят. Колбэк в serp.py стоит ВНЕ try/except конкретной корзины — иначе generic обработчик проглотил бы исключение из него. Это же даёт кооперативную отмену по корзинам, которой у ДомКлика не было вовсе: is_cancelled проверялся только перед SERP-фазой, и повисший свип нельзя было снять до watchdog'а, все три часа держа один из двух узлов affinity='any'. Sentinel'ы RuntimeError("cancelled") и ("shutdown") и разбор ветки — тот же приём, что в run_cian_full_load. watchdog_sec — явный override формулы, читается из default_params в ОБОИХ хендлерах (kit-native и продуктовом, который перекрывает его ради кук Sber ID). None сохраняет прежнюю формулу байт-в-байт, как и on_bucket=None. Отдельно поправлены пять тестов, ломавшихся на переезде: их стабы подменяли fetch_city и не звали on_bucket, поэтому после правки проверяли мёртвую ветку — save_listings не вызывался вовсе. Теперь стабы вызывают колбэк, и проверяется реальный путь. Прогон: полный бэкенд-набор 6469 passed, 70 skipped. ruff check и format чисто. Claude-Session: https://claude.ai/code/session_01NQb6WeJtagZwZnUsSjDizs --- .../backend/app/services/product_handlers.py | 3 + .../tests/test_2406_sweep_edge_cases.py | 23 +- .../test_3369_domclick_cancel_checkpoint.py | 10 +- .../tests/test_domclick_incremental_save.py | 229 ++++++++++++++++++ .../test_scraper_kit_pipeline_parity2.py | 13 +- .../src/scraper_kit/orchestration/pipeline.py | 167 ++++++++++--- .../scraper_kit/orchestration/scheduler.py | 3 + .../scraper_kit/providers/domclick/serp.py | 18 ++ 8 files changed, 419 insertions(+), 47 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_domclick_incremental_save.py diff --git a/tradein-mvp/backend/app/services/product_handlers.py b/tradein-mvp/backend/app/services/product_handlers.py index 7cc7d6d6..2f29ba8f 100644 --- a/tradein-mvp/backend/app/services/product_handlers.py +++ b/tradein-mvp/backend/app/services/product_handlers.py @@ -544,6 +544,9 @@ async def _job_domclick_city_sweep( region_code=kit_resolve_region_code(params), resume_run_id=kit_pick_resume(db, run_id), cookies=cookies, + watchdog_sec=( + int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None + ), ) diff --git a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py index feeac535..fb8709ad 100644 --- a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py +++ b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py @@ -173,10 +173,14 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None: async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None: """Корзина без сохранённых строк не попадает в чекпоинт. - Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин. - Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из - корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы, - что следующий прогон пропустит её через skip_buckets навсегда (миграция 308). + С инкрементальным сохранением (fix/domclick-incremental-save) save_listings зовётся из + колбэка on_bucket ПОСЛЕ каждой корзины, а чекпоинт пишется там же и только + ПОСЛЕ успешного save. Поэтому корзина, чей save упал, в done_buckets не попадает, + хотя скрейпер её фетч прошёл (_s.completed_buckets её содержит). Снятая + watchdog'ом фаза — тот же инвариант: до on_bucket она не дошла. + + Отметить такую корзину пройденной значило бы, что следующий прогон пропустит её + через skip_buckets навсегда (миграция 308). """ class _Dc(_Scraper): @@ -185,10 +189,17 @@ async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> Non buckets_completed, buckets_total = 1, 6 completed_buckets = ["st"] # noqa: RUF012 - async def fetch_city(self, **_kw: Any) -> list[Any]: + async def fetch_city(self, **kw: Any) -> list[Any]: if broken == "fetch_timeout": raise TimeoutError - return [object()] + # Заглушка обязана ВЫЗВАТЬ колбэк: путь сохранения переехал внутрь цикла + # по корзинам (serp.py fetch_city), и стаб, который просто возвращает лоты, + # проверял бы мёртвую ветку — save_listings не был бы вызван вовсе. + on_bucket = kw.get("on_bucket") + lots = [object()] + if on_bucket is not None: + on_bucket("st", lots) + return lots saved: list[int] = [] diff --git a/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py index 9effa69d..9338047c 100644 --- a/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py +++ b/tradein-mvp/backend/tests/test_3369_domclick_cancel_checkpoint.py @@ -153,7 +153,15 @@ class _CleanSweepScraper: async def __aexit__(self, *_e: Any) -> None: return None - async def fetch_city(self, **_kw: Any) -> list[Any]: + async def fetch_city(self, **kw: Any) -> list[Any]: + # Чекпоинт с fix/domclick-incremental-save набирается ИЗ колбэка on_bucket, а не + # мержится в конце из completed_buckets — стаб обязан его вызвать, иначе + # проверялась бы мёртвая ветка. Лотов нет: корзина пройдена, но пустая — + # это законный случай, save_listings для неё не зовётся, чекпоинт пишется. + on_bucket = kw.get("on_bucket") + if on_bucket is not None: + for bucket in self.completed_buckets: + on_bucket(bucket, []) return [] diff --git a/tradein-mvp/backend/tests/test_domclick_incremental_save.py b/tradein-mvp/backend/tests/test_domclick_incremental_save.py new file mode 100644 index 00000000..ec969d4a --- /dev/null +++ b/tradein-mvp/backend/tests/test_domclick_incremental_save.py @@ -0,0 +1,229 @@ +"""fix/domclick-incremental-save: большая выдача (Москва ≈23 690 лотов против ЕКБ +≈6 300) теряла ВСЁ собранное при снятии свипа по watchdog — сохранение было ОДНО, +в самом конце `run_domclick_city_sweep`. Прод run 7344 (`domclick_city_sweep_moskva`): +за 2ч watchdog'а (11 100с) пройдено 2 бакета из 6. + +Три дефекта, три слоя тестов ниже: +1. Инкрементальное сохранение по бакетам (`fetch_city.on_bucket`) — секция 1. +2. Чекпоинт врал ("пройдено" без "сохранено") — секция 1, тест + `test_on_bucket_exception_interrupts_bucket_loop` документирует нюанс, на + котором строится фикс в pipeline.py: scraper.completed_buckets отмечает бакет + как "фетч прошёл" ДО вызова on_bucket, поэтому пайплайн больше не берёт + чекпоинт оттуда — только из факта успешного on_bucket (см. комментарий у + "Перенести счётчики" в run_domclick_city_sweep). +3. Watchdog не знал о размере выдачи — секция 2, `watchdog_sec` override. +""" + +from __future__ import annotations + +import os +from types import SimpleNamespace +from typing import Any +from unittest.mock import patch + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +import pytest +from scraper_kit.orchestration.pipeline import run_domclick_city_sweep +from scraper_kit.providers.domclick.serp import ( + ROOM_BUCKETS, + DomClickBlockedError, + DomClickScraper, +) + +PFX = "scraper_kit.orchestration.pipeline" + + +# ── Секция 1: DomClickScraper.fetch_city(on_bucket=...) ────────────────────── + + +def _scraper() -> DomClickScraper: + return DomClickScraper(SimpleNamespace(scraper_proxy_url=None)) + + +class _FakeFetcherCtx: + async def __aenter__(self) -> SimpleNamespace: + return SimpleNamespace(report_ban=lambda *_a, **_k: None) + + async def __aexit__(self, *_exc: Any) -> None: + return None + + +async def _run_fetch_city( + scraper: DomClickScraper, + *, + sweep_bucket: Any, + on_bucket: Any = None, + start: int = 0, + skip: set[str] | None = None, +) -> list[str]: + with ( + patch.object(scraper, "_sweep_bucket", sweep_bucket), + patch( + "scraper_kit.providers._base.build_browser_fetcher", + lambda *_a, **_k: _FakeFetcherCtx(), + ), + ): + return await scraper.fetch_city( + city_id=1, + pages=1, + start_bucket_index=start, + skip_buckets=skip, + on_bucket=on_bucket, + ) + + +async def test_on_bucket_called_once_per_bucket_with_only_its_own_lots() -> None: + """on_bucket зовётся по разу на каждый успешный бакет с лотами ИМЕННО его, + а не накопленным out_lots (главная регрессия #1 issue — было "одно сохранение + в конце").""" + + async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None: + out_lots.append(f"{rooms}-a") + out_lots.append(f"{rooms}-b") + + calls: list[tuple[str, list[str]]] = [] + s = _scraper() + result = await _run_fetch_city( + s, sweep_bucket=_fake_sweep, on_bucket=lambda n, lots: calls.append((n, list(lots))) + ) + + assert [c[0] for c in calls] == list(ROOM_BUCKETS) + for bucket, lots in calls: + assert lots == [f"{bucket}-a", f"{bucket}-b"], (bucket, lots, "получил чужие лоты") + assert len(result) == len(ROOM_BUCKETS) * 2 + + +async def test_on_bucket_skips_failed_and_blocked_buckets() -> None: + """Битый бакет (generic Exception) и бакет с QRATOR-блоком не отдают лоты в + on_bucket — там ничего не собрано/не гарантированно собрано.""" + + async def _fake_sweep_fail_middle(*, rooms: str, out_lots: list[str], **_kw: Any) -> None: + out_lots.append(f"{rooms}-x") + if rooms == "1": + raise ValueError("boom") + + calls_a: list[str] = [] + s_a = _scraper() + await _run_fetch_city( + s_a, sweep_bucket=_fake_sweep_fail_middle, on_bucket=lambda n, _lots: calls_a.append(n) + ) + assert "1" not in calls_a + # continue идёт дальше — остальные бакеты всё равно получают on_bucket. + assert calls_a == [b for b in ROOM_BUCKETS if b != "1"], calls_a + + async def _fake_sweep_block(*, rooms: str, out_lots: list[str], **_kw: Any) -> None: + out_lots.append(f"{rooms}-x") + if rooms == "1": + raise DomClickBlockedError("QRATOR") + + calls_b: list[str] = [] + s_b = _scraper() + await _run_fetch_city( + s_b, sweep_bucket=_fake_sweep_block, on_bucket=lambda n, _lots: calls_b.append(n) + ) + # break останавливает обход целиком — после блока ни один бакет не пробуется. + assert calls_b == ["st"], calls_b + + +async def test_on_bucket_exception_interrupts_bucket_loop() -> None: + """Исключение из on_bucket (канал кооперативной отмены) прерывает цикл по + ROOM_BUCKETS целиком — остальные бакеты не идут.""" + + async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None: + out_lots.append(f"{rooms}-x") + + visited: list[str] = [] + + def _on_bucket(name: str, lots: list[str]) -> None: + visited.append(name) + if name == "1": + raise RuntimeError("cancelled") + + s = _scraper() + with pytest.raises(RuntimeError, match="cancelled"): + await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=_on_bucket) + + assert visited == ["st", "1"], visited + # Скрейпер-уровень успел зафетчить оба бакета ДО того, как on_bucket поднял + # исключение — на этом нюансе строится фикс чекпоинта в pipeline.py: pipeline + # больше не берёт done_buckets из scraper.completed_buckets, только из факта + # успешного on_bucket (см. run_domclick_city_sweep._on_bucket/_checkpoint). + assert s.completed_buckets == ["st", "1"], s.completed_buckets + + +async def test_on_bucket_none_preserves_return_value() -> None: + """on_bucket=None → поведение прежнее, байт-в-байт: лоты возвращаются из + fetch_city как раньше, никаких промежуточных вызовов.""" + + async def _fake_sweep(*, rooms: str, out_lots: list[str], **_kw: Any) -> None: + out_lots.append(f"{rooms}-only") + + s = _scraper() + result = await _run_fetch_city(s, sweep_bucket=_fake_sweep, on_bucket=None) + + assert result == [f"{b}-only" for b in ROOM_BUCKETS] + assert s.completed_buckets == list(ROOM_BUCKETS) + + +# ── Секция 2: run_domclick_city_sweep(watchdog_sec=...) ────────────────────── + + +class _NoOpRuns: + """Достаточно методов, чтобы SERP-фаза дошла до asyncio.wait_for и честно + финализировалась после симулированного TimeoutError — значения не важны, + важен ТОЛЬКО timeout, с которым позвали wait_for.""" + + def is_cancelled(self, db: Any, run_id: int) -> bool: + return False + + def update_heartbeat(self, db: Any, run_id: int, counters: dict[str, Any]) -> None: + return None + + def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None: + return None + + def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None: + return None + + def mark_banned( + self, db: Any, run_id: int, error: str, counters: dict[str, Any], **kw: Any + ) -> None: + return None + + +async def _drive_and_capture_timeout(**kwargs: Any) -> float: + captured: dict[str, float] = {} + + async def _fake_wait_for(coro: Any, timeout: float) -> None: + captured["timeout"] = timeout + coro.close() + raise TimeoutError() + + with ( + patch(f"{PFX}.asyncio.wait_for", _fake_wait_for), + patch(f"{PFX}.runs", _NoOpRuns()), + ): + await run_domclick_city_sweep( + object(), # type: ignore[arg-type] + config=SimpleNamespace(browser_http_endpoint="http://x:9000"), + matcher=object(), + run_id=1, + city_id=4, + **kwargs, + ) + return captured["timeout"] + + +async def test_watchdog_sec_default_formula_matches_prod_11100() -> None: + """watchdog_sec=None (дефолт) → прежняя формула байт-в-байт. pages=100, + delay=6.0 — те же параметры, что дали 11 100с в run 7344 (Москва).""" + timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0) + assert timeout == 11100 + + +async def test_watchdog_sec_override_bypasses_formula() -> None: + """watchdog_sec задан явно → формула не считается вовсе, идёт ровно override + — даже с теми же pages/delay, что в тесте формулы выше.""" + timeout = await _drive_and_capture_timeout(pages=100, request_delay_sec=6.0, watchdog_sec=777) + assert timeout == 777 diff --git a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py index aaa1db08..7a4d112f 100644 --- a/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py +++ b/tradein-mvp/backend/tests/test_scraper_kit_pipeline_parity2.py @@ -411,8 +411,19 @@ async def _drive_domclick( recorder = _RunsRecorder() db = MagicMock() lots = [MagicMock() for _ in range(lots_n)] + + async def _fetch_city(**kw: Any) -> list[Any]: + # fix/domclick-incremental-save: save_listings переехал внутрь цикла по корзинам и + # зовётся из колбэка on_bucket. Стаб, который просто возвращает лоты, не + # вызвал бы сохранение вовсе — фикстура проверяла бы мёртвую ветку. + # Одна корзина со всеми лотами: ровно один save_listings, как и было. + on_bucket = kw.get("on_bucket") + if on_bucket is not None and lots: + on_bucket(ROOM_BUCKETS[0], lots) + return lots + scraper = _ctx_scraper( - fetch_city=AsyncMock(return_value=lots), + fetch_city=_fetch_city, blocked=blocked, geo_filtered=0, fetch_errors=fetch_errors, 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 0c69733b..2db4fd81 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 @@ -4868,6 +4868,7 @@ async def run_domclick_city_sweep( region_code: int = DEFAULT_REGION_CODE, resume_run_id: int | None = None, cookies: dict[str, str] | None = None, + watchdog_sec: int | None = None, ) -> DomClickCitySweepCounters: """DomClick citywide sweep через BFF JSON API. @@ -4896,6 +4897,15 @@ async def run_domclick_city_sweep( зависли на challenge). None (дефолт, сессии в БД нет/протухла) — прежнее поведение, без инъекции. + watchdog_sec (incremental-save): явный override расчётной формулы watchdog'а. Формула + ниже (buckets × pages × per_fetch + budget) не знает про бисекцию по цене + (ДомКлик режет offset на 2000, каждый лист пагинируется отдельно отдельным + деревом сплитов) и на больших городах (Москва ≈23 690 лотов против ЕКБ + ≈6 300) занижена в разы — прод run 7344: за 2ч watchdog'а (11 100с) пройдено + 2 бакета из 6. С инкрементальным сохранением (on_bucket, см. ниже) ранний + снос по watchdog больше не теряет собранное, поэтому вместо более точной + оценки — простой override: None (дефолт) — прежняя формула байт-в-байт. + Возвращает DomClickCitySweepCounters. """ # Гео-скоуп свипа больше не зашит в ЕКБ: его задаёт профиль региона (address @@ -4909,12 +4919,18 @@ async def run_domclick_city_sweep( _resolved_delay = request_delay_sec if request_delay_sec is not None else 6.0 counters = DomClickCitySweepCounters() - # Watchdog: 6 buckets × pages × per_fetch + budget. + # Watchdog: 6 buckets × pages × per_fetch + budget. watchdog_sec (incremental-save) — явный + # override для больших городов, где формула занижена (см. докстринг выше). _num_fetches = _DOMCLICK_NUM_BUCKETS * max(1, pages) - _sweep_timeout = max( - ANCHOR_TIMEOUT_SEC, - int(_num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S), - ) + if watchdog_sec is not None: + _sweep_timeout = int(watchdog_sec) + else: + _sweep_timeout = max( + ANCHOR_TIMEOUT_SEC, + int( + _num_fetches * (_resolved_delay + _DOMCLICK_PER_FETCH_S) + _DOMCLICK_SWEEP_BUDGET_S + ), + ) # Мутируемый контейнер для захвата scraper-ссылки из замыкания. _scraper_ref: list[DomClickScraper] = [] @@ -4991,14 +5007,68 @@ async def run_domclick_city_sweep( _sweep_timeout, ) - lots: list[ScrapedLot] = [] - # #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому - # корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже). - _saved = False + # incremental-save: сохранение инкрементальное — save_listings зовётся из _on_bucket ПОСЛЕ + # КАЖДОГО room-бакета, а не одним save_listings в самом конце. Раньше снятие фазы + # по watchdog'у (asyncio.wait_for TimeoutError) теряло ВСЁ собранное — на большой + # выдаче (Москва ≈23 690 лотов vs ЕКБ ≈6 300) формула watchdog'а систематически + # не укладывалась в отведённое время (прод run 7344: 2 бакета из 6 за 2ч). + _cancel_reason: str | None = None + + def _on_bucket(bucket_key: str, bucket_lots: list[ScrapedLot]) -> None: + """Инкрементальный save сразу после того, как room-бакет отфетчился. + + fetch_city (serp.py) зовёт колбэк ВНЕ try/except конкретного бакета — + исключение отсюда прерывает обход ROOM_BUCKETS целиком (канал кооперативной + отмены/SIGTERM-дрейна), а не проглатывается generic except'ом бакета. + Sentinel-приём (RuntimeError("cancelled")/("shutdown")) — тот же, что в + run_cian_full_load._on_bucket (#1182 Phase 3a). + """ + nonlocal _checkpoint + if runs.is_cancelled(db, run_id): + logger.info( + "domclick-sweep run_id=%d: cancel detected in on_bucket (%s)", + run_id, + bucket_key, + ) + raise RuntimeError("cancelled") + elif shutdown_requested(): + logger.info( + "domclick-sweep run_id=%d: SIGTERM-drain — stopping at bucket %s", + run_id, + bucket_key, + ) + raise RuntimeError("shutdown") + if bucket_lots: + # Имя города — из профиля региона, а не из сравнения с vestigial + # city_id: гео-скоп задаёт регион, он же знает, какой город штамповать. + # У области (50) city_name=None — одного города нет, угадывать нечего. + inserted, updated = save_listings( + db, + bucket_lots, + matcher=matcher, + region_code=region_code, + run_id=run_id, + city=_geo_profile.city_name, + ) + counters.lots_fetched += len(bucket_lots) + counters.lots_inserted += inserted + counters.lots_updated += updated + # #3118, теперь на уровне бакета: done_buckets означает "собрано И + # сохранено" — чекпоинт пишем ТОЛЬКО пройдя cancel/shutdown-гейт выше и + # save_listings этого бакета, не из scraper.completed_buckets (см. + # комментарий у "Перенести счётчики" ниже — там раньше был баг #2 issue). + _checkpoint = sorted(set(_checkpoint) | {bucket_key}) + runs.update_heartbeat(db, run_id, _payload()) + logger.info( + "domclick-sweep run_id=%d: bucket %s saved lots=%d total=%d", + run_id, + bucket_key, + len(bucket_lots), + counters.lots_fetched, + ) async def _domclick_phase() -> None: - """Единственная citywide-фаза: fetch_city + save.""" - nonlocal lots, _saved + """Единственная citywide-фаза: fetch_city с инкрементальным save по бакетам.""" async with DomClickScraper( config, proxy_provider=proxy_provider, @@ -5021,36 +5091,21 @@ async def run_domclick_city_sweep( # источники и растёт неравномерно, — но за 30 суток каждая корзина # получает порядка пяти стартов, чего достаточно для критерия приёмки # «объявления с rooms >= 2 появились». - lots = await _scraper.fetch_city( + await _scraper.fetch_city( city_id=city_id, rooms=rooms, pages=pages, start_bucket_index=run_id % len(ROOM_BUCKETS), skip_buckets=skip_buckets or None, + on_bucket=_on_bucket, ) - counters.lots_fetched += len(lots) - if lots: - # Имя города — из профиля региона, а не из сравнения с vestigial - # city_id: гео-скоп задаёт регион, он же знает, какой город штамповать. - # У области (50) city_name=None — одного города нет, угадывать нечего. - _dc_city = _geo_profile.city_name - inserted, updated = save_listings( - db, - lots, - matcher=matcher, - region_code=region_code, - run_id=run_id, - city=_dc_city, - ) - counters.lots_inserted += inserted - counters.lots_updated += updated - _saved = True try: await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout) except TimeoutError: logger.warning( - "domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results", + "domclick-sweep run_id=%d: SERP phase timed out after %ds — partial results " + "(incremental-save: уже сохранены инкрементально, бакет за бакетом)", run_id, _sweep_timeout, ) @@ -5080,6 +5135,16 @@ async def run_domclick_city_sweep( ban_kind=ban_kind_of_exception(exc), ) return counters + except RuntimeError as exc: + # on_bucket кидает RuntimeError("cancelled") при кооперативной отмене, + # RuntimeError("shutdown") при SIGTERM-дрейне — тот же sentinel-приём, что в + # run_cian_full_load._on_bucket (#1182 Phase 3a). Прочие RuntimeError — + # обычная поломка фазы, ведём себя как под generic except ниже. + if str(exc) in ("cancelled", "shutdown"): + _cancel_reason = str(exc) + else: + logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id) + counters.errors_count += 1 except Exception: logger.exception("domclick-sweep run_id=%d: SERP phase failed", run_id) counters.errors_count += 1 @@ -5095,15 +5160,13 @@ async def run_domclick_city_sweep( # скрейпер живая, а его счётчик показывает, докуда прогон дошёл. counters.buckets_completed = _s.buckets_completed counters.buckets_total = _s.buckets_total - # #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне. - # Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не - # затирают, и оборванный болезнью финализации прогон его не теряет. - # #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза - # не дошла до save_listings — у её корзин в БД ноль строк, а отметка - # «пройдена» заставила бы следующий прогон пропустить их навсегда - # (механизм разобран в миграции 308, из-за него выключены свипы 77/50). - if _saved: - _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) + # incremental-save: _checkpoint СЮДА больше не мержится из _s.completed_buckets — он + # пишется инкрементально внутри _on_bucket, СРАЗУ после save_listings этого + # бакета. _s.completed_buckets на уровне скрейпера отмечает бакет как "фетч + # прошёл" ДО вызова on_bucket (см. serp.py): если on_bucket поймал + # cancel/shutdown ДО save для этого самого бакета, _s.completed_buckets + # включил бы его, а _checkpoint — честно нет (лоты не сохранены). Мердж + # отсюда воспроизвёл бы старый баг — чекпоинт врёт про несохранённые бакеты. runs.update_heartbeat(db, run_id, _payload()) counters.bucket_start_index = _s.bucket_start_index @@ -5111,6 +5174,32 @@ async def run_domclick_city_sweep( counters.pages_fetched = _num_fetches runs.update_heartbeat(db, run_id, _payload()) + # incremental-save: кооперативная отмена/SIGTERM-дрейн, пойманные в on_bucket — партиал уже + # сохранён инкрементально, финализируем как run_cian_full_load (mark_done partial, + # не honest-status ниже: обрыв тут known-signal, а не "прогон не доделал сам"). + if _cancel_reason == "cancelled": + logger.info( + "domclick-sweep run_id=%d: cancelled — partial results lots=%d (ins=%d/upd=%d)", + run_id, + counters.lots_fetched, + counters.lots_inserted, + counters.lots_updated, + ) + runs.mark_done(db, run_id, _payload()) + return counters + if _cancel_reason == "shutdown": + logger.info( + "domclick-sweep run_id=%d: SIGTERM-drain — partial results lots=%d (ins=%d/upd=%d)", + run_id, + counters.lots_fetched, + counters.lots_inserted, + counters.lots_updated, + ) + _drain = {**_payload(), "interrupted": 1} + runs.update_heartbeat(db, run_id, _drain) + runs.mark_done(db, run_id, _drain) + return counters + # ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ─────────────────────────── # Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно # отличается от cian/yandex (#2625/#2642): там независимые anchor'ы, и 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 51bd6702..4cea925f 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 @@ -1270,6 +1270,9 @@ async def _job_domclick_city_sweep( request_delay_sec=float(params.get("request_delay_sec", 6.0)), region_code=_resolve_region_code(params), resume_run_id=_pick_resume(db, run_id), + watchdog_sec=( + int(params["watchdog_sec"]) if params.get("watchdog_sec") is not None else None + ), ) diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py index 5f8bbc97..9700761a 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/providers/domclick/serp.py @@ -459,6 +459,7 @@ class DomClickScraper(BaseScraper): pages: int = 100, start_bucket_index: int = 0, skip_buckets: set[str] | None = None, + on_bucket: Callable[[str, list[ScrapedLot]], None] | None = None, ) -> list[ScrapedLot]: """Citywide sweep через BFF JSON API. @@ -492,6 +493,17 @@ class DomClickScraper(BaseScraper): buckets_total при этом = числу корзин В ЭТОМ прогоне (без скипнутых) — иначе honest-status читал бы возобновлённый прогон как вечно-частичный. + on_bucket: колбэк инкрементального сохранения — большая выдача теряла + ВСЁ собранное при снятии по watchdog, единственный save был в самом + конце (fix/domclick-incremental-save). Зовётся СИНХРОННО сразу после того, как + бакет отработал успешно (после DomClickBlockedError/generic + Exception — НЕ зовётся), аргументы: имя бакета + лоты ИМЕННО + этого бакета (не накопленный out_lots). Вызов стоит ВНЕ + try/except этого бакета: исключение из колбэка (кооперативная + отмена/SIGTERM-дрейн — см. run_domclick_city_sweep) обязано + прервать цикл по ROOM_BUCKETS, а не быть проглоченным generic + except'ом. on_bucket=None (дефолт) — поведение прежнее, + байт-в-байт. Returns: Дедуплицированный по source_id список ScrapedLot. @@ -548,6 +560,7 @@ class DomClickScraper(BaseScraper): city_id, pages, ) + _bucket_start_len = len(out_lots) try: await self._sweep_bucket( fetcher=fetcher, @@ -591,6 +604,11 @@ class DomClickScraper(BaseScraper): continue self.buckets_completed += 1 self.completed_buckets.append(bucket) + # incremental-save: колбэк ВНЕ try/except этого бакета — исключение (канал + # кооперативной отмены/SIGTERM-дрейна в pipeline.run_domclick_city_sweep) + # обязано прервать цикл, а не попасть в generic except выше. + if on_bucket is not None: + on_bucket(bucket, out_lots[_bucket_start_len:]) logger.info( "domklik: fetch_city done city_id=%d total=%d buckets=%d/%d "