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:
@@ -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 (закрыто):
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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).
|
||||
|
||||
|
||||
@@ -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(
|
||||
|
||||
+11
-2
@@ -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 <service>`
|
||||
- `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
|
||||
|
||||
@@ -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-процесса
|
||||
|
||||
@@ -67,7 +67,9 @@ make up
|
||||
- `generator_state` — слепок состояния генератора на правой границе стартовой истории.
|
||||
По нему live-продолжение понимает, откуда продолжать тот же мир;
|
||||
- `generator_startup_history_manifest` — паспорт стартовой истории: seed, границы
|
||||
модельного времени и контрольные числа.
|
||||
модельного времени и контрольные числа;
|
||||
- `generator_batch_history` — журнал батчей генератора: сколько сообщений он
|
||||
отправил и чем закончилась каждая пачка записи.
|
||||
|
||||
Служебные топики нужны стенду, но в упражнениях курса мы их не меняем. У каждого нашего
|
||||
топика событий в колонке с партициями стоит **1**: топик маленький, делить не на что.
|
||||
|
||||
@@ -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. При расхождении генератор громко
|
||||
|
||||
+16
-9
@@ -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` запускаются отдельно.
|
||||
|
||||
### Структура тестов
|
||||
|
||||
```
|
||||
|
||||
@@ -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 = (
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 = (
|
||||
|
||||
@@ -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):
|
||||
|
||||
@@ -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 ""
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
Reference in New Issue
Block a user