refactor(tradein): одна реализация scrape_runs — kit orchestration/runs.py, app-модуль = алиас; counters везде мержатся (#3390) #3400

Merged
bot-backend merged 2 commits from refactor/3390-single-runs-module into main 2026-09-06 08:11:53 +00:00
9 changed files with 523 additions and 1225 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 пульсом. Имя оставлено как есть: теперь это
# один объект, и переименование в 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 "

File diff suppressed because it is too large Load diff

View file

@ -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'ом.

View file

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

View file

@ -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), а не факт вызова.

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

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

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]