From 8e5f097c5eaa929b216f4993ac091c91f401718d Mon Sep 17 00:00:00 2001 From: bot-backend Date: Thu, 17 Sep 2026 12:47:05 +0500 Subject: [PATCH] =?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()