fix(tradein/scraper): фильтр skipped в админке + освежение схлопнутой строки (#2658)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 7s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Successful in 1m3s
CI Trade-In / backend-tests (pull_request) Successful in 2m46s
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 7s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Successful in 1m3s
CI Trade-In / backend-tests (pull_request) Successful in 2m46s
Правки по ревью PR #2662. Фильтр статуса. `GET /admin/scrape/runs?status=skipped` отдавал 422 — 'skipped' не было в Literal, а во фронте не было чипа. Строки рисовались, но задать вопрос «что сейчас пропускается» на единственной поверхности, построенной ровно для этого, было нельзя. Добавлено в оба места (translateStatus «пропущено» и нейтральный бейдж уже умели). Схлопывание освежает строку. UPDATE двигал только finished_at/heartbeat_at, из-за чего живой стрик замерзал: списки прогонов сортируют ORDER BY started_at DESC и берут limit=20, поэтому 37-дневный пропуск утонул бы под свежими прогонами других источников — след в базе есть, на экране нет. Теперь started_at = NOW(), а начало стрика переезжает в counters.first_skip_at; сортировку общего списка не трогаем (она про все источники, чинить надо было одну строку). Там же обновляется counters.detail — иначе в строке 37 дней висел текст «протухли 1 день назад», хотя именно эта цифра и есть предмет issue. jsonb_set заменён на `||` + jsonb_build_object: три вложенных jsonb_set читать в 3 ночи невозможно, а NULL в jsonb_set обнуляет весь counters. Поиск последней строки. `ORDER BY id DESC` не ложится на индекс (source, started_at DESC) из миграции 015 — для unknown_source (тикает каждые 60 с бессрочно) это отбор всех строк источника с сортировкой раз в минуту. Теперь ORDER BY started_at DESC, id DESC. session_expires_at получил valid_only: предупреждение «скоро протухнут» считает срок ИМЕННО той записи, которую взял load_session — при нескольких аккаунтах свежайшая-любая может быть чужой протухшей строкой. Диагностика после None по-прежнему смотрит на свежайшую любую (валидных там нет по определению). Запись пропуска намеренно НЕ обёрнута в свой try/except: если db.execute падает, то падает и claim следующего расписания в этом же тике — тик срывается в любом случае, а глушить исключение здесь значило бы вернуть ровно тот немой пропуск, ради которого заведён #2658. Самовосстановление через 60 с.
This commit is contained in:
parent
0b54b96984
commit
b800760c24
6 changed files with 109 additions and 18 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.
|
||||
"""
|
||||
|
|
|
|||
|
|
@ -300,22 +300,30 @@ def load_session(db: Session) -> dict[str, str] | None:
|
|||
return cookies
|
||||
|
||||
|
||||
def session_expires_at(db: Session) -> datetime | None:
|
||||
"""Когда протухают самые свежезагруженные куки — БЕЗ фильтра валидности (#2658).
|
||||
def session_expires_at(db: Session, *, valid_only: bool = False) -> datetime | None:
|
||||
"""Когда протухают самые свежезагруженные куки (#2658).
|
||||
|
||||
`load_session` отбирает только ещё валидные записи (expires_at_estimate > NOW()) и на
|
||||
протухших отдаёт None — вызывающий не мог отличить «кук никогда не загружали» от
|
||||
«протухли позавчера» и не мог предупредить ЗАРАНЕЕ. Здесь фильтра нет: 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
|
||||
|
|
|
|||
|
|
@ -121,7 +121,9 @@ async def _cian_pre_claim(db: Session, schedule_row: dict[str, Any], ctx: Schedu
|
|||
return False
|
||||
|
||||
# Куки рабочие — предупреждаем, пока есть время их обновить без простоя сбора.
|
||||
expires_at = session_expires_at(db)
|
||||
# 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 дн.) — обновите заранее, "
|
||||
|
|
|
|||
|
|
@ -115,6 +115,40 @@ def test_mark_skipped_collapses_consecutive_same_reason() -> None:
|
|||
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 ──────────────────────────────────────────
|
||||
|
||||
|
||||
|
|
@ -394,6 +428,9 @@ async def test_cian_pre_claim_warns_before_expiry_not_after() -> None:
|
|||
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)
|
||||
|
||||
|
|
@ -469,6 +506,24 @@ def test_zero_result_monitor_ignores_skipped_rows() -> None:
|
|||
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)
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -248,9 +248,17 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
`details` — человеческое пояснение, кладётся в `counters.detail`.
|
||||
|
||||
Схлопывание подряд идущих одинаковых пропусков: если ПОСЛЕДНЯЯ строка прогона этого
|
||||
source — уже 'skipped' с тем же reason, новая не создаётся; у существующей
|
||||
обновляется finished_at и счётчик `counters.skips`. Без этого «уже бежит» и
|
||||
«неизвестный source» плодили бы строку каждый тик (60 с), пока держится причина.
|
||||
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 строки (новой или обновлённой).
|
||||
"""
|
||||
|
|
@ -260,16 +268,19 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
WITH latest AS (
|
||||
SELECT id FROM scrape_runs
|
||||
WHERE source = :source
|
||||
ORDER BY id DESC
|
||||
ORDER BY started_at DESC, id DESC
|
||||
LIMIT 1
|
||||
)
|
||||
UPDATE scrape_runs r
|
||||
SET heartbeat_at = NOW(),
|
||||
started_at = NOW(),
|
||||
finished_at = NOW(),
|
||||
counters = jsonb_set(
|
||||
COALESCE(r.counters, '{}'::jsonb),
|
||||
'{skips}',
|
||||
to_jsonb(COALESCE(CAST(r.counters ->> 'skips' AS int), 0) + 1)
|
||||
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
|
||||
|
|
@ -278,7 +289,7 @@ def mark_skipped(db: Session, *, source: str, reason: str, details: str | None =
|
|||
RETURNING r.id
|
||||
"""
|
||||
),
|
||||
{"source": source, "reason": reason},
|
||||
{"source": source, "reason": reason, "details": details},
|
||||
).fetchone()
|
||||
|
||||
if row is None:
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue