diff --git a/.scratch/generator-model-time-startup-history/issues/07-startup-history-portable-artifact-and-usage-docs.md b/.scratch/generator-model-time-startup-history/issues/07-startup-history-portable-artifact-and-usage-docs.md index ec07139..733b5b9 100644 --- a/.scratch/generator-model-time-startup-history/issues/07-startup-history-portable-artifact-and-usage-docs.md +++ b/.scratch/generator-model-time-startup-history/issues/07-startup-history-portable-artifact-and-usage-docs.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: done # Портативный артефакт стартовой истории и runbook по стенду diff --git a/Makefile b/Makefile index 02cbd5f..7fb09e0 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,6 @@ .PHONY: up down clean ddl data transform logs \ generated-history-analytics generated-history-check \ + startup-history-export startup-history-import startup-history-check \ reload-monitoring recover-monitoring \ superset-init superset-dashboard superset-ui superset-restart \ generator-up generator-down generator-logs generator-restart \ @@ -46,6 +47,18 @@ generated-history-analytics: generated-history-check: COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh +# Сгенерировать стартовую историю и сохранить портативный артефакт +startup-history-export: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/export_startup_history_artifact.sh + +# Воспроизвести портативный артефакт стартовой истории в Kafka +startup-history-import: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/import_startup_history_artifact.sh + +# Сверить ClickHouse после импорта с manifest артефакта +startup-history-check: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_startup_history_manifest.sh + # Перезагрузка конфигурации мониторинга (после изменений в provisioning) reload-monitoring: @echo "=== Перезагрузка сервисов мониторинга ===" diff --git a/README.md b/README.md index 5a77ee7..abbb004 100644 --- a/README.md +++ b/README.md @@ -50,6 +50,9 @@ GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics make generated-history-check ``` +Сохранить стартовую историю в файл и восстановить её без новой генерации можно +по [runbook стартовой истории](./docs/runbooks/startup-history.md). + Проверить, что данные дошли до витрин: ```bash @@ -116,6 +119,8 @@ flowchart LR обоснование решений. - [Запуск и эксплуатация](./docs/OPERATIONS.md) — сценарий запуска, параметры DAG-ов, мониторинг, частые проблемы. +- [Runbook стартовой истории](./docs/runbooks/startup-history.md) — экспорт, + импорт и live-продолжение из готового артефакта. - [Карта репозитория](./docs/REPO_MAP.md) — где какие файлы и что менять. - [Курс «Кликстрим на ClickHouse»](./docs/course/README.md) — учебная программа на этом стенде. diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index ac0b235..200dc5d 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -118,6 +118,7 @@ make generator-logs | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | | `GEN_RUN_MODE` | Режим генератора | `live` | +| `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | пусто | | `GEN_STATE_ENABLED` | Сохранять state v2 между рестартами | `true` | | `GEN_STATE_RESET` | Сбросить state при старте | `false` | @@ -163,6 +164,9 @@ GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 make generated-history-analytics контрольными числами. При live-запуске с теми же настройками генератор видит, что state совпадает с manifest, и стартует ровно с `T_end` без настенной дельты. +Если историю нужно сохранить в файл и восстановить на чистом стенде без новой +генерации, используйте [runbook стартовой истории](./runbooks/startup-history.md). + Для чистого повтора пересоздавайте volumes. Это сбрасывает ClickHouse, Kafka-топики данных, state и manifest генератора. `make generated-history-analytics` делает это по умолчанию (`CLEAN_START=1`). diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md new file mode 100644 index 0000000..a386db1 --- /dev/null +++ b/docs/runbooks/startup-history.md @@ -0,0 +1,91 @@ +# Runbook: стартовая история стенда + +Этот runbook нужен, чтобы один раз создать стартовую историю генератора, сохранить +её в файл и быстро восстановить на чистом стенде. + +За устройством генератора см. [`generator/README.md`](../../generator/README.md) +и [спеку модельного времени](../specs/2026-06-14-generator-model-time-and-startup-history.md). + +## Что дёшево + +- Перезапуск без очистки volumes: Kafka хранит `generator_state` и manifest. +- Live-продолжение после импорта: генератор стартует с `T_end`, если настройки + совпадают. +- Восстановление чистого стенда из готового файла: события повторно пишутся в + Kafka, ClickHouse наполняется штатным путём. + +## Что требует нового артефакта + +- Другая длительность истории (`GEN_MODEL_T_END`). +- Другой `GEN_SEED`, `GEN_MODEL_T0`, часовой пояс, скорость или настройки + генерации. +- Осознанный новый мир после несовместимого state: сначала сбросьте state через + `GEN_STATE_RESET=true` или чистые volumes. + +## Экспорт + +По умолчанию создаётся быстрый 6-часовой артефакт: + +```bash +make startup-history-export +``` + +Файл по умолчанию: `/tmp/clickstream-startup-history.json`. + +Суточная история: + +```bash +ARTIFACT=/tmp/clickstream-startup-history-1d.json \ +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 \ +make startup-history-export +``` + +Команда делает чистый backfill и пишет в файл один связный набор: события Kafka, +state и manifest. + +## Импорт на чистый стенд + +```bash +make clean +docker compose up -d clickhouse kafka +make ddl + +ARTIFACT=/tmp/clickstream-startup-history.json make startup-history-import + +sleep 10 +make transform +ARTIFACT=/tmp/clickstream-startup-history.json make startup-history-check +make generated-history-check +``` + +Импорт не пишет напрямую в ClickHouse. Он воспроизводит события и служебные +compact-топики в Kafka. ClickHouse читает данные через свои Kafka-таблицы и +Materialized View, затем batch строит ODS, DDS и DM. + +`make startup-history-check` сверяет контрольные числа DM-витрины с manifest +артефакта: события, визиты, пользователей и диапазон `event_timestamp`. Если +data-топики Kafka уже непустые, импорт остановится до публикации событий. + +В текущем стеке `kafka-python` не даёт транзакционный producer для нескольких +топиков. Поэтому импорт остаётся clean-stand операцией: при ошибке записи он +удаляет import-топики Kafka, чтобы повторный импорт не дописал дубли. Если +ClickHouse уже успел прочитать частичные сообщения, очистите стенд через +`make clean` и повторите импорт. + +## Live-продолжение + +После импорта запускайте live с теми же настройками, что были в артефакте: + +```bash +GEN_STATE_RESET=false \ +GEN_SEED=4242 \ +GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ +GEN_MODEL_T_END=2026-01-01T06:00:00+00:00 \ +GEN_MODEL_TIMEZONE=UTC \ +GEN_MODEL_TIME_SPEED=1 \ +docker compose up -d generator +``` + +Если читаемый state есть, но настройки не совпадают, генератор падает с перечнем +полей. Это защита от смешения разных миров. Для намеренного нового мира +используйте `GEN_STATE_RESET=true` или `make clean`. diff --git a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md index 5d89088..081ae0e 100644 --- a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md +++ b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md @@ -90,6 +90,8 @@ ADR-0005 решил отвязать время генератора от реа одну настенную секунду. Формат: положительное число, по умолчанию `1`. - `GEN_RUN_MODE` — режим запуска: `live` или `backfill`. Значение по умолчанию — `live`. +- `GEN_STARTUP_HISTORY_ARTIFACT` — путь к JSON-файлу, куда `backfill` дополнительно + пишет портативный артефакт стартовой истории. В обычном live-запуске не нужен. - `GEN_SEED`, `GEN_TICK_SECONDS` и остальные настройки генерации остаются частью контракта повторяемости. Если они отличаются, артефакт стартовой истории считается другим. @@ -162,12 +164,35 @@ resume_model_at = `GEN_MODEL_TIME_SPEED` короткая настенная пауза может стать долгой модельной паузой, и тогда просроченные активные визиты закрываются. +Если state отсутствует или повреждён, генератор стартует чисто и пишет +предупреждение. Если state читается, `GEN_STATE_RESET=false`, но настройки +продолжения несовместимы (`GEN_SEED`, `GEN_MODEL_T0`, `GEN_MODEL_TIMEZONE`, +`GEN_MODEL_TIME_SPEED`, а для стартовой истории ещё и `GEN_MODEL_T_END` или +manifest), генератор должен упасть с перечнем разошедшихся полей и подсказкой +использовать `GEN_STATE_RESET=true` для осознанного нового мира. + #### Манифест стартовой истории Стартовая история состоит из трёх частей: события, слепок состояния и манифест. -Манифест хранится как JSON в Kafka compact-topic +В рабочем стенде манифест хранится как JSON в Kafka compact-topic `generator_startup_history_manifest`, ключ `default`. Слепок состояния хранится -в `generator_state`, ключ `default`. Манифест минимум содержит: +в `generator_state`, ключ `default`. + +Портативный файл-артефакт хранит тот же связный набор: сообщения топиков +`browser_events`, `location_events`, `device_events`, `geo_events`, слепок state +и manifest. Для событий хранится raw JSON value, чтобы импорт мог воспроизвести +Kafka-сообщения без повторной сериализации dict. Импорт артефакта воспроизводит +сообщения в Kafka и записывает state с manifest в служебные compact-топики. +Напрямую в ClickHouse импорт не пишет: ClickHouse наполняется штатным путём через +Kafka engine и Materialized View. + +Импорт рассчитан на чистый стенд. Перед записью он проверяет, что data-топики +Kafka пустые. Так как текущий Python-клиент Kafka не даёт транзакционный producer +для нескольких топиков, при ошибке записи импорт удаляет import-топики Kafka, +чтобы повторный импорт не дописал дубли; если ClickHouse уже успел прочитать +частичные сообщения, стенд очищается как clean-stand сценарий. + +Манифест минимум содержит: - `manifest_version`; - `generated_at` — настенная UTC-метка создания артефакта; @@ -186,8 +211,8 @@ resume_model_at = только если `state.last_batch_id`, `state.model_timestamp`, `GEN_SEED`, `GEN_MODEL_T0`, `GEN_MODEL_T_END`, `GEN_MODEL_TIMEZONE`, `GEN_MODEL_TIME_SPEED` и `generation_settings` совпадают. Startup-history state без подходящего -manifest считается несовместимым и ведёт к чистому старту, а не к восстановлению -по правилу live-сбоя. +manifest считается несовместимым читаемым state и даёт жёсткий отказ при +`GEN_STATE_RESET=false`. #### Повторяемая проверка в ClickHouse diff --git a/generator/README.md b/generator/README.md index f5a75b0..5837633 100644 --- a/generator/README.md +++ b/generator/README.md @@ -78,6 +78,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | | `GEN_RUN_MODE` | Режим генератора | `live` | +| `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | — | | `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` | | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | @@ -104,6 +105,8 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator на `T_end` в `generator_state` и пишет manifest в compact-topic `generator_startup_history_manifest`. Live-запуск с теми же настройками использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты. +Экспорт и импорт портативного файла описаны в +[`docs/runbooks/startup-history.md`](../docs/runbooks/startup-history.md). Для одноразового backfill-запуска через compose используйте `docker compose run --rm generator`, а не `docker compose up generator`: у штатного сервиса включён restart policy. diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index aaf33b0..a0efe22 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -92,6 +92,11 @@ class Config: run_mode: str = field( default_factory=lambda: os.getenv("GEN_RUN_MODE", "live") ) + startup_history_artifact: Path | None = field( + default_factory=lambda: Path(os.getenv("GEN_STARTUP_HISTORY_ARTIFACT")) + if os.getenv("GEN_STARTUP_HISTORY_ARTIFACT") + else None + ) def __post_init__(self): if self.tick_seconds < 1: diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py index 1774211..3bd80ba 100644 --- a/generator/src/clickstream_generator/kafka_io.py +++ b/generator/src/clickstream_generator/kafka_io.py @@ -62,6 +62,16 @@ def _kafka_importer(): return _facade_attr("_import_kafka", _import_kafka) +def kafka_event_value_json(event: dict) -> str: + """Возвращает JSON value ровно в формате Kafka producer генератора.""" + return json.dumps(event) + + +def kafka_event_value_bytes(event: dict) -> bytes: + """Возвращает Kafka value bytes для события генератора.""" + return kafka_event_value_json(event).encode("utf-8") + + @dataclass class BatchRecord: """Запись об отправленном батче.""" @@ -380,7 +390,7 @@ class KafkaPublisher: logger.info(f"Connecting to Kafka at {self.bootstrap_servers}") self.producer = KafkaProducerCls( bootstrap_servers=self.bootstrap_servers, - value_serializer=lambda v: json.dumps(v).encode("utf-8"), + value_serializer=kafka_event_value_bytes, key_serializer=lambda k: k.encode("utf-8") if k else None, batch_size=16384, linger_ms=100, diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 270e34f..715ed80 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -1,7 +1,5 @@ """Основной сервисный цикл генератора.""" -import hashlib -import json import logging import sys import time @@ -27,108 +25,19 @@ from clickstream_generator.metrics import ( METRICS_TICK_DURATION, ) from clickstream_generator.runtime import TickStreamGenerator +from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + generation_settings_from_config, + write_startup_history_artifact, +) logger = logging.getLogger("generator") -class _TopicManifestStats: - """Накопительные счётчики одного Kafka-топика для manifest.""" - - def __init__(self): - self.rows = 0 - self.min_event_timestamp: str | None = None - self.max_event_timestamp: str | None = None - self._checksum = hashlib.sha256() - - @property - def checksum(self) -> str: - return self._checksum.hexdigest() - - def add(self, event: dict, event_timestamp: str | None = None) -> None: - self.rows += 1 - self._checksum.update( - json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") - ) - timestamp = event_timestamp if event_timestamp is not None else event.get("event_timestamp") - self.add_timestamp(timestamp) - - def add_timestamp(self, timestamp: str | None) -> None: - if timestamp is None: - return - if self.min_event_timestamp is None or timestamp < self.min_event_timestamp: - self.min_event_timestamp = timestamp - if self.max_event_timestamp is None or timestamp > self.max_event_timestamp: - self.max_event_timestamp = timestamp - - def to_dict(self) -> dict: - return { - "rows": self.rows, - "min_event_timestamp": self.min_event_timestamp, - "max_event_timestamp": self.max_event_timestamp, - "checksum_sha256": self.checksum, - } - - -class _ManifestCounters: - """Счётчики стартовой истории для manifest.""" - - def __init__(self): - self.topic_stats = { - topic: _TopicManifestStats() - for topic in ("browser_events", "location_events", "device_events", "geo_events") - } - self.click_ids: set[str] = set() - self.user_ids: set[str] = set() - - def add_batch(self, batch: dict[str, list[dict]]) -> None: - browser_events = batch.get("browser_events", []) - event_timestamps = { - event["event_id"]: event.get("event_timestamp") - for event in browser_events - if event.get("event_id") and event.get("event_timestamp") - } - click_timestamps: dict[str, list[str]] = {} - for event in browser_events: - click_id = event.get("click_id") - timestamp = event.get("event_timestamp") - if click_id and timestamp: - click_timestamps.setdefault(click_id, []).append(timestamp) - - for topic, events in batch.items(): - stats = self.topic_stats[topic] - for event in events: - event_timestamp = event.get("event_timestamp") - if event_timestamp is None and topic == "location_events": - event_timestamp = event_timestamps.get(event.get("event_id")) - if event_timestamp is None and topic in ("device_events", "geo_events"): - timestamps = click_timestamps.get(event.get("click_id")) - if timestamps: - event_timestamp = min(timestamps) - stats.add_timestamp(max(timestamps)) - stats.add(event, event_timestamp=event_timestamp) - click_id = event.get("click_id") - if topic == "browser_events" and click_id: - self.click_ids.add(click_id) - user_id = event.get("user_domain_id") - if topic == "device_events" and user_id: - self.user_ids.add(user_id) - - def to_manifest_topics(self) -> dict: - return { - topic: stats.to_dict() - for topic, stats in self.topic_stats.items() - } - - def to_manifest_totals(self) -> dict: - browser_stats = self.topic_stats["browser_events"] - return { - "events": browser_stats.rows, - "visits": len(self.click_ids), - "users": len(self.user_ids), - "min_event_timestamp": browser_stats.min_event_timestamp, - "max_event_timestamp": browser_stats.max_event_timestamp, - } +class IncompatibleStateError(ValueError): + """Читаемый state относится к другому миру генерации.""" class GeneratorService: @@ -201,8 +110,15 @@ class GeneratorService: model_t_end=model_t_end, ) elif self._is_startup_history_marker(restored_state): - raise ValueError( - "startup-history state without matching manifest" + mismatches = self._startup_history_mismatch_fields( + restored_state, + manifest, + ) + raise IncompatibleStateError( + "startup-history state mismatch: " + + ", ".join(mismatches) + + "; set GEN_STATE_RESET=true to start a new " + "world intentionally" ) else: self._restore_live_state( @@ -214,6 +130,13 @@ class GeneratorService: f"model_time={self._model_time.isoformat()}, " f"last_batch_id={restored_state.last_batch_id}" ) + except IncompatibleStateError: + logger.error( + "Readable generator state is incompatible with " + "current settings; set GEN_STATE_RESET=true to " + "start a new world intentionally." + ) + raise except Exception as e: logger.warning( f"State data was invalid, starting fresh: {e}" @@ -290,27 +213,36 @@ class GeneratorService: """Проверяет, что state относится к текущей конфигурации live-запуска.""" mismatches = [] if state.gen_seed != self.config.seed: - mismatches.append("gen_seed") + mismatches.append("GEN_SEED") if self._as_aware_utc(state.model_t0) != self.config.model_t0: - mismatches.append("model_t0") + mismatches.append("GEN_MODEL_T0") if state.model_timezone != self.config.model_timezone: - mismatches.append("model_timezone") + mismatches.append("GEN_MODEL_TIMEZONE") if abs(state.model_time_speed - self.config.model_time_speed) > 1e-9: - mismatches.append("model_time_speed") + mismatches.append("GEN_MODEL_TIME_SPEED") if mismatches: - raise ValueError( - "state config mismatch: " + ", ".join(mismatches) + raise IncompatibleStateError( + "state config mismatch: " + + ", ".join(mismatches) + + "; set GEN_STATE_RESET=true to start a new world intentionally" ) def _is_startup_history_state(self, state, manifest: dict | None) -> bool: """Проверяет, что state совпадает со слепком стартовой истории.""" + return not self._startup_history_mismatch_fields(state, manifest) + + def _startup_history_mismatch_fields(self, state, manifest: dict | None) -> list[str]: + """Возвращает поля, по которым startup-history state не совпал.""" + mismatches = [] if not manifest or manifest.get("run_mode") != "backfill": - return False + return ["generator_startup_history_manifest"] expected_state = manifest.get("state") or {} + if manifest.get("state_version") != state.version: + mismatches.append("state_version") if expected_state.get("last_batch_id") != state.last_batch_id: - return False + mismatches.append("state.last_batch_id") try: model_t_end = self._as_aware_utc( @@ -320,32 +252,34 @@ class GeneratorService: datetime.fromisoformat(manifest["model_t0"]) ) except (KeyError, TypeError, ValueError): - return False + return ["manifest.model_t0", "manifest.model_t_end"] if self._as_aware_utc(state.model_timestamp) != model_t_end: - return False + mismatches.append("GEN_MODEL_T_END") if self._as_aware_utc(state.model_t0) != model_t0: - return False + mismatches.append("GEN_MODEL_T0") if state.gen_seed != manifest.get("gen_seed"): - return False + mismatches.append("GEN_SEED") if state.model_timezone != manifest.get("model_timezone"): - return False + mismatches.append("GEN_MODEL_TIMEZONE") if ( self.config.model_t_end is not None and self.config.model_t_end != model_t_end ): - return False + mismatches.append("GEN_MODEL_T_END") if model_t0 != self.config.model_t0: - return False + mismatches.append("GEN_MODEL_T0") if manifest.get("gen_seed") != self.config.seed: - return False + mismatches.append("GEN_SEED") if manifest.get("model_timezone") != self.config.model_timezone: - return False + mismatches.append("GEN_MODEL_TIMEZONE") settings = manifest.get("generation_settings") or {} if abs(state.model_time_speed - settings.get("model_time_speed", -1)) > 1e-9: - return False - return self._generation_settings() == settings + mismatches.append("GEN_MODEL_TIME_SPEED") + if self._generation_settings() != settings: + mismatches.append("generation_settings") + return list(dict.fromkeys(mismatches)) def _is_startup_history_marker(self, state) -> bool: """Отличает state стартовой истории от обычного live-state.""" @@ -400,7 +334,7 @@ class GeneratorService: self.config.model_t0.isoformat(), self.config.model_t_end.isoformat(), ) - counters = _ManifestCounters() + artifact_builder = StartupHistoryArtifactBuilder() while self._model_time < self.config.model_t_end: self._tick += 1 @@ -415,7 +349,7 @@ class GeneratorService: ) total_sent, sent_counts, status = self._publish_batch(batch) self._raise_on_backfill_publish_error(batch_id, status, sent_counts) - counters.add_batch(batch) + artifact_builder.add_batch(batch) self._write_batch_history( batch_id=batch_id, started_at=started_at, @@ -441,7 +375,7 @@ class GeneratorService: final_status, sent_counts, ) - counters.add_batch(final_batch) + artifact_builder.add_batch(final_batch) self._write_batch_history( batch_id=batch_id, started_at=started_at, @@ -455,7 +389,7 @@ class GeneratorService: self.publisher.flush() self._model_time = self.config.model_t_end - state_batch_id = self._startup_state_batch_id(counters) + state_batch_id = self._startup_state_batch_id(artifact_builder.counters) state = self.stream.to_state( tick=self._tick, rng_state=self.generator.rng.getstate(), @@ -469,8 +403,9 @@ class GeneratorService: gen_seed=self.config.seed, ) - manifest = self._build_startup_history_manifest( - counters=counters, + manifest = build_manifest( + config=self.config, + counters=artifact_builder.counters, state=state, ) if self.manifest_manager: @@ -481,6 +416,17 @@ class GeneratorService: self.state_manager.save(state) self.state_manager.flush() + if self.config.startup_history_artifact is not None: + artifact = artifact_builder.to_artifact(manifest=manifest, state=state) + write_startup_history_artifact( + self.config.startup_history_artifact, + artifact, + ) + logger.info( + "Startup history artifact written: %s", + self.config.startup_history_artifact, + ) + logger.info( "Backfill completed: events=%s, visits=%s, users=%s, final_sent=%s", manifest["totals"]["events"], @@ -565,40 +511,10 @@ class GeneratorService: return f"startup-history-{digest}" def _build_startup_history_manifest(self, counters, state) -> dict: - return { - "manifest_version": "1.0", - "generated_at": datetime.now(timezone.utc).isoformat(), - "gen_seed": self.config.seed, - "model_t0": self.config.model_t0.isoformat(), - "model_t_end": self.config.model_t_end.isoformat(), - "model_timezone": self.config.model_timezone, - "run_mode": "backfill", - "generation_settings": self._generation_settings(), - "state_version": state.version, - "state": { - "topic": KafkaStateManager.STATE_TOPIC, - "key": KafkaStateManager.STATE_KEY, - "last_batch_id": state.last_batch_id, - "model_timestamp": state.model_timestamp.isoformat(), - }, - "topics": counters.to_manifest_topics(), - "totals": counters.to_manifest_totals(), - } + return build_manifest(config=self.config, counters=counters, state=state) def _generation_settings(self) -> dict: - return { - "tick_seconds": self.config.tick_seconds, - "lambda_base_per_min": self.config.lambda_base_per_min, - "jitter_pct": self.config.jitter_pct, - "min_events_per_tick": self.config.min_events_per_tick, - "max_events_per_tick": self.config.max_events_per_tick, - "max_session_events": self.config.max_session_events, - "max_active_sessions": self.config.max_active_sessions, - "population_max": self.config.population_max, - "p_new_user": self.config.p_new_user, - "min_return_minutes": self.config.min_return_minutes, - "model_time_speed": self.config.model_time_speed, - } + return generation_settings_from_config(self.config) def _main_loop(self): """Основной цикл тиков.""" diff --git a/generator/src/clickstream_generator/startup_history_artifact.py b/generator/src/clickstream_generator/startup_history_artifact.py new file mode 100644 index 0000000..e4bba3c --- /dev/null +++ b/generator/src/clickstream_generator/startup_history_artifact.py @@ -0,0 +1,642 @@ +"""Портативный артефакт стартовой истории.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import logging +from copy import deepcopy +from datetime import datetime, timezone +from pathlib import Path +from typing import Any + +from clickstream_generator.config import Config +from clickstream_generator.kafka_io import ( + KafkaStateManager, + KafkaStartupHistoryManifest, + ensure_topics, + kafka_event_value_json, + _kafka_importer, +) +from clickstream_generator.state import GeneratorState + + +logger = logging.getLogger("generator") + +ARTIFACT_KIND = "clickstream_generator_startup_history" +ARTIFACT_VERSION = "1.0" +TOPICS = ("browser_events", "location_events", "device_events", "geo_events") +IMPORT_TOPICS = ( + "browser_events", + "location_events", + "device_events", + "geo_events", + KafkaStateManager.STATE_TOPIC, + KafkaStartupHistoryManifest.MANIFEST_TOPIC, +) + + +class _TopicManifestStats: + """Накопительные счётчики одного Kafka-топика для manifest.""" + + def __init__(self): + self.rows = 0 + self.min_event_timestamp: str | None = None + self.max_event_timestamp: str | None = None + self._checksum = hashlib.sha256() + + @property + def checksum(self) -> str: + return self._checksum.hexdigest() + + def add(self, event: dict, event_timestamp: str | None = None) -> None: + self.rows += 1 + self._checksum.update( + json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") + ) + timestamp = ( + event_timestamp + if event_timestamp is not None + else event.get("event_timestamp") + ) + self.add_timestamp(timestamp) + + def add_timestamp(self, timestamp: str | None) -> None: + if timestamp is None: + return + if self.min_event_timestamp is None or timestamp < self.min_event_timestamp: + self.min_event_timestamp = timestamp + if self.max_event_timestamp is None or timestamp > self.max_event_timestamp: + self.max_event_timestamp = timestamp + + def to_dict(self) -> dict: + return { + "rows": self.rows, + "min_event_timestamp": self.min_event_timestamp, + "max_event_timestamp": self.max_event_timestamp, + "checksum_sha256": self.checksum, + } + + +class ManifestCounters: + """Счётчики стартовой истории для manifest.""" + + def __init__(self): + self.topic_stats = {topic: _TopicManifestStats() for topic in TOPICS} + self.click_ids: set[str] = set() + self.user_ids: set[str] = set() + + def add_batch(self, batch: dict[str, list[dict]]) -> None: + browser_events = batch.get("browser_events", []) + event_timestamps = { + event["event_id"]: event.get("event_timestamp") + for event in browser_events + if event.get("event_id") and event.get("event_timestamp") + } + click_timestamps: dict[str, list[str]] = {} + for event in browser_events: + click_id = event.get("click_id") + timestamp = event.get("event_timestamp") + if click_id and timestamp: + click_timestamps.setdefault(click_id, []).append(timestamp) + + for topic, events in batch.items(): + if topic not in self.topic_stats: + raise ValueError(f"unknown startup history topic: {topic}") + stats = self.topic_stats[topic] + for event in events: + event_timestamp = event.get("event_timestamp") + if event_timestamp is None and topic == "location_events": + event_timestamp = event_timestamps.get(event.get("event_id")) + if event_timestamp is None and topic in ("device_events", "geo_events"): + timestamps = click_timestamps.get(event.get("click_id")) + if timestamps: + event_timestamp = min(timestamps) + stats.add_timestamp(max(timestamps)) + stats.add(event, event_timestamp=event_timestamp) + click_id = event.get("click_id") + if topic == "browser_events" and click_id: + self.click_ids.add(click_id) + user_id = event.get("user_domain_id") + if topic == "device_events" and user_id: + self.user_ids.add(user_id) + + def to_manifest_topics(self) -> dict: + return {topic: stats.to_dict() for topic, stats in self.topic_stats.items()} + + def to_manifest_totals(self) -> dict: + browser_stats = self.topic_stats["browser_events"] + return { + "events": browser_stats.rows, + "visits": len(self.click_ids), + "users": len(self.user_ids), + "min_event_timestamp": browser_stats.min_event_timestamp, + "max_event_timestamp": browser_stats.max_event_timestamp, + } + + +class StartupHistoryArtifactBuilder: + """Собирает сообщения топиков для портативного файла.""" + + def __init__(self): + self.topics = {topic: [] for topic in TOPICS} + self.raw_topics = {topic: [] for topic in TOPICS} + self.counters = ManifestCounters() + + def add_batch(self, batch: dict[str, list[dict]]) -> None: + self.counters.add_batch(batch) + for topic in TOPICS: + for event in batch.get(topic, []): + self.topics[topic].append(deepcopy(event)) + self.raw_topics[topic].append( + { + "key": _event_key(event), + "value_json": _event_value_json(event), + } + ) + + def to_artifact(self, manifest: dict, state: GeneratorState) -> dict: + return { + "artifact_version": ARTIFACT_VERSION, + "kind": ARTIFACT_KIND, + "created_at": datetime.now(timezone.utc).isoformat(), + "manifest": deepcopy(manifest), + "state": state.to_dict(), + "topics": deepcopy(self.topics), + "raw_topics": deepcopy(self.raw_topics), + } + + +class KafkaRawPublisher: + """Публикует сохранённые Kafka value bytes без повторной JSON-сериализации.""" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + KafkaProducerCls, _ = _kafka_importer()() + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + + def publish_records(self, topic: str, records: list[dict]) -> tuple[int, int]: + sent = 0 + errors = 0 + futures = [] + for record in records: + future = self.producer.send( + topic, + key=record.get("key"), + value=record["value_json"].encode("utf-8"), + ) + futures.append(future) + for future in futures: + try: + future.get(timeout=10) + sent += 1 + except Exception: + errors += 1 + return sent, errors + + def flush(self) -> None: + self.producer.flush() + + def close(self) -> None: + self.producer.close() + + +class KafkaTopicInspector: + """Проверяет Kafka-топики перед clean-stand импортом.""" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + + def assert_data_topics_empty(self) -> None: + """Падает, если в data-топиках уже есть сообщения.""" + from kafka import KafkaConsumer, TopicPartition + + consumer = KafkaConsumer( + bootstrap_servers=self.bootstrap_servers, + enable_auto_commit=False, + consumer_timeout_ms=1000, + ) + try: + dirty_topics = [] + for topic in TOPICS: + partitions = consumer.partitions_for_topic(topic) + if not partitions: + continue + topic_partitions = [ + TopicPartition(topic, partition) + for partition in partitions + ] + end_offsets = consumer.end_offsets(topic_partitions) + if any(offset > 0 for offset in end_offsets.values()): + dirty_topics.append(topic) + if dirty_topics: + raise RuntimeError( + "Kafka data topics are not empty; clean the stand before " + "startup history import: " + ", ".join(dirty_topics) + ) + finally: + consumer.close() + + def snapshot_import_topics(self) -> dict: + """Снимок для отката clean-stand импорта.""" + return {"topics": list(IMPORT_TOPICS)} + + def rollback_import_topics(self, snapshot: dict) -> None: + """Удаляет топики, которые могли получить частичный импорт.""" + from kafka import KafkaAdminClient + from kafka.errors import UnknownTopicOrPartitionError + + admin = KafkaAdminClient(bootstrap_servers=self.bootstrap_servers) + try: + existing_topics = set(admin.list_topics()) + topics = [ + topic + for topic in snapshot.get("topics", IMPORT_TOPICS) + if topic in existing_topics + ] + if not topics: + return + try: + admin.delete_topics(topics) + except UnknownTopicOrPartitionError: + return + finally: + admin.close() + + +def generation_settings_from_config(config: Config) -> dict: + """Возвращает настройки, влияющие на поток генерации.""" + return { + "tick_seconds": config.tick_seconds, + "lambda_base_per_min": config.lambda_base_per_min, + "jitter_pct": config.jitter_pct, + "min_events_per_tick": config.min_events_per_tick, + "max_events_per_tick": config.max_events_per_tick, + "max_session_events": config.max_session_events, + "max_active_sessions": config.max_active_sessions, + "population_max": config.population_max, + "p_new_user": config.p_new_user, + "min_return_minutes": config.min_return_minutes, + "model_time_speed": config.model_time_speed, + } + + +def build_manifest(config: Config, counters: ManifestCounters, state: GeneratorState) -> dict: + """Строит manifest стартовой истории.""" + if config.model_t_end is None: + raise ValueError("GEN_MODEL_T_END is required for startup history manifest") + return { + "manifest_version": "1.0", + "generated_at": datetime.now(timezone.utc).isoformat(), + "gen_seed": config.seed, + "model_t0": config.model_t0.isoformat(), + "model_t_end": config.model_t_end.isoformat(), + "model_timezone": config.model_timezone, + "run_mode": "backfill", + "generation_settings": generation_settings_from_config(config), + "state_version": state.version, + "state": { + "topic": KafkaStateManager.STATE_TOPIC, + "key": KafkaStateManager.STATE_KEY, + "last_batch_id": state.last_batch_id, + "model_timestamp": state.model_timestamp.isoformat(), + }, + "topics": counters.to_manifest_topics(), + "totals": counters.to_manifest_totals(), + } + + +def write_startup_history_artifact(path: str | Path, artifact: dict) -> None: + """Пишет артефакт в JSON-файл.""" + target = Path(path) + target.parent.mkdir(parents=True, exist_ok=True) + target.write_text( + json.dumps(artifact, ensure_ascii=True, indent=2), + encoding="utf-8", + ) + + +def load_startup_history_artifact(path: str | Path) -> dict: + """Читает артефакт из JSON-файла.""" + return json.loads(Path(path).read_text(encoding="utf-8")) + + +def validate_startup_history_artifact( + artifact: dict, + expected_config: Config | None = None, +) -> tuple[GeneratorState, dict, dict[str, list[dict]], dict[str, list[dict]]]: + """Проверяет связку manifest, state и сообщений топиков.""" + if not isinstance(artifact, dict): + raise ValueError("artifact must be an object") + if artifact.get("kind") != ARTIFACT_KIND: + raise ValueError("artifact kind mismatch") + if artifact.get("artifact_version") != ARTIFACT_VERSION: + raise ValueError("artifact version mismatch") + + manifest = artifact.get("manifest") + state_payload = artifact.get("state") + topics = artifact.get("topics") + raw_topics = artifact.get("raw_topics") + if not isinstance(manifest, dict): + raise ValueError("artifact manifest must be an object") + if not isinstance(state_payload, dict): + raise ValueError("artifact state must be an object") + if not isinstance(topics, dict): + raise ValueError("artifact topics must be an object") + if not isinstance(raw_topics, dict): + raise ValueError("artifact raw_topics must be an object") + for topic in TOPICS: + if not isinstance(topics.get(topic), list): + raise ValueError(f"artifact topic {topic} must be a list") + if not isinstance(raw_topics.get(topic), list): + raise ValueError(f"artifact raw topic {topic} must be a list") + + state = GeneratorState.from_dict(state_payload) + _validate_manifest_state(manifest, state) + _validate_manifest_topics(manifest, topics) + _validate_raw_topics(topics, raw_topics) + if expected_config is not None: + _validate_manifest_config(manifest, expected_config) + return ( + state, + manifest, + {topic: topics[topic] for topic in TOPICS}, + {topic: raw_topics[topic] for topic in TOPICS}, + ) + + +def import_startup_history_artifact( + artifact: dict, + *, + publisher, + state_manager, + manifest_manager, + expected_config: Config | None = None, + topic_inspector=None, +) -> dict: + """Воспроизводит артефакт в Kafka и compact-топиках генератора.""" + state, manifest, topics, raw_topics = validate_startup_history_artifact( + artifact, + expected_config=expected_config, + ) + if topic_inspector is not None: + topic_inspector.assert_data_topics_empty() + snapshot = ( + topic_inspector.snapshot_import_topics() + if hasattr(topic_inspector, "snapshot_import_topics") + else None + ) + else: + snapshot = None + + if not hasattr(publisher, "publish_records"): + raise TypeError( + "startup history import requires publisher.publish_records " + "for raw Kafka values" + ) + + try: + sent_total = 0 + sent_by_topic = {} + + for topic in TOPICS: + raw_records = raw_topics[topic] + if not raw_records: + sent_by_topic[topic] = 0 + continue + sent, errors = publisher.publish_records(topic, raw_records) + if errors: + raise RuntimeError( + f"Startup history import failed for {topic}: sent={sent}, errors={errors}" + ) + sent_total += sent + sent_by_topic[topic] = sent + + publisher.flush() + manifest_manager.save(manifest) + manifest_manager.flush() + state_manager.save(state) + state_manager.flush() + except Exception: + if ( + topic_inspector is not None + and snapshot is not None + and hasattr(topic_inspector, "rollback_import_topics") + ): + topic_inspector.rollback_import_topics(snapshot) + raise + + return {"events": sent_total, "topics": sent_by_topic} + + +def compare_clickhouse_stats_to_manifest( + manifest: dict, + stats: dict[str, str], +) -> list[str]: + """Возвращает поля ClickHouse, которые не совпали с manifest.""" + totals = manifest.get("totals") or {} + expected = { + "events": str(totals.get("events")), + "visits": str(totals.get("visits")), + "users": str(totals.get("users")), + "min_event_timestamp": str(totals.get("min_event_timestamp")), + "max_event_timestamp": str(totals.get("max_event_timestamp")), + } + return [ + field + for field, expected_value in expected.items() + if str(stats.get(field)) != expected_value + ] + + +def _validate_manifest_state(manifest: dict, state: GeneratorState) -> None: + mismatches = [] + expected_state = manifest.get("state") or {} + if manifest.get("run_mode") != "backfill": + mismatches.append("run_mode") + if manifest.get("state_version") != state.version: + mismatches.append("state_version") + if expected_state.get("last_batch_id") != state.last_batch_id: + mismatches.append("state.last_batch_id") + if expected_state.get("model_timestamp") != state.model_timestamp.isoformat(): + mismatches.append("state.model_timestamp") + try: + model_t0 = _parse_manifest_timestamp(manifest["model_t0"]) + model_t_end = _parse_manifest_timestamp(manifest["model_t_end"]) + except (KeyError, TypeError, ValueError) as exc: + raise ValueError(f"manifest time fields are invalid: {exc}") from exc + if state.model_t0 != model_t0: + mismatches.append("GEN_MODEL_T0") + if state.model_timestamp != model_t_end: + mismatches.append("GEN_MODEL_T_END") + if state.gen_seed != manifest.get("gen_seed"): + mismatches.append("GEN_SEED") + if state.model_timezone != manifest.get("model_timezone"): + mismatches.append("GEN_MODEL_TIMEZONE") + settings = manifest.get("generation_settings") or {} + if abs(state.model_time_speed - settings.get("model_time_speed", -1)) > 1e-9: + mismatches.append("GEN_MODEL_TIME_SPEED") + if mismatches: + raise ValueError("startup history artifact mismatch: " + ", ".join(mismatches)) + + +def _validate_manifest_topics(manifest: dict, topics: dict[str, list[dict]]) -> None: + counters = ManifestCounters() + counters.add_batch({topic: topics[topic] for topic in TOPICS}) + expected_topics = manifest.get("topics") or {} + expected_totals = manifest.get("totals") or {} + if counters.to_manifest_topics() != expected_topics: + raise ValueError("startup history artifact mismatch: topic checksums") + if counters.to_manifest_totals() != expected_totals: + raise ValueError("startup history artifact mismatch: totals") + + +def _validate_raw_topics( + topics: dict[str, list[dict]], + raw_topics: dict[str, list[dict]], +) -> None: + for topic in TOPICS: + if len(raw_topics[topic]) != len(topics[topic]): + raise ValueError(f"artifact raw topic {topic} row count mismatch") + for event, record in zip(topics[topic], raw_topics[topic]): + if record.get("key") != _event_key(event): + raise ValueError(f"artifact raw topic {topic} key mismatch") + if _raw_value_to_event(record.get("value_json"), topic) != event: + raise ValueError(f"artifact raw topic {topic} value mismatch") + + +def _event_key(event: dict) -> str | None: + return event.get("event_id") or event.get("click_id") + + +def _event_value_json(event: dict) -> str: + return kafka_event_value_json(event) + + +def _raw_value_to_event(value: Any, topic: str) -> dict: + if not isinstance(value, str): + raise ValueError(f"artifact raw topic {topic} value must be a string") + try: + event = json.loads(value) + except json.JSONDecodeError as exc: + raise ValueError(f"artifact raw topic {topic} value must be valid JSON") from exc + if not isinstance(event, dict): + raise ValueError(f"artifact raw topic {topic} value must be a JSON object") + return event + + +def _validate_manifest_config(manifest: dict, config: Config) -> None: + mismatches = [] + if manifest.get("gen_seed") != config.seed: + mismatches.append("GEN_SEED") + if _parse_manifest_timestamp(manifest.get("model_t0")) != config.model_t0: + mismatches.append("GEN_MODEL_T0") + if config.model_t_end is None: + mismatches.append("GEN_MODEL_T_END") + elif _parse_manifest_timestamp(manifest.get("model_t_end")) != config.model_t_end: + mismatches.append("GEN_MODEL_T_END") + if manifest.get("model_timezone") != config.model_timezone: + mismatches.append("GEN_MODEL_TIMEZONE") + if manifest.get("generation_settings") != generation_settings_from_config(config): + mismatches.append("generation_settings") + if mismatches: + raise ValueError( + "startup history artifact is incompatible with current settings: " + + ", ".join(mismatches) + ) + + +def _parse_manifest_timestamp(value: Any) -> datetime: + if not isinstance(value, str): + raise ValueError("timestamp must be a string") + timestamp = datetime.fromisoformat(value.replace("Z", "+00:00")) + if timestamp.tzinfo is None: + raise ValueError("timestamp must include timezone") + return timestamp.astimezone(timezone.utc) + + +def _run_import(args: argparse.Namespace) -> int: + config = Config() + artifact = load_startup_history_artifact(args.artifact) + ensure_topics(config.kafka_bootstrap_servers) + publisher = KafkaRawPublisher(config.kafka_bootstrap_servers) + state_manager = KafkaStateManager(config.kafka_bootstrap_servers) + manifest_manager = KafkaStartupHistoryManifest(config.kafka_bootstrap_servers) + topic_inspector = KafkaTopicInspector(config.kafka_bootstrap_servers) + try: + result = import_startup_history_artifact( + artifact, + publisher=publisher, + state_manager=state_manager, + manifest_manager=manifest_manager, + expected_config=config, + topic_inspector=topic_inspector, + ) + finally: + publisher.close() + state_manager.close() + manifest_manager.close() + logger.info( + "Startup history artifact imported: events=%s, topics=%s", + result["events"], + result["topics"], + ) + return 0 + + +def _run_validate(args: argparse.Namespace) -> int: + config = Config() + artifact = load_startup_history_artifact(args.artifact) + validate_startup_history_artifact(artifact, expected_config=config) + logger.info("Startup history artifact is valid: %s", args.artifact) + return 0 + + +def _run_manifest_summary(args: argparse.Namespace) -> int: + artifact = load_startup_history_artifact(args.artifact) + _, manifest, _, _ = validate_startup_history_artifact(artifact) + totals = manifest["totals"] + print( + "\t".join( + [ + str(totals["events"]), + str(totals["visits"]), + str(totals["users"]), + str(totals["min_event_timestamp"]), + str(totals["max_event_timestamp"]), + str(manifest["model_t0"]), + str(manifest["model_t_end"]), + ] + ) + ) + return 0 + + +def main(argv: list[str] | None = None) -> int: + parser = argparse.ArgumentParser(description="Startup history artifact tools") + subparsers = parser.add_subparsers(dest="command", required=True) + + validate_parser = subparsers.add_parser("validate") + validate_parser.add_argument("--artifact", required=True) + validate_parser.set_defaults(func=_run_validate) + + import_parser = subparsers.add_parser("import") + import_parser.add_argument("--artifact", required=True) + import_parser.set_defaults(func=_run_import) + + summary_parser = subparsers.add_parser("manifest-summary") + summary_parser.add_argument("--artifact", required=True) + summary_parser.set_defaults(func=_run_manifest_summary) + + args = parser.parse_args(argv) + return args.func(args) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index c03e515..d461c60 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -369,6 +369,46 @@ class TestGeneratorServiceBackfill: assert first["digest"] == second["digest"] assert first["manifest_digest"] == second["manifest_digest"] + def test_backfill_writes_portable_artifact_when_requested( + self, base_config, tmp_path + ): + """Backfill пишет переносимый файл с событиями, state и manifest.""" + from clickstream_generator.startup_history_artifact import ( + load_startup_history_artifact, + validate_startup_history_artifact, + ) + + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=1) + artifact_path = tmp_path / "startup-history.json" + config = replace( + base_config, + run_mode="backfill", + model_t0=model_t0, + model_t_end=model_t_end, + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + state_enabled=True, + startup_history_artifact=artifact_path, + ) + + self._run_backfill(config) + + artifact = load_startup_history_artifact(artifact_path) + state, manifest, topics, _ = validate_startup_history_artifact( + artifact, + expected_config=config, + ) + assert topics["browser_events"] + assert state.last_batch_id == manifest["state"]["last_batch_id"] + assert manifest["totals"]["events"] == len(topics["browser_events"]) + def test_live_start_uses_startup_manifest_without_wall_delta( self, base_config, event_dictionary ): @@ -403,6 +443,7 @@ class TestGeneratorServiceBackfill: "model_t_end": model_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), + "state_version": state.version, "state": {"last_batch_id": state.last_batch_id}, } state_manager = MagicMock() @@ -468,6 +509,7 @@ class TestGeneratorServiceBackfill: "model_t_end": manifest_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), + "state_version": state.version, "state": {"last_batch_id": state.last_batch_id}, } state_manager = MagicMock() @@ -506,12 +548,61 @@ class TestGeneratorServiceBackfill: patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(config) - service.start() + with pytest.raises(ValueError, match="GEN_MODEL_T_END.*GEN_STATE_RESET=true"): + service.start() restore_from_startup_history.assert_not_called() restore_live_state.assert_not_called() - assert service._tick == 0 - assert service._model_time == config.model_t0 + + def test_live_start_fails_loudly_on_readable_incompatible_state( + self, base_config, event_dictionary + ): + """Читаемый state другого мира при продолжении даёт жёсткий отказ.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + config = replace( + base_config, + model_t0=model_t0, + tick_seconds=60, + model_time_speed=10, + state_reset=False, + ) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="live-state", + last_timestamp=model_t0, + model_timestamp=model_t0, + wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed + 1, + ) + state_manager = MagicMock() + state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = None + + with patch("clickstream_generator.service.start_http_server"), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch("clickstream_generator.service.KafkaPublisher"), \ + patch("clickstream_generator.service.KafkaBatchHistory"), \ + patch( + "clickstream_generator.service.KafkaStateManager", + return_value=state_manager, + ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ + patch.object(GeneratorService, "_main_loop", return_value=None): + + service = GeneratorService(config) + with pytest.raises(ValueError, match="GEN_SEED.*GEN_STATE_RESET=true"): + service.start() def test_live_start_rejects_orphan_startup_state_without_live_restore( self, base_config, event_dictionary, caplog @@ -564,13 +655,14 @@ class TestGeneratorServiceBackfill: caplog.at_level(logging.WARNING, logger="generator"): service = GeneratorService(config) - service.start() + with pytest.raises( + ValueError, + match="generator_startup_history_manifest.*GEN_STATE_RESET=true", + ): + service.start() restore_from_startup_history.assert_not_called() restore_live_state.assert_not_called() - assert service._tick == 0 - assert service._model_time == model_t0 - assert "startup-history state without matching manifest" in caplog.text def test_startup_history_state_checks_state_fields( self, base_config, event_dictionary @@ -601,6 +693,44 @@ class TestGeneratorServiceBackfill: "model_t_end": model_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), + "state_version": state.version, + "state": {"last_batch_id": state.last_batch_id}, + } + + service = GeneratorService(config) + + assert not service._is_startup_history_state(state, manifest) + + def test_startup_history_state_checks_state_version( + self, base_config, event_dictionary + ): + """Startup-history state сверяется с версией state из manifest.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=5) + config = replace(base_config, model_t0=model_t0, model_time_speed=10) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="startup-history-state", + last_timestamp=model_t_end, + model_timestamp=model_t_end, + wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed, + ) + manifest = { + "run_mode": "backfill", + "gen_seed": config.seed, + "model_t0": model_t0.isoformat(), + "model_t_end": model_t_end.isoformat(), + "model_timezone": config.model_timezone, + "generation_settings": GeneratorService(config)._generation_settings(), + "state_version": "1.0", "state": {"last_batch_id": state.last_batch_id}, } @@ -930,8 +1060,8 @@ class TestGeneratorServiceStateV2: assert saved_state.active_visits service.state_manager.flush.assert_called_once() - def test_incompatible_seed_state_starts_fresh(self, base_config, event_dictionary, caplog): - """State от другого GEN_SEED не смешивается с текущим запуском.""" + def test_incompatible_seed_state_fails_loudly(self, base_config, event_dictionary, caplog): + """State от другого GEN_SEED даёт жёсткий отказ при продолжении.""" source_config = replace(base_config, seed=7) source_generator = EventGenerator(event_dictionary, source_config) source_stream = TickStreamGenerator(source_generator) @@ -970,11 +1100,10 @@ class TestGeneratorServiceStateV2: caplog.at_level(logging.WARNING, logger="generator"): service = GeneratorService(base_config) - service.start() + with pytest.raises(ValueError, match="GEN_SEED.*GEN_STATE_RESET=true"): + service.start() - assert service._tick == 0 - assert service._model_time == base_config.model_t0 - assert "state config mismatch: gen_seed" in caplog.text + assert "Readable generator state is incompatible" in caplog.text def test_state_reset_skips_loading_saved_state(self, base_config): """GEN_STATE_RESET=true запускает сервис с чистого состояния.""" diff --git a/generator/tests/test_startup_history_artifact.py b/generator/tests/test_startup_history_artifact.py new file mode 100644 index 0000000..b47bf16 --- /dev/null +++ b/generator/tests/test_startup_history_artifact.py @@ -0,0 +1,518 @@ +""" +Тесты портативного артефакта стартовой истории. +""" + +from dataclasses import replace +from datetime import datetime, timezone + +import pytest + +from generator import GeneratorState + + +def _state() -> GeneratorState: + import random + + now = datetime(2026, 1, 1, 1, 0, tzinfo=timezone.utc) + return GeneratorState( + tick=1, + rng_state=random.Random(42).getstate(), + last_batch_id="startup-history-abc", + last_timestamp=now, + model_timestamp=now, + wall_timestamp=now, + model_time_speed=1, + model_timezone="UTC", + model_t0=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + gen_seed=42, + population=[ + { + "user_domain_id": "user-1", + "seed_click_id": "seed-1", + "active_click_id": None, + "last_finished_at": None, + } + ], + ) + + +def _batch() -> dict[str, list[dict]]: + return { + "browser_events": [ + { + "event_id": "event-1", + "click_id": "click-1", + "user_domain_id": "user-1", + "event_timestamp": "2026-01-01 00:00:00.000000", + } + ], + "location_events": [ + { + "event_id": "event-1", + "page_url_path": "/home", + } + ], + "device_events": [ + { + "click_id": "click-1", + "user_domain_id": "user-1", + } + ], + "geo_events": [ + { + "click_id": "click-1", + "country": "RU", + } + ], + } + + +def test_artifact_roundtrip_keeps_events_state_and_manifest(base_config): + """Артефакт хранит события, state и manifest одним проверяемым набором.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + load_startup_history_artifact, + validate_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, + ), + counters=builder.counters, + state=state, + ) + + artifact = builder.to_artifact(manifest=manifest, state=state) + ( + loaded_state, + loaded_manifest, + loaded_topics, + loaded_raw_topics, + ) = validate_startup_history_artifact( + artifact, + expected_config=replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ), + ) + + assert artifact["artifact_version"] == "1.0" + assert loaded_state.last_batch_id == state.last_batch_id + assert loaded_manifest["state"]["last_batch_id"] == state.last_batch_id + assert loaded_topics["browser_events"][0]["event_id"] == "event-1" + assert loaded_raw_topics["browser_events"][0]["value_json"].encode("utf-8") == ( + b'{"event_id": "event-1", "click_id": "click-1", ' + b'"user_domain_id": "user-1", ' + b'"event_timestamp": "2026-01-01 00:00:00.000000"}' + ) + + path = base_config.data_dir.parent / "tmp-startup-history-artifact.json" + try: + write_startup_history_artifact(path, artifact) + loaded = load_startup_history_artifact(path) + assert loaded["manifest"]["totals"]["events"] == 1 + finally: + path.unlink(missing_ok=True) + + +def test_artifact_validation_rejects_mismatched_manifest(base_config): + """Валидация отвергает артефакт, где manifest не совпадает с событиями.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + validate_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + manifest = build_manifest(config=config, counters=builder.counters, state=state) + artifact = builder.to_artifact(manifest=manifest, state=state) + artifact["manifest"]["gen_seed"] = 43 + + with pytest.raises(ValueError, match="GEN_SEED"): + validate_startup_history_artifact(artifact, expected_config=config) + + +def test_manifest_clickhouse_stats_detect_mismatch(): + """Сверка ClickHouse с manifest ловит расхождение контрольных чисел.""" + from clickstream_generator.startup_history_artifact import ( + compare_clickhouse_stats_to_manifest, + ) + + manifest = { + "totals": { + "events": 10, + "visits": 4, + "users": 3, + "min_event_timestamp": "2026-01-01 00:00:00.000000", + "max_event_timestamp": "2026-01-01 00:10:00.000000", + } + } + stats = { + "events": "9", + "visits": "4", + "users": "3", + "min_event_timestamp": "2026-01-01 00:00:00.000000", + "max_event_timestamp": "2026-01-01 00:10:00.000000", + } + + assert compare_clickhouse_stats_to_manifest(manifest, stats) == ["events"] + + +def test_raw_topic_value_can_keep_original_json_bytes(base_config): + """Raw value проверяется как JSON, но импортируется без пересборки строки.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + raw_value = ( + '{"event_timestamp":"2026-01-01 00:00:00.000000",' + '"user_domain_id":"user-1","click_id":"click-1","event_id":"event-1"}' + ) + artifact["raw_topics"]["browser_events"][0]["value_json"] = raw_value + + class Publisher: + def __init__(self): + self.records = None + + def publish_records(self, topic, records): + if topic == "browser_events": + self.records = records + return len(records), 0 + + def flush(self): + return None + + class CompactWriter: + def save(self, value): + self.saved = value + + def flush(self): + return None + + publisher = Publisher() + import_startup_history_artifact( + artifact, + publisher=publisher, + state_manager=CompactWriter(), + manifest_manager=CompactWriter(), + expected_config=config, + ) + + assert publisher.records[0]["value_json"] == raw_value + + +def test_artifact_v1_requires_raw_topics_for_import(base_config): + """Artifact v1.0 без raw value не может обещать byte-for-byte импорт.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + del artifact["raw_topics"] + + class Publisher: + def publish_records(self, topic, records): + raise AssertionError("publish must not be called") + + with pytest.raises(ValueError, match="raw_topics"): + import_startup_history_artifact( + artifact, + publisher=Publisher(), + state_manager=None, + manifest_manager=None, + expected_config=config, + ) + + +def test_import_rejects_publisher_without_raw_publish_records(base_config): + """Импорт не должен пересобирать dict, если publisher не умеет raw records.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + artifact["raw_topics"]["browser_events"][0]["value_json"] = ( + '{"event_timestamp":"2026-01-01 00:00:00.000000",' + '"user_domain_id":"user-1","click_id":"click-1","event_id":"event-1"}' + ) + + class DictPublisher: + def publish(self, topic, events): + raise AssertionError("dict publish must not be called") + + def flush(self): + return None + + with pytest.raises(TypeError, match="publish_records"): + import_startup_history_artifact( + artifact, + publisher=DictPublisher(), + state_manager=None, + manifest_manager=None, + expected_config=config, + ) + + +def test_import_rejects_dirty_kafka_topics_before_publish(base_config): + """Повторный импорт не должен дописывать дубли в непустые data-топики.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + + class Publisher: + def publish(self, topic, events): + raise AssertionError("publish must not be called") + + class DirtyTopics: + def assert_data_topics_empty(self): + raise RuntimeError("Kafka data topics are not empty") + + with pytest.raises(RuntimeError, match="not empty"): + import_startup_history_artifact( + artifact, + publisher=Publisher(), + state_manager=None, + manifest_manager=None, + expected_config=config, + topic_inspector=DirtyTopics(), + ) + + +def test_import_rolls_back_kafka_topics_after_partial_publish(base_config): + """При частичной публикации импорт откатывает Kafka-топики до повтора.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + + class PartialPublisher: + def __init__(self): + self.calls = [] + + def publish_records(self, topic, records): + self.calls.append(topic) + if topic == "location_events": + raise RuntimeError("location down") + return len(records), 0 + + class CompactWriter: + def save(self, value): + raise AssertionError("compact topics must not be written") + + class Inspector: + def __init__(self): + self.snapshot = None + self.rollback = None + + def assert_data_topics_empty(self): + return None + + def snapshot_import_topics(self): + self.snapshot = {"topics": ["browser_events", "location_events"]} + return self.snapshot + + def rollback_import_topics(self, snapshot): + self.rollback = snapshot + + publisher = PartialPublisher() + inspector = Inspector() + + with pytest.raises(RuntimeError, match="location down"): + import_startup_history_artifact( + artifact, + publisher=publisher, + state_manager=CompactWriter(), + manifest_manager=CompactWriter(), + expected_config=config, + topic_inspector=inspector, + ) + + assert publisher.calls == ["browser_events", "location_events"] + assert inspector.rollback == inspector.snapshot + + +def test_import_replays_events_and_compact_topics(base_config): + """Импорт пишет события в Kafka и служебные compact-топики, не трогая ClickHouse.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + import_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + artifact = builder.to_artifact( + manifest=build_manifest(config=config, counters=builder.counters, state=state), + state=state, + ) + + class Publisher: + def __init__(self): + self.calls = [] + self.flushed = False + + def publish_records(self, topic, records): + self.calls.append((topic, records)) + return len(records), 0 + + def flush(self): + self.flushed = True + + class CompactWriter: + def __init__(self): + self.saved = None + self.flushed = False + + def save(self, value): + self.saved = value + + def flush(self): + self.flushed = True + + publisher = Publisher() + state_manager = CompactWriter() + manifest_manager = CompactWriter() + + result = import_startup_history_artifact( + artifact, + publisher=publisher, + state_manager=state_manager, + manifest_manager=manifest_manager, + expected_config=config, + ) + + assert [topic for topic, _ in publisher.calls] == [ + "browser_events", + "location_events", + "device_events", + "geo_events", + ] + assert publisher.calls[0][1][0]["value_json"].encode("utf-8") == ( + b'{"event_id": "event-1", "click_id": "click-1", ' + b'"user_domain_id": "user-1", ' + b'"event_timestamp": "2026-01-01 00:00:00.000000"}' + ) + assert result["events"] == 4 + assert state_manager.saved.last_batch_id == state.last_batch_id + assert manifest_manager.saved["state"]["last_batch_id"] == state.last_batch_id + assert publisher.flushed + assert state_manager.flushed + assert manifest_manager.flushed diff --git a/scripts/check_startup_history_manifest.sh b/scripts/check_startup_history_manifest.sh new file mode 100755 index 0000000..06b9ef0 --- /dev/null +++ b/scripts/check_startup_history_manifest.sh @@ -0,0 +1,75 @@ +#!/usr/bin/env bash +# +# Сверяет DM-витрину ClickHouse с manifest портативной стартовой истории. + +set -euo pipefail + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" +ARTIFACT="${ARTIFACT:-/tmp/clickstream-startup-history.json}" + +fail() { + echo "Ошибка: $*" >&2 + exit 1 +} + +clickhouse_datetime_literal() { + local value="$1" + value="${value/T/ }" + value="${value%Z}" + if [[ "${value}" =~ ^(.*)[+-][0-9]{2}:[0-9]{2}$ ]]; then + value="${BASH_REMATCH[1]}" + fi + echo "${value}" +} + +[[ -s "${ARTIFACT}" ]] || fail "артефакт не найден или пуст: ${ARTIFACT}" + +ARTIFACT_DIR="$(dirname "${ARTIFACT}")" +ARTIFACT_FILE="$(basename "${ARTIFACT}")" + +manifest_summary="$(${COMPOSE_BIN} run --rm --no-deps \ + -v "${ARTIFACT_DIR}:/startup-artifacts:ro" \ + generator python -m clickstream_generator.startup_history_artifact manifest-summary \ + --artifact "/startup-artifacts/${ARTIFACT_FILE}")" + +IFS=$'\t' read -r expected_events expected_visits expected_users expected_min_ts expected_max_ts model_t0 model_t_end <<< "${manifest_summary}" + +[[ "${expected_events}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать manifest из артефакта" + +ch_model_t0="$(clickhouse_datetime_literal "${model_t0}")" +ch_model_t_end="$(clickhouse_datetime_literal "${model_t_end}")" + +actual="$(${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --query " +WITH + toDateTime64('${ch_model_t0}', 6) AS t0, + toDateTime64('${ch_model_t_end}', 6) AS t_end +SELECT + count() AS events, + uniqExact(click_id) AS visits, + uniqExact(user_domain_id) AS users, + toString(min(event_ts)) AS min_event_timestamp, + toString(max(event_ts)) AS max_event_timestamp +FROM dm.v_events_enriched +WHERE event_ts >= t0 AND event_ts < t_end +FORMAT TabSeparated")" + +IFS=$'\t' read -r actual_events actual_visits actual_users actual_min_ts actual_max_ts <<< "${actual}" + +[[ "${actual_events}" == "${expected_events}" ]] || fail "events: ClickHouse=${actual_events}, manifest=${expected_events}" +[[ "${actual_visits}" == "${expected_visits}" ]] || fail "visits: ClickHouse=${actual_visits}, manifest=${expected_visits}" +[[ "${actual_users}" == "${expected_users}" ]] || fail "users: ClickHouse=${actual_users}, manifest=${expected_users}" +[[ "${actual_min_ts}" == "${expected_min_ts}" ]] || fail "min_event_timestamp: ClickHouse=${actual_min_ts}, manifest=${expected_min_ts}" +[[ "${actual_max_ts}" == "${expected_max_ts}" ]] || fail "max_event_timestamp: ClickHouse=${actual_max_ts}, manifest=${expected_max_ts}" + +echo "manifest_events=${expected_events}" +echo "manifest_visits=${expected_visits}" +echo "manifest_users=${expected_users}" +echo "manifest_min_event_timestamp=${expected_min_ts}" +echo "manifest_max_event_timestamp=${expected_max_ts}" +echo "ClickHouse совпадает с manifest артефакта." diff --git a/scripts/export_startup_history_artifact.sh b/scripts/export_startup_history_artifact.sh new file mode 100755 index 0000000..f26523f --- /dev/null +++ b/scripts/export_startup_history_artifact.sh @@ -0,0 +1,70 @@ +#!/usr/bin/env bash +# +# Экспортирует стартовую историю в портативный JSON-артефакт. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLEAN_START="${CLEAN_START:-1}" +ARTIFACT="${ARTIFACT:-/tmp/clickstream-startup-history.json}" + +GEN_SEED="${GEN_SEED:-4242}" +GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}" +GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}" +GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE:-UTC}" +GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED:-1}" +GEN_TICK_SECONDS="${GEN_TICK_SECONDS:-60}" +GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN:-60}" +GEN_JITTER_PCT="${GEN_JITTER_PCT:-0}" +GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK:-1}" +GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK:-1000}" + +cd "${REPO_ROOT}" + +ARTIFACT_DIR="$(dirname "${ARTIFACT}")" +ARTIFACT_FILE="$(basename "${ARTIFACT}")" +mkdir -p "${ARTIFACT_DIR}" + +echo "=== Экспорт стартовой истории ===" +echo "ARTIFACT=${ARTIFACT}" +echo "GEN_SEED=${GEN_SEED}" +echo "GEN_MODEL_T0=${GEN_MODEL_T0}" +echo "GEN_MODEL_T_END=${GEN_MODEL_T_END}" +echo "" + +if [[ "${CLEAN_START}" == "1" ]]; then + echo "Шаг 0: очистка Kafka/state перед повторяемым экспортом" + ${COMPOSE_BIN} down -v --remove-orphans +else + echo "Шаг 0: CLEAN_START=0, очистка пропущена" +fi + +echo "Шаг 1: запуск Kafka" +${COMPOSE_BIN} up -d kafka + +echo "Шаг 2: сборка образа генератора" +${COMPOSE_BIN} build generator + +echo "Шаг 3: backfill и запись артефакта" +${COMPOSE_BIN} run --rm --no-deps \ + -v "${ARTIFACT_DIR}:/startup-artifacts" \ + -e GEN_RUN_MODE=backfill \ + -e GEN_STATE_RESET=true \ + -e GEN_STARTUP_HISTORY_ARTIFACT="/startup-artifacts/${ARTIFACT_FILE}" \ + -e GEN_SEED="${GEN_SEED}" \ + -e GEN_MODEL_T0="${GEN_MODEL_T0}" \ + -e GEN_MODEL_T_END="${GEN_MODEL_T_END}" \ + -e GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE}" \ + -e GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED}" \ + -e GEN_TICK_SECONDS="${GEN_TICK_SECONDS}" \ + -e GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN}" \ + -e GEN_JITTER_PCT="${GEN_JITTER_PCT}" \ + -e GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK}" \ + -e GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK}" \ + generator + +test -s "${ARTIFACT}" +echo "Готово: ${ARTIFACT}" diff --git a/scripts/import_startup_history_artifact.sh b/scripts/import_startup_history_artifact.sh new file mode 100755 index 0000000..4ce4eb0 --- /dev/null +++ b/scripts/import_startup_history_artifact.sh @@ -0,0 +1,63 @@ +#!/usr/bin/env bash +# +# Импортирует портативный артефакт стартовой истории в Kafka. + +set -euo pipefail + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO_ROOT="$(cd "${SCRIPT_DIR}/.." && pwd)" + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +ARTIFACT="${ARTIFACT:-/tmp/clickstream-startup-history.json}" + +GEN_SEED="${GEN_SEED:-4242}" +GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}" +GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}" +GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE:-UTC}" +GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED:-1}" +GEN_TICK_SECONDS="${GEN_TICK_SECONDS:-60}" +GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN:-60}" +GEN_JITTER_PCT="${GEN_JITTER_PCT:-0}" +GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK:-1}" +GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK:-1000}" + +cd "${REPO_ROOT}" + +if [[ ! -s "${ARTIFACT}" ]]; then + echo "Ошибка: артефакт не найден или пуст: ${ARTIFACT}" >&2 + exit 1 +fi + +ARTIFACT_DIR="$(dirname "${ARTIFACT}")" +ARTIFACT_FILE="$(basename "${ARTIFACT}")" + +echo "=== Импорт стартовой истории ===" +echo "ARTIFACT=${ARTIFACT}" +echo "GEN_SEED=${GEN_SEED}" +echo "GEN_MODEL_T0=${GEN_MODEL_T0}" +echo "GEN_MODEL_T_END=${GEN_MODEL_T_END}" +echo "" + +echo "Шаг 1: запуск Kafka" +${COMPOSE_BIN} up -d kafka + +echo "Шаг 2: сборка образа генератора" +${COMPOSE_BIN} build generator + +echo "Шаг 3: воспроизведение артефакта в Kafka" +${COMPOSE_BIN} run --rm --no-deps \ + -v "${ARTIFACT_DIR}:/startup-artifacts:ro" \ + -e GEN_SEED="${GEN_SEED}" \ + -e GEN_MODEL_T0="${GEN_MODEL_T0}" \ + -e GEN_MODEL_T_END="${GEN_MODEL_T_END}" \ + -e GEN_MODEL_TIMEZONE="${GEN_MODEL_TIMEZONE}" \ + -e GEN_MODEL_TIME_SPEED="${GEN_MODEL_TIME_SPEED}" \ + -e GEN_TICK_SECONDS="${GEN_TICK_SECONDS}" \ + -e GEN_LAMBDA_BASE_PER_MIN="${GEN_LAMBDA_BASE_PER_MIN}" \ + -e GEN_JITTER_PCT="${GEN_JITTER_PCT}" \ + -e GEN_MIN_EVENTS_PER_TICK="${GEN_MIN_EVENTS_PER_TICK}" \ + -e GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK}" \ + generator python -m clickstream_generator.startup_history_artifact import \ + --artifact "/startup-artifacts/${ARTIFACT_FILE}" + +echo "Готово: артефакт воспроизведён в Kafka."