Сбор МЕРА: упавший прогон не висит zombie 6 часов, провалившаяся фаза не попадает в чекпоинт, проверка отмены не держит транзакцию часами #3561

Merged
bot-backend merged 8 commits from fix/kit-orchestration into main 2026-09-17 09:23:44 +00:00
4 changed files with 187 additions and 88 deletions
Showing only changes of commit 2e4a82dd46 - Show all commits

View file

@ -1252,8 +1252,11 @@ async def start_avito_city_sweep(
request_delay_sec=payload.request_delay_sec, request_delay_sec=payload.request_delay_sec,
enrich_imv=payload.enrich_imv, enrich_imv=payload.enrich_imv,
) )
except Exception: except Exception as exc:
logger.exception("city-sweep background task run_id=%d crashed", run_id) 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: finally:
sweep_db.close() sweep_db.close()
@ -1344,8 +1347,11 @@ async def start_cian_city_sweep(
detail_top_n=payload.detail_top_n, detail_top_n=payload.detail_top_n,
enrich_houses=payload.enrich_houses, enrich_houses=payload.enrich_houses,
) )
except Exception: except Exception as exc:
logger.exception("cian-sweep background task run_id=%d crashed", run_id) 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: finally:
sweep_db.close() sweep_db.close()
@ -1487,8 +1493,11 @@ async def start_cian_full_load(
resume_run_id=payload.resume_run_id, resume_run_id=payload.resume_run_id,
secondary_only=payload.secondary_only, secondary_only=payload.secondary_only,
) )
except Exception: except Exception as exc:
logger.exception("cian-full-load background task run_id=%d crashed", run_id) 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: finally:
task_db.close() task_db.close()
@ -1590,8 +1599,11 @@ async def start_yandex_full_load(
concurrency=payload.concurrency, concurrency=payload.concurrency,
resume_run_id=payload.resume_run_id, 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) 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: finally:
task_db.close() task_db.close()
@ -1658,8 +1670,11 @@ async def start_yandex_city_sweep(
request_delay_sec=payload.request_delay_sec, request_delay_sec=payload.request_delay_sec,
enrich_address=payload.enrich_address, enrich_address=payload.enrich_address,
) )
except Exception: except Exception as exc:
logger.exception("yandex-sweep background task run_id=%d crashed", run_id) 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: finally:
sweep_db.close() sweep_db.close()

View file

