fix(generator): усилены проверки startup-history

- Зачем:
  - коммит-гейт не запускал корневые контрактные тесты, а часть подтверждённых обходов могла снова смешать разные миры генератора.
- Что:
  - добавлены цели 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.
This commit is contained in:
2026-07-05 22:46:53 +03:00
parent 0faa5cc219
commit 9b0b063fed
18 changed files with 346 additions and 59 deletions
@@ -91,14 +91,17 @@ Status: Implemented
startup-history/Superset, курс на чистом стенде. Подробности и коммиты см. в startup-history/Superset, курс на чистом стенде. Подробности и коммиты см. в
issue-файлах и `coordinator-journal.md`. issue-файлах и `coordinator-journal.md`.
Дальнейшая задача вне текущей очереди: Оставшаяся задача вне текущей очереди:
- `issues/13-...` — доливка истории от слепка; без приоритета, - `issues/13-...` — доливка истории от слепка; без приоритета,
`needs-triage`. После кросс-линейного ревью 2026-07-05 блокируется `needs-triage`. После кросс-линейного ревью 2026-07-05 блокируется
задачами 15 и 17: доливка тиражирует стыки мира и должна опираться на задачами 15 и 17: доливка тиражирует стыки мира и должна опираться на
доверенную проверку фактической границы. доверенную проверку фактической границы.
Дополнительная задача закрыта:
- `issues/19-test-and-lint-targets.md` — единые `make test` и `make lint` для - `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 (закрыто): Доработки по кросс-линейному ревью 2026-07-05 (закрыто):
@@ -1,4 +1,4 @@
Status: ready-for-agent Status: done
# Единые make test и make lint для коммит-гейта # Единые make test и make lint для коммит-гейта
@@ -24,17 +24,17 @@ config --quiet`, `py_compile`, `bash -n`, `git diff --check` и отдельны
## Acceptance criteria ## Acceptance criteria
- [ ] `make test` существует и запускает полный быстрый набор проверок, нужный - [x] `make test` существует и запускает полный быстрый набор проверок, нужный
перед коммитом: тесты генератора, контрактные тесты верхнего уровня и перед коммитом: тесты генератора, контрактные тесты верхнего уровня и
ключевые проверки конфигурации. ключевые проверки конфигурации.
- [ ] `make lint` существует и запускает статические проверки, которые сейчас - [x] `make lint` существует и запускает статические проверки, которые сейчас
выполняются вручную: синтаксис Python, `bash -n`, `docker compose config выполняются вручную: синтаксис Python, `bash -n`, `docker compose config
--quiet`, `git diff --check` или обоснованно выбранный эквивалент. --quiet`, `git diff --check` или обоснованно выбранный эквивалент.
- [ ] Цели используют `uv` там, где нужен Python вне Docker, и не зависят от - [x] Цели используют `uv` там, где нужен Python вне Docker, и не зависят от
локального `pytest`, случайно установленного на хосте. локального `pytest`, случайно установленного на хосте.
- [ ] Шумные команды не вываливают полный список тестов при успешном прогоне; - [x] Шумные команды не вываливают полный список тестов при успешном прогоне;
подробный вывод доступен через явный режим или отдельную команду. подробный вывод доступен через явный режим или отдельную команду.
- [ ] `docs/OPERATIONS.md` или `README.md` кратко объясняет, когда запускать - [x] `docs/OPERATIONS.md` или `README.md` кратко объясняет, когда запускать
`make test`, `make lint` и какие более дорогие стендовые проверки остаются `make test`, `make lint` и какие более дорогие стендовые проверки остаются
отдельными. отдельными.
@@ -51,3 +51,13 @@ config --quiet`, `py_compile`, `bash -n`, `git diff --check` и отдельны
`uv run ... pytest` падал до запуска тестов из-за `snap-confine`, поэтому `uv run ... pytest` падал до запуска тестов из-за `snap-confine`, поэтому
цель должна либо идти через устойчивый Docker-путь, либо явно диагностировать цель должна либо идти через устойчивый 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.
+31 -3
View File
@@ -5,9 +5,14 @@
superset-init superset-dashboard superset-ui superset-restart \ superset-init superset-dashboard superset-ui superset-restart \
generator-up generator-down generator-logs generator-restart \ generator-up generator-down generator-logs generator-restart \
generator-backfill generator-continue generator-reset \ 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 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: 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 # Быстрый runtime gate issue 17: короткий daily-wave backfill + live seam без Superset UI
generated-history-runtime-check: generated-history-runtime-check:
@@ -154,9 +159,32 @@ generator-test-build:
# Запустить тесты генератора (pytest) # Запустить тесты генератора (pytest)
generator-test: generator-test-build generator-test: generator-test-build
@echo "=== Запуск тестов генератора ===" @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 generator-test-cov: generator-test-build
@echo "=== Запуск тестов с покрытием ===" @echo "=== Запуск тестов с покрытием ==="
docker run --rm -v $(PWD):/workspace -w /workspace/generator generator:test pytest tests/ -v --cov=. --cov-report=term-missing 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
+10
View File
@@ -82,6 +82,16 @@ CHECK_LIVE_SEAM=0 make generated-history-check
make generated-history-runtime-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). по [runbook стартовой истории](./docs/runbooks/startup-history.md).
+4 -10
View File
@@ -30,6 +30,7 @@ from clickstream_generator.airflow_control import (
load_manifest_from_kafka, load_manifest_from_kafka,
run_backfill, run_backfill,
run_import, run_import,
target_dag_trigger_error,
validate_import_artifact, validate_import_artifact,
) )
from clickstream_generator.launch import PROFILES 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. fail_when_dag_is_paused, поэтому паузу проверяем заранее через DagModel.
""" """
dag_model = session.query(DagModel).filter(DagModel.dag_id == dag_id).one_or_none() dag_model = session.query(DagModel).filter(DagModel.dag_id == dag_id).one_or_none()
if dag_model is None: error = target_dag_trigger_error(dag_id, dag_model)
raise AirflowException( if error:
f"DAG {dag_id} ещё не найден Airflow. Подождите парсинга DAG-файлов " raise AirflowException(error)
"или перезапустите Airflow через make up."
)
if dag_model.is_paused:
raise AirflowException(
f"DAG {dag_id} стоит на паузе: снимите паузу в Airflow UI, "
"затем повторите generator_control."
)
with DAG( with DAG(
+11 -2
View File
@@ -18,11 +18,19 @@
- `make data` (архивный путь: заливает `data/*.jsonl` в Kafka; не основной источник аналитики) - `make data` (архивный путь: заливает `data/*.jsonl` в Kafka; не основной источник аналитики)
- `make transform` (запускает batch-процесс ODS -> DDS -> DM) - `make transform` (запускает batch-процесс ODS -> DDS -> DM)
- `make superset-init` (повторная инициализация Superset: подключение к ClickHouse, датасеты, дашборд) - `make superset-init` (повторная инициализация Superset: подключение к ClickHouse, датасеты, дашборд)
- `make test` (быстрый предкоммитный набор: тесты генератора, верхние
контрактные тесты и проверка compose-конфигурации)
- `make lint` (статические проверки: синтаксис Python и Bash, compose-конфиг,
пробелы в diff)
- `docker compose ps` - `docker compose ps`
- `docker compose logs -f --tail=200 <service>` - `docker compose logs -f --tail=200 <service>`
- `docker compose down` (сохраняет named volumes, включая `clickhouse-data`) - `docker compose down` (сохраняет named volumes, включая `clickhouse-data`)
- `docker compose down -v` (удаляет named volumes, использовать осознанно) - `docker compose down -v` (удаляет named volumes, использовать осознанно)
`make test` и `make lint` не заменяют стендовые проверки, которые управляют
volumes или live-генератором. Для стыка backfill/live отдельно запускайте
`make generated-history-runtime-check`.
## Порты ## Порты
Порты задаются в `docker-compose.yml`: Порты задаются в `docker-compose.yml`:
@@ -278,8 +286,9 @@ docker compose exec -T kafka /opt/kafka/bin/kafka-console-consumer.sh \
Live-продолжение стартует с этого state. Используйте тот же профиль или ту же Live-продолжение стартует с этого state. Используйте тот же профиль или ту же
длительность, что были у backfill. Команда сама выставит `GEN_STATE_RESET=false`. длительность, что были у backfill. Команда сама выставит `GEN_STATE_RESET=false`.
Если state записан старой версией генератора, запуск теперь падает громко: Если state записан старой версией генератора, запуск теперь падает громко:
сначала очистите стенд через `make clean` и пересоздайте историю либо явно сначала очистите стенд через `make clean` и пересоздайте историю. Команда
начните новый мир через `make generator-reset`. `make generator-reset` подходит только для осознанного нового live-мира на
чистом стенде: она тоже проверяет, что Kafka data-топики и STG пустые.
```bash ```bash
PROFILE=daily-wave make generator-continue PROFILE=daily-wave make generator-continue
+1 -1
View File
@@ -145,7 +145,7 @@ make logs service=superset # Логи сервиса
# ETL # ETL
make generated-history-analytics # Чистый прогон генерации до Superset 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 ddl # Применение DDL в ClickHouse
make data # Архивная загрузка data/*.jsonl в Kafka make data # Архивная загрузка data/*.jsonl в Kafka
make transform # Запуск batch-процесса make transform # Запуск batch-процесса
+3 -1
View File
@@ -67,7 +67,9 @@ make up
- `generator_state` — слепок состояния генератора на правой границе стартовой истории. - `generator_state` — слепок состояния генератора на правой границе стартовой истории.
По нему live-продолжение понимает, откуда продолжать тот же мир; По нему live-продолжение понимает, откуда продолжать тот же мир;
- `generator_startup_history_manifest` — паспорт стартовой истории: seed, границы - `generator_startup_history_manifest` — паспорт стартовой истории: seed, границы
модельного времени и контрольные числа. модельного времени и контрольные числа;
- `generator_batch_history` — журнал батчей генератора: сколько сообщений он
отправил и чем закончилась каждая пачка записи.
Служебные топики нужны стенду, но в упражнениях курса мы их не меняем. У каждого нашего Служебные топики нужны стенду, но в упражнениях курса мы их не меняем. У каждого нашего
топика событий в колонке с партициями стоит **1**: топик маленький, делить не на что. топика событий в колонке с партициями стоит **1**: топик маленький, делить не на что.
+11 -6
View File
@@ -40,7 +40,9 @@
1. Поднимите стенд: `make up`. 1. Поднимите стенд: `make up`.
2. Если DDL ещё не применён, запустите `ddl_init`. 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. Долгий простой стенда создаёт большую дыру в У `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 есть, но настройки не совпадают, генератор падает с перечнем Если читаемый state есть, но настройки не совпадают, генератор падает с перечнем
полей. Это защита от смешения разных миров. Для намеренного нового мира полей. Это защита от смешения разных миров. Для намеренного нового мира
используйте `make generator-reset` или `make clean`. сначала очистите стенд через `make clean`, затем запускайте нужный сценарий.
После обновления кода старый state может оказаться в старом формате. При После обновления кода старый state может оказаться в старом формате. При
`GEN_STATE_RESET=false` это теперь громкий отказ, а не тихий старт с нуля поверх `GEN_STATE_RESET=false` это теперь громкий отказ, а не тихий старт с нуля поверх
старой истории. Оператору нужно выбрать одно из двух: очистить стенд через старой истории. Оператору нужно выбрать одно из двух: очистить стенд через
`make clean` и заново создать стартовую историю, либо осознанно начать новый мир `make clean` и заново создать стартовую историю, либо осознанно начать новый
через `GEN_STATE_RESET=true` / `make generator-reset`. live-мир через `make generator-reset` на чистом стенде.
После нестандартного мира `make generator-continue` нужно запускать с теми же После нестандартного мира `make generator-continue` нужно запускать с теми же
настройками, что были у backfill/import. При расхождении генератор громко настройками, что были у backfill/import. При расхождении генератор громко
+16 -9
View File
@@ -176,24 +176,25 @@ make generator-backfill
# Продолжить live-поток из state # Продолжить live-поток из state
make generator-continue make generator-continue
# Начать live-поток как новый мир # Начать live-поток как новый мир на чистом стенде
make generator-reset make generator-reset
# Чистый аналитический прогон всего стенда # Чистый аналитический прогон всего стенда
make generated-history-analytics make generated-history-analytics
# Запуск тестов # Быстрые проверки перед коммитом
make generator-test make test
make lint
``` ```
Основной ручной путь запуска — глаголы `generator-backfill`, `generator-continue` Основной ручной путь запуска — глаголы `generator-backfill`, `generator-continue`
и `generator-reset`. Старые `GEN_RUN_MODE`, `GEN_STATE_RESET` и и `generator-reset`. Старые `GEN_RUN_MODE`, `GEN_STATE_RESET` и
`GEN_MODEL_T_END` остаются низкоуровневым способом для отладки. `GEN_MODEL_T_END` остаются низкоуровневым способом для отладки.
`generator-backfill` и `startup-history-import` перед записью проверяют, что `generator-backfill`, `startup-history-import` и `generator-reset` перед записью
Kafka data-топики и STG пустые. `generator-continue` продолжает только проверяют, что Kafka data-топики и STG пустые. `generator-continue` продолжает
совместимый state; старый формат state при `GEN_STATE_RESET=false` даёт отказ с только совместимый state; старый формат state при `GEN_STATE_RESET=false` даёт
подсказкой очистить стенд или явно начать новый мир. отказ с подсказкой очистить стенд или явно начать новый мир на чистом стенде.
## Метрики Prometheus ## Метрики Prometheus
@@ -351,16 +352,22 @@ GEN_STATE_ENABLED=false docker compose up -d generator
```bash ```bash
# Через Makefile (рекомендуется) # Через Makefile (рекомендуется)
make generator-test make test
make lint
# Вручную через Docker # Вручную через Docker
docker build -t generator:test . 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 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` запускаются отдельно.
### Структура тестов ### Структура тестов
``` ```
@@ -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: def default_artifact_path(operation: str) -> str:
"""Возвращает путь артефакта по умолчанию в общем томе data.""" """Возвращает путь артефакта по умолчанию в общем томе data."""
filename = ( filename = (
+4 -9
View File
@@ -11,7 +11,6 @@ logger = logging.getLogger("generator")
STATE_VERSION = "3.0" STATE_VERSION = "3.0"
KNOWN_OLD_STATE_VERSIONS = {"2.0"}
class UnsupportedStateVersionError(ValueError): class UnsupportedStateVersionError(ValueError):
@@ -245,15 +244,9 @@ class GeneratorState:
try: try:
if not isinstance(data, dict): if not isinstance(data, dict):
raise ValueError("state must be an object") 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 version != STATE_VERSION:
if str(version) in KNOWN_OLD_STATE_VERSIONS: raise UnsupportedStateVersionError(version)
raise UnsupportedStateVersionError(str(version))
logger.warning(
"Unsupported generator state version %s, will start fresh",
version,
)
raise ValueError(f"unsupported state version: {version}")
_validate_v3_payload(data) _validate_v3_payload(data)
model_timestamp = _parse_aware_utc( model_timestamp = _parse_aware_utc(
data["model_timestamp"], data["model_timestamp"],
@@ -309,5 +302,7 @@ class GeneratorState:
"""Безопасная загрузка state с graceful degradation.""" """Безопасная загрузка state с graceful degradation."""
try: try:
return cls.from_dict(data) return cls.from_dict(data)
except UnsupportedStateVersionError:
raise
except Exception: except Exception:
return None return None
+16
View File
@@ -5,6 +5,7 @@
from dataclasses import replace from dataclasses import replace
from datetime import datetime, timezone from datetime import datetime, timezone
from pathlib import Path from pathlib import Path
from types import SimpleNamespace
import pytest import pytest
@@ -49,6 +50,21 @@ def test_import_env_requires_artifact_path_and_uses_backfill_contract():
assert env["GEN_HISTORY_DURATION"] == "2d" 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(): def test_assert_stand_clean_rejects_non_empty_stg_before_writes():
"""Backfill/import не стартуют на непустом STG.""" """Backfill/import не стартуют на непустом STG."""
from clickstream_generator.airflow_control import assert_stg_tables_empty from clickstream_generator.airflow_control import assert_stg_tables_empty
@@ -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 'op_kwargs={"dag_id": ETL_DAG_ID}' in text
assert "session.query(DagModel)" in text assert "session.query(DagModel)" in text
assert "DagModel.dag_id == dag_id" 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 "Airflow 2.10.5" in text
assert "fail_when_dag_is_paused" 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_backfill"' in text
assert 'return "check_etl_not_paused_before_import"' in text assert 'return "check_etl_not_paused_before_import"' in text
assert "check_etl_not_paused_before_backfill >> precheck_backfill_task" 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(): def test_clean_start_paths_include_live_generator_profile():
"""Все чистые сбросы видят профильный live-генератор.""" """Все чистые сбросы видят профильный live-генератор."""
generated_history = ( generated_history = (
+27 -7
View File
@@ -432,6 +432,22 @@ class TestGeneratorStateValidation:
with pytest.raises(UnsupportedStateVersionError, match="2.0.*3.0"): with pytest.raises(UnsupportedStateVersionError, match="2.0.*3.0"):
GeneratorState.from_dict(data) 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: class TestKafkaStateManager:
"""Тесты менеджера состояния.""" """Тесты менеджера состояния."""
@@ -536,10 +552,14 @@ class TestKafkaStateManager:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
# Мокаем consumer с невалидным сообщением # Мокаем consumer с невалидным сообщением текущей версии.
mock_message = MagicMock() mock_message = MagicMock()
mock_message.key = b"default" 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 = MagicMock()
mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message]))
@@ -749,8 +769,8 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_version_1_state_returns_none_with_warning(self, caplog): def test_load_version_1_state_fails_loudly(self, caplog):
"""Старое state v1 не восстанавливается и даёт чистый старт.""" """Старое state v1 не скрывается за чистым стартом."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
@@ -774,10 +794,10 @@ class TestKafkaStateManager:
mock_consumer_class.return_value = mock_consumer mock_consumer_class.return_value = mock_consumer
manager = KafkaStateManager("kafka:29092") manager = KafkaStateManager("kafka:29092")
with caplog.at_level(logging.WARNING, logger="generator"): with caplog.at_level(logging.ERROR, logger="generator"):
result = manager.load() with pytest.raises(UnsupportedStateVersionError, match="1.0.*3.0"):
manager.load()
assert result is None
assert "version" in caplog.text assert "version" in caplog.text
def test_load_ignores_wrong_key(self): def test_load_ignores_wrong_key(self):
+6 -1
View File
@@ -520,12 +520,17 @@ FROM
ORDER BY source ORDER BY source
FORMAT TabSeparated" 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 while IFS=$'\t' read -r source rows; do
[[ -n "${source}" ]] || continue [[ -n "${source}" ]] || continue
[[ "${rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число строк для ${source}" [[ "${rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать число строк для ${source}"
(( rows > 0 )) || fail "${source} пустая" (( rows > 0 )) || fail "${source} пустая"
echo "${source}=${rows}" echo "${source}=${rows}"
done < <(ch_query "${views_query}") done <<< "${views_rows}"
if [[ "${REQUIRE_SUPERSET}" != "1" ]]; then if [[ "${REQUIRE_SUPERSET}" != "1" ]]; then
echo "" echo ""
+7 -1
View File
@@ -87,6 +87,12 @@ elif [[ "${VERB}" == "continue" ]]; then
echo "Шаг 3: запуск live-генератора" echo "Шаг 3: запуск live-генератора"
${COMPOSE_BIN} up -d --build generator ${COMPOSE_BIN} up -d --build generator
else 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 ${COMPOSE_BIN} up -d --build generator
fi fi
@@ -1,7 +1,110 @@
from pathlib import Path from pathlib import Path
import os
import subprocess
import textwrap
REPO_ROOT = Path(__file__).resolve().parents[1] 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(): 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 "kafka-manifest-summary" in script
assert "manifest_model_t_end" in script assert "manifest_model_t_end" in script
assert "CH_MODEL_T_END=\"$(clickhouse_datetime_literal \"${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 "GEN_LAUNCH_PROFILE" in makefile
assert "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 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(): def test_clean_generated_history_run_explicitly_skips_live_seam():
"""Чистый backfill без live-продолжения отключает проверку стыка явно.""" """Чистый backfill без live-продолжения отключает проверку стыка явно."""
script = (REPO_ROOT / "scripts" / "run_generated_history_analytics.sh").read_text( script = (REPO_ROOT / "scripts" / "run_generated_history_analytics.sh").read_text(