refactor(tradein/runs): одна реализация scrape_runs — kit, семантика counters мерж (#3390)
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m58s
CI / 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

Две живые копии одного модуля с противоположной семантикой counters: app
`mark_done`/`mark_failed`/`mark_banned`/`update_heartbeat` ЗАМЕНЯЛИ
(`counters = CAST(:counters AS jsonb)`), kit — МЕРЖИЛИ
(`COALESCE(counters,'{}') || …`). Расхождение дважды за сутки дало ложные
выводы на ревью (#3388 «отдать только флаг, остальное домержится» — на
replace это стёрло бы измеренное; #3355). Разошлись и другие места: гейт
статуса, `honors_cancel` у mark_cancelled (был только в app), `mark_skipped`
(только в kit), `mark_backfill_finished`/`distinct_sources` (только в app).

Реализация теперь одна — `scraper_kit.orchestration.runs`; в неё перенесены
app-only функции. `app.services.scrape_runs` — алиас kit-модуля через
sys.modules, а не реэкспорт имён: реэкспорт разводит патч-цели
(`patch("app.services.scrape_runs.sentry_sdk")`, `patch.object(runs_mod,
"mark_done")` правили бы глобаль модуля-обёртки, а тело функции читает
глобаль kit'а) — тест остался бы зелёным, не подменив ничего. С алиасом оба
имени ведут в единственную реализацию, и ни один из ~40 вызывающих и ~30
патч-сайтов в тестах не правится.

Победила семантика мержа: у строки прогона несколько писателей (пульс,
финализатор, дрейн), каждый знает лишь свои ключи, и замена теряла чужие —
чекпоинт done_buckets (#930), метку interrupted (#3391), замер из пульса
(#3384). Обратной зависимости («вызывающий рассчитывает, что финализатор
УДАЛИТ ключ заменой») нет: строка создаётся пустой в create_run, резюм читает
counters ПРЕДЫДУЩЕГО прогона по его id.

Тесты по значению на обоих путях импорта (двойник сессии читает SQL: `||`
против CAST, WHERE-гейт из текста): пульс {a:5} + mark_failed {b:1} → {a,b};
пульс/финализатор по финализированной строке — no-op; mark_cancelled
отказывает источнику, который отмену не опрашивает. На main эти тесты
красные для app-пути.

Комментарии в app/services/scheduler.py и kit/pipeline.py, утверждавшие про
живого «перезаписывающего двойника», приведены в соответствие.
This commit is contained in:
bot-backend 2026-09-06 11:45:56 +05:00
parent 075cec4b57
commit 36f2429fbe
5 changed files with 484 additions and 1180 deletions

View file

@ -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 пульсом. Имя оставлено: на него смотрит
# текстовый гейт test_3168 (test_cursor_write_uses_merge_not_replace_heartbeat).
from scraper_kit.orchestration import runs as kit_runs
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
@ -610,8 +612,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 "

File diff suppressed because it is too large Load diff

View file

@ -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

View file

@ -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]:

View file

@ -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]