fix(domclick): чекпоинт done_buckets во всех payload'ах свипа
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m5s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m5s
cancel-ветка run_domclick_city_sweep писала голый counters.to_dict(); писатель — app-level scrape_runs с полной перезаписью (CAST(:counters AS jsonb)), а 'cancelled' входит в _RESUME_STATUSES → отменённый прогон закрывался с пустым чекпоинтом и резюм пересобирал все корзины. Те же потери были у финального heartbeat и у mark_banned/mark_failed/mark_done (banned/failed тоже резюмируемы): они затирали чекпоинт, записанный после SERP-фазы. Один _payload() на функцию — done_buckets едет в каждой записи counters. Refs #3355, PR #3363. Closes #3369
This commit is contained in:
parent
7afaa12d75
commit
b796702d9a
2 changed files with 156 additions and 9 deletions
|
|
@ -0,0 +1,133 @@
|
||||||
|
"""Cancel-ветка domclick-свипа обязана нести чекпоинт `done_buckets` (#3369).
|
||||||
|
|
||||||
|
Доводка #3355/PR #3363: дрейн- и бан-ветки `run_domclick_city_sweep` кладут
|
||||||
|
`done_buckets` в свой payload явно, а cancel-ветка писала голый
|
||||||
|
`counters.to_dict()`. Мержа jsonb тут нет: писатель — app-level
|
||||||
|
`app/services/scrape_runs.py` (`counters = CAST(:counters AS jsonb)`, ПОЛНАЯ
|
||||||
|
перезапись), а scheduler при claim `done_buckets` не наследует. `cancelled`
|
||||||
|
входит в `_RESUME_STATUSES` → отменённый прогон закрывался с пустым чекпоинтом,
|
||||||
|
и резюм пересобирал все шесть корзин заново.
|
||||||
|
|
||||||
|
Проверка по ЗНАЧЕНИЮ (не по факту записи): унаследованные две корзины обязаны
|
||||||
|
дословно оказаться в payload'е отмены.
|
||||||
|
"""
|
||||||
|
|
||||||
|
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: чекпоинт потерян ({counters.get('done_buckets')!r} вместо "
|
||||||
|
f"{inherited!r}) — резюм пересоберёт уже собранные корзины"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@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}"
|
||||||
|
)
|
||||||
|
|
@ -4426,6 +4426,16 @@ async def run_domclick_city_sweep(
|
||||||
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
# Мутируемый контейнер для захвата scraper-ссылки из замыкания.
|
||||||
_scraper_ref: list[DomClickScraper] = []
|
_scraper_ref: list[DomClickScraper] = []
|
||||||
|
|
||||||
|
# #3369: чекпоинт (done_buckets) обязан ехать в КАЖДОЙ записи counters этой
|
||||||
|
# функции. Писатель — app-level scrape_runs: `counters = CAST(:counters AS jsonb)`,
|
||||||
|
# полная перезапись без jsonb-мержа, а scheduler при claim done_buckets не
|
||||||
|
# наследует. Любая запись без ключа стирает чекпоинт, и резюм (cancelled/banned/
|
||||||
|
# failed ∈ _RESUME_STATUSES) пересобирает уже собранные корзины заново.
|
||||||
|
_checkpoint: list[str] = []
|
||||||
|
|
||||||
|
def _payload() -> dict[str, Any]:
|
||||||
|
return {**counters.to_dict(), "done_buckets": _checkpoint}
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
|
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
|
||||||
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
|
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
|
||||||
|
|
@ -4451,10 +4461,14 @@ async def run_domclick_city_sweep(
|
||||||
len(skip_buckets),
|
len(skip_buckets),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
_checkpoint = sorted(skip_buckets)
|
||||||
|
|
||||||
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
|
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
|
||||||
if runs.is_cancelled(db, run_id):
|
if runs.is_cancelled(db, run_id):
|
||||||
logger.info("domclick-sweep run_id=%d: cancelled before SERP", 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 без done_buckets стирал
|
||||||
|
# унаследованный чекпоинт, и резюм пересобирал все корзины заново.
|
||||||
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
return counters
|
return counters
|
||||||
elif shutdown_requested():
|
elif shutdown_requested():
|
||||||
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
|
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
|
||||||
|
|
@ -4588,13 +4602,13 @@ async def run_domclick_city_sweep(
|
||||||
# #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне.
|
# #3118: чекпоинт = унаследованное ∪ завершённое в этом прогоне.
|
||||||
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
|
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
|
||||||
# затирают, и оборванный болезнью финализации прогон его не теряет.
|
# затирают, и оборванный болезнью финализации прогон его не теряет.
|
||||||
_done_now = sorted(skip_buckets | set(_s.completed_buckets))
|
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
|
||||||
runs.update_heartbeat(db, run_id, {**counters.to_dict(), "done_buckets": _done_now})
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
counters.bucket_start_index = _s.bucket_start_index
|
counters.bucket_start_index = _s.bucket_start_index
|
||||||
|
|
||||||
# pages_fetched: worst-case число страниц (buckets × pages cap).
|
# pages_fetched: worst-case число страниц (buckets × pages cap).
|
||||||
counters.pages_fetched = _num_fetches
|
counters.pages_fetched = _num_fetches
|
||||||
runs.update_heartbeat(db, run_id, counters.to_dict())
|
runs.update_heartbeat(db, run_id, _payload())
|
||||||
|
|
||||||
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
# ── ЧЕСТНЫЙ СТАТУС (#1968, ужесточён #2657) ───────────────────────────
|
||||||
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
# Распознанный QRATOR-блок НИКОГДА не даёт done. Домклик тут структурно
|
||||||
|
|
@ -4621,7 +4635,7 @@ async def run_domclick_city_sweep(
|
||||||
run_id,
|
run_id,
|
||||||
f"QRATOR block aborted sweep — {counters.lots_fetched} listings "
|
f"QRATOR block aborted sweep — {counters.lots_fetched} listings "
|
||||||
"collected before abort (#2657)",
|
"collected before abort (#2657)",
|
||||||
counters.to_dict(),
|
_payload(),
|
||||||
# #2687: эта ветка входится ТОЛЬКО при распознанном QRATOR-блоке
|
# #2687: эта ветка входится ТОЛЬКО при распознанном QRATOR-блоке
|
||||||
# (counters.blocked == 1, DomClickBlockedError — площадка показала
|
# (counters.blocked == 1, DomClickBlockedError — площадка показала
|
||||||
# block-страницу). Это ровно определение BAN_KIND_PLATFORM
|
# block-страницу). Это ровно определение BAN_KIND_PLATFORM
|
||||||
|
|
@ -4651,7 +4665,7 @@ async def run_domclick_city_sweep(
|
||||||
f"sweep оборван: пройдено {counters.buckets_completed} комнатных "
|
f"sweep оборван: пройдено {counters.buckets_completed} комнатных "
|
||||||
f"бакетов из {counters.buckets_total}, собрано "
|
f"бакетов из {counters.buckets_total}, собрано "
|
||||||
f"{counters.lots_fetched} лотов (#2670)",
|
f"{counters.lots_fetched} лотов (#2670)",
|
||||||
counters.to_dict(),
|
_payload(),
|
||||||
)
|
)
|
||||||
elif counters.lots_fetched == 0 and counters.errors_count > 0:
|
elif counters.lots_fetched == 0 and counters.errors_count > 0:
|
||||||
logger.error(
|
logger.error(
|
||||||
|
|
@ -4663,10 +4677,10 @@ async def run_domclick_city_sweep(
|
||||||
db,
|
db,
|
||||||
run_id,
|
run_id,
|
||||||
"fetch errors — 0 listings",
|
"fetch errors — 0 listings",
|
||||||
counters.to_dict(),
|
_payload(),
|
||||||
)
|
)
|
||||||
else:
|
else:
|
||||||
runs.mark_done(db, run_id, counters.to_dict())
|
runs.mark_done(db, run_id, _payload())
|
||||||
|
|
||||||
logger.info(
|
logger.info(
|
||||||
"domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) "
|
"domclick-sweep run_id=%d done: lots=%d (ins=%d/upd=%d) "
|
||||||
|
|
@ -4686,5 +4700,5 @@ async def run_domclick_city_sweep(
|
||||||
|
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.exception("domclick-sweep run_id=%d: fatal error", run_id)
|
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
|
raise
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue