fix(tradein/scraper-kit): бакет в done_buckets только если каждая фаза дала удачу (#3415)

cian_city_sweep 6179 (06.09) — `failed` по honest-status «фаза 'houses' отказала
полностью — 30 из 30», но все пять якорей записаны в counters.done_buckets.
Прогон 6276 (07.09) унаследовал точку, пропустил все якоря и отчитался `done` с
houses_attempted=0, lots_fetched=0 — зелёный по построению: фаза не исполнялась.
Из-за этого фикс #3396 (houses через пул прокси) простоял на проде сутки, ни разу
не отработав.

Корень: отметка «якорь пройден» опиралась на поток управления (исключение/таймаут
→ continue), а фаза отказывает и БЕЗ исключения — fetch_newbuilding возвращает
None, растёт только счётчик. _bucket_phase_totally_failed сверяет прирост пары
X_attempted/X_failed за якорь (та же конвенция, что у runs._phase_totally_failed,
но порог одна попытка, а не три: здесь решается «собирать ли заново», и цена
ошибок несимметрична).

Форма чекпоинта не менялась — плоский список имён якорей: его читают ещё три
свипа и scheduler._resume_decision, а ключ `bucket:phase` потребовал бы новых
читателей ради того же решения; резюм со старой точкой работает как прежде.

Тот же гейт — в run_avito_city_sweep (та же схема, отдельная функция; прод-следа
за 45 суток у него нет, у cian — 4 прогона). Ветка «houses DB query failed» теперь
растит houses_attempted вместе с houses_failed: пара attempted=0/failed=N не
видна ни run-level правилу, ни гейту бакета.

Соседи без правки: yandex combo-чекпоинт (единица пишется SERP-слоем до фазы,
пробел покрывает yandex_address_backfill, 0 прогонов на проде), domclick (одна
фаза), full-load'ы (бакет = SERP-страница, detail тянется из БД, не из бакета).

Closes #3415
This commit is contained in:
bot-backend 2026-09-07 16:47:08 +05:00
parent 278f8055a4
commit 19a335f4de
2 changed files with 363 additions and 2 deletions

View file

@ -0,0 +1,271 @@
"""Бакет уходит в чекпоинт только если ВСЕ его фазы что-то дали (#3415).
Прод-факт. `cian_city_sweep` 6179 (06.09) статус `failed` с причиной
«phase-honest-status: фаза 'houses' отказала полностью 30 из 30 попыток
неудачны, обогащено 0», а `counters.done_buckets` при этом несёт ВСЕ пять якорей
(«Академический», «Пионерский», «Уралмаш», «Центр», «ЮЗ»). Следующий прогон 6276
(07.09) резюмировался от него (`resume_from: 6179`, `resume_buckets: 5`),
пропустил все якоря и отчитался `done` с `errors_count=0`, `houses_attempted=0`,
`lots_fetched=0` зелёный по построению: фаза не исполнялась. Из-за этого фикс
#3396 (houses через пул прокси) простоял на проде с 06.09 и ни разу не отработал.
ИНВАРИАНТ. Отказ ЦЕЛОЙ фазы внутри якоря это «якорь пройден не был», ровно как
исключение или таймаут (#3074/#3319). Разница только в том, что отказ фазы виден
не потоком управления, а бухгалтерией счётчиков: `houses_attempted` попыток,
`houses_failed` отказов, обогащено ноль. Порог здесь одна попытка, а не три,
как у run-level `_phase_totally_failed` (#2700): тот СТАВИТ ДИАГНОЗ прогону и
обязан отсеивать шум, а этот решает «собирать ли якорь заново», и цена ошибок
несимметрична лишний повтор якоря стоит нескольких запросов, а ложная отметка
«пройден» теряет часть города навсегда и молча.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url.
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost:5432/test")
import json
import types
from typing import Any
from unittest.mock import MagicMock, patch
import pytest
ANCHOR_A = (56.83, 60.60, "ekb-center")
ANCHOR_B = (56.79, 60.63, "ekb-south")
# Два дома на якорь — НИЖЕ порога run-level правила (_PHASE_MIN_ATTEMPTS = 3):
# гейт бакета обязан сработать на своей, более строгой границе.
HOUSE_ROWS = [
{"id": 101, "cian_zhk_url": "https://cian.ru/zhk/101"},
{"id": 102, "cian_zhk_url": "https://cian.ru/zhk/102"},
]
DETAIL_ROWS = [
{"id": 201, "source_url": "/ekb/kvartiry/201"},
{"id": 202, "source_url": "/ekb/kvartiry/202"},
]
class _FakeDb:
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
self.prev_counters = prev_counters or {}
self.heartbeats: list[dict[str, Any]] = []
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
self.heartbeats.append(json.loads(params["counters"]))
return MagicMock()
if params and "rid" in params:
return MagicMock(fetchone=lambda: types.SimpleNamespace(counters=self.prev_counters))
if params and "ids" in params:
# houses-фаза: дома, у которых есть cian_zhk_url.
return MagicMock(mappings=lambda: MagicMock(all=lambda: HOUSE_ROWS))
if params and "limit" in params:
# detail-фаза avito: карточки-кандидаты на обогащение.
return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS))
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
def _lot() -> types.SimpleNamespace:
"""Новостроечный лот с привязкой к ЖК — только такой доходит до houses-фазы."""
return types.SimpleNamespace(
listing_segment="novostroyki",
house_source="cian_newbuilding",
house_ext_id="nb-1",
)
class _FakeScraper:
"""Двойник CianScraper: помнит, за какими якорями реально ходили."""
visited: list[str] = [] # noqa: RUF012 — тестовый сборник
def __init__(self, *_a: Any, **_kw: Any) -> None:
self.state_extraction_attempts = 1
self.state_extraction_failures = 0
self.request_delay_sec = 0.0
async def __aenter__(self) -> _FakeScraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
async def fetch_around_multi_room(
self, lat: float, lon: float, *_a: Any, **_kw: Any
) -> list[types.SimpleNamespace]:
_FakeScraper.visited.append(f"{lat},{lon}")
return [_lot()]
def _config() -> types.SimpleNamespace:
return types.SimpleNamespace(
scraper_proxy_url=None,
scraper_fetch_mode="cffi",
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
)
async def _run(
prev: dict[str, Any] | None, *, houses_ok: bool, run_id: int = 8401
) -> tuple[_FakeDb, Any]:
from scraper_kit.orchestration import pipeline as pl
_FakeScraper.visited = []
db = _FakeDb(prev)
async def _fake_newbuilding(*_a: Any, **_kw: Any) -> Any:
# None — ровно тот отказ, что был на проде 06.09: исключения нет,
# счётчик houses_failed растёт, фаза возвращается штатно.
return object() if houses_ok else None
with (
patch.object(pl, "CianScraper", _FakeScraper),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl, "fetch_newbuilding", _fake_newbuilding),
patch.object(pl, "save_newbuilding_enrichment", lambda *_a, **_kw: None),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_cian_city_sweep(
db, # type: ignore[arg-type]
run_id=run_id,
config=_config(),
matcher=MagicMock(),
anchors=[ANCHOR_A, ANCHOR_B],
enrich_houses=True,
detail_top_n=0,
request_delay_sec=0.0,
resume_run_id=(run_id - 1) if prev is not None else None,
)
return db, counters
def _last_checkpoint(db: _FakeDb) -> list[str]:
with_ckpt = [hb for hb in db.heartbeats if "done_buckets" in hb]
assert with_ckpt, "ни один heartbeat не унёс done_buckets — чекпоинт не персистится"
return with_ckpt[-1]["done_buckets"]
@pytest.mark.asyncio
async def test_anchor_with_totally_failed_houses_is_not_checkpointed() -> None:
"""(а) Якорь, где houses отказала на ВСЕХ домах, не попадает в чекпоинт."""
db, counters = await _run(None, houses_ok=False)
assert counters.houses_attempted > 0, "houses-фаза вообще не исполнялась — тест не о том"
assert counters.houses_enriched == 0, "houses-фаза что-то обогатила — отказ не полный"
assert _last_checkpoint(db) == [], (
"якорь с полностью отказавшей houses-фазой помечен пройденным — "
f"чекпоинт {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_next_run_re_enters_houses() -> None:
"""(б) Следующий прогон с этим чекпоинтом снова заходит в houses."""
db1, _ = await _run(None, houses_ok=False)
_db2, counters2 = await _run(
{"done_buckets": _last_checkpoint(db1)}, houses_ok=True, run_id=8402
)
assert counters2.houses_attempted > 0, (
"прогон-наследник пропустил houses-фазу: чекпоинт предшественника "
"объявил якоря пройденными, хотя фаза в них не дала ничего"
)
assert counters2.houses_enriched > 0, "houses-фаза исполнилась, но ничего не обогатила"
@pytest.mark.asyncio
async def test_healthy_anchor_is_still_checkpointed() -> None:
"""(в) Контроль: якорь, где все фазы отработали, помечается как раньше."""
db, counters = await _run(None, houses_ok=True)
assert counters.houses_enriched > 0
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], (
"перестали помечать пройденные якоря вовсе — резюм (#3074) сломан"
)
class _FakeAvitoScraper:
"""Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
async def fetch_around(self, *_a: Any, **_kw: Any) -> list[types.SimpleNamespace]:
return [types.SimpleNamespace(house_url=None)]
class _FakeAsyncSession:
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
async def __aenter__(self) -> _FakeAsyncSession:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
@pytest.mark.asyncio
async def test_avito_anchor_with_totally_failed_detail_is_not_checkpointed() -> None:
"""Сосед: у avito-свипа тот же гейт — фаза без единой удачи не даёт отметки.
Кода общего у свипов нет (две отдельные функции), общий только гейт. Прод-следа
у avito за 45 суток нет (0 прогонов против 4 у циана) тест сторожит вторую
копию проводки, а не найденный в данных дефект.
"""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
async def _boom(*_a: Any, **_kw: Any) -> Any:
raise RuntimeError("detail отдал 403")
with (
patch.object(pl, "AvitoScraper", _FakeAvitoScraper),
patch.object(pl, "AsyncSession", _FakeAsyncSession),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl, "fetch_detail", _boom),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_avito_city_sweep(
db, # type: ignore[arg-type]
run_id=8404,
config=types.SimpleNamespace(
scraper_fetch_mode="cffi",
scraper_proxy_url=None,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
avito_serp_ok_not_banned=True,
),
matcher=MagicMock(),
enrichment=MagicMock(),
anchors=[ANCHOR_A],
enrich_houses=False,
enrich_imv=False,
detail_top_n=len(DETAIL_ROWS),
request_delay_sec=0.0,
)
assert counters.detail_attempted == len(DETAIL_ROWS)
assert counters.detail_enriched == 0
assert _last_checkpoint(db) == [], (
f"якорь с полностью отказавшей detail-фазой помечен пройденным: {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_old_format_checkpoint_still_resumes() -> None:
"""(г) Чекпоинт СТАРОГО формата (плоский список имён якорей) читается как раньше."""
db, _ = await _run({"done_buckets": ["ekb-center"]}, houses_ok=True, run_id=8403)
assert f"{ANCHOR_A[0]},{ANCHOR_A[1]}" not in _FakeScraper.visited, (
"якорь из унаследованного чекпоинта всё-таки опрашивали"
)
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], (
"унаследованный якорь потерян — формат чекпоинта разъехался со старыми прогонами"
)

View file

@ -36,7 +36,7 @@ from __future__ import annotations
import asyncio
import logging
import random
from collections.abc import Callable
from collections.abc import Callable, Mapping
from contextlib import AsyncExitStack
from dataclasses import dataclass, field, fields
from datetime import date, timedelta
@ -1100,6 +1100,43 @@ async def run_avito_pipeline(
await browser_fetcher.__aexit__(None, None, None)
def _bucket_phase_totally_failed(before: Mapping[str, int], after: Mapping[str, int]) -> str | None:
"""Фаза бакета, где отказала КАЖДАЯ попытка (#3415). Имя фазы или None.
Считает ПРИРОСТ счётчиков за один бакет (якорь): counters у свипа run-level,
так что «фаза этого якоря» существует только как разница до/после.
Прод-факт, ради которого гейт заведён: `cian_city_sweep` 6179 (06.09)
`houses_attempted=30, houses_failed=30, houses_enriched=0`, статус `failed` по
run-level правилу (`_phase_totally_failed`, #2700), а в `done_buckets` записаны
все пять якорей. Прогон 6276 (07.09) унаследовал точку, пропустил все якоря и
отчитался `done` с `houses_attempted=0` зелёный по построению, потому что фаза
не исполнялась. Фикс #3396 (houses через пул прокси) простоял на проде сутки, ни
разу не отработав.
Пары ищутся В САМИХ counters (ключ `X_attempted` со спутником `X_failed`) тот же
приём и та же причина, что у `runs._phase_totally_failed`: зашитый список фаз это
ровно то место, куда забывают дописать новую. Фазы без счётчика попыток гейт НЕ
видит у avito houses роль attempted играет `unique_houses` (нет `houses_attempted`),
и притягивать её сюда значило бы блокировать отметку якоря из-за ОДНОГО упавшего
дома из десяти.
Порог одна попытка, а не три, как у run-level правила: тот ставит прогону
диагноз и обязан отсеивать шум, а этот решает «собирать ли якорь заново». Цена
ошибок несимметрична: лишний повтор якоря стоит нескольких запросов, ложная
отметка «пройден» теряет часть города навсегда и молча (прогон-то `done`).
"""
for key in sorted(after):
if not key.endswith("_attempted"):
continue
phase = key[: -len("_attempted")]
attempted = after[key] - before.get(key, 0)
failed = after.get(f"{phase}_failed", 0) - before.get(f"{phase}_failed", 0)
if attempted >= 1 and failed >= attempted:
return phase
return None
@dataclass
class CitySweepCounters:
"""Aggregate counters для full city sweep run."""
@ -1349,6 +1386,12 @@ async def run_avito_city_sweep(
# #3074: сбрасывается на КАЖДОЙ итерации — иначе один упавший якорь
# заразил бы все последующие, и чекпоинт не пополнялся бы вовсе.
_anchor_ok = True
# #3415: снимок счётчиков ДО якоря — по разнице видно, дала ли
# каждая фаза ЭТОГО якоря хоть одну удачную попытку (см.
# _bucket_phase_totally_failed). Тот же дефект, что у cian-свипа:
# detail-фаза, отказавшая на всех карточках без исключения, доходила
# до отметки «якорь пройден» наравне с целым якорем.
_before_anchor = counters.to_dict()
# Capture loop variables in default args (B023): prevents stale binding
# if the coroutine is scheduled after the loop variable changes.
@ -1854,6 +1897,22 @@ async def run_avito_city_sweep(
# его как пройденный значило бы, что следующий прогон его пропустит и
# объявления оттуда не соберутся НИКОГДА, причём молча: прогон
# завершится штатно. Тот же инвариант, что у combo в yandex-свипе.
#
# #3415: флага мало — фаза отказывает и БЕЗ исключения (каждая
# detail-карточка отдала 403, счётчик вырос, корутина завершилась
# штатно). Замер на проде за 45 суток: у avito-свипов такой якорь
# пока не встречался (0 прогонов), у cian — 4; гейт стоит здесь
# потому, что дефект один и тот же, а не по следу в данных.
_failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict())
if _failed_phase is not None:
logger.warning(
"city-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' "
"не дала ни одной удачной попытки; следующий прогон соберёт его заново",
run_id,
name,
_failed_phase,
)
_anchor_ok = False
if _anchor_ok:
_done_anchors.add(name)
runs.update_heartbeat(db, run_id, _ckpt())
@ -3028,6 +3087,10 @@ async def run_cian_city_sweep(
)
# ── Per-anchor phases под watchdog-таймаутом ────────────────────
# #3415: counters у свипа run-level — снимок ДО якоря даёт бухгалтерию
# его собственных фаз (разницей), по которой ниже решается, считать ли
# якорь пройденным.
_before_anchor = counters.to_dict()
anchor_lots: list[ScrapedLot] = []
_c_lat, _c_lon, _c_name = lat, lon, name
@ -3243,6 +3306,12 @@ async def run_cian_city_sweep(
len(nb_id_list),
exc,
)
# #3415: попытки считаем ВМЕСТЕ с отказами. Раньше здесь рос
# только `houses_failed`, и пара получалась несуществующей
# (attempted=0, failed=N): и run-level `_phase_totally_failed`
# (#2700), и гейт бакета сверяют `failed` с `attempted`, так что
# полный отказ фазы на этой ветке не видел ни один из них.
counters.houses_attempted += len(nb_id_list)
counters.houses_failed += len(nb_id_list)
try:
db.rollback()
@ -3434,7 +3503,28 @@ async def run_cian_city_sweep(
# Записать упавший якорь пройденным значило бы, что следующий прогон
# пропустит его навсегда — молча, потому что прогон завершится
# штатно, просто часть города не соберётся.
_done_anchors.add(name)
#
# #3415: потока управления мало. Фаза внутри якоря отказывает БЕЗ
# исключения — `fetch_newbuilding` возвращает None, счётчик растёт,
# `_cian_anchor_phases` завершается штатно, — и такой якорь доходил
# сюда наравне с целым. Дальше ровно тот же исход, что и у записанного
# упавшего: следующий прогон пропускает якорь и рапортует `done` с
# нулевой фазой. Форма отметки прежняя (плоский список имён): чекпоинт
# читают ещё три свипа и `scheduler._resume_decision`, а ключ
# `bucket:phase` потребовал бы новых читателей ради того же решения.
_failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict())
if _failed_phase is None:
_done_anchors.add(name)
else:
logger.warning(
"cian-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' "
"не дала ни одной удачной попытки; следующий прогон соберёт его заново",
run_id,
name,
_failed_phase,
)
# Heartbeat пишем в любом случае: без него reap_zombies посчитает живой
# прогон мёртвым, а jsonb-мерж сохранит унаследованную часть точки.
runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)}
)