Merge pull request 'fix(scraper-kit): domclick city sweep — done_buckets во всех семи выходах, а не только в дрейне и одном бане' (#3372) from fix/3369-domclick-cancel-checkpoint into main
Some checks failed
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Successful in 11s
Deploy Trade-In / build-frontend (push) Has been skipped
Deploy Trade-In / build-browser (push) Has been skipped
Deploy Trade-In / test (push) Has been cancelled

This commit is contained in:
bot-backend 2026-09-05 20:19:44 +00:00
commit 1c8035934d
2 changed files with 227 additions and 26 deletions

View file

@ -0,0 +1,190 @@
"""Каждая запись counters domclick-свипа несёт чекпоинт `done_buckets` (#3369).
Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут
`done_buckets` в payload явно, а cancel-ветка, финальный heartbeat и финализаторы
писали голый `counters.to_dict()`. Это НЕ было потерей данных: writer
`scraper_kit.orchestration.runs`, все его писатели мержат jsonb
(`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`), а `_pick_resume`
наследует ключ при claim (#3074). Инвариант ставится ради единообразия и
независимости веток от того, какой писатель окажется на другом конце
(перезаписывающий двойник `app/services/scrape_runs.py` существует).
Проверка по ЗНАЧЕНИЮ, а не по факту записи: чекпоинт обязан дословно лежать в
payload'е отмены, финального heartbeat'а и `mark_done`.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url — как в
# test_3355_drain_mark_full_loads.py, до остальных импортов.
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
class _FakeDb:
"""Запоминает каждый UPDATE с counters; SELECT по `rid` отдаёт чекпоинт."""
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
self.writes: list[tuple[str, dict[str, Any]]] = []
self._prev_counters = prev_counters
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
self.writes.append((str(stmt), json.loads(params["counters"])))
return MagicMock()
if params and "rid" in params:
row = types.SimpleNamespace(counters=self._prev_counters)
return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None})
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
class _NeverCalledScraper:
"""Отмена срабатывает ДО скрапера: обращение сюда — сломанный порядок проверок."""
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
async def __aenter__(self) -> _NeverCalledScraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
def __getattr__(self, name: str) -> Any:
raise AssertionError(f"cancel-ветка не сработала: тронут скрапер ({name})")
def _config() -> types.SimpleNamespace:
return types.SimpleNamespace(
scraper_fetch_mode="cffi",
scraper_proxy_url=None,
scraper_skip_seen_today=False,
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
avito_serp_ok_not_banned=True,
cian_full_load_per_fetch_timeout_s=0.0,
)
@pytest.mark.asyncio
async def test_domclick_sweep_cancel_keeps_inherited_checkpoint() -> None:
"""Отмена до SERP-фазы: payload несёт ровно унаследованные две корзины."""
from scraper_kit.orchestration import pipeline as pl
inherited = ["1:0:5000000", "2:0:5000000"]
db = _FakeDb(prev_counters={"done_buckets": inherited})
with (
patch.object(pl, "DomClickScraper", _NeverCalledScraper),
patch.object(pl.runs, "is_cancelled", lambda *_a: True),
):
await pl.run_domclick_city_sweep(
db, # type: ignore[arg-type]
run_id=3369,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: False,
resume_run_id=3368,
)
assert db.writes, "cancel не оставил ни одной записи counters"
_sql, counters = db.writes[-1]
assert counters.get("done_buckets") == inherited, (
f"cancel: payload без чекпоинта ({counters.get('done_buckets')!r} вместо "
f"{inherited!r}) — ветка зависит от того, мержит ли писатель jsonb"
)
@pytest.mark.asyncio
async def test_domclick_sweep_cancel_marker_is_not_resume_status_only() -> None:
"""Отмена без чекпоинта в предшественнике не выдумывает корзин (пустой список)."""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
with (
patch.object(pl, "DomClickScraper", _NeverCalledScraper),
patch.object(pl.runs, "is_cancelled", lambda *_a: True),
):
await pl.run_domclick_city_sweep(
db, # type: ignore[arg-type]
run_id=3369,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: False,
)
assert db.writes, "cancel не оставил ни одной записи counters"
_sql, counters = db.writes[-1]
assert counters.get("done_buckets") == [], (
f"cancel без предшественника: ожидался пустой чекпоинт, получено "
f"{counters.get('done_buckets')!r}"
)
class _CleanSweepScraper:
"""Свип прошёл все корзины без блока: ветка честного статуса → mark_done."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self.blocked = False
self.geo_filtered = 0
self.fetch_errors = 0
self.buckets_completed = 1
self.buckets_total = 1
self.completed_buckets = ["3:0:5000000"]
self.bucket_start_index = 0
self.request_delay_sec = 0.0
async def __aenter__(self) -> _CleanSweepScraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
async def fetch_city(self, **_kw: Any) -> list[Any]:
return []
@pytest.mark.asyncio
async def test_domclick_sweep_final_heartbeat_and_mark_done_carry_checkpoint() -> None:
"""Финальный heartbeat и mark_done несут унаследованное пройденное этим прогоном."""
from scraper_kit.orchestration import pipeline as pl
inherited = ["1:0:5000000", "2:0:5000000"]
expected = sorted([*inherited, "3:0:5000000"])
db = _FakeDb(prev_counters={"done_buckets": inherited})
with (
patch.object(pl, "DomClickScraper", _CleanSweepScraper),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_domclick_city_sweep(
db, # type: ignore[arg-type]
run_id=3369,
config=_config(),
matcher=MagicMock(),
request_delay_sec=0.0,
shutdown_requested=lambda: False,
resume_run_id=3368,
)
assert len(db.writes) >= 2, f"ожидались heartbeat + финализатор, есть {len(db.writes)}"
sql_done, counters_done = db.writes[-1]
assert "finished_at" in sql_done, f"последняя запись — не финализатор: {sql_done!r}"
_sql_hb, counters_hb = db.writes[-2]
for label, payload in (("финальный heartbeat", counters_hb), ("mark_done", counters_done)):
assert payload.get("done_buckets") == expected, (
f"{label}: payload без чекпоинта ({payload.get('done_buckets')!r} вместо "
f"{expected!r}) — ветка зависит от того, мержит ли писатель jsonb"
)

View file

@ -4426,6 +4426,21 @@ async def run_domclick_city_sweep(
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
_scraper_ref: list[DomClickScraper] = []
# #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции —
# это единообразие и защита от смены писателя, а НЕ восстановление потери.
# Как есть сегодня: пишет kit-овый scraper_kit.orchestration.runs (импорт выше),
# и все четыре его писателя МЕРЖАТ jsonb — `counters = COALESCE(counters,'{}')
# || CAST(:counters AS jsonb)`, так что ключ, записанный раньше, переживает
# payload без него; плюс _pick_resume наследует done_buckets при claim (#3074).
# То есть чекпоинт в БД не терялся. Перезаписывающий двойник существует —
# app/services/scrape_runs.py (`counters = CAST(:counters AS jsonb)`), — но
# этой функцией не вызывается. Полный payload делает ветки нечувствительными
# к тому, какой из двух писателей окажется на другом конце.
_checkpoint: list[str] = []
def _payload() -> dict[str, Any]:
return {**counters.to_dict(), "done_buckets": _checkpoint}
try:
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
@ -4451,28 +4466,24 @@ async def run_domclick_city_sweep(
len(skip_buckets),
)
_checkpoint = sorted(skip_buckets)
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
if runs.is_cancelled(db, run_id):
logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id)
runs.update_heartbeat(db, run_id, counters.to_dict())
# #3369: 'cancelled' ∈ _RESUME_STATUSES, прогон резюмируем — пишем
# полный payload, как и все прочие ветки (см. _payload выше).
runs.update_heartbeat(db, run_id, _payload())
return counters
elif shutdown_requested():
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
# минуя honest-status (это чистый drain, не QRATOR-блок).
logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id)
# #3355: метка дрейна (см. city-свипы #3330/#3333) И явный перенос
# чекпоинта. Никакого jsonb-мержа тут нет: domclick пишет counters
# через app-level scrape_runs (update_heartbeat/mark_done делают
# `counters = CAST(:counters AS jsonb)` — ПОЛНАЯ перезапись), а
# scheduler при claim done_buckets не наследует. Без явного
# done_buckets дрейн-прогон закрылся бы с пустым чекпоинтом и резюм
# пересобирал бы все шесть корзин заново. Тот же приём, что в
# ban-ветке ниже (except NoProxyAvailableError).
_drain = {
**counters.to_dict(),
"done_buckets": sorted(skip_buckets),
"interrupted": 1,
}
# #3355: метка дрейна (см. city-свипы #3330/#3333) поверх обычного
# payload'а. Чекпоинт тут не «спасается»: kit-писатель мержит jsonb
# (см. _payload выше) — done_buckets в payload'е нужен для
# единообразия, а interrupted=1 отличает дрейн от честного done.
_drain = {**_payload(), "interrupted": 1}
runs.update_heartbeat(db, run_id, _drain)
runs.mark_done(db, run_id, _drain)
return counters
@ -4557,16 +4568,16 @@ async def run_domclick_city_sweep(
# хотя площадка вообще не была затронута. Тот же диагноз и тот же по духу
# обработчик, что у run_avito_full_load/run_cian_full_load/run_yandex_full_load
# (см. except NoProxyAvailableError там же) — ban_kind_of_exception() относит тип
# исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). done_buckets
# сохраняем как унаследованные skip_buckets — этот прогон новых не завершил, но
# терять уже собранный чекпоинт (#2687-класс дефекта) resume не должен.
# исключения к BAN_KIND_INFRA («наша инфраструктура», не площадка). В payload'е
# done_buckets == унаследованные skip_buckets: этот прогон новых корзин не
# завершил (_checkpoint обновляется только после SERP-фазы).
logger.error("domclick-sweep run_id=%d: no proxy available — %s", run_id, exc)
counters.errors_count += 1
runs.mark_banned(
db,
run_id,
f"domclick sweep aborted: {exc}",
{**counters.to_dict(), "done_buckets": sorted(skip_buckets)},
_payload(),
ban_kind=ban_kind_of_exception(exc),
)
return counters
@ -4588,13 +4599,13 @@ async def run_domclick_city_sweep(
# #3118: чекпоинт = унаследованное завершённое в этом прогоне.
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
# затирают, и оборванный болезнью финализации прогон его не теряет.
_done_now = sorted(skip_buckets | set(_s.completed_buckets))
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": _done_now})
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
runs.update_heartbeat(db, run_id, _payload())
counters.bucket_start_index = _s.bucket_start_index
# pages_fetched: worst-case число страниц (buckets × pages cap).
counters.pages_fetched = _num_fetches
runs.update_heartbeat(db, run_id, counters.to_dict())
runs.update_heartbeat(db, run_id, _payload())
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
@ -4621,7 +4632,7 @@ async def run_domclick_city_sweep(
run_id,
f"QRATOR block aborted sweep — {counters.lots_fetched} listings "
"collected before abort (#2657)",
counters.to_dict(),
_payload(),
# #2687: эта ветка входится ТОЛЬКО при распознанном QRATOR-блоке
# (counters.blocked == 1, DomClickBlockedError — площадка показала
# block-страницу). Это ровно определение BAN_KIND_PLATFORM
@ -4651,7 +4662,7 @@ async def run_domclick_city_sweep(
f"sweep оборван: пройдено {counters.buckets_completed} комнатных "
f"бакетов из {counters.buckets_total}, собрано "
f"{counters.lots_fetched} лотов (#2670)",
counters.to_dict(),
_payload(),
)
elif counters.lots_fetched == 0 and counters.errors_count > 0:
logger.error(
@ -4663,10 +4674,10 @@ async def run_domclick_city_sweep(
db,
run_id,
"fetch errors — 0 listings",
counters.to_dict(),
_payload(),
)
else:
runs.mark_done(db, run_id, counters.to_dict())
runs.mark_done(db, run_id, _payload())
logger.info(
"domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) "
@ -4686,5 +4697,5 @@ async def run_domclick_city_sweep(
except Exception as exc:
logger.exception("domclick-sweep run_id=%d: fatal error", run_id)
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
runs.mark_failed(db, run_id, str(exc), _payload())
raise