From 9b0b063fed95de2d4ca44f1b238e416563721651 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 5 Jul 2026 22:46:53 +0300 Subject: [PATCH] =?UTF-8?q?fix(generator):=20=D1=83=D1=81=D0=B8=D0=BB?= =?UTF-8?q?=D0=B5=D0=BD=D1=8B=20=D0=BF=D1=80=D0=BE=D0=B2=D0=B5=D1=80=D0=BA?= =?UTF-8?q?=D0=B8=20startup-history?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - коммит-гейт не запускал корневые контрактные тесты, а часть подтверждённых обходов могла снова смешать разные миры генератора. - Что: - добавлены цели make test, make lint и contract-test с тихим pytest-выводом через Docker. - закрыты обходы через generator-reset, неизвестную версию state и fail-open проверку DM-витрин. - усилены поведенческие контракты CHECK_LIVE_SEAM, профиля manifest и pause-check etl_pipeline; обновлены документы и issue 19. - Проверка: - make test; make lint; git diff --check. --- .../PRD.md | 7 +- .../issues/19-test-and-lint-targets.md | 22 ++- Makefile | 34 +++- README.md | 10 ++ airflow/dags/generator_control_dag.py | 14 +- docs/OPERATIONS.md | 13 +- docs/SUPERSET_DASHBOARD.md | 2 +- docs/course/lessons/00_kafka_intro.md | 4 +- docs/runbooks/startup-history.md | 17 +- generator/README.md | 25 +-- .../clickstream_generator/airflow_control.py | 15 ++ generator/src/clickstream_generator/state.py | 13 +- generator/tests/test_airflow_control.py | 16 ++ .../test_generator_control_dag_contract.py | 15 +- generator/tests/test_state.py | 34 +++- scripts/check_generated_analytics.sh | 7 +- scripts/run_generator.sh | 8 +- ...test_generated_analytics_check_contract.py | 149 ++++++++++++++++++ 18 files changed, 346 insertions(+), 59 deletions(-) diff --git a/.scratch/generator-model-time-startup-history/PRD.md b/.scratch/generator-model-time-startup-history/PRD.md index 1d25797..be3e4d1 100644 --- a/.scratch/generator-model-time-startup-history/PRD.md +++ b/.scratch/generator-model-time-startup-history/PRD.md @@ -91,14 +91,17 @@ Status: Implemented startup-history/Superset, курс на чистом стенде. Подробности и коммиты см. в issue-файлах и `coordinator-journal.md`. -Дальнейшая задача вне текущей очереди: +Оставшаяся задача вне текущей очереди: - `issues/13-...` — доливка истории от слепка; без приоритета, `needs-triage`. После кросс-линейного ревью 2026-07-05 блокируется задачами 15 и 17: доливка тиражирует стыки мира и должна опираться на доверенную проверку фактической границы. + +Дополнительная задача закрыта: + - `issues/19-test-and-lint-targets.md` — единые `make test` и `make lint` для - коммит-гейта; `ready-for-agent`. Возникло из coordinator-loop 2026-07-05: + коммит-гейта; `done`. Возникло из coordinator-loop 2026-07-05: координатору пришлось вручную собирать тестовый набор из частных команд. Доработки по кросс-линейному ревью 2026-07-05 (закрыто): diff --git a/.scratch/generator-model-time-startup-history/issues/19-test-and-lint-targets.md b/.scratch/generator-model-time-startup-history/issues/19-test-and-lint-targets.md index 54707d0..78aba8a 100644 --- a/.scratch/generator-model-time-startup-history/issues/19-test-and-lint-targets.md +++ b/.scratch/generator-model-time-startup-history/issues/19-test-and-lint-targets.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: done # Единые make test и make lint для коммит-гейта @@ -24,17 +24,17 @@ config --quiet`, `py_compile`, `bash -n`, `git diff --check` и отдельны ## Acceptance criteria -- [ ] `make test` существует и запускает полный быстрый набор проверок, нужный +- [x] `make test` существует и запускает полный быстрый набор проверок, нужный перед коммитом: тесты генератора, контрактные тесты верхнего уровня и ключевые проверки конфигурации. -- [ ] `make lint` существует и запускает статические проверки, которые сейчас +- [x] `make lint` существует и запускает статические проверки, которые сейчас выполняются вручную: синтаксис Python, `bash -n`, `docker compose config --quiet`, `git diff --check` или обоснованно выбранный эквивалент. -- [ ] Цели используют `uv` там, где нужен Python вне Docker, и не зависят от +- [x] Цели используют `uv` там, где нужен Python вне Docker, и не зависят от локального `pytest`, случайно установленного на хосте. -- [ ] Шумные команды не вываливают полный список тестов при успешном прогоне; +- [x] Шумные команды не вываливают полный список тестов при успешном прогоне; подробный вывод доступен через явный режим или отдельную команду. -- [ ] `docs/OPERATIONS.md` или `README.md` кратко объясняет, когда запускать +- [x] `docs/OPERATIONS.md` или `README.md` кратко объясняет, когда запускать `make test`, `make lint` и какие более дорогие стендовые проверки остаются отдельными. @@ -51,3 +51,13 @@ config --quiet`, `py_compile`, `bash -n`, `git diff --check` и отдельны `uv run ... pytest` падал до запуска тестов из-за `snap-confine`, поэтому цель должна либо идти через устойчивый Docker-путь, либо явно диагностировать такую ошибку. +- По ревью Claude дополнительно подключены корневые контрактные тесты, усилены + ключевые проверки поведения `CHECK_LIVE_SEAM`/профиля/DM-витрин, а также + закрыты подтверждённые обходы через `generator-reset` и безверсионный state. + +## Проверка + +- `make contract-test` — 16 passed. +- `make generator-test` — 190 passed. +- `make lint` — passed. +- `make test` — 190 generator tests, 16 contract tests, compose config passed. diff --git a/Makefile b/Makefile index 966ae54..01e0384 100644 --- a/Makefile +++ b/Makefile @@ -5,9 +5,14 @@ superset-init superset-dashboard superset-ui superset-restart \ generator-up generator-down generator-logs generator-restart \ generator-backfill generator-continue generator-reset \ - generator-test generator-test-build generator-test-cov + generator-test generator-test-build generator-test-cov \ + contract-test config-test test lint COMPOSE ?= docker compose +PROFILE ?= ci +PYTEST_ARGS ?= -q -p no:cacheprovider +PYTHON_CHECK_PATHS := airflow/dags configs generator scripts superset tests +BASH_SCRIPTS := $(shell find scripts -type f -name '*.sh' | sort) # ============================================================================ # Основные команды @@ -46,7 +51,7 @@ generated-history-analytics: # Повторяемая проверка после прогона стартовой истории generated-history-check: - PROFILE="$${PROFILE:-$(PROFILE)}" GEN_LAUNCH_PROFILE="$${GEN_LAUNCH_PROFILE:-$${PROFILE:-$(PROFILE)}}" CHECK_LIVE_SEAM="$${CHECK_LIVE_SEAM:-1}" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh + PROFILE="$${PROFILE:-$${GEN_LAUNCH_PROFILE:-$(PROFILE)}}" GEN_LAUNCH_PROFILE="$${GEN_LAUNCH_PROFILE:-$${PROFILE:-$(PROFILE)}}" CHECK_LIVE_SEAM="$${CHECK_LIVE_SEAM:-1}" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh # Быстрый runtime gate issue 17: короткий daily-wave backfill + live seam без Superset UI generated-history-runtime-check: @@ -154,9 +159,32 @@ generator-test-build: # Запустить тесты генератора (pytest) generator-test: generator-test-build @echo "=== Запуск тестов генератора ===" - docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -v + docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ $(PYTEST_ARGS) # Запустить тесты с покрытием generator-test-cov: generator-test-build @echo "=== Запуск тестов с покрытием ===" docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -v --cov=. --cov-report=term-missing + +# Запустить контрактные тесты верхнего уровня +contract-test: generator-test-build + @echo "=== Запуск контрактных тестов репозитория ===" + docker run --rm -v $(PWD):/workspace -w /workspace generator:test pytest tests/ $(PYTEST_ARGS) + +# Проверить конфигурацию compose без запуска сервисов +config-test: + $(COMPOSE) config --quiet + +# Быстрый предкоммитный набор: тесты генератора, контракты и конфигурация +test: generator-test contract-test config-test + +# Статические проверки без запуска стенда +lint: generator-test-build + @echo "=== Проверка синтаксиса Python ===" + docker run --rm -e PYTHONPYCACHEPREFIX=/tmp/pycache -v $(PWD):/workspace -w /workspace generator:test python -m compileall -q $(PYTHON_CHECK_PATHS) + @echo "=== Проверка синтаксиса Bash ===" + bash -n $(BASH_SCRIPTS) + @echo "=== Проверка docker compose ===" + $(COMPOSE) config --quiet + @echo "=== Проверка пробелов в diff ===" + git diff --check diff --git a/README.md b/README.md index 5bed3e4..d312d5c 100644 --- a/README.md +++ b/README.md @@ -82,6 +82,16 @@ CHECK_LIVE_SEAM=0 make generated-history-check make generated-history-runtime-check ``` +Перед коммитом используйте быстрые проверки: + +```bash +make test +make lint +``` + +Они не чистят volumes и не запускают долгие стендовые сценарии. Полная проверка +стыка backfill/live остаётся отдельной командой `make generated-history-runtime-check`. + Сохранить стартовую историю в файл и восстановить её без новой генерации можно по [runbook стартовой истории](./docs/runbooks/startup-history.md). diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/generator_control_dag.py index d24ec60..519f3d5 100644 --- a/airflow/dags/generator_control_dag.py +++ b/airflow/dags/generator_control_dag.py @@ -30,6 +30,7 @@ from clickstream_generator.airflow_control import ( load_manifest_from_kafka, run_backfill, run_import, + target_dag_trigger_error, validate_import_artifact, ) from clickstream_generator.launch import PROFILES @@ -158,16 +159,9 @@ def assert_target_dag_not_paused(dag_id: str, session=None) -> None: fail_when_dag_is_paused, поэтому паузу проверяем заранее через DagModel. """ dag_model = session.query(DagModel).filter(DagModel.dag_id == dag_id).one_or_none() - if dag_model is None: - raise AirflowException( - f"DAG {dag_id} ещё не найден Airflow. Подождите парсинга DAG-файлов " - "или перезапустите Airflow через make up." - ) - if dag_model.is_paused: - raise AirflowException( - f"DAG {dag_id} стоит на паузе: снимите паузу в Airflow UI, " - "затем повторите generator_control." - ) + error = target_dag_trigger_error(dag_id, dag_model) + if error: + raise AirflowException(error) with DAG( diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 85c4203..4f0384d 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -18,11 +18,19 @@ - `make data` (архивный путь: заливает `data/*.jsonl` в Kafka; не основной источник аналитики) - `make transform` (запускает batch-процесс ODS -> DDS -> DM) - `make superset-init` (повторная инициализация Superset: подключение к ClickHouse, датасеты, дашборд) +- `make test` (быстрый предкоммитный набор: тесты генератора, верхние + контрактные тесты и проверка compose-конфигурации) +- `make lint` (статические проверки: синтаксис Python и Bash, compose-конфиг, + пробелы в diff) - `docker compose ps` - `docker compose logs -f --tail=200 ` - `docker compose down` (сохраняет named volumes, включая `clickhouse-data`) - `docker compose down -v` (удаляет named volumes, использовать осознанно) +`make test` и `make lint` не заменяют стендовые проверки, которые управляют +volumes или live-генератором. Для стыка backfill/live отдельно запускайте +`make generated-history-runtime-check`. + ## Порты Порты задаются в `docker-compose.yml`: @@ -278,8 +286,9 @@ docker compose exec -T kafka /opt/kafka/bin/kafka-console-consumer.sh \ Live-продолжение стартует с этого state. Используйте тот же профиль или ту же длительность, что были у backfill. Команда сама выставит `GEN_STATE_RESET=false`. Если state записан старой версией генератора, запуск теперь падает громко: -сначала очистите стенд через `make clean` и пересоздайте историю либо явно -начните новый мир через `make generator-reset`. +сначала очистите стенд через `make clean` и пересоздайте историю. Команда +`make generator-reset` подходит только для осознанного нового live-мира на +чистом стенде: она тоже проверяет, что Kafka data-топики и STG пустые. ```bash PROFILE=daily-wave make generator-continue diff --git a/docs/SUPERSET_DASHBOARD.md b/docs/SUPERSET_DASHBOARD.md index 4397705..a27022b 100644 --- a/docs/SUPERSET_DASHBOARD.md +++ b/docs/SUPERSET_DASHBOARD.md @@ -145,7 +145,7 @@ make logs service=superset # Логи сервиса # ETL make generated-history-analytics # Чистый прогон генерации до Superset -make generated-history-check # Проверка DM и Superset metadata +CHECK_LIVE_SEAM=0 make generated-history-check # Проверка DM и Superset после backfill make ddl # Применение DDL в ClickHouse make data # Архивная загрузка data/*.jsonl в Kafka make transform # Запуск batch-процесса diff --git a/docs/course/lessons/00_kafka_intro.md b/docs/course/lessons/00_kafka_intro.md index df24bff..53544f6 100644 --- a/docs/course/lessons/00_kafka_intro.md +++ b/docs/course/lessons/00_kafka_intro.md @@ -67,7 +67,9 @@ make up - `generator_state` — слепок состояния генератора на правой границе стартовой истории. По нему live-продолжение понимает, откуда продолжать тот же мир; - `generator_startup_history_manifest` — паспорт стартовой истории: seed, границы - модельного времени и контрольные числа. + модельного времени и контрольные числа; +- `generator_batch_history` — журнал батчей генератора: сколько сообщений он + отправил и чем закончилась каждая пачка записи. Служебные топики нужны стенду, но в упражнениях курса мы их не меняем. У каждого нашего топика событий в колонке с партициями стоит **1**: топик маленький, делить не на что. diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 00d839c..81cc5c3 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -40,7 +40,9 @@ 1. Поднимите стенд: `make up`. 2. Если DDL ещё не применён, запустите `ddl_init`. -3. Откройте `generator_control` и выберите `operation`. +3. Снимите паузу с `etl_pipeline`, если он ещё paused: + `docker compose exec -T airflow-webserver airflow dags unpause etl_pipeline`. +4. Откройте `generator_control` и выберите `operation`. Операции: @@ -148,18 +150,21 @@ PROFILE=daily-wave make generator-continue У `daily-wave` скорость ×60. Долгий простой стенда создаёт большую дыру в модельном времени: ночь простоя может стать десятками модельных суток без -событий. Для чистой демонстрации лучше запустите `make generator-reset` или -повторите импорт стартовой истории через `make startup-history-import`. +событий. Для чистой демонстрации лучше очистите стенд, затем запустите +`make generator-reset` или повторите импорт стартовой истории через +`make startup-history-import`. +`make generator-reset` начнёт новый live-мир только на чистом стенде; если в +Kafka data-топиках или STG уже есть строки, команда попросит `make clean`. Если читаемый state есть, но настройки не совпадают, генератор падает с перечнем полей. Это защита от смешения разных миров. Для намеренного нового мира -используйте `make generator-reset` или `make clean`. +сначала очистите стенд через `make clean`, затем запускайте нужный сценарий. После обновления кода старый state может оказаться в старом формате. При `GEN_STATE_RESET=false` это теперь громкий отказ, а не тихий старт с нуля поверх старой истории. Оператору нужно выбрать одно из двух: очистить стенд через -`make clean` и заново создать стартовую историю, либо осознанно начать новый мир -через `GEN_STATE_RESET=true` / `make generator-reset`. +`make clean` и заново создать стартовую историю, либо осознанно начать новый +live-мир через `make generator-reset` на чистом стенде. После нестандартного мира `make generator-continue` нужно запускать с теми же настройками, что были у backfill/import. При расхождении генератор громко diff --git a/generator/README.md b/generator/README.md index 078cc33..389ab04 100644 --- a/generator/README.md +++ b/generator/README.md @@ -176,24 +176,25 @@ make generator-backfill # Продолжить live-поток из state make generator-continue -# Начать live-поток как новый мир +# Начать live-поток как новый мир на чистом стенде make generator-reset # Чистый аналитический прогон всего стенда make generated-history-analytics -# Запуск тестов -make generator-test +# Быстрые проверки перед коммитом +make test +make lint ``` Основной ручной путь запуска — глаголы `generator-backfill`, `generator-continue` и `generator-reset`. Старые `GEN_RUN_MODE`, `GEN_STATE_RESET` и `GEN_MODEL_T_END` остаются низкоуровневым способом для отладки. -`generator-backfill` и `startup-history-import` перед записью проверяют, что -Kafka data-топики и STG пустые. `generator-continue` продолжает только -совместимый state; старый формат state при `GEN_STATE_RESET=false` даёт отказ с -подсказкой очистить стенд или явно начать новый мир. +`generator-backfill`, `startup-history-import` и `generator-reset` перед записью +проверяют, что Kafka data-топики и STG пустые. `generator-continue` продолжает +только совместимый state; старый формат state при `GEN_STATE_RESET=false` даёт +отказ с подсказкой очистить стенд или явно начать новый мир на чистом стенде. ## Метрики Prometheus @@ -351,16 +352,22 @@ GEN_STATE_ENABLED=false docker compose up -d generator ```bash # Через Makefile (рекомендуется) -make generator-test +make test +make lint # Вручную через Docker docker build -t generator:test . -docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -v +docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -q # Конкретный файл тестов docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/test_generation.py -v ``` +`make test` запускает тесты генератора, контрактные тесты верхнего уровня и +проверку compose-конфигурации. `make lint` проверяет синтаксис Python и Bash, +compose-конфиг и пробелы в diff. Долгие стендовые проверки вроде +`make generated-history-runtime-check` запускаются отдельно. + ### Структура тестов ``` diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index c6e6fae..d127976 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -199,6 +199,21 @@ def assert_clickhouse_matches_manifest(manifest: dict, clickhouse_hook) -> None: ) +def target_dag_trigger_error(dag_id: str, dag_model) -> str | None: + """Возвращает причину, почему зависимый DAG нельзя запускать.""" + if dag_model is None: + return ( + f"DAG {dag_id} ещё не найден Airflow. Подождите парсинга DAG-файлов " + "или перезапустите Airflow через make up." + ) + if dag_model.is_paused: + return ( + f"DAG {dag_id} стоит на паузе: снимите паузу в Airflow UI, " + "затем повторите generator_control." + ) + return None + + def default_artifact_path(operation: str) -> str: """Возвращает путь артефакта по умолчанию в общем томе data.""" filename = ( diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py index 1780ba8..4d70463 100644 --- a/generator/src/clickstream_generator/state.py +++ b/generator/src/clickstream_generator/state.py @@ -11,7 +11,6 @@ logger = logging.getLogger("generator") STATE_VERSION = "3.0" -KNOWN_OLD_STATE_VERSIONS = {"2.0"} class UnsupportedStateVersionError(ValueError): @@ -245,15 +244,9 @@ class GeneratorState: try: if not isinstance(data, dict): raise ValueError("state must be an object") - version = data.get("version", "1.0") + version = str(data.get("version", "1.0")) if version != STATE_VERSION: - if str(version) in KNOWN_OLD_STATE_VERSIONS: - raise UnsupportedStateVersionError(str(version)) - logger.warning( - "Unsupported generator state version %s, will start fresh", - version, - ) - raise ValueError(f"unsupported state version: {version}") + raise UnsupportedStateVersionError(version) _validate_v3_payload(data) model_timestamp = _parse_aware_utc( data["model_timestamp"], @@ -309,5 +302,7 @@ class GeneratorState: """Безопасная загрузка state с graceful degradation.""" try: return cls.from_dict(data) + except UnsupportedStateVersionError: + raise except Exception: return None diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py index 19f4ce3..994496d 100644 --- a/generator/tests/test_airflow_control.py +++ b/generator/tests/test_airflow_control.py @@ -5,6 +5,7 @@ from dataclasses import replace from datetime import datetime, timezone from pathlib import Path +from types import SimpleNamespace import pytest @@ -49,6 +50,21 @@ def test_import_env_requires_artifact_path_and_uses_backfill_contract(): assert env["GEN_HISTORY_DURATION"] == "2d" +def test_target_dag_trigger_error_covers_missing_paused_and_ready_states(): + """Проверка зависимого DAG различает три реальные ветки.""" + from clickstream_generator.airflow_control import target_dag_trigger_error + + assert "ещё не найден" in target_dag_trigger_error("etl_pipeline", None) + assert "стоит на паузе" in target_dag_trigger_error( + "etl_pipeline", + SimpleNamespace(is_paused=True), + ) + assert target_dag_trigger_error( + "etl_pipeline", + SimpleNamespace(is_paused=False), + ) is None + + def test_assert_stand_clean_rejects_non_empty_stg_before_writes(): """Backfill/import не стартуют на непустом STG.""" from clickstream_generator.airflow_control import assert_stg_tables_empty diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py index c05c3cb..f746785 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_generator_control_dag_contract.py @@ -57,9 +57,9 @@ def test_generator_control_prechecks_etl_dag_not_paused_before_waiting(): assert 'op_kwargs={"dag_id": ETL_DAG_ID}' in text assert "session.query(DagModel)" in text assert "DagModel.dag_id == dag_id" in text + assert "target_dag_trigger_error" in text assert "Airflow 2.10.5" in text assert "fail_when_dag_is_paused" in text - assert "снимите паузу" in text assert 'return "check_etl_not_paused_before_backfill"' in text assert 'return "check_etl_not_paused_before_import"' in text assert "check_etl_not_paused_before_backfill >> precheck_backfill_task" in text @@ -177,6 +177,19 @@ def test_console_continue_checks_state_or_clean_stand_before_live_start(): ) +def test_console_reset_checks_clean_stand_before_live_start(): + """Reset не пишет новый мир поверх старых данных.""" + run_generator = (REPO_ROOT / "scripts" / "run_generator.sh").read_text( + encoding="utf-8" + ) + reset_branch = run_generator.split("else\n", maxsplit=1)[1] + + assert "bash \"${SCRIPT_DIR}/assert_stand_clean.sh\" clean" in reset_branch + assert reset_branch.index("assert_stand_clean.sh\" clean") < reset_branch.index( + "up -d --build generator" + ) + + def test_clean_start_paths_include_live_generator_profile(): """Все чистые сбросы видят профильный live-генератор.""" generated_history = ( diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index 3dc6873..e3b2228 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -432,6 +432,22 @@ class TestGeneratorStateValidation: with pytest.raises(UnsupportedStateVersionError, match="2.0.*3.0"): GeneratorState.from_dict(data) + def test_from_dict_rejects_future_state_version_loudly(self): + """Любая неизвестная версия state требует решения оператора.""" + data = _make_valid_state_data() + data["version"] = "4.0" + + with pytest.raises(UnsupportedStateVersionError, match="4.0.*3.0"): + GeneratorState.from_dict(data) + + def test_from_dict_safe_does_not_hide_old_state_version(self): + """Безопасная загрузка не превращает старый state в тихий fresh-start.""" + data = _make_valid_state_data() + data["version"] = "2.0" + + with pytest.raises(UnsupportedStateVersionError, match="2.0.*3.0"): + GeneratorState.from_dict_safe(data) + class TestKafkaStateManager: """Тесты менеджера состояния.""" @@ -536,10 +552,14 @@ class TestKafkaStateManager: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - # Мокаем consumer с невалидным сообщением + # Мокаем consumer с невалидным сообщением текущей версии. mock_message = MagicMock() mock_message.key = b"default" - mock_message.value = {"tick": 42, "rng_state": "invalid"} + mock_message.value = { + "tick": 42, + "rng_state": "invalid", + "version": "3.0", + } mock_consumer = MagicMock() mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) @@ -749,8 +769,8 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_version_1_state_returns_none_with_warning(self, caplog): - """Старое state v1 не восстанавливается и даёт чистый старт.""" + def test_load_version_1_state_fails_loudly(self, caplog): + """Старое state v1 не скрывается за чистым стартом.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: @@ -774,10 +794,10 @@ class TestKafkaStateManager: mock_consumer_class.return_value = mock_consumer manager = KafkaStateManager("kafka:29092") - with caplog.at_level(logging.WARNING, logger="generator"): - result = manager.load() + with caplog.at_level(logging.ERROR, logger="generator"): + with pytest.raises(UnsupportedStateVersionError, match="1.0.*3.0"): + manager.load() - assert result is None assert "version" in caplog.text def test_load_ignores_wrong_key(self): diff --git a/scripts/check_generated_analytics.sh b/scripts/check_generated_analytics.sh index f5ab409..72cd11c 100644 --- a/scripts/check_generated_analytics.sh +++ b/scripts/check_generated_analytics.sh @@ -520,12 +520,17 @@ FROM ORDER BY source FORMAT TabSeparated" +if ! views_rows="$(ch_query "${views_query}")"; then + fail "не удалось прочитать основные DM-витрины" +fi +[[ -n "${views_rows}" ]] || fail "ClickHouse не вернул строки основных DM-витрин" + while IFS=$'\t' read -r source rows; do [[ -n "${source}" ]] || continue [[ "${rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число строк для ${source}" (( rows > 0 )) || fail "${source} пустая" echo "${source}=${rows}" -done < <(ch_query "${views_query}") +done <<< "${views_rows}" if [[ "${REQUIRE_SUPERSET}" != "1" ]]; then echo "" diff --git a/scripts/run_generator.sh b/scripts/run_generator.sh index 66598ab..c908d87 100644 --- a/scripts/run_generator.sh +++ b/scripts/run_generator.sh @@ -87,6 +87,12 @@ elif [[ "${VERB}" == "continue" ]]; then echo "Шаг 3: запуск live-генератора" ${COMPOSE_BIN} up -d --build generator else - echo "Шаг 1: запуск live-генератора" + echo "Шаг 1: запуск Kafka и ClickHouse" + ${COMPOSE_BIN} up -d kafka clickhouse + + echo "Шаг 2: предпроверка чистого стенда" + bash "${SCRIPT_DIR}/assert_stand_clean.sh" clean + + echo "Шаг 3: запуск live-генератора" ${COMPOSE_BIN} up -d --build generator fi diff --git a/tests/test_generated_analytics_check_contract.py b/tests/test_generated_analytics_check_contract.py index 85d7ca0..ad53ed7 100644 --- a/tests/test_generated_analytics_check_contract.py +++ b/tests/test_generated_analytics_check_contract.py @@ -1,7 +1,110 @@ from pathlib import Path +import os +import subprocess +import textwrap REPO_ROOT = Path(__file__).resolve().parents[1] +CHECK_SCRIPT = REPO_ROOT / "scripts" / "check_generated_analytics.sh" + + +def _fake_compose( + tmp_path, + *, + manifest_profile="ci", + live_rows="3", + views_mode="ok", +): + fake = tmp_path / "docker-compose" + fake.write_text( + textwrap.dedent( + f"""\ + #!/usr/bin/env bash + set -euo pipefail + + args="$*" + query="" + previous="" + for arg in "$@"; do + if [[ "$previous" == "--query" ]]; then + query="$arg" + break + fi + previous="$arg" + done + + if [[ "$args" == *"kafka-manifest-summary"* ]]; then + printf '%s\\n' '16054\t2800\t930\t2026-01-01T00:00:00+00:00\t2026-01-01T05:59:00+00:00\t2026-01-01T00:00:00+00:00\t2026-01-01T06:00:00+00:00\t{manifest_profile}' + exit 0 + fi + + if [[ "$args" == *"clickhouse-client"* ]]; then + if [[ "$query" == *"hex(sipHash128"* ]]; then + printf '%s\\n' '16054\t2800\t930\t2026-01-01 00:00:00.000000\t2026-01-01 05:59:00.000000\t1\t1\t1\tdeadbeef' + elif [[ "$query" == *"returning_users / users"* ]]; then + printf '%s\\n' '930\t120\t0.129\t4' + elif [[ "$query" == *"short_visit_share"* ]]; then + printf '%s\\n' '2800\t0.2\t5\t5.7\t0\t120\t900\t18' + elif [[ "$query" == *"minIf(event_ts"* ]]; then + printf '%s\\n' '2800\t2100\t1500\t900\t650\t1\t0.23' + elif [[ "$query" == *"has_home"* ]]; then + printf '%s\\n' '2800\t2200\t1600\t950\t700\t1\t0.25' + elif [[ "$query" == *"SELECT count()"* && "$query" == *"FROM dds.event"* ]]; then + printf '%s\\n' '{live_rows}' + elif [[ "$query" == *"uniqExact(event_id)"* ]]; then + printf '%s\\n' '16200\t16200\t0' + elif [[ "$query" == *"per_event_homogeneous_visits"* ]]; then + printf '%s\\n' '19\t19\t19' + elif [[ "$query" == *"ods_device_rows"* ]]; then + printf '%s\\n' '19\t19\t19\t0\t0' + elif [[ "$query" == *"SELECT source, rows"* ]]; then + if [[ "{views_mode}" == "fail" ]]; then + exit 42 + fi + printf '%s\\n' \ + 'dm.dq_summary\t1' \ + 'dm.v_daily_traffic\t1' \ + 'dm.v_events_enriched\t16054' \ + 'dm.v_session_overview\t1' \ + 'dm.v_top_pages_daily\t1' \ + 'dm.v_utm_effectiveness\t1' + else + echo "unexpected ClickHouse query: $query" >&2 + exit 91 + fi + exit 0 + fi + + echo "unexpected compose call: $args" >&2 + exit 92 + """ + ), + encoding="utf-8", + ) + fake.chmod(0o755) + return fake + + +def _run_generated_history_check(tmp_path, *, fake_compose, **env_overrides): + env = os.environ.copy() + env.update( + { + "COMPOSE_BIN": str(fake_compose), + "REQUIRE_SUPERSET": "0", + "WAIT_LIVE_ROWS_SECONDS": "0", + } + ) + env.update(env_overrides) + return subprocess.run( + ["bash", str(CHECK_SCRIPT)], + cwd=REPO_ROOT, + env=env, + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + timeout=30, + check=False, + ) def test_generated_history_check_uses_actual_manifest_boundary_and_profile(): @@ -14,6 +117,7 @@ def test_generated_history_check_uses_actual_manifest_boundary_and_profile(): assert "kafka-manifest-summary" in script assert "manifest_model_t_end" in script assert "CH_MODEL_T_END=\"$(clickhouse_datetime_literal \"${manifest_model_t_end}\")\"" in script + assert "PROFILE ?= ci" in makefile assert "GEN_LAUNCH_PROFILE" in makefile assert "PROFILE" in makefile @@ -31,6 +135,51 @@ def test_generated_history_check_does_not_skip_required_live_seam(): assert 'CHECK_LIVE_SEAM="$${CHECK_LIVE_SEAM:-1}"' in makefile +def test_generated_history_check_rejects_manifest_profile_mismatch(tmp_path): + """Проверка падает, если ожидаемый профиль не совпал с manifest.""" + fake_compose = _fake_compose(tmp_path, manifest_profile="ci") + + result = _run_generated_history_check( + tmp_path, + fake_compose=fake_compose, + PROFILE="daily-wave", + CHECK_LIVE_SEAM="0", + ) + + assert result.returncode != 0 + assert "профиля daily-wave, но manifest от ci" in result.stderr + + +def test_generated_history_check_requires_live_rows_when_seam_is_required(tmp_path): + """CHECK_LIVE_SEAM=1 падает, если live-продолжение не записало строк.""" + fake_compose = _fake_compose(tmp_path, live_rows="0") + + result = _run_generated_history_check( + tmp_path, + fake_compose=fake_compose, + PROFILE="ci", + CHECK_LIVE_SEAM="1", + ) + + assert result.returncode != 0 + assert "live-продолжение не записало строки" in result.stderr + + +def test_generated_history_check_fails_when_dm_views_query_fails(tmp_path): + """Ошибка запроса DM-витрин не превращается в пустой успешный цикл.""" + fake_compose = _fake_compose(tmp_path, views_mode="fail") + + result = _run_generated_history_check( + tmp_path, + fake_compose=fake_compose, + PROFILE="ci", + CHECK_LIVE_SEAM="0", + ) + + assert result.returncode != 0 + assert "не удалось прочитать основные DM-витрины" in result.stderr + + def test_clean_generated_history_run_explicitly_skips_live_seam(): """Чистый backfill без live-продолжения отключает проверку стыка явно.""" script = (REPO_ROOT / "scripts" / "run_generated_history_analytics.sh").read_text(