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)],
|
db: Annotated[Session, Depends(get_db)],
|
||||||
source: Annotated[str | None, Query()] = None,
|
source: Annotated[str | None, Query()] = None,
|
||||||
status: Annotated[
|
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,
|
] = None,
|
||||||
limit: Annotated[int, Query(ge=1, le=200)] = 50,
|
limit: Annotated[int, Query(ge=1, le=200)] = 50,
|
||||||
offset: Annotated[int, Query(ge=0)] = 0,
|
offset: Annotated[int, Query(ge=0)] = 0,
|
||||||
|
|
@ -2274,7 +2278,7 @@ def list_scrape_runs_unified(
|
||||||
|
|
||||||
Query:
|
Query:
|
||||||
source — опц. фильтр по source (avito_city_sweep / cian_city_sweep / ...).
|
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.
|
limit — default 50, max 200.
|
||||||
offset — default 0.
|
offset — default 0.
|
||||||
"""
|
"""
|
||||||
|
|
|
||||||
|
|
@ -8,6 +8,7 @@ from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
|
from datetime import datetime
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from curl_cffi.requests import AsyncSession
|
from curl_cffi.requests import AsyncSession
|
||||||
|
|
@ -23,6 +24,11 @@ from app.core.config import settings
|
||||||
|
|
||||||
logger = logging.getLogger(__name__)
|
logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
# За сколько дней до протухания кук предупреждать (#2658). Обновление кук — РУЧНАЯ
|
||||||
|
# операция (залить дамп через админку), человеку нужен запас: алерт по факту протухания
|
||||||
|
# приходит, когда сбор уже встал. save_session ставит ttl 30 дней, так что окно широкое.
|
||||||
|
COOKIE_EXPIRY_WARN_DAYS = 5
|
||||||
|
|
||||||
# Cookies критичные для Cian auth — фильтр перед сохранением.
|
# Cookies критичные для Cian auth — фильтр перед сохранением.
|
||||||
# Список обновлён по реальному DevTools-дампу из logged-in сессии cian.ru (2026-05-23).
|
# Список обновлён по реальному DevTools-дампу из logged-in сессии cian.ru (2026-05-23).
|
||||||
# Старые записи оставлены как fallback (backward compat).
|
# Старые записи оставлены как fallback (backward compat).
|
||||||
|
|
@ -294,6 +300,37 @@ def load_session(db: Session) -> dict[str, str] | None:
|
||||||
return cookies
|
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:
|
def mark_session_invalid(db: Session, account_user_id: int) -> None:
|
||||||
"""Flag session как expired/invalid (например после 401 во время scrape)."""
|
"""Flag session как expired/invalid (например после 401 во время scrape)."""
|
||||||
db.execute(
|
db.execute(
|
||||||
|
|
|
||||||
|
|
@ -21,8 +21,10 @@ from __future__ import annotations
|
||||||
|
|
||||||
import asyncio
|
import asyncio
|
||||||
import logging
|
import logging
|
||||||
|
from datetime import UTC, datetime, timedelta
|
||||||
from typing import TYPE_CHECKING, Any
|
from typing import TYPE_CHECKING, Any
|
||||||
|
|
||||||
|
from scraper_kit.orchestration import runs as kit_runs
|
||||||
from scraper_kit.orchestration.scheduler import (
|
from scraper_kit.orchestration.scheduler import (
|
||||||
Handler,
|
Handler,
|
||||||
reschedule_after_minutes,
|
reschedule_after_minutes,
|
||||||
|
|
@ -39,39 +41,97 @@ logger = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
# ── cian_history_backfill — cookie-gated backfill ────────────────────────────
|
# ── 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:
|
async def _cian_pre_claim(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext) -> bool:
|
||||||
"""Pre-claim gate: проверить наличие/валидность cian-cookies ДО claim (#1522).
|
"""Pre-claim gate: проверить наличие/валидность cian-cookies ДО claim (#1522).
|
||||||
|
|
||||||
Cookies отсутствуют/протухли → defer next_run_at на следующее окно и skip
|
Cookies отсутствуют/протухли → пишем строку прогона status='skipped' с причиной,
|
||||||
(иначе get_due_schedules переотбирает schedule каждые 60с и verify_session
|
двигаем next_run_at на следующее окно и skip (иначе get_due_schedules переотбирает
|
||||||
долбит Cian круглосуточно). Дословно из боевого trigger_cian_backfill_run.
|
schedule каждые 60с и verify_session долбит Cian круглосуточно).
|
||||||
"""
|
|
||||||
import sentry_sdk
|
|
||||||
|
|
||||||
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)
|
cookies = load_session(db)
|
||||||
if cookies is None:
|
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)
|
kit_defer_next_run_at(db, schedule_row)
|
||||||
return False
|
return False
|
||||||
|
|
||||||
state = await verify_session(cookies)
|
state = await verify_session(cookies)
|
||||||
if state is None:
|
if state is None:
|
||||||
logger.warning(
|
# verify вернул именно None (401 / isAuthenticated=false) — куки числятся
|
||||||
"scheduler: cian_history_backfill — cookies expired or invalid, skipping run"
|
# валидными по сроку, но Циан их не принимает. Sentinel-ответы (бан / источник
|
||||||
)
|
# недоступен / сменилась вёрстка) сюда НЕ попадают, они truthy — см. cian_session.
|
||||||
try:
|
detail = "Циан не принимает куки (разлогин)"
|
||||||
sentry_sdk.capture_message(
|
_alert_cian_cookies(source, detail)
|
||||||
"cian_history_backfill skipped: Cian session cookies expired — "
|
kit_runs.mark_skipped(db, source=source, reason=SKIP_CIAN_COOKIES_INVALID, details=detail)
|
||||||
"please re-upload via admin UI",
|
|
||||||
level="warning",
|
|
||||||
)
|
|
||||||
except Exception:
|
|
||||||
pass # sentry_sdk not initialised in dev
|
|
||||||
kit_defer_next_run_at(db, schedule_row)
|
kit_defer_next_run_at(db, schedule_row)
|
||||||
return False
|
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
|
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.committed = False
|
||||||
self.rolled_back = False
|
self.rolled_back = False
|
||||||
self.update_calls = 0
|
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:
|
def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult:
|
||||||
sql = str(stmt)
|
sql = str(stmt)
|
||||||
|
|
@ -249,6 +251,13 @@ class _FakeClaimDB:
|
||||||
if "UPDATE scrape_schedules" in sql:
|
if "UPDATE scrape_schedules" in sql:
|
||||||
self.update_calls += 1
|
self.update_calls += 1
|
||||||
return _FakeResult()
|
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()
|
return _FakeResult()
|
||||||
|
|
||||||
def commit(self) -> None:
|
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)
|
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
|
||||||
assert run_id is None
|
assert run_id is None
|
||||||
ctx.runs.create_run.assert_not_called()
|
ctx.runs.create_run.assert_not_called()
|
||||||
|
assert db.skip_rows == 1 # #2658: пропуск оставляет строку, а не только лог
|
||||||
|
|
||||||
|
|
||||||
def test_claim_run_skip_lock_busy() -> None:
|
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)
|
run_id = _claim_run(db, _make_sched("avito_city_sweep"), ctx)
|
||||||
assert run_id is None
|
assert run_id is None
|
||||||
ctx.runs.create_run.assert_not_called()
|
ctx.runs.create_run.assert_not_called()
|
||||||
|
assert db.skip_rows == 1 # #2658
|
||||||
|
|
||||||
|
|
||||||
def test_claim_run_running_appeared_under_lock() -> None:
|
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 run_id is None
|
||||||
assert db.rolled_back is True
|
assert db.rolled_back is True
|
||||||
ctx.runs.create_run.assert_not_called()
|
ctx.runs.create_run.assert_not_called()
|
||||||
|
assert db.skip_rows == 1 # #2658
|
||||||
|
|
||||||
|
|
||||||
# ── 3. reap_zombies ──────────────────────────────────────────────────────────
|
# ── 3. reap_zombies ──────────────────────────────────────────────────────────
|
||||||
|
|
|
||||||
|
|
@ -34,7 +34,18 @@ interface RunsListResp {
|
||||||
|
|
||||||
// ── Hook ───────────────────────────────────────────────────────────────────
|
// ── 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];
|
type RunStatusFilter = (typeof RUN_STATUS_ALL)[number];
|
||||||
|
|
||||||
// "" means "all sources"; otherwise a specific source prefix (avito / cian / yandex)
|
// "" 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).
|
"""scrape_runs helpers — tracking long-running pipeline runs (strangler-копия #2135).
|
||||||
|
|
||||||
Байт-эквивалент `app.services.scrape_runs` — чистые SQL-хелперы поверх таблицы
|
Байт-эквивалент `app.services.scrape_runs` — чистые SQL-хелперы поверх таблицы
|
||||||
`scrape_runs` (миграции 015 + 051). Развязка от `app.*`: единственное намеренное
|
`scrape_runs` (миграции 015 + 051). Развязка от `app.*`. Намеренные отличия от
|
||||||
отличие — `sentry_sdk` импортируется опционально (kit standalone-импортируем, а
|
app-копии:
|
||||||
|
1. `sentry_sdk` импортируется опционально (kit standalone-импортируем, а
|
||||||
`sentry-sdk` не входит в его зависимости). Если пакет не установлен — alert-хук
|
`sentry-sdk` не входит в его зависимости). Если пакет не установлен — alert-хук
|
||||||
best-effort no-op, поведение SQL-финализаторов идентично старому.
|
best-effort no-op, поведение SQL-финализаторов идентично старому.
|
||||||
|
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
|
||||||
|
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
|
||||||
|
app-копии эта функция не нужна.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -230,6 +234,95 @@ def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
|
||||||
return int(row.id)
|
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:
|
def update_heartbeat(db: Session, run_id: int, counters: dict[str, int]) -> None:
|
||||||
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
"""UPDATE heartbeat_at=NOW(), counters=:counters + total_seen/new_count колонки.
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -69,6 +69,15 @@ ZOMBIE_THRESHOLD_HOURS = 6
|
||||||
# (см. #1182 P2 — идентично боевому scheduler'у).
|
# (см. #1182 P2 — идентично боевому scheduler'у).
|
||||||
_CHILD_DRAIN_TIMEOUT_S = 80.0
|
_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/handler ─────────────────────────────────────────────────────────
|
||||||
# Job получает свежую сессию (открыта `_dispatch`), run_id, params и весь контекст
|
# Job получает свежую сессию (открыта `_dispatch`), run_id, params и весь контекст
|
||||||
# (config/matcher/enrichment/runs) — чтобы иметь доступ к инжектированным зависимостям.
|
# (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.
|
Returns run_id или None если уже есть running run для этого source.
|
||||||
Общий helper для всех source-handler'ов. Advisory-lock concurrency-логика перенесена
|
Общий helper для всех source-handler'ов. Advisory-lock concurrency-логика перенесена
|
||||||
ДОСЛОВНО из боевого scheduler'а (создание run — через инжектированный ctx.runs).
|
ДОСЛОВНО из боевого scheduler'а (создание run — через инжектированный ctx.runs).
|
||||||
|
|
||||||
|
#2658: каждая ветка «пропустили окно» пишет строку scrape_runs(status='skipped')
|
||||||
|
с машиночитаемой причиной — раньше пропуск был виден только в docker-логах, которые
|
||||||
|
теряются при редеплое. Пишем через kit-копию runs (у инжектированной app-копии
|
||||||
|
mark_skipped нет — строки-пропуски создаёт только планировщик).
|
||||||
"""
|
"""
|
||||||
source = schedule_row["source"]
|
source = schedule_row["source"]
|
||||||
if has_running_run(db, source):
|
if has_running_run(db, source):
|
||||||
logger.info("scheduler: skip — already running for source=%s", 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
|
return None
|
||||||
|
|
||||||
# #750: атомарный claim через transaction-scoped advisory lock. Сериализует
|
# #750: атомарный claim через transaction-scoped advisory lock. Сериализует
|
||||||
|
|
@ -272,6 +292,12 @@ def _claim_run(db: Session, schedule_row: dict[str, Any], ctx: SchedulerContext)
|
||||||
).scalar()
|
).scalar()
|
||||||
if not got_lock:
|
if not got_lock:
|
||||||
logger.info("scheduler: skip — concurrent claim in progress for source=%s", source)
|
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
|
return None
|
||||||
|
|
||||||
# Double-checked под локом: конкурентный тик мог закоммитить running-run МЕЖДУ
|
# 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):
|
if has_running_run(db, source):
|
||||||
logger.info("scheduler: skip — running appeared under lock for source=%s", source)
|
logger.info("scheduler: skip — running appeared under lock for source=%s", source)
|
||||||
db.rollback() # освобождаем advisory lock (claim не состоялся)
|
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
|
return None
|
||||||
|
|
||||||
params = schedule_row.get("default_params") or {}
|
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)
|
handler = resolve_handler(source, registry)
|
||||||
if handler is None:
|
if handler is None:
|
||||||
logger.warning("scheduler: unknown source=%s, skip", source)
|
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:
|
else:
|
||||||
await _dispatch(handler, db, sch, ctx)
|
await _dispatch(handler, db, sch, ctx)
|
||||||
|
|
||||||
|
|
@ -730,6 +771,10 @@ async def scheduler_loop(ctx: SchedulerContext, registry: Mapping[str, Handler])
|
||||||
|
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"SCHEDULER_TICK_SEC",
|
"SCHEDULER_TICK_SEC",
|
||||||
|
"SKIP_ALREADY_RUNNING",
|
||||||
|
"SKIP_CONCURRENT_CLAIM",
|
||||||
|
"SKIP_RUNNING_UNDER_LOCK",
|
||||||
|
"SKIP_UNKNOWN_SOURCE",
|
||||||
"ZOMBIE_THRESHOLD_HOURS",
|
"ZOMBIE_THRESHOLD_HOURS",
|
||||||
"Handler",
|
"Handler",
|
||||||
"SchedulerContext",
|
"SchedulerContext",
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue