fix(tradein/scraper): пропуск расписания пишет строку прогона со статусом skipped (#2658) #2662
8 changed files with 822 additions and 25 deletions
|
|
@ -2262,7 +2262,11 @@ def list_scrape_runs_unified(
|
|||
db: Annotated[Session, Depends(get_db)],
|
||||
source: Annotated[str | None, Query()] = None,
|
||||
status: Annotated[
|
||||
Literal["done", "running", "banned", "zombie", "failed", "cancelled"] | None, Query()
|
||||
# 'skipped' (#2658) — пропущенное расписание; без него оператор не может
|
||||
# спросить «что сейчас пропускается» (фильтр отдавал 422 на единственной
|
||||
# поверхности, построенной ровно для этого вопроса).
|
||||
Literal["done", "running", "banned", "zombie", "failed", "cancelled", "skipped"] | None,
|
||||
Query(),
|
||||
] = None,
|
||||
limit: Annotated[int, Query(ge=1, le=200)] = 50,
|
||||
offset: Annotated[int, Query(ge=0)] = 0,
|
||||
|
|
@ -2274,7 +2278,7 @@ def list_scrape_runs_unified(
|
|||
|
||||
Query:
|
||||
source — опц. фильтр по source (avito_city_sweep / cian_city_sweep / ...).
|
||||
status — опц. фильтр (done/running/banned/zombie/failed/cancelled).
|
||||
status — опц. фильтр (done/running/banned/zombie/failed/cancelled/skipped).
|
||||
limit — default 50, max 200.
|
||||
offset — default 0.
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -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,37 @@ def load_session(db: Session) -> dict[str, str] | None:
|
|||
return cookies
|
||||
|
||||
|
||||
def session_expires_at(db: Session, *, valid_only: bool = False) -> datetime | None:
|
||||
"""Когда протухают самые свежезагруженные куки (#2658).
|
||||
|
||||
`load_session` отбирает только ещё валидные записи (expires_at_estimate > NOW()) и на
|
||||
протухших отдаёт None — вызывающий не мог отличить «кук никогда не загружали» от
|
||||
«протухли позавчера» и не мог предупредить ЗАРАНЕЕ.
|
||||
|
||||
valid_only=False (диагностика после None от load_session) — свежайшая запись любая:
|
||||
валидных по определению нет, нужен именно срок протухшей. valid_only=True — та же
|
||||
запись, которую взял бы load_session: для предупреждения «скоро протухнут» нужен срок
|
||||
ИМЕННО используемых кук, иначе при нескольких аккаунтах посчитаем по чужой строке.
|
||||
"""
|
||||
row = db.execute(
|
||||
text(
|
||||
"""
|
||||
SELECT expires_at_estimate FROM cian_session_cookies
|
||||
WHERE NOT CAST(:valid_only AS boolean)
|
||||
OR (expires_at_estimate > NOW()
|
||||
AND (last_invalid_at IS NULL OR last_invalid_at < uploaded_at))
|
||||
ORDER BY uploaded_at DESC
|
||||
LIMIT 1
|
||||
"""
|
||||
),
|
||||
{"valid_only": valid_only},
|
||||
).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(
|
||||
|
|
|
|||
|
|
@ -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,97 @@ 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
|
||||
|
||||
# Куки рабочие — предупреждаем, пока есть время их обновить без простоя сбора.
|
||||
# valid_only=True: срок ИМЕННО той записи, которую взял load_session (при нескольких
|
||||
# аккаунтах свежайшая-любая может быть чужой протухшей строкой).
|
||||
expires_at = session_expires_at(db, valid_only=True)
|
||||
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
|
||||
|
||||
|
||||
|
|
|
|||
535
tradein-mvp/backend/tests/test_scrape_skip_visibility.py
Normal file
535
tradein-mvp/backend/tests/test_scrape_skip_visibility.py
Normal file
|
|
@ -0,0 +1,535 @@
|
|||
"""Пропуск расписания оставляет след — строку 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 # счётчик повторов растёт вместо новой строки
|
||||
|
||||
|
||||
def test_mark_skipped_collapse_refreshes_detail_and_started_at() -> None:
|
||||
"""Схлопывание освежает строку: detail, started_at (иначе живой стрик тонет в списке).
|
||||
|
||||
Список прогонов сортирует ORDER BY started_at DESC с limit=20 — замороженный
|
||||
started_at утопил бы 37-дневный стрик под свежими прогонами других источников.
|
||||
Начало стрика сохраняется в counters.first_skip_at.
|
||||
"""
|
||||
db = _FakeSkipDB(collapse=True)
|
||||
kit_runs.mark_skipped(
|
||||
db, source="cian_history_backfill", reason="cian_cookies_expired", details="37 дн. назад"
|
||||
)
|
||||
update = db.sql_of("UPDATE scrape_runs r")
|
||||
assert update is not None
|
||||
sql, params = update
|
||||
assert "started_at = NOW()" in sql
|
||||
assert "'detail', CAST(:details AS text)" in sql, "detail замерзает от первого пропуска"
|
||||
assert "first_skip_at" in sql, "начало стрика потеряно"
|
||||
assert params["details"] == "37 дн. назад"
|
||||
|
||||
|
||||
def test_mark_skipped_latest_lookup_uses_indexed_order() -> None:
|
||||
"""Поиск последней строки идёт по (source, started_at DESC) — индекс из миграции 015.
|
||||
|
||||
ORDER BY id DESC этот индекс не использует: для unknown_source (тик каждые 60 с
|
||||
бессрочно) это был бы отбор всех строк источника с сортировкой раз в минуту.
|
||||
"""
|
||||
db = _FakeSkipDB(collapse=True)
|
||||
kit_runs.mark_skipped(db, source="avito_full_load", reason=SKIP_ALREADY_RUNNING)
|
||||
update = db.sql_of("UPDATE scrape_runs r")
|
||||
assert update is not None
|
||||
sql, _params = update
|
||||
assert "ORDER BY started_at DESC, id DESC" 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()
|
||||
)
|
||||
# Срок считаем по ЗАПИСИ, которую взял load_session: при нескольких аккаунтах
|
||||
# свежайшая-любая может быть чужой протухшей строкой (валидности она не знает).
|
||||
assert cian_session.session_expires_at.call_args.kwargs == {"valid_only": True}
|
||||
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_admin_runs_filter_accepts_skipped_status() -> None:
|
||||
"""GET /admin/scrape/runs?status=skipped не должен отдавать 422.
|
||||
|
||||
Единственная поверхность, где оператор спрашивает «что сейчас пропускается» — фильтр
|
||||
статуса в таблице прогонов. Без 'skipped' в Literal строки видны, а вопрос задать
|
||||
нельзя, то есть цель #2658 на UI достигнута наполовину.
|
||||
"""
|
||||
import typing
|
||||
|
||||
from app.api.v1.admin import list_scrape_runs_unified
|
||||
|
||||
annotation = typing.get_type_hints(list_scrape_runs_unified, include_extras=True)["status"]
|
||||
allowed: set[str] = set()
|
||||
for arg in typing.get_args(typing.get_args(annotation)[0]):
|
||||
allowed.update(typing.get_args(arg))
|
||||
assert "skipped" in allowed
|
||||
|
||||
|
||||
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
|
||||
|
|
@ -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 ──────────────────────────────────────────────────────────
|
||||
|
|
|
|||
|
|
@ -34,7 +34,18 @@ interface RunsListResp {
|
|||
|
||||
// ── Hook ───────────────────────────────────────────────────────────────────
|
||||
|
||||
const RUN_STATUS_ALL = ["", "running", "done", "failed", "cancelled", "zombie", "banned"] as const;
|
||||
// "skipped" (#2658) — пропущенное расписание (нет кук / уже бежит / нет handler'а);
|
||||
// translateStatus уже знает «пропущено», бейдж падает в нейтральный вариант.
|
||||
const RUN_STATUS_ALL = [
|
||||
"",
|
||||
"running",
|
||||
"done",
|
||||
"failed",
|
||||
"cancelled",
|
||||
"zombie",
|
||||
"banned",
|
||||
"skipped",
|
||||
] as const;
|
||||
type RunStatusFilter = (typeof RUN_STATUS_ALL)[number];
|
||||
|
||||
// "" means "all sources"; otherwise a specific source prefix (avito / cian / yandex)
|
||||
|
|
|
|||
|
|
@ -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,95 @@ 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, новая не создаётся; у существующей обновляются
|
||||
started_at/finished_at, счётчик `counters.skips` и свежий `counters.detail`. Без этого
|
||||
«уже бежит» и «неизвестный source» плодили бы строку каждый тик (60 с), пока держится
|
||||
причина.
|
||||
|
||||
started_at двигаем намеренно: списки прогонов (`list_all`/`list_recent`) сортируют
|
||||
ORDER BY started_at DESC и берутся с limit=20, поэтому замороженный started_at утопил
|
||||
бы ЖИВОЙ стрик под свежими прогонами других источников — ровно тот сценарий, из-за
|
||||
которого заведён #2658 (в базе след есть, на экране нет). Начало стрика при этом не
|
||||
теряется: первый started_at переезжает в `counters.first_skip_at`. Сортировку общего
|
||||
списка не трогаем — она про все источники, а чинить надо было одну строку.
|
||||
|
||||
Возвращает id строки (новой или обновлённой).
|
||||
"""
|
||||
row = db.execute(
|
||||
text(
|
||||
"""
|
||||
WITH latest AS (
|
||||
SELECT id FROM scrape_runs
|
||||
WHERE source = :source
|
||||
ORDER BY started_at DESC, id DESC
|
||||
LIMIT 1
|
||||
)
|
||||
UPDATE scrape_runs r
|
||||
SET heartbeat_at = NOW(),
|
||||
started_at = NOW(),
|
||||
finished_at = NOW(),
|
||||
counters = COALESCE(r.counters, '{}'::jsonb) || jsonb_build_object(
|
||||
'skips', COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1,
|
||||
'detail', CAST(:details AS text),
|
||||
'first_skip_at', COALESCE(
|
||||
r.counters ->> 'first_skip_at', CAST(r.started_at AS text)
|
||||
)
|
||||
)
|
||||
FROM latest
|
||||
WHERE r.id = latest.id
|
||||
AND r.status = 'skipped'
|
||||
AND r.error = :reason
|
||||
RETURNING r.id
|
||||
"""
|
||||
),
|
||||
{"source": source, "reason": reason, "details": details},
|
||||
).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 колонки.
|
||||
|
||||
|
|
|
|||
|
|
@ -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",
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue