gendesign/tradein-mvp/backend/app/services/scrape_runs.py
bot-backend e9ca744e85
All checks were successful
CI / changes (pull_request) Successful in 9s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI Trade-In / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 5m4s
fix(tradein/scrapers): не путать rows_inserted/processed с честным результатным ключом
Ревью честного run-status нашло, что _RESULT_COUNTER_KEYS ловил не только целевой
yandex_newbuilding_sweep, но и rosreestr_dkp_import (rows_inserted, 66 из 67 прод-
прогонов = здоровый ноль догнавшего инкрементального импорта) и newbuilding_enrich
(processed — счётчик попыток, ==limit даже при частичном провале). Первое завело бы
практически непрерываемый ложный zero-стрик у здорового источника, второе маскировало
бы реальные отказы под measured-N.

Проверено по прод-БД (2026-08-15): "succeeded" пишут ТОЛЬКО yandex_newbuilding_sweep
(42 прогона/90д) и newbuilding_enrich (65/90д) — ни разу rosreestr_dkp_import; у
yandex_newbuilding_sweep succeeded численно совпадает с rows_inserted на всех 42/42
прогонах. Заменил "rows_inserted"+"processed" на "succeeded" в _RESULT_COUNTER_KEYS
(app-копия и byte-эквивалентная kit-копия) — цель (b) исходной правки сохранена, ложный
стрик у rosreestr_dkp_import снят, попутно newbuilding_enrich получает честное
измерение вместо счётчика попыток.

Также поправлены докстринги test_backfill_honest_status.py — два кейса (76%/72%
отказов -> 'done') проверяют только выбор финализатора mark_backfill_finished
(mark_done там замокан); реальный mark_done с honest-run-status переквалифицирует их
в 'failed' через _failed_ratio_too_high — это не документировалось явно.
2026-08-15 18:49:37 +03:00

1020 lines
62 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

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

"""Utility functions для scrape_runs table — tracking long-running pipeline runs.
Таблица scrape_runs создана в 015_scrape_runs.sql.
Расширена в 051_scrape_runs_extend.sql: params/counters/error/finished_at/cancelled.
ВРЕМЯ ПИШЕТСЯ 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`).
Признак — собственная бухгалтерия фазы: `<phase>_failed == <phase>_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 (нужна пара
# "<phase>_attempted"/"<phase>_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]