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:
+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):
|
||||
|
||||
Reference in New Issue
Block a user