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()