Merge pull request 'Сбор МЕРА: упавший прогон не висит zombie 6 часов, провалившаяся фаза не попадает в чекпоинт, проверка отмены не держит транзакцию часами' (#3561) from fix/kit-orchestration into main
Some checks failed
Deploy Trade-In / test (push) Blocked by required conditions
Deploy Trade-In / build-backend (push) Blocked by required conditions
Deploy Trade-In / deploy (push) Blocked by required conditions
Deploy Trade-In / perimeter-smoke (push) Blocked by required conditions
Deploy Trade-In / build-frontend (push) Blocked by required conditions
Deploy Trade-In / build-browser (push) Blocked by required conditions
Deploy Trade-In / deploy-status (push) Blocked by required conditions
Deploy Trade-In / changes (push) Has been cancelled

This commit is contained in:
bot-backend 2026-09-17 09:23:42 +00:00
commit 1e99d91afe
8 changed files with 1030 additions and 11 deletions

View file

@ -1252,8 +1252,11 @@ async def start_avito_city_sweep(
request_delay_sec=payload.request_delay_sec,
enrich_imv=payload.enrich_imv,
)
except Exception:
except Exception as exc:
logger.exception("city-sweep background task run_id=%d crashed", run_id)
# #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший
# до финализатора пайплайна прогон висел 'running' до zombie.
runs_mod.mark_crashed(sweep_db, run_id, exc)
finally:
sweep_db.close()
@ -1344,8 +1347,11 @@ async def start_cian_city_sweep(
detail_top_n=payload.detail_top_n,
enrich_houses=payload.enrich_houses,
)
except Exception:
except Exception as exc:
logger.exception("cian-sweep background task run_id=%d crashed", run_id)
# #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший
# до финализатора пайплайна прогон висел 'running' до zombie.
runs_mod.mark_crashed(sweep_db, run_id, exc)
finally:
sweep_db.close()
@ -1487,8 +1493,11 @@ async def start_cian_full_load(
resume_run_id=payload.resume_run_id,
secondary_only=payload.secondary_only,
)
except Exception:
except Exception as exc:
logger.exception("cian-full-load background task run_id=%d crashed", run_id)
# #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший
# до финализатора пайплайна прогон висел 'running' до zombie.
runs_mod.mark_crashed(task_db, run_id, exc)
finally:
task_db.close()
@ -1590,8 +1599,11 @@ async def start_yandex_full_load(
concurrency=payload.concurrency,
resume_run_id=payload.resume_run_id,
)
except Exception:
except Exception as exc:
logger.exception("yandex-full-load background task run_id=%d crashed", run_id)
# #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший
# до финализатора пайплайна прогон висел 'running' до zombie.
runs_mod.mark_crashed(task_db, run_id, exc)
finally:
task_db.close()
@ -1658,8 +1670,11 @@ async def start_yandex_city_sweep(
request_delay_sec=payload.request_delay_sec,
enrich_address=payload.enrich_address,
)
except Exception:
except Exception as exc:
logger.exception("yandex-sweep background task run_id=%d crashed", run_id)
# #1940: ручной запуск идёт мимо scheduler._dispatch — без этого упавший
# до финализатора пайплайна прогон висел 'running' до zombie.
runs_mod.mark_crashed(sweep_db, run_id, exc)
finally:
sweep_db.close()

View file

@ -0,0 +1,202 @@
"""Упавший хендлер не оставляет строку scrape_runs в 'running' до zombie-reaper'а (#1940).
Прод 17.09 06:58: avito_city_sweep run 7349 через 10 мс после claim BrowserFetcher.__aenter__
бросил NoProxyAvailableError, планировщик записал «crashed run_id=7349» и всё. Строка
осталась 'running' с heartbeat_at == started_at; 21 сутки до этого так ушли в 'zombie'
18 avito-свипов, каждый через 6 ч. Проверяется значение статуса в строке прогона.
Финализатор `runs.mark_crashed`. Его зовут оба запускающих: планировщик и ручные
запуски из админки (они идут мимо `_dispatch`). Здесь боевой `mark_crashed` над
двойником одной строки scrape_runs: `_Row` отвечает на те же SQL, что пишет runs.py,
с гейтом `WHERE status = 'running'`.
"""
from __future__ import annotations
import asyncio
import os
import re
import types
from typing import Any
from unittest.mock import AsyncMock, MagicMock, patch
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
import pytest
from scraper_kit.orchestration import runs as kit_runs
from scraper_kit.orchestration.runs import BAN_KIND_INFRA, CONSECUTIVE_FAILURE_ALERT_THRESHOLD
from scraper_kit.orchestration.scheduler import Handler, SchedulerContext, _dispatch
from scraper_kit.proxy_errors import NoProxyAvailableError
RUN_ID = 7349
class _Row:
"""Сессия над одной строкой scrape_runs (status/ban_kind/error/counters).
История источника для стрик-алерта `CONSECUTIVE_FAILURE_ALERT_THRESHOLD` неудач:
ровно веха лестницы `_streak_alert_due`, на которой дубль и уходил бы в Sentry.
"""
def __init__(self) -> None:
self.status = "running"
self.ban_kind: str | None = None
self.error: str | None = None
self.updates = 0
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
sql = str(stmt)
p = params or {}
res = MagicMock()
res.fetchone.return_value = None # гейт _claim_run: «не running»
res.scalar.return_value = True # advisory-lock взят
if "INSERT INTO scrape_runs" in sql:
res.fetchone.return_value = types.SimpleNamespace(id=RUN_ID)
elif "SELECT status FROM scrape_runs WHERE id" in sql:
res.fetchone.return_value = types.SimpleNamespace(status=self.status)
elif "UPDATE scrape_runs" in sql and "SET status" in sql:
hit = self.status == "running"
if hit:
m = re.search(r"SET status = '(\w+)'", sql)
assert m is not None
self.status = m.group(1)
self.ban_kind = p.get("ban_kind")
self.error = p.get("error")
self.updates += 1
res.first.return_value = types.SimpleNamespace(id=RUN_ID) if hit else None
elif "SELECT source FROM scrape_runs" in sql:
res.fetchone.return_value = types.SimpleNamespace(source="avito_city_sweep")
elif "SELECT status, counters FROM scrape_runs" in sql:
res.fetchall.return_value = [
types.SimpleNamespace(status="failed", counters={})
] * CONSECUTIVE_FAILURE_ALERT_THRESHOLD
return res
def commit(self) -> None: ...
def rollback(self) -> None: ...
def close(self) -> None: ...
async def _dispatch_and_wait(job: Any, row: _Row) -> None:
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=lambda: row,
)
sched = {
"id": 1,
"source": "avito_city_sweep",
"window_start_hour": 6,
"window_end_hour": 7,
"default_params": {},
}
run_id = await _dispatch(Handler(job, "avito_city_sweep"), row, sched, ctx) # type: ignore[arg-type]
assert run_id == RUN_ID
await asyncio.gather(*ctx._inflight_tasks, return_exceptions=True)
# ── планировщик ─────────────────────────────────────────────────────────────────
async def test_no_proxy_crash_before_try_marks_run_banned_infra() -> None:
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
raise NoProxyAvailableError("avito")
row = _Row()
await _dispatch_and_wait(_job, row)
assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA)
assert row.error is not None and "no proxy available" in row.error
async def test_wrapped_no_proxy_is_still_infra() -> None:
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
try:
raise NoProxyAvailableError("cian")
except NoProxyAvailableError as exc:
raise RuntimeError("sidecar fetch failed") from exc
row = _Row()
await _dispatch_and_wait(_job, row)
assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA)
async def test_other_crash_marks_run_failed() -> None:
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
raise KeyError("anchors")
row = _Row()
await _dispatch_and_wait(_job, row)
assert row.status == "failed"
assert row.error is not None and "KeyError" in row.error
async def test_handler_own_final_status_survives() -> None:
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
kit_runs.mark_done(db, run_id, {"done_buckets": ["center"]})
raise RuntimeError("after mark_done")
row = _Row()
await _dispatch_and_wait(_job, row)
assert (row.status, row.updates) == ("done", 1)
async def test_pipeline_failure_is_alerted_once_not_twice() -> None:
"""Обычный путь пайплайна: mark_failed и raise. Повторная финализация не шлёт дубль.
Стрик после первой финализации стоит на вехе лестницы; no-op mark_failed из
планировщика позвал бы `_alert_on_run_id` ещё раз с тем же стриком.
"""
async def _job(db: Any, run_id: int, params: dict[str, Any], c: Any) -> None:
kit_runs.mark_failed(db, run_id, "cian-sweep fatal", {})
raise RuntimeError("cian-sweep fatal")
row = _Row()
sentry = MagicMock()
with patch.object(kit_runs, "sentry_sdk", sentry):
await _dispatch_and_wait(_job, row)
assert (row.status, row.updates) == ("failed", 1)
assert sentry.capture_message.call_count == 1
# ── ручные запуски из админки (мимо _dispatch) ─────────────────────────────────
@pytest.mark.parametrize(
("path", "runner"),
[
("avito-city-sweep", "run_avito_city_sweep"),
("cian-city-sweep", "run_cian_city_sweep"),
("cian-full-load", "run_cian_full_load"),
("yandex-full-load", "run_yandex_full_load"),
("yandex-city-sweep", "run_yandex_city_sweep"),
],
)
def test_admin_manual_launch_crash_is_finalized(path: str, runner: str) -> None:
from fastapi import FastAPI
from fastapi.testclient import TestClient
from app.api.v1 import admin as admin_module
from app.core.db import get_db
app = FastAPI()
app.include_router(admin_module.router, prefix="/api/v1/admin")
row = _Row()
app.dependency_overrides[get_db] = lambda: row
with (
patch.object(kit_runs, "sentry_sdk", None),
patch.object(admin_module, "has_running_run", return_value=False),
patch.object(admin_module, "SessionLocal", lambda: row),
patch.object(admin_module, runner, AsyncMock(side_effect=NoProxyAvailableError("avito"))),
):
r = TestClient(app).post(f"/api/v1/admin/scrape/{path}", json={})
assert r.status_code == 200, r.text
assert (row.status, row.ban_kind) == ("banned", BAN_KIND_INFRA)

