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
Две живые копии одного модуля с противоположной семантикой 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, утверждавшие про
живого «перезаписывающего двойника», приведены в соответствие.
205 lines
9.8 KiB
Python
205 lines
9.8 KiB
Python
"""#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
|