Планировщик: упавший до финализатора прогон получает статус сразу, а не zombie через 6 часов (#1940)
Исключение, вылетевшее из хендлера до его 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 <noreply@anthropic.com>
This commit is contained in:
parent
34642e1dd5
commit
8e5f097c5e
2 changed files with 154 additions and 1 deletions
133
tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py
Normal file
133
tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py
Normal file
|
|
@ -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"]}
|
||||
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue