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
Ревью опровергло формулировку «точку терял каждый финализатор»: все четыре писателя в 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) — метка дрейна не меняет вердикт «нечего подхватывать».
225 lines
9.5 KiB
Python
225 lines
9.5 KiB
Python
"""Чекпоинт 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"
|