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
|
import logging
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
# kit_runs.update_heartbeat (в отличие от локального runs_mod.update_heartbeat) мержит
|
# kit_runs — ТОТ ЖЕ модуль, что и runs_mod ниже: с #3390 `app.services.scrape_runs`
|
||||||
# counters (`counters || :counters`) вместо замены — нужен для чекпоинта курсора
|
# его алиас, реализация одна и counters везде МЕРЖАТСЯ (`counters || :counters`). До
|
||||||
# import_rosreestr_dkp (issue #3168), чтобы resume-вердикт, записанный на старте, не
|
# #3390 копии было две, и app-копия counters ЗАМЕНЯЛА — тогда чекпоинт курсора
|
||||||
# затирался последующими per-batch heartbeat'ами того же прогона.
|
# import_rosreestr_dkp (#3168) обязан был писаться именно kit-именем, иначе resume-вердикт
|
||||||
|
# со старта затирался первым же per-batch пульсом. Имя оставлено как есть: теперь это
|
||||||
|
# один объект, и переименование в runs_mod ничего не чинит и ничего не ломает.
|
||||||
from scraper_kit.orchestration import runs as kit_runs
|
from scraper_kit.orchestration import runs as kit_runs
|
||||||
|
|
||||||
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
|
# compute_next_run_at жил здесь ВТОРОЙ, побайтово одинаковой копией kit-версии (#2674).
|
||||||
|
|
@ -129,12 +131,10 @@ async def _execute_cian_backfill(
|
||||||
поведению (пометка 'zombie' на 6-м часу)."""
|
поведению (пометка 'zombie' на 6-м часу)."""
|
||||||
nonlocal counters
|
nonlocal counters
|
||||||
# #3384: снимок измеренного едет не только в БД, но и в `counters` — этот словарь
|
# #3384: снимок измеренного едет не только в БД, но и в `counters` — этот словарь
|
||||||
# уезжает в mark_failed из общего except ниже, а тот counters ЗАМЕНЯЕТ
|
# уезжает в mark_failed из общего except ниже. Мерж (#3390) спасает лишь ключи,
|
||||||
# (scrape_runs.py:738 `counters = CAST(:counters AS jsonb)`; мерж `||` — только у
|
# которых в payload нет; одноимённые он ПЕРЕЗАПИСЫВАЕТ, поэтому предынициализированные
|
||||||
# kit-копии, которую этот путь не зовёт). Пока снимок сюда не доезжал, любой отказ
|
# нули без этого присваивания легли бы поверх измеренного, и SQL-разбор простоя
|
||||||
# ПОСЛЕ пройденной стадии (пул опустел между стадиями, упал SELECT домов) писал
|
# (#3288/#3367) прочитал бы «к площадке не ходили» про прогон, который ходил.
|
||||||
# поверх измеренного предынициализированные нули, и SQL-разбор простоя
|
|
||||||
# (#3288/#3367) читал «к площадке не ходили» про прогон, который ходил.
|
|
||||||
# Присваивание ДО записи в БД: сбой heartbeat'а не должен стирать сам факт замера.
|
# Присваивание ДО записи в БД: сбой heartbeat'а не должен стирать сам факт замера.
|
||||||
counters = _counters(progress)
|
counters = _counters(progress)
|
||||||
try:
|
try:
|
||||||
|
|
@ -610,8 +610,9 @@ def import_rosreestr_dkp(
|
||||||
counters["last_id"] = last_id # type: ignore[assignment]
|
counters["last_id"] = last_id # type: ignore[assignment]
|
||||||
|
|
||||||
# Heartbeat = checkpoint: allows zombie detection + resume visibility.
|
# Heartbeat = checkpoint: allows zombie detection + resume visibility.
|
||||||
# kit_runs (merge, не замена) — иначе этот REPLACE стёр бы resume_verdict,
|
# Пульс МЕРЖИТ counters (`counters || :counters`), поэтому resume_verdict,
|
||||||
# записанный _resume_dkp_cursor'ом перед циклом (issue #3168).
|
# записанный _resume_dkp_cursor'ом перед циклом, переживает per-batch запись
|
||||||
|
# (issue #3168; с #3390 мерж — единственная семантика, см. scrape_runs).
|
||||||
kit_runs.update_heartbeat(db, run_id, counters)
|
kit_runs.update_heartbeat(db, run_id, counters)
|
||||||
logger.info(
|
logger.info(
|
||||||
"rosreestr_dkp_import run_id=%d: batch=%d fetched=%d "
|
"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): старше — курсор
|
- Потолок возраста чекпоинта — 24ч (_DKP_CHECKPOINT_STALE_HOURS): старше — курсор
|
||||||
считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'.
|
считается негодным, прогон стартует с last_id=0, причина 'checkpoint_stale'.
|
||||||
- Вердикт пишется в scrape_runs.counters через kit_runs.update_heartbeat
|
- Вердикт пишется в scrape_runs.counters через update_heartbeat, а тот counters
|
||||||
(`counters || :counters` — merge), а не локальный runs_mod.update_heartbeat
|
МЕРЖИТ (`counters || :counters`) — иначе первый же per-batch heartbeat после старта
|
||||||
(`CAST(:counters AS jsonb)` — полная замена): иначе первый же per-batch heartbeat
|
стёр бы resume-вердикт. На момент #3168 мерж был только у kit-копии, поэтому запись
|
||||||
после старта стирает resume-вердикт.
|
шла именно kit-именем; с #3390 копия одна и семантика мержа — единственная.
|
||||||
|
|
||||||
Обратимость (см. PR summary): временный откат last_id на литерал 0 красит
|
Обратимость (см. PR summary): временный откат last_id на литерал 0 красит
|
||||||
test_resume_continues_from_saved_last_id (last_id == 123456 не совпадает с 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
|
from __future__ import annotations
|
||||||
|
|
||||||
import inspect
|
|
||||||
import json
|
import json
|
||||||
import os
|
import os
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
@ -163,20 +162,15 @@ def test_checkpoint_just_under_ceiling_is_still_accepted() -> None:
|
||||||
|
|
||||||
|
|
||||||
# ── 3. Запись курсора не затирает посторонние ключи в counters (мерж, не замена) ────
|
# ── 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:
|
def test_update_heartbeat_merges_into_existing_counters() -> 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:
|
|
||||||
"""Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая
|
"""Сама примитива слияния: второй write добавляет rows_fetched/last_id, НЕ стирая
|
||||||
resume_from/resume_reason, записанные первым write'ом.
|
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:
|
async def test_domclick_sweep_drain_keeps_inherited_checkpoint() -> None:
|
||||||
"""domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload.
|
"""domclick: дрейн ОБЯЗАН переписать унаследованный чекпоинт в свой payload.
|
||||||
|
|
||||||
Мержа jsonb тут нет: app-level `scrape_runs.update_heartbeat`/`mark_done` пишут
|
Мерж jsonb (#3390) тут не спасает: он сливает payload со СВОЕЙ строкой прогона, а та
|
||||||
`counters = CAST(:counters AS jsonb)` — полная перезапись, а scheduler при claim
|
создана пустой (`create_run`) — унаследованный чекпоинт лежит в counters ПРЕДЫДУЩЕГО
|
||||||
done_buckets не наследует. Без явного ключа дрейн-прогон закрывался бы с пустым
|
прогона, и scheduler при claim его не переносит. Без явного ключа дрейн-прогон
|
||||||
чекпоинтом, и резюм пересобирал бы все шесть корзин заново.
|
закрывался бы с пустым чекпоинтом, и резюм пересобирал бы все шесть корзин заново.
|
||||||
"""
|
"""
|
||||||
from scraper_kit.orchestration import pipeline as pl
|
from scraper_kit.orchestration import pipeline as pl
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -16,8 +16,9 @@ break`), которая и пишет `counters.no_proxy_stop`. Общий `exce
|
||||||
`httpx.AsyncClient.post` — «к площадке не ходили» проверяется, а не предполагается.
|
`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()`) в общий
|
У циана (в отличие от avito/домклика с их живым `counters.to_dict()`) в общий
|
||||||
`except` приходит СТАРЫЙ словарь: реальные значения присваиваются уже после возврата
|
`except` приходит СТАРЫЙ словарь: реальные значения присваиваются уже после возврата
|
||||||
из `backfill_cian_history`, а отказ бывает и посреди неё — пул опустел между
|
из `backfill_cian_history`, а отказ бывает и посреди неё — пул опустел между
|
||||||
стадиями, упал SELECT домов. `runs_mod` здесь настоящий
|
стадиями, упал SELECT домов. `runs_mod` здесь настоящий (`app.services.scrape_runs`,
|
||||||
(`app.services.scrape_runs`), и его `mark_failed` counters ЗАМЕНЯЕТ
|
с #3390 — алиас kit'а), и его `mark_failed` counters МЕРЖИТ; мерж спасает только
|
||||||
(`counters = CAST(:counters AS jsonb)`, scrape_runs.py:738) — то есть в записи
|
ключи, которых в payload'е нет, а одноимённые ПЕРЕЗАПИСЫВАЕТ — нули поверх замера
|
||||||
прогона остаётся ровно то, что уехало последним аргументом.
|
всё равно недопустимы, и проверять их надо у вызывающего.
|
||||||
|
|
||||||
Проверка по значению: смотрим jsonb-payload обоих UPDATE'ов (heartbeat и
|
Проверка по значению: смотрим jsonb-payload обоих UPDATE'ов (heartbeat и
|
||||||
mark_failed), а не факт вызова.
|
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` дрейна) даёт тот
|
3. hard-cancel из scheduler_main (CancelledError прямо в `asyncio.wait` дрейна) даёт тот
|
||||||
же результат — это ровно прод-путь 07.09, ветка таймаута его не покрывает;
|
же результат — это ровно прод-путь 07.09, ветка таймаута его не покрывает;
|
||||||
4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly»;
|
4. `scheduler_main` после hard-cancel'а больше не печатает «drained cleanly»;
|
||||||
5. пульс НЕ пишет по финализированной строке (обе копии `update_heartbeat`) — иначе
|
5. пульс НЕ пишет по финализированной строке (оба пути импорта `update_heartbeat`) —
|
||||||
помеченная, но ещё живая задача стирает метку следующим же ударом;
|
иначе помеченная, но ещё живая задача стирает метку следующим же ударом;
|
||||||
6. отмена, прилетевшая в тело тика (не в дрейн), помечает in-flight так же;
|
6. отмена, прилетевшая в тело тика (не в дрейн), помечает in-flight так же;
|
||||||
7. отказ SQL на одном run_id не уносит остальные, а WARNING перечисляет ровно
|
7. отказ SQL на одном run_id не уносит остальные, а WARNING перечисляет ровно
|
||||||
помеченных.
|
помеченных.
|
||||||
|
|
||||||
Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — так ведёт себя боевая app-копия
|
Фейковый `mark_done` counters ЗАМЕНЯЕТ, а не мержит — намеренно строже боевого (тот с
|
||||||
(`app.services.scrape_runs`, #3390), которая и инжектируется в прод-контекст. Голый
|
#3390 мержит везде): голый `{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1)
|
||||||
`{"interrupted": 1}` на этом фейке стёр бы чекпоинт, и (1) покраснел бы.
|
покраснел бы. Дрейн обязан слать чекпоинт явно, а не полагаться на мерж в БД.
|
||||||
"""
|
"""
|
||||||
|
|
||||||
from __future__ import annotations
|
from __future__ import annotations
|
||||||
|
|
@ -49,8 +49,8 @@ import app.scheduler_main as sm
|
||||||
from app.core import shutdown as sd
|
from app.core import shutdown as sd
|
||||||
from app.services import scrape_runs as app_runs
|
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}
|
_RUNS_MODULES = {"kit": kit_runs, "app": app_runs}
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -131,7 +131,7 @@ class _FakeRuns:
|
||||||
if row is None or row["status"] != "running":
|
if row is None or row["status"] != "running":
|
||||||
return # боевой UPDATE ... WHERE status = 'running' — no-op
|
return # боевой UPDATE ... WHERE status = 'running' — no-op
|
||||||
row["status"] = "done"
|
row["status"] = "done"
|
||||||
row["counters"] = dict(counters) # app-копия ЗАМЕНЯЕТ counters (#3390)
|
row["counters"] = dict(counters) # строже боевого мержа (#3390) — см. докстринг модуля
|
||||||
|
|
||||||
|
|
||||||
def _make_sched(source: str) -> dict[str, Any]:
|
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
|
assert "drained and exited cleanly" not in caplog.text
|
||||||
|
|
||||||
|
|
||||||
# ── 5. Пульс не пишет по финализированной строке (обе копии update_heartbeat) ────
|
# ── 5. Пульс не пишет по финализированной строке (оба пути импорта) ─────────────
|
||||||
|
|
||||||
|
|
||||||
class _RunRowDb:
|
class _RunRowDb:
|
||||||
|
|
@ -293,8 +293,8 @@ class _RunRowDb:
|
||||||
Допустимые статусы вычитываются из SQL (`status = 'x'` / `status IN ('x', 'y')`), а не
|
Допустимые статусы вычитываются из SQL (`status = 'x'` / `status IN ('x', 'y')`), а не
|
||||||
зашиты ожиданием теста: на коде без гейта (WHERE только по id) апдейт проходит по
|
зашиты ожиданием теста: на коде без гейта (WHERE только по id) апдейт проходит по
|
||||||
строке ЛЮБОГО статуса — и тест краснеет по значению, а не по отсутствию подстроки.
|
строке ЛЮБОГО статуса — и тест краснеет по значению, а не по отсутствию подстроки.
|
||||||
Мерж jsonb (`||`) против замены (`CAST(:counters AS jsonb)`) — тоже по тексту: у
|
Мерж jsonb (`||`) против замены — тоже по тексту statement'а, а не по ожиданию:
|
||||||
kit- и app-копии он разный, а гейт нужен обеим.
|
вернётся замена (было до #3390) — двойник это отразит, и тест покраснеет по значению.
|
||||||
|
|
||||||
SELECT'ы alert-хуков (`_alert_on_run_id`) не моделируются — они best-effort и
|
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`
|
Ровно эта последовательность достижима на проде: ветка таймаута `drain_inflight`
|
||||||
оставляет задачу внешнему hard-cancel'у, и до него она успевает несколько итераций
|
оставляет задачу внешнему hard-cancel'у, и до него она успевает несколько итераций
|
||||||
(`app/services/scheduler.py:141` — пульс на каждый батч). Без гейта по статусу
|
(`app/services/scheduler.py:141` — пульс на каждый батч). Без гейта по статусу
|
||||||
app-копия ЗАМЕНЯЛА counters и стирала метку: оборванный прогон снова читался как
|
пульс лёг бы поверх метки (одноимённые ключи мерж перезаписывает, а до #3390
|
||||||
|
app-копия и вовсе ЗАМЕНЯЛА весь словарь): оборванный прогон снова читался бы как
|
||||||
полный проход, а резюм его не подхватывал.
|
полный проход, а резюм его не подхватывал.
|
||||||
"""
|
"""
|
||||||
mod = _RUNS_MODULES[name]
|
mod = _RUNS_MODULES[name]
|
||||||
|
|
|
||||||
|
|
@ -4529,14 +4529,14 @@ async def run_domclick_city_sweep(
|
||||||
|
|
||||||
# #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции —
|
# #3369: чекпоинт (done_buckets) едет в КАЖДОЙ записи counters этой функции —
|
||||||
# это единообразие и защита от смены писателя, а НЕ восстановление потери.
|
# это единообразие и защита от смены писателя, а НЕ восстановление потери.
|
||||||
# Как есть сегодня: пишет kit-овый scraper_kit.orchestration.runs (импорт выше),
|
# Как есть сегодня: пишет scraper_kit.orchestration.runs (импорт выше) — с #3390
|
||||||
# и все четыре его писателя МЕРЖАТ jsonb — `counters = COALESCE(counters,'{}')
|
# единственная реализация, `app/services/scrape_runs.py` её алиас, — и все его
|
||||||
# || CAST(:counters AS jsonb)`, так что ключ, записанный раньше, переживает
|
# писатели МЕРЖАТ jsonb (`counters = COALESCE(counters,'{}') || CAST(:counters AS
|
||||||
# payload без него; плюс _pick_resume наследует done_buckets при claim (#3074).
|
# jsonb)`), так что ключ, записанный раньше, переживает payload без него; плюс
|
||||||
# То есть чекпоинт в БД не терялся. Перезаписывающий двойник существует —
|
# _pick_resume наследует done_buckets при claim (#3074). То есть чекпоинт в БД не
|
||||||
# app/services/scrape_runs.py (`counters = CAST(:counters AS jsonb)`), — но
|
# терялся. Перезаписывающий двойник (`counters = CAST(:counters AS jsonb)`) жил в
|
||||||
# этой функцией не вызывается. Полный payload делает ветки нечувствительными
|
# app-копии до #3390; полный payload оставлен и делает ветки нечувствительными к
|
||||||
# к тому, какой из двух писателей окажется на другом конце.
|
# тому, какой писатель окажется на другом конце.
|
||||||
_checkpoint: list[str] = []
|
_checkpoint: list[str] = []
|
||||||
|
|
||||||
def _payload() -> dict[str, Any]:
|
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-хелперы поверх таблицы
|
ЕДИНСТВЕННАЯ реализация (#3390). `app.services.scrape_runs` — тонкий алиас этого
|
||||||
`scrape_runs` (миграции 015 + 051). Развязка от `app.*`. Намеренные отличия от
|
модуля, а не вторая копия: до #3390 копий было две, и «байт-эквивалентными» они не
|
||||||
app-копии:
|
остались. Разошлись, в частности, семантика counters (здесь МЕРЖ `counters || :counters`,
|
||||||
1. `sentry_sdk` импортируется опционально (kit standalone-импортируем, а
|
там была ЗАМЕНА `counters = CAST(:counters AS jsonb)`), гейт по статусу, `honors_cancel`
|
||||||
`sentry-sdk` не входит в его зависимости). Если пакет не установлен — alert-хук
|
у `mark_cancelled` и сам состав функций. Расхождение стоило дважды за сутки ложных
|
||||||
best-effort no-op, поведение SQL-финализаторов идентично старому.
|
выводов на ревью (#3388: «отдать только флаг, остальное домержится» — на копии-заменителе
|
||||||
2. `mark_skipped` (#2658) есть только здесь: строки-пропуски создаёт исключительно
|
это стёрло бы измеренное; #3355).
|
||||||
планировщик (kit `_claim_run`/`scheduler_loop` + продуктовый cian pre_claim),
|
|
||||||
app-копии эта функция не нужна.
|
Мерж, а не замена, потому что писателей у одной строки прогона несколько (пульс задачи,
|
||||||
|
её финализатор, дрейн), и каждый знает лишь СВОИ ключи: чекпоинт `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 —
|
ВРЕМЯ ПИШЕТСЯ clock_timestamp(), А НЕ now() (#2702). `now()` в PostgreSQL —
|
||||||
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы
|
синоним `transaction_timestamp()`: он замерзает на СТАРТЕ транзакции. Финализаторы
|
||||||
|
|
@ -17,15 +26,20 @@ app-копии:
|
||||||
`finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с
|
`finished_at` получал время начала работы. Прод 2026-08-06: 133 прогона из 487 с
|
||||||
окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в
|
окном finished_at − started_at меньше секунды при работе дольше 10 с; 126 из них в
|
||||||
диапазоне 9-64 мс (= от коммита claim'а до первого запроса рабочей транзакции).
|
диапазоне 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
|
from __future__ import annotations
|
||||||
|
|
||||||
import json
|
import json
|
||||||
import logging
|
import logging
|
||||||
|
from collections import Counter
|
||||||
|
from collections.abc import Callable, Collection, Mapping
|
||||||
from functools import lru_cache
|
from functools import lru_cache
|
||||||
from collections.abc import Callable, Mapping
|
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
from sqlalchemy import text
|
from sqlalchemy import text
|
||||||
|
|
@ -338,10 +352,8 @@ def _phase_totally_failed(counters: Mapping[str, Any]) -> str | None:
|
||||||
# та же 'failed', но с другой формулировкой причины ("деградировал", не "провалился"), чтобы
|
# та же 'failed', но с другой формулировкой причины ("деградировал", не "провалился"), чтобы
|
||||||
# оператор видел разницу читая error, не только status.
|
# оператор видел разницу читая error, не только status.
|
||||||
#
|
#
|
||||||
# mark_backfill_finished (единственный писатель "attempted"/"failed" на верхнем уровне
|
# Единственный писатель "attempted"/"failed" на верхнем уровне counters —
|
||||||
# counters) живёт только в app.services.scrape_runs — здесь эта проверка сейчас неактивна
|
# mark_backfill_finished (ниже в этом же модуле); его зовут четыре detail-backfill'а.
|
||||||
# ни для одного реального вызывающего, но kit-копия держится байт-эквивалентной app-копии
|
|
||||||
# (см. docstring модуля), и будущий kit-native job с тем же словарём получит её даром.
|
|
||||||
FAILED_RATIO_FAILED_THRESHOLD = 0.5
|
FAILED_RATIO_FAILED_THRESHOLD = 0.5
|
||||||
FAILED_RATIO_DEGRADED_THRESHOLD = 0.15
|
FAILED_RATIO_DEGRADED_THRESHOLD = 0.15
|
||||||
# Минимум попыток, при котором доля вообще что-то значит — иначе 1 отказ из 2 (=0.5)
|
# Минимум попыток, при котором доля вообще что-то значит — иначе 1 отказ из 2 (=0.5)
|
||||||
|
|
@ -771,6 +783,27 @@ def update_heartbeat(db: Session, run_id: int, counters: dict[str, Any]) -> None
|
||||||
db.commit()
|
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:
|
def is_cancelled(db: Session, run_id: int) -> bool:
|
||||||
"""Проверить status='cancelled' (cooperative cancel в long-running pipeline)."""
|
"""Проверить status='cancelled' (cooperative cancel в long-running pipeline)."""
|
||||||
row = db.execute(
|
row = db.execute(
|
||||||
|
|
@ -959,8 +992,175 @@ def mark_banned(
|
||||||
_alert_on_run_id(db, run_id)
|
_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:
|
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(
|
result = db.execute(
|
||||||
text(
|
text(
|
||||||
"""
|
"""
|
||||||
|
|
@ -1048,3 +1248,22 @@ def list_all(
|
||||||
.all()
|
.all()
|
||||||
)
|
)
|
||||||
return total, [dict(r) for r in rows]
|
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