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