Сбор МЕРА: упавший прогон не висит zombie 6 часов, провалившаяся фаза не попадает в чекпоинт, проверка отмены не держит транзакцию часами #3561
8 changed files with 1030 additions and 11 deletions
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
202
tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py
Normal file
202
tradein-mvp/backend/tests/test_1940_crashed_run_is_finalized.py
Normal 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)
|
||||
248
tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py
Normal file
248
tradein-mvp/backend/tests/test_2406_sweep_edge_cases.py
Normal 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"
|
||||
|
|
@ -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"], (
|
||||
"унаследованный якорь потерян — формат чекпоинта разъехался со старыми прогонами"
|
||||
)
|
||||
|
|
@ -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]
|
||||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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).
|
||||
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue