diff --git a/tradein-mvp/backend/app/services/scheduler.py b/tradein-mvp/backend/app/services/scheduler.py index 5c483c1a..0f67306d 100644 --- a/tradein-mvp/backend/app/services/scheduler.py +++ b/tradein-mvp/backend/app/services/scheduler.py @@ -28,10 +28,12 @@ from __future__ import annotations import logging from typing import Any -# kit_runs.update_heartbeat (в отличие от локального runs_mod.update_heartbeat) мержит -# counters (`counters || :counters`) вместо замены — нужен для чекпоинта курсора -# import_rosreestr_dkp (issue #3168), чтобы resume-вердикт, записанный на старте, не -# затирался последующими per-batch heartbeat'ами того же прогона. +# kit_runs — ТОТ ЖЕ модуль, что и runs_mod ниже: с #3390 `app.services.scrape_runs` +# его алиас, реализация одна и counters везде МЕРЖАТСЯ (`counters || :counters`). До +# #3390 копии было две, и app-копия counters ЗАМЕНЯЛА — тогда чекпоинт курсора +# import_rosreestr_dkp (#3168) обязан был писаться именно kit-именем, иначе resume-вердикт +# со старта затирался первым же per-batch пульсом. Имя оставлено как есть: теперь это +# один объект, и переименование в runs_mod ничего не чинит и ничего не ломает. from scraper_kit.orchestration import runs as kit_runs # compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674). @@ -129,12 +131,10 @@ async def _execute_cian_backfill( поведению (пометка 'zombie' на 6-м часу).""" nonlocal counters # #3384: снимок измеренного едет не только в БД, но и в `counters` — этот словарь - # уезжает в mark_failed из общего except ниже, а тот counters ЗАМЕНЯЕТ - # (scrape_runs.py:738 `counters = CAST(:counters AS jsonb)`; мерж `||` — только у - # kit-копии, которую этот путь не зовёт). Пока снимок сюда не доезжал, любой отказ - # ПОСЛЕ пройденной стадии (пул опустел между стадиями, упал SELECT домов) писал - # поверх измеренного предынициализированные нули, и SQL-разбор простоя - # (#3288/#3367) читал «к площадке не ходили» про прогон, который ходил. + # уезжает в mark_failed из общего except ниже. Мерж (#3390) спасает лишь ключи, + # которых в payload нет; одноимённые он ПЕРЕЗАПИСЫВАЕТ, поэтому предынициализированные + # нули без этого присваивания легли бы поверх измеренного, и SQL-разбор простоя + # (#3288/#3367) прочитал бы «к площадке не ходили» про прогон, который ходил. # Присваивание ДО записи в БД: сбой heartbeat'а не должен стирать сам факт замера. counters = _counters(progress) try: @@ -610,8 +610,9 @@ def import_rosreestr_dkp( counters["last_id"] = last_id # type: ignore[assignment] # Heartbeat = checkpoint: allows zombie detection + resume visibility. - # kit_runs (merge, не замена) — иначе этот REPLACE стёр бы resume_verdict, - # записанный _resume_dkp_cursor'ом перед циклом (issue #3168). + # Пульс МЕРЖИТ counters (`counters || :counters`), поэтому resume_verdict, + # записанный _resume_dkp_cursor'ом перед циклом, переживает per-batch запись + # (issue #3168; с #3390 мерж — единственная семантика, см. scrape_runs). kit_runs.update_heartbeat(db, run_id, counters) logger.info( "rosreestr_dkp_import run_id=%d: batch=%d fetched=%d " diff --git a/tradein-mvp/backend/app/services/scrape_runs.py b/tradein-mvp/backend/app/services/scrape_runs.py index 6e4b5450..949db937 100644 --- a/tradein-mvp/backend/app/services/scrape_runs.py +++ b/tradein-mvp/backend/app/services/scrape_runs.py @@ -1,1159 +1,36 @@ -"""Utility functions для scrape_runs table — tracking long-running pipeline runs. +"""`app.services.scrape_runs` — АЛИАС `scraper_kit.orchestration.runs` (#3390). -Таблица scrape_runs создана в 015_scrape_runs.sql. -Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled. +Реализация одна, и она в kit. Здесь нет ни функций, ни констант: модуль подменяет +себя kit-модулем в `sys.modules`, поэтому `from app.services import scrape_runs as +runs_mod` и `from scraper_kit.orchestration import runs` дают ОДИН И ТОТ ЖЕ объект. -ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL — -синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции и не двигается, -сколько бы та ни жила. Финализаторы (mark_done/mark_failed/mark_banned) выполняются -ТОЙ ЖЕ сессией, что и работа задачи, — и если рабочая транзакция всё это время -оставалась открытой (задача ничего не коммитила: нечего было сохранять, батч читающий, -сохранение шло чужой сессией), их UPDATE попадал ВНУТРЬ неё, и `finished_at` получал -время НАЧАЛА работы, а не её конца. +Зачем так, а не «тонкий модуль с реэкспортом имён»: реэкспорт разводит патч-цели. +`patch("app.services.scrape_runs.sentry_sdk")` или `patch.object(runs_mod, "mark_done")` +правит глобаль ЭТОГО модуля, а тело реэкспортированной функции читает глобаль СВОЕГО — +kit'а. Тест остался бы зелёным, не подменив ничего (или покраснел бы на пустом месте), +а прод-путь оказался бы вне проверки: ровно тот класс дефекта, из-за которого заведена +#3390. С алиасом патч по любому из двух имён попадает в единственную реализацию. -Замер на проде 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 мс. +Что здесь было до #3390: полная вторая копия модуля (1159 строк), которая разошлась с +kit'ом минимум четырежды — counters ЗАМЕНЯЛА (`counters = CAST(:counters AS jsonb)`) +против мержа (`COALESCE(counters,'{}') || …`) у kit, свой гейт по статусу, `honors_cancel` +только здесь, `mark_skipped` только там. Ревью дважды за сутки делало из этого ложные +выводы (#3388, #3355), потому что «вызывающий пишет через runs» не отвечало на вопрос, +какая из двух семантик применится. Победила семантика kit'а — мерж: у одной строки +прогона несколько писателей (пульс, финализатор, дрейн), каждый знает лишь свои ключи, +и замена теряла чужие (чекпоинт `done_buckets` #930, метка `interrupted` #3391, замер из +пульса #3384). Зависимости «финализатор обязан УДАЛИТЬ ключ заменой» нет ни у одного +вызывающего: строка прогона создаётся пустой (`create_run`), а резюм читает counters +ПРЕДЫДУЩЕГО прогона по его id. -Дефект был не сплошной ровно потому, что зависел от того, коммитила ли задача перед -финалом: 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 ч). +НЕ добавляй сюда код: всё, что написано ниже подмены, недостижимо — импортирующий +получает kit-модуль. Новые функции — в `scraper_kit/orchestration/runs.py`. """ from __future__ import annotations -import json -import logging -from collections import Counter -from collections.abc import Callable, Collection, Mapping -from functools import cache -from typing import Any +import sys -import sentry_sdk -from sqlalchemy import text -from sqlalchemy.orm import Session +from scraper_kit.orchestration import runs as _kit_runs -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 - - -def _was_interrupted(counters: Mapping[str, Any] | None) -> bool: - """Прогон оборван SIGTERM-дрейном (#3391/#3363): counters ЧАСТИЧНЫ по построению. - - До #3392 деплой оставлял такую строку в 'running' → 'zombie', и сторожей она не - касалась. Теперь дрейн финализирует её штатным `mark_done`/`mark_failed` с - `interrupted=1` — и в лестницах стриков (#3393) её надо считать так, как будто - прогона не было: судить по частичным счётчикам нельзя ни в одну сторону, поэтому - он стрик не продлевает и не обнуляет. Обнуление и есть задокументированный вред - (kit `orchestration/scheduler.py:99`: «5 банов подряд обнулил один "cancelled" - 09.08»). 'cancelled' оставлен как был: там прогон прервал человек, и разрыв - стрика — осознанное решение, а не побочный эффект деплоя. - """ - return bool((counters or {}).get("interrupted")) - - -# #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). - """ - # #3044: detail-бэкфиллы (avito/yandex/domclick_detail_backfill) пишут свой - # результат ТОЛЬКО ключом 'enriched' — колонка total_seen у них оставалась 0 - # навсегда, и замер в issue прочитал «обогащено 0 за 7 дней» при реальных 801 - # (listings.detail_enriched_at). 'enriched' — фолбэк ИМЕННО ЗДЕСЬ, а не в - # _RESULT_COUNTER_KEYS: тот список кормит ещё и zero-result-стрик, где - # «догнавший очередь» бэкфилл (attempted=0, enriched=0) стал бы измеренным - # нулём и копил бы непрерываемый стрик — ровно та ловушка, за которую ревью - # выкинуло из списка голый 'rows_inserted' (см. test_rosreestr_dkp_import_ - # healthy_zero_stays_unmeasured). Витринной колонке фолбэк безопасен: она не - # участвует в стриках (#2703 читает counters, не колонку). - total_seen = _run_result_count(counters) - if total_seen is None: - total_seen = _pick_int(counters, "enriched") - return total_seen, _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, 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() - - # #3393: оборванного деплоем прогона в популяции стрика нет. getattr — - # функция целиком под `except Exception: pass`, и строка без колонки - # выключила бы сторож молча (тот же класс, что #2703). - scanned = len(rows) - rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] - streak = _leading_streak(rows, lambda r: r.status in ("failed", "banned")) - # Потолок — по СКАНИРОВАННОМУ окну, а не по отфильтрованному: по второму - # одна interrupted-строка внутри 500 опускает измеримый максимум до 499, и - # `capped` становится недостижим по построению. Источник, сломанный наглухо, - # замолкал бы после 384-й неудачи навсегда (следующая веха 768 больше окна), - # а новую interrupted-строку в окно подкладывает каждый деплой — ровно - # анти-спам «один раз навсегда» из #2670/#2703. Условие читается так: окно - # было полным И стрик покрывает всё, что мы вообще могли судить. - capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows) - 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() - - # #3393: та же популяция, что у сторожа неудач — оборванный деплоем прогон - # не судим (в т.ч. не он решает, «измерен» ли свежайший результат ниже). - scanned = len(rows) - rows = [r for r in rows if not _was_interrupted(getattr(r, "counters", None))] - 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) - # Потолок — по сканированному окну, как у _alert_if_consecutive_failures: - # по отфильтрованному списку одна interrupted-строка делает `capped` - # недостижимым, и лестница за последней вехой в пределах окна молчит навсегда. - capped = scanned >= STREAK_SCAN_LIMIT and streak >= len(rows) - 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 — старое значение колонки сохраняется. - - #3391: гейт по статусу — тот же, что у mark_done/mark_failed. Пульс по УЖЕ - финализированной строке не просто холостой: эта копия counters ЗАМЕНЯЕТ (#3390), - поэтому задача, помеченная дрейном как `interrupted`, но ещё живая (ветка таймаута - drain_inflight отдаёт её на внешний hard-cancel — несколько итераций спустя), - следующим же пульсом стирала метку, и оборванный прогон снова читался как полный - проход. - - 'cancelled' в гейте НАМЕРЕННО, это не «running и всё финализированное»: user-cancel - финализирует строку, но задача останавливается лишь на ближайшей границе якоря, и её - последний пульс — ЕДИНСТВЕННЫЙ писатель чекпоинта в этот момент (pipeline.py:1308 / - 2368 / 2986 / 4488; mark_done там уже no-op по своему гейту, а 'cancelled' входит в - _RESUME_STATUSES). Сузить гейт до 'running' — молча потерять точку возобновления - у каждой отмены. - """ - total_seen, new_count = _column_counts(counters) - row = 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 AND status IN ('running', 'cancelled') - 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( - "update_heartbeat no-op: run_id=%d not in 'running'/'cancelled' state", run_id - ) - 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. - - #3393: все три гейта пропускаются, если прогон оборван дрейном (`interrupted=1`). - Их вход — счётчики ЗАВЕРШЁННОЙ работы; у оборванного они частичны, и диагноз про - площадку («сбор деградировал») был бы выдуман: бэкфилл, убитый деплоем на 20% - отказов, — это про деплой, а не про площадку. Статус остаётся штатным 'done' с - меткой interrupted (конвенция #3319/#3333/#3355), причина — в логе. - """ - if _was_interrupted(counters): - logger.warning( - "mark_done: run_id=%d оборван деплоем (SIGTERM-drain), частичный результат — " - "honest-status-гейты пропущены (#3393)", - run_id, - ) - honest_bad = None - else: - honest_bad = ( - _sweep_run_did_nothing(counters) - or _phase_totally_failed(counters) - or _failed_ratio_too_high(counters) - ) - if honest_bad is not None: - logger.error("%s run_id=%d", honest_bad, run_id) - mark_failed(db, run_id, honest_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 _dominant_ban_kind(census: Mapping[str, int]) -> str: - """Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178). - - Раньше вызывающий код терял кратности до вызова (`set()`), поэтому 4 блока - 'platform' + 1 'infra' и 2+2 давали функции один и тот же вход {'platform', - 'infra'} — неотличимые случаи, хотя первый явно платформенный, а второй - действительно спорный. Перепись приходит уже с кратностями (Counter), здесь — - только выбор: - - пусто → 'unknown' (диагнозов не было вовсе); - - один вид → он, независимо от количества; - - несколько видов, но один строго больше половины всех блоков → он - (доминирующий диагноз, единичные выбросы других типов его не размывают); - - иначе (нет строгого большинства) → 'unknown' — по какой причине оборвался - именно этот прогон, честно не знаем. - """ - if not census: - return BAN_KIND_UNKNOWN - if len(census) == 1: - return next(iter(census)) - total = sum(census.values()) - kind, count = max(census.items(), key=lambda kv: kv[1]) - if count > total / 2: - return kind - return BAN_KIND_UNKNOWN - - -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] | Mapping[str, int] = (), -) -> 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) ВСЕХ блоков, которые задача - поймала за прогон; пустой (дефолт) = задача типы не различает. Принимает либо - Collection[str] (старые вызовы — список/set диагнозов, кратности не несут) либо - уже готовую перепись Mapping[str, int] (kind -> сколько раз). Раньше здесь стоял - set(ban_kinds) — терял кратности ДО решения: 4 блока 'platform' + 1 'infra' - схлопывались в тот же вход {'platform', 'infra'}, что и настоящие 2+2, и оба - давали 'unknown' (#3178, прод: 5 прогонов подряд 4×platform+1×infra → unknown, - один прогон 5/5 одного вида → platform — при том же исключении на каждом блоке, - AvitoBlockedError firewall/soft-block). Перепись кладём в - counters["ban_kinds"] (kind -> count) — переживает финализацию наравне с - остальными counters, диагноз строки прогона выбирает _dominant_ban_kind: один - вид → он; явное большинство (строго > половины блоков) → он; иначе — '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): - # Counter() принимает и Collection (считает элементы — старые set/list-вызовы), - # и Mapping (копирует кратности как есть — census от вызывающего) одним и тем - # же конструктором. - census = Counter(ban_kinds) - if census: - counters["ban_kinds"] = dict(census) # type: ignore[assignment] - dominant = _dominant_ban_kind(census) - - if dominant == BAN_KIND_INFRA and produced > 0: - # #3288: 'banned' означает «площадка нас заблокировала» — и читается так - # же (триггер ротации IP, алерты, разбор простоя). Прогон 5425 при 41 - # infra из 48 «блоков» честно обогатил 41 карточку — его оборвал брейкер - # по доле, а не площадка, — и всё равно рапортовал «остановлен блоками - # источника». Понижаем ровно этот случай: диагноз infra И прогон работу - # сделал → 'done'. - # - # Нулевой прогон с infra остаётся 'banned' — контракт #2764/#3196: - # там диагноз несёт ban_kind строки ('infra'), а статус говорит «прогон - # оборван отказами». Понижать его до 'failed' по одному лишь диагнозу - # опаснее исходного дефекта: при пустом/отсутствующем census (источник - # видов не различает) dominant='unknown', а настоящий бан площадки, - # опознанный как infra по 5xx, спрятался бы под «нашей поломкой». - reason = ( - f"backfill-honest-status: {source} оборван брейкером на отказах НАШЕГО " - f"тракта — blocked={blocked} (диагноз '{BAN_KIND_INFRA}' у большинства), " - f"обогащено {enriched} из {attempted} попыток{hint} (#3288)" - ) - logger.error("%s run_id=%d", reason, run_id) - mark_done(db, run_id, counters) - return - - 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, - ban_kind=dominant, - ) - 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] +sys.modules[__name__] = _kit_runs diff --git a/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py index 25e7f499..eedd35ba 100644 --- a/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py +++ b/tradein-mvp/backend/tests/test_3168_backfill_cursor_resume.py @@ -23,10 +23,10 @@ import_rosreestr_dkp (source='rosreestr_dkp_import'), шестой backfill, н штатным поведением). - Потолок возраста чекпоинта — 24ч (_DKP_CHECKPOINT_STALE_HOURS): старше — курсор считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'. - - Вердикт пишется в scrape_runs.counters через kit_runs.update_heartbeat - (`counters || :counters` — merge), а не локальный runs_mod.update_heartbeat - (`CAST(:counters AS jsonb)` — полная замена): иначе первый же per-batch heartbeat - после старта стирает resume-вердикт. + - Вердикт пишется в scrape_runs.counters через update_heartbeat, а тот counters + МЕРЖИТ (`counters || :counters`) — иначе первый же per-batch heartbeat после старта + стёр бы resume-вердикт. На момент #3168 мерж был только у kit-копии, поэтому запись + шла именно kit-именем; с #3390 копия одна и семантика мержа — единственная. Обратимость (см. PR summary): временный откат last_id на литерал 0 красит test_resume_continues_from_saved_last_id (last_id == 123456 не совпадает с 0). @@ -34,7 +34,6 @@ test_resume_continues_from_saved_last_id (last_id == 123456 не совпада from __future__ import annotations -import inspect import json import os from typing import Any @@ -163,20 +162,15 @@ def test_checkpoint_just_under_ceiling_is_still_accepted() -> None: # ── 3. Запись курсора не затирает посторонние ключи в counters (мерж, не замена) ──── +# +# Текстового гейта на имя `kit_runs.update_heartbeat` здесь больше нет (#3390): пока копий +# было две, имя выбирало семантику, и проверять его по тексту имело смысл. Теперь +# `app.services.scrape_runs` — алиас kit'а, оба имени дают ОДИН объект, и гейт краснел бы +# на переименовании, ничего при этом не защищая. Мерж проверяется по значению — тестом +# ниже и test_3390_single_runs_module.py (оба пути импорта, heartbeat + финализатор). -def test_cursor_write_uses_merge_not_replace_heartbeat() -> None: - """import_rosreestr_dkp обязан писать чекпоинт через kit_runs.update_heartbeat - (merge: `counters || :counters`), а не локальный runs_mod.update_heartbeat (замена: - `CAST(:counters AS jsonb)`) — иначе resume-вердикт, записанный ДО цикла, стирается - первым же per-batch heartbeat'ом того же прогона. - """ - src = inspect.getsource(sched.import_rosreestr_dkp) - assert "kit_runs.update_heartbeat" in src - assert "runs_mod.update_heartbeat" not in src - - -def test_kit_runs_update_heartbeat_merges_into_existing_counters() -> None: +def test_update_heartbeat_merges_into_existing_counters() -> None: """Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая resume_from/resume_reason, записанные первым write'ом. diff --git a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py index a5aaa286..5bc7bbde 100644 --- a/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py +++ b/tradein-mvp/backend/tests/test_3355_drain_mark_full_loads.py @@ -247,10 +247,10 @@ async def test_domclick_sweep_drain_is_marked_interrupted() -> None: async def test_domclick_sweep_drain_keeps_inherited_checkpoint() -> None: """domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload. - Мержа jsonb тут нет: app-level `scrape_runs.update_heartbeat`/`mark_done` пишут - `counters = CAST(:counters AS jsonb)` — полная перезапись, а scheduler при claim - done_buckets не наследует. Без явного ключа дрейн-прогон закрывался бы с пустым - чекпоинтом, и резюм пересобирал бы все шесть корзин заново. + Мерж jsonb (#3390) тут не спасает: он сливает payload со СВОЕЙ строкой прогона, а та + создана пустой (`create_run`) — унаследованный чекпоинт лежит в counters ПРЕДЫДУЩЕГО + прогона, и scheduler при claim его не переносит. Без явного ключа дрейн-прогон + закрывался бы с пустым чекпоинтом, и резюм пересобирал бы все шесть корзин заново. """ from scraper_kit.orchestration import pipeline as pl diff --git a/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py b/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py index c0e36231..5d2b9ea3 100644 --- a/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py +++ b/tradein-mvp/backend/tests/test_3384_no_proxy_at_batch_start.py @@ -16,8 +16,9 @@ break`), которая и пишет `counters.no_proxy_stop`. Общий `exce `httpx.AsyncClient.post` — «к площадке не ходили» проверяется, а не предполагается. Третий тест — про соседний случай: пул опустел МЕЖДУ стадиями, то есть уже ПОСЛЕ -реальной работы. Там проверяется, что в запись прогона не уезжают нули поверх -измеренного (`mark_failed` у `app.services.scrape_runs` counters ЗАМЕНЯЕТ, а не мержит). +реальной работы. Там проверяется, что в запись прогона не уезжают нули поверх измеренного: +`mark_failed` counters МЕРЖИТ (#3390), но одноимённые ключи при мерже перезаписываются, +поэтому нули в payload'е финализатора всё так же стирают замер пульса. Сеть/БД замоканы; в сеть тест не ходит. """ @@ -199,10 +200,10 @@ async def test_cian_empty_pool_between_stages_keeps_measured_counters() -> None: У циана (в отличие от avito/домклика с их живым `counters.to_dict()`) в общий `except` приходит СТАРЫЙ словарь: реальные значения присваиваются уже после возврата из `backfill_cian_history`, а отказ бывает и посреди неё — пул опустел между - стадиями, упал SELECT домов. `runs_mod` здесь настоящий - (`app.services.scrape_runs`), и его `mark_failed` counters ЗАМЕНЯЕТ - (`counters = CAST(:counters AS jsonb)`, scrape_runs.py:738) — то есть в записи - прогона остаётся ровно то, что уехало последним аргументом. + стадиями, упал SELECT домов. `runs_mod` здесь настоящий (`app.services.scrape_runs`, + с #3390 — алиас kit'а), и его `mark_failed` counters МЕРЖИТ; мерж спасает только + ключи, которых в payload'е нет, а одноимённые ПЕРЕЗАПИСЫВАЕТ — нули поверх замера + всё равно недопустимы, и проверять их надо у вызывающего. Проверка по значению: смотрим jsonb-payload обоих UPDATE'ов (heartbeat и mark_failed), а не факт вызова. diff --git a/tradein-mvp/backend/tests/test_3390_single_runs_module.py b/tradein-mvp/backend/tests/test_3390_single_runs_module.py new file mode 100644 index 00000000..40a41742 --- /dev/null +++ b/tradein-mvp/backend/tests/test_3390_single_runs_module.py @@ -0,0 +1,205 @@ +"""#3390: у runs-модуля ОДНА реализация, и её семантика counters — мерж. + +До этой правки жили две копии одного модуля с ПРОТИВОПОЛОЖНОЙ семантикой: +`app.services.scrape_runs` counters ЗАМЕНЯЛ (`counters = CAST(:counters AS jsonb)`), +`scraper_kit.orchestration.runs` — МЕРЖИЛ (`COALESCE(counters,'{}') || …`). Разошлись +не только они: гейт `status`, `honors_cancel` у `mark_cancelled`, набор функций. +Ревью дважды за сутки делало из этого ложные выводы (#3388: «отдать только флаг, +остальное домержится» — на копии-заменителе это стёрло бы измеренное; #3355). + +Проверки здесь — ПО ЗНАЧЕНИЮ, через двойник сессии, который читает SQL: мерж (`||`) +против замены и WHERE-гейт по статусу берутся из текста самого statement'а, а не +зашиты ожиданием теста. На коде без гейта (WHERE только по id) апдейт проходит по +строке любого статуса — и тест краснеет по значению, а не по отсутствию подстроки. +""" + +from __future__ import annotations + +import json +import os +import re +from typing import Any + +os.environ.setdefault("DATABASE_URL", "postgresql+psycopg://test:test@localhost/test_db") + +import pytest +from scraper_kit.orchestration import runs as kit_runs + +from app.services import scrape_runs as app_runs + +# Оба пути импорта, которыми пользуется прод: kit-scheduler/pipeline ходят через kit, +# app-задачи (scheduler.py, avito/domclick_detail_backfill, admin API) — через app. +_RUNS_MODULES = {"app": app_runs, "kit": kit_runs} + + +class _FakeResult: + def __init__(self, row: Any = None) -> None: + self._row = row + + def first(self) -> Any: + return self._row + + def fetchone(self) -> Any: + return self._row + + def fetchall(self) -> list[Any]: + return [] + + +class _RunRowDb: + """Мини-Postgres на одну строку scrape_runs. + + UPDATE применяется, только если строка проходит WHERE из ТЕКСТА statement'а; + counters мержатся при `||` и заменяются при `CAST(:counters AS jsonb)` — тоже по + тексту. SELECT'ы alert-хуков (`_alert_on_run_id`) не моделируются: они best-effort + и возвращают пусто. + """ + + def __init__(self, *, status: str = "running", counters: dict[str, Any] | None = None) -> None: + self.row: dict[str, Any] = {"status": status, "counters": dict(counters or {})} + + @staticmethod + def _allowed_statuses(sql: str) -> set[str] | None: + """Статусы из WHERE. None — гейта нет, UPDATE бьёт по строке любого статуса.""" + where = sql.rsplit("WHERE", 1)[-1] + in_list = re.search(r"status\s+IN\s*\(([^)]*)\)", where) + if in_list is not None: + return {s.strip().strip("'") for s in in_list.group(1).split(",")} + eq = re.search(r"status\s*=\s*'(\w+)'", where) + return {eq.group(1)} if eq is not None else None + + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = " ".join(str(stmt).split()) + if not sql.startswith("UPDATE scrape_runs"): + return _FakeResult() + allowed = self._allowed_statuses(sql) + if allowed is not None and self.row["status"] not in allowed: + return _FakeResult(None) # WHERE не пропустил — 0 строк, RETURNING пуст + raw = (params or {}).get("counters") + if raw is not None: # mark_cancelled counters не пишет вовсе + payload = json.loads(raw) + self.row["counters"] = ( + {**self.row["counters"], **payload} if "||" in sql else dict(payload) + ) + new_status = re.search(r"SET status = '(\w+)'", sql) + if new_status is not None: + self.row["status"] = new_status.group(1) + return _FakeResult((1,)) + + def commit(self) -> None: + pass + + def rollback(self) -> None: + pass + + +# ── 1. Реализация одна ────────────────────────────────────────────────────── + + +@pytest.mark.parametrize( + "name", + [ + "create_run", + "update_heartbeat", + "is_cancelled", + "mark_done", + "mark_failed", + "mark_banned", + "mark_cancelled", + "mark_backfill_finished", + "mark_skipped", + "honors_cancel", + "list_recent", + "list_all", + "distinct_sources", + ], +) +def test_app_and_kit_expose_the_same_object(name: str) -> None: + """Публичное имя из app-копии — ТОТ ЖЕ объект, что и в kit (#3390). + + Не «эквивалентный текст», а идентичность: пока это две функции, любая правка + обязана попасть в обе, и следующее расхождение — вопрос времени (их было + минимум четыре: counters, гейт статуса, honors_cancel, состав функций). + """ + assert getattr(app_runs, name) is getattr(kit_runs, name), ( + f"{name}: app-копия и kit-копия — разные объекты, реализация снова раздвоена" + ) + + +# ── 2. Семантика counters: мерж, а не замена ──────────────────────────────── + + +@pytest.mark.parametrize("mod_name", list(_RUNS_MODULES)) +def test_finalizer_keeps_measured_counters_of_heartbeat(mod_name: str) -> None: + """heartbeat записал замер → mark_failed с ДРУГИМ ключом его не стирает. + + Прод-повод (#3384/#3388): задача бьёт пульс живыми счётчиками, а в общий `except` + приходит частичный/старый словарь. На копии-заменителе финализатор клал его ПОВЕРХ + всего, и измеренная работа исчезала из строки прогона. + """ + mod = _RUNS_MODULES[mod_name] + db = _RunRowDb() + + mod.update_heartbeat(db, 3390, {"lots_fetched": 5}) + mod.mark_failed(db, 3390, "boom", {"no_proxy_stop": 1}) + + assert db.row["counters"] == {"lots_fetched": 5, "no_proxy_stop": 1}, ( + f"{mod_name}: финализатор ЗАМЕНИЛ counters — замер heartbeat'а потерян" + ) + + +@pytest.mark.parametrize("mod_name", list(_RUNS_MODULES)) +def test_done_keeps_checkpoint_written_by_heartbeat(mod_name: str) -> None: + """Чекпоинт `done_buckets` от пульса переживает mark_done без этого ключа (#930).""" + mod = _RUNS_MODULES[mod_name] + db = _RunRowDb() + + mod.update_heartbeat(db, 3390, {"done_buckets": ["1:0:5"], "lots_fetched": 7}) + mod.mark_done(db, 3390, {"lots_fetched": 9}) + + assert db.row["counters"]["done_buckets"] == ["1:0:5"], ( + f"{mod_name}: точка возобновления стёрта финализатором" + ) + assert db.row["counters"]["lots_fetched"] == 9, "свежее значение обязано перекрывать старое" + + +# ── 3. Гейт по статусу — в единственной реализации ────────────────────────── + + +@pytest.mark.parametrize("mod_name", list(_RUNS_MODULES)) +@pytest.mark.parametrize("writer", ["update_heartbeat", "mark_done", "mark_failed"]) +def test_writers_do_not_touch_finalized_row(mod_name: str, writer: str) -> None: + """Пульс/финализатор по УЖЕ завершённой строке — no-op, а не затирание метки. + + Задача переживает собственную финализацию (дрейн пометил `interrupted`, а ветка + таймаута отдала её внешнему hard-cancel) и продолжает слать прогресс. + """ + mod = _RUNS_MODULES[mod_name] + db = _RunRowDb(status="done", counters={"interrupted": 1}) + + args: tuple[Any, ...] = ("boom", {"lots_fetched": 1}) if writer == "mark_failed" else ({},) + getattr(mod, writer)(db, 3390, *args) + + assert db.row["counters"] == {"interrupted": 1}, ( + f"{mod_name}.{writer}: запись прошла по финализированной строке" + ) + + +# ── 4. Отказ отменять то, что не опрашивает отмену — тоже в единственной ──── + + +@pytest.mark.parametrize("mod_name", list(_RUNS_MODULES)) +def test_mark_cancelled_refuses_source_that_ignores_cancel(mod_name: str) -> None: + """`honors_cancel`-гейт был только в app-копии: kit пометил бы 'cancelled' любой + прогон, а задача продолжила бы работать — второй свип на том же IP (инцидент + 2026-05-31, runs #26+#27).""" + mod = _RUNS_MODULES[mod_name] + + class _SourceDb(_RunRowDb): + def execute(self, stmt: Any, params: dict[str, Any] | None = None) -> _FakeResult: + sql = " ".join(str(stmt).split()) + if sql.startswith("SELECT source"): + return _FakeResult(type("R", (), {"source": "yandex_newbuilding_sweep"})()) + return super().execute(stmt, params) + + assert mod.mark_cancelled(_SourceDb(), 3390) is False diff --git a/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py b/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py index 631c3043..daf4d42f 100644 --- a/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py +++ b/tradein-mvp/backend/tests/test_3391_drain_marks_inflight_app_tasks.py @@ -15,15 +15,15 @@ cian_history_backfill 6173) остались в scrape_runs со статусо 3. hard-cancel из scheduler_main (CancelledError прямо в `asyncio.wait` дрейна) даёт тот же результат — это ровно прод-путь 07.09, ветка таймаута его не покрывает; 4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly»; - 5. пульс НЕ пишет по финализированной строке (обе копии `update_heartbeat`) — иначе - помеченная, но ещё живая задача стирает метку следующим же ударом; + 5. пульс НЕ пишет по финализированной строке (оба пути импорта `update_heartbeat`) — + иначе помеченная, но ещё живая задача стирает метку следующим же ударом; 6. отмена, прилетевшая в тело тика (не в дрейн), помечает in-flight так же; 7. отказ SQL на одном run_id не уносит остальные, а WARNING перечисляет ровно помеченных. -Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — так ведёт себя боевая app-копия -(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый -`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы. +Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — намеренно строже боевого (тот с +#3390 мержит везде): голый `{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) +покраснел бы. Дрейн обязан слать чекпоинт явно, а не полагаться на мерж в БД. """ from __future__ import annotations @@ -49,8 +49,8 @@ import app.scheduler_main as sm from app.core import shutdown as sd from app.services import scrape_runs as app_runs -# Обе копии runs-модуля: у app counters ЗАМЕНЯЮТСЯ, у kit мержатся (#3390) — гейт по -# статусу нужен обеим, и проверяется на обеих одним и тем же телом теста. +# Оба пути импорта runs-модуля (с #3390 это ОДИН объект: app.services.scrape_runs — +# алиас kit'а): гейт по статусу проверяется через каждый из них одним телом теста. _RUNS_MODULES = {"kit": kit_runs, "app": app_runs} @@ -131,7 +131,7 @@ class _FakeRuns: if row is None or row["status"] != "running": return # боевой UPDATE ... WHERE status = 'running' — no-op row["status"] = "done" - row["counters"] = dict(counters) # app-копия ЗАМЕНЯЕТ counters (#3390) + row["counters"] = dict(counters) # строже боевого мержа (#3390) — см. докстринг модуля def _make_sched(source: str) -> dict[str, Any]: @@ -283,7 +283,7 @@ async def test_hard_cancel_does_not_log_drained_cleanly( assert "drained and exited cleanly" not in caplog.text -# ── 5. Пульс не пишет по финализированной строке (обе копии update_heartbeat) ──── +# ── 5. Пульс не пишет по финализированной строке (оба пути импорта) ───────────── class _RunRowDb: @@ -293,8 +293,8 @@ class _RunRowDb: Допустимые статусы вычитываются из SQL (`status = 'x'` / `status IN ('x', 'y')`), а не зашиты ожиданием теста: на коде без гейта (WHERE только по id) апдейт проходит по строке ЛЮБОГО статуса — и тест краснеет по значению, а не по отсутствию подстроки. - Мерж jsonb (`||`) против замены (`CAST(:counters AS jsonb)`) — тоже по тексту: у - kit- и app-копии он разный, а гейт нужен обеим. + Мерж jsonb (`||`) против замены — тоже по тексту statement'а, а не по ожиданию: + вернётся замена (было до #3390) — двойник это отразит, и тест покраснеет по значению. SELECT'ы alert-хуков (`_alert_on_run_id`) не моделируются — они best-effort и возвращают пусто. @@ -348,7 +348,8 @@ def test_heartbeat_does_not_erase_drain_mark_of_finalized_run(name: str) -> None Ровно эта последовательность достижима на проде: ветка таймаута `drain_inflight` оставляет задачу внешнему hard-cancel'у, и до него она успевает несколько итераций (`app/services/scheduler.py:141` — пульс на каждый батч). Без гейта по статусу - app-копия ЗАМЕНЯЛА counters и стирала метку: оборванный прогон снова читался как + пульс лёг бы поверх метки (одноимённые ключи мерж перезаписывает, а до #3390 + app-копия и вовсе ЗАМЕНЯЛА весь словарь): оборванный прогон снова читался бы как полный проход, а резюм его не подхватывал. """ mod = _RUNS_MODULES[name] diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py index 2f7e204c..59bafba9 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/pipeline.py @@ -4529,14 +4529,14 @@ async def run_domclick_city_sweep( # #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции — # это единообразие и защита от смены писателя, а НЕ восстановление потери. - # Как есть сегодня: пишет kit-овый scraper_kit.orchestration.runs (импорт выше), - # и все четыре его писателя МЕРЖАТ jsonb — `counters = COALESCE(counters,'{}') - # || CAST(:counters AS jsonb)`, так что ключ, записанный раньше, переживает - # payload без него; плюс _pick_resume наследует done_buckets при claim (#3074). - # То есть чекпоинт в БД не терялся. Перезаписывающий двойник существует — - # app/services/scrape_runs.py (`counters = CAST(:counters AS jsonb)`), — но - # этой функцией не вызывается. Полный payload делает ветки нечувствительными - # к тому, какой из двух писателей окажется на другом конце. + # Как есть сегодня: пишет scraper_kit.orchestration.runs (импорт выше) — с #3390 + # единственная реализация, `app/services/scrape_runs.py` её алиас, — и все его + # писатели МЕРЖАТ jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS + # jsonb)`), так что ключ, записанный раньше, переживает payload без него; плюс + # _pick_resume наследует done_buckets при claim (#3074). То есть чекпоинт в БД не + # терялся. Перезаписывающий двойник (`counters = CAST(:counters AS jsonb)`) жил в + # app-копии до #3390; полный payload оставлен и делает ветки нечувствительными к + # тому, какой писатель окажется на другом конце. _checkpoint: list[str] = [] def _payload() -> dict[str, Any]: diff --git a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py index 67915221..e2d8d009 100644 --- a/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py +++ b/tradein-mvp/packages/scraper-kit/src/scraper_kit/orchestration/runs.py @@ -1,14 +1,23 @@ -"""scrape_runs helpers — tracking long-running pipeline runs (strangler-копия #2135). +"""scrape_runs helpers — tracking long-running pipeline runs. -Байт-эквивалент `app.services.scrape_runs` — чистые SQL-хелперы поверх таблицы -`scrape_runs` (миграции 015 + 051). Развязка от `app.*`. Намеренные отличия от -app-копии: - 1. `sentry_sdk` импортируется опционально (kit standalone-импортируем, а - `sentry-sdk` не входит в его зависимости). Если пакет не установлен — alert-хук - best-effort no-op, поведение SQL-финализаторов идентично старому. - 2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно - планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim), - app-копии эта функция не нужна. +ЕДИНСТВЕННАЯ реализация (#3390). `app.services.scrape_runs` — тонкий алиас этого +модуля, а не вторая копия: до #3390 копий было две, и «байт-эквивалентными» они не +остались. Разошлись, в частности, семантика counters (здесь МЕРЖ `counters || :counters`, +там была ЗАМЕНА `counters = CAST(:counters AS jsonb)`), гейт по статусу, `honors_cancel` +у `mark_cancelled` и сам состав функций. Расхождение стоило дважды за сутки ложных +выводов на ревью (#3388: «отдать только флаг, остальное домержится» — на копии-заменителе +это стёрло бы измеренное; #3355). + +Мерж, а не замена, потому что писателей у одной строки прогона несколько (пульс задачи, +её финализатор, дрейн), и каждый знает лишь СВОИ ключи: чекпоинт `done_buckets` писали +только знающие о нём сайты (#930), метку `interrupted` ставит дрейн (#3391), а замер +уезжает в пульсе. Замена делала запись прогона равной последнему payload'у — то есть +теряла всё, чего в нём случайно не оказалось. Обратной зависимости — «вызывающий +рассчитывает, что финализатор УДАЛИТ ключ заменой» — нет ни одной: run-строка создаётся +пустой (`create_run`), резюм читает counters ПРЕДЫДУЩЕГО прогона по его id, а не свои. + +`sentry_sdk` импортируется опционально: kit standalone-импортируем, а `sentry-sdk` не +входит в его зависимости. Если пакет не установлен — alert-хук best-effort no-op. ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL — синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы @@ -17,15 +26,20 @@ app-копии: `finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в диапазоне 9-64 мс (= от коммита claim'а до первого запроса рабочей транзакции). -Полный разбор — в docstring app-копии `app/services/scrape_runs.py`. +Дефект был не сплошной: он зависел от того, коммитила ли задача перед финалом +(cadastral_geo_match / house_imv_backfill / avito_detail_backfill коммитят поштучно, у +них окно совпадало с работой; yandex_address_backfill / newbuilding_enrich / +cian_history_backfill — нет). Побочно чинится и `heartbeat_at`, на котором стоит поиск +зависших прогонов (reap_zombies, порог 6 ч). """ from __future__ import annotations import json import logging +from collections import Counter +from collections.abc import Callable, Collection, Mapping from functools import lru_cache -from collections.abc import Callable, Mapping from typing import Any from sqlalchemy import text @@ -338,10 +352,8 @@ def _phase_totally_failed(counters: Mapping[str, Any]) -> str | None: # та же 'failed', но с другой формулировкой причины ("деградировал", не "провалился"), чтобы # оператор видел разницу читая error, не только status. # -# mark_backfill_finished (единственный писатель "attempted"/"failed" на верхнем уровне -# counters) живёт только в app.services.scrape_runs — здесь эта проверка сейчас неактивна -# ни для одного реального вызывающего, но kit-копия держится байт-эквивалентной app-копии -# (см. docstring модуля), и будущий kit-native job с тем же словарём получит её даром. +# Единственный писатель "attempted"/"failed" на верхнем уровне counters — +# mark_backfill_finished (ниже в этом же модуле); его зовут четыре detail-backfill'а. FAILED_RATIO_FAILED_THRESHOLD = 0.5 FAILED_RATIO_DEGRADED_THRESHOLD = 0.15 # Минимум попыток, при котором доля вообще что-то значит — иначе 1 отказ из 2 (=0.5) @@ -771,6 +783,27 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None 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 +# (app 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( @@ -959,8 +992,175 @@ def mark_banned( _alert_on_run_id(db, run_id) +def _dominant_ban_kind(census: Mapping[str, int]) -> str: + """Диагноз по переписи блоков прогона: kind -> сколько раз он встретился (#3178). + + Раньше вызывающий код терял кратности до вызова (`set()`), поэтому 4 блока + 'platform' + 1 'infra' и 2+2 давали функции один и тот же вход {'platform', + 'infra'} — неотличимые случаи, хотя первый явно платформенный, а второй + действительно спорный. Перепись приходит уже с кратностями (Counter), здесь — + только выбор: + - пусто → 'unknown' (диагнозов не было вовсе); + - один вид → он, независимо от количества; + - несколько видов, но один строго больше половины всех блоков → он + (доминирующий диагноз, единичные выбросы других типов его не размывают); + - иначе (нет строгого большинства) → 'unknown' — по какой причине оборвался + именно этот прогон, честно не знаем. + """ + if not census: + return BAN_KIND_UNKNOWN + if len(census) == 1: + return next(iter(census)) + total = sum(census.values()) + kind, count = max(census.items(), key=lambda kv: kv[1]) + if count > total / 2: + return kind + return BAN_KIND_UNKNOWN + + +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] | Mapping[str, int] = (), +) -> 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) ВСЕХ блоков, которые задача + поймала за прогон; пустой (дефолт) = задача типы не различает. Принимает либо + Collection[str] (старые вызовы — список/set диагнозов, кратности не несут) либо + уже готовую перепись Mapping[str, int] (kind -> сколько раз). Раньше здесь стоял + set(ban_kinds) — терял кратности ДО решения: 4 блока 'platform' + 1 'infra' + схлопывались в тот же вход {'platform', 'infra'}, что и настоящие 2+2, и оба + давали 'unknown' (#3178, прод: 5 прогонов подряд 4×platform+1×infra → unknown, + один прогон 5/5 одного вида → platform — при том же исключении на каждом блоке, + AvitoBlockedError firewall/soft-block). Перепись кладём в + counters["ban_kinds"] (kind -> count) — переживает финализацию наравне с + остальными counters, диагноз строки прогона выбирает _dominant_ban_kind: один + вид → он; явное большинство (строго > половины блоков) → он; иначе — '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): + # Counter() принимает и Collection (считает элементы — старые set/list-вызовы), + # и Mapping (копирует кратности как есть — census от вызывающего) одним и тем + # же конструктором. + census = Counter(ban_kinds) + if census: + counters["ban_kinds"] = dict(census) # type: ignore[assignment] + dominant = _dominant_ban_kind(census) + + if dominant == BAN_KIND_INFRA and produced > 0: + # #3288: 'banned' означает «площадка нас заблокировала» — и читается так + # же (триггер ротации IP, алерты, разбор простоя). Прогон 5425 при 41 + # infra из 48 «блоков» честно обогатил 41 карточку — его оборвал брейкер + # по доле, а не площадка, — и всё равно рапортовал «остановлен блоками + # источника». Понижаем ровно этот случай: диагноз infra И прогон работу + # сделал → 'done'. + # + # Нулевой прогон с infra остаётся 'banned' — контракт #2764/#3196: + # там диагноз несёт ban_kind строки ('infra'), а статус говорит «прогон + # оборван отказами». Понижать его до 'failed' по одному лишь диагнозу + # опаснее исходного дефекта: при пустом/отсутствующем census (источник + # видов не различает) dominant='unknown', а настоящий бан площадки, + # опознанный как infra по 5xx, спрятался бы под «нашей поломкой». + reason = ( + f"backfill-honest-status: {source} оборван брейкером на отказах НАШЕГО " + f"тракта — blocked={blocked} (диагноз '{BAN_KIND_INFRA}' у большинства), " + f"обогащено {enriched} из {attempted} попыток{hint} (#3288)" + ) + logger.error("%s run_id=%d", reason, run_id) + mark_done(db, run_id, counters) + return + + 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, + ban_kind=dominant, + ) + 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.""" + """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( """ @@ -1048,3 +1248,22 @@ def list_all( .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]