From 8e5f097c5eaa929b216f4993ac091c91f401718d Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 12:47:05 +0500 Subject: [PATCH 1/7] =?UTF-8?q?=D0=9F=D0=BB=D0=B0=D0=BD=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D1=89=D0=B8=D0=BA:=20=D1=83=D0=BF=D0=B0=D0=B2=D1=88?= =?UTF-8?q?=D0=B8=D0=B9=20=D0=B4=D0=BE=20=D1=84=D0=B8=D0=BD=D0=B0=D0=BB?= =?UTF-8?q?=D0=B8=D0=B7=D0=B0=D1=82=D0=BE=D1=80=D0=B0=20=D0=BF=D1=80=D0=BE?= =?UTF-8?q?=D0=B3=D0=BE=D0=BD=20=D0=BF=D0=BE=D0=BB=D1=83=D1=87=D0=B0=D0=B5?= =?UTF-8?q?=D1=82=20=D1=81=D1=82=D0=B0=D1=82=D1=83=D1=81=20=D1=81=D1=80?= =?UTF-8?q?=D0=B0=D0=B7=D1=83,=20=D0=B0=20=D0=BD=D0=B5=20zombie=20=D1=87?= =?UTF-8?q?=D0=B5=D1=80=D0=B5=D0=B7=206=20=D1=87=D0=B0=D1=81=D0=BE=D0=B2?= =?UTF-8?q?=20(#1940)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Исключение, вылетевшее из хендлера до его try с mark_failed (прод: пустой пул прокси в BrowserFetcher.__aenter__ через 10 мс после claim), планировщик только логировал. Строка оставалась 'running' с heartbeat_at == started_at, reaper через 6 ч ставил 'zombie' — за 21 сутки так 18 avito-свипов (7349 и 7350 висят сейчас). _dispatch._run теперь финализирует прогон сам: пустой пул прокси в цепочке причин -> banned/ban_kind=infra (как в пайплайнах), остальное -> failed. Уже финализированную хендлером строку не трогает (WHERE status='running'), пустые counters мержатся и чекпоинт не стирают. Co-Authored-By: Claude Opus 5 --- .../test_1940_crashed_run_is_finalized.py | 133 ++++++++++++++++++ .../scraper_kit/orchestration/scheduler.py | 22 ++- 2 files changed, 154 insertions(+), 1 deletion(-) create mode 100644 tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py diff --git a/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py b/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py new file mode 100644 index 00000000..447f6eb1 --- /dev/null +++ b/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py @@ -0,0 +1,133 @@ +"""Упавший хендлер не оставляет строку scrape_runs в 'running' до zombie-reaper'а (#1940). + +Прод 17.09 06:58: avito_city_sweep run 7349 — через 10 мс после claim BrowserFetcher.__aenter__ +бросил NoProxyAvailableError, планировщик записал «crashed run_id=7349» и всё. Строка +осталась 'running' с heartbeat_at == started_at; 21 сутки до этого так ушли в 'zombie' +18 avito-свипов, каждый через 6 ч. Проверяется значение статуса в строке прогона. +""" + +from __future__ import annotations + +import asyncio +import os +from typing import Any +from unittest.mock import MagicMock + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from scraper_kit.orchestration.runs import BAN_KIND_INFRA +from scraper_kit.orchestration.scheduler import Handler, SchedulerContext, _dispatch +from scraper_kit.proxy_errors import NoProxyAvailableError + + +class _Db: + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: + res = MagicMock() + res.fetchone.return_value = None # «не running» для гейта _claim_run + res.scalar.return_value = True # advisory-lock взят + return res + + def commit(self) -> None: + pass + + def rollback(self) -> None: + pass + + def close(self) -> None: + pass + + +class _Runs: + """ctx.runs с боевым гейтом `WHERE status = 'running'` и мержем counters.""" + + def __init__(self) -> None: + self.rows: dict[int, dict[str, Any]] = {} + + def create_run(self, db: Any, *, source: str, params: dict[str, Any]) -> int: + run_id = 7349 + len(self.rows) + self.rows[run_id] = {"status": "running", "counters": {}, "ban_kind": None} + return run_id + + def _finish(self, run_id: int, status: str, counters: dict[str, Any], **extra: Any) -> None: + row = self.rows[run_id] + if row["status"] != "running": + return + row.update(status=status, **extra) + row["counters"] = {**row["counters"], **counters} + + def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None: + self._finish(run_id, "done", counters) + + def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None: + self._finish(run_id, "failed", counters, error=error) + + def mark_banned( + self, db: Any, run_id: int, error: str, counters: dict[str, Any], *, ban_kind: str + ) -> None: + self._finish(run_id, "banned", counters, error=error, ban_kind=ban_kind) + + +async def _dispatch_and_wait(job: Any) -> dict[str, Any]: + runs = _Runs() + ctx = SchedulerContext( + config=MagicMock(), + matcher=MagicMock(), + enrichment=MagicMock(), + session_factory=_Db, + runs=runs, + ) + sched = { + "id": 1, + "source": "avito_city_sweep", + "window_start_hour": 6, + "window_end_hour": 7, + "default_params": {}, + } + run_id = await _dispatch(Handler(job, "avito_city_sweep"), _Db(), sched, ctx) + assert run_id is not None + await asyncio.gather(*ctx._inflight_tasks, return_exceptions=True) + return runs.rows[run_id] + + +async def test_no_proxy_crash_before_try_marks_run_banned_infra() -> None: + async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: + raise NoProxyAvailableError("avito") + + row = await _dispatch_and_wait(_job) + + assert row["status"] == "banned" + assert row["ban_kind"] == BAN_KIND_INFRA + assert "no proxy available" in row["error"] + + +async def test_wrapped_no_proxy_is_still_infra() -> None: + async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: + try: + raise NoProxyAvailableError("cian") + except NoProxyAvailableError as exc: + raise RuntimeError("sidecar fetch failed") from exc + + row = await _dispatch_and_wait(_job) + + assert (row["status"], row["ban_kind"]) == ("banned", BAN_KIND_INFRA) + + +async def test_other_crash_marks_run_failed() -> None: + async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: + raise KeyError("anchors") + + row = await _dispatch_and_wait(_job) + + assert row["status"] == "failed" + assert "KeyError" in row["error"] + + +async def test_handler_own_final_status_and_checkpoint_survive() -> None: + async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: + c.runs.mark_done(db, run_id, {"done_buckets": ["center"]}) + raise RuntimeError("after mark_done") + + row = await _dispatch_and_wait(_job) + + assert row["status"] == "done" + assert row["counters"] == {"done_buckets": ["center"]} 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 7923efbe..47932091 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 @@ -57,6 +57,8 @@ from scraper_kit.orchestration.pipeline import ( run_domclick_city_sweep, run_yandex_city_sweep, ) +from scraper_kit.orchestration.runs import BAN_KIND_INFRA +from scraper_kit.proxy_errors import caused_by_no_proxy if TYPE_CHECKING: from sqlalchemy.orm import Session @@ -1024,8 +1026,26 @@ async def _dispatch( run_db = ctx.session_factory() try: await handler.job(run_db, run_id, params, ctx) - except Exception: + except Exception as exc: logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id) + # #1940: исключение, вылетевшее ДО try-финализатора хендлера (прод: пустой + # пул прокси в BrowserFetcher.__aenter__ через 10 мс после claim), оставляло + # строку 'running' без единого пульса — через 6 ч её снимал reaper как + # 'zombie', источник терял день. Финализируем здесь, в единственной точке, + # через которую идёт любой хендлер. Уже финализированную хендлером строку + # не трогаем: mark_* пишут только `WHERE status = 'running'`, а {} counters + # мержится и чекпоинт не стирает. + try: + if caused_by_no_proxy(exc): + ctx.runs.mark_banned( + run_db, run_id, f"crashed: {exc}", {}, ban_kind=BAN_KIND_INFRA + ) + else: + ctx.runs.mark_failed( + run_db, run_id, f"crashed: {type(exc).__name__}: {exc}", {} + ) + except Exception: + logger.exception("scheduler: не удалось финализировать run_id=%d", run_id) finally: run_db.close() -- 2.45.3 From facdbb256626c1a308558b82bbd3658cc7884de4 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 12:48:20 +0500 Subject: [PATCH 2/7] =?UTF-8?q?=D0=A1=D0=B2=D0=B8=D0=BF=D1=8B=20=D0=A6?= =?UTF-8?q?=D0=B8=D0=B0=D0=BD=D0=B0=20=D0=B8=20=D0=90=D0=B2=D0=B8=D1=82?= =?UTF-8?q?=D0=BE:=20=D1=8F=D0=BA=D0=BE=D1=80=D1=8C=20=D1=81=20=D0=BF?= =?UTF-8?q?=D0=BE=D0=BB=D0=BD=D0=BE=D1=81=D1=82=D1=8C=D1=8E=20=D0=BE=D1=82?= =?UTF-8?q?=D0=BA=D0=B0=D0=B7=D0=B0=D0=B2=D1=88=D0=B5=D0=B9=20=D1=84=D0=B0?= =?UTF-8?q?=D0=B7=D0=BE=D0=B9=20=D0=BD=D0=B5=20=D0=BF=D0=BE=D0=BF=D0=B0?= =?UTF-8?q?=D0=B4=D0=B0=D0=B5=D1=82=20=D0=B2=20=D1=87=D0=B5=D0=BA=D0=BF?= =?UTF-8?q?=D0=BE=D0=B8=D0=BD=D1=82=20(#3415)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit cian_city_sweep 6179 (06.09) — failed по «фаза houses отказала полностью, 30 из 30», но в done_buckets все пять якорей. 6276 резюмировался от него, пропустил все якоря и отчитался done с houses_attempted=0, lots_fetched=0. Отметка «якорь пройден» опиралась только на поток управления, а фаза отказывает и без исключения (fetch_newbuilding -> None, растёт счётчик). Гейт _bucket_phase_totally_failed сверяет прирост пар X_attempted/X_failed за якорь; порог одна попытка. Тот же гейт в run_avito_city_sweep (detail). Ветка «houses DB query failed» теперь растит houses_attempted вместе с houses_failed. Форма чекпоинта прежняя (плоский список имён). Перенос ветки fix/3415-bucket-done-only-if-phases-ok (19a335f4) на свежий main, применилась без конфликтов. Co-Authored-By: Claude Opus 5 --- ...test_3415_bucket_done_only_if_phases_ok.py | 271 ++++++++++++++++++ .../src/scraper_kit/orchestration/pipeline.py | 94 +++++- 2 files changed, 363 insertions(+), 2 deletions(-) create mode 100644 tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py diff --git a/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py b/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py new file mode 100644 index 00000000..f74aaaa9 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py @@ -0,0 +1,271 @@ +"""Бакет уходит в чекпоинт только если ВСЕ его фазы что-то дали (#3415). + +Прод-факт. `cian_city_sweep` 6179 (06.09) — статус `failed` с причиной +«phase-honest-status: фаза 'houses' отказала полностью — 30 из 30 попыток +неудачны, обогащено 0», а `counters.done_buckets` при этом несёт ВСЕ пять якорей +(«Академический», «Пионерский», «Уралмаш», «Центр», «ЮЗ»). Следующий прогон 6276 +(07.09) резюмировался от него (`resume_from: 6179`, `resume_buckets: 5`), +пропустил все якоря и отчитался `done` с `errors_count=0`, `houses_attempted=0`, +`lots_fetched=0` — зелёный по построению: фаза не исполнялась. Из-за этого фикс +#3396 (houses через пул прокси) простоял на проде с 06.09 и ни разу не отработал. + +ИНВАРИАНТ. Отказ ЦЕЛОЙ фазы внутри якоря — это «якорь пройден не был», ровно как +исключение или таймаут (#3074/#3319). Разница только в том, что отказ фазы виден +не потоком управления, а бухгалтерией счётчиков: `houses_attempted` попыток, +`houses_failed` отказов, обогащено ноль. Порог здесь — одна попытка, а не три, +как у run-level `_phase_totally_failed` (#2700): тот СТАВИТ ДИАГНОЗ прогону и +обязан отсеивать шум, а этот решает «собирать ли якорь заново», и цена ошибок +несимметрична — лишний повтор якоря стоит нескольких запросов, а ложная отметка +«пройден» теряет часть города навсегда и молча. +""" + +from __future__ import annotations + +import os + +# Settings собирается автофикстурой conftest'а и требует database_url. +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +import types +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest + +ANCHOR_A = (56.83, 60.60, "ekb-center") +ANCHOR_B = (56.79, 60.63, "ekb-south") + +# Два дома на якорь — НИЖЕ порога run-level правила (_PHASE_MIN_ATTEMPTS = 3): +# гейт бакета обязан сработать на своей, более строгой границе. +HOUSE_ROWS = [ + {"id": 101, "cian_zhk_url": "https://cian.ru/zhk/101"}, + {"id": 102, "cian_zhk_url": "https://cian.ru/zhk/102"}, +] +DETAIL_ROWS = [ + {"id": 201, "source_url": "/ekb/kvartiry/201"}, + {"id": 202, "source_url": "/ekb/kvartiry/202"}, +] + + +class _FakeDb: + def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: + self.prev_counters = prev_counters or {} + self.heartbeats: list[dict[str, Any]] = [] + + def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: + if params and "counters" in params: + self.heartbeats.append(json.loads(params["counters"])) + return MagicMock() + if params and "rid" in params: + return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters)) + if params and "ids" in params: + # houses-фаза: дома, у которых есть cian_zhk_url. + return MagicMock(mappings=lambda: MagicMock(all=lambda: HOUSE_ROWS)) + if params and "limit" in params: + # detail-фаза avito: карточки-кандидаты на обогащение. + return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS)) + return MagicMock() + + def commit(self) -> None: ... + def rollback(self) -> None: ... + + +def _lot() -> types.SimpleNamespace: + """Новостроечный лот с привязкой к ЖК — только такой доходит до houses-фазы.""" + return types.SimpleNamespace( + listing_segment="novostroyki", + house_source="cian_newbuilding", + house_ext_id="nb-1", + ) + + +class _FakeScraper: + """Двойник CianScraper: помнит, за какими якорями реально ходили.""" + + visited: list[str] = [] # noqa: RUF012 — тестовый сборник + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self.state_extraction_attempts = 1 + self.state_extraction_failures = 0 + self.request_delay_sec = 0.0 + + async def __aenter__(self) -> _FakeScraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + async def fetch_around_multi_room( + self, lat: float, lon: float, *_a: Any, **_kw: Any + ) -> list[types.SimpleNamespace]: + _FakeScraper.visited.append(f"{lat},{lon}") + return [_lot()] + + +def _config() -> types.SimpleNamespace: + return types.SimpleNamespace( + scraper_proxy_url=None, + scraper_fetch_mode="cffi", + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + ) + + +async def _run( + prev: dict[str, Any] | None, *, houses_ok: bool, run_id: int = 8401 +) -> tuple[_FakeDb, Any]: + from scraper_kit.orchestration import pipeline as pl + + _FakeScraper.visited = [] + db = _FakeDb(prev) + + async def _fake_newbuilding(*_a: Any, **_kw: Any) -> Any: + # None — ровно тот отказ, что был на проде 06.09: исключения нет, + # счётчик houses_failed растёт, фаза возвращается штатно. + return object() if houses_ok else None + + with ( + patch.object(pl, "CianScraper", _FakeScraper), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl, "fetch_newbuilding", _fake_newbuilding), + patch.object(pl, "save_newbuilding_enrichment", lambda *_a, **_kw: None), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + counters = await pl.run_cian_city_sweep( + db, # type: ignore[arg-type] + run_id=run_id, + config=_config(), + matcher=MagicMock(), + anchors=[ANCHOR_A, ANCHOR_B], + enrich_houses=True, + detail_top_n=0, + request_delay_sec=0.0, + resume_run_id=(run_id - 1) if prev is not None else None, + ) + return db, counters + + +def _last_checkpoint(db: _FakeDb) -> list[str]: + with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb] + assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится" + return with_ckpt[-1]["done_buckets"] + + +@pytest.mark.asyncio +async def test_anchor_with_totally_failed_houses_is_not_checkpointed() -> None: + """(а) Якорь, где houses отказала на ВСЕХ домах, не попадает в чекпоинт.""" + db, counters = await _run(None, houses_ok=False) + assert counters.houses_attempted > 0, "houses-фаза вообще не исполнялась — тест не о том" + assert counters.houses_enriched == 0, "houses-фаза что-то обогатила — отказ не полный" + assert _last_checkpoint(db) == [], ( + "якорь с полностью отказавшей houses-фазой помечен пройденным — " + f"чекпоинт {_last_checkpoint(db)}" + ) + + +@pytest.mark.asyncio +async def test_next_run_re_enters_houses() -> None: + """(б) Следующий прогон с этим чекпоинтом снова заходит в houses.""" + db1, _ = await _run(None, houses_ok=False) + _db2, counters2 = await _run( + {"done_buckets": _last_checkpoint(db1)}, houses_ok=True, run_id=8402 + ) + assert counters2.houses_attempted > 0, ( + "прогон-наследник пропустил houses-фазу: чекпоинт предшественника " + "объявил якоря пройденными, хотя фаза в них не дала ничего" + ) + assert counters2.houses_enriched > 0, "houses-фаза исполнилась, но ничего не обогатила" + + +@pytest.mark.asyncio +async def test_healthy_anchor_is_still_checkpointed() -> None: + """(в) Контроль: якорь, где все фазы отработали, помечается как раньше.""" + db, counters = await _run(None, houses_ok=True) + assert counters.houses_enriched > 0 + assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], ( + "перестали помечать пройденные якоря вовсе — резюм (#3074) сломан" + ) + + +class _FakeAvitoScraper: + """Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза.""" + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + + async def fetch_around(self, *_a: Any, **_kw: Any) -> list[types.SimpleNamespace]: + return [types.SimpleNamespace(house_url=None)] + + +class _FakeAsyncSession: + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + + async def __aenter__(self) -> _FakeAsyncSession: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + +@pytest.mark.asyncio +async def test_avito_anchor_with_totally_failed_detail_is_not_checkpointed() -> None: + """Сосед: у avito-свипа тот же гейт — фаза без единой удачи не даёт отметки. + + Кода общего у свипов нет (две отдельные функции), общий только гейт. Прод-следа + у avito за 45 суток нет (0 прогонов против 4 у циана) — тест сторожит вторую + копию проводки, а не найденный в данных дефект. + """ + from scraper_kit.orchestration import pipeline as pl + + db = _FakeDb() + + async def _boom(*_a: Any, **_kw: Any) -> Any: + raise RuntimeError("detail отдал 403") + + with ( + patch.object(pl, "AvitoScraper", _FakeAvitoScraper), + patch.object(pl, "AsyncSession", _FakeAsyncSession), + patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)), + patch.object(pl, "fetch_detail", _boom), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + counters = await pl.run_avito_city_sweep( + db, # type: ignore[arg-type] + run_id=8404, + config=types.SimpleNamespace( + scraper_fetch_mode="cffi", + scraper_proxy_url=None, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + ), + matcher=MagicMock(), + enrichment=MagicMock(), + anchors=[ANCHOR_A], + enrich_houses=False, + enrich_imv=False, + detail_top_n=len(DETAIL_ROWS), + request_delay_sec=0.0, + ) + + assert counters.detail_attempted == len(DETAIL_ROWS) + assert counters.detail_enriched == 0 + assert _last_checkpoint(db) == [], ( + f"якорь с полностью отказавшей detail-фазой помечен пройденным: {_last_checkpoint(db)}" + ) + + +@pytest.mark.asyncio +async def test_old_format_checkpoint_still_resumes() -> None: + """(г) Чекпоинт СТАРОГО формата (плоский список имён якорей) читается как раньше.""" + db, _ = await _run({"done_buckets": ["ekb-center"]}, houses_ok=True, run_id=8403) + assert f"{ANCHOR_A[0]},{ANCHOR_A[1]}" not in _FakeScraper.visited, ( + "якорь из унаследованного чекпоинта всё-таки опрашивали" + ) + assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], ( + "унаследованный якорь потерян — формат чекпоинта разъехался со старыми прогонами" + ) 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 f9ad6e60..4fae93a3 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 @@ -36,7 +36,7 @@ from __future__ import annotations import asyncio import logging import random -from collections.abc import Callable +from collections.abc import Callable, Mapping from contextlib import AsyncExitStack from dataclasses import dataclass, field, fields from datetime import date, timedelta @@ -1346,6 +1346,43 @@ async def run_avito_pipeline( await browser_fetcher.__aexit__(None, None, None) +def _bucket_phase_totally_failed(before: Mapping[str, int], after: Mapping[str, int]) -> str | None: + """Фаза бакета, где отказала КАЖДАЯ попытка (#3415). Имя фазы или None. + + Считает ПРИРОСТ счётчиков за один бакет (якорь): counters у свипа run-level, + так что «фаза этого якоря» существует только как разница до/после. + + Прод-факт, ради которого гейт заведён: `cian_city_sweep` 6179 (06.09) — + `houses_attempted=30, houses_failed=30, houses_enriched=0`, статус `failed` по + run-level правилу (`_phase_totally_failed`, #2700), а в `done_buckets` записаны + все пять якорей. Прогон 6276 (07.09) унаследовал точку, пропустил все якоря и + отчитался `done` с `houses_attempted=0` — зелёный по построению, потому что фаза + не исполнялась. Фикс #3396 (houses через пул прокси) простоял на проде сутки, ни + разу не отработав. + + Пары ищутся В САМИХ counters (ключ `X_attempted` со спутником `X_failed`) — тот же + приём и та же причина, что у `runs._phase_totally_failed`: зашитый список фаз это + ровно то место, куда забывают дописать новую. Фазы без счётчика попыток гейт НЕ + видит — у avito houses роль attempted играет `unique_houses` (нет `houses_attempted`), + и притягивать её сюда значило бы блокировать отметку якоря из-за ОДНОГО упавшего + дома из десяти. + + Порог — одна попытка, а не три, как у run-level правила: тот ставит прогону + диагноз и обязан отсеивать шум, а этот решает «собирать ли якорь заново». Цена + ошибок несимметрична: лишний повтор якоря стоит нескольких запросов, ложная + отметка «пройден» теряет часть города навсегда и молча (прогон-то `done`). + """ + for key in sorted(after): + if not key.endswith("_attempted"): + continue + phase = key[: -len("_attempted")] + attempted = after[key] - before.get(key, 0) + failed = after.get(f"{phase}_failed", 0) - before.get(f"{phase}_failed", 0) + if attempted >= 1 and failed >= attempted: + return phase + return None + + @dataclass class CitySweepCounters: """Aggregate counters для full city sweep run.""" @@ -1629,6 +1666,12 @@ async def run_avito_city_sweep( # #3074: сбрасывается на КАЖДОЙ итерации — иначе один упавший якорь # заразил бы все последующие, и чекпоинт не пополнялся бы вовсе. _anchor_ok = True + # #3415: снимок счётчиков ДО якоря — по разнице видно, дала ли + # каждая фаза ЭТОГО якоря хоть одну удачную попытку (см. + # _bucket_phase_totally_failed). Тот же дефект, что у cian-свипа: + # detail-фаза, отказавшая на всех карточках без исключения, доходила + # до отметки «якорь пройден» наравне с целым якорем. + _before_anchor = counters.to_dict() # Capture loop variables in default args (B023): prevents stale binding # if the coroutine is scheduled after the loop variable changes. @@ -2137,6 +2180,22 @@ async def run_avito_city_sweep( # его как пройденный значило бы, что следующий прогон его пропустит и # объявления оттуда не соберутся НИКОГДА, причём молча: прогон # завершится штатно. Тот же инвариант, что у combo в yandex-свипе. + # + # #3415: флага мало — фаза отказывает и БЕЗ исключения (каждая + # detail-карточка отдала 403, счётчик вырос, корутина завершилась + # штатно). Замер на проде за 45 суток: у avito-свипов такой якорь + # пока не встречался (0 прогонов), у cian — 4; гейт стоит здесь + # потому, что дефект один и тот же, а не по следу в данных. + _failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict()) + if _failed_phase is not None: + logger.warning( + "city-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' " + "не дала ни одной удачной попытки; следующий прогон соберёт его заново", + run_id, + name, + _failed_phase, + ) + _anchor_ok = False if _anchor_ok: _done_anchors.add(name) runs.update_heartbeat(db, run_id, _ckpt()) @@ -3313,6 +3372,10 @@ async def run_cian_city_sweep( ) # ── Per-anchor phases под watchdog-таймаутом ──────────────────── + # #3415: counters у свипа run-level — снимок ДО якоря даёт бухгалтерию + # его собственных фаз (разницей), по которой ниже решается, считать ли + # якорь пройденным. + _before_anchor = counters.to_dict() anchor_lots: list[ScrapedLot] = [] _c_lat, _c_lon, _c_name = lat, lon, name @@ -3529,6 +3592,12 @@ async def run_cian_city_sweep( len(nb_id_list), exc, ) + # #3415: попытки считаем ВМЕСТЕ с отказами. Раньше здесь рос + # только `houses_failed`, и пара получалась несуществующей + # (attempted=0, failed=N): и run-level `_phase_totally_failed` + # (#2700), и гейт бакета сверяют `failed` с `attempted`, так что + # полный отказ фазы на этой ветке не видел ни один из них. + counters.houses_attempted += len(nb_id_list) counters.houses_failed += len(nb_id_list) try: db.rollback() @@ -3720,7 +3789,28 @@ async def run_cian_city_sweep( # Записать упавший якорь пройденным значило бы, что следующий прогон # пропустит его навсегда — молча, потому что прогон завершится # штатно, просто часть города не соберётся. - _done_anchors.add(name) + # + # #3415: потока управления мало. Фаза внутри якоря отказывает БЕЗ + # исключения — `fetch_newbuilding` возвращает None, счётчик растёт, + # `_cian_anchor_phases` завершается штатно, — и такой якорь доходил + # сюда наравне с целым. Дальше ровно тот же исход, что и у записанного + # упавшего: следующий прогон пропускает якорь и рапортует `done` с + # нулевой фазой. Форма отметки прежняя (плоский список имён): чекпоинт + # читают ещё три свипа и `scheduler._resume_decision`, а ключ + # `bucket:phase` потребовал бы новых читателей ради того же решения. + _failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict()) + if _failed_phase is None: + _done_anchors.add(name) + else: + logger.warning( + "cian-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' " + "не дала ни одной удачной попытки; следующий прогон соберёт его заново", + run_id, + name, + _failed_phase, + ) + # Heartbeat пишем в любом случае: без него reap_zombies посчитает живой + # прогон мёртвым, а jsonb-мерж сохранит унаследованную часть точки. runs.update_heartbeat( db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)} ) -- 2.45.3 From fb02e41d77c1727b55ca9e21f5274c1b83742507 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 12:54:04 +0500 Subject: [PATCH 3/7] =?UTF-8?q?=D0=9F=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA?= =?UTF-8?q?=D0=B0=20=D0=BE=D1=82=D0=BC=D0=B5=D0=BD=D1=8B=20=D0=B7=D0=B0?= =?UTF-8?q?=D0=BA=D1=80=D1=8B=D0=B2=D0=B0=D0=B5=D1=82=20=D1=81=D0=B2=D0=BE?= =?UTF-8?q?=D1=8E=20=D1=82=D1=80=D0=B0=D0=BD=D0=B7=D0=B0=D0=BA=D1=86=D0=B8?= =?UTF-8?q?=D1=8E:=20=D1=81=D0=B2=D0=B8=D0=BF=20=D0=94=D0=BE=D0=BC=D0=9A?= =?UTF-8?q?=D0=BB=D0=B8=D0=BA=D0=B0=20=D0=B1=D0=BE=D0=BB=D1=8C=D1=88=D0=B5?= =?UTF-8?q?=20=D0=BD=D0=B5=20=D0=B4=D0=B5=D1=80=D0=B6=D0=B8=D1=82=20=D0=B3?= =?UTF-8?q?=D0=BE=D1=80=D0=B8=D0=B7=D0=BE=D0=BD=D1=82=20vacuum=20=D0=BD?= =?UTF-8?q?=D0=B0=20=D0=B2=D0=B5=D1=81=D1=8C=20=D0=BE=D0=B1=D1=85=D0=BE?= =?UTF-8?q?=D0=B4=20(#3480)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Виновник по живому снимку pg_stat_activity 17.09 07:48 UTC: pid 1933154 из tradein-scraper, 'idle in transaction' 1 ч 35 мин, последний запрос `SELECT status FROM scrape_runs WHERE id = $1` (runs.is_cancelled), xact_start через 30 мс после старта domclick_city_sweep_moskva 7344. По Prometheus за 03-17.09 10 из 12 окон «транзакция > 1 ч» совпадают по началу и концу со свипами ДомКлика. is_cancelled делал SELECT без commit, а вызывающие сразу уходят в сетевую фазу (у ДомКлика — весь fetch_city, до 3 ч). Теперь commit после чтения — в общей точке для всех 12 вызовов; незакоммиченных записей, рассчитанных на откат, перед вызовами нет (перед каждым — update_heartbeat/save_listings с commit или чтения). Co-Authored-By: Claude Opus 5 --- ...test_3480_no_open_transaction_over_http.py | 90 +++++++++++++++++++ .../src/scraper_kit/orchestration/runs.py | 10 ++- 2 files changed, 99 insertions(+), 1 deletion(-) create mode 100644 tradein-mvp/backend/tests/test_3480_no_open_transaction_over_http.py diff --git a/tradein-mvp/backend/tests/test_3480_no_open_transaction_over_http.py b/tradein-mvp/backend/tests/test_3480_no_open_transaction_over_http.py new file mode 100644 index 00000000..8d1e8011 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3480_no_open_transaction_over_http.py @@ -0,0 +1,90 @@ +"""Проверка отмены не оставляет транзакцию открытой на сетевую фазу свипа (#3480). + +Прод 17.09 07:48 UTC, pg_stat_activity БД tradein: pid 1933154 из tradein-scraper, +'idle in transaction' 1 ч 35 мин, последний запрос `SELECT status FROM scrape_runs +WHERE id = $1` (runs.is_cancelled), xact_start через 30 мс после старта +domclick_city_sweep_moskva 7344. За 14 суток 10 из 12 окон «транзакция > 1 ч» +совпали со свипами ДомКлика. + +Сессия настоящая (SQLAlchemy + SQLite): значение `in_transaction()` — то же, что видит +Postgres как 'idle in transaction'. +""" + +from __future__ import annotations + +import os +from types import SimpleNamespace +from typing import Any +from unittest.mock import MagicMock, patch + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +from scraper_kit.orchestration import pipeline as pl +from scraper_kit.orchestration import runs +from sqlalchemy import create_engine, text +from sqlalchemy.orm import Session + + +def _session() -> Session: + db = Session(create_engine("sqlite://")) + db.execute(text("CREATE TABLE scrape_runs (id INTEGER, status TEXT, counters TEXT)")) + db.execute( + text( + "INSERT INTO scrape_runs VALUES " + "(7333, 'done', NULL), (7344, 'running', NULL), (7345, 'cancelled', NULL)" + ) + ) + db.commit() + return db + + +def test_is_cancelled_closes_its_transaction() -> None: + db = _session() + + assert runs.is_cancelled(db, 7344) is False + assert db.in_transaction() is False + assert runs.is_cancelled(db, 7345) is True + assert db.in_transaction() is False + + +async def test_domclick_sweep_http_phase_runs_outside_transaction() -> None: + db = _session() + seen: list[bool] = [] + + class _Scraper: + blocked = False + geo_filtered = fetch_errors = buckets_completed = buckets_total = 0 + bucket_start_index = 0 + completed_buckets: list[str] = [] # noqa: RUF012 + + def __init__(self, *_a: Any, **_kw: Any) -> None: ... + + async def __aenter__(self) -> _Scraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + async def fetch_city(self, **_kw: Any) -> list[Any]: + seen.append(db.in_transaction()) # момент сетевого обхода + return [] + + with ( + patch.object(pl, "DomClickScraper", _Scraper), + patch.object(runs, "update_heartbeat", MagicMock()), + patch.object(runs, "mark_done", MagicMock()), + patch.object(runs, "mark_failed", MagicMock()), + patch.object(runs, "mark_banned", MagicMock()), + ): + await pl.run_domclick_city_sweep( + db, + run_id=7344, + config=SimpleNamespace(browser_http_endpoint="http://x:9000"), + matcher=MagicMock(), + city_id=4, + pages=1, + request_delay_sec=0.0, + resume_run_id=7333, # резюм-SELECT тоже открывает транзакцию + ) + + assert seen == [False] 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 4b38558a..b2d32d91 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 @@ -815,11 +815,19 @@ def honors_cancel(source: str) -> bool: def is_cancelled(db: Session, run_id: int) -> bool: - """Проверить status='cancelled' (cooperative cancel в long-running pipeline).""" + """Проверить status='cancelled' (cooperative cancel в long-running pipeline). + + #3480: SELECT открывает транзакцию (autobegin), а вызывающие сразу уходят в сетевую + фазу. Без commit сессия висела 'idle in transaction' весь обход и держала горизонт + vacuum всей БД: прод 17.09, domclick_city_sweep_moskva 7344 — транзакция 1.6 ч с + последним запросом ровно этим SELECT; 10 из 12 окон > 1 ч за 14 суток — свипы + ДомКлика. Закрываем здесь: проверку отмены зовут перед каждым долгим await. + """ row = db.execute( text("SELECT status FROM scrape_runs WHERE id = :id"), {"id": run_id}, ).fetchone() + db.commit() return row is not None and row.status == "cancelled" -- 2.45.3 From d2a31dff0807175a923621577403649c02189979 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 12:57:13 +0500 Subject: [PATCH 4/7] =?UTF-8?q?=D0=A2=D0=B5=D1=81=D1=82=D1=8B=20kit-=D1=81?= =?UTF-8?q?=D0=B2=D0=B8=D0=BF=D0=BE=D0=B2:=20=D1=82=D0=B0=D0=B9=D0=BC?= =?UTF-8?q?=D0=B0=D1=83=D1=82=20=D1=8F=D0=BA=D0=BE=D1=80=D1=8F,=20=D0=B4?= =?UTF-8?q?=D1=80=D0=B5=D0=B9=D0=BD=20=D0=B8=20=D0=BE=D1=82=D0=BC=D0=B5?= =?UTF-8?q?=D0=BD=D0=B0=20=D0=BF=D0=BE=D1=81=D1=80=D0=B5=D0=B4=D0=B8=20?= =?UTF-8?q?=D0=BE=D0=B1=D1=85=D0=BE=D0=B4=D0=B0=20=D1=83=20=D0=B2=D1=81?= =?UTF-8?q?=D0=B5=D1=85=20=D1=81=D0=B2=D0=B8=D0=BF=D0=BE=D0=B2,=20=D0=B0?= =?UTF-8?q?=20=D0=BD=D0=B5=20=D1=82=D0=BE=D0=BB=D1=8C=D0=BA=D0=BE=20=D1=83?= =?UTF-8?q?=20=D0=90=D0=B2=D0=B8=D1=82=D0=BE=20(#2406)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit После удаления legacy-пайплайна эти сценарии были заассерчены только у avito (test_3319) и дрейн до первого якоря у yandex/cian (test_3333). Добавлено: - таймаут якоря у cian и yandex: якорь не в чекпоинте, следующий обработан, errors_count растёт; таймаут фазы у domclick: завершённые корзины в чекпоинте, статус failed, а не done; - SIGTERM-дрейн после первого якоря у cian и yandex: interrupted=1, в чекпоинте ровно первый якорь, второй не посещён; - пользовательская отмена (is_cancelled=True) после первого якоря у avito, cian и yandex: второй якорь не посещён, чекпоинт на месте, без mark_done. Проверка по строке прогона: фейковая БД мержит counters всех записей, как jsonb-мерж в runs.py. Мутации (снятый continue/return/interrupted/errors_count в pipeline.py) роняют ровно свой тест, 6 из 6. Co-Authored-By: Claude Opus 5 --- .../tests/test_2406_sweep_edge_cases.py | 224 ++++++++++++++++++ 1 file changed, 224 insertions(+) create mode 100644 tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py diff --git a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py new file mode 100644 index 00000000..8b15a43d --- /dev/null +++ b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py @@ -0,0 +1,224 @@ +"""Крайние случаи city-свипов kit'а, не покрытые после удаления legacy-пайплайна (#2406). + +Уже покрыто: avito — таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО +первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369). +Здесь остаток: + (а) таймаут якоря у cian и yandex, таймаут фазы у domclick; + (б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex; + (в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito, + cian и yandex. + +Проверяется строка прогона, а не вызовы: _FakeDb мержит counters всех записей по +порядку (как `counters || :counters` в runs.py) и запоминает финальный статус. +""" + +from __future__ import annotations + +import os + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test") + +import json +import types +from typing import Any +from unittest.mock import MagicMock, patch + +import pytest +from scraper_kit.orchestration import pipeline as pl + +ANCHOR_A = (56.83, 60.60, "ekb-center") +ANCHOR_B = (56.79, 60.63, "ekb-south") + + +class _FakeDb: + def __init__(self) -> None: + self.counters: dict[str, Any] = {} + self.status = "running" + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: + sql = str(stmt) + if params and "counters" in params: + self.counters.update(json.loads(params["counters"])) + for st in ("done", "failed", "banned"): + if f"SET status = '{st}'" in sql and self.status == "running": + self.status = st + return MagicMock() + + def commit(self) -> None: ... + def rollback(self) -> None: ... + + +class _Scraper: + """Двойник Cian/Yandex/AvitoScraper: помнит якоря, роняет заданный TimeoutError'ом.""" + + visited: list[float] = [] # noqa: RUF012 + timeout_on: float | None = None + + gate_fetch_attempts = state_extraction_attempts = 1 + gate_fetch_failures = state_extraction_failures = 0 + request_delay_sec = 0.0 + + def __init__(self, *_a: Any, **_kw: Any) -> None: + self._browser = None + self._cffi = None + + async def __aenter__(self) -> _Scraper: + return self + + async def __aexit__(self, *_e: Any) -> None: + return None + + def _visit(self, lat: float) -> None: + _Scraper.visited.append(lat) + if lat == _Scraper.timeout_on: + raise TimeoutError + + async def fetch_around_multi_room(self, lat: float, *_a: Any, **kw: Any) -> list[Any]: + self._visit(lat) + if kw.get("on_combo") is not None: # yandex: единица чекпоинта — combo + kw["on_combo"](f"combo@{lat}", []) + return [] + return [types.SimpleNamespace(listing_segment="vtorichka", house_ext_id=None)] + + async def fetch_around(self, lat: float, *_a: Any, **_kw: Any) -> list[Any]: + self._visit(lat) + return [] + + +def _config() -> types.SimpleNamespace: + return types.SimpleNamespace( + scraper_fetch_mode="cffi", + scraper_proxy_url=None, + use_proxy_pool_browser=False, + browser_http_endpoint=None, + environment="test", + avito_serp_ok_not_banned=True, + ) + + +async def _sweep( + kind: str, + *, + timeout_on: float | None = None, + drain_after_first: bool = False, + cancel_after_first: bool = False, +) -> _FakeDb: + _Scraper.visited = [] + _Scraper.timeout_on = timeout_on + db = _FakeDb() + common: dict[str, Any] = { + "run_id": 2406, + "config": _config(), + "matcher": MagicMock(), + "anchors": [ANCHOR_A, ANCHOR_B], + "request_delay_sec": 0.0, + "shutdown_requested": lambda: drain_after_first and bool(_Scraper.visited), + } + with ( + patch.object(pl, "CianScraper", _Scraper), + patch.object(pl, "YandexRealtyScraper", _Scraper), + patch.object(pl, "AvitoScraper", _Scraper), + patch.object(pl, "AsyncSession", _Scraper), + patch.object(pl, "save_listings", lambda *_a, **_kw: (1, 0)), + patch.object( + pl.runs, "is_cancelled", lambda *_a: cancel_after_first and bool(_Scraper.visited) + ), + ): + if kind == "cian": + await pl.run_cian_city_sweep( + db, # type: ignore[arg-type] + enrich_houses=False, + detail_top_n=0, + **common, + ) + elif kind == "yandex": + enrichment = MagicMock() + enrichment.record_yandex_price_history.return_value = 0 + await pl.run_yandex_city_sweep( + db, # type: ignore[arg-type] + enrichment=enrichment, + enrich_address=False, + **common, + ) + else: + await pl.run_avito_city_sweep( + db, # type: ignore[arg-type] + enrichment=MagicMock(), + enrich_houses=False, + enrich_imv=False, + detail_top_n=0, + **common, + ) + return db + + +def _done_key(kind: str, anchor: tuple[float, float, str]) -> str: + return f"combo@{anchor[0]}" if kind == "yandex" else anchor[2] + + +# ── (а) таймаут якоря: якорь не засчитан, следующий обработан ────────────────── + + +@pytest.mark.parametrize("kind", ["cian", "yandex"]) +async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None: + db = await _sweep(kind, timeout_on=ANCHOR_A[0]) + + assert _Scraper.visited == [ANCHOR_A[0], ANCHOR_B[0]] + assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_B)] + assert db.counters["errors_count"] >= 1 + + +async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() -> None: + class _Dc(_Scraper): + blocked = False + geo_filtered = fetch_errors = bucket_start_index = 0 + buckets_completed, buckets_total = 1, 6 + completed_buckets = ["st"] # noqa: RUF012 + + async def fetch_city(self, **_kw: Any) -> list[Any]: + raise TimeoutError + + db = _FakeDb() + with ( + patch.object(pl, "DomClickScraper", _Dc), + patch.object(pl.runs, "is_cancelled", lambda *_a: False), + ): + counters = await pl.run_domclick_city_sweep( + db, # type: ignore[arg-type] + run_id=2406, + config=types.SimpleNamespace(browser_http_endpoint="http://x:9000"), + matcher=MagicMock(), + pages=1, + request_delay_sec=0.0, + ) + + assert counters.errors_count >= 1 + assert db.counters["done_buckets"] == ["st"] + assert db.status == "failed" + + +# ── (б) дрейн после первого якоря: interrupted=1, в чекпоинте ровно первый ────── + + +@pytest.mark.parametrize("kind", ["cian", "yandex"]) +async def test_drain_after_first_anchor_is_partial_and_interrupted(kind: str) -> None: + db = await _sweep(kind, drain_after_first=True) + + assert _Scraper.visited == [ANCHOR_A[0]] + assert db.counters["interrupted"] == 1 + assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)] + assert db.status == "done" + + +# ── (в) пользовательская отмена после первого якоря ──────────────────────────── + + +@pytest.mark.parametrize("kind", ["avito", "cian", "yandex"]) +async def test_user_cancel_after_first_anchor_stops_without_finalizing(kind: str) -> None: + db = await _sweep(kind, cancel_after_first=True) + + assert _Scraper.visited == [ANCHOR_A[0]] + assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)] + assert "interrupted" not in db.counters + # Строку финализирует сам mark_cancelled (UI); свип только останавливается. + assert db.status == "running" -- 2.45.3 From c7a03495e2b5ae42bfe03f176121424e0052b1bf Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 13:54:55 +0500 Subject: [PATCH 5/7] =?UTF-8?q?=D0=A1=D0=B2=D0=B8=D0=BF=20=D0=94=D0=BE?= =?UTF-8?q?=D0=BC=D0=9A=D0=BB=D0=B8=D0=BA=D0=B0:=20=D0=BA=D0=BE=D1=80?= =?UTF-8?q?=D0=B7=D0=B8=D0=BD=D0=B0=20=D0=B1=D0=B5=D0=B7=20=D1=81=D0=BE?= =?UTF-8?q?=D1=85=D1=80=D0=B0=D0=BD=D1=91=D0=BD=D0=BD=D1=8B=D1=85=20=D1=81?= =?UTF-8?q?=D1=82=D1=80=D0=BE=D0=BA=20=D0=BD=D0=B5=20=D0=BF=D0=BE=D0=BF?= =?UTF-8?q?=D0=B0=D0=B4=D0=B0=D0=B5=D1=82=20=D0=B2=20=D1=87=D0=B5=D0=BA?= =?UTF-8?q?=D0=BF=D0=BE=D0=B8=D0=BD=D1=82=20(#2406)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью PR #3561: тест таймаута фазы закреплял как контракт done_buckets == ['st']. Лоты ДомКлика копятся в памяти и пишутся одним save_listings после всех корзин, поэтому снятая watchdog'ом (или упавшая на save) фаза не сохраняет ничего, а чекпоинт всё равно получал completed_buckets живого скрейпера — следующий прогон пропускал корзину навсегда (механизм миграции 308). Чекпоинт пополняется только если фаза дошла до конца (флаг _saved после save). Тест развёрнут: при таймауте fetch_city и при падении save_listings done_buckets == [], статус failed. Контроль «полный проход несёт унаследованное ∪ пройденное» — test_3369. Co-Authored-By: Claude Opus 5 --- .../tests/test_2406_sweep_edge_cases.py | 32 ++++++++++++++++--- .../src/scraper_kit/orchestration/pipeline.py | 13 ++++++-- 2 files changed, 39 insertions(+), 6 deletions(-) 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 8b15a43d..feeac535 100644 --- a/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py +++ b/tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py @@ -3,7 +3,8 @@ Уже покрыто: avito — таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369). Здесь остаток: - (а) таймаут якоря у cian и yandex, таймаут фазы у domclick; + (а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает + корзины, строки которых не сохранены; (б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex; (в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito, cian и yandex. @@ -168,7 +169,16 @@ async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None: assert db.counters["errors_count"] >= 1 -async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() -> None: +@pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"]) +async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None: + """Корзина без сохранённых строк не попадает в чекпоинт. + + Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин. + Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из + корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы, + что следующий прогон пропустит её через skip_buckets навсегда (миграция 308). + """ + class _Dc(_Scraper): blocked = False geo_filtered = fetch_errors = bucket_start_index = 0 @@ -176,11 +186,22 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() completed_buckets = ["st"] # noqa: RUF012 async def fetch_city(self, **_kw: Any) -> list[Any]: - raise TimeoutError + if broken == "fetch_timeout": + raise TimeoutError + return [object()] + + saved: list[int] = [] + + def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]: + if broken == "save_failed": + raise RuntimeError("save_listings упал") + saved.append(len(lots)) + return len(lots), 0 db = _FakeDb() with ( patch.object(pl, "DomClickScraper", _Dc), + patch.object(pl, "save_listings", _save), patch.object(pl.runs, "is_cancelled", lambda *_a: False), ): counters = await pl.run_domclick_city_sweep( @@ -192,8 +213,11 @@ async def test_domclick_phase_timeout_keeps_completed_buckets_and_is_not_done() request_delay_sec=0.0, ) + assert saved == [], "тест не о том: строки сохранились" assert counters.errors_count >= 1 - assert db.counters["done_buckets"] == ["st"] + assert db.counters["done_buckets"] == [], ( + f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}" + ) assert db.status == "failed" 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 4fae93a3..0c69733b 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 @@ -4992,10 +4992,13 @@ async def run_domclick_city_sweep( ) lots: list[ScrapedLot] = [] + # #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому + # корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже). + _saved = False async def _domclick_phase() -> None: """Единственная citywide-фаза: fetch_city + save.""" - nonlocal lots + nonlocal lots, _saved async with DomClickScraper( config, proxy_provider=proxy_provider, @@ -5041,6 +5044,7 @@ async def run_domclick_city_sweep( ) counters.lots_inserted += inserted counters.lots_updated += updated + _saved = True try: await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout) @@ -5094,7 +5098,12 @@ async def run_domclick_city_sweep( # #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне. # Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не # затирают, и оборванный болезнью финализации прогон его не теряет. - _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) + # #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза + # не дошла до save_listings — у её корзин в БД ноль строк, а отметка + # «пройдена» заставила бы следующий прогон пропустить их навсегда + # (механизм разобран в миграции 308, из-за него выключены свипы 77/50). + if _saved: + _checkpoint = sorted(skip_buckets | set(_s.completed_buckets)) runs.update_heartbeat(db, run_id, _payload()) counters.bucket_start_index = _s.bucket_start_index -- 2.45.3 From 2e4a82dd469674a296dad4ffd038619797ef6f5e Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 13:57:34 +0500 Subject: [PATCH 6/7] =?UTF-8?q?=D0=A3=D0=BF=D0=B0=D0=B2=D1=88=D0=B8=D0=B9?= =?UTF-8?q?=20=D0=BF=D1=80=D0=BE=D0=B3=D0=BE=D0=BD=20=D1=84=D0=B8=D0=BD?= =?UTF-8?q?=D0=B0=D0=BB=D0=B8=D0=B7=D0=B8=D1=80=D1=83=D0=B5=D1=82=20=D0=BE?= =?UTF-8?q?=D0=B1=D1=89=D0=B8=D0=B9=20runs.mark=5Fcrashed:=20=D0=B8=20?= =?UTF-8?q?=D1=80=D1=83=D1=87=D0=BD=D1=8B=D0=B5=20=D0=B7=D0=B0=D0=BF=D1=83?= =?UTF-8?q?=D1=81=D0=BA=D0=B8=20=D0=B8=D0=B7=20=D0=B0=D0=B4=D0=BC=D0=B8?= =?UTF-8?q?=D0=BD=D0=BA=D0=B8,=20=D0=B8=20=D0=B1=D0=B5=D0=B7=20=D0=B4?= =?UTF-8?q?=D1=83=D0=B1=D0=BB=D1=8F=20=D0=B0=D0=BB=D0=B5=D1=80=D1=82=D0=B0?= =?UTF-8?q?=20(#1940)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью PR #3561, две находки. 1. Ручные запуски из админки (avito/cian city sweep, cian_full_load, yandex_full_load, yandex_city_sweep) идут мимо scheduler._dispatch: в except был только logger.exception, и отказ «пул пуст» до try-финализатора пайплайна оставлял строку running до zombie. Логика финализации переехала из планировщика в runs.mark_crashed; её зовут планировщик и все пять ручек. 2. Обычный путь пайплайна — mark_failed и raise; планировщик звал mark_failed второй раз. UPDATE — no-op, но _alert_on_run_id срабатывал снова с тем же стриком, и на вехе лестницы в Sentry уходил дубль (+ WARNING «no-op»). mark_crashed сначала читает статус и финализирует только running. Тесты по значению над двойником строки scrape_runs и боевыми mark_*: статус banned/infra или failed после планировщика и после каждой из пяти ручек, одна отправка в Sentry на неудаче пайплайна. Co-Authored-By: Claude Opus 5 --- tradein-mvp/backend/app/api/v1/admin.py | 25 ++- .../test_1940_crashed_run_is_finalized.py | 195 ++++++++++++------ .../src/scraper_kit/orchestration/runs.py | 32 +++ .../scraper_kit/orchestration/scheduler.py | 23 +-- 4 files changed, 187 insertions(+), 88 deletions(-) diff --git a/tradein-mvp/backend/app/api/v1/admin.py b/tradein-mvp/backend/app/api/v1/admin.py index aa3abd64..e973965e 100644 --- a/tradein-mvp/backend/app/api/v1/admin.py +++ b/tradein-mvp/backend/app/api/v1/admin.py @@ -1252,8 +1252,11 @@ async def start_avito_city_sweep( request_delay_sec=payload.request_delay_sec, enrich_imv=payload.enrich_imv, ) - except Exception: + except Exception as exc: logger.exception("city-sweep background task run_id=%d crashed", run_id) + # #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший + # до финализатора пайплайна прогон висел 'running' до zombie. + runs_mod.mark_crashed(sweep_db, run_id, exc) finally: sweep_db.close() @@ -1344,8 +1347,11 @@ async def start_cian_city_sweep( detail_top_n=payload.detail_top_n, enrich_houses=payload.enrich_houses, ) - except Exception: + except Exception as exc: logger.exception("cian-sweep background task run_id=%d crashed", run_id) + # #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший + # до финализатора пайплайна прогон висел 'running' до zombie. + runs_mod.mark_crashed(sweep_db, run_id, exc) finally: sweep_db.close() @@ -1487,8 +1493,11 @@ async def start_cian_full_load( resume_run_id=payload.resume_run_id, secondary_only=payload.secondary_only, ) - except Exception: + except Exception as exc: logger.exception("cian-full-load background task run_id=%d crashed", run_id) + # #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший + # до финализатора пайплайна прогон висел 'running' до zombie. + runs_mod.mark_crashed(task_db, run_id, exc) finally: task_db.close() @@ -1590,8 +1599,11 @@ async def start_yandex_full_load( concurrency=payload.concurrency, resume_run_id=payload.resume_run_id, ) - except Exception: + except Exception as exc: logger.exception("yandex-full-load background task run_id=%d crashed", run_id) + # #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший + # до финализатора пайплайна прогон висел 'running' до zombie. + runs_mod.mark_crashed(task_db, run_id, exc) finally: task_db.close() @@ -1658,8 +1670,11 @@ async def start_yandex_city_sweep( request_delay_sec=payload.request_delay_sec, enrich_address=payload.enrich_address, ) - except Exception: + except Exception as exc: logger.exception("yandex-sweep background task run_id=%d crashed", run_id) + # #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший + # до финализатора пайплайна прогон висел 'running' до zombie. + runs_mod.mark_crashed(sweep_db, run_id, exc) finally: sweep_db.close() diff --git a/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py b/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py index 447f6eb1..676f03df 100644 --- a/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py +++ b/tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py @@ -4,77 +4,85 @@ бросил NoProxyAvailableError, планировщик записал «crashed run_id=7349» и всё. Строка осталась 'running' с heartbeat_at == started_at; 21 сутки до этого так ушли в 'zombie' 18 avito-свипов, каждый через 6 ч. Проверяется значение статуса в строке прогона. + +Финализатор — `runs.mark_crashed`. Его зовут оба запускающих: планировщик и ручные +запуски из админки (они идут мимо `_dispatch`). Здесь боевой `mark_crashed` над +двойником одной строки scrape_runs: `_Row` отвечает на те же SQL, что пишет runs.py, +с гейтом `WHERE status = 'running'`. """ from __future__ import annotations import asyncio import os +import re +import types from typing import Any -from unittest.mock import MagicMock +from unittest.mock import AsyncMock, MagicMock, patch os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") -from scraper_kit.orchestration.runs import BAN_KIND_INFRA +import pytest +from scraper_kit.orchestration import runs as kit_runs +from scraper_kit.orchestration.runs import BAN_KIND_INFRA, CONSECUTIVE_FAILURE_ALERT_THRESHOLD from scraper_kit.orchestration.scheduler import Handler, SchedulerContext, _dispatch from scraper_kit.proxy_errors import NoProxyAvailableError - -class _Db: - def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: - res = MagicMock() - res.fetchone.return_value = None # «не running» для гейта _claim_run - res.scalar.return_value = True # advisory-lock взят - return res - - def commit(self) -> None: - pass - - def rollback(self) -> None: - pass - - def close(self) -> None: - pass +RUN_ID = 7349 -class _Runs: - """ctx.runs с боевым гейтом `WHERE status = 'running'` и мержем counters.""" +class _Row: + """Сессия над одной строкой scrape_runs (status/ban_kind/error/counters). + + История источника для стрик-алерта — `CONSECUTIVE_FAILURE_ALERT_THRESHOLD` неудач: + ровно веха лестницы `_streak_alert_due`, на которой дубль и уходил бы в Sentry. + """ def __init__(self) -> None: - self.rows: dict[int, dict[str, Any]] = {} + self.status = "running" + self.ban_kind: str | None = None + self.error: str | None = None + self.updates = 0 - def create_run(self, db: Any, *, source: str, params: dict[str, Any]) -> int: - run_id = 7349 + len(self.rows) - self.rows[run_id] = {"status": "running", "counters": {}, "ban_kind": None} - return run_id + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any: + sql = str(stmt) + p = params or {} + res = MagicMock() + res.fetchone.return_value = None # гейт _claim_run: «не running» + res.scalar.return_value = True # advisory-lock взят + if "INSERT INTO scrape_runs" in sql: + res.fetchone.return_value = types.SimpleNamespace(id=RUN_ID) + elif "SELECT status FROM scrape_runs WHERE id" in sql: + res.fetchone.return_value = types.SimpleNamespace(status=self.status) + elif "UPDATE scrape_runs" in sql and "SET status" in sql: + hit = self.status == "running" + if hit: + m = re.search(r"SET status = '(\w+)'", sql) + assert m is not None + self.status = m.group(1) + self.ban_kind = p.get("ban_kind") + self.error = p.get("error") + self.updates += 1 + res.first.return_value = types.SimpleNamespace(id=RUN_ID) if hit else None + elif "SELECT source FROM scrape_runs" in sql: + res.fetchone.return_value = types.SimpleNamespace(source="avito_city_sweep") + elif "SELECT status, counters FROM scrape_runs" in sql: + res.fetchall.return_value = [ + types.SimpleNamespace(status="failed", counters={}) + ] * CONSECUTIVE_FAILURE_ALERT_THRESHOLD + return res - def _finish(self, run_id: int, status: str, counters: dict[str, Any], **extra: Any) -> None: - row = self.rows[run_id] - if row["status"] != "running": - return - row.update(status=status, **extra) - row["counters"] = {**row["counters"], **counters} - - def mark_done(self, db: Any, run_id: int, counters: dict[str, Any]) -> None: - self._finish(run_id, "done", counters) - - def mark_failed(self, db: Any, run_id: int, error: str, counters: dict[str, Any]) -> None: - self._finish(run_id, "failed", counters, error=error) - - def mark_banned( - self, db: Any, run_id: int, error: str, counters: dict[str, Any], *, ban_kind: str - ) -> None: - self._finish(run_id, "banned", counters, error=error, ban_kind=ban_kind) + def commit(self) -> None: ... + def rollback(self) -> None: ... + def close(self) -> None: ... -async def _dispatch_and_wait(job: Any) -> dict[str, Any]: - runs = _Runs() +async def _dispatch_and_wait(job: Any, row: _Row) -> None: ctx = SchedulerContext( config=MagicMock(), matcher=MagicMock(), enrichment=MagicMock(), - session_factory=_Db, - runs=runs, + session_factory=lambda: row, ) sched = { "id": 1, @@ -83,21 +91,23 @@ async def _dispatch_and_wait(job: Any) -> dict[str, Any]: "window_end_hour": 7, "default_params": {}, } - run_id = await _dispatch(Handler(job, "avito_city_sweep"), _Db(), sched, ctx) - assert run_id is not None + run_id = await _dispatch(Handler(job, "avito_city_sweep"), row, sched, ctx) # type: ignore[arg-type] + assert run_id == RUN_ID await asyncio.gather(*ctx._inflight_tasks, return_exceptions=True) - return runs.rows[run_id] + + +# ── планировщик ───────────────────────────────────────────────────────────────── async def test_no_proxy_crash_before_try_marks_run_banned_infra() -> None: async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: raise NoProxyAvailableError("avito") - row = await _dispatch_and_wait(_job) + row = _Row() + await _dispatch_and_wait(_job, row) - assert row["status"] == "banned" - assert row["ban_kind"] == BAN_KIND_INFRA - assert "no proxy available" in row["error"] + assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA) + assert row.error is not None and "no proxy available" in row.error async def test_wrapped_no_proxy_is_still_infra() -> None: @@ -107,27 +117,86 @@ async def test_wrapped_no_proxy_is_still_infra() -> None: except NoProxyAvailableError as exc: raise RuntimeError("sidecar fetch failed") from exc - row = await _dispatch_and_wait(_job) + row = _Row() + await _dispatch_and_wait(_job, row) - assert (row["status"], row["ban_kind"]) == ("banned", BAN_KIND_INFRA) + assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA) async def test_other_crash_marks_run_failed() -> None: async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: raise KeyError("anchors") - row = await _dispatch_and_wait(_job) + row = _Row() + await _dispatch_and_wait(_job, row) - assert row["status"] == "failed" - assert "KeyError" in row["error"] + assert row.status == "failed" + assert row.error is not None and "KeyError" in row.error -async def test_handler_own_final_status_and_checkpoint_survive() -> None: +async def test_handler_own_final_status_survives() -> None: async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: - c.runs.mark_done(db, run_id, {"done_buckets": ["center"]}) + kit_runs.mark_done(db, run_id, {"done_buckets": ["center"]}) raise RuntimeError("after mark_done") - row = await _dispatch_and_wait(_job) + row = _Row() + await _dispatch_and_wait(_job, row) - assert row["status"] == "done" - assert row["counters"] == {"done_buckets": ["center"]} + assert (row.status, row.updates) == ("done", 1) + + +async def test_pipeline_failure_is_alerted_once_not_twice() -> None: + """Обычный путь пайплайна: mark_failed и raise. Повторная финализация не шлёт дубль. + + Стрик после первой финализации стоит на вехе лестницы; no-op mark_failed из + планировщика позвал бы `_alert_on_run_id` ещё раз с тем же стриком. + """ + + async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: + kit_runs.mark_failed(db, run_id, "cian-sweep fatal", {}) + raise RuntimeError("cian-sweep fatal") + + row = _Row() + sentry = MagicMock() + with patch.object(kit_runs, "sentry_sdk", sentry): + await _dispatch_and_wait(_job, row) + + assert (row.status, row.updates) == ("failed", 1) + assert sentry.capture_message.call_count == 1 + + +# ── ручные запуски из админки (мимо _dispatch) ───────────────────────────────── + + +@pytest.mark.parametrize( + ("path", "runner"), + [ + ("avito-city-sweep", "run_avito_city_sweep"), + ("cian-city-sweep", "run_cian_city_sweep"), + ("cian-full-load", "run_cian_full_load"), + ("yandex-full-load", "run_yandex_full_load"), + ("yandex-city-sweep", "run_yandex_city_sweep"), + ], +) +def test_admin_manual_launch_crash_is_finalized(path: str, runner: str) -> None: + from fastapi import FastAPI + from fastapi.testclient import TestClient + + from app.api.v1 import admin as admin_module + from app.core.db import get_db + + app = FastAPI() + app.include_router(admin_module.router, prefix="/api/v1/admin") + row = _Row() + app.dependency_overrides[get_db] = lambda: row + + with ( + patch.object(kit_runs, "sentry_sdk", None), + patch.object(admin_module, "has_running_run", return_value=False), + patch.object(admin_module, "SessionLocal", lambda: row), + patch.object(admin_module, runner, AsyncMock(side_effect=NoProxyAvailableError("avito"))), + ): + r = TestClient(app).post(f"/api/v1/admin/scrape/{path}", json={}) + + assert r.status_code == 200, r.text + assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA) 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 b2d32d91..2053530b 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 @@ -46,6 +46,7 @@ from sqlalchemy import text from sqlalchemy.orm import Session from scraper_kit.orchestration.run_context import current_run_id +from scraper_kit.proxy_errors import caused_by_no_proxy try: # sentry опционален — kit не тянет его в зависимостях import sentry_sdk @@ -1017,6 +1018,37 @@ def mark_banned( _alert_on_run_id(db, run_id) +def mark_crashed(db: Session, run_id: int, exc: BaseException) -> None: + """Финализировать прогон, чья задача вылетела исключением наружу (#1940). + + Зовут те, кто запускает пайплайн: планировщик (`scheduler._dispatch`) и ручные + запуски из админки. Исключение, брошенное ДО try-финализатора пайплайна (прод: + пустой пул прокси в `BrowserFetcher.__aenter__` через 10 мс после claim), + оставляло строку 'running' без пульса до zombie-reaper'а через 6 ч. + + Строку, которую пайплайн уже финализировал сам (обычный путь: mark_failed и + raise), не трогаем и ПОВТОРНО mark_* не зовём: UPDATE был бы no-op, но mark_* + всё равно позвал бы `_alert_on_run_id` — стрик тот же, и на вехе лестницы + (`_streak_alert_due`) в Sentry ушёл бы дубль, плюс WARNING «no-op» на каждом + таком падении. Best-effort: сбой финализации логируется, не бросается. + """ + try: + db.rollback() # сессия упавшей задачи может быть в aborted-транзакции + row = db.execute( + text("SELECT status FROM scrape_runs WHERE id = :id"), + {"id": run_id}, + ).fetchone() + db.rollback() # чтение не держит транзакцию (#3480) + if row is None or row.status != "running": + return + if caused_by_no_proxy(exc): + mark_banned(db, run_id, f"crashed: {exc}", {}, ban_kind=BAN_KIND_INFRA) + else: + mark_failed(db, run_id, f"crashed: {type(exc).__name__}: {exc}", {}) + except Exception: + logger.exception("mark_crashed: не удалось финализировать run_id=%d", run_id) + + def _dominant_ban_kind(census: Mapping[str, int]) -> str: """Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178). 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 47932091..51bd6702 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 @@ -57,8 +57,6 @@ from scraper_kit.orchestration.pipeline import ( run_domclick_city_sweep, run_yandex_city_sweep, ) -from scraper_kit.orchestration.runs import BAN_KIND_INFRA -from scraper_kit.proxy_errors import caused_by_no_proxy if TYPE_CHECKING: from sqlalchemy.orm import Session @@ -1028,24 +1026,9 @@ async def _dispatch( await handler.job(run_db, run_id, params, ctx) except Exception as exc: logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id) - # #1940: исключение, вылетевшее ДО try-финализатора хендлера (прод: пустой - # пул прокси в BrowserFetcher.__aenter__ через 10 мс после claim), оставляло - # строку 'running' без единого пульса — через 6 ч её снимал reaper как - # 'zombie', источник терял день. Финализируем здесь, в единственной точке, - # через которую идёт любой хендлер. Уже финализированную хендлером строку - # не трогаем: mark_* пишут только `WHERE status = 'running'`, а {} counters - # мержится и чекпоинт не стирает. - try: - if caused_by_no_proxy(exc): - ctx.runs.mark_banned( - run_db, run_id, f"crashed: {exc}", {}, ban_kind=BAN_KIND_INFRA - ) - else: - ctx.runs.mark_failed( - run_db, run_id, f"crashed: {type(exc).__name__}: {exc}", {} - ) - except Exception: - logger.exception("scheduler: не удалось финализировать run_id=%d", run_id) + # #1940: исключение мимо try-финализатора хендлера оставляло строку + # 'running' до zombie через 6 ч. См. runs.mark_crashed. + ctx.runs.mark_crashed(run_db, run_id, exc) finally: run_db.close() -- 2.45.3 From 69e70473985d0043c78a25d6ca562a8109deb817 Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 13:58:43 +0500 Subject: [PATCH 7/7] =?UTF-8?q?=D0=A2=D0=B5=D1=81=D1=82=D1=8B=20=D0=B3?= =?UTF-8?q?=D0=B5=D0=B9=D1=82=D0=B0=20=D1=87=D0=B5=D0=BA=D0=BF=D0=BE=D0=B8?= =?UTF-8?q?=D0=BD=D1=82=D0=B0:=20=D1=87=D0=B0=D1=81=D1=82=D0=B8=D1=87?= =?UTF-8?q?=D0=BD=D1=8B=D0=B9=20=D0=BE=D1=82=D0=BA=D0=B0=D0=B7=20=D1=84?= =?UTF-8?q?=D0=B0=D0=B7=D1=8B,=20=D0=BF=D0=BE=D1=80=D0=BE=D0=B3=20=D0=B2?= =?UTF-8?q?=20=D0=BE=D0=B4=D0=BD=D1=83=20=D0=BF=D0=BE=D0=BF=D1=8B=D1=82?= =?UTF-8?q?=D0=BA=D1=83,=20=D1=83=D0=BF=D0=B0=D0=B2=D1=88=D0=B8=D0=B9=20?= =?UTF-8?q?=D0=B7=D0=B0=D0=BF=D1=80=D0=BE=D1=81=20=D0=B4=D0=BE=D0=BC=D0=BE?= =?UTF-8?q?=D0=B2=20(#3415)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Ревью PR #3561: три мутации гейта оставались зелёными — порог `attempted >= 2`, жадный гейт `failed >= 1` и снятый `houses_attempted += len(nb_id_list)` в ветке «houses DB query failed». Добавлены проверки по чекпоинту: - 1 из 2 домов отказал в каждом якоре — оба якоря в чекпоинте (контроль жадности); - единственный дом якоря отказал — якорь не в чекпоинте (порог); - упал запрос домов — якоря не в чекпоинте, houses_attempted == houses_failed == 2. Код не менялся. Co-Authored-By: Claude Opus 5 --- ...test_3415_bucket_done_only_if_phases_ok.py | 61 +++++++++++++++++-- 1 file changed, 56 insertions(+), 5 deletions(-) diff --git a/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py b/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py index f74aaaa9..6b0155a8 100644 --- a/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py +++ b/tradein-mvp/backend/tests/test_3415_bucket_done_only_if_phases_ok.py @@ -49,8 +49,14 @@ DETAIL_ROWS = [ class _FakeDb: - def __init__(self, prev_counters: dict[str, Any] | None = None) -> None: + def __init__( + self, + prev_counters: dict[str, Any] | None = None, + house_rows: list[dict[str, Any]] | None = None, + ) -> None: self.prev_counters = prev_counters or {} + # [] — запрос домов падает (ветка «houses DB query failed»). + self.house_rows = HOUSE_ROWS if house_rows is None else house_rows self.heartbeats: list[dict[str, Any]] = [] def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any: @@ -61,7 +67,9 @@ class _FakeDb: return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters)) if params and "ids" in params: # houses-фаза: дома, у которых есть cian_zhk_url. - return MagicMock(mappings=lambda: MagicMock(all=lambda: HOUSE_ROWS)) + if not self.house_rows: + raise RuntimeError("houses DB query failed") + return MagicMock(mappings=lambda: MagicMock(all=lambda: self.house_rows)) if params and "limit" in params: # detail-фаза avito: карточки-кандидаты на обогащение. return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS)) @@ -114,16 +122,23 @@ def _config() -> types.SimpleNamespace: async def _run( - prev: dict[str, Any] | None, *, houses_ok: bool, run_id: int = 8401 + prev: dict[str, Any] | None, + *, + houses_ok: bool, + run_id: int = 8401, + house_rows: list[dict[str, Any]] | None = None, + failing_urls: frozenset[str] | None = None, ) -> tuple[_FakeDb, Any]: from scraper_kit.orchestration import pipeline as pl _FakeScraper.visited = [] - db = _FakeDb(prev) + db = _FakeDb(prev, house_rows) - async def _fake_newbuilding(*_a: Any, **_kw: Any) -> Any: + async def _fake_newbuilding(zhk_url: str, *_a: Any, **_kw: Any) -> Any: # None — ровно тот отказ, что был на проде 06.09: исключения нет, # счётчик houses_failed растёт, фаза возвращается штатно. + if failing_urls is not None: + return None if zhk_url in failing_urls else object() return object() if houses_ok else None with ( @@ -189,6 +204,42 @@ async def test_healthy_anchor_is_still_checkpointed() -> None: ) +@pytest.mark.asyncio +async def test_anchor_with_partially_failed_houses_is_still_checkpointed() -> None: + """Контроль жадности: один отказавший дом из двух — якорь пройден. + + Гейт про ПОЛНЫЙ отказ фазы. Жадный гейт (любой отказ) собирал бы такой якорь + заново каждым прогоном — на проде houses почти всегда теряет пару домов. + """ + db, counters = await _run( + None, houses_ok=True, failing_urls=frozenset({HOUSE_ROWS[0]["cian_zhk_url"]}) + ) + assert (counters.houses_attempted, counters.houses_failed) == (4, 2) + assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], ( + f"якорь с частичным отказом houses не помечен пройденным: {_last_checkpoint(db)}" + ) + + +@pytest.mark.asyncio +async def test_single_failed_attempt_blocks_checkpoint() -> None: + """Порог — одна попытка: единственный дом якоря отказал — якорь не пройден.""" + db, counters = await _run(None, houses_ok=False, house_rows=HOUSE_ROWS[:1]) + assert (counters.houses_attempted, counters.houses_enriched) == (2, 0) + assert _last_checkpoint(db) == [], ( + f"якорь с 1 из 1 отказавшим домом помечен пройденным: {_last_checkpoint(db)}" + ) + + +@pytest.mark.asyncio +async def test_houses_db_query_failure_blocks_checkpoint() -> None: + """Упал запрос домов: попытки считаются вместе с отказами, якорь не пройден.""" + db, counters = await _run(None, houses_ok=True, house_rows=[]) + assert _last_checkpoint(db) == [], ( + f"якорь с упавшим запросом домов помечен пройденным: {_last_checkpoint(db)}" + ) + assert (counters.houses_attempted, counters.houses_failed) == (2, 2) + + class _FakeAvitoScraper: """Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза.""" -- 2.45.3