fix(tradein/scraper): сводка «что сейчас не собирает» — лестница вех не видит стрик 0 (#2670) #2806

Merged
bot-backend merged 1 commit from fix/2670-stale-digest into main 2026-08-10 08:29:27 +00:00
2 changed files with 347 additions and 2 deletions

View 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()

View file

@ -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: