gendesign/tradein-mvp/backend/tests/test_3319_citysweep_checkpoint.py
bot-backend 41f21c4969
All checks were successful
CI Trade-In / changes (pull_request) Successful in 11s
CI / changes (pull_request) Successful in 13s
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 / 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 5m13s
docs(scraper-kit): обоснование #3319 — что именно теряло чекпоинт
Ревью опровергло формулировку «точку терял каждый финализатор»: все четыре
писателя в runs.py мержат jsonb (`counters || :counters`), записанный ключ
переживал mark_done/mark_banned/mark_failed. Правка закрывает выходы РАНЬШЕ
первого end-of-anchor heartbeat (cancel/дрейн на первом якоре, ранний done
#1950 на якоре №1) — комментарий и докстринг переписаны на это.

Замер «0 из 67 за 60 дней» назван тем, чем он является: запись появилась
26.08.2026 (#3074) при такте avito 7 суток, выборка почти вся из эры без
механизма. Плюс тест на прогон без ключей (эра до #3074) — метка дрейна не
меняет вердикт «нечего подхватывать».
2026-09-02 14:53:02 +05:00

225 lines
9.5 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Чекпоинт avito_city_sweep доживает до финализатора (#3319).
Точку писала ровно одна строка — end-of-anchor heartbeat в конце итерации цикла
якорей. Финализаторы её не стирали (все писатели в runs.py мержат jsonb:
`counters || :counters`), дыра в другом: выходы, случившиеся РАНЬШЕ первой такой
записи, точки не оставляли вовсе — cancel/SIGTERM-дрейн на границе первого якоря
и ранний done #1950 («SERP собран, detail заблокирован») на якоре №1. Ими и
кончается типичный прод-прогон с `anchors_done: 1` из 5.
Замер «0 из 67 прогонов за 60 дней несут done_buckets» тут НЕ доказательство:
строка записи появилась только 26.08.2026 (#3074) при такте avito 7 суток —
выборка почти целиком из эры, где механизма не существовало.
Три инварианта, ради которых тест:
1. done-выход несёт done_buckets — иначе точка существует только в логе.
2. Якорь, умерший по таймауту, НЕ пройден: SERP мог успеть, detail нет.
Пройденным его записать = резюм пропустит его навсегда и молча.
3. SIGTERM-дрейн отличим от полного обхода (counters.interrupted=1) и
участвует в резюме — статус у обоих 'done', счётчики частичные.
"""
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
ANCHOR_A = (56.83, 60.60, "ekb-center")
ANCHOR_B = (56.79, 60.63, "ekb-south")
class _FakeDb:
"""Все UPDATE'ы с counters (heartbeat И финализаторы) складываются по порядку."""
def __init__(self) -> None:
self.writes: list[dict[str, Any]] = []
def execute(self, _stmt: Any, params: dict[str, Any] | None = None) -> Any:
if params and "counters" in params:
self.writes.append(json.loads(params["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
exc: type[BaseException] | None = None
lots_per_anchor: int = 0
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) and _FakeScraper.exc is not None:
raise _FakeScraper.exc("якорь сорвался")
return [MagicMock() for _ in range(_FakeScraper.lots_per_anchor)]
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(
*,
raise_on: tuple[float, float] | None = None,
exc: type[BaseException] | None = None,
shutdown_after_first: bool = False,
saved: tuple[int, int] = (0, 0),
lots_per_anchor: int = 0,
) -> _FakeDb:
from scraper_kit.orchestration import pipeline as pl
_FakeScraper.visited = []
_FakeScraper.raise_on = raise_on
_FakeScraper.exc = exc
_FakeScraper.lots_per_anchor = lots_per_anchor
db = _FakeDb()
def _shutdown() -> bool:
return shutdown_after_first and bool(_FakeScraper.visited)
with (
patch.object(pl, "AvitoScraper", _FakeScraper),
patch.object(pl, "AsyncSession", _FakeAsyncSession),
patch.object(pl, "save_listings", lambda *_a, **_kw: saved),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
await pl.run_avito_city_sweep(
db, # type: ignore[arg-type]
run_id=3319,
config=_config(),
matcher=MagicMock(),
enrichment=MagicMock(),
anchors=[ANCHOR_A, ANCHOR_B],
enrich_houses=False,
enrich_imv=False,
detail_top_n=0,
shutdown_requested=_shutdown,
)
return db
@pytest.mark.asyncio
async def test_done_exit_carries_checkpoint() -> None:
"""Финализатор полного обхода несёт done_buckets, а не голые счётчики."""
db = await _run()
assert db.writes[-1].get("done_buckets") == ["ekb-center", "ekb-south"], (
"финальный (done) выход отдал counters без чекпоинта — точки в прогоне нет"
)
@pytest.mark.asyncio
async def test_serp_ok_done_exit_carries_checkpoint() -> None:
"""Ранний done-выход #1950 («SERP собран, detail заблокирован») — тоже.
Именно этим выходом кончается типичный прод-прогон, и он происходит РАНЬШЕ
единственной строки, которая писала точку.
"""
from scraper_kit.orchestration import pipeline as pl
db = await _run(
raise_on=(ANCHOR_B[0], ANCHOR_B[1]),
exc=pl.AvitoBlockedError,
saved=(1, 0), # SERP intake > 0 → ветка ставит 'done', а не 'banned'
lots_per_anchor=1,
)
last = db.writes[-1]
assert "enrichment_abort_note" in last, "сработала не та ветка выхода"
assert last.get("done_buckets") == ["ekb-center"], (
"ранний done-выход потерял якорь, пройденный до блокировки"
)
@pytest.mark.asyncio
async def test_timed_out_anchor_is_not_checkpointed() -> None:
"""Якорь, умерший по таймауту, не считается пройденным.
Иначе резюм пропустит его навсегда, и это будет незаметно: прогон
завершается штатно, просто часть города не собирается никогда.
"""
db = await _run(raise_on=(ANCHOR_A[0], ANCHOR_A[1]), exc=TimeoutError)
ckpt = db.writes[-1].get("done_buckets")
assert "ekb-center" not in ckpt, "якорь-таймаут попал в чекпоинт"
assert "ekb-south" in ckpt, "исправный якорь не зафиксирован"
@pytest.mark.asyncio
async def test_drain_exit_is_distinguishable_from_full_done() -> None:
"""SIGTERM-дрейн помечен interrupted=1; полный обход — нет."""
drained = await _run(shutdown_after_first=True)
full = await _run()
assert drained.writes[-1].get("interrupted") == 1, (
"оборванный дрейном прогон неотличим от полного обхода"
)
assert drained.writes[-1].get("done_buckets") == ["ekb-center"]
assert "interrupted" not in full.writes[-1], "полный обход помечен как оборванный"
def _prev_run(counters: dict[str, Any]) -> types.SimpleNamespace:
return types.SimpleNamespace(
prev_id=4707,
prev_status="done",
prev_counters=counters,
same_params=True,
age_h=2.0,
interval_days="7",
)
def test_drained_done_is_resumable_but_clean_done_is_not() -> None:
"""Метка дрейна доходит до решения о резюме — иначе она диагностика ради себя."""
from scraper_kit.orchestration.scheduler import _resume_decision
ckpt = {"done_buckets": ["ekb-center"], "resume_chain": 0}
resume_from, verdict = _resume_decision(_prev_run({**ckpt, "interrupted": 1}))
assert resume_from == 4707, f"дрейн не подхвачен: {verdict}"
resume_from, verdict = _resume_decision(_prev_run(ckpt))
assert resume_from is None, "полный обход подхватывать нечего"
assert verdict["resume_reason"] == "status_done"
# Прогон из эры до #3074: ключей нет вовсе — метка дрейна не должна менять
# вердикт «нечего подхватывать» на что-то другое.
_, verdict = _resume_decision(_prev_run({}))
assert verdict["resume_reason"] == "status_done"
_, verdict = _resume_decision(_prev_run({"interrupted": 1}))
assert verdict["resume_reason"] == "no_checkpoint"