@ -4,77 +4,85 @@
бросил NoProxyAvailableError, планировщик записал «crashed run_id=7349» и всё. Строка бросил NoProxyAvailableError, планировщик записал «crashed run_id=7349» и всё. Строка
осталась 'running' с heartbeat_at == started_at; 21 сутки до этого так ушли в 'zombie' осталась 'running' с heartbeat_at == started_at; 21 сутки до этого так ушли в 'zombie'
18 avito-свипов, каждый через 6 ч. Проверяется значение статуса в строке прогона. 18 avito-свипов, каждый через 6 ч. Проверяется значение статуса в строке прогона.
Финализатор `runs.mark_crashed`. Его зовут оба запускающих: планировщик и ручные
запуски из админки (они идут мимо `_dispatch`). Здесь боевой `mark_crashed` над
двойником одной строки scrape_runs: `_Row` отвечает на те же SQL, что пишет runs.py,
с гейтом `WHERE status = 'running'`.
""" """
from __future__ import annotations from __future__ import annotations
import asyncio import asyncio
import os import os
import re
import types
from typing import Any 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") 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.orchestration.scheduler import Handler, SchedulerContext, _dispatch
from scraper_kit.proxy_errors import NoProxyAvailableError from scraper_kit.proxy_errors import NoProxyAvailableError
RUN_ID = 7349
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: class _Row:
"""ctx.runs с боевым гейтом `WHERE status = 'running'` и мержем counters.""" """Сессия над одной строкой scrape_runs (status/ban_kind/error/counters).
История источника для стрик-алерта `CONSECUTIVE_FAILURE_ALERT_THRESHOLD` неудач:
ровно веха лестницы `_streak_alert_due`, на которой дубль и уходил бы в Sentry.
"""
def __init__(self) -> None: 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: def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
run_id = 7349 + len(self.rows) sql = str(stmt)
self.rows[run_id] = {"status": "running", "counters": {}, "ban_kind": None} p = params or {}
return run_id 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: def commit(self) -> None: ...
row = self.rows[run_id] def rollback(self) -> None: ...
if row["status"] != "running": def close(self) -> None: ...
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]: async def _dispatch_and_wait(job: Any, row: _Row) -> None:
runs = _Runs()
ctx = SchedulerContext( ctx = SchedulerContext(
config=MagicMock(), config=MagicMock(),
matcher=MagicMock(), matcher=MagicMock(),
enrichment=MagicMock(), enrichment=MagicMock(),
session_factory=_Db, session_factory=lambda: row,
runs=runs,
) )
sched = { sched = {
"id": 1, "id": 1,
@ -83,21 +91,23 @@ async def _dispatch_and_wait(job: Any) -> dict[str, Any]:
"window_end_hour": 7, "window_end_hour": 7,
"default_params": {}, "default_params": {},
} }
run_id = await _dispatch(Handler(job, "avito_city_sweep"), _Db(), sched, ctx) run_id = await _dispatch(Handler(job, "avito_city_sweep"), row, sched, ctx) # type: ignore[arg-type]
assert run_id is not None assert run_id == RUN_ID
await asyncio.gather(*ctx._inflight_tasks, return_exceptions=True) 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 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: async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
raise NoProxyAvailableError("avito") raise NoProxyAvailableError("avito")
row = await _dispatch_and_wait(_job) row = _Row()
await _dispatch_and_wait(_job, row)
assert row["status"] == "banned" assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA)
assert row["ban_kind"] == BAN_KIND_INFRA assert row.error is not None and "no proxy available" in row.error
assert "no proxy available" in row["error"]
async def test_wrapped_no_proxy_is_still_infra() -> None: 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: except NoProxyAvailableError as exc:
raise RuntimeError("sidecar fetch failed") from 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 test_other_crash_marks_run_failed() -> None:
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None: async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
raise KeyError("anchors") raise KeyError("anchors")
row = await _dispatch_and_wait(_job) row = _Row()
await _dispatch_and_wait(_job, row)
assert row["status"] == "failed" assert row.status == "failed"
assert "KeyError" in row["error"] 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: 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") 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.status, row.updates) == ("done", 1)
assert row["counters"] == {"done_buckets": ["center"]}
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)

View file

@ -46,6 +46,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
from scraper_kit.orchestration.run_context import current_run_id from scraper_kit.orchestration.run_context import current_run_id
from scraper_kit.proxy_errors import caused_by_no_proxy
try: # sentry опционален — kit не тянет его в зависимостях try: # sentry опционален — kit не тянет его в зависимостях
import sentry_sdk import sentry_sdk
@ -1017,6 +1018,37 @@ def mark_banned(
_alert_on_run_id(db, run_id) _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: def _dominant_ban_kind(census: Mapping[str, int]) -> str:
"""Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178). """Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178).

View file

@ -57,8 +57,6 @@ from scraper_kit.orchestration.pipeline import (
run_domclick_city_sweep, run_domclick_city_sweep,
run_yandex_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: if TYPE_CHECKING:
from sqlalchemy.orm import Session from sqlalchemy.orm import Session
@ -1028,24 +1026,9 @@ async def _dispatch(
await handler.job(run_db, run_id, params, ctx) await handler.job(run_db, run_id, params, ctx)
except Exception as exc: except Exception as exc:
logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id) logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id)
# #1940: исключение, вылетевшее ДО try-финализатора хендлера (прод: пустой # #1940: исключение мимо try-финализатора хендлера оставляло строку
# пул прокси в BrowserFetcher.__aenter__ через 10 мс после claim), оставляло # 'running' до zombie через 6 ч. См. runs.mark_crashed.
# строку 'running' без единого пульса — через 6 ч её снимал reaper как ctx.runs.mark_crashed(run_db, run_id, exc)
# '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: finally:
run_db.close() run_db.close()