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

Merged
bot-backend merged 8 commits from fix/kit-orchestration into main 2026-09-17 09:23:44 +00:00
Showing only changes of commit 69e7047398 - Show all commits

View file

@ -49,8 +49,14 @@ DETAIL_ROWS = [
class _FakeDb: 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 {} 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]] = [] self.heartbeats: list[dict[str, Any]] = []
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> 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)) return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters))
if params and "ids" in params: if params and "ids" in params:
# houses-фаза: дома, у которых есть cian_zhk_url. # 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: if params and "limit" in params:
# detail-фаза avito: карточки-кандидаты на обогащение. # detail-фаза avito: карточки-кандидаты на обогащение.
return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS)) return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS))
@ -114,16 +122,23 @@ def _config() -> types.SimpleNamespace:
async def _run( 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]: ) -> tuple[_FakeDb, Any]:
from scraper_kit.orchestration import pipeline as pl from scraper_kit.orchestration import pipeline as pl
_FakeScraper.visited = [] _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: исключения нет, # None — ровно тот отказ, что был на проде 06.09: исключения нет,
# счётчик houses_failed растёт, фаза возвращается штатно. # счётчик 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 return object() if houses_ok else None
with ( 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: class _FakeAvitoScraper:
"""Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза.""" """Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза."""