Compare commits
3 commits
6ced16618c
...
83946e5d4f
| Author | SHA1 | Date | |
|---|---|---|---|
| 83946e5d4f | |||
| 12945d7b2f | |||
| f6cf948358 |
2 changed files with 442 additions and 29 deletions
369
tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py
Normal file
369
tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py
Normal file
|
|
@ -0,0 +1,369 @@
|
||||||
|
"""SIGTERM-дрейн помечен `interrupted=1` и у full-load'ов с domclick-свипом (#3355).
|
||||||
|
|
||||||
|
#3330/#3346 закрыли city-свипы, но четыре функции с ЖИВЫМ чекпоинтом `done_buckets`
|
||||||
|
продолжали финализировать дрейн чистым `done`: cian/yandex/avito full-load (ветка
|
||||||
|
`RuntimeError("shutdown")`) и domclick city sweep (дрейн до SERP-фазы). Резюм
|
||||||
|
(`_drained_done` в scheduler._resume_decision) такой прогон не берёт — чекпоинт
|
||||||
|
пишется в никуда.
|
||||||
|
|
||||||
|
Кто чекпоинт ЧИТАЕТ (метка не диагностическая):
|
||||||
|
* cian_full_load — `skip_set` из `done_buckets` (ключ "room:lo:hi"),
|
||||||
|
планировщик зовёт с `resume_run_id=_pick_resume(...)`;
|
||||||
|
* avito_full_load — то же, ДВА job'а (обычный и exhaustive, #3315);
|
||||||
|
* domclick_city_sweep— `skip_buckets` из `done_buckets` (ROOM_BUCKETS), тоже
|
||||||
|
`_pick_resume`;
|
||||||
|
* yandex_full_load — читает `done_buckets` так же, НО у него нет job'а в
|
||||||
|
scheduler.py: единственный вход с `resume_run_id` —
|
||||||
|
админский POST. Метка там честно работает только при
|
||||||
|
ручном подхвате; автоматического потребителя нет.
|
||||||
|
|
||||||
|
Статус финализации нулевого дрейна фиксируется тестом отдельно: honest-status
|
||||||
|
гейты (runs.mark_done → _sweep_run_did_nothing/_phase_totally_failed/
|
||||||
|
_failed_ratio_too_high) могут увести прогон в 'failed'. Это резюмируемо
|
||||||
|
('failed' ∈ _RESUME_STATUSES), но статус — факт, а не догадка.
|
||||||
|
|
||||||
|
Анти-цикл («чистый done не резюмируется») уже закрыт в
|
||||||
|
test_3333_drain_mark_all_sweeps.py — здесь не дублируется.
|
||||||
|
"""
|
||||||
|
|
||||||
|
from __future__ import annotations
|
||||||
|
|
||||||
|
import os
|
||||||
|
|
||||||
|
# Settings собирается автофикстурой conftest'а и требует database_url — как в
|
||||||
|
# test_3333_drain_mark_all_sweeps.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 запоминается вместе с текстом SQL (для статуса)."""
|
||||||
|
|
||||||
|
def __init__(
|
||||||
|
self,
|
||||||
|
prev_counters: dict[str, Any] | None = None,
|
||||||
|
detail_rows: list[dict[str, Any]] | None = None,
|
||||||
|
) -> None:
|
||||||
|
self.writes: list[tuple[str, dict[str, Any]]] = []
|
||||||
|
self._prev_counters = prev_counters
|
||||||
|
self._detail_rows = detail_rows or []
|
||||||
|
|
||||||
|
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:
|
||||||
|
# SELECT counters предыдущего прогона (resume_run_id → чекпоинт).
|
||||||
|
row = types.SimpleNamespace(counters=self._prev_counters)
|
||||||
|
return MagicMock(**{"fetchone.return_value": row if self._prev_counters else None})
|
||||||
|
if params and "lim" in params:
|
||||||
|
# SELECT кандидатов detail-обогащения (cian full-load).
|
||||||
|
return MagicMock(**{"mappings.return_value.all.return_value": self._detail_rows})
|
||||||
|
return MagicMock()
|
||||||
|
|
||||||
|
def commit(self) -> None: ...
|
||||||
|
def rollback(self) -> None: ...
|
||||||
|
|
||||||
|
|
||||||
|
def _final(db: _FakeDb) -> tuple[str, dict[str, Any]]:
|
||||||
|
"""Последняя запись counters = финализатор."""
|
||||||
|
assert db.writes, "дрейн не оставил ни одной записи counters"
|
||||||
|
sql, counters = db.writes[-1]
|
||||||
|
if "status = 'done'" in sql:
|
||||||
|
return "done", counters
|
||||||
|
if "status = 'failed'" in sql:
|
||||||
|
return "failed", counters
|
||||||
|
return "heartbeat", counters
|
||||||
|
|
||||||
|
|
||||||
|
class _DrainAtFirstBucket:
|
||||||
|
"""Любой fetch-метод сразу зовёт on_bucket — настоящий `_on_bucket` и роняет дрейн."""
|
||||||
|
|
||||||
|
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||||||
|
self._browser = None
|
||||||
|
self._cffi = None
|
||||||
|
self.request_delay_sec = 0.0
|
||||||
|
|
||||||
|
async def __aenter__(self) -> _DrainAtFirstBucket:
|
||||||
|
return self
|
||||||
|
|
||||||
|
async def __aexit__(self, *_e: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
def __getattr__(self, name: str) -> Any:
|
||||||
|
async def _fetch(*_a: Any, **kw: Any) -> Any:
|
||||||
|
on_bucket = kw.get("on_bucket")
|
||||||
|
assert on_bucket is not None, f"{name} вызван без on_bucket — тест мимо дрейна"
|
||||||
|
on_bucket("2:0:5000000", [])
|
||||||
|
raise AssertionError("_on_bucket не оборвал прогон при shutdown_requested()")
|
||||||
|
|
||||||
|
return _fetch
|
||||||
|
|
||||||
|
|
||||||
|
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:
|
||||||
|
async def _boom(*_a: Any, **_kw: Any) -> Any:
|
||||||
|
raise AssertionError(f"скрапер вызван при дрейне: {name}")
|
||||||
|
|
||||||
|
return _boom
|
||||||
|
|
||||||
|
|
||||||
|
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_cian_full_load_drain_is_marked_interrupted() -> None:
|
||||||
|
"""cian full-load: дрейн на границе бакета → counters.interrupted == 1."""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
db = _FakeDb()
|
||||||
|
with (
|
||||||
|
patch.object(pl, "CianScraper", _DrainAtFirstBucket),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_cian_full_load(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
_status, counters = _final(db)
|
||||||
|
assert counters.get("interrupted") == 1, (
|
||||||
|
"cian full-load: оборванный дрейном прогон неотличим от полного обхода — "
|
||||||
|
"резюм его чекпоинт не подхватит"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_yandex_full_load_drain_is_marked_interrupted() -> None:
|
||||||
|
"""yandex full-load: дрейн на границе бакета → counters.interrupted == 1."""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
db = _FakeDb()
|
||||||
|
with (
|
||||||
|
patch.object(pl, "YandexRealtyScraper", _DrainAtFirstBucket),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_yandex_full_load(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
enrichment=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
_status, counters = _final(db)
|
||||||
|
assert counters.get("interrupted") == 1, (
|
||||||
|
"yandex full-load: оборванный дрейном прогон неотличим от полного обхода"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_avito_full_load_drain_is_marked_interrupted() -> None:
|
||||||
|
"""avito full-load (exhaustive, #3315): дрейн на границе бакета → interrupted == 1."""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
db = _FakeDb()
|
||||||
|
with (
|
||||||
|
patch.object(pl, "AvitoScraper", _DrainAtFirstBucket),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_avito_full_load(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
_status, counters = _final(db)
|
||||||
|
assert counters.get("interrupted") == 1, (
|
||||||
|
"avito full-load: дрейн без метки — бан-бюджет тратится на уже собранные бакеты"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_domclick_sweep_drain_is_marked_interrupted() -> None:
|
||||||
|
"""domclick city sweep: дрейн до SERP-фазы → counters.interrupted == 1."""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
db = _FakeDb()
|
||||||
|
with (
|
||||||
|
patch.object(pl, "DomClickScraper", _NeverCalledScraper),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_domclick_city_sweep(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
_status, counters = _final(db)
|
||||||
|
assert counters.get("interrupted") == 1, (
|
||||||
|
"domclick: оборванный дрейном прогон неотличим от полного обхода"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_domclick_sweep_drain_keeps_inherited_checkpoint() -> None:
|
||||||
|
"""domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload.
|
||||||
|
|
||||||
|
Мержа jsonb тут нет: app-level `scrape_runs.update_heartbeat`/`mark_done` пишут
|
||||||
|
`counters = CAST(:counters AS jsonb)` — полная перезапись, а scheduler при claim
|
||||||
|
done_buckets не наследует. Без явного ключа дрейн-прогон закрывался бы с пустым
|
||||||
|
чекпоинтом, и резюм пересобирал бы все шесть корзин заново.
|
||||||
|
"""
|
||||||
|
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: False),
|
||||||
|
):
|
||||||
|
await pl.run_domclick_city_sweep(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
resume_run_id=3354,
|
||||||
|
)
|
||||||
|
|
||||||
|
assert len(db.writes) >= 2, "дрейн обязан писать и heartbeat, и финализатор"
|
||||||
|
for label, (_sql, counters) in (("heartbeat", db.writes[-2]), ("финализатор", db.writes[-1])):
|
||||||
|
assert counters.get("done_buckets") == inherited, (
|
||||||
|
f"{label}: чекпоинт потерян ({counters.get('done_buckets')!r} вместо "
|
||||||
|
f"{inherited!r}) — резюм пересоберёт уже собранные корзины"
|
||||||
|
)
|
||||||
|
assert counters.get("interrupted") == 1, f"{label}: дрейн без метки"
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_cian_full_load_drain_in_detail_phase_is_marked_interrupted() -> None:
|
||||||
|
"""cian full-load: дрейн в detail-фазе — SERP целый, но прогон НЕ полный.
|
||||||
|
|
||||||
|
Ранний выход из detail-цикла делал `break` и уходил в обычный `mark_done`:
|
||||||
|
статус 'done' без метки, хотя обогащение обрезано на первой же записи.
|
||||||
|
"""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
class _SerpDoneScraper:
|
||||||
|
"""SERP отработал целиком (on_bucket не зовётся) — дрейн ловит detail-фаза."""
|
||||||
|
|
||||||
|
def __init__(self, *_a: Any, **_kw: Any) -> None:
|
||||||
|
self._browser = None
|
||||||
|
self.request_delay_sec = 0.0
|
||||||
|
self.last_dropped_nb = 0
|
||||||
|
self.state_extraction_attempts = 0
|
||||||
|
self.state_extraction_failures = 0
|
||||||
|
|
||||||
|
async def __aenter__(self) -> _SerpDoneScraper:
|
||||||
|
return self
|
||||||
|
|
||||||
|
async def __aexit__(self, *_e: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
async def fetch_all_secondary(self, **_kw: Any) -> None:
|
||||||
|
return None
|
||||||
|
|
||||||
|
db = _FakeDb(detail_rows=[{"id": 1, "source_url": "https://cian.ru/1"}])
|
||||||
|
with (
|
||||||
|
patch.object(pl, "CianScraper", _SerpDoneScraper),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_cian_full_load(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
enrich_detail=True,
|
||||||
|
detail_top_n=5,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
_status, counters = _final(db)
|
||||||
|
assert counters.get("interrupted") == 1, (
|
||||||
|
"cian full-load: дрейн в detail-фазе финализирован как полный прогон"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
@pytest.mark.asyncio
|
||||||
|
async def test_zero_drain_status_is_resumable() -> None:
|
||||||
|
"""Минор ревью #3346: какой СТАТУС даёт нулевой дрейн — и берёт ли его резюм.
|
||||||
|
|
||||||
|
honest-status гейты в runs.mark_done (#2625/#2700/honest-run-status) вправе
|
||||||
|
увести прогон в 'failed'. Тест фиксирует ФАКТ, а не ожидание: обе ветки
|
||||||
|
резюмируемы ('failed' — через _RESUME_STATUSES, 'done' — через метку), но
|
||||||
|
третьего исхода быть не должно.
|
||||||
|
"""
|
||||||
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
from scraper_kit.orchestration.scheduler import _RESUME_STATUSES, _resume_decision
|
||||||
|
|
||||||
|
db = _FakeDb()
|
||||||
|
with (
|
||||||
|
patch.object(pl, "AvitoScraper", _DrainAtFirstBucket),
|
||||||
|
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
|
||||||
|
):
|
||||||
|
await pl.run_avito_full_load(
|
||||||
|
db, # type: ignore[arg-type]
|
||||||
|
run_id=3355,
|
||||||
|
config=_config(),
|
||||||
|
matcher=MagicMock(),
|
||||||
|
request_delay_sec=0.0,
|
||||||
|
shutdown_requested=lambda: True,
|
||||||
|
)
|
||||||
|
|
||||||
|
status, counters = _final(db)
|
||||||
|
assert status == "done", f"нулевой дрейн финализирован как {status!r} (было 'done')"
|
||||||
|
|
||||||
|
prev = types.SimpleNamespace(
|
||||||
|
prev_id=4707,
|
||||||
|
prev_status=status,
|
||||||
|
prev_counters={**counters, "done_buckets": ["2:0:5000000"], "resume_chain": 0},
|
||||||
|
same_params=True,
|
||||||
|
age_h=2.0,
|
||||||
|
interval_days="7",
|
||||||
|
)
|
||||||
|
resume_from, verdict = _resume_decision(prev)
|
||||||
|
assert resume_from == 4707, (
|
||||||
|
f"дрейн со статусом {status!r} не подхвачен резюмом: {verdict} "
|
||||||
|
f"(_RESUME_STATUSES={sorted(_RESUME_STATUSES)})"
|
||||||
|
)
|
||||||
|
|
@ -3648,7 +3648,12 @@ async def run_cian_full_load(
|
||||||
didx + 1,
|
didx + 1,
|
||||||
len(priority_rows),
|
len(priority_rows),
|
||||||
)
|
)
|
||||||
break
|
# #3355: не break. SERP-фаза тут ЦЕЛАЯ (done_buckets полные),
|
||||||
|
# обрезано только обогащение — но break уводил в обычный
|
||||||
|
# mark_done ниже, и оборванный прогон выглядел полным. Тот же
|
||||||
|
# sentinel, что в _on_bucket → handler `RuntimeError("shutdown")`
|
||||||
|
# ставит interrupted=1 и сохраняет done_buckets.
|
||||||
|
raise RuntimeError("shutdown")
|
||||||
listing_id: int = row["id"]
|
listing_id: int = row["id"]
|
||||||
source_url: str = row["source_url"]
|
source_url: str = row["source_url"]
|
||||||
counters.detail_attempted += 1
|
counters.detail_attempted += 1
|
||||||
|
|
@ -3755,7 +3760,15 @@ async def run_cian_full_load(
|
||||||
counters.saved_inserted,
|
counters.saved_inserted,
|
||||||
counters.saved_updated,
|
counters.saved_updated,
|
||||||
)
|
)
|
||||||
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
|
||||||
|
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
|
||||||
|
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
|
||||||
|
# и done_buckets этого прогона не подхватывал никто.
|
||||||
|
runs.mark_done(
|
||||||
|
db,
|
||||||
|
run_id,
|
||||||
|
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
|
||||||
|
)
|
||||||
return counters
|
return counters
|
||||||
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
logger.exception("cian-full-load run_id=%d: fatal error", run_id)
|
||||||
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
|
|
@ -4019,7 +4032,15 @@ async def run_yandex_full_load(
|
||||||
counters.saved_inserted,
|
counters.saved_inserted,
|
||||||
counters.saved_updated,
|
counters.saved_updated,
|
||||||
)
|
)
|
||||||
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
|
||||||
|
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
|
||||||
|
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
|
||||||
|
# и done_buckets этого прогона не подхватывал никто.
|
||||||
|
runs.mark_done(
|
||||||
|
db,
|
||||||
|
run_id,
|
||||||
|
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
|
||||||
|
)
|
||||||
return counters
|
return counters
|
||||||
logger.exception("yandex-full-load run_id=%d: fatal error", run_id)
|
logger.exception("yandex-full-load run_id=%d: fatal error", run_id)
|
||||||
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
|
|
@ -4285,7 +4306,15 @@ async def run_avito_full_load(
|
||||||
counters.saved_inserted,
|
counters.saved_inserted,
|
||||||
counters.saved_updated,
|
counters.saved_updated,
|
||||||
)
|
)
|
||||||
runs.mark_done(db, run_id, {**counters.to_dict(), "done_buckets": sorted(done)})
|
# #3355: та же метка дрейна, что у city-свипов (#3330/#3333). Без неё
|
||||||
|
# оборванный дрейном обход неотличим от полного: статус тот же 'done',
|
||||||
|
# а `_drained_done` в scheduler._resume_decision смотрит именно на неё —
|
||||||
|
# и done_buckets этого прогона не подхватывал никто.
|
||||||
|
runs.mark_done(
|
||||||
|
db,
|
||||||
|
run_id,
|
||||||
|
{**counters.to_dict(), "done_buckets": sorted(done), "interrupted": 1},
|
||||||
|
)
|
||||||
return counters
|
return counters
|
||||||
logger.exception("avito-full-load run_id=%d: fatal error", run_id)
|
logger.exception("avito-full-load run_id=%d: fatal error", run_id)
|
||||||
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
runs.mark_failed(db, run_id, str(exc), counters.to_dict())
|
||||||
|
|
@ -4398,35 +4427,12 @@ async def run_domclick_city_sweep(
|
||||||
_scraper_ref: list[DomClickScraper] = []
|
_scraper_ref: list[DomClickScraper] = []
|
||||||
|
|
||||||
try:
|
try:
|
||||||
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
|
|
||||||
if runs.is_cancelled(db, run_id):
|
|
||||||
logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id)
|
|
||||||
runs.update_heartbeat(db, run_id, counters.to_dict())
|
|
||||||
return counters
|
|
||||||
elif shutdown_requested():
|
|
||||||
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
|
|
||||||
# минуя honest-status (это чистый drain, не QRATOR-блок).
|
|
||||||
logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id)
|
|
||||||
runs.update_heartbeat(db, run_id, counters.to_dict())
|
|
||||||
runs.mark_done(db, run_id, counters.to_dict())
|
|
||||||
return counters
|
|
||||||
|
|
||||||
logger.info(
|
|
||||||
"domclick-sweep run_id=%d: BFF citywide sweep city_id=%d "
|
|
||||||
"buckets=%d pages_cap=%d (watchdog %ds)",
|
|
||||||
run_id,
|
|
||||||
city_id,
|
|
||||||
len(ROOM_BUCKETS),
|
|
||||||
pages,
|
|
||||||
_sweep_timeout,
|
|
||||||
)
|
|
||||||
|
|
||||||
lots: list[ScrapedLot] = []
|
|
||||||
|
|
||||||
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
|
# ── Checkpoint/resume (#3118): корзина ROOM_BUCKETS — единица обхода ──
|
||||||
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
|
# Вместе со сдвигом #2854 подхват превращает случайную ротацию в
|
||||||
# систематический обход: банимый на 1-й корзине источник закрывает все
|
# систематический обход: банимый на 1-й корзине источник закрывает все
|
||||||
# шесть за несколько прогонов вместо повторов случайных.
|
# шесть за несколько прогонов вместо повторов случайных.
|
||||||
|
# Читаем ДО ранних выходов: дрейн-ветка ниже обязана положить чекпоинт в
|
||||||
|
# свой payload сама (#3355), см. комментарий там.
|
||||||
skip_buckets: set[str] = set()
|
skip_buckets: set[str] = set()
|
||||||
if resume_run_id is not None:
|
if resume_run_id is not None:
|
||||||
_prev_row = db.execute(
|
_prev_row = db.execute(
|
||||||
|
|
@ -4445,6 +4451,44 @@ async def run_domclick_city_sweep(
|
||||||
len(skip_buckets),
|
len(skip_buckets),
|
||||||
)
|
)
|
||||||
|
|
||||||
|
# ── Cooperative cancel перед SERP-фазой ──────────────────────────────
|
||||||
|
if runs.is_cancelled(db, run_id):
|
||||||
|
logger.info("domclick-sweep run_id=%d: cancelled before SERP", run_id)
|
||||||
|
runs.update_heartbeat(db, run_id, counters.to_dict())
|
||||||
|
return counters
|
||||||
|
elif shutdown_requested():
|
||||||
|
# SIGTERM-drain до SERP-фазы: финализируем run как mark_done(partial),
|
||||||
|
# минуя honest-status (это чистый drain, не QRATOR-блок).
|
||||||
|
logger.info("domclick-sweep run_id=%d: SIGTERM-drain — stopping before SERP", run_id)
|
||||||
|
# #3355: метка дрейна (см. city-свипы #3330/#3333) И явный перенос
|
||||||
|
# чекпоинта. Никакого jsonb-мержа тут нет: domclick пишет counters
|
||||||
|
# через app-level scrape_runs (update_heartbeat/mark_done делают
|
||||||
|
# `counters = CAST(:counters AS jsonb)` — ПОЛНАЯ перезапись), а
|
||||||
|
# scheduler при claim done_buckets не наследует. Без явного
|
||||||
|
# done_buckets дрейн-прогон закрылся бы с пустым чекпоинтом и резюм
|
||||||
|
# пересобирал бы все шесть корзин заново. Тот же приём, что в
|
||||||
|
# ban-ветке ниже (except NoProxyAvailableError).
|
||||||
|
_drain = {
|
||||||
|
**counters.to_dict(),
|
||||||
|
"done_buckets": sorted(skip_buckets),
|
||||||
|
"interrupted": 1,
|
||||||
|
}
|
||||||
|
runs.update_heartbeat(db, run_id, _drain)
|
||||||
|
runs.mark_done(db, run_id, _drain)
|
||||||
|
return counters
|
||||||
|
|
||||||
|
logger.info(
|
||||||
|
"domclick-sweep run_id=%d: BFF citywide sweep city_id=%d "
|
||||||
|
"buckets=%d pages_cap=%d (watchdog %ds)",
|
||||||
|
run_id,
|
||||||
|
city_id,
|
||||||
|
len(ROOM_BUCKETS),
|
||||||
|
pages,
|
||||||
|
_sweep_timeout,
|
||||||
|
)
|
||||||
|
|
||||||
|
lots: list[ScrapedLot] = []
|
||||||
|
|
||||||
async def _domclick_phase() -> None:
|
async def _domclick_phase() -> None:
|
||||||
"""Единственная citywide-фаза: fetch_city + save."""
|
"""Единственная citywide-фаза: fetch_city + save."""
|
||||||
nonlocal lots
|
nonlocal lots
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue