fix(tradein/scraper): сводка «что сейчас не собирает» — лестница вех не видит стрик 0 (#2670) #2806
2 changed files with 347 additions and 2 deletions
199
tradein-mvp/backend/tests/test_2670_stale_source_digest.py
Normal file
199
tradein-mvp/backend/tests/test_2670_stale_source_digest.py
Normal file
|
|
@ -0,0 +1,199 @@
|
|||
"""#2670 (остаток): лестница напоминаний не отвечает на вопрос «что сломано сейчас».
|
||||
|
||||
#2720 вылечил «алерт ровно один раз за серию»: теперь вехи 3, 6, 12, 24, 48… Но лестница
|
||||
шагает по ПОДРЯД ИДУЩИМ завершённым failed/banned прогонам, а на проде 2026-08-10 три
|
||||
самых залежавшихся источника из шести просроченных ей недоступны — и лишь один из трёх
|
||||
из-за редких вех:
|
||||
|
||||
┌────────────────────────────┬────────┬───────┬──────────────────────────────────────┐
|
||||
│ источник │ стрик │ сут. │ когда напомнит лестница │
|
||||
├────────────────────────────┼────────┼───────┼──────────────────────────────────────┤
|
||||
│ avito_full_load_exhaustive │ 0 │ 49.5 │ никогда: 5 банов обнулил 'cancelled' │
|
||||
│ cian_history_backfill │ 0 │ 42.1 │ никогда: прогонов нет с 30.06 │
|
||||
│ avito_full_load │ 31 │ 37.7 │ веха 48 → +17 прогонов × 7 сут = 119 │
|
||||
│ avito_detail_backfill │ 5 │ 5.2 │ веха 6 → завтра │
|
||||
│ domclick_city_sweep │ 5 │ 5.1 │ веха 6 → завтра │
|
||||
│ domclick_detail_backfill │ 4 │ 5.0 │ веха 6 → послезавтра │
|
||||
└────────────────────────────┴────────┴───────┴──────────────────────────────────────┘
|
||||
|
||||
Уплотнение вех (3,4,5,6…) чинит ТОЛЬКО третью строку: у первых двух стрик равен нулю,
|
||||
уплотнять нечего — «замолчал» там означает «перестал производить прогоны», а не «серия
|
||||
длиннее последней вехи». Поэтому остаток задачи закрывает сводка, считающая КАЛЕНДАРНЫЙ
|
||||
возраст последнего успеха, а не длину серии.
|
||||
|
||||
Фальсификация: на коде до этой правки `emit_stale_digest`/`stale_sources` не существует
|
||||
(ImportError на сборе тестов) — сводки нет ни в каком виде. Тест
|
||||
`test_ladder_is_silent_for_the_worst_two` — КОНТРОЛЬ: он зелёный и до, и после правки и
|
||||
показывает ровно то, чего сводка не заменяет, а добавляет: лестница на этих двух молчит.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
from datetime import UTC, datetime, timedelta
|
||||
from types import SimpleNamespace
|
||||
from typing import Any
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
import pytest
|
||||
|
||||
os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db")
|
||||
|
||||
from scraper_kit.orchestration import runs as kit_runs
|
||||
from scraper_kit.orchestration import scheduler as sched
|
||||
|
||||
NOW = datetime(2026, 8, 10, 8, 0, tzinfo=UTC)
|
||||
|
||||
|
||||
def _row(source: str, interval_days: Any, age_days: float, never_ok: bool = False) -> Any:
|
||||
"""Строка `_STALE_SOURCES_SQL`: last_ok уже схлопнут в `since` через COALESCE."""
|
||||
return SimpleNamespace(
|
||||
source=source,
|
||||
interval_days=interval_days,
|
||||
since=NOW - timedelta(days=age_days),
|
||||
never_ok=never_ok,
|
||||
)
|
||||
|
||||
|
||||
# Снимок прода 2026-08-10 08:00 UTC: все 52 включённых расписания не влезают, взяты все
|
||||
# просроченные + четыре контрольных, каждое из которых мимо порога по своей причине.
|
||||
PROD_ROWS = [
|
||||
_row("cian_history_backfill", None, 42.1), # такт по умолчанию (daily)
|
||||
_row("avito_full_load_exhaustive", 7, 49.5),
|
||||
_row("avito_full_load", 7, 37.7),
|
||||
_row("avito_detail_backfill", None, 5.2),
|
||||
_row("domclick_city_sweep", None, 5.1),
|
||||
_row("domclick_detail_backfill", None, 5.0),
|
||||
# ── контроль: НЕ просрочены ──
|
||||
_row("rosreestr_quarter_poll", 28, 24.0), # 24 сут при такте 28 — норма
|
||||
_row("sber_index_pull", 7, 4.1),
|
||||
_row("avito_city_sweep", None, 1.1),
|
||||
_row("proxy_healthcheck", None, 0.02),
|
||||
]
|
||||
|
||||
# Порядок — по числу ПРОПУЩЕННЫХ ТАКТОВ (age/interval), а не по календарю: 42 суток
|
||||
# у суточного backfill'а = 42 пропущенных такта, 49.5 у недельного = 7.
|
||||
PROD_STALE = [
|
||||
"cian_history_backfill", # 42.1 / 1
|
||||
"avito_full_load_exhaustive", # 49.5 / 7 = 7.07
|
||||
"avito_full_load", # 37.7 / 7 = 5.39
|
||||
"avito_detail_backfill", # 5.2 / 1
|
||||
"domclick_city_sweep", # 5.1 / 1
|
||||
"domclick_detail_backfill", # 5.0 / 1
|
||||
]
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True)
|
||||
def _reset_digest_clock() -> Any:
|
||||
"""Выпуск сводки помнится в памяти модуля — сбрасываем между тестами."""
|
||||
sched._last_stale_digest_at = None
|
||||
yield
|
||||
sched._last_stale_digest_at = None
|
||||
|
||||
|
||||
def _db(rows: list[Any]) -> MagicMock:
|
||||
db = MagicMock()
|
||||
db.execute.return_value.fetchall.return_value = rows
|
||||
return db
|
||||
|
||||
|
||||
# ── 1. Чистая логика порога ──────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_stale_sources_names_exactly_the_prod_six() -> None:
|
||||
"""Шесть просроченных из десяти, порядок — по числу пропущенных ТАКТОВ, не суток."""
|
||||
stale = sched.stale_sources(PROD_ROWS, NOW)
|
||||
assert [s.source for s in stale] == PROD_STALE
|
||||
|
||||
|
||||
def test_quarterly_source_is_not_stale_at_24_days() -> None:
|
||||
"""Порог считается в тактах: 24 сут для 28-суточного poll'а — не просрочка."""
|
||||
assert sched.stale_sources([_row("rosreestr_quarter_poll", 28, 24.0)], NOW) == []
|
||||
# …а 85 суток (>3×28) — уже просрочка.
|
||||
assert [s.source for s in sched.stale_sources([_row("q", 28, 85.0)], NOW)] == ["q"]
|
||||
|
||||
|
||||
@pytest.mark.parametrize("raw", [None, "null", "", "abc", 0, -5])
|
||||
def test_broken_interval_falls_back_to_daily(raw: Any) -> None:
|
||||
"""`interval_days: null` и мусор → такт 1 сут, как у compute_next_run_at."""
|
||||
assert sched._schedule_interval_days(raw) == 1
|
||||
|
||||
|
||||
def test_never_successful_source_is_reported_with_a_flag() -> None:
|
||||
"""Расписание без единого 'done' считается от created_at и помечается явно."""
|
||||
(only,) = sched.stale_sources([_row("brand_new", 1, 9.0, never_ok=True)], NOW)
|
||||
assert only.never_ok is True
|
||||
|
||||
|
||||
# ── 2. Выпуск сводки ─────────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def test_digest_emits_one_event_listing_all_stale_sources() -> None:
|
||||
sentry = MagicMock()
|
||||
with patch.object(sched, "sentry_sdk", sentry):
|
||||
stale = sched.emit_stale_digest(_db(PROD_ROWS), now=NOW)
|
||||
assert [s.source for s in stale] == PROD_STALE
|
||||
sentry.capture_message.assert_called_once()
|
||||
msg = sentry.capture_message.call_args[0][0]
|
||||
assert msg.startswith("6 scraper sources are stale")
|
||||
for name in PROD_STALE:
|
||||
assert name in msg
|
||||
assert "rosreestr_quarter_poll" not in msg
|
||||
|
||||
|
||||
def test_digest_covers_the_two_sources_the_ladder_cannot_reach() -> None:
|
||||
"""Главное свойство: стрик 0 не мешает сводке — она меряет календарь, а не серию."""
|
||||
sentry = MagicMock()
|
||||
with patch.object(sched, "sentry_sdk", sentry):
|
||||
stale = sched.emit_stale_digest(_db(PROD_ROWS), now=NOW)
|
||||
zero_streak = {"avito_full_load_exhaustive", "cian_history_backfill"}
|
||||
assert zero_streak <= {s.source for s in stale}
|
||||
|
||||
|
||||
def test_digest_is_quiet_when_everything_is_fresh() -> None:
|
||||
sentry = MagicMock()
|
||||
fresh = [_row("avito_city_sweep", None, 1.1), _row("sber_index_pull", 7, 4.1)]
|
||||
with patch.object(sched, "sentry_sdk", sentry):
|
||||
assert sched.emit_stale_digest(_db(fresh), now=NOW) == []
|
||||
sentry.capture_message.assert_not_called()
|
||||
|
||||
|
||||
def test_digest_is_daily_not_per_tick() -> None:
|
||||
"""Планировщик тикает раз в минуту; сводка обязана выходить раз в сутки."""
|
||||
sentry = MagicMock()
|
||||
db = _db(PROD_ROWS)
|
||||
with patch.object(sched, "sentry_sdk", sentry):
|
||||
sched.emit_stale_digest(db, now=NOW)
|
||||
sched.emit_stale_digest(db, now=NOW + timedelta(minutes=1))
|
||||
sched.emit_stale_digest(db, now=NOW + timedelta(hours=23))
|
||||
assert sentry.capture_message.call_count == 1
|
||||
sched.emit_stale_digest(db, now=NOW + timedelta(hours=24, minutes=1))
|
||||
assert sentry.capture_message.call_count == 2
|
||||
|
||||
|
||||
def test_digest_failure_never_breaks_the_tick() -> None:
|
||||
"""Сводка — best-effort: упавший запрос не имеет права уронить тик планировщика."""
|
||||
db = MagicMock()
|
||||
db.execute.side_effect = RuntimeError("db down")
|
||||
with patch.object(sched, "sentry_sdk", MagicMock()):
|
||||
assert sched.emit_stale_digest(db, now=NOW) == []
|
||||
|
||||
|
||||
# ── 3. Контроль: что именно сводка ДОБАВЛЯЕТ к лестнице ──────────────────────
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("name", "streak"),
|
||||
[("avito_full_load_exhaustive", 0), ("cian_history_backfill", 0), ("avito_full_load", 31)],
|
||||
)
|
||||
def test_ladder_is_silent_for_the_worst_two(name: str, streak: int) -> None:
|
||||
"""КОНТРОЛЬ (зелёный и до правки): у трёх худших источников лестница молчит.
|
||||
|
||||
Стрик 0 — прогонов нет / серию обнулил 'cancelled'; стрик 31 — между вехами 24 и 48.
|
||||
"""
|
||||
rows = [SimpleNamespace(status="banned") for _ in range(streak)]
|
||||
rows += [SimpleNamespace(status="done") for _ in range(3)]
|
||||
sentry = MagicMock()
|
||||
with patch.object(kit_runs, "sentry_sdk", sentry):
|
||||
kit_runs._alert_if_consecutive_failures(_db(rows), name)
|
||||
sentry.capture_message.assert_not_called()
|
||||
|
|
@ -17,8 +17,11 @@ Kit-native (sweep-оркестраторы, уже перенесённые в `
|
|||
осталось в `app` (rosreestr_dkp / sber_index / deactivate_stale / *_backfill / …),
|
||||
инжектируются извне как `Handler` через `build_registry(product_handlers=...)`.
|
||||
|
||||
Боевой рантайм (scraper-контейнер) по-прежнему крутит старый `app.services.scheduler` —
|
||||
это COPY, не MOVE. Переключение — отдельный поздний strangler-шаг.
|
||||
Боевой рантайм (scraper-контейнер, `python -m app.scheduler_main`) крутит ИМЕННО ЭТОТ
|
||||
loop: #2397 Part C удалил legacy-ветку `app.services.scheduler.scheduler_loop`, и
|
||||
`_run_kit_scheduler` остался единственным путём. Строка «по-прежнему крутит старый
|
||||
app.services.scheduler» жила здесь после того, как перестала быть правдой, и посылала
|
||||
правку сторожей не в тот файл (#2670).
|
||||
|
||||
Критичная concurrency-логика (`_claim_run` advisory-lock + double-check, `reap_zombies`
|
||||
порог, heartbeat, SIGTERM-drain) перенесена ДОСЛОВНО — тот же SQL, то же ветвление.
|
||||
|
|
@ -36,6 +39,11 @@ from typing import TYPE_CHECKING, Any
|
|||
|
||||
from sqlalchemy import text
|
||||
|
||||
try: # sentry опционален — kit standalone-импортируем, sentry-sdk не в его зависимостях
|
||||
import sentry_sdk
|
||||
except ImportError: # pragma: no cover - в проде backend-env sentry_sdk присутствует
|
||||
sentry_sdk = None # type: ignore[assignment]
|
||||
|
||||
from scraper_kit.orchestration import runs as _kit_runs
|
||||
from scraper_kit.orchestration.pipeline import (
|
||||
get_city_anchors,
|
||||
|
|
@ -78,6 +86,141 @@ SKIP_CONCURRENT_CLAIM = "concurrent_claim"
|
|||
SKIP_RUNNING_UNDER_LOCK = "running_appeared_under_lock"
|
||||
SKIP_UNKNOWN_SOURCE = "unknown_source"
|
||||
|
||||
# ── сводка «что сейчас не собирает» (#2670, второй пункт задачи) ─────────────
|
||||
# Лестница напоминаний из #2720 считает ПОДРЯД ИДУЩИЕ неудачные ПРОГОНЫ. Прод
|
||||
# 2026-08-10: шесть источников не имели успешного прогона дольше 3× своего такта, и
|
||||
# трое худших из них лестнице недоступны ПО ПОСТРОЕНИЮ, а не из-за редких вех:
|
||||
#
|
||||
# cian_history_backfill 42.1 сут без успеха, стрик 0 — с 30.06 прогонов нет
|
||||
# вовсе (сегодняшний единственный — 'skipped',
|
||||
# cian_cookies_expired), а лестница шагает только по
|
||||
# завершённым failed/banned;
|
||||
# avito_full_load_exhaustive 49.5 сут без успеха, стрик 0 — 5 банов подряд обнулил
|
||||
# один 'cancelled' 09.08 (деплой убил бегущий прогон);
|
||||
# avito_full_load 37.7 сут без успеха, стрик 31 — веха 48 при такте
|
||||
# interval_days=7 наступит через 17 прогонов ≈ 119 суток.
|
||||
#
|
||||
# Уплотнение вех чинит только третий случай: у первых двух стрик равен нулю, уплотнять
|
||||
# нечего. Поэтому сводка не «ещё один сторож помельче», а ЕДИНСТВЕННЫЙ ответ на вопрос
|
||||
# «что сломано сейчас»: она считает КАЛЕНДАРНЫЙ возраст последнего успеха, поэтому
|
||||
# видит и молчащий источник, и обнулённый стрик, и редкую веху. Лестница остаётся как
|
||||
# была — она отвечает на другой вопрос («что сломалось только что») и стоит дёшево.
|
||||
STALE_DIGEST_INTERVAL_FACTOR = 3
|
||||
STALE_DIGEST_PERIOD_H = 24
|
||||
|
||||
# Возраст последнего УСПЕШНОГО ('done') прогона на каждое включённое расписание.
|
||||
# COALESCE(last_ok, created_at): у расписания без единого успеха отсчёт идёт от его
|
||||
# создания — иначе «никогда не собирал» выглядело бы как «нет данных, судить нечем».
|
||||
_STALE_SOURCES_SQL = text("""
|
||||
SELECT sch.source,
|
||||
sch.default_params->>'interval_days' AS interval_days,
|
||||
COALESCE(
|
||||
(SELECT max(r.finished_at) FROM scrape_runs r
|
||||
WHERE r.source = sch.source AND r.status = 'done'),
|
||||
sch.created_at
|
||||
) AS since,
|
||||
(NOT EXISTS (SELECT 1 FROM scrape_runs r
|
||||
WHERE r.source = sch.source AND r.status = 'done')) AS never_ok
|
||||
FROM scrape_schedules sch
|
||||
WHERE sch.enabled
|
||||
""")
|
||||
|
||||
# ponytail: последний выпуск сводки помнится В ПАМЯТИ процесса, поэтому рестарт
|
||||
# scheduler'а (деплой) даёт лишний выпуск. Осознанный размен: альтернатива — таблица
|
||||
# состояния (миграция) ради анти-спама у механизма, который и заводится ПРОТИВ
|
||||
# молчания. Понадобится точность — переносить в scrape_runs строкой своего source'а.
|
||||
_last_stale_digest_at: datetime | None = None
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class StaleSource:
|
||||
"""Источник, не собиравший дольше STALE_DIGEST_INTERVAL_FACTOR× своего такта."""
|
||||
|
||||
source: str
|
||||
interval_days: int
|
||||
age_days: float
|
||||
never_ok: bool
|
||||
|
||||
|
||||
def _schedule_interval_days(raw: Any) -> int:
|
||||
"""default_params.interval_days → такт в сутках; всё непонятное → 1 (как у claim'а).
|
||||
|
||||
Тот же дефолт, что у `compute_next_run_at` (interval_days=1 == daily): порог сводки
|
||||
обязан считаться из ТОГО ЖЕ числа, которым расписание себя двигает, иначе «просрочен»
|
||||
будет мерить не тот такт. `"interval_days": null` в jsonb приезжает сюда None.
|
||||
"""
|
||||
try:
|
||||
return max(1, int(raw))
|
||||
except (TypeError, ValueError):
|
||||
return 1
|
||||
|
||||
|
||||
def stale_sources(rows: list[Any], now: datetime) -> list[StaleSource]:
|
||||
"""Чистая часть сводки: какие расписания просрочены и на сколько (свежие — внизу).
|
||||
|
||||
Просрочка меряется в ТАКТАХ, а не в сутках: у rosreestr_quarter_poll такт 28 суток,
|
||||
и 24 суток без сбора для него норма, а для суточного domclick_city_sweep — авария.
|
||||
Сортировка по числу пропущенных тактов, а не по календарю, по той же причине.
|
||||
"""
|
||||
stale: list[StaleSource] = []
|
||||
for row in rows:
|
||||
interval = _schedule_interval_days(row.interval_days)
|
||||
age_days = (now - row.since).total_seconds() / 86400.0
|
||||
if age_days > STALE_DIGEST_INTERVAL_FACTOR * interval:
|
||||
stale.append(
|
||||
StaleSource(
|
||||
source=row.source,
|
||||
interval_days=interval,
|
||||
age_days=age_days,
|
||||
never_ok=bool(row.never_ok),
|
||||
)
|
||||
)
|
||||
stale.sort(key=lambda s: s.age_days / s.interval_days, reverse=True)
|
||||
return stale
|
||||
|
||||
|
||||
def emit_stale_digest(db: Session, *, now: datetime | None = None) -> list[StaleSource]:
|
||||
"""Раз в STALE_DIGEST_PERIOD_H часов — одно событие «что сейчас не собирает».
|
||||
|
||||
Возвращает список просроченных источников (пустой — либо всё свежо, либо выпуск ещё
|
||||
не подошёл по времени). Best-effort, как и оба сторожа в runs.py: сводка не имеет
|
||||
права уронить тик планировщика.
|
||||
"""
|
||||
global _last_stale_digest_at
|
||||
now = now or datetime.now(UTC)
|
||||
if _last_stale_digest_at is not None and now - _last_stale_digest_at < timedelta(
|
||||
hours=STALE_DIGEST_PERIOD_H
|
||||
):
|
||||
return []
|
||||
try:
|
||||
stale = stale_sources(list(db.execute(_STALE_SOURCES_SQL).fetchall()), now)
|
||||
_last_stale_digest_at = now
|
||||
if not stale:
|
||||
logger.info("scheduler: stale-digest — просроченных источников нет (#2670)")
|
||||
return []
|
||||
details = ", ".join(
|
||||
f"{s.source} {s.age_days:.1f}d/{s.interval_days}d"
|
||||
+ (" (успеха не было ни разу)" if s.never_ok else "")
|
||||
for s in stale
|
||||
)
|
||||
logger.error(
|
||||
"scheduler: %d источников не собирают дольше %d× своего такта — %s (#2670)",
|
||||
len(stale),
|
||||
STALE_DIGEST_INTERVAL_FACTOR,
|
||||
details,
|
||||
)
|
||||
if sentry_sdk is not None:
|
||||
sentry_sdk.capture_message(
|
||||
f"{len(stale)} scraper sources are stale (no successful run for more than "
|
||||
f"{STALE_DIGEST_INTERVAL_FACTOR}× their schedule interval): {details}",
|
||||
level="error",
|
||||
)
|
||||
return stale
|
||||
except Exception:
|
||||
logger.exception("scheduler: stale-digest failed")
|
||||
return []
|
||||
|
||||
|
||||
# ── типы job/handler ─────────────────────────────────────────────────────────
|
||||
# Job получает свежую сессию (открыта `_dispatch`), run_id, params и весь контекст
|
||||
# (config/matcher/enrichment/runs) — чтобы иметь доступ к инжектированным зависимостям.
|
||||
|
|
@ -759,6 +902,9 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
|
|||
try:
|
||||
# Reap zombies first
|
||||
reap_zombies(db)
|
||||
# #2670: раз в сутки — сводка «что сейчас не собирает». Календарная, а
|
||||
# не по стрику: два самых залежавшихся источника прода имеют стрик 0.
|
||||
emit_stale_digest(db)
|
||||
# Process due schedules
|
||||
due = get_due_schedules(db)
|
||||
for sch in due:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue