feat(generator): добавлены профили запуска
- Зачем: - запуск генератора должен быть понятным перед будущим DAG-пультом. - Что: - добавлены глаголы запуска backfill, continue и reset. - добавлены профили ci и daily-wave с расчётом длительности истории. - обновлены runbook и документы запуска под профильный интерфейс. - Проверка: - uv run --with-requirements generator/requirements.txt pytest generator/tests -q. - bash -n scripts/run_generator.sh scripts/export_startup_history_artifact.sh scripts/import_startup_history_artifact.sh scripts/run_generated_history_analytics.sh. - PROFILE=daily-wave COMPOSE_BIN=true bash scripts/run_generator.sh backfill.
This commit is contained in:
+25
-6
@@ -78,6 +78,7 @@ generator-service -> Kafka topics -> (потребители отдельно)
|
||||
| `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` |
|
||||
| `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` |
|
||||
| `GEN_RUN_MODE` | Режим генератора | `live` |
|
||||
| `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` |
|
||||
| `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | — |
|
||||
| `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` |
|
||||
| `GEN_SEED` | Сид для воспроизводимости | — |
|
||||
@@ -125,12 +126,17 @@ make generated-history-analytics
|
||||
проверяет Superset metadata. Файлы `data/*.jsonl` при этом не грузятся в Kafka:
|
||||
они пока используются только как фактура для генератора.
|
||||
|
||||
По умолчанию используется быстрый проверочный профиль на 6 часов модельного
|
||||
времени (`GEN_MODEL_T_END=2026-01-01T06:00:00+00:00`). Суточный прогон доступен
|
||||
явно:
|
||||
По умолчанию используется быстрый профиль `ci`: 6 часов модельного времени.
|
||||
Профиль `daily-wave` даёт 2 суток, чтобы была видна суточная волна:
|
||||
|
||||
```bash
|
||||
GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics
|
||||
PROFILE=daily-wave make generated-history-analytics
|
||||
```
|
||||
|
||||
Разовую длительность можно задать без ручного расчёта `GEN_MODEL_T_END`:
|
||||
|
||||
```bash
|
||||
GEN_HISTORY_DURATION=2d make generated-history-analytics
|
||||
```
|
||||
|
||||
### Режим "раз в минуту" (для демо)
|
||||
@@ -162,13 +168,26 @@ make generator-logs
|
||||
# Перезапуск с пересборкой
|
||||
make generator-restart
|
||||
|
||||
# Запуск тестов
|
||||
make generator-test
|
||||
# Промотать стартовую историю
|
||||
make generator-backfill
|
||||
|
||||
# Продолжить live-поток из state
|
||||
make generator-continue
|
||||
|
||||
# Начать live-поток как новый мир
|
||||
make generator-reset
|
||||
|
||||
# Чистый аналитический прогон всего стенда
|
||||
make generated-history-analytics
|
||||
|
||||
# Запуск тестов
|
||||
make generator-test
|
||||
```
|
||||
|
||||
Основной ручной путь запуска — глаголы `generator-backfill`, `generator-continue`
|
||||
и `generator-reset`. Старые `GEN_RUN_MODE`, `GEN_STATE_RESET` и
|
||||
`GEN_MODEL_T_END` остаются низкоуровневым способом для отладки.
|
||||
|
||||
## Метрики Prometheus
|
||||
|
||||
Генератор экспортирует метрики на `:9109/metrics`:
|
||||
|
||||
@@ -97,6 +97,9 @@ class Config:
|
||||
if os.getenv("GEN_STARTUP_HISTORY_ARTIFACT")
|
||||
else None
|
||||
)
|
||||
launch_profile: str = field(
|
||||
default_factory=lambda: os.getenv("GEN_LAUNCH_PROFILE", "ci")
|
||||
)
|
||||
|
||||
def __post_init__(self):
|
||||
if self.tick_seconds < 1:
|
||||
|
||||
@@ -0,0 +1,175 @@
|
||||
"""Глаголы и профили запуска генератора."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import os
|
||||
import shlex
|
||||
import sys
|
||||
from dataclasses import dataclass
|
||||
from datetime import datetime, timedelta, timezone
|
||||
|
||||
|
||||
LAUNCH_ENV_KEYS = (
|
||||
"GEN_SEED",
|
||||
"GEN_MODEL_T0",
|
||||
"GEN_MODEL_TIMEZONE",
|
||||
"GEN_MODEL_TIME_SPEED",
|
||||
"GEN_TICK_SECONDS",
|
||||
"GEN_LAMBDA_BASE_PER_MIN",
|
||||
"GEN_JITTER_PCT",
|
||||
"GEN_MIN_EVENTS_PER_TICK",
|
||||
"GEN_MAX_EVENTS_PER_TICK",
|
||||
)
|
||||
|
||||
|
||||
@dataclass(frozen=True)
|
||||
class LaunchProfile:
|
||||
"""Именованный набор настроек запуска."""
|
||||
|
||||
duration: str
|
||||
env: dict[str, str]
|
||||
|
||||
|
||||
PROFILES = {
|
||||
"ci": LaunchProfile(
|
||||
duration="6h",
|
||||
env={
|
||||
"GEN_SEED": "4242",
|
||||
"GEN_MODEL_T0": "2026-01-01T00:00:00+00:00",
|
||||
"GEN_MODEL_TIMEZONE": "UTC",
|
||||
"GEN_MODEL_TIME_SPEED": "1",
|
||||
"GEN_TICK_SECONDS": "60",
|
||||
"GEN_LAMBDA_BASE_PER_MIN": "60",
|
||||
"GEN_JITTER_PCT": "0",
|
||||
"GEN_MIN_EVENTS_PER_TICK": "1",
|
||||
"GEN_MAX_EVENTS_PER_TICK": "1000",
|
||||
},
|
||||
),
|
||||
"daily-wave": LaunchProfile(
|
||||
duration="2d",
|
||||
env={
|
||||
"GEN_SEED": "4242",
|
||||
"GEN_MODEL_T0": "2026-01-01T00:00:00+00:00",
|
||||
"GEN_MODEL_TIMEZONE": "UTC",
|
||||
"GEN_MODEL_TIME_SPEED": "1",
|
||||
"GEN_TICK_SECONDS": "60",
|
||||
"GEN_LAMBDA_BASE_PER_MIN": "60",
|
||||
"GEN_JITTER_PCT": "0",
|
||||
"GEN_MIN_EVENTS_PER_TICK": "1",
|
||||
"GEN_MAX_EVENTS_PER_TICK": "1000",
|
||||
},
|
||||
),
|
||||
}
|
||||
|
||||
|
||||
def parse_duration(value: str) -> timedelta:
|
||||
"""Разбирает короткую длительность: 30m, 6h, 2d."""
|
||||
if len(value) < 2:
|
||||
raise ValueError("Duration must look like 6h or 2d")
|
||||
amount_text = value[:-1]
|
||||
unit = value[-1]
|
||||
if not amount_text.isdigit():
|
||||
raise ValueError("Duration amount must be a positive integer")
|
||||
amount = int(amount_text)
|
||||
if amount <= 0:
|
||||
raise ValueError("Duration amount must be > 0")
|
||||
if unit == "s":
|
||||
return timedelta(seconds=amount)
|
||||
if unit == "m":
|
||||
return timedelta(minutes=amount)
|
||||
if unit == "h":
|
||||
return timedelta(hours=amount)
|
||||
if unit == "d":
|
||||
return timedelta(days=amount)
|
||||
raise ValueError("Duration unit must be one of: s, m, h, d")
|
||||
|
||||
|
||||
def build_launch_env(
|
||||
verb: str,
|
||||
*,
|
||||
profile_name: str = "ci",
|
||||
duration: str | None = None,
|
||||
overrides: dict[str, str] | None = None,
|
||||
) -> dict[str, str]:
|
||||
"""Возвращает старые env-переменные для нового глагола запуска."""
|
||||
if verb not in {"backfill", "continue", "reset"}:
|
||||
raise ValueError("verb must be backfill, continue or reset")
|
||||
if profile_name not in PROFILES:
|
||||
known = ", ".join(sorted(PROFILES))
|
||||
raise ValueError(f"Unknown launch profile {profile_name!r}; known: {known}")
|
||||
|
||||
profile = PROFILES[profile_name]
|
||||
overrides = overrides or {}
|
||||
env = dict(profile.env)
|
||||
for key in LAUNCH_ENV_KEYS:
|
||||
if key in overrides and overrides[key] != "":
|
||||
env[key] = overrides[key]
|
||||
|
||||
selected_duration = (
|
||||
duration
|
||||
or overrides.get("GEN_HISTORY_DURATION")
|
||||
or profile.duration
|
||||
)
|
||||
env["GEN_LAUNCH_VERB"] = verb
|
||||
env["GEN_LAUNCH_PROFILE"] = profile_name
|
||||
|
||||
if verb == "backfill":
|
||||
env["GEN_RUN_MODE"] = "backfill"
|
||||
env["GEN_STATE_RESET"] = "true"
|
||||
env["GEN_HISTORY_DURATION"] = selected_duration
|
||||
env["GEN_MODEL_T_END"] = _model_t_end(env["GEN_MODEL_T0"], selected_duration)
|
||||
elif verb == "continue":
|
||||
env["GEN_RUN_MODE"] = "live"
|
||||
env["GEN_STATE_RESET"] = "false"
|
||||
env["GEN_HISTORY_DURATION"] = selected_duration
|
||||
env["GEN_MODEL_T_END"] = _model_t_end(env["GEN_MODEL_T0"], selected_duration)
|
||||
else:
|
||||
env["GEN_RUN_MODE"] = "live"
|
||||
env["GEN_STATE_RESET"] = "true"
|
||||
|
||||
return env
|
||||
|
||||
|
||||
def shell_assignments(env: dict[str, str]) -> str:
|
||||
"""Печатает env в виде, пригодном для eval в Bash."""
|
||||
return "\n".join(
|
||||
f"export {key}={shlex.quote(value)}"
|
||||
for key, value in sorted(env.items())
|
||||
)
|
||||
|
||||
|
||||
def _model_t_end(model_t0: str, duration: str) -> str:
|
||||
t0 = datetime.fromisoformat(model_t0.replace("Z", "+00:00"))
|
||||
if t0.tzinfo is None:
|
||||
raise ValueError("GEN_MODEL_T0 must include timezone")
|
||||
return (t0.astimezone(timezone.utc) + parse_duration(duration)).isoformat()
|
||||
|
||||
|
||||
def _parse_args(argv: list[str]) -> argparse.Namespace:
|
||||
parser = argparse.ArgumentParser(description="Глаголы запуска генератора")
|
||||
parser.add_argument("verb", choices=("backfill", "continue", "reset"))
|
||||
parser.add_argument(
|
||||
"--profile",
|
||||
default=os.getenv("PROFILE") or os.getenv("GEN_LAUNCH_PROFILE") or "ci",
|
||||
choices=sorted(PROFILES),
|
||||
)
|
||||
parser.add_argument("--duration", default=os.getenv("GEN_HISTORY_DURATION"))
|
||||
return parser.parse_args(argv)
|
||||
|
||||
|
||||
def main(argv: list[str] | None = None) -> int:
|
||||
"""CLI для Bash-скриптов запуска."""
|
||||
args = _parse_args(sys.argv[1:] if argv is None else argv)
|
||||
env = build_launch_env(
|
||||
args.verb,
|
||||
profile_name=args.profile,
|
||||
duration=args.duration,
|
||||
overrides=os.environ,
|
||||
)
|
||||
print(shell_assignments(env))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
raise SystemExit(main())
|
||||
@@ -70,6 +70,7 @@ class GeneratorService:
|
||||
f"Configuration: tick={self.config.tick_seconds}s, "
|
||||
f"lambda_base={self.config.lambda_base_per_min}/min, "
|
||||
f"jitter={self.config.jitter_pct}%, "
|
||||
f"launch_profile={self.config.launch_profile}, "
|
||||
f"state_enabled={self.config.state_enabled}, "
|
||||
f"state_reset={self.config.state_reset}"
|
||||
)
|
||||
|
||||
@@ -299,6 +299,7 @@ def build_manifest(config: Config, counters: ManifestCounters, state: GeneratorS
|
||||
"model_t_end": config.model_t_end.isoformat(),
|
||||
"model_timezone": config.model_timezone,
|
||||
"run_mode": "backfill",
|
||||
"launch_profile": config.launch_profile,
|
||||
"generation_settings": generation_settings_from_config(config),
|
||||
"state_version": state.version,
|
||||
"state": {
|
||||
|
||||
@@ -118,3 +118,21 @@ class TestConfigDefaults:
|
||||
finally:
|
||||
if orig_value is not None:
|
||||
os.environ["GEN_TICK_SECONDS"] = orig_value
|
||||
|
||||
def test_launch_profile_defaults_to_ci(self, monkeypatch, data_dir):
|
||||
"""По умолчанию выбран быстрый профиль запуска ci."""
|
||||
monkeypatch.setenv("GEN_DATA_DIR", str(data_dir))
|
||||
monkeypatch.delenv("GEN_LAUNCH_PROFILE", raising=False)
|
||||
|
||||
config = Config()
|
||||
|
||||
assert config.launch_profile == "ci"
|
||||
|
||||
def test_launch_profile_can_be_read_from_env(self, monkeypatch, data_dir):
|
||||
"""Выбранный профиль запуска попадает в конфигурацию для логов."""
|
||||
monkeypatch.setenv("GEN_DATA_DIR", str(data_dir))
|
||||
monkeypatch.setenv("GEN_LAUNCH_PROFILE", "daily-wave")
|
||||
|
||||
config = Config()
|
||||
|
||||
assert config.launch_profile == "daily-wave"
|
||||
|
||||
@@ -0,0 +1,66 @@
|
||||
"""
|
||||
Тесты профилей и глаголов запуска генератора.
|
||||
"""
|
||||
|
||||
from datetime import datetime, timezone
|
||||
|
||||
|
||||
def test_duration_2d_sets_model_t_end_from_t0():
|
||||
"""Длительность 2d считается от T0 без ручного GEN_MODEL_T_END."""
|
||||
from clickstream_generator.launch import build_launch_env
|
||||
|
||||
env = build_launch_env("backfill", profile_name="ci", duration="2d")
|
||||
|
||||
assert env["GEN_MODEL_T0"] == "2026-01-01T00:00:00+00:00"
|
||||
assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00"
|
||||
|
||||
|
||||
def test_daily_wave_profile_uses_two_days_by_default():
|
||||
"""Профиль daily-wave даёт историю с суточной волной."""
|
||||
from clickstream_generator.launch import build_launch_env
|
||||
|
||||
env = build_launch_env("backfill", profile_name="daily-wave")
|
||||
|
||||
assert env["GEN_LAUNCH_PROFILE"] == "daily-wave"
|
||||
assert env["GEN_HISTORY_DURATION"] == "2d"
|
||||
assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00"
|
||||
|
||||
|
||||
def test_backfill_verb_maps_to_existing_low_level_flags():
|
||||
"""Глагол backfill оставляет прежнюю механику запуска под капотом."""
|
||||
from clickstream_generator.launch import build_launch_env
|
||||
|
||||
env = build_launch_env("backfill", profile_name="ci")
|
||||
|
||||
assert env["GEN_RUN_MODE"] == "backfill"
|
||||
assert env["GEN_STATE_RESET"] == "true"
|
||||
|
||||
|
||||
def test_continue_and_reset_verbs_map_to_existing_low_level_flags():
|
||||
"""Глаголы continue и reset задают только старый контракт state."""
|
||||
from clickstream_generator.launch import build_launch_env
|
||||
|
||||
continue_env = build_launch_env("continue", profile_name="ci")
|
||||
reset_env = build_launch_env("reset", profile_name="ci")
|
||||
|
||||
assert continue_env["GEN_RUN_MODE"] == "live"
|
||||
assert continue_env["GEN_STATE_RESET"] == "false"
|
||||
assert continue_env["GEN_MODEL_T_END"] == "2026-01-01T06:00:00+00:00"
|
||||
assert reset_env["GEN_RUN_MODE"] == "live"
|
||||
assert reset_env["GEN_STATE_RESET"] == "true"
|
||||
assert "GEN_MODEL_T_END" not in reset_env
|
||||
|
||||
|
||||
def test_ci_profile_keeps_previous_default_backfill_window():
|
||||
"""CI-профиль сохраняет прежний 6-часовой проверочный запуск."""
|
||||
from clickstream_generator.launch import build_launch_env
|
||||
|
||||
env = build_launch_env("backfill", profile_name="ci")
|
||||
|
||||
t0 = datetime.fromisoformat(env["GEN_MODEL_T0"])
|
||||
t_end = datetime.fromisoformat(env["GEN_MODEL_T_END"])
|
||||
assert t0 == datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc)
|
||||
assert (t_end - t0).total_seconds() == 6 * 60 * 60
|
||||
assert env["GEN_TICK_SECONDS"] == "60"
|
||||
assert env["GEN_LAMBDA_BASE_PER_MIN"] == "60"
|
||||
assert env["GEN_JITTER_PCT"] == "0"
|
||||
@@ -130,6 +130,43 @@ def test_artifact_roundtrip_keeps_events_state_and_manifest(base_config):
|
||||
path.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def test_manifest_and_artifact_show_launch_profile(base_config):
|
||||
"""Manifest и файл артефакта показывают выбранный профиль запуска."""
|
||||
from clickstream_generator.startup_history_artifact import (
|
||||
StartupHistoryArtifactBuilder,
|
||||
build_manifest,
|
||||
load_startup_history_artifact,
|
||||
write_startup_history_artifact,
|
||||
)
|
||||
|
||||
state = _state()
|
||||
builder = StartupHistoryArtifactBuilder()
|
||||
builder.add_batch(_batch())
|
||||
manifest = build_manifest(
|
||||
config=replace(
|
||||
base_config,
|
||||
model_t0=state.model_t0,
|
||||
model_t_end=state.model_timestamp,
|
||||
model_time_speed=1,
|
||||
model_timezone="UTC",
|
||||
seed=42,
|
||||
launch_profile="daily-wave",
|
||||
),
|
||||
counters=builder.counters,
|
||||
state=state,
|
||||
)
|
||||
artifact = builder.to_artifact(manifest=manifest, state=state)
|
||||
|
||||
path = base_config.data_dir.parent / "tmp-startup-history-profile.json"
|
||||
try:
|
||||
write_startup_history_artifact(path, artifact)
|
||||
loaded = load_startup_history_artifact(path)
|
||||
assert manifest["launch_profile"] == "daily-wave"
|
||||
assert loaded["manifest"]["launch_profile"] == "daily-wave"
|
||||
finally:
|
||||
path.unlink(missing_ok=True)
|
||||
|
||||
|
||||
def test_artifact_validation_rejects_mismatched_manifest(base_config):
|
||||
"""Валидация отвергает артефакт, где manifest не совпадает с событиями."""
|
||||
from clickstream_generator.startup_history_artifact import (
|
||||
|
||||
Reference in New Issue
Block a user