refactor(tradein): одна реализация scrape_runs — kit orchestration/runs.py, app-модуль = алиас; counters везде мержатся (#3390) #3400
9 changed files with 523 additions and 1225 deletions
|
|
@ -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
|
|
@ -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'ом.
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
||||
|
|
|
|||
|
|
@ -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), а не факт вызова.
|
||||
|
|
|
|||
205
tradein-mvp/backend/tests/test_3390_single_runs_module.py
Normal file
205
tradein-mvp/backend/tests/test_3390_single_runs_module.py
Normal 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
|
||||
|
|
@ -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]
|
||||
|
|
|
|||
|
|
@ -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]:
|
||||
|
|
|
|||
|
|
@ -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]
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue