All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / changes (pull_request) Successful in 11s
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 5m3s
Правка ревью к PR #3372: комментарии и текст коммита обещали восстановление потери, которой нет. Факт: pipeline.py импортирует scraper_kit.orchestration.runs, и все четыре его писателя МЕРЖАТ jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS jsonb)`, runs.py:708/776/817/880), а _pick_resume наследует done_buckets при claim (scheduler.py:741-748, #3074) — чекпоинт domclick в БД не терялся. Перезаписывающий двойник app/services/scrape_runs.py этой функцией не вызывается. Поэтому done_buckets в каждом payload — единообразие и защита от смены писателя (если запись пойдёт через app-копию или kit-писатель станет перезаписывающим), а не спасение данных. Комментарии переписаны под этот факт. Дрейн-ветка и NoProxyAvailableError строили dict вручную — переведены на {**_payload(), "interrupted": 1} и _payload(), чтобы «один _payload() на функцию» было правдой. Тест: к cancel добавлены ассерты на финальный heartbeat и mark_done — оба payload'а несут унаследованное ∪ пройденное этим прогоном. Refs #3355, PR #3363. Closes #3369
190 lines
7.9 KiB
Python
190 lines
7.9 KiB
Python
"""Каждая запись 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"
|
||
)
|