"""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. ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции и не двигается, сколько бы та ни жила. Финализаторы (mark_done/mark_failed/mark_banned) выполняются ТОЙ ЖЕ сессией, что и работа задачи, — и если рабочая транзакция всё это время оставалась открытой (задача ничего не коммитила: нечего было сохранять, батч читающий, сохранение шло чужой сессией), их UPDATE попадал ВНУТРЬ неё, и `finished_at` получал время НАЧАЛА работы, а не её конца. Замер на проде 2026-08-06 (487 прогонов, у которых есть и finished_at, и счётчик counters.duration_sec): у 153 заявленная длительность превышала собственное окно finished_at − started_at более чем в 1.5 раза, у 133 окно было меньше секунды при работе дольше 10 с. 126 из этих 133 окон лежат в диапазоне 9-64 мс — это не разброс, а подпись механизма: столько проходит от коммита claim'а до первого запроса рабочей транзакции. Крайний случай — прогон 346 (cian_history_backfill): 18230 с работы, окно 32 мс. Дефект был не сплошной ровно потому, что зависел от того, коммитила ли задача перед финалом: cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят поштучно, у них окно совпадало с работой; yandex_address_backfill (45 из 50 прогонов), newbuilding_enrich, cian_history_backfill — нет. Побочно это чинит и `heartbeat_at`: он писался тем же `now()` и по той же причине отставал от реальности на возраст открытой транзакции, а на нём стоит поиск зависших прогонов (reap_zombies, порог 6 ч). """ from __future__ import annotations import json import logging from collections.abc import Callable, Mapping from functools import cache 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 # #2670: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик прерывается # не только в принципе, но и на практике. Оба сторожа ниже слали алерт ровно на N-й # подряд неудаче и дальше молчали навсегда — а у постоянно сломанного источника # «дальше» длится месяцами. Прод 2026-08-06: у avito_full_load 31 неудача подряд, # последний успешный прогон 03.07 (34 дня без сбора), алерт был ровно один — на # третьей; у avito_full_load_exhaustive 5 подряд. Тишина при этом неотличима от # «всё хорошо» — ровно та ловушка, из-за которой #2574 месяц выглядела как норма. # # Вместо «ровно N» — разреженная лестница напоминаний: N, 2N, 4N, 8N…, а дальше не # реже, чем раз в STREAK_ALERT_MAX_PERIOD×N прогонов. Лестница по ПРОГОНАМ, а не # «раз в сутки», потому что источники идут разным тактом: domclick_city_sweep — раз # в день, proxy_healthcheck — раз в полчаса; календарное разрежение для одного из # них всегда будет либо спамом, либо молчанием. STREAK_ALERT_MAX_PERIOD = 16 # Потолок сканирования истории источника при подсчёте стрика. Достигнутый потолок # сам по себе повод для алерта (стрик заведомо огромен) — так «замолчать навсегда» # невозможно по построению, а не по счастливому совпадению чисел. STREAK_SCAN_LIMIT = 500 def _streak_alert_due(streak: int, threshold: int) -> bool: """Достиг ли стрик очередной вехи напоминания (#2670). True на threshold, 2×, 4×, 8×… и дальше на каждом кратном STREAK_ALERT_MAX_PERIOD×threshold. Первый алерт приходит там же, где и раньше — на N-й подряд неудаче; меняется только то, что он не последний. """ if streak < threshold or streak % threshold: return False mult = streak // threshold if mult % STREAK_ALERT_MAX_PERIOD == 0: return True return mult & (mult - 1) == 0 def _leading_streak(rows: list[Any], is_bad: Callable[[Any], bool]) -> int: """Длина серии подряд идущих «плохих» строк с начала списка (свежие — первыми).""" streak = 0 for row in rows: if not is_bad(row): break streak += 1 return streak # #2686: диагноз оборванного прогона. Пишется в scrape_runs.ban_kind (миграция 218) # РЯДОМ со status='banned', а не ВМЕСТО него — сознательный выбор между «новый # статус» и «явное поле причины»: # 1. Побочная функция 'banned' — сохранение done_buckets-чекпоинта (mark_failed # его теряет) — нужна ОБОИМ исходам. Оставив статус, получаем её даром; расщепив # статус, пришлось бы дублировать её в каждом потребителе. # 2. Новое значение статуса пришлось бы доучить пяти местам, каждое из которых # молча даёт неверный ответ, если про него забыть: CHECK-констрейнт схемы, # IN-списки обоих сторожей (_alert_if_consecutive_failures / _zero_results), # Literal-фильтр admin API и хардкод-список статусов во фронте. Это ровно тот # класс оборванной проводки, из-за которого задача и появилась. # 3. Прогон в обоих случаях требует одного и того же обращения (оборвать, сохранить # частичное); различается только ДИАГНОЗ — то есть метаданное, не состояние. BAN_KIND_PLATFORM = "platform" # площадка показала firewall/403/captcha — внешнее BAN_KIND_INFRA = "infra" # наш сайдкар/прокси не отдал страницу — внутреннее def _pick_int(counters: Mapping[str, Any], *keys: str) -> int | None: """Первое присутствующее из ``keys`` как 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 # #2703: ключи, которыми задача сообщает СВОЙ бизнес-результат. Список намеренно # короткий и состоит из синонимов ОДНОЙ величины — «сколько объявлений отдала выдача»: # total_seen — если задача посчитала сама; # lots_fetched — все city/newbuilding-sweep'ы (21 источник, 455 прогонов на проде); # unique_fetched — full-load'ы avito/cian/yandex (4 источника, 133 прогона) — раньше # сторож их не видел, хотя у cian_full_load 6 из 38 успешных прогонов # реально дали ноль. # Сводить сюда счётчики ОСТАЛЬНЫХ задач бессмысленно: на проде 28 источников (2650 # прогонов) не имеют общего результатного ключа вовсе — у каждого свой словарь # (deactivated / rows_written / poi_loaded / snapshotted / upserted / listings_matched # …), а у refresh_search_matview counters пусты буквально ({} во всех 55 строках) и у # трёх мониторов результата нет по смыслу. Ноль у них — часто ЗДОРОВЫЙ ответ # (deactivate_stale_* без протухших объявлений). Поэтому сторож не угадывает их # словарь, а честно признаёт, что мерить нечем — см. _run_result_count. _RESULT_COUNTER_KEYS = ("total_seen", "lots_fetched", "unique_fetched") def _run_result_count(counters: Mapping[str, Any] | None) -> int | None: """Бизнес-результат прогона; **None = прогон его не сообщил** (≠ ноль). Ровно это различие и было потеряно: сторож читал колонку ``total_seen``, у которой DEFAULT 0, поэтому «не измерено» и «измерено, ноль» выглядели одинаково. """ return _pick_int(counters or {}, *_RESULT_COUNTER_KEYS) @cache def _warn_source_has_no_result_metric(source: str, keys: tuple[str, ...]) -> None: """Один раз на процесс: у источника нет ключа, по которому сторож судит (#2703). Не алерт — алертить не о чем, судить не о чем тоже. Это делает слепую зону ВИДИМОЙ: раньше её признаком был вечно молчащий сторож, выглядящий настроенным. """ logger.warning( "zero-result watchdog неприменим к source=%s: counters не содержат ни одного " "результатного ключа %s (есть: %s) — прогоны этого источника больше не считаются " "нулевыми по умолчанию (#2703)", source, _RESULT_COUNTER_KEYS, ", ".join(keys) or "<пусто>", ) 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 ← _RESULT_COUNTER_KEYS (total_seen / lots_fetched / unique_fetched) - new_count ← 'new_count' (если уже есть) иначе 'lots_inserted' Возвращает (total_seen, new_count); None для ключа, которого нет в counters — тогда соответствующая колонка не перезаписывается (COALESCE-семантика в UPDATE). """ return _run_result_count(counters), _pick_int(counters, "new_count", "lots_inserted") def _alert_if_consecutive_failures(db: Session, source: str) -> None: """Sentry alert на серию из CONSECUTIVE_FAILURE_ALERT_THRESHOLD неудач подряд (статусы 'failed'/'banned') у данного source. Anti-spam: не на каждой неудаче, а по разреженной лестнице вех (см. _streak_alert_due). До #2670 алерт приходил РОВНО на N-й неудаче и дальше не повторялся никогда: серия, ставшая длиннее порога, замолкала навсегда. На проде это дало avito_full_load — 31 неудача подряд, 34 дня без единого успешного прогона, один алерт за всё время. Стрик прерывается любым завершением, кроме failed/banned, — по данным прода это достижимо и достигается (у domclick_city_sweep текущий стрик равен 1 при 47 завершённых прогонах), поэтому лестница не вырождается в постоянный алерт. Best-effort: весь блок обёрнут в try/except — сбой запроса или неинициализированный Sentry НЕ должен нарушать вызывающий mark_* путь. """ if sentry_sdk is None: return n = CONSECUTIVE_FAILURE_ALERT_THRESHOLD try: # Завершённые (non-running) прогоны источника, самые свежие первыми. 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": STREAK_SCAN_LIMIT}, ).fetchall() streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) capped = streak >= STREAK_SCAN_LIMIT if not capped and not _streak_alert_due(streak, n): return sentry_sdk.capture_message( f"Scraper source '{source}' has {streak} 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, streak ) 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 дали ИЗМЕРЕННЫЙ нулевой результат (#2625). Отличается от _alert_if_consecutive_failures: статус здесь формально 'done' (errors_count=0) — деградация невидима существующему failed/banned алерту. Причина обычно капча/пустая выдача источника, у которого нет (или не сработал) детект блокировки (см. providers/cian/serp.py, providers/yandex/serp.py). Anti-spam: та же разреженная лестница вех, что у _alert_if_consecutive_failures (#2670) — N, 2N, 4N…, а не «ровно N и дальше тишина». #2703: анти-спам «один раз на стрик» безопасен ТОЛЬКО там, где стрик может прерваться. Сторож читал колонку total_seen (DEFAULT 0), которой у 28 из 53 источников не заполняет ничто — значит у них он читал 0 ВСЕГДА, в том числе у полностью успешного прогона, стрик не прерывался никогда, и после первого события сторож замолкал навсегда, продолжая выглядеть настроенным. Теперь признак берётся из counters, а «не измерено» (None) стрик ПРЕРЫВАЕТ — ложный вечный стрик стал невозможен по построению, а слепая зона логируется явно. 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/прогон без результатной метрики), не только успешным # сбором. counters, а НЕ колонка total_seen: у колонки DEFAULT 0, по ней # «не измерено» неотличимо от «ноль» (#2703). rows = db.execute( text( """ SELECT status, counters 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": STREAK_SCAN_LIMIT}, ).fetchall() if not rows: return def _is_zero_done(r: Any) -> bool: """Только ИЗМЕРЕННЫЙ ноль. Прогон без результатной метрики стрик ПРЕРЫВАЕТ. Так недостижимое условие прерывания невозможно по построению: источник, чей словарь счётчиков сторожу неизвестен, не копит ложный стрик и не запирает анти-спам «один раз на стрик» в «один раз навсегда». """ return r.status == "done" and _run_result_count(r.counters) == 0 if _run_result_count(rows[0].counters) is None: # Свежайший завершённый прогон не сообщил результата — судить нечем. # Логируем (один раз на источник за процесс) вместо молчаливого нуля. _warn_source_has_no_result_metric(source, tuple(sorted(rows[0].counters or {}))) return streak = _leading_streak(rows, _is_zero_done) capped = streak >= STREAK_SCAN_LIMIT if not capped and not _streak_alert_due(streak, n): return sentry_sdk.capture_message( f"Scraper source '{source}' has {streak} 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, streak, ) 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=clock_timestamp()). started_at пишется СВОЕЙ транзакцией (db.commit() ниже) — откат рабочей транзакции задачи его уже не достаёт (#2702). Вид прогона несёт сам `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), clock_timestamp(), clock_timestamp() ) 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 + 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 = clock_timestamp(), 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 + 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 = clock_timestamp(), heartbeat_at = clock_timestamp(), 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 = clock_timestamp(), heartbeat_at = clock_timestamp(), 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], *, ban_kind: str = BAN_KIND_PLATFORM, ) -> None: """Финализация run: status='banned' + диагноз ban_kind (#2686). Per migration 015 — 'banned' задокументирован как 'Avito вернул 403/captcha'. Отличается от 'failed': прогон оборван внешним/блокирующим условием, а не нашим багом, и — важно — СОХРАНЯЕТ done_buckets-чекпоинт в counters (mark_failed его теряет). Cooldown 2-4 часа. `ban_kind` разводит два исхода, которые раньше схлопывались в один статус: - BAN_KIND_PLATFORM — площадка нас заблокировала (firewall/403/captcha); - BAN_KIND_INFRA — упала НАША инфраструктура (браузерный сайдкар/прокси). Значение приходит от места ПОРОЖДЕНИЯ отказа (тип исключения), а не из разбора текста ошибки. Default 'platform' = историческая семантика статуса, поэтому вызывающие, которым разводить нечего, не меняются. Оба исхода одинаково сохраняют чекпоинт — они отличаются только диагнозом. 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 = clock_timestamp(), heartbeat_at = clock_timestamp(), error = :error, counters = CAST(:counters AS jsonb), ban_kind = :ban_kind, 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), "ban_kind": ban_kind, "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_backfill_finished( db: Session, run_id: int, counters: dict[str, int], *, source: str, aborted_by_blocks: bool = False, fail_hint: str | None = None, ) -> None: """Честный финал detail-backfill'а (#2674): нулевой прогон ≠ 'done'. Все три detail-backfill'а (avito/yandex/domclick) финализировались ОДНИМ mark_done: прогон, который сделал N попыток и не обогатил НИ ОДНОГО объявления, отчитывался успехом. На проде (2026-08-06) это 78 прогонов из 158 — avito 23/76 (в т.ч. 5 прогонов по 1500-1600 попыток с нулём обогащений), yandex 31/52 (все attempted=5 failed=5), domclick 24/30 (494 попытки → 0). Существующие алерты этот класс не ловили: _alert_if_consecutive_failures считает только failed/banned, а _alert_if_consecutive_zero_results смотрит total_seen, которого в counters backfill'ов нет вовсе (всегда 0 → стрик не прерывается никогда → анти-спам молчит после первого раза). Правила (порядок важен), по образцу #2657 для domclick_city_sweep: - попыток не было (attempted=0) → 'done', честная пустота: кандидатов нет; - есть блоки источника И (прогон оборван брейкером ИЛИ ноль результата) → 'banned': external constraint, не наш баг (и триггер ротации IP #2611); - ноль результата без блоков → 'failed': это наша поломка (парсер/сеть/БД); - иначе (обогатили хоть что-то) → 'done', в т.ч. частичный прогон. `gone` (404 у avito) считается результатом наравне с `enriched`: прогон, который подтвердил снятие объявлений, работу сделал. `fail_hint` — самая частая причина отказа этого прогона (задача считает её сама, см. avito_detail_backfill._failure_signature). Дописывается в текст статуса, потому что «blocked=5, обогащено 0» не отвечает на единственный вопрос, ради которого статус и читают: отказала площадка или наш тракт (#2686, #2698). Логи контейнера на этот вопрос отвечать не могут — они исчезают при пересоздании контейнера, то есть на первом же деплое после ночного прогона. """ attempted = int(counters.get("attempted") or 0) enriched = int(counters.get("enriched") or 0) blocked = int(counters.get("blocked") or 0) produced = enriched + int(counters.get("gone") or 0) hint = f"; причина: {fail_hint}" if fail_hint else "" if attempted == 0: mark_done(db, run_id, counters) return if blocked and (aborted_by_blocks or produced == 0): reason = ( f"backfill-honest-status: {source} остановлен блоками источника — " f"blocked={blocked}, обогащено {enriched} из {attempted} попыток{hint} (#2674)" ) logger.error("%s run_id=%d", reason, run_id) mark_banned(db, run_id, reason, counters) return if produced == 0: reason = ( f"backfill-honest-status: {source} без результата — 0 обогащено из " f"{attempted} попыток (failed={counters.get('failed', 0)}, " f"blocked={blocked}){hint} (#2674)" ) logger.error("%s run_id=%d", reason, run_id) mark_failed(db, run_id, reason, counters) return mark_done(db, run_id, counters) 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 = clock_timestamp() 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, ban_kind, 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]