gendesign/tradein-mvp/backend/app/services/scrape_runs.py
bot-backend 3c5f535e6c
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI / changes (pull_request) Successful in 8s
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / frontend-checks (pull_request) Successful in 1m2s
CI Trade-In / backend-tests (pull_request) Successful in 2m58s
fix(tradein/admin): гейт отмены по источнику, честный комментарий view, лимит 50 (#2674)
Ревью PR #2684 — четыре MINOR.

1. Починка фильтра открыла кнопку отмены на все 53 источника. Раньше таблица была
   пуста на каждой вкладке, поэтому кнопка не рендерилась НИ РАЗУ и дыра не
   проявлялась: ручки отмены source не проверяют вовсе. Оператор на вкладке Авито
   мог бы «отменить» refresh_search_matview — задача продолжила бы работать под
   статусом 'cancelled' (ещё один врущий статус ровно в тот день, когда их
   вычищаем), а has_running_run перестал бы держать single-run guard, который
   существует из-за инцидента с двойным свипом и баном (2026-05-31).

   Гейт поставлен на общем узле всех пяти ручек — scrape_runs.honors_cancel +
   отказ в mark_cancelled, — а не в UI: иначе ручной POST по-прежнему снимал бы
   guard. Флаг cancellable отдаётся в строке, UI по нему прячет кнопку.
   Состав набора выведен из call-site'ов runs.is_cancelled: city-sweep'ы (все
   площадки и города), full-load'ы, avito_newbuilding_sweep, rosreestr_dkp_import.
   Правило НЕ «любой *_sweep»: yandex_newbuilding_sweep отмену не опрашивает.

2. Комментарий пересозданного v_data_quality утверждал, что его обновляет
   /api/v1/admin/data-quality. Читателей у view нет ни одного — живая ручка строит
   свой запрос. PR с тезисом «ложный показатель хуже отсутствующего» не имеет права
   переносить в прод ложное утверждение о читателе.

3. Лимит выдачи 20 → 50: первые 20 строк по started_at на три четверти —
   сердцебиение proxy_healthcheck (1631 из 3245), часовой сбор мог не поместиться.
   Привязка к вкладке НЕ возвращается.

4. Тест «действующее определение view» искал маркер подстрокой с OR REPLACE —
   миграция с обычным CREATE VIEW или парой DROP+CREATE была бы невидима, и тест
   проверял бы 214, пока показатель уже вернулся. Заменено регуляркой на обе формы.

Фальсификация трёх новых тестов патч-методом — все три красные. Полный прогон
3490 passed / 9 skipped, tsc --noEmit чистый.
2026-08-06 03:57:46 +05:00

522 lines
24 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

"""Utility functions для scrape_runs table — tracking long-running pipeline runs.
Таблица scrape_runs создана в 015_scrape_runs.sql.
Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled.
"""
from __future__ import annotations
import json
import logging
from collections.abc import Callable
from typing import Any
import sentry_sdk
from sqlalchemy import text
from sqlalchemy.orm import Session
logger = logging.getLogger(__name__)
# Количество последовательных неудачных запусков (failed/banned), при достижении
# которого отправляется Sentry alert. Алерт срабатывает ровно при N-й подряд ошибке
# (anti-spam: не на каждой последующей).
CONSECUTIVE_FAILURE_ALERT_THRESHOLD = 3
# #2625: количество последовательных 'done' запусков с нулевым бизнес-результатом
# (total_seen=0), при достижении которого отправляется Sentry alert. Статус 'done'
# формально успешен (errors_count=0), но капча/пустая выдача/смена вёрстки источника
# без детекта (см. providers/cian, providers/yandex) деградируют молча — этот класс
# невидим для CONSECUTIVE_FAILURE_ALERT_THRESHOLD (тот считает только failed/banned).
CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD = 3
def _column_counts(counters: dict[str, int]) -> tuple[int | None, int | None]:
"""Извлечь значения для dedicated-колонок total_seen / new_count из jsonb-counters.
Все sweep-counters (YandexCitySweepCounters / CitySweepCounters / …) пишут число
объявлений в выдаче как ``lots_fetched`` и число впервые вставленных как
``lots_inserted``. Раньше эти значения попадали ТОЛЬКО в counters jsonb, а
выделенные колонки total_seen/new_count оставались 0 → admin/observability
показывала total_seen=0 при реально сохранённых строках (audit #1871/#1926).
Приоритет ключей:
- total_seen ← 'total_seen' (если уже есть в counters) иначе 'lots_fetched'
- new_count ← 'new_count' (если уже есть) иначе 'lots_inserted'
Возвращает (total_seen, new_count); None для ключа, которого нет в counters —
тогда соответствующая колонка не перезаписывается (COALESCE-семантика в UPDATE).
"""
def _pick(*keys: str) -> int | None:
for key in keys:
val = counters.get(key)
if val is not None:
try:
return int(val)
except (TypeError, ValueError):
return None
return None
return _pick("total_seen", "lots_fetched"), _pick("new_count", "lots_inserted")
def _alert_if_consecutive_failures(db: Session, source: str) -> None:
"""Отправить Sentry alert если последние CONSECUTIVE_FAILURE_ALERT_THRESHOLD
завершённых запусков для данного source имеют статус 'failed' или 'banned'.
Anti-spam: алерт срабатывает ТОЛЬКО когда стрик РОВНО равен порогу — т.е. запрос
возвращает ровно N последних (failed|banned) и (N+1)-й, если существует, НЕ является
failed/banned. Это предотвращает повторный алерт на каждой ошибке сверх порога.
Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный
Sentry НЕ должен нарушать вызывающий mark_* путь.
"""
n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD
try:
# Берём последние N+1 завершённых (non-running) запусков по source.
# Сортируем по finished_at DESC чтобы самые свежие шли первыми.
rows = db.execute(
text(
"""
SELECT status FROM scrape_runs
WHERE source = :source
AND status IN ('failed', 'banned', 'done', 'cancelled')
ORDER BY finished_at DESC NULLS LAST
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
).fetchall()
if len(rows) < n:
# Ещё не набралось N завершённых запусков вообще — алерт не нужен.
return
# Первые N должны быть все failed/banned.
first_n = rows[:n]
if not all(r.status in ("failed", "banned") for r in first_n):
return
# (N+1)-й запуск, если есть, тоже должен НЕ быть failed/banned — иначе мы уже
# должны были отправить алерт раньше и не стоит дублировать.
if len(rows) > n and rows[n].status in ("failed", "banned"):
return
# Стрик ровно достиг порога — отправляем алерт.
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive failed/banned runs — "
"manual intervention may be required (expired cookies / ban / broken parser).",
level="error",
)
logger.error(
"sentry alert sent: source=%s has %d consecutive failed/banned runs", source, n
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only
def _alert_if_consecutive_zero_results(db: Session, source: str) -> None:
"""Отправить Sentry alert если последние CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
завершённых 'done' запусков для source имеют total_seen=0 (#2625).
Отличается от _alert_if_consecutive_failures: статус здесь формально 'done'
(errors_count=0) — деградация невидима существующему failed/banned алерту.
Причина обычно капча/пустая выдача источника, у которого нет (или не сработал)
детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py).
Anti-spam: тот же N-й-стрик паттерн, что у _alert_if_consecutive_failures —
алерт срабатывает ровно когда стрик достигает порога, не на каждом запуске сверх.
Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный
Sentry НЕ должен нарушать вызывающий mark_done путь.
"""
n = CONSECUTIVE_ZERO_RESULT_ALERT_THRESHOLD
try:
# Те же non-running статусы, что у _alert_if_consecutive_failures — стрик
# 'done'-с-нулём прерывается ЛЮБЫМ другим завершением (failed/banned/done-
# с-результатом/cancelled), не только успешным сбором.
rows = db.execute(
text(
"""
SELECT status, total_seen FROM scrape_runs
WHERE source = :source
AND status IN ('failed', 'banned', 'done', 'cancelled')
ORDER BY finished_at DESC NULLS LAST
LIMIT :limit
"""
),
{"source": source, "limit": n + 1},
).fetchall()
if len(rows) < n:
return
def _is_zero_done(r: Any) -> bool:
return r.status == "done" and (r.total_seen or 0) == 0
first_n = rows[:n]
if not all(_is_zero_done(r) for r in first_n):
return
if len(rows) > n and _is_zero_done(rows[n]):
return
sentry_sdk.capture_message(
f"Scraper source '{source}' has {n} consecutive 'done' runs with zero "
"lots fetched — captcha/layout-change likely undetected "
"(manual check recommended).",
level="error",
)
logger.error(
"sentry alert sent: source=%s has %d consecutive zero-result 'done' runs",
source,
n,
)
except Exception:
pass # sentry_sdk not initialised in dev, or query failed — best-effort only
def _alert_on_run_id(
db: Session,
run_id: int,
*,
checker: Callable[[Session, str], None] = _alert_if_consecutive_failures,
) -> None:
"""Вспомогательная обёртка: извлекает source по run_id и вызывает `checker`
(default _alert_if_consecutive_failures; mark_done передаёт
_alert_if_consecutive_zero_results — #2625). Best-effort — не бросает исключений.
"""
try:
row = db.execute(
text("SELECT source FROM scrape_runs WHERE id = :run_id"),
{"run_id": run_id},
).fetchone()
if row is None:
return
checker(db, str(row.source))
except Exception:
pass
def create_run(db: Session, *, source: str, params: dict[str, Any]) -> int:
"""INSERT scrape_runs(source, status='running', params, started_at=NOW()).
Вид прогона несёт сам `source` (avito_city_sweep / domclick_detail_backfill / …);
отдельной колонки run_type больше нет — она 3244 прогона подряд молчала
дефолтом 'city_sweep' и подписывала им, например, proxy_healthcheck (#2674).
Returns run_id (bigint).
"""
row = db.execute(
text(
"""
INSERT INTO scrape_runs (source, status, params, started_at, heartbeat_at)
VALUES (:source, 'running', CAST(:params AS jsonb), NOW(), NOW())
RETURNING id
"""
),
{"source": source, "params": json.dumps(params, ensure_ascii=False)},
).fetchone()
db.commit()
assert row is not None, "scrape_runs INSERT returned no id"
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 колонки.
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и
пишутся в выделенные колонки, чтобы observability не показывала 0 (audit #1926).
COALESCE: если ключа нет в counters — старое значение колонки сохраняется.
"""
total_seen, new_count = _column_counts(counters)
db.execute(
text(
"""
UPDATE scrape_runs
SET heartbeat_at = NOW(),
counters = CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id
"""
),
{
"run_id": run_id,
"counters": json.dumps(counters),
"total_seen": total_seen,
"new_count": new_count,
},
)
db.commit()
# Источники, чей джоб РЕАЛЬНО опрашивает status='cancelled' в своём цикле.
# Всё остальное отменить нельзя: строка стала бы 'cancelled', а задача продолжила бы
# работать — это, во-первых, ещё один врущий статус, во-вторых (хуже) обход guard'а
# has_running_run: он перестанет видеть прогон как running и пустит второй свип на том
# же прокси-IP → бан (инцидент 2026-05-31, runs #26+#27).
# Состав проверен по call-site'ам runs.is_cancelled: kit pipeline (city-sweep'ы всех
# площадок и городов, full-load'ы, avito_newbuilding_sweep) + rosreestr_dkp_import
# (scheduler.py). yandex_newbuilding_sweep отмену НЕ опрашивает — поэтому правило не
# «любой *_sweep». Актуально с #2674: до починки фильтра таблица прогонов была пуста
# на всех вкладках, кнопка отмены не рендерилась ни разу и дыра не проявлялась.
_CANCEL_HONORING_EXACT = frozenset({"avito_newbuilding_sweep", "rosreestr_dkp_import"})
_CANCEL_HONORING_SUBSTRINGS = ("city_sweep", "full_load")
def honors_cancel(source: str) -> bool:
"""True, если джоб этого source опрашивает отмену и реально остановится."""
return source in _CANCEL_HONORING_EXACT or any(
key in source for key in _CANCEL_HONORING_SUBSTRINGS
)
def is_cancelled(db: Session, run_id: int) -> bool:
"""Проверить status='cancelled' (cooperative cancel в long-running pipeline)."""
row = db.execute(
text("SELECT status FROM scrape_runs WHERE id = :id"),
{"id": run_id},
).fetchone()
return row is not None and row.status == "cancelled"
def mark_done(db: Session, run_id: int, counters: dict[str, int]) -> None:
"""Финализация run: status='done', finished_at=NOW(), counters + total_seen/new_count.
total_seen/new_count извлекаются из counters (lots_fetched/lots_inserted) и пишутся
в выделенные колонки — иначе admin/observability показывает 0 (audit #1926).
"""
total_seen, new_count = _column_counts(counters)
row = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'done', finished_at = NOW(), heartbeat_at = NOW(),
counters = CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id AND status = 'running'
RETURNING id
"""
),
{
"run_id": run_id,
"counters": json.dumps(counters),
"total_seen": total_seen,
"new_count": new_count,
},
).first()
if row is None:
logger.warning("mark_done no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# #2625: N подряд 'done' с нулевым бизнес-результатом — деградация, невидимая
# для failed/banned алерта (капча/пустая выдача под видом успеха). Best-effort,
# после коммита — статус уже персистирован в БД.
_alert_on_run_id(db, run_id, checker=_alert_if_consecutive_zero_results)
def mark_failed(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None:
"""Финализация run: status='failed', error (первые 1000 символов).
Defensive rollback: если до этого вызова в той же транзакции был ошибочный UPDATE,
он мог оставить сессию в error state — rollback сбрасывает состояние.
"""
try:
db.rollback() # defensive: cascade-safe if prior UPDATE left txn in error state
except Exception:
pass
total_seen, new_count = _column_counts(counters)
row = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'failed', finished_at = NOW(), heartbeat_at = NOW(),
error = :error, counters = CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id AND status = 'running'
RETURNING id
"""
),
{
"run_id": run_id,
"error": error[:1000],
"counters": json.dumps(counters),
"total_seen": total_seen,
"new_count": new_count,
},
).first()
if row is None:
logger.warning("mark_failed no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
_alert_on_run_id(db, run_id)
def mark_banned(db: Session, run_id: int, error: str, counters: dict[str, int]) -> None:
"""Финализация run: status='banned' (IP заблокирован Avito — 403/captcha).
Per migration 015 — 'banned' задокументирован как 'Avito вернул 403/captcha'.
Отличается от 'failed': это external constraint, не наш bug. Cooldown 2-4 часа.
Defensive rollback: если до этого вызова в той же транзакции был ошибочный UPDATE,
он мог оставить сессию в error state — rollback сбрасывает состояние.
"""
try:
db.rollback() # defensive: cascade-safe if prior UPDATE left txn in error state
except Exception:
pass
total_seen, new_count = _column_counts(counters)
row = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'banned', finished_at = NOW(), heartbeat_at = NOW(),
error = :error, counters = CAST(:counters AS jsonb),
total_seen = COALESCE(CAST(:total_seen AS int), total_seen),
new_count = COALESCE(CAST(:new_count AS int), new_count)
WHERE id = :run_id AND status = 'running'
RETURNING id
"""
),
{
"run_id": run_id,
"error": error[:1000],
"counters": json.dumps(counters),
"total_seen": total_seen,
"new_count": new_count,
},
).first()
if row is None:
logger.warning("mark_banned no-op: run_id=%d not in 'running' state", run_id)
db.commit()
# Алерт после коммита — статус уже персистирован в БД; best-effort (не бросает).
_alert_on_run_id(db, run_id)
def mark_cancelled(db: Session, run_id: int) -> bool:
"""Set status='cancelled' если currently 'running'. Returns True если cancelled.
Отказ (False) для source'ов, чей джоб отмену не опрашивает — см. honors_cancel:
там 'cancelled' был бы враньём в статусе и снял бы has_running_run-guard.
Ручки отмены source не проверяют (любая из пяти принимает любой run_id), поэтому
гейт стоит здесь — на общем узле всех пяти.
"""
row = db.execute(
text("SELECT source FROM scrape_runs WHERE id = :run_id"),
{"run_id": run_id},
).fetchone()
if row is not None and not honors_cancel(str(row.source)):
logger.warning(
"mark_cancelled отказ: run_id=%d source=%s не опрашивает отмену — "
"задача продолжила бы работать под статусом 'cancelled'",
run_id,
row.source,
)
return False
result = db.execute(
text(
"""
UPDATE scrape_runs
SET status = 'cancelled', finished_at = NOW()
WHERE id = :run_id AND status = 'running'
RETURNING id
"""
),
{"run_id": run_id},
).fetchone()
db.commit()
return result is not None
def list_recent(db: Session, *, source: str, limit: int = 10) -> list[dict[str, Any]]:
"""Список последних N runs по source, в порядке убывания started_at."""
rows = (
db.execute(
text(
"""
SELECT id AS run_id, source, status, params, counters, error,
started_at, finished_at, heartbeat_at
FROM scrape_runs
WHERE source = :source
ORDER BY started_at DESC NULLS LAST
LIMIT :limit
"""
),
{"source": source, "limit": limit},
)
.mappings()
.all()
)
return [dict(r) for r in rows]
def list_all(
db: Session,
*,
source: str | None = None,
status: str | None = None,
limit: int = 50,
offset: int = 0,
) -> tuple[int, list[dict[str, Any]]]:
"""Unified-список runs по всем source'ам с опц. фильтрами source/status.
Возвращает (total, rows): total — кол-во строк под фильтром (без limit/offset),
rows — текущая страница (ORDER BY started_at DESC, LIMIT/OFFSET).
Используется единой scrapers-страницей вместо 3 per-source /runs.
"""
where = ["1=1"]
params: dict[str, Any] = {"limit": limit, "offset": offset}
if source is not None:
where.append("source = :source")
params["source"] = source
if status is not None:
where.append("status = :status")
params["status"] = status
where_sql = " AND ".join(where)
total_row = db.execute(
text(f"SELECT COUNT(*) AS n FROM scrape_runs WHERE {where_sql}"),
params,
).fetchone()
total = int(total_row.n) if total_row is not None else 0
rows = (
db.execute(
text(
f"""
SELECT id AS run_id, source, status, params, counters,
total_seen, new_count, started_at, finished_at,
heartbeat_at, error AS error_text
FROM scrape_runs
WHERE {where_sql}
ORDER BY started_at DESC NULLS LAST
LIMIT CAST(:limit AS int) OFFSET CAST(:offset AS int)
"""
),
params,
)
.mappings()
.all()
)
return total, [dict(r) for r in rows]
def distinct_sources(db: Session) -> list[str]:
"""Все значения source, которые РЕАЛЬНО есть в scrape_runs (по алфавиту).
#2674: фильтр источников в админке был захардкожен тремя площадками
(avito/cian/yandex), а в таблице 53 разных source и ни одной строки с таким
точным значением — все три пункта фильтра давали пустую выдачу, а 76%
прогонов (включая всю площадку Домклик) отфильтровать было нечем.
Список обязан приходить из данных: новый source появляется в фильтре сам,
без правки кода.
Игнорирует фильтры /scrape/runs — иначе выбор источника вырезал бы из
выпадающего списка все остальные.
"""
rows = db.execute(
text("SELECT DISTINCT source FROM scrape_runs WHERE source IS NOT NULL ORDER BY source")
).fetchall()
return [str(r.source) for r in rows]