gendesign/tradein-mvp/backend/tests/test_3390_single_runs_module.py
bot-backend 36f2429fbe
All checks were successful
CI Trade-In / changes (pull_request) Successful in 7s
CI Trade-In / browser-tests (pull_request) Has been skipped
CI / openapi-codegen-check (pull_request) Has been skipped
CI Trade-In / backend-tests (pull_request) Successful in 4m58s
CI / changes (pull_request) Successful in 9s
CI Trade-In / frontend-checks (pull_request) Has been skipped
CI / backend-tests (pull_request) Has been skipped
CI / frontend-tests (pull_request) Has been skipped
refactor(tradein/runs): одна реализация scrape_runs — kit, семантика counters мерж (#3390)
Две живые копии одного модуля с противоположной семантикой counters: app
`mark_done`/`mark_failed`/`mark_banned`/`update_heartbeat` ЗАМЕНЯЛИ
(`counters = CAST(:counters AS jsonb)`), kit — МЕРЖИЛИ
(`COALESCE(counters,'{}') || …`). Расхождение дважды за сутки дало ложные
выводы на ревью (#3388 «отдать только флаг, остальное домержится» — на
replace это стёрло бы измеренное; #3355). Разошлись и другие места: гейт
статуса, `honors_cancel` у mark_cancelled (был только в app), `mark_skipped`
(только в kit), `mark_backfill_finished`/`distinct_sources` (только в app).

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

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

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

Комментарии в app/services/scheduler.py и kit/pipeline.py, утверждавшие про
живого «перезаписывающего двойника», приведены в соответствие.
2026-09-06 11:45:56 +05:00

205 lines
9.8 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters

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

"""#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