fix(tradein/scraper): пропуск расписания пишет строку прогона со статусом skipped (#2658)
All checks were successful
CI / changes (pull_request) Successful in 7s
CI Trade-In / changes (pull_request) Successful in 8s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 2m41s

Пропуск наступившего окна был немым: logger + сдвиг next_run_at, ни строки в
scrape_runs, ни изменения last_run_at. cian_history_backfill так простоял 37 дней
на протухших куках Циана и снаружи выглядел работающим — next_run_at исправно
двигался вперёд, а docker-логи с warning'ом терялись на каждом редеплое.

Статус 'skipped' заведён ещё миграцией 015 и локализован во фронте («пропущено»),
но в проде имел 0 строк — механизм построен и ни разу не использован. Задействуем
его во всех пяти местах, где расписание пропускалось без следа: kit `_claim_run`
(already_running / concurrent_claim / running_appeared_under_lock), kit
`scheduler_loop` (unknown_source) и продуктовый cian `pre_claim`. Причина — слаг в
`error`, по нему «нет кук» отличается от «уже бежит» запросом, а не грепом логов.

Подряд идущие одинаковые пропуски схлопываются в одну строку со счётчиком
`counters.skips`: «уже бежит» и «неизвестный source» не двигают next_run_at и
иначе плодили бы строку каждый тик (60 с).

Алерт про куки жил в недостижимой ветке: он стоял там, где verify_session вернул
None, а на протухших куках load_session сам фильтрует expires_at_estimate > NOW()
и отдаёт None ещё в первой, немой ветке. Теперь алерт в обеих ветках и через
logger.error — в scraper-контейнере GlitchTip поднят с LoggingIntegration
(event_level=ERROR), поэтому прежний capture_message(level="warning") событием не
становился. Плюс предупреждение ЗАРАНЕЕ (COOKIE_EXPIRY_WARN_DAYS=5) в том же
pre_claim: обновление кук — ручная операция, алерт по факту протухания приходит,
когда сбор уже встал.

Монитор нулевых прогонов (#2625) не трогаем: обе alert-выборки отбирают
failed/banned/done/cancelled, поэтому 'skipped' в стрик не попадает и его не
прерывает — пропуск не «прогон вернул ноль лотов», смешивать нельзя.
This commit is contained in:
bot-backend 2026-08-05 22:35:51 +05:00
parent 0a001ee3f7
commit 0b54b96984
6 changed files with 728 additions and 22 deletions

View file

@ -8,6 +8,7 @@ from __future__ import annotations
import json
import logging
from datetime import datetime
from typing import Any
from curl_cffi.requests import AsyncSession
@ -23,6 +24,11 @@ from app.core.config import settings
logger = logging.getLogger(__name__)
# За сколько дней до протухания кук предупреждать (#2658). Обновление кук — РУЧНАЯ
# операция (залить дамп через админку), человеку нужен запас: алерт по факту протухания
# приходит, когда сбор уже встал. save_session ставит ttl 30 дней, так что окно широкое.
COOKIE_EXPIRY_WARN_DAYS = 5
# Cookies критичные для Cian auth — фильтр перед сохранением.
# Список обновлён по реальному DevTools-дампу из logged-in сессии cian.ru (2026-05-23).
# Старые записи оставлены как fallback (backward compat).
@ -294,6 +300,29 @@ def load_session(db: Session) -> dict[str, str] | None:
return cookies
def session_expires_at(db: Session) -> datetime | None:
"""Когда протухают самые свежезагруженные куки — БЕЗ фильтра валидности (#2658).
`load_session` отбирает только ещё валидные записи (expires_at_estimate > NOW()) и на
протухших отдаёт None вызывающий не мог отличить «кук никогда не загружали» от
«протухли позавчера» и не мог предупредить ЗАРАНЕЕ. Здесь фильтра нет: None означает
ровно «записей нет вовсе».
"""
row = db.execute(
text(
"""
SELECT expires_at_estimate FROM cian_session_cookies
ORDER BY uploaded_at DESC
LIMIT 1
"""
)
).first()
if row is None:
return None
expires_at: datetime | None = row[0]
return expires_at
def mark_session_invalid(db: Session, account_user_id: int) -> None:
"""Flag session как expired/invalid (например после 401 во время scrape)."""
db.execute(

View file

@ -21,8 +21,10 @@ from __future__ import annotations
import asyncio
import logging
from datetime import UTC, datetime, timedelta
from typing import TYPE_CHECKING, Any
from scraper_kit.orchestration import runs as kit_runs
from scraper_kit.orchestration.scheduler import (
Handler,
reschedule_after_minutes,
@ -39,39 +41,95 @@ logger = logging.getLogger(__name__)
# ── cian_history_backfill — cookie-gated backfill ────────────────────────────
# Машиночитаемые причины пропуска (#2658) — пишутся в scrape_runs.error строки
# со status='skipped'. Отделены от kit-причин (already_running и т.п.): по слагу
# видно, встал ли сбор из-за кук или из-за конкурентного прогона.
SKIP_CIAN_COOKIES_MISSING = "cian_cookies_missing"
SKIP_CIAN_COOKIES_EXPIRED = "cian_cookies_expired"
SKIP_CIAN_COOKIES_INVALID = "cian_cookies_invalid"
def _alert_cian_cookies(source: str, detail: str) -> None:
"""Громкий алерт «сбор встал из-за кук» — logger.error, НЕ capture_message(warning).
В scraper-контейнере GlitchTip поднят с LoggingIntegration(event_level=ERROR)
(scheduler_main.py) ERROR-запись сама становится событием, а прежний
`capture_message(..., level="warning")` до этого уровня не дотягивал (и стоял в
недостижимой ветке, см. докстринг _cian_pre_claim). Заодно причина остаётся в
docker-логах и в строке scrape_runs, которая переживает редеплой.
"""
logger.error(
"scheduler: %s пропущен — %s. Перезалейте куки Циана через админку "
"(до этого backfill истории стоит)",
source,
detail,
)
async def _cian_pre_claim(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> bool:
"""Pre-claim gate: проверить наличие/валидность cian-cookies ДО claim (#1522).
Cookies отсутствуют/протухли defer next_run_at на следующее окно и skip
(иначе get_due_schedules переотбирает schedule каждые 60с и verify_session
долбит Cian круглосуточно). Дословно из боевого trigger_cian_backfill_run.
"""
import sentry_sdk
Cookies отсутствуют/протухли пишем строку прогона status='skipped' с причиной,
двигаем next_run_at на следующее окно и skip (иначе get_due_schedules переотбирает
schedule каждые 60с и verify_session долбит Cian круглосуточно).
from app.services.cian_session import load_session, verify_session
#2658 — что было не так. Первая ветка (load_session вернул None) молчала: warning в
docker-лог, сдвиг next_run_at, `return False`. Ни строки в scrape_runs, ни изменения
last_run_at снаружи 37 дней простоя выглядели как «всё по расписанию». Sentry-алерт
стоял во ВТОРОЙ ветке (verify_session вернул None), до которой на протухших куках
исполнение не доходит НИКОГДА: load_session сам фильтрует expires_at_estimate > NOW()
и отдаёт None ещё в первой. Теперь громко в обеих + предупреждение ЗАРАНЕЕ, пока куки
ещё валидны (COOKIE_EXPIRY_WARN_DAYS) обновление кук ручное, ему нужен запас.
"""
from app.services.cian_session import (
COOKIE_EXPIRY_WARN_DAYS,
load_session,
session_expires_at,
verify_session,
)
source: str = schedule_row["source"]
now = datetime.now(tz=UTC)
cookies = load_session(db)
if cookies is None:
logger.warning("scheduler: cian_history_backfill skipped — no valid session cookies in DB")
expires_at = session_expires_at(db)
if expires_at is None:
reason, detail = SKIP_CIAN_COOKIES_MISSING, "кук Циана нет в БД"
elif expires_at <= now:
reason = SKIP_CIAN_COOKIES_EXPIRED
detail = (
f"куки Циана протухли {expires_at:%Y-%m-%d} ({(now - expires_at).days} дн. назад)"
)
else:
reason = SKIP_CIAN_COOKIES_INVALID
detail = "куки Циана помечены невалидными (last_invalid_at)"
_alert_cian_cookies(source, detail)
kit_runs.mark_skipped(db, source=source, reason=reason, details=detail)
kit_defer_next_run_at(db, schedule_row)
return False
state = await verify_session(cookies)
if state is None:
logger.warning(
"scheduler: cian_history_backfill — cookies expired or invalid, skipping run"
)
try:
sentry_sdk.capture_message(
"cian_history_backfill skipped: Cian session cookies expired — "
"please re-upload via admin UI",
level="warning",
)
except Exception:
pass # sentry_sdk not initialised in dev
# verify вернул именно None (401 / isAuthenticated=false) — куки числятся
# валидными по сроку, но Циан их не принимает. Sentinel-ответы (бан / источник
# недоступен / сменилась вёрстка) сюда НЕ попадают, они truthy — см. cian_session.
detail = "Циан не принимает куки (разлогин)"
_alert_cian_cookies(source, detail)
kit_runs.mark_skipped(db, source=source, reason=SKIP_CIAN_COOKIES_INVALID, details=detail)
kit_defer_next_run_at(db, schedule_row)
return False
# Куки рабочие — предупреждаем, пока есть время их обновить без простоя сбора.
expires_at = session_expires_at(db)
if expires_at is not None and expires_at - now <= timedelta(days=COOKIE_EXPIRY_WARN_DAYS):
logger.error(
"scheduler: куки Циана протухнут %s (осталось %.1f дн.) — обновите заранее, "
"иначе %s встанет молча",
expires_at.date().isoformat(),
(expires_at - now).total_seconds() / 86400,
source,
)
return True

View file

@ -0,0 +1,480 @@
"""Пропуск расписания оставляет след — строку scrape_runs(status='skipped') (#2658).
До #2658 планировщик пропускал наступившее окно НЕМО: logger + сдвиг next_run_at, ни
строки прогона, ни изменения last_run_at. cian_history_backfill так простоял 37 дней на
протухших куках Циана, и снаружи это выглядело как «всё по расписанию».
Покрываем:
1. mark_skipped INSERT новой строки и схлопывание подряд идущих одинаковых причин.
2. Все ПЯТЬ мест, которые раньше пропускали молча: kit `_claim_run` ×3 (already_running /
concurrent_claim / running_appeared_under_lock), kit `scheduler_loop` (unknown_source)
и продуктовый cian `pre_claim`.
3. Достижимость алерта: протухшие куки (load_session None) дают ERROR-запись, которая
в scraper-контейнере становится событием GlitchTip (event_level=ERROR), раньше
алерт стоял во второй ветке, куда на протухших куках исполнение не доходит.
4. Неломание монитора нулевых прогонов (#2625): 'skipped' — не «прогон вернул ноль
лотов», обе alert-выборки его не видят и стрик им не прерывается.
Без сети, без БД.
"""
from __future__ import annotations
import logging
import os
from datetime import UTC, datetime, timedelta
from typing import Any
from unittest.mock import AsyncMock, 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 kit_sched
from scraper_kit.orchestration.scheduler import (
SKIP_ALREADY_RUNNING,
SKIP_CONCURRENT_CLAIM,
SKIP_RUNNING_UNDER_LOCK,
SKIP_UNKNOWN_SOURCE,
SchedulerContext,
_claim_run,
)
def _make_sched(source: str) -> dict[str, Any]:
return {
"id": 1,
"source": source,
"enabled": True,
"window_start_hour": 2,
"window_end_hour": 5,
"default_params": {},
"last_run_id": None,
"last_run_at": None,
"next_run_at": None,
}
# ── 1. mark_skipped: INSERT + схлопывание ────────────────────────────────────
class _FakeSkipDB:
"""Session-мок для mark_skipped: UPDATE-ветка (схлопывание) vs INSERT-ветка."""
def __init__(self, *, collapse: bool) -> None:
self._collapse = collapse
self.statements: list[tuple[str, dict[str, Any]]] = []
self.committed = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
sql = str(stmt)
self.statements.append((sql, params or {}))
result = MagicMock()
if "UPDATE scrape_runs r" in sql:
result.fetchone.return_value = MagicMock(id=7) if self._collapse else None
else:
result.fetchone.return_value = MagicMock(id=8)
return result
def commit(self) -> None:
self.committed = True
def sql_of(self, needle: str) -> tuple[str, dict[str, Any]] | None:
for sql, params in self.statements:
if needle in sql:
return sql, params
return None
def test_mark_skipped_inserts_row_with_machine_readable_reason() -> None:
db = _FakeSkipDB(collapse=False)
run_id = kit_runs.mark_skipped(
db, source="cian_history_backfill", reason="cian_cookies_expired", details="протухли"
)
assert run_id == 8
insert = db.sql_of("INSERT INTO scrape_runs")
assert insert is not None, "новая строка прогона не создана"
sql, params = insert
assert "'skipped'" in sql
assert params["source"] == "cian_history_backfill"
assert params["reason"] == "cian_cookies_expired" # слаг, не человеческий текст
assert "протухли" in params["counters"]
assert db.committed is True
def test_mark_skipped_collapses_consecutive_same_reason() -> None:
"""Подряд идущие одинаковые пропуски не плодят строки — иначе тик 60с = строка/мин."""
db = _FakeSkipDB(collapse=True)
run_id = kit_runs.mark_skipped(db, source="avito_full_load", reason=SKIP_ALREADY_RUNNING)
assert run_id == 7
assert db.sql_of("INSERT INTO scrape_runs") is None, "схлопывание не сработало"
update = db.sql_of("UPDATE scrape_runs r")
assert update is not None
sql, _params = update
assert "skips" in sql # счётчик повторов растёт вместо новой строки
# ── 2. пять мест: kit _claim_run ×3 ──────────────────────────────────────────
class _FakeResult:
def __init__(self, *, scalar: Any = None, fetchone: Any = None) -> None:
self._scalar = scalar
self._fetchone = fetchone
def scalar(self) -> Any:
return self._scalar
def fetchone(self) -> Any:
return self._fetchone
class _FakeClaimDB:
"""Мок Session для _claim_run (тот же сценарный контракт, что в parity-тестах)."""
def __init__(self, *, running_states: list[bool], lock: bool = True) -> None:
self._running = list(running_states)
self._lock = lock
self.rolled_back = False
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
if "SELECT 1 FROM scrape_runs" in sql:
return _FakeResult(fetchone=(1,) if self._running.pop(0) else None)
if "pg_try_advisory_xact_lock" in sql:
return _FakeResult(scalar=self._lock)
return _FakeResult()
def commit(self) -> None:
pass
def rollback(self) -> None:
self.rolled_back = True
def _ctx() -> SchedulerContext:
runs = MagicMock()
runs.create_run = MagicMock(return_value=42)
return SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(),
runs=runs,
)
@pytest.mark.parametrize(
("running_states", "lock", "expected_reason"),
[
([True], True, SKIP_ALREADY_RUNNING),
([False], False, SKIP_CONCURRENT_CLAIM),
([False, True], True, SKIP_RUNNING_UNDER_LOCK),
],
)
def test_claim_run_skip_writes_skipped_row(
running_states: list[bool], lock: bool, expected_reason: str
) -> None:
"""Каждая из трёх skip-веток _claim_run пишет строку прогона с своей причиной."""
db = _FakeClaimDB(running_states=running_states, lock=lock)
with patch.object(kit_runs, "mark_skipped") as mark:
run_id = _claim_run(db, _make_sched("avito_city_sweep"), _ctx())
assert run_id is None
mark.assert_called_once()
assert mark.call_args.kwargs["reason"] == expected_reason
assert mark.call_args.kwargs["source"] == "avito_city_sweep"
def test_claim_run_under_lock_rolls_back_before_writing() -> None:
"""Строка пишется ПОСЛЕ rollback'а — иначе её INSERT улетел бы в откат лока."""
db = _FakeClaimDB(running_states=[False, True], lock=True)
order: list[str] = []
real_rollback = db.rollback
def _rollback() -> None:
order.append("rollback")
real_rollback()
db.rollback = _rollback # type: ignore[method-assign]
with patch.object(kit_runs, "mark_skipped", side_effect=lambda *a, **k: order.append("skip")):
_claim_run(db, _make_sched("avito_city_sweep"), _ctx())
assert order == ["rollback", "skip"]
def test_claim_run_happy_path_writes_no_skip_row() -> None:
"""Успешный claim не должен оставлять skip-строк (иначе счётчики мусорные)."""
db = _FakeClaimDB(running_states=[False, False], lock=True)
with patch.object(kit_runs, "mark_skipped") as mark:
run_id = _claim_run(db, _make_sched("avito_city_sweep"), _ctx())
assert run_id == 42
mark.assert_not_called()
# ── 2b. пятое место: scheduler_loop с неизвестным source ─────────────────────
async def test_scheduler_loop_unknown_source_writes_skipped_row() -> None:
"""enabled-расписание без handler'а: раньше — warning каждый тик и ноль следов."""
ctx = SchedulerContext(
config=MagicMock(),
matcher=MagicMock(),
enrichment=MagicMock(),
session_factory=MagicMock(return_value=MagicMock()),
runs=MagicMock(),
# первый вызов (верх тика) — работаем, дальше — drain, чтобы выйти из while True
shutdown_requested=MagicMock(side_effect=[False, True, True]),
)
with (
patch.object(kit_sched.asyncio, "sleep", AsyncMock()),
patch.object(kit_sched, "reap_zombies", MagicMock(return_value=0)),
patch.object(
kit_sched,
"get_due_schedules",
MagicMock(return_value=[_make_sched("source_from_mars")]),
),
patch.object(kit_runs, "mark_skipped") as mark,
):
await kit_sched.scheduler_loop(ctx, registry={})
mark.assert_called_once()
assert mark.call_args.kwargs["reason"] == SKIP_UNKNOWN_SOURCE
assert mark.call_args.kwargs["source"] == "source_from_mars"
# ── 3. cian pre_claim: строка + достижимый алерт ─────────────────────────────
def _patch_cian(
*,
cookies: dict[str, str] | None,
expires_at: datetime | None,
verify: Any = None,
) -> Any:
from app.services import cian_session
return (
patch.object(cian_session, "load_session", MagicMock(return_value=cookies)),
patch.object(cian_session, "session_expires_at", MagicMock(return_value=expires_at)),
patch.object(cian_session, "verify_session", AsyncMock(return_value=verify)),
)
async def test_cian_pre_claim_expired_cookies_logs_error() -> None:
"""Фальсификация #2658-2: на протухших куках ДОЛЖНА быть ERROR-запись.
Именно ERROR в scraper-контейнере GlitchTip поднят с
LoggingIntegration(event_level=ERROR) (scheduler_main.py), warning событием не станет.
Старый код в этой ветке писал только logger.warning, а capture_message стоял во
второй ветке (verify_session None), недостижимой при протухании: load_session сам
фильтрует expires_at_estimate > NOW().
"""
from app.services import product_handlers
expired = datetime.now(tz=UTC) - timedelta(days=37)
p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=expired)
records: list[logging.LogRecord] = []
class _Collector(logging.Handler):
def emit(self, record: logging.LogRecord) -> None:
records.append(record)
handler = _Collector()
product_handlers.logger.addHandler(handler)
try:
with (
p_load,
p_exp,
p_verify,
patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()),
patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()),
):
proceed = await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
finally:
product_handlers.logger.removeHandler(handler)
assert proceed is False
errors = [r for r in records if r.levelno >= logging.ERROR]
assert errors, "протухание кук не породило ERROR-запись → события GlitchTip не будет"
async def test_cian_pre_claim_expired_cookies_writes_skipped_row() -> None:
"""Немой `return False` заменён строкой прогона с причиной cian_cookies_expired."""
from app.services import product_handlers
expired = datetime.now(tz=UTC) - timedelta(days=37)
p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=expired)
with (
p_load,
p_exp,
p_verify,
patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()) as defer,
patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark,
):
proceed = await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
assert proceed is False
mark.assert_called_once()
assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_EXPIRED
assert mark.call_args.kwargs["source"] == "cian_history_backfill"
defer.assert_called_once() # next_run_at по-прежнему двигаем (не долбим Циан каждые 60с)
async def test_cian_pre_claim_missing_cookies_reason_differs_from_expired() -> None:
"""«Кук нет вовсе» и «протухли» — разные слаги: причина машиночитаема."""
from app.services import product_handlers
p_load, p_exp, p_verify = _patch_cian(cookies=None, expires_at=None)
with (
p_load,
p_exp,
p_verify,
patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()),
patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark,
):
await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_MISSING
assert product_handlers.SKIP_CIAN_COOKIES_MISSING != product_handlers.SKIP_CIAN_COOKIES_EXPIRED
async def test_cian_pre_claim_rejected_cookies_writes_skipped_row() -> None:
"""Вторая ветка (Циан не принимает куки) — тоже строка, а не только алерт."""
from app.services import product_handlers
future = datetime.now(tz=UTC) + timedelta(days=20)
p_load, p_exp, p_verify = _patch_cian(
cookies={"DMIR_AUTH": "x"}, expires_at=future, verify=None
)
with (
p_load,
p_exp,
p_verify,
patch.object(product_handlers, "kit_defer_next_run_at", MagicMock()),
patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark,
):
proceed = await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
assert proceed is False
assert mark.call_args.kwargs["reason"] == product_handlers.SKIP_CIAN_COOKIES_INVALID
async def test_cian_pre_claim_warns_before_expiry_not_after() -> None:
"""Предупреждаем ЗАРАНЕЕ: куки ещё рабочие, но жить им меньше COOKIE_EXPIRY_WARN_DAYS.
Обновление кук ручная операция; алерт по факту протухания приходит, когда сбор уже
встал. Гейт при этом пропускает прогон (proceed=True) предупреждение, не блокировка.
"""
from app.services import cian_session, product_handlers
soon = datetime.now(tz=UTC) + timedelta(days=cian_session.COOKIE_EXPIRY_WARN_DAYS - 1)
p_load, p_exp, p_verify = _patch_cian(
cookies={"DMIR_AUTH": "x"}, expires_at=soon, verify={"user": {"isAuthenticated": True}}
)
records: list[logging.LogRecord] = []
class _Collector(logging.Handler):
def emit(self, record: logging.LogRecord) -> None:
records.append(record)
handler = _Collector()
product_handlers.logger.addHandler(handler)
try:
with p_load, p_exp, p_verify:
proceed = await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
finally:
product_handlers.logger.removeHandler(handler)
assert proceed is True
assert [r for r in records if r.levelno >= logging.ERROR], "не предупредили заранее"
async def test_cian_pre_claim_fresh_cookies_are_silent() -> None:
"""Свежие куки — ни алерта, ни skip-строки (иначе алерт-усталость)."""
from app.services import product_handlers
far = datetime.now(tz=UTC) + timedelta(days=25)
p_load, p_exp, p_verify = _patch_cian(
cookies={"DMIR_AUTH": "x"}, expires_at=far, verify={"user": {"isAuthenticated": True}}
)
records: list[logging.LogRecord] = []
class _Collector(logging.Handler):
def emit(self, record: logging.LogRecord) -> None:
records.append(record)
handler = _Collector()
product_handlers.logger.addHandler(handler)
try:
with (
p_load,
p_exp,
p_verify,
patch.object(product_handlers.kit_runs, "mark_skipped", MagicMock()) as mark,
):
proceed = await product_handlers._cian_pre_claim(
MagicMock(), _make_sched("cian_history_backfill"), MagicMock()
)
finally:
product_handlers.logger.removeHandler(handler)
assert proceed is True
mark.assert_not_called()
assert not [r for r in records if r.levelno >= logging.ERROR]
# ── 4. монитор нулевых прогонов не смешивается с пропусками ──────────────────
class _RecordingDB:
"""Возвращает заданные строки на SELECT и запоминает SQL (для проверки фильтров)."""
def __init__(self, rows: list[Any]) -> None:
self.rows = rows
self.sql: list[str] = []
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> Any:
self.sql.append(str(stmt))
result = MagicMock()
result.fetchall.return_value = self.rows
return result
def test_zero_result_monitor_ignores_skipped_rows() -> None:
"""'skipped' НЕ участвует в стрике нулевых прогонов — это не «прогон вернул ноль».
Смешать их в одном счётчике нельзя в обе стороны: пропуск не должен ни считаться
нулевым прогоном, ни прерывать стрик реальных нулевых. Отсечка делается в SQL
проверяем, что 'skipped' не попал в список статусов выборки.
"""
db = _RecordingDB(rows=[])
kit_runs._alert_if_consecutive_zero_results(db, "cian_city_sweep")
kit_runs._alert_if_consecutive_failures(db, "cian_city_sweep")
assert db.sql, "alert-выборка не выполнилась"
for sql in db.sql:
assert "status IN ('failed', 'banned', 'done', 'cancelled')" in sql
assert "skipped" not in sql
def test_mark_skipped_status_is_not_a_failure_status() -> None:
"""Строка-пропуск не попадает и в failed/banned-стрик (алерт «3 подряд ошибки»)."""
db = _FakeSkipDB(collapse=False)
kit_runs.mark_skipped(db, source="cian_history_backfill", reason="cian_cookies_expired")
insert = db.sql_of("INSERT INTO scrape_runs")
assert insert is not None
sql, _params = insert
assert "'skipped'" in sql
assert "'failed'" not in sql and "'banned'" not in sql

View file

@ -238,6 +238,8 @@ class _FakeClaimDB:
self.committed = False
self.rolled_back = False
self.update_calls = 0
# #2658: строки-пропуски (scrape_runs status='skipped'), которые пишет mark_skipped.
self.skip_rows = 0
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
sql = str(stmt)
@ -249,6 +251,13 @@ class _FakeClaimDB:
if "UPDATE scrape_schedules" in sql:
self.update_calls += 1
return _FakeResult()
# #2658 mark_skipped: сначала пробует схлопнуть последнюю skip-строку (UPDATE
# ... RETURNING), не нашёл — вставляет новую.
if "UPDATE scrape_runs r" in sql:
return _FakeResult(fetchone=None)
if "INSERT INTO scrape_runs" in sql:
self.skip_rows += 1
return _FakeResult(fetchone=MagicMock(id=999))
return _FakeResult()
def commit(self) -> None:
@ -286,6 +295,7 @@ def test_claim_run_skip_already_running() -> None:
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id is None
ctx.runs.create_run.assert_not_called()
assert db.skip_rows == 1 # #2658: пропуск оставляет строку, а не только лог
def test_claim_run_skip_lock_busy() -> None:
@ -295,6 +305,7 @@ def test_claim_run_skip_lock_busy() -> None:
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
assert run_id is None
ctx.runs.create_run.assert_not_called()
assert db.skip_rows == 1 # #2658
def test_claim_run_running_appeared_under_lock() -> None:
@ -305,6 +316,7 @@ def test_claim_run_running_appeared_under_lock() -> None:
assert run_id is None
assert db.rolled_back is True
ctx.runs.create_run.assert_not_called()
assert db.skip_rows == 1 # #2658
# ── 3. reap_zombies ──────────────────────────────────────────────────────────

View file

@ -1,10 +1,14 @@
"""scrape_runs helpers — tracking long-running pipeline runs (strangler-копия #2135).
Байт-эквивалент `app.services.scrape_runs` чистые SQL-хелперы поверх таблицы
`scrape_runs` (миграции 015 + 051). Развязка от `app.*`: единственное намеренное
отличие `sentry_sdk` импортируется опционально (kit standalone-импортируем, а
`sentry-sdk` не входит в его зависимости). Если пакет не установлен alert-хук
best-effort no-op, поведение SQL-финализаторов идентично старому.
`scrape_runs` (миграции 015 + 051). Развязка от `app.*`. Намеренные отличия от
app-копии:
1. `sentry_sdk` импортируется опционально (kit standalone-импортируем, а
`sentry-sdk` не входит в его зависимости). Если пакет не установлен alert-хук
best-effort no-op, поведение SQL-финализаторов идентично старому.
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
app-копии эта функция не нужна.
"""
from __future__ import annotations
@ -230,6 +234,84 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
return int(row.id)
def mark_skipped(db: Session, *, source: str, reason: str, details: str | None = None) -> int:
"""INSERT scrape_runs(status='skipped') — след пропущенного расписания (#2658).
Раньше планировщик пропускал наступившее окно НЕМО: logger + сдвиг next_run_at, ни
строки прогона, ни изменения last_run_at снаружи «всё по расписанию»
(cian_history_backfill так простоял 37 дней на протухших куках). Статус 'skipped'
заведён ещё миграцией 015 и локализован во фронте («пропущено»), но до #2658 не
использовался ни разу.
`reason` машиночитаемый слаг (already_running / cian_cookies_expired / ), пишется
в `error`: по нему пропуск «нет кук» отличается от «уже бежит» без разбора текста.
`details` человеческое пояснение, кладётся в `counters.detail`.
Схлопывание подряд идущих одинаковых пропусков: если ПОСЛЕДНЯЯ строка прогона этого
source уже 'skipped' с тем же reason, новая не создаётся; у существующей
обновляется finished_at и счётчик `counters.skips`. Без этого «уже бежит» и
«неизвестный source» плодили бы строку каждый тик (60 с), пока держится причина.
Возвращает id строки (новой или обновлённой).
"""
row = db.execute(
text(
"""
WITH latest AS (
SELECT id FROM scrape_runs
WHERE source = :source
ORDER BY id DESC
LIMIT 1
)
UPDATE scrape_runs r
SET heartbeat_at = NOW(),
finished_at = NOW(),
counters = jsonb_set(
COALESCE(r.counters, '{}'::jsonb),
'{skips}',
to_jsonb(COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1)
)
FROM latest
WHERE r.id = latest.id
AND r.status = 'skipped'
AND r.error = :reason
RETURNING r.id
"""
),
{"source": source, "reason": reason},
).fetchone()
if row is None:
row = db.execute(
text(
"""
INSERT INTO scrape_runs (
source, status, error, counters, started_at, heartbeat_at, finished_at
)
VALUES (
:source, 'skipped', :reason, CAST(:counters AS jsonb), NOW(), NOW(), NOW()
)
RETURNING id
"""
),
{
"source": source,
"reason": reason,
"counters": json.dumps({"skips": 1, "detail": details}, ensure_ascii=False),
},
).fetchone()
db.commit()
assert row is not None, "scrape_runs skipped-row INSERT returned no id"
logger.warning(
"scheduler: skipped run source=%s reason=%s%s",
source,
reason,
f" ({details})" if details else "",
)
return int(row.id)
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.

View file

@ -69,6 +69,15 @@ ZOMBIE_THRESHOLD_HOURS = 6
# (см. #1182 P2 — идентично боевому scheduler'у).
_CHILD_DRAIN_TIMEOUT_S = 80.0
# Машиночитаемые причины пропуска расписания (#2658) — пишутся в scrape_runs.error
# строки со status='skipped'. Слаг, а не человеческий текст: по нему «уже бежит»
# отличается от «нет кук» (продуктовые причины — в app.services.product_handlers)
# запросом, а не грепом логов.
SKIP_ALREADY_RUNNING = "already_running"
SKIP_CONCURRENT_CLAIM = "concurrent_claim"
SKIP_RUNNING_UNDER_LOCK = "running_appeared_under_lock"
SKIP_UNKNOWN_SOURCE = "unknown_source"
# ── типы job/handler ─────────────────────────────────────────────────────────
# Job получает свежую сессию (открыта `_dispatch`), run_id, params и весь контекст
# (config/matcher/enrichment/runs) — чтобы иметь доступ к инжектированным зависимостям.
@ -254,10 +263,21 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
Returns run_id или None если уже есть running run для этого source.
Общий helper для всех source-handler'ов. Advisory-lock concurrency-логика перенесена
ДОСЛОВНО из боевого scheduler'а (создание run — через инжектированный ctx.runs).
#2658: каждая ветка «пропустили окно» пишет строку scrape_runs(status='skipped')
с машиночитаемой причиной раньше пропуск был виден только в docker-логах, которые
теряются при редеплое. Пишем через kit-копию runs (у инжектированной app-копии
mark_skipped нет строки-пропуски создаёт только планировщик).
"""
source = schedule_row["source"]
if has_running_run(db, source):
logger.info("scheduler: skip — already running for source=%s", source)
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_ALREADY_RUNNING,
details="предыдущий прогон ещё идёт",
)
return None
# #750: атомарный claim через transaction-scoped advisory lock. Сериализует
@ -272,6 +292,12 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
).scalar()
if not got_lock:
logger.info("scheduler: skip — concurrent claim in progress for source=%s", source)
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_CONCURRENT_CLAIM,
details="конкурентный тик уже клеймит этот source",
)
return None
# Double-checked под локом: конкурентный тик мог закоммитить running-run МЕЖДУ
@ -280,6 +306,13 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
if has_running_run(db, source):
logger.info("scheduler: skip — running appeared under lock for source=%s", source)
db.rollback() # освобождаем advisory lock (claim не состоялся)
# mark_skipped ПОСЛЕ rollback'а: его INSERT + commit иначе улетели бы в откат.
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_RUNNING_UNDER_LOCK,
details="running-прогон появился под локом",
)
return None
params = schedule_row.get("default_params") or {}
@ -700,6 +733,14 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
handler = resolve_handler(source, registry)
if handler is None:
logger.warning("scheduler: unknown source=%s, skip", source)
# #2658: enabled-расписание без handler'а молча не выполнялось
# бы вечно (next_run_at не двигается → warning каждый тик).
_kit_runs.mark_skipped(
db,
source=source,
reason=SKIP_UNKNOWN_SOURCE,
details="нет handler'а в реестре",
)
else:
await _dispatch(handler, db, sch, ctx)
@ -730,6 +771,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
__all__ = [
"SCHEDULER_TICK_SEC",
"SKIP_ALREADY_RUNNING",
"SKIP_CONCURRENT_CLAIM",
"SKIP_RUNNING_UNDER_LOCK",
"SKIP_UNKNOWN_SOURCE",
"ZOMBIE_THRESHOLD_HOURS",
"Handler",
"SchedulerContext",