All checks were successful
CI Trade-In / changes (pull_request) Successful in 8s
CI / changes (pull_request) Successful in 10s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m41s
Последний длинный свип без возобновления: при обрыве прогон начинался с первой страницы, а собранное терялось целиком — save_listings вызывался один раз на весь sweep. Единица возобновления — страница выдачи, по образцу якорей в city sweep. `_paginate_sweep`/`fetch_newbuildings` получили `start_page` (уже собранные страницы не запрашиваются) и колбэк `on_page`, который вызывается только после того, как страница пройдена до конца. Сохранение стало постраничным, номера пройденных страниц копятся в `scrape_runs.counters.done_buckets` мержем через `update_heartbeat`. Два инварианта, без которых фича вредна: 1. В чекпоинт попадает только страница, чьи лоты СОХРАНЕНЫ. Отказ save_listings перехвачен и прогон продолжается, но отметить такую страницу пройденной значило бы, что следующий прогон её пропустит и объявления оттуда не соберутся никогда — молча, потому что прогон завершится штатно. 2. Подхват начинается с ПЕРВОЙ несобранной страницы, а не с max+1. Дыра в чекпоинте возможна ровно из-за п.1, и max+1 перепрыгнул бы её навсегда. Страницы после дыры перечитаются — это дешевле потери и безопасно, повторная запись схлопывается по dedup_hash. Оба инварианта закрыты тестами, которые падают при их нарушении.
266 lines
11 KiB
Python
266 lines
11 KiB
Python
"""Чекпоинт по страницам для avito_newbuilding_sweep (#3074).
|
||
|
||
Последний длинный свип без чекпоинтов: 27.08 прогон убит деплоем на 343-й
|
||
минуте, собранное потеряно целиком — save_listings был ОДИН на весь sweep,
|
||
ждать в конце было нечего. Единица возобновления — СТРАНИЦА выдачи (у функции
|
||
уже есть `pages`; цикл по страницам живёт в `AvitoScraper._paginate_sweep`).
|
||
|
||
Ключ чекпоинта — номер страницы. В отличие от якорей/combo, страницы строго
|
||
последовательны (break-on-empty), поэтому resume — это `start_page =
|
||
max(done_buckets) + 1`, без skip-набора произвольных элементов.
|
||
|
||
ИНВАРИАНТЫ, РАДИ КОТОРЫХ ТЕСТ:
|
||
1. Резюм пропускает уже собранные страницы БЕЗ единого запроса к источнику.
|
||
2. В чекпоинт попадает только страница, пройденная до конца — оборванная
|
||
исключением/break страница НЕ фиксируется, иначе следующий прогон
|
||
пропустил бы её навсегда, причём молча (прогон завершится штатно).
|
||
3. Запись идёт мержем (`counters || :counters`), а не заменой — посторонний
|
||
ключ в `counters` текущего run'а переживает запись чекпоинта.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import os
|
||
|
||
# Settings собирается автофикстурой conftest'а и требует database_url. Выставляем
|
||
# до остальных импортов — так же, как в test_3074_avito_anchor_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
|
||
|
||
|
||
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 _MergingFakeDb(_FakeDb):
|
||
"""Как _FakeDb, но heartbeat-запись мержит в текущую строку (jsonb `||`),
|
||
а не просто копится списком — для теста инварианта #3 (посторонний ключ)."""
|
||
|
||
def __init__(self, prev_counters: dict[str, Any] | None = None) -> None:
|
||
super().__init__(prev_counters)
|
||
self.row: dict[str, Any] = {}
|
||
|
||
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
|
||
if params and "counters" in params:
|
||
incoming = json.loads(params["counters"])
|
||
self.row = {**self.row, **incoming}
|
||
self.heartbeats.append(dict(self.row))
|
||
return MagicMock()
|
||
return super().execute(_stmt, params)
|
||
|
||
|
||
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
|
||
|
||
|
||
def _lot(tag: str) -> MagicMock:
|
||
m = MagicMock()
|
||
m.source_id = tag
|
||
return m
|
||
|
||
|
||
class _FakeScraper:
|
||
"""Двойник AvitoScraper: помнит реально пройденные страницы, умеет упасть."""
|
||
|
||
visited_pages: list[int] = [] # noqa: RUF012 — тестовый сборник
|
||
fail_on: int | None = None
|
||
|
||
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||
self._browser = None
|
||
self._cffi = None
|
||
|
||
async def fetch_newbuildings(
|
||
self,
|
||
*,
|
||
pages: int,
|
||
start_page: int = 1,
|
||
on_page: Any = None,
|
||
delay_override_sec: float | None = None,
|
||
) -> list[Any]:
|
||
all_lots: list[Any] = []
|
||
for page in range(start_page, pages + 1):
|
||
_FakeScraper.visited_pages.append(page)
|
||
if _FakeScraper.fail_on == page:
|
||
# Страница оборвана (сеть/парсинг) — break, on_page НЕ зовём.
|
||
break
|
||
new_lots = [_lot(f"p{page}-{i}") for i in range(2)]
|
||
all_lots.extend(new_lots)
|
||
if on_page is not None:
|
||
on_page(page, new_lots)
|
||
return all_lots
|
||
|
||
|
||
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",
|
||
)
|
||
|
||
|
||
def _make_save(fail_on_call: int | None):
|
||
"""save_listings, падающий на N-м вызове — имитация отказа БД на одной странице."""
|
||
calls = {"n": 0}
|
||
|
||
def _save(*_a, **_kw):
|
||
calls["n"] += 1
|
||
if fail_on_call is not None and calls["n"] == fail_on_call:
|
||
raise RuntimeError(f"save_listings упал на странице {calls['n']}")
|
||
return (2, 0)
|
||
|
||
return _save
|
||
|
||
|
||
async def _run(
|
||
db: _FakeDb,
|
||
*,
|
||
resume_run_id: int | None,
|
||
pages: int = 4,
|
||
fail_on: int | None = None,
|
||
save_fail_on_call: int | None = None,
|
||
):
|
||
from scraper_kit.orchestration import pipeline as pl
|
||
|
||
_FakeScraper.visited_pages = []
|
||
_FakeScraper.fail_on = fail_on
|
||
|
||
with (
|
||
patch.object(pl, "AvitoScraper", _FakeScraper),
|
||
patch.object(pl, "AsyncSession", _FakeAsyncSession),
|
||
patch.object(pl, "save_listings", _make_save(save_fail_on_call)),
|
||
):
|
||
await pl.run_avito_newbuilding_sweep(
|
||
db, # type: ignore[arg-type]
|
||
run_id=8001,
|
||
config=_config(),
|
||
matcher=MagicMock(),
|
||
pages=pages,
|
||
request_delay_sec=0.0,
|
||
resume_run_id=resume_run_id,
|
||
)
|
||
return db
|
||
|
||
|
||
def _last_checkpoint(db: _FakeDb) -> list[int]:
|
||
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_resume_skips_done_pages_and_fetches_the_rest() -> None:
|
||
"""Резюм с чекпоинтом [1,2] на pages=4 обходит только страницы 3 и 4 —
|
||
ни одного запроса к уже собранным."""
|
||
db = _FakeDb({"done_buckets": [1, 2]})
|
||
await _run(db, resume_run_id=7999, pages=4)
|
||
|
||
assert _FakeScraper.visited_pages == [3, 4], (
|
||
"уже собранные страницы 1/2 не должны опрашиваться повторно"
|
||
)
|
||
assert _last_checkpoint(db) == [1, 2, 3, 4], "чекпоинт дописывается поверх унаследованного"
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_without_resume_starts_from_page_one() -> None:
|
||
"""Без resume_run_id поведение прежнее — обход с первой страницы."""
|
||
db = _FakeDb(None)
|
||
await _run(db, resume_run_id=None, pages=3)
|
||
|
||
assert _FakeScraper.visited_pages == [1, 2, 3]
|
||
assert _last_checkpoint(db) == [1, 2, 3]
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_failed_page_does_not_enter_checkpoint() -> None:
|
||
"""Страница, оборванная на середине (fail_on=2), НЕ считается пройденной.
|
||
|
||
Иначе следующий прогон пропустит её навсегда молча — sweep завершится
|
||
штатно, просто часть выдачи не соберётся никогда.
|
||
"""
|
||
db = _FakeDb(None)
|
||
await _run(db, resume_run_id=None, pages=3, fail_on=2)
|
||
|
||
# Страница 2 была АТАКОВАНА (попытка была — visited), но не пройдена до конца.
|
||
assert _FakeScraper.visited_pages == [1, 2]
|
||
ckpt = _last_checkpoint(db)
|
||
assert ckpt == [1], "упавшая страница 2 (и не начатая 3) не должны попасть в чекпоинт"
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_checkpoint_write_merges_and_preserves_foreign_key() -> None:
|
||
"""Запись чекпоинта — мерж (`counters || :counters`), не замена.
|
||
|
||
Симулируем текущую строку run'а с посторонним ключом (например, оставленным
|
||
другим писателем/предыдущим heartbeat'ом) — он обязан пережить наши записи.
|
||
"""
|
||
db = _MergingFakeDb(None)
|
||
db.row = {"foreign_key": "survives-me"}
|
||
|
||
await _run(db, resume_run_id=None, pages=1)
|
||
|
||
assert db.row.get("foreign_key") == "survives-me", (
|
||
"посторонний ключ в counters не пережил запись чекпоинта — запись была заменой, не мержем"
|
||
)
|
||
assert "done_buckets" in db.row
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_page_whose_save_failed_does_not_enter_checkpoint() -> None:
|
||
"""Страница собрана, но её save упал — в чекпоинт она попасть не должна.
|
||
|
||
Отказ save_listings перехватывается и прогон продолжается (это осознанно:
|
||
одна упавшая страница не должна ронять весь sweep). Но если отметить её
|
||
пройденной, следующий прогон её пропустит, и объявления оттуда не соберутся
|
||
НИКОГДА — молча, потому что прогон завершится штатно.
|
||
"""
|
||
db = _FakeDb(None)
|
||
await _run(db, resume_run_id=None, pages=3, save_fail_on_call=2)
|
||
|
||
assert _FakeScraper.visited_pages == [1, 2, 3]
|
||
ckpt = _last_checkpoint(db)
|
||
assert 2 not in ckpt, "страница с упавшим save попала в чекпоинт — покрытие потеряно молча"
|
||
assert ckpt == [1, 3]
|
||
|
||
|
||
@pytest.mark.asyncio
|
||
async def test_resume_starts_at_first_gap_not_after_the_last_page() -> None:
|
||
"""Дыра в чекпоинте перечитывается, а не перепрыгивается.
|
||
|
||
`max(done)+1` пропустил бы страницу 2 навсегда. Продолжаем с первой
|
||
несобранной: страницы после дыры перечитаются, что дешевле потери и
|
||
безопасно — повторная запись схлопывается по dedup_hash.
|
||
"""
|
||
db = _FakeDb({"done_buckets": [1, 3]})
|
||
await _run(db, resume_run_id=7777, pages=4)
|
||
|
||
assert _FakeScraper.visited_pages[0] == 2, (
|
||
"подхват начался не с дыры — пропущенная страница не соберётся никогда"
|
||
)
|