View file

@ -0,0 +1,248 @@
"""Крайние случаи city-свипов kit'а, не покрытые после удаления legacy-пайплайна (#2406).
Уже покрыто: avito таймаут якоря и дрейн после первого якоря (test_3319); дрейн ДО
первого якоря у yandex/cian/newbuilding (test_3333); отмена у domclick (test_3369).
Здесь остаток:
(а) таймаут якоря у cian и yandex; у domclick снятая/упавшая фаза не отмечает
корзины, строки которых не сохранены;
(б) SIGTERM-дрейн ПОСЛЕ первого якоря у cian и yandex;
(в) пользовательская отмена (runs.is_cancelled=True) после первого якоря у avito,
cian и yandex.
Проверяется строка прогона, а не вызовы: _FakeDb мержит counters всех записей по
порядку (как `counters || :counters` в runs.py) и запоминает финальный статус.
"""
from __future__ import annotations
import os
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
from scraper_kit.orchestration import pipeline as pl
ANCHOR_A = (56.83, 60.60, "ekb-center")
ANCHOR_B = (56.79, 60.63, "ekb-south")
class _FakeDb:
def __init__(self) -> None:
self.counters: dict[str, Any] = {}
self.status = "running"
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
sql = str(stmt)
if params and "counters" in params:
self.counters.update(json.loads(params["counters"]))
for st in ("done", "failed", "banned"):
if f"SET status = '{st}'" in sql and self.status == "running":
self.status = st
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
class _Scraper:
"""Двойник Cian/Yandex/AvitoScraper: помнит якоря, роняет заданный TimeoutError'ом."""
visited: list[float] = [] # noqa: RUF012
timeout_on: float | None = None
gate_fetch_attempts = state_extraction_attempts = 1
gate_fetch_failures = state_extraction_failures = 0
request_delay_sec = 0.0
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
async def __aenter__(self) -> _Scraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
def _visit(self, lat: float) -> None:
_Scraper.visited.append(lat)
if lat == _Scraper.timeout_on:
raise TimeoutError
async def fetch_around_multi_room(self, lat: float, *_a: Any, **kw: Any) -> list[Any]:
self._visit(lat)
if kw.get("on_combo") is not None: # yandex: единица чекпоинта — combo
kw["on_combo"](f"combo@{lat}", [])
return []
return [types.SimpleNamespace(listing_segment="vtorichka", house_ext_id=None)]
async def fetch_around(self, lat: float, *_a: Any, **_kw: Any) -> list[Any]:
self._visit(lat)
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 _sweep(
kind: str,
*,
timeout_on: float | None = None,
drain_after_first: bool = False,
cancel_after_first: bool = False,
) -> _FakeDb:
_Scraper.visited = []
_Scraper.timeout_on = timeout_on
db = _FakeDb()
common: dict[str, Any] = {
"run_id": 2406,
"config": _config(),
"matcher": MagicMock(),
"anchors": [ANCHOR_A, ANCHOR_B],
"request_delay_sec": 0.0,
"shutdown_requested": lambda: drain_after_first and bool(_Scraper.visited),
}
with (
patch.object(pl, "CianScraper", _Scraper),
patch.object(pl, "YandexRealtyScraper", _Scraper),
patch.object(pl, "AvitoScraper", _Scraper),
patch.object(pl, "AsyncSession", _Scraper),
patch.object(pl, "save_listings", lambda *_a, **_kw: (1, 0)),
patch.object(
pl.runs, "is_cancelled", lambda *_a: cancel_after_first and bool(_Scraper.visited)
),
):
if kind == "cian":
await pl.run_cian_city_sweep(
db, # type: ignore[arg-type]
enrich_houses=False,
detail_top_n=0,
**common,
)
elif kind == "yandex":
enrichment = MagicMock()
enrichment.record_yandex_price_history.return_value = 0
await pl.run_yandex_city_sweep(
db, # type: ignore[arg-type]
enrichment=enrichment,
enrich_address=False,
**common,
)
else:
await pl.run_avito_city_sweep(
db, # type: ignore[arg-type]
enrichment=MagicMock(),
enrich_houses=False,
enrich_imv=False,
detail_top_n=0,
**common,
)
return db
def _done_key(kind: str, anchor: tuple[float, float, str]) -> str:
return f"combo@{anchor[0]}" if kind == "yandex" else anchor[2]
# ── (а) таймаут якоря: якорь не засчитан, следующий обработан ──────────────────
@pytest.mark.parametrize("kind", ["cian", "yandex"])
async def test_anchor_timeout_skips_anchor_and_continues(kind: str) -> None:
db = await _sweep(kind, timeout_on=ANCHOR_A[0])
assert _Scraper.visited == [ANCHOR_A[0], ANCHOR_B[0]]
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_B)]
assert db.counters["errors_count"] >= 1
@pytest.mark.parametrize("broken", ["fetch_timeout", "save_failed"])
async def test_domclick_unsaved_buckets_are_not_checkpointed(broken: str) -> None:
"""Корзина без сохранённых строк не попадает в чекпоинт.
Лоты Домклика копятся в памяти и пишутся ОДНИМ save_listings после всех корзин.
Фаза, снятая watchdog'ом (или упавшая на save), не сохранила ничего — даже из
корзины 'st', которую скрейпер успел пройти. Отметить её пройденной значило бы,
что следующий прогон пропустит её через skip_buckets навсегда (миграция 308).
"""
class _Dc(_Scraper):
blocked = False
geo_filtered = fetch_errors = bucket_start_index = 0
buckets_completed, buckets_total = 1, 6
completed_buckets = ["st"] # noqa: RUF012
async def fetch_city(self, **_kw: Any) -> list[Any]:
if broken == "fetch_timeout":
raise TimeoutError
return [object()]
saved: list[int] = []
def _save(_db: Any, lots: list[Any], **_kw: Any) -> tuple[int, int]:
if broken == "save_failed":
raise RuntimeError("save_listings упал")
saved.append(len(lots))
return len(lots), 0
db = _FakeDb()
with (
patch.object(pl, "DomClickScraper", _Dc),
patch.object(pl, "save_listings", _save),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_domclick_city_sweep(
db, # type: ignore[arg-type]
run_id=2406,
config=types.SimpleNamespace(browser_http_endpoint="http://x:9000"),
matcher=MagicMock(),
pages=1,
request_delay_sec=0.0,
)
assert saved == [], "тест не о том: строки сохранились"
assert counters.errors_count >= 1
assert db.counters["done_buckets"] == [], (
f"корзины без единой сохранённой строки в чекпоинте: {db.counters['done_buckets']}"
)
assert db.status == "failed"
# ── (б) дрейн после первого якоря: interrupted=1, в чекпоинте ровно первый ──────
@pytest.mark.parametrize("kind", ["cian", "yandex"])
async def test_drain_after_first_anchor_is_partial_and_interrupted(kind: str) -> None:
db = await _sweep(kind, drain_after_first=True)
assert _Scraper.visited == [ANCHOR_A[0]]
assert db.counters["interrupted"] == 1
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)]
assert db.status == "done"
# ── (в) пользовательская отмена после первого якоря ────────────────────────────
@pytest.mark.parametrize("kind", ["avito", "cian", "yandex"])
async def test_user_cancel_after_first_anchor_stops_without_finalizing(kind: str) -> None:
db = await _sweep(kind, cancel_after_first=True)
assert _Scraper.visited == [ANCHOR_A[0]]
assert db.counters["done_buckets"] == [_done_key(kind, ANCHOR_A)]
assert "interrupted" not in db.counters
# Строку финализирует сам mark_cancelled (UI); свип только останавливается.
assert db.status == "running"

