"""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, Collection, 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" # наш сайдкар/прокси не отдал страницу — внутреннее # #2764: причина НЕ установлена. Дефолт mark_banned — именно он, а не 'platform': # на проде оба прогона, помеченных после миграции 218, получили 'platform' по # умолчанию (ни один их не передавал), то есть метка выглядела доказательством, не # будучи им. 'unknown' делает пробел измеримым (SELECT ban_kind, count(*)), а # 'platform'/'infra' начинают означать ровно то, что доказано типом исключения. BAN_KIND_UNKNOWN = "unknown" 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 успешных прогонов # реально дали ноль. # succeeded — yandex_newbuilding_sweep (42 прогона/90д) и newbuilding_enrich # (65 прогонов/90д, единственные два писателя ключа на проде, # проверено 2026-08-15). НЕ 'rows_inserted': тот ключ пишет ЕЩЁ и # rosreestr_dkp_import (67 прогонов/90д) — у него rows_inserted=0 в # 66 из 67 это ЗДОРОВЫЙ ответ догнавшего инкрементального импорта # (rows_fetched=rows_skipped=96974, last_id не двигается неделями), # а не отказ; если бы 'rows_inserted' попал в этот список, сторож # зачитывал бы этот здоровый ноль как измеренный провал и копил бы # практически непрерываемый стрик (rosreestr_dkp_import не # прерывается другим статусом — импорт либо 'done', либо не бежал). # НЕ 'processed' по той же причине с другой стороны: это счётчик # ПОПЫТОК (у newbuilding_enrich processed==attempted==limit даже # когда succeeded меньше — прод-факт 09.08: processed=25 succeeded=14, # 44% отказов замаскировались бы под measured-25) — сторож нулевого # результата на нём молчал бы ровно там, где должен сработать, а на # будущем опустении очереди домов (cian_houses_pending) создал бы # свой вечный ложный zero-стрик. 'succeeded' у yandex_newbuilding_sweep # численно совпадает с 'rows_inserted' на всех 42/42 прод-прогонах — # замена не теряет исходную цель (десять прогонов подряд 26.07-10.08, # все 'done', succeeded=0 rows_inserted=0 failed_resolve=4-5 — раньше # ни total_seen/lots_fetched/unique_fetched не было, и # _run_result_count всегда возвращал None (honest-run-status)). # Сводить сюда счётчики ОСТАЛЬНЫХ задач бессмысленно: на проде 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", "succeeded", ) 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 _sweep_run_did_nothing(counters: Mapping[str, Any]) -> str | None: """Развёртка, у которой КАЖДЫЙ якорь кончился отказом и не принесла ничего (#2625). Возвращает текст причины (для error) либо None, если прогон таким не является. Третий исход, у которого не было терминального статуса. Развёртка различает: 1. «площадка отбила» — попытки разбора были, структура не извлеклась ни разу → `mark_banned` в самих sweep'ах (#2642, cian/yandex); 2. «площадка честно отдала пустоту» — валидный ответ, ноль предложений → `done` с нулём, это здоровый результат (в Серове реально 10 объявлений); 3. «мы не дошли» — якорь упал по таймауту или исключению ДО того, как что-либо стало разбирать. Ровно этот случай в счётчики бана не попадает НАМЕРЕННО (#2600 п.1: transport_error не должен выглядеть баном площадки), и статуса ему никто не выдал — прогон уходил в `done`. Признак — собственная бухгалтерия прогона, а не список известных антибот-маркеров: `errors_count >= anchors_total` при нулевом ИЗМЕРЕННОМ результате означает, что отказом кончился каждый якорь, который у прогона был, и собрано ноль. Это НЕ доказывает, КТО виноват (капча площадки / наш прокси / наш баг), поэтому статус 'failed' без диагноза, а не 'banned' с 'platform' (#2764: диагноз не назначается по умолчанию). Что признак НЕ ловит: прогон, где часть якорей отдала данные, а часть отказала — `errors_count < anchors_total`, статус остаётся 'done' (частичный сбор — сбор). Замер на проде 2026-08-10 за 90 суток: под правило попадают 28 прогонов (yandex_city_sweep_nizhniy_tagil 16 подряд по 15-30.07 — каждый ровно 240 с, таймаут якоря, 0 лотов, 'done'; yandex_city_sweep 6; avito_city_sweep 5; yandex_city_sweep_pervouralsk 1 от 09.08 — 155 мс, исключение до первого запроса). НЕ затронуты: 132 прогона с отказами, но ненулевым сбором, и 37 прогонов честной пустоты (errors_count=0) — они остаются 'done'. """ anchors = _pick_int(counters, "anchors_total") errors = _pick_int(counters, "errors_count") if not anchors or anchors <= 0 or errors is None or errors < anchors: return None if _run_result_count(counters) != 0: # None (не измерено) сюда тоже НЕ попадает return None return ( f"sweep-honest-status: отказом кончились все {anchors} якорей прогона " f"(errors_count={errors}), собрано 0 — работа не сделана. Причина НЕ " f"установлена: якорь мог упасть по таймауту, из-за нашего прокси или " f"блокировкой площадки — статус 'failed' без диагноза (#2625)" ) # #2700: сколько попыток фазы должно быть, чтобы «отказали все» что-то значило. # 3 — не круглое число, а порог, на котором сам сбор уже сдаётся: столько подряд # неудачных detail'ов достаточно оркестратору, чтобы ротировать прокси и оборвать фазу # (_cian_detail_abort в orchestration/pipeline.py). Замер на проде 2026-08-10 за 90 # суток: порог отсекает 2 прогона с ЕДИНСТВЕННОЙ попыткой (одиночный отказ — шум, не # диагноз) и оставляет 50 прогонов, где отказали 3-50 попыток подряд. _PHASE_MIN_ATTEMPTS = 3 def _phase_totally_failed(counters: Mapping[str, Any]) -> str | None: """Фаза прогона, у которой отказала КАЖДАЯ попытка (#2700). Текст причины или None. Прогон состоит из фаз, а статус у него один. `_sweep_run_did_nothing` (#2625) ловит случай, когда не сделано НИЧЕГО; этот — когда целое направление работы отказало на сто процентов, а соседнее сработало, и суммарный ненулевой сбор прячет отказ. Живой повод (#2700): `cian_city_sweep` 15 суток подряд писал `detail_attempted=50, detail_failed=50, errors_count=0, status=done` — каждая detail-страница отдавала HTTP 403. Ноль обогащённых при 1 680 собранных лотах внешне неотличим от здорового прогона: результатный счётчик (lots_fetched) ненулевой, а до `errors_count` отказ подзадачи не доходил вовсе (403 гасился внутри провайдера в `return None`). Признак — собственная бухгалтерия фазы: `_failed == _attempted` при `attempted >= _PHASE_MIN_ATTEMPTS`. Пары ищутся В САМИХ counters (любой ключ `X_attempted` со спутником `X_failed`), а не по зашитому списку фаз: список — это ровно то место, куда забывают дописать новую фазу, и тогда сторож молчит, выглядя настроенным. На проде за 90 суток таких пар четыре: detail/houses/address/imv. Что признак НЕ доказывает: КТО виноват (площадка, наш прокси, наш парсер) — поэтому 'failed' без диагноза, как и в #2625/#2764, а не 'banned'/'platform'. Замер на проде 2026-08-10 за 90 суток, ПРОГНАННЫЙ УЖЕ ДЕПЛОЙНУТОЙ функцией по боевым counters (3 574 прогона, из них 3 293 'done'): правило переводит в 'failed' 42 прогона (1.3%) — 31 cian_city_sweep* и 11 avito_city_sweep*; про вторые никто не знал. Остальные 3 251 остаются 'done'. Первая версия этого абзаца называла 52 — это было число ПАР «прогон × фаза» из SQL-замера, а не прогонов: у 10 прогонов отказали обе фазы (detail и houses) сразу, и они посчитались дважды. """ for key in sorted(counters): if not key.endswith("_attempted"): continue phase = key[: -len("_attempted")] attempted = _pick_int(counters, key) failed = _pick_int(counters, f"{phase}_failed") if attempted is None or failed is None: continue if attempted >= _PHASE_MIN_ATTEMPTS and failed == attempted: return ( f"phase-honest-status: фаза '{phase}' отказала полностью — " f"{failed} из {attempted} попыток неудачны, обогащено 0. Остальные фазы " f"прогона могли отработать, поэтому ненулевой сбор это НЕ опровергает. " f"Причина НЕ установлена: блок площадки, наш прокси или разбор — статус " f"'failed' без диагноза (#2700)" ) return None # honest-run-status (2026-08-15): доля отказов, которая обесценивает формально ненулевой # сбор. Прод-факт avito_detail_backfill 15.08: {"attempted":64,"failed":57,"enriched":6, # "blocked":1} — 89% попыток отказали, а mark_backfill_finished всё равно звал mark_done, # потому что "produced != 0" (6 обогащено). Ни _sweep_run_did_nothing (нужны # anchors_total/errors_count, у backfill'ов их нет), ни _phase_totally_failed (нужна пара # "_attempted"/"_failed" — здесь голые "attempted"/"failed" без фазового # префикса, `"attempted".endswith("_attempted")` не матчит) эту форму counters не ловят — # обе проверки написаны под СВОИ формы, а не под backfill'овскую. # # Порог 'failed' — половина и больше отказов: сбор для практических целей провалился, # даже если несколько записей всё же обогатились. Порог 'partial' НЕ заведён отдельным # статусом scrape_runs.status — это потребовало бы миграции (DROP+ADD CHECK constraint, # 051_scrape_runs_extend.sql) и обучило бы новому значению ещё 4 места (Literal-фильтр # admin API, хардкод статусов фронта, оба IN-списка сторожей) — тот же класс "оборванной # проводки", из-за которого заведён #2686/ban_kind. Вместо статуса — тот же диагноз, что и # у ban_kind: causa в тексте `error`, терминальный статус один ('failed'). 0.15..0.5 — # та же 'failed', но с другой формулировкой причины ("деградировал", не "провалился"), чтобы # оператор видел разницу читая error, не только status. FAILED_RATIO_FAILED_THRESHOLD = 0.5 FAILED_RATIO_DEGRADED_THRESHOLD = 0.15 # Минимум попыток, при котором доля вообще что-то значит — иначе 1 отказ из 2 (=0.5) # палит статус на шуме единичного случая. То же рассуждение и то же число, что у # _PHASE_MIN_ATTEMPTS (см. выше). _FAILED_RATIO_MIN_ATTEMPTS = _PHASE_MIN_ATTEMPTS def _failed_ratio_too_high(counters: Mapping[str, Any]) -> str | None: """Прогон, у которого доля отказов слишком велика, даже если что-то собрано. Возвращает текст причины (для error) либо None. Читает ГОЛЫЕ ключи "attempted"/ "failed" (без фазового префикса) — сейчас это словарь только у четырёх detail-backfill'ов (avito/yandex/domclick/newbuilding_enrich), все идут через mark_backfill_finished → mark_done. `attempted < _FAILED_RATIO_MIN_ATTEMPTS` или отсутствие любого из ключей → None (нечем/не о чём судить — счётчики либо не заполнены, либо принадлежат другому источнику со своим словарём). Что признак НЕ доказывает: КТО виноват (площадка, наш прокси, наш парсер) — поэтому 'failed' без диагноза, как и у #2625/#2700/#2764. """ attempted = _pick_int(counters, "attempted") failed = _pick_int(counters, "failed") if attempted is None or failed is None or attempted < _FAILED_RATIO_MIN_ATTEMPTS: return None ratio = failed / max(attempted, 1) if ratio >= FAILED_RATIO_FAILED_THRESHOLD: verb = "провалился" elif ratio >= FAILED_RATIO_DEGRADED_THRESHOLD: verb = "деградировал" else: return None return ( f"failed-ratio-honest-status: сбор {verb} — {failed} из {attempted} попыток " f"отказали (доля {ratio:.0%}); формально ненулевой результат этого не искупает. " f"Причина НЕ установлена — статус 'failed' без диагноза" ) 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 / succeeded) - new_count ← 'new_count' / 'lots_inserted' / 'saved_inserted' / 'rows_inserted' (первый присутствующий). 'saved_inserted' — full-load'ы (cian/avito/yandex, CianFullLoadCounters и аналоги в pipeline.py): на проде витрина показывала new_count=0 у трёх подряд cian_full_load при реально сохранённых saved_inserted=482/214/239 (honest-run-status) — ключ 'new_count'/'lots_inserted' у full-load'ов в counters не пишется вовсе. 'rows_inserted' — тот же ключ, которым yandex_newbuilding_sweep и rosreestr_dkp_import сообщают число upsert'ов; здесь (для витринной колонки new_count) это безопасно — в отличие от _RESULT_COUNTER_KEYS этот список не участвует в подсчёте zero-result-стрика. Возвращает (total_seen, new_count); None для ключа, которого нет в counters — тогда соответствующая колонка не перезаписывается (COALESCE-семантика в UPDATE). """ return _run_result_count(counters), _pick_int( counters, "new_count", "lots_inserted", "saved_inserted", "rows_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). #2625: сюда же сведён отказ называть успехом прогон, у которого отказом кончился каждый якорь и собрано ноль — см. _sweep_run_did_nothing. Проверка стоит здесь, а не в каждом sweep'е, ровно потому, что вызывающих у mark_done четыре десятка: страж, который надо не забыть позвать, — это тот же дефект оборванной проводки, из-за которого задача и появилась. #2700: там же — отказ называть успехом прогон, у которого отказала КАЖДАЯ попытка целой фазы (см. _phase_totally_failed). Отличие от #2625: тот случай про «не сделано ничего», этот — про «одно направление работы мертво, а суммарный сбор это прячет». honest-run-status: там же — отказ называть успехом прогон с высокой долей отказов, даже если собрано > 0 (см. _failed_ratio_too_high). Отличие от #2625/#2700: те два смотрят на «всё или ничего» (все якоря / вся фаза), этот — на ДОЛЮ отказов у detail-backfill'ов, где ни один из первых двух признаков не матчит форму counters. """ did_nothing = _sweep_run_did_nothing(counters) if did_nothing is not None: logger.error("%s run_id=%d", did_nothing, run_id) mark_failed(db, run_id, did_nothing, counters) return phase_dead = _phase_totally_failed(counters) if phase_dead is not None: logger.error("%s run_id=%d", phase_dead, run_id) mark_failed(db, run_id, phase_dead, counters) return ratio_bad = _failed_ratio_too_high(counters) if ratio_bad is not None: logger.error("%s run_id=%d", ratio_bad, run_id) mark_failed(db, run_id, ratio_bad, counters) return 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_UNKNOWN, ) -> None: """Финализация run: status='banned' + диагноз ban_kind (#2686, дефолт — #2764). 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 — упала НАША инфраструктура (браузерный сайдкар/прокси); - BAN_KIND_UNKNOWN (дефолт) — причина не установлена. Значение приходит от места ПОРОЖДЕНИЯ отказа (тип исключения), а не из разбора текста ошибки. Дефолт 'unknown', а НЕ 'platform' (#2764): вызывающий, которому разводить нечего, ничего и не знает — а не «знает, что виновата площадка». Оба исхода одинаково сохраняют чекпоинт — они отличаются только диагнозом. 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, ban_kinds: Collection[str] = (), ) -> 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). Логи контейнера на этот вопрос отвечать не могут — они исчезают при пересоздании контейнера, то есть на первом же деплое после ночного прогона. `ban_kinds` — диагнозы (ban_kind_of_exception) ВСЕХ блоков, которые задача поймала за прогон; пустой (дефолт) = задача типы не различает. Схлопываем сами, в одном месте на все три backfill'а: все блоки сошлись в одном диагнозе → он и пишется; разошлись (или их типы ничего не доказывают) → 'unknown'. Смешанный прогон честнее пометить неизвестным, чем выбрать из двух причин ту, что попалась последней — какая из них оборвала прогон, мы не знаем (#2764). """ 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) kinds = set(ban_kinds) mark_banned( db, run_id, reason, counters, ban_kind=kinds.pop() if len(kinds) == 1 else BAN_KIND_UNKNOWN, ) 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]