All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m47s
Продолжение после yandex-свипа. Выбор источника — по замеру, а не «для
полноты»: за 60 дней avito_city_sweep дал 68 прогонов, 25 банов и 3 отмены
деплоем при среднем времени 15 мин и максимуме 97.
Прод-факт, который решает дело. Типичный итог свипа:
{"anchors_done": 1, "anchors_total": 5, ...,
"enrichment_abort_note": "detail enrichment aborted (Avito detail
firewall/soft-block ...)"}
Прогон срывается блокировкой на ПЕРВОМ из пяти якорей. Без чекпоинта
следующий прогон снова идёт в первый якорь, упирается в ту же стену, и якоря
2-5 не собираются никогда.
Ключ чекпоинта — ИМЯ якоря, а не индекс: состав списка зависит от city_slug,
позиция в нём между городами не устойчива. По той же причине здесь не нужна
гарда по числу якорей, которая есть у combo-чекпоинта яндекса: там ключ
якоря не содержал, здесь якорь и есть ключ.
Инвариант, ради которого отдельный флаг _anchor_ok: в чекпоинт попадает
только якорь, пройденный до конца. Ветка блокировки делает return и до записи
не доходит, а generic-except доходит — якорь упал, но цикл продолжается.
Записать такой якорь пройденным значило бы, что следующий прогон пропустит
его навсегда, причём молча: прогон завершится штатно, просто часть города не
соберётся. Флаг сбрасывается на каждой итерации, иначе один упавший якорь
заразил бы все последующие.
Пропущенный якорь двигает anchors_done — чтобы счётчик продолжал означать
«докуда дошли по списку», а не «сколько собрал именно этот прогон».
domclick_city_sweep намеренно НЕ трогаю: 54 прогона, ноль отмен деплоем,
среднее время 3 минуты — чекпоинт там не окупается. cian_city_sweep (среднее
35 мин, 2 отмены) — следующий шард.
Тесты (4) поведенческие: якорь из чекпоинта не опрашивается вовсе; пройденный
дописывается поверх унаследованных; без чекпоинта обходятся все; упавший в
чекпоинт НЕ попадает. Оговорка: на исходном коде они падают по сигнатуре
(unexpected keyword argument), то есть доказывают отсутствие параметра, а не
поведение — поведенческую часть держат сами проверки. Весь набор #3074 —
10 passed.
175 lines
7.4 KiB
Python
175 lines
7.4 KiB
Python
"""Чекпоинт по якорям для avito_city_sweep (#3074).
|
||
|
||
Прод-факт, из которого выросла задача. Типичный итог свипа:
|
||
|
||
{"anchors_done": 1, "anchors_total": 5, ..., "enrichment_abort_note":
|
||
"detail enrichment aborted (Avito detail firewall/soft-block ...)"}
|
||
|
||
То есть прогон срывается блокировкой на ПЕРВОМ из пяти якорей — 25 банов за
|
||
60 дней. Без чекпоинта следующий прогон снова начинает с первого якоря,
|
||
упирается в ту же стену, и якоря 2-5 не собираются никогда.
|
||
|
||
Ключ чекпоинта — ИМЯ якоря, а не его индекс: состав списка зависит от
|
||
`city_slug`, и позиция в нём не устойчива между городами.
|
||
|
||
ИНВАРИАНТ, РАДИ КОТОРОГО ТЕСТ. В чекпоинт попадает только якорь, пройденный до
|
||
конца. Ветка блокировки делает `return` и до записи не доходит, а вот
|
||
`except Exception` — доходит: якорь упал, но цикл продолжается. Записать такой
|
||
якорь как пройденный значило бы, что следующий прогон его пропустит и
|
||
объявления оттуда не соберутся НИКОГДА, причём молча — прогон завершится
|
||
штатно. Ровно та же граница, что у combo в yandex-свипе.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import os
|
||
|
||
# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем
|
||
# до остальных импортов — так же, как в test_3074_yandex_sweep_checkpoint.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
|
||
|
||
ANCHOR_A = (56.83, 60.60, "ekb-center")
|
||
ANCHOR_B = (56.79, 60.63, "ekb-south")
|
||
|
||
|
||
class _FakeDb:
|
||
"""Резюм-SELECT отдаёт counters предшественника; heartbeat'ы записываются."""
|
||
|
||
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)
|
||
)
|
||
return MagicMock()
|
||
|
||
def commit(self) -> None: ...
|
||
def rollback(self) -> 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
|
||
|
||
|
||
class _FakeScraper:
|
||
"""Двойник AvitoScraper: помнит, за какими якорями реально ходили."""
|
||
|
||
visited: list[tuple[float, float]] = [] # noqa: RUF012 — тестовый сборник
|
||
raise_on: tuple[float, float] | None = None
|
||
|
||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||
self._browser = None
|
||
self._cffi = None
|
||
|
||
async def fetch_around(self, lat: float, lon: float, *_a: Any, **_kw: Any) -> list:
|
||
_FakeScraper.visited.append((lat, lon))
|
||
if _FakeScraper.raise_on == (lat, lon):
|
||
raise RuntimeError("якорь упал по не-баново́й причине")
|
||
return []
|
||
|
||
|
||
def _config() -> types.SimpleNamespace:
|
||
return 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,
|
||
)
|
||
|
||
|
||
async def _run(prev: dict[str, Any] | None, raise_on: tuple[float, float] | None = None):
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
_FakeScraper.visited = []
|
||
_FakeScraper.raise_on = raise_on
|
||
db = _FakeDb(prev)
|
||
|
||
with (
|
||
patch.object(pl, "AvitoScraper", _FakeScraper),
|
||
patch.object(pl, "AsyncSession", _FakeAsyncSession),
|
||
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
|
||
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||
):
|
||
await pl.run_avito_city_sweep(
|
||
db, # type: ignore[arg-type]
|
||
run_id=7001,
|
||
config=_config(),
|
||
matcher=MagicMock(),
|
||
enrichment=MagicMock(),
|
||
anchors=[ANCHOR_A, ANCHOR_B],
|
||
enrich_houses=False,
|
||
enrich_imv=False,
|
||
detail_top_n=0,
|
||
resume_run_id=6999 if prev is not None else None,
|
||
)
|
||
return db
|
||
|
||
|
||
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_checkpointed_anchor_is_skipped_without_a_single_request() -> None:
|
||
"""Якорь из чекпоинта не опрашивается вовсе — ни одного обращения к источнику.
|
||
|
||
Ядро задачи: до фикса повторный прогон снова шёл в первый якорь и снова
|
||
получал там бан.
|
||
"""
|
||
await _run({"done_buckets": ["ekb-center"]})
|
||
|
||
assert (ANCHOR_A[0], ANCHOR_A[1]) not in _FakeScraper.visited, (
|
||
"якорь из чекпоинта всё-таки опрашивали"
|
||
)
|
||
assert (ANCHOR_B[0], ANCHOR_B[1]) in _FakeScraper.visited, "второй якорь не обошли"
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_checkpoint_accumulates_over_inherited() -> None:
|
||
"""Пройденный якорь дописывается поверх унаследованных, а не затирает их."""
|
||
db = await _run({"done_buckets": ["ekb-center"]})
|
||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_without_resume_all_anchors_are_visited() -> None:
|
||
"""Без чекпоинта поведение прежнее — обходятся все якоря."""
|
||
db = await _run(None)
|
||
assert len(_FakeScraper.visited) == 2
|
||
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"]
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_failed_anchor_does_not_enter_checkpoint() -> None:
|
||
"""Упавший якорь НЕ считается пройденным.
|
||
|
||
Иначе следующий прогон пропустит его навсегда, и это будет незаметно:
|
||
прогон завершается штатно, просто часть города не собирается никогда.
|
||
"""
|
||
db = await _run(None, raise_on=(ANCHOR_A[0], ANCHOR_A[1]))
|
||
|
||
ckpt = _last_checkpoint(db)
|
||
assert "ekb-center" not in ckpt, "упавший якорь попал в чекпоинт"
|
||
assert "ekb-south" in ckpt, "исправный якорь не зафиксирован"
|