View file

@ -0,0 +1,322 @@
"""Бакет уходит в чекпоинт только если ВСЕ его фазы что-то дали (#3415).
Прод-факт. `cian_city_sweep` 6179 (06.09) статус `failed` с причиной
«phase-honest-status: фаза 'houses' отказала полностью 30 из 30 попыток
неудачны, обогащено 0», а `counters.done_buckets` при этом несёт ВСЕ пять якорей
(«Академический», «Пионерский», «Уралмаш», «Центр», «ЮЗ»). Следующий прогон 6276
(07.09) резюмировался от него (`resume_from: 6179`, `resume_buckets: 5`),
пропустил все якоря и отчитался `done` с `errors_count=0`, `houses_attempted=0`,
`lots_fetched=0` зелёный по построению: фаза не исполнялась. Из-за этого фикс
#3396 (houses через пул прокси) простоял на проде с 06.09 и ни разу не отработал.
ИНВАРИАНТ. Отказ ЦЕЛОЙ фазы внутри якоря это «якорь пройден не был», ровно как
исключение или таймаут (#3074/#3319). Разница только в том, что отказ фазы виден
не потоком управления, а бухгалтерией счётчиков: `houses_attempted` попыток,
`houses_failed` отказов, обогащено ноль. Порог здесь одна попытка, а не три,
как у run-level `_phase_totally_failed` (#2700): тот СТАВИТ ДИАГНОЗ прогону и
обязан отсеивать шум, а этот решает «собирать ли якорь заново», и цена ошибок
несимметрична лишний повтор якоря стоит нескольких запросов, а ложная отметка
«пройден» теряет часть города навсегда и молча.
"""
from __future__ import annotations
import os
# Settings собирается автофикстурой conftest'а и требует database_url.
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")
# Два дома на якорь — НИЖЕ порога run-level правила (_PHASE_MIN_ATTEMPTS = 3):
# гейт бакета обязан сработать на своей, более строгой границе.
HOUSE_ROWS = [
{"id": 101, "cian_zhk_url": "https://cian.ru/zhk/101"},
{"id": 102, "cian_zhk_url": "https://cian.ru/zhk/102"},
]
DETAIL_ROWS = [
{"id": 201, "source_url": "/ekb/kvartiry/201"},
{"id": 202, "source_url": "/ekb/kvartiry/202"},
]
class _FakeDb:
def __init__(
self,
prev_counters: dict[str, Any] | None = None,
house_rows: list[dict[str, Any]] | None = None,
) -> None:
self.prev_counters = prev_counters or {}
# [] — запрос домов падает (ветка «houses DB query failed»).
self.house_rows = HOUSE_ROWS if house_rows is None else house_rows
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))
if params and "ids" in params:
# houses-фаза: дома, у которых есть cian_zhk_url.
if not self.house_rows:
raise RuntimeError("houses DB query failed")
return MagicMock(mappings=lambda: MagicMock(all=lambda: self.house_rows))
if params and "limit" in params:
# detail-фаза avito: карточки-кандидаты на обогащение.
return MagicMock(mappings=lambda: MagicMock(all=lambda: DETAIL_ROWS))
return MagicMock()
def commit(self) -> None: ...
def rollback(self) -> None: ...
def _lot() -> types.SimpleNamespace:
"""Новостроечный лот с привязкой к ЖК — только такой доходит до houses-фазы."""
return types.SimpleNamespace(
listing_segment="novostroyki",
house_source="cian_newbuilding",
house_ext_id="nb-1",
)
class _FakeScraper:
"""Двойник CianScraper: помнит, за какими якорями реально ходили."""
visited: list[str] = [] # noqa: RUF012 — тестовый сборник
def __init__(self, *_a: Any, **_kw: Any) -> None:
self.state_extraction_attempts = 1
self.state_extraction_failures = 0
self.request_delay_sec = 0.0
async def __aenter__(self) -> _FakeScraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
async def fetch_around_multi_room(
self, lat: float, lon: float, *_a: Any, **_kw: Any
) -> list[types.SimpleNamespace]:
_FakeScraper.visited.append(f"{lat},{lon}")
return [_lot()]
def _config() -> types.SimpleNamespace:
return types.SimpleNamespace(
scraper_proxy_url=None,
scraper_fetch_mode="cffi",
use_proxy_pool_browser=False,
browser_http_endpoint=None,
environment="test",
)
async def _run(
prev: dict[str, Any] | None,
*,
houses_ok: bool,
run_id: int = 8401,
house_rows: list[dict[str, Any]] | None = None,
failing_urls: frozenset[str] | None = None,
) -> tuple[_FakeDb, Any]:
from scraper_kit.orchestration import pipeline as pl
_FakeScraper.visited = []
db = _FakeDb(prev, house_rows)
async def _fake_newbuilding(zhk_url: str, *_a: Any, **_kw: Any) -> Any:
# None — ровно тот отказ, что был на проде 06.09: исключения нет,
# счётчик houses_failed растёт, фаза возвращается штатно.
if failing_urls is not None:
return None if zhk_url in failing_urls else object()
return object() if houses_ok else None
with (
patch.object(pl, "CianScraper", _FakeScraper),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl, "fetch_newbuilding", _fake_newbuilding),
patch.object(pl, "save_newbuilding_enrichment", lambda *_a, **_kw: None),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_cian_city_sweep(
db, # type: ignore[arg-type]
run_id=run_id,
config=_config(),
matcher=MagicMock(),
anchors=[ANCHOR_A, ANCHOR_B],
enrich_houses=True,
detail_top_n=0,
request_delay_sec=0.0,
resume_run_id=(run_id - 1) if prev is not None else None,
)
return db, counters
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_anchor_with_totally_failed_houses_is_not_checkpointed() -> None:
"""(а) Якорь, где houses отказала на ВСЕХ домах, не попадает в чекпоинт."""
db, counters = await _run(None, houses_ok=False)
assert counters.houses_attempted > 0, "houses-фаза вообще не исполнялась — тест не о том"
assert counters.houses_enriched == 0, "houses-фаза что-то обогатила — отказ не полный"
assert _last_checkpoint(db) == [], (
"якорь с полностью отказавшей houses-фазой помечен пройденным — "
f"чекпоинт {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_next_run_re_enters_houses() -> None:
"""(б) Следующий прогон с этим чекпоинтом снова заходит в houses."""
db1, _ = await _run(None, houses_ok=False)
_db2, counters2 = await _run(
{"done_buckets": _last_checkpoint(db1)}, houses_ok=True, run_id=8402
)
assert counters2.houses_attempted > 0, (
"прогон-наследник пропустил houses-фазу: чекпоинт предшественника "
"объявил якоря пройденными, хотя фаза в них не дала ничего"
)
assert counters2.houses_enriched > 0, "houses-фаза исполнилась, но ничего не обогатила"
@pytest.mark.asyncio
async def test_healthy_anchor_is_still_checkpointed() -> None:
"""(в) Контроль: якорь, где все фазы отработали, помечается как раньше."""
db, counters = await _run(None, houses_ok=True)
assert counters.houses_enriched > 0
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], (
"перестали помечать пройденные якоря вовсе — резюм (#3074) сломан"
)
@pytest.mark.asyncio
async def test_anchor_with_partially_failed_houses_is_still_checkpointed() -> None:
"""Контроль жадности: один отказавший дом из двух — якорь пройден.
Гейт про ПОЛНЫЙ отказ фазы. Жадный гейт (любой отказ) собирал бы такой якорь
заново каждым прогоном на проде houses почти всегда теряет пару домов.
"""
db, counters = await _run(
None, houses_ok=True, failing_urls=frozenset({HOUSE_ROWS[0]["cian_zhk_url"]})
)
assert (counters.houses_attempted, counters.houses_failed) == (4, 2)
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], (
f"якорь с частичным отказом houses не помечен пройденным: {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_single_failed_attempt_blocks_checkpoint() -> None:
"""Порог — одна попытка: единственный дом якоря отказал — якорь не пройден."""
db, counters = await _run(None, houses_ok=False, house_rows=HOUSE_ROWS[:1])
assert (counters.houses_attempted, counters.houses_enriched) == (2, 0)
assert _last_checkpoint(db) == [], (
f"якорь с 1 из 1 отказавшим домом помечен пройденным: {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_houses_db_query_failure_blocks_checkpoint() -> None:
"""Упал запрос домов: попытки считаются вместе с отказами, якорь не пройден."""
db, counters = await _run(None, houses_ok=True, house_rows=[])
assert _last_checkpoint(db) == [], (
f"якорь с упавшим запросом домов помечен пройденным: {_last_checkpoint(db)}"
)
assert (counters.houses_attempted, counters.houses_failed) == (2, 2)
class _FakeAvitoScraper:
"""Двойник AvitoScraper: SERP отдаёт лоты, дальше работает detail-фаза."""
def __init__(self, *_a: Any, **_kw: Any) -> None:
self._browser = None
self._cffi = None
async def fetch_around(self, *_a: Any, **_kw: Any) -> list[types.SimpleNamespace]:
return [types.SimpleNamespace(house_url=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
@pytest.mark.asyncio
async def test_avito_anchor_with_totally_failed_detail_is_not_checkpointed() -> None:
"""Сосед: у avito-свипа тот же гейт — фаза без единой удачи не даёт отметки.
Кода общего у свипов нет (две отдельные функции), общий только гейт. Прод-следа
у avito за 45 суток нет (0 прогонов против 4 у циана) тест сторожит вторую
копию проводки, а не найденный в данных дефект.
"""
from scraper_kit.orchestration import pipeline as pl
db = _FakeDb()
async def _boom(*_a: Any, **_kw: Any) -> Any:
raise RuntimeError("detail отдал 403")
with (
patch.object(pl, "AvitoScraper", _FakeAvitoScraper),
patch.object(pl, "AsyncSession", _FakeAsyncSession),
patch.object(pl, "save_listings", lambda *_a, **_kw: (0, 0)),
patch.object(pl, "fetch_detail", _boom),
patch.object(pl.runs, "is_cancelled", lambda *_a: False),
):
counters = await pl.run_avito_city_sweep(
db, # type: ignore[arg-type]
run_id=8404,
config=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,
),
matcher=MagicMock(),
enrichment=MagicMock(),
anchors=[ANCHOR_A],
enrich_houses=False,
enrich_imv=False,
detail_top_n=len(DETAIL_ROWS),
request_delay_sec=0.0,
)
assert counters.detail_attempted == len(DETAIL_ROWS)
assert counters.detail_enriched == 0
assert _last_checkpoint(db) == [], (
f"якорь с полностью отказавшей detail-фазой помечен пройденным: {_last_checkpoint(db)}"
)
@pytest.mark.asyncio
async def test_old_format_checkpoint_still_resumes() -> None:
"""(г) Чекпоинт СТАРОГО формата (плоский список имён якорей) читается как раньше."""
db, _ = await _run({"done_buckets": ["ekb-center"]}, houses_ok=True, run_id=8403)
assert f"{ANCHOR_A[0]},{ANCHOR_A[1]}" not in _FakeScraper.visited, (
"якорь из унаследованного чекпоинта всё-таки опрашивали"
)
assert _last_checkpoint(db) == ["ekb-center", "ekb-south"], (
"унаследованный якорь потерян — формат чекпоинта разъехался со старыми прогонами"
)

View file

@ -0,0 +1,90 @@
"""Проверка отмены не оставляет транзакцию открытой на сетевую фазу свипа (#3480).
Прод 17.09 07:48 UTC, pg_stat_activity БД tradein: pid 1933154 из tradein-scraper,
'idle in transaction' 1 ч 35 мин, последний запрос `SELECT status FROM scrape_runs
WHERE id = $1` (runs.is_cancelled), xact_start через 30 мс после старта
domclick_city_sweep_moskva 7344. За 14 суток 10 из 12 окон «транзакция > 1 ч»
совпали со свипами ДомКлика.
Сессия настоящая (SQLAlchemy + SQLite): значение `in_transaction()` то же, что видит
Postgres как 'idle in transaction'.
"""
from __future__ import annotations
import os
from types import SimpleNamespace
from typing import Any
from unittest.mock import MagicMock, patch
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
from scraper_kit.orchestration import pipeline as pl
from scraper_kit.orchestration import runs
from sqlalchemy import create_engine, text
from sqlalchemy.orm import Session
def _session() -> Session:
db = Session(create_engine("sqlite://"))
db.execute(text("CREATE TABLE scrape_runs (id INTEGER, status TEXT, counters TEXT)"))
db.execute(
text(
"INSERT INTO scrape_runs VALUES "
"(7333, 'done', NULL), (7344, 'running', NULL), (7345, 'cancelled', NULL)"
)
)
db.commit()
return db
def test_is_cancelled_closes_its_transaction() -> None:
db = _session()
assert runs.is_cancelled(db, 7344) is False
assert db.in_transaction() is False
assert runs.is_cancelled(db, 7345) is True
assert db.in_transaction() is False
async def test_domclick_sweep_http_phase_runs_outside_transaction() -> None:
db = _session()
seen: list[bool] = []
class _Scraper:
blocked = False
geo_filtered = fetch_errors = buckets_completed = buckets_total = 0
bucket_start_index = 0
completed_buckets: list[str] = [] # noqa: RUF012
def __init__(self, *_a: Any, **_kw: Any) -> None: ...
async def __aenter__(self) -> _Scraper:
return self
async def __aexit__(self, *_e: Any) -> None:
return None
async def fetch_city(self, **_kw: Any) -> list[Any]:
seen.append(db.in_transaction()) # момент сетевого обхода
return []
with (
patch.object(pl, "DomClickScraper", _Scraper),
patch.object(runs, "update_heartbeat", MagicMock()),
patch.object(runs, "mark_done", MagicMock()),
patch.object(runs, "mark_failed", MagicMock()),
patch.object(runs, "mark_banned", MagicMock()),
):
await pl.run_domclick_city_sweep(
db,
run_id=7344,
config=SimpleNamespace(browser_http_endpoint="http://x:9000"),
matcher=MagicMock(),
city_id=4,
pages=1,
request_delay_sec=0.0,
resume_run_id=7333, # резюм-SELECT тоже открывает транзакцию
)
assert seen == [False]

View file

@ -36,7 +36,7 @@ from __future__ import annotations
import asyncio
import logging
import random
from collections.abc import Callable
from collections.abc import Callable, Mapping
from contextlib import AsyncExitStack
from dataclasses import dataclass, field, fields
from datetime import date, timedelta
@ -1346,6 +1346,43 @@ async def run_avito_pipeline(
await browser_fetcher.__aexit__(None, None, None)
def _bucket_phase_totally_failed(before: Mapping[str, int], after: Mapping[str, int]) -> str | None:
"""Фаза бакета, где отказала КАЖДАЯ попытка (#3415). Имя фазы или None.
Считает ПРИРОСТ счётчиков за один бакет (якорь): counters у свипа run-level,
так что «фаза этого якоря» существует только как разница до/после.
Прод-факт, ради которого гейт заведён: `cian_city_sweep` 6179 (06.09)
`houses_attempted=30, houses_failed=30, houses_enriched=0`, статус `failed` по
run-level правилу (`_phase_totally_failed`, #2700), а в `done_buckets` записаны
все пять якорей. Прогон 6276 (07.09) унаследовал точку, пропустил все якоря и
отчитался `done` с `houses_attempted=0` зелёный по построению, потому что фаза
не исполнялась. Фикс #3396 (houses через пул прокси) простоял на проде сутки, ни
разу не отработав.
Пары ищутся В САМИХ counters (ключ `X_attempted` со спутником `X_failed`) тот же
приём и та же причина, что у `runs._phase_totally_failed`: зашитый список фаз это
ровно то место, куда забывают дописать новую. Фазы без счётчика попыток гейт НЕ
видит у avito houses роль attempted играет `unique_houses` (нет `houses_attempted`),
и притягивать её сюда значило бы блокировать отметку якоря из-за ОДНОГО упавшего
дома из десяти.
Порог одна попытка, а не три, как у run-level правила: тот ставит прогону
диагноз и обязан отсеивать шум, а этот решает «собирать ли якорь заново». Цена
ошибок несимметрична: лишний повтор якоря стоит нескольких запросов, ложная
отметка «пройден» теряет часть города навсегда и молча (прогон-то `done`).
"""
for key in sorted(after):
if not key.endswith("_attempted"):
continue
phase = key[: -len("_attempted")]
attempted = after[key] - before.get(key, 0)
failed = after.get(f"{phase}_failed", 0) - before.get(f"{phase}_failed", 0)
if attempted >= 1 and failed >= attempted:
return phase
return None
@dataclass
class CitySweepCounters:
"""Aggregate counters для full city sweep run."""
@ -1629,6 +1666,12 @@ async def run_avito_city_sweep(
# #3074: сбрасывается на КАЖДОЙ итерации — иначе один упавший якорь
# заразил бы все последующие, и чекпоинт не пополнялся бы вовсе.
_anchor_ok = True
# #3415: снимок счётчиков ДО якоря — по разнице видно, дала ли
# каждая фаза ЭТОГО якоря хоть одну удачную попытку (см.
# _bucket_phase_totally_failed). Тот же дефект, что у cian-свипа:
# detail-фаза, отказавшая на всех карточках без исключения, доходила
# до отметки «якорь пройден» наравне с целым якорем.
_before_anchor = counters.to_dict()
# Capture loop variables in default args (B023): prevents stale binding
# if the coroutine is scheduled after the loop variable changes.
@ -2137,6 +2180,22 @@ async def run_avito_city_sweep(
# его как пройденный значило бы, что следующий прогон его пропустит и
# объявления оттуда не соберутся НИКОГДА, причём молча: прогон
# завершится штатно. Тот же инвариант, что у combo в yandex-свипе.
#
# #3415: флага мало — фаза отказывает и БЕЗ исключения (каждая
# detail-карточка отдала 403, счётчик вырос, корутина завершилась
# штатно). Замер на проде за 45 суток: у avito-свипов такой якорь
# пока не встречался (0 прогонов), у cian — 4; гейт стоит здесь
# потому, что дефект один и тот же, а не по следу в данных.
_failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict())
if _failed_phase is not None:
logger.warning(
"city-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' "
"не дала ни одной удачной попытки; следующий прогон соберёт его заново",
run_id,
name,
_failed_phase,
)
_anchor_ok = False
if _anchor_ok:
_done_anchors.add(name)
runs.update_heartbeat(db, run_id, _ckpt())
@ -3313,6 +3372,10 @@ async def run_cian_city_sweep(
)
# ── Per-anchor phases под watchdog-таймаутом ────────────────────
# #3415: counters у свипа run-level — снимок ДО якоря даёт бухгалтерию
# его собственных фаз (разницей), по которой ниже решается, считать ли
# якорь пройденным.
_before_anchor = counters.to_dict()
anchor_lots: list[ScrapedLot] = []
_c_lat, _c_lon, _c_name = lat, lon, name
@ -3529,6 +3592,12 @@ async def run_cian_city_sweep(
len(nb_id_list),
exc,
)
# #3415: попытки считаем ВМЕСТЕ с отказами. Раньше здесь рос
# только `houses_failed`, и пара получалась несуществующей
# (attempted=0, failed=N): и run-level `_phase_totally_failed`
# (#2700), и гейт бакета сверяют `failed` с `attempted`, так что
# полный отказ фазы на этой ветке не видел ни один из них.
counters.houses_attempted += len(nb_id_list)
counters.houses_failed += len(nb_id_list)
try:
db.rollback()
@ -3720,7 +3789,28 @@ async def run_cian_city_sweep(
# Записать упавший якорь пройденным значило бы, что следующий прогон
# пропустит его навсегда — молча, потому что прогон завершится
# штатно, просто часть города не соберётся.
_done_anchors.add(name)
#
# #3415: потока управления мало. Фаза внутри якоря отказывает БЕЗ
# исключения — `fetch_newbuilding` возвращает None, счётчик растёт,
# `_cian_anchor_phases` завершается штатно, — и такой якорь доходил
# сюда наравне с целым. Дальше ровно тот же исход, что и у записанного
# упавшего: следующий прогон пропускает якорь и рапортует `done` с
# нулевой фазой. Форма отметки прежняя (плоский список имён): чекпоинт
# читают ещё три свипа и `scheduler._resume_decision`, а ключ
# `bucket:phase` потребовал бы новых читателей ради того же решения.
_failed_phase = _bucket_phase_totally_failed(_before_anchor, counters.to_dict())
if _failed_phase is None:
_done_anchors.add(name)
else:
logger.warning(
"cian-sweep run_id=%d: якорь %s НЕ помечен пройденным — фаза '%s' "
"не дала ни одной удачной попытки; следующий прогон соберёт его заново",
run_id,
name,
_failed_phase,
)
# Heartbeat пишем в любом случае: без него reap_zombies посчитает живой
# прогон мёртвым, а jsonb-мерж сохранит унаследованную часть точки.
runs.update_heartbeat(
db, run_id, {**counters.to_dict(), "done_buckets": sorted(_done_anchors)}
)
@ -4902,10 +4992,13 @@ async def run_domclick_city_sweep(
)
lots: list[ScrapedLot] = []
# #2406: собранное сохраняется ОДНИМ save_listings после всех корзин, поэтому
# корзина считается пройденной только если фаза дошла до конца (см. чекпоинт ниже).
_saved = False
async def _domclick_phase() -> None:
"""Единственная citywide-фаза: fetch_city + save."""
nonlocal lots
nonlocal lots, _saved
async with DomClickScraper(
config,
proxy_provider=proxy_provider,
@ -4951,6 +5044,7 @@ async def run_domclick_city_sweep(
)
counters.lots_inserted += inserted
counters.lots_updated += updated
_saved = True
try:
await asyncio.wait_for(_domclick_phase(), timeout=_sweep_timeout)
@ -5004,7 +5098,12 @@ async def run_domclick_city_sweep(
# #3118: чекпоинт = унаследованное завершённое в этом прогоне.
# Пишем heartbeat'ом СЕЙЧАС (мерж jsonb) — финализаторы ключ не
# затирают, и оборванный болезнью финализации прогон его не теряет.
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
# #2406: только если лоты сохранены. Снятая watchdog'ом или упавшая фаза
# не дошла до save_listings — у её корзин в БД ноль строк, а отметка
# «пройдена» заставила бы следующий прогон пропустить их навсегда
# (механизм разобран в миграции 308, из-за него выключены свипы 77/50).
if _saved:
_checkpoint = sorted(skip_buckets | set(_s.completed_buckets))
runs.update_heartbeat(db, run_id, _payload())
counters.bucket_start_index = _s.bucket_start_index

View file

@ -46,6 +46,7 @@ from sqlalchemy import text
from sqlalchemy.orm import Session
from scraper_kit.orchestration.run_context import current_run_id
from scraper_kit.proxy_errors import caused_by_no_proxy
try: # sentry опционален — kit не тянет его в зависимостях
import sentry_sdk
@ -815,11 +816,19 @@ def honors_cancel(source: str) -> bool:
def is_cancelled(db: Session, run_id: int) -> bool:
"""Проверить status='cancelled' (cooperative cancel в long-running pipeline)."""
"""Проверить status='cancelled' (cooperative cancel в long-running pipeline).
#3480: SELECT открывает транзакцию (autobegin), а вызывающие сразу уходят в сетевую
фазу. Без commit сессия висела 'idle in transaction' весь обход и держала горизонт
vacuum всей БД: прод 17.09, domclick_city_sweep_moskva 7344 транзакция 1.6 ч с
последним запросом ровно этим SELECT; 10 из 12 окон > 1 ч за 14 суток свипы
ДомКлика. Закрываем здесь: проверку отмены зовут перед каждым долгим await.
"""
row = db.execute(
text("SELECT status FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
db.commit()
return row is not None and row.status == "cancelled"
@ -1009,6 +1018,37 @@ def mark_banned(
_alert_on_run_id(db, run_id)
def mark_crashed(db: Session, run_id: int, exc: BaseException) -> None:
"""Финализировать прогон, чья задача вылетела исключением наружу (#1940).
Зовут те, кто запускает пайплайн: планировщик (`scheduler._dispatch`) и ручные
запуски из админки. Исключение, брошенное ДО try-финализатора пайплайна (прод:
пустой пул прокси в `BrowserFetcher.__aenter__` через 10 мс после claim),
оставляло строку 'running' без пульса до zombie-reaper'а через 6 ч.
Строку, которую пайплайн уже финализировал сам (обычный путь: mark_failed и
raise), не трогаем и ПОВТОРНО mark_* не зовём: UPDATE был бы no-op, но mark_*
всё равно позвал бы `_alert_on_run_id` стрик тот же, и на вехе лестницы
(`_streak_alert_due`) в Sentry ушёл бы дубль, плюс WARNING «no-op» на каждом
таком падении. Best-effort: сбой финализации логируется, не бросается.
"""
try:
db.rollback() # сессия упавшей задачи может быть в aborted-транзакции
row = db.execute(
text("SELECT status FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
db.rollback() # чтение не держит транзакцию (#3480)
if row is None or row.status != "running":
return
if caused_by_no_proxy(exc):
mark_banned(db, run_id, f"crashed: {exc}", {}, ban_kind=BAN_KIND_INFRA)
else:
mark_failed(db, run_id, f"crashed: {type(exc).__name__}: {exc}", {})
except Exception:
logger.exception("mark_crashed: не удалось финализировать run_id=%d", run_id)
def _dominant_ban_kind(census: Mapping[str, int]) -> str:
"""Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178).

View file

@ -1024,8 +1024,11 @@ async def _dispatch(
run_db = ctx.session_factory()
try:
await handler.job(run_db, run_id, params, ctx)
except Exception:
except Exception as exc:
logger.exception("scheduler: %s crashed run_id=%d", handler.log_name, run_id)
# #1940: исключение мимо try-финализатора хендлера оставляло строку
# 'running' до zombie через 6 ч. См. runs.mark_crashed.
ctx.runs.mark_crashed(run_db, run_id, exc)
finally:
run_db.close()