diff --git a/Makefile b/Makefile index f6c6427..6b96431 100644 --- a/Makefile +++ b/Makefile @@ -46,7 +46,7 @@ generated-history-analytics: # Повторяемая проверка после прогона стартовой истории generated-history-check: - COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh + CHECK_LIVE_SEAM="$${CHECK_LIVE_SEAM:-auto}" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh # Сгенерировать стартовую историю и сохранить портативный артефакт startup-history-export: diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 0d211b5..fb6adb8 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -147,7 +147,7 @@ make generator-logs | `GEN_RUN_MODE` | Режим генератора | `live` | | `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` | | `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | пусто | -| `GEN_STATE_ENABLED` | Сохранять state v2 между рестартами | `true` | +| `GEN_STATE_ENABLED` | Сохранять state v3 между рестартами | `true` | | `GEN_STATE_RESET` | Сбросить state при старте | `false` | В `docker-compose.yml` через окружение переопределяются демо-параметры модели и @@ -177,6 +177,13 @@ CI короткая повторная проверка после уже гот make generated-history-check ``` +После live-продолжения из `T_end` эта же команда автоматически включает проверку +стыка. Для принудительной проверки: + +```bash +CHECK_LIVE_SEAM=1 GEN_LIVE_CHECK_MINUTES=10 make generated-history-check +``` + По умолчанию команда использует быстрый профиль `ci`: 6 часов модельного времени. Историю на 2 суток с суточной волной можно получить одной командой. В live-продолжении `daily-wave` идёт с ×60 и тикает раз в секунду, поэтому 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 ab79c9f..9954e44 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 @@ -174,6 +174,11 @@ state-записей остаётся прежним контрактом воз - `wall_timestamp` — настенная UTC-метка, когда это состояние было сохранено; - `model_time_speed`, `model_timezone` и `model_t0`. +State v3 также хранит для каждого активного визита `base_click_id` — click_id +донора браузерной и source-фактуры из статического сида. При восстановлении +активный визит берёт per-event поля от этого донора; неизвестный донор считается +битым state, а не поводом выбрать запасную фактуру. + При восстановлении живого режима после сбоя модельная точка считается так: ```text diff --git a/generator/KNOWN_ISSUES.md b/generator/KNOWN_ISSUES.md index 31b2c64..3ee5751 100644 --- a/generator/KNOWN_ISSUES.md +++ b/generator/KNOWN_ISSUES.md @@ -3,7 +3,7 @@ > **Статус (2026-06-11):** исторический дефект старой плоской генерации закрыт > для режима `steady-stream`. Генератор строит визиты с общим `click_id`, > монотонным временем событий, путём по страницам воронки, популяцией -> возвращающихся пользователей и состоянием v2 для активных визитов. +> возвращающихся пользователей и состоянием для активных визитов. > > Эта заметка больше не является предупреждением «генератор концептуально > сломан». Она оставлена как учебный разбор старого дефекта и как место для @@ -122,7 +122,7 @@ device_event = {**base_device, "click_id": new_click_id} # user_domain_i времени. 3. **Время событий** внутри визита строго растёт и не прилипает к одному `now()` для всего батча. -4. **Состояние v2** сохраняет популяцию, активные визиты, накопленный бюджет +4. **Состояние** сохраняет популяцию, активные визиты, накопленный бюджет рождения визитов, номер тика и состояние ГПСЧ. После этой переделки поток на длинном окне и при штатных параметрах даёт diff --git a/generator/README.md b/generator/README.md index e076857..388edda 100644 --- a/generator/README.md +++ b/generator/README.md @@ -30,7 +30,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `src/clickstream_generator/intensity.py` | расчёт событийного бюджета тика | | `src/clickstream_generator/runtime.py` | тиковый слой: активные визиты и выпуск созревших событий | | `src/clickstream_generator/kafka_io.py` | Kafka publisher, история batch, Kafka-state и служебные топики | -| `src/clickstream_generator/state.py` | сериализуемое состояние генератора v2 | +| `src/clickstream_generator/state.py` | сериализуемое состояние генератора v3 | | `src/clickstream_generator/metrics.py` | Prometheus-метрики | | `src/clickstream_generator/service.py` | основной цикл сервиса | | `generator.py` | запуск сервиса и совместимый фасад | @@ -102,7 +102,7 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator реальному часу запуска процесса. В режиме `backfill` генератор без сна проходит от `GEN_MODEL_T0` до -`GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v2 +`GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v3 на `T_end` в `generator_state` и пишет manifest в compact-topic `generator_startup_history_manifest`. Live-запуск с теми же настройками использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты. @@ -286,11 +286,12 @@ docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \ ### Как работает -1. После каждого успешного тика состояние v2 сохраняется в `generator_state`. +1. После каждого успешного тика состояние v3 сохраняется в `generator_state`. 2. При старте генератор читает последнее состояние из топика. 3. Если состояние найдено, сервис восстанавливает номер тика, состояние ГПСЧ, популяцию пользователей, накопленный бюджет рождения визитов и активные - визиты. + визиты. Для активного визита state хранит `base_click_id` — донора браузерной + и source-фактуры из статического сида. 4. Если состояния нет или оно невалидно, генератор начинает с чистого листа. Активный визит после простоя до 30 минут продолжается со своими исходными diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py index 1db9454..9f63f51 100644 --- a/generator/src/clickstream_generator/generation.py +++ b/generator/src/clickstream_generator/generation.py @@ -142,13 +142,27 @@ class EventGenerator: user_profile: dict[str, dict] | None = None, ) -> dict[str, list[dict]]: """Генерирует один визит с сохранением связей.""" + batch, _base_click_id = self.generate_visit_batch( + batch_size, + planned_start_at=planned_start_at, + user_profile=user_profile, + ) + return batch + + def generate_visit_batch( + self, + batch_size: int, + planned_start_at: datetime | None = None, + user_profile: dict[str, dict] | None = None, + ) -> tuple[dict[str, list[dict]], str | None]: + """Генерирует визит и возвращает click_id донора фактуры.""" if not self.dictionary.browser_events: return { "browser_events": [], "location_events": [], "device_events": [], "geo_events": [], - } + }, None batch = { "browser_events": [], @@ -158,7 +172,7 @@ class EventGenerator: } if batch_size <= 0: - return batch + return batch, None max_visit_events = min(batch_size, self.config.max_session_events) min_visit_events = min(2, max_visit_events) @@ -181,8 +195,18 @@ class EventGenerator: base_browser_events = self.dictionary.browser_by_click_id[base_click_id][:len(visit_path)] else: base_browser = self.rng.choice(self.dictionary.browser_events) + # Запасная ветка тоже восстановима: state хранит донора, а индекс идёт по кругу. base_click_id = base_browser["click_id"] - base_browser_events = [base_browser for _ in range(len(visit_path))] + source_events = self.dictionary.browser_by_click_id[base_click_id] + if not all( + event["event_id"] in self.dictionary.location_by_event_id + for event in source_events + ): + raise ValueError(f"Fixture base_click_id has incomplete locations: {base_click_id}") + base_browser_events = [ + source_events[event_index % len(source_events)] + for event_index in range(len(visit_path)) + ] base_device = ( user_profile["device"] @@ -229,4 +253,4 @@ class EventGenerator: planned_timestamp += timedelta(seconds=self._visit_pause_seconds()) - return batch + return batch, base_click_id diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py index 3072c14..9d2ce28 100644 --- a/generator/src/clickstream_generator/runtime.py +++ b/generator/src/clickstream_generator/runtime.py @@ -35,6 +35,7 @@ class ActiveVisit: batch: dict[str, list[dict]] timestamps: list[datetime] + base_click_id: str user: UserProfile | None = None next_index: int = 0 @@ -130,7 +131,12 @@ def _datetime_to_state(value: datetime | None) -> str | None: def _timestamp_to_state_offset(started_at: datetime, timestamp: datetime) -> int: - return int((timestamp - started_at).total_seconds() * 1_000_000) + delta = timestamp - started_at + return ( + delta.days * 86_400_000_000 + + delta.seconds * 1_000_000 + + delta.microseconds + ) def _format_event_timestamp(timestamp: datetime) -> str: @@ -217,6 +223,7 @@ class TickStreamGenerator: return { "user_domain_id": visit.user.user_domain_id if visit.user else None, "click_id": browser_events[0]["click_id"], + "base_click_id": visit.base_click_id, "next_index": visit.next_index, "started_at": started_at.isoformat(), "offsets_us": [ @@ -235,7 +242,7 @@ class TickStreamGenerator: resume_model_at: datetime | None = None, restarted_at: datetime | None = None, ) -> None: - """Восстанавливает популяцию и активные визиты из state v2.""" + """Восстанавливает популяцию и активные визиты из state.""" users = [ self._user_from_state(item) for item in state.population @@ -302,6 +309,7 @@ class TickStreamGenerator: ] batch = self._compact_visit_batch( click_id=item["click_id"], + base_click_id=item["base_click_id"], user=user, timestamps=timestamps, page_url_paths=item["page_url_paths"], @@ -309,6 +317,7 @@ class TickStreamGenerator: return ActiveVisit( batch=batch, timestamps=timestamps, + base_click_id=item["base_click_id"], user=user, next_index=item["next_index"], ) @@ -316,24 +325,26 @@ class TickStreamGenerator: def _compact_visit_batch( self, click_id: str, + base_click_id: str, user: UserProfile, timestamps: list[datetime], page_url_paths: list[str], ) -> dict[str, list[dict]]: batch = _empty_batch() - browser_templates = self.generator.dictionary.browser_by_click_id.get( - user.seed_click_id, - self.generator.dictionary.browser_events, - ) + if base_click_id not in self.generator.dictionary.browser_by_click_id: + raise ValueError(f"Unknown fixture base_click_id: {base_click_id}") + browser_templates = self.generator.dictionary.browser_by_click_id[base_click_id] + if not browser_templates: + raise ValueError(f"Fixture base_click_id has no browser events: {base_click_id}") for event_index, (timestamp, page_url_path) in enumerate( zip(timestamps, page_url_paths) ): browser_template = browser_templates[event_index % len(browser_templates)] - location_template = self.generator.dictionary.location_by_event_id.get( - browser_template["event_id"], - self.generator.dictionary.location_events[0], - ) + source_event_id = browser_template["event_id"] + if source_event_id not in self.generator.dictionary.location_by_event_id: + raise ValueError(f"Unknown fixture location event_id: {source_event_id}") + location_template = self.generator.dictionary.location_by_event_id[source_event_id] event_id = _stable_event_id(click_id, event_index) batch["browser_events"].append( { @@ -428,7 +439,7 @@ class TickStreamGenerator: self._pending_visit_births = 0.0 return - visit_batch = self.generator.generate_batch( + visit_batch, base_click_id = self.generator.generate_visit_batch( self.generator.config.max_session_events, planned_start_at=tick_time, user_profile={"device": user.device, "geo": user.geo}, @@ -439,11 +450,18 @@ class TickStreamGenerator: ] if not timestamps: break + if base_click_id is None: + raise ValueError("Generated active visit is missing fixture base_click_id") click_id = visit_batch["browser_events"][0]["click_id"] self.population.start_visit(user, click_id) self.active_visits.append( - ActiveVisit(batch=visit_batch, timestamps=timestamps, user=user) + ActiveVisit( + batch=visit_batch, + timestamps=timestamps, + base_click_id=base_click_id, + user=user, + ) ) self._pending_visit_births -= 1.0 diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index f6a182f..3ac0f81 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -205,7 +205,7 @@ class GeneratorService: state, wall_now_utc: datetime, ) -> datetime: - """Считает модельную точку live-восстановления по state v2.""" + """Считает модельную точку live-восстановления по state.""" wall_now_utc = self._as_aware_utc(wall_now_utc) wall_saved_at = self._as_aware_utc(state.wall_timestamp) idle_seconds = max(0.0, (wall_now_utc - wall_saved_at).total_seconds()) diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py index fc68431..93029df 100644 --- a/generator/src/clickstream_generator/state.py +++ b/generator/src/clickstream_generator/state.py @@ -10,7 +10,7 @@ from zoneinfo import ZoneInfo, ZoneInfoNotFoundError logger = logging.getLogger("generator") -STATE_VERSION = "2.0" +STATE_VERSION = "3.0" def _nested_list_to_tuple(obj): @@ -74,7 +74,7 @@ def _validate_resume_fields(data: dict) -> None: raise ValueError("gen_seed must be an integer or null") -def _validate_v2_payload(data: dict) -> None: +def _validate_v3_payload(data: dict) -> None: _validate_resume_fields(data) population = data.get("population") active_visits = data.get("active_visits") @@ -123,6 +123,7 @@ def _validate_v2_payload(data: dict) -> None: ( "user_domain_id", "click_id", + "base_click_id", "next_index", "started_at", "offsets_us", @@ -132,6 +133,7 @@ def _validate_v2_payload(data: dict) -> None: ) user_domain_id = visit["user_domain_id"] click_id = visit["click_id"] + base_click_id = visit["base_click_id"] offsets = visit["offsets_us"] page_url_paths = visit["page_url_paths"] next_index = visit["next_index"] @@ -139,6 +141,8 @@ def _validate_v2_payload(data: dict) -> None: raise ValueError(f"active_visits[{index}].user_domain_id is unknown") if not isinstance(click_id, str) or not click_id: raise ValueError(f"active_visits[{index}].click_id must be a string") + if not isinstance(base_click_id, str) or not base_click_id: + raise ValueError(f"active_visits[{index}].base_click_id must be a string") if not isinstance(offsets, list) or not offsets: raise ValueError(f"active_visits[{index}].offsets_us must be a non-empty list") if not all(isinstance(offset, int) and offset >= 0 for offset in offsets): @@ -235,7 +239,7 @@ class GeneratorState: version, ) raise ValueError(f"unsupported state version: {version}") - _validate_v2_payload(data) + _validate_v3_payload(data) model_timestamp = _parse_aware_utc( data["model_timestamp"], "model_timestamp", diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index d8458cd..efd78df 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -18,6 +18,7 @@ from generator import ( generate_tick_batch, hour_factor, ) +from clickstream_generator.runtime import _timestamp_to_state_offset ALLOWED_PAGE_PATHS = { @@ -37,6 +38,24 @@ def _parse_event_timestamps(batch): ] +def _without_fields(batch, excluded_fields): + return { + topic: [ + { + key: value + for key, value in event.items() + if key not in excluded_fields + } + for event in events + ] + for topic, events in batch.items() + } + + +def _without_event_ids(batch): + return _without_fields(batch, {"event_id"}) + + def _page_path_visits(generator, visits_count: int): return [ [ @@ -297,7 +316,7 @@ class TestEventGeneration: def test_tick_stream_state_stays_compact_at_active_session_limit( self, event_dictionary, base_config ): - """State v2 не хранит полные события активных визитов.""" + """State не хранит полные события активных визитов.""" config = replace( base_config, max_session_events=30, @@ -386,6 +405,102 @@ class TestEventGeneration: for timestamp in resumed_timestamps ) + def test_tick_stream_restored_visit_keeps_per_event_fixture( + self, event_dictionary, base_config + ): + """Восстановленный визит продолжает ту же фактуру событий.""" + config = replace(base_config, max_session_events=8) + tick_at = datetime(2026, 6, 11, 12, 0) + resume_at = tick_at + timedelta(hours=2) + + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + stream.generate_tick(event_budget=10, tick_started_at=tick_at) + state = stream.to_state( + tick=1, + rng_state=generator.rng.getstate(), + last_batch_id="batch-1", + last_timestamp=tick_at, + ) + restored_state = type(state).from_dict(json.loads(json.dumps(state.to_dict()))) + + restored_generator = EventGenerator(event_dictionary, config) + restored_stream = TickStreamGenerator(restored_generator) + restored_stream.restore_state(restored_state, resume_model_at=tick_at) + + uninterrupted = stream.generate_tick(event_budget=0, tick_started_at=resume_at) + restored = restored_stream.generate_tick(event_budget=0, tick_started_at=resume_at) + + assert restored["browser_events"] + assert _without_event_ids(restored) == _without_event_ids(uninterrupted) + + def test_tick_stream_restored_visit_keeps_cyclic_fallback_fixture( + self, event_dictionary, base_config + ): + """Короткий донор фактуры после восстановления повторяется так же.""" + config = replace(base_config, max_session_events=30) + tick_at = datetime(2026, 6, 11, 12, 0) + resume_at = tick_at + timedelta(hours=12) + + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + stream.generate_tick(event_budget=100, tick_started_at=tick_at) + state = stream.to_state( + tick=1, + rng_state=generator.rng.getstate(), + last_batch_id="batch-1", + last_timestamp=tick_at, + ) + assert any( + len(event_dictionary.browser_by_click_id[visit["base_click_id"]]) + < len(visit["offsets_us"]) + for visit in state.active_visits + ) + restored_state = type(state).from_dict(json.loads(json.dumps(state.to_dict()))) + + restored_generator = EventGenerator(event_dictionary, config) + restored_stream = TickStreamGenerator(restored_generator) + restored_stream.restore_state(restored_state, resume_model_at=tick_at) + + uninterrupted = stream.generate_tick(event_budget=0, tick_started_at=resume_at) + restored = restored_stream.generate_tick(event_budget=0, tick_started_at=resume_at) + + assert restored["browser_events"] + assert _without_event_ids(restored) == _without_event_ids(uninterrupted) + + def test_state_offset_preserves_microseconds_without_float_rounding(self): + """Смещение state не теряет микросекунду на float-округлении.""" + started_at = datetime(2026, 6, 11, 12, 0, 0) + timestamp = datetime(2026, 6, 11, 12, 8, 34, 130779) + + assert _timestamp_to_state_offset(started_at, timestamp) == 514_130_779 + + def test_tick_stream_restore_rejects_unknown_fixture_donor( + self, event_dictionary, base_config + ): + """Неизвестный донор фактуры в state даёт ошибку восстановления.""" + config = replace(base_config, max_session_events=5) + tick_at = datetime(2026, 6, 11, 12, 0) + + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + stream.generate_tick(event_budget=10, tick_started_at=tick_at) + state = stream.to_state( + tick=1, + rng_state=generator.rng.getstate(), + last_batch_id="batch-1", + last_timestamp=tick_at, + ) + state_data = json.loads(json.dumps(state.to_dict())) + state_data["active_visits"][0]["base_click_id"] = "missing-click-id" + restored_state = type(state).from_dict(state_data) + + restored_generator = EventGenerator(event_dictionary, config) + restored_stream = TickStreamGenerator(restored_generator) + + with pytest.raises(ValueError, match="Unknown fixture base_click_id"): + restored_stream.restore_state(restored_state, resume_model_at=tick_at) + def test_tick_stream_restored_after_long_idle_closes_overdue_visit_without_replay( self, event_dictionary, base_config ): diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 7dfeb9e..f02b238 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -926,11 +926,11 @@ class TestGeneratorServiceBackfill: } -class TestGeneratorServiceStateV2: - """Тесты подключения state v2 к сервисному запуску.""" +class TestGeneratorServiceState: + """Тесты подключения state к сервисному запуску.""" - def test_start_restores_tick_stream_state_v2(self, base_config, event_dictionary): - """Сервис восстанавливает популяцию и активные визиты из state v2.""" + def test_start_restores_tick_stream_state(self, base_config, event_dictionary): + """Сервис восстанавливает популяцию и активные визиты из state.""" source_generator = EventGenerator(event_dictionary, base_config) source_stream = TickStreamGenerator(source_generator) tick_at = datetime.now(timezone.utc).replace(tzinfo=None) @@ -1073,8 +1073,8 @@ class TestGeneratorServiceStateV2: assert service._model_time == model_t_end assert service.stream.active_visit_count == source_stream.active_visit_count - def test_save_state_writes_tick_stream_state_v2(self, base_config): - """Сервис сохраняет v2-снимок тикового слоя.""" + def test_save_state_writes_tick_stream_state(self, base_config): + """Сервис сохраняет снимок тикового слоя.""" service = GeneratorService(base_config) service.state_manager = MagicMock() tick_at = datetime.now(timezone.utc).replace(tzinfo=None) @@ -1085,7 +1085,7 @@ class TestGeneratorServiceStateV2: service._save_state("batch-1") saved_state = service.state_manager.save.call_args.args[0] - assert saved_state.version == "2.0" + assert saved_state.version == "3.0" assert saved_state.model_timestamp == tick_at.replace(tzinfo=timezone.utc) assert saved_state.wall_timestamp.tzinfo is not None assert saved_state.model_time_speed == base_config.model_time_speed @@ -1162,13 +1162,13 @@ class TestGeneratorServiceStateV2: state_manager.load.assert_not_called() assert service._tick == 0 - def test_invalid_restored_v2_state_starts_fresh(self, base_config, caplog): - """Сервис не падает, если v2 state ссылается на неизвестный профиль.""" + def test_invalid_restored_state_starts_fresh(self, base_config, caplog): + """Сервис не падает, если state ссылается на неизвестный профиль.""" state_manager = MagicMock() state_manager.load.return_value = GeneratorState( tick=9, rng_state=random.Random(42).getstate(), - last_batch_id="bad-v2", + last_batch_id="bad-state", last_timestamp=datetime.now(timezone.utc), model_timestamp=base_config.model_t0, wall_timestamp=datetime.now(timezone.utc), diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index fc283a5..1a0cebc 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -18,8 +18,8 @@ def _make_valid_rng_state(seed: int = 42): return rng.getstate() -def _make_valid_v2_state_data() -> dict: - """Создаёт минимальный валидный state v2 для тестов загрузки.""" +def _make_valid_state_data() -> dict: + """Создаёт минимальный валидный state для тестов загрузки.""" return { "tick": 42, "rng_state": list(_make_valid_rng_state(42)), @@ -31,7 +31,7 @@ def _make_valid_v2_state_data() -> dict: "model_timezone": "UTC", "model_t0": "2026-01-01T00:00:00+00:00", "gen_seed": 42, - "version": "2.0", + "version": "3.0", "population": [ { "user_domain_id": "user-1", @@ -44,6 +44,7 @@ def _make_valid_v2_state_data() -> dict: { "user_domain_id": "user-1", "click_id": "visit-1", + "base_click_id": "seed-1", "next_index": 1, "started_at": "2026-06-11T12:00:00", "offsets_us": [0, 60_000_000], @@ -66,7 +67,7 @@ def _minimal_population() -> list[dict]: def _with_resume_fields(data: dict) -> dict: - """Добавляет обязательные поля state v2, не связанные с проверяемой ошибкой.""" + """Добавляет обязательные поля state, не связанные с проверяемой ошибкой.""" return { **data, "model_timestamp": "2026-01-01T10:00:00+00:00", @@ -97,7 +98,7 @@ class TestGeneratorState: model_timezone="UTC", model_t0=now, gen_seed=42, - version="2.0", + version="3.0", ) assert state.tick == 42 @@ -110,7 +111,7 @@ class TestGeneratorState: assert state.model_timezone == "UTC" assert state.model_t0 == now assert state.gen_seed == 42 - assert state.version == "2.0" + assert state.version == "3.0" def test_default_version(self): """Новые состояния по умолчанию пишутся в версии 2.""" @@ -124,7 +125,7 @@ class TestGeneratorState: last_timestamp=now, ) - assert state.version == "2.0" + assert state.version == "3.0" def test_to_dict_serialization(self): """Сериализация в словарь (JSON-safe, без pickle).""" @@ -150,7 +151,7 @@ class TestGeneratorState: assert data["model_timezone"] == "UTC" assert data["model_t0"] == now.isoformat() assert data["gen_seed"] is None - assert data["version"] == "2.0" + assert data["version"] == "3.0" # Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle) assert "rng_state" in data @@ -225,8 +226,8 @@ class TestGeneratorState: assert next_values == values_after - def test_version_2_roundtrip_keeps_population_and_active_visits(self): - """State v2 хранит популяцию и активные визиты в JSON.""" + def test_state_roundtrip_keeps_population_and_active_visits(self): + """State хранит популяцию и активные визиты в JSON.""" rng_state = _make_valid_rng_state(42) state = GeneratorState( tick=7, @@ -239,7 +240,7 @@ class TestGeneratorState: model_timezone="Europe/Moscow", model_t0=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), gen_seed=42, - version="2.0", + version="3.0", population=[ { "user_domain_id": "user-1", @@ -252,6 +253,7 @@ class TestGeneratorState: { "user_domain_id": "user-1", "click_id": "visit-1", + "base_click_id": "seed-1", "next_index": 1, "started_at": "2026-06-11T12:00:00", "offsets_us": [0, 60_000_000], @@ -263,7 +265,7 @@ class TestGeneratorState: restored = GeneratorState.from_dict(json.loads(json.dumps(state.to_dict()))) - assert restored.version == "2.0" + assert restored.version == "3.0" assert restored.tick == state.tick assert restored.rng_state == rng_state assert restored.population == state.population @@ -281,7 +283,7 @@ class TestGeneratorStateValidation: """Тесты валидации состояния и graceful degradation.""" def test_from_dict_missing_version_raises(self): - """from_dict выбрасывает исключение при отсутствии версии v2.""" + """from_dict выбрасывает исключение при отсутствии версии.""" data = { "tick": 42, "last_batch_id": "test", @@ -298,7 +300,7 @@ class TestGeneratorStateValidation: "rng_state": "not_a_tuple", "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], }) @@ -313,7 +315,7 @@ class TestGeneratorStateValidation: "rng_state": [1], # Слишком короткий "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], }) @@ -328,7 +330,7 @@ class TestGeneratorStateValidation: "rng_state": [999, [1, 2, 3], None], # Невалидный state "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], }) @@ -343,7 +345,7 @@ class TestGeneratorStateValidation: "rng_state": "invalid", "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], }) @@ -365,7 +367,7 @@ class TestGeneratorStateValidation: "model_timezone": "UTC", "model_t0": "2026-01-01T00:00:00+00:00", "gen_seed": 42, - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], } @@ -378,7 +380,7 @@ class TestGeneratorStateValidation: def test_from_dict_safe_returns_none_on_invalid_gen_seed(self): """gen_seed в JSON state должен быть числом или null.""" - data = _make_valid_v2_state_data() + data = _make_valid_state_data() data["gen_seed"] = "42" result = GeneratorState.from_dict_safe(data) @@ -387,7 +389,7 @@ class TestGeneratorStateValidation: def test_from_dict_safe_returns_none_on_bool_model_time_speed(self): """model_time_speed не принимает bool как числовую скорость.""" - data = _make_valid_v2_state_data() + data = _make_valid_state_data() data["model_time_speed"] = True result = GeneratorState.from_dict_safe(data) @@ -395,14 +397,14 @@ class TestGeneratorStateValidation: assert result is None def test_from_dict_safe_returns_none_without_model_resume_fields(self): - """State v2 без связки модельного и настенного времени несовместим.""" + """State без связки модельного и настенного времени несовместим.""" rng = random.Random(42) data = { "tick": 42, "rng_state": list(rng.getstate()), "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": _minimal_population(), "active_visits": [], } @@ -497,7 +499,7 @@ class TestKafkaStateManager: rng_state=_make_valid_rng_state(100), last_batch_id="xyz789", last_timestamp=now, - version="2.0", + version="3.0", population=_minimal_population(), ) @@ -540,8 +542,8 @@ class TestKafkaStateManager: # Должно вернуть None из-за невалидного state assert result is None - def test_load_invalid_v2_nested_state_returns_none(self, caplog): - """Битое state v2 с валидным rng_state даёт чистый старт.""" + def test_load_invalid_nested_state_returns_none(self, caplog): + """Битое state с валидным rng_state даёт чистый старт.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: @@ -555,12 +557,13 @@ class TestKafkaStateManager: "rng_state": list(_make_valid_rng_state(42)), "last_batch_id": "bad-v2", "last_timestamp": "2026-06-11T12:00:00+00:00", - "version": "2.0", + "version": "3.0", "population": [{"user_domain_id": "user-1"}], "active_visits": [ { "user_domain_id": "user-1", "click_id": "visit-1", + "base_click_id": "seed-1", "next_index": 1, "started_at": "2026-06-11T12:00:00", "offsets_us": [0], @@ -579,15 +582,15 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_empty_population_v2_returns_none(self, caplog): - """Пустая популяция в state v2 не восстанавливается.""" + def test_load_empty_population_returns_none(self, caplog): + """Пустая популяция в state не восстанавливается.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - bad_state = _make_valid_v2_state_data() + bad_state = _make_valid_state_data() bad_state["population"] = [] bad_state["active_visits"] = [] @@ -606,15 +609,15 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_bad_pending_births_v2_returns_none(self, caplog): - """Нечисловой pending_visit_births в state v2 не восстанавливается.""" + def test_load_bad_pending_births_returns_none(self, caplog): + """Нечисловой pending_visit_births в state не восстанавливается.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - bad_state = _make_valid_v2_state_data() + bad_state = _make_valid_state_data() bad_state["pending_visit_births"] = "bad" mock_message = MagicMock() @@ -632,7 +635,33 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_active_visit_with_unknown_user_v2_returns_none(self, caplog): + def test_load_active_visit_without_base_click_id_returns_none(self, caplog): + """Активный визит без донора фактуры не восстанавливается.""" + with patch("generator._import_kafka") as mock_import, \ + patch("kafka.KafkaConsumer") as mock_consumer_class: + + mock_producer_class = MagicMock() + mock_import.return_value = (mock_producer_class, None) + + bad_state = _make_valid_state_data() + del bad_state["active_visits"][0]["base_click_id"] + + mock_message = MagicMock() + mock_message.key = b"default" + mock_message.value = bad_state + + mock_consumer = MagicMock() + mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message])) + mock_consumer_class.return_value = mock_consumer + + manager = KafkaStateManager("kafka:29092") + with caplog.at_level(logging.WARNING, logger="generator"): + result = manager.load() + + assert result is None + assert "base_click_id" in caplog.text + + def test_load_active_visit_with_unknown_user_returns_none(self, caplog): """Активный визит должен ссылаться на пользователя из популяции.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: @@ -640,7 +669,7 @@ class TestKafkaStateManager: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - bad_state = _make_valid_v2_state_data() + bad_state = _make_valid_state_data() bad_state["active_visits"][0]["user_domain_id"] = "missing-user" mock_message = MagicMock() @@ -658,7 +687,7 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_active_visit_with_conflicting_click_id_v2_returns_none(self, caplog): + def test_load_active_visit_with_conflicting_click_id_returns_none(self, caplog): """active_click_id пользователя не должен противоречить визиту.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: @@ -666,7 +695,7 @@ class TestKafkaStateManager: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - bad_state = _make_valid_v2_state_data() + bad_state = _make_valid_state_data() bad_state["population"][0]["active_click_id"] = "other-visit" mock_message = MagicMock() @@ -684,7 +713,7 @@ class TestKafkaStateManager: assert result is None assert "Invalid state" in caplog.text - def test_load_population_ghost_active_click_id_v2_returns_none(self, caplog): + def test_load_population_ghost_active_click_id_returns_none(self, caplog): """active_click_id пользователя должен иметь соответствующий активный визит.""" with patch("generator._import_kafka") as mock_import, \ patch("kafka.KafkaConsumer") as mock_consumer_class: @@ -692,7 +721,7 @@ class TestKafkaStateManager: mock_producer_class = MagicMock() mock_import.return_value = (mock_producer_class, None) - bad_state = _make_valid_v2_state_data() + bad_state = _make_valid_state_data() bad_state["population"][0]["active_click_id"] = "ghost" bad_state["active_visits"] = [] diff --git a/scripts/check_generated_analytics.sh b/scripts/check_generated_analytics.sh index 83c87b0..00d82e9 100644 --- a/scripts/check_generated_analytics.sh +++ b/scripts/check_generated_analytics.sh @@ -12,6 +12,8 @@ CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" 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}" REQUIRE_SUPERSET="${REQUIRE_SUPERSET:-1}" +CHECK_LIVE_SEAM="${CHECK_LIVE_SEAM:-auto}" +GEN_LIVE_CHECK_MINUTES="${GEN_LIVE_CHECK_MINUTES:-10}" fail() { echo "Ошибка: $*" >&2 @@ -46,6 +48,13 @@ pg_query() { CH_MODEL_T0="$(clickhouse_datetime_literal "${GEN_MODEL_T0}")" CH_MODEL_T_END="$(clickhouse_datetime_literal "${GEN_MODEL_T_END}")" +[[ "${GEN_LIVE_CHECK_MINUTES}" =~ ^[0-9]+$ ]] \ + || fail "GEN_LIVE_CHECK_MINUTES должен быть целым числом минут" +case "${CHECK_LIVE_SEAM}" in + 0|1|auto) ;; + *) fail "CHECK_LIVE_SEAM должен быть 0, 1 или auto" ;; +esac + echo "=== Проверка ClickHouse: данные генерации в DM ===" stats_query=" @@ -261,6 +270,181 @@ echo "confirmation=${contains_confirmation}" echo "monotonic_ok=${contains_monotonic_ok}" echo "confirmation_share=${contains_confirmation_share}" +should_check_live_seam="${CHECK_LIVE_SEAM}" +if [[ "${should_check_live_seam}" == "auto" ]]; then + live_rows_query=" +WITH + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end +SELECT count() +FROM dds.event +WHERE event_ts >= t_end AND event_ts < t_live_end +FORMAT TabSeparated" + live_rows="$(ch_query "${live_rows_query}")" + [[ "${live_rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать live-строки DDS после GEN_MODEL_T_END" + if (( live_rows > 0 )); then + should_check_live_seam="1" + else + should_check_live_seam="0" + fi +fi + +if [[ "${should_check_live_seam}" == "1" ]]; then + echo "" + echo "=== Проверка ClickHouse: стык backfill/live в DDS ===" + + seam_duplicates_query=" +WITH + toDateTime64('${CH_MODEL_T0}', 6) AS t0, + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end +SELECT + count() AS events, + uniqExact(event_id) AS unique_events, + events - unique_events AS duplicate_events +FROM dds.event +WHERE event_ts >= t0 AND event_ts < t_live_end +FORMAT TabSeparated" + + seam_duplicates="$(ch_query "${seam_duplicates_query}")" + IFS=$'\t' read -r seam_events seam_unique_events seam_duplicate_events <<< "${seam_duplicates}" + + [[ "${seam_events}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать события DDS на стыке" + (( seam_events > 0 )) || fail "нет событий DDS в диапазоне проверки стыка" + (( seam_duplicate_events == 0 )) || fail "на стыке найдены дубли event_id: ${seam_duplicate_events}" + + echo "events=${seam_events}" + echo "unique_events=${seam_unique_events}" + echo "duplicate_events=${seam_duplicate_events}" + + seam_context_query=" +WITH + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end, + crossing AS ( + SELECT + e.click_id, + min(e.event_ts) AS first_ts, + max(e.event_ts) AS last_ts, + groupUniqArray(coalesce(toString(e.browser_name), '__NULL__')) AS browser_names, + groupUniqArray(coalesce(toString(e.browser_language), '__NULL__')) AS browser_languages, + groupUniqArray(coalesce(toString(e.browser_user_agent), '__NULL__')) AS browser_user_agents, + groupUniqArray(coalesce(toString(e.referer_url), '__NULL__')) AS referer_urls, + groupUniqArray(coalesce(toString(e.referer_medium), '__NULL__')) AS referer_mediums, + groupUniqArray(coalesce(toString(e.utm_medium), '__NULL__')) AS utm_mediums, + groupUniqArray(coalesce(toString(e.utm_source), '__NULL__')) AS utm_sources, + groupUniqArray(coalesce(toString(e.utm_content), '__NULL__')) AS utm_contents, + groupUniqArray(coalesce(toString(e.utm_campaign), '__NULL__')) AS utm_campaigns, + groupUniqArray(coalesce(toString(c.device_type), '__NULL__')) AS device_types, + groupUniqArray(coalesce(toString(c.os_name), '__NULL__')) AS os_names, + groupUniqArray(coalesce(toString(c.geo_country), '__NULL__')) AS geo_countries + FROM dds.event AS e + LEFT JOIN dds.click AS c ON c.click_id = e.click_id + WHERE e.event_ts >= t_end - INTERVAL 30 MINUTE + AND e.event_ts < t_live_end + GROUP BY e.click_id + HAVING first_ts < t_end AND last_ts >= t_end + ) +SELECT + count() AS crossing_visits, + countIf( + length(browser_names) = 1 + AND length(browser_languages) = 1 + AND length(browser_user_agents) = 1 + AND length(referer_urls) = 1 + AND length(referer_mediums) = 1 + AND length(utm_mediums) = 1 + AND length(utm_sources) = 1 + AND length(utm_contents) = 1 + AND length(utm_campaigns) = 1 + ) AS per_event_homogeneous_visits, + countIf( + length(device_types) = 1 + AND length(os_names) = 1 + AND length(geo_countries) = 1 + ) AS click_context_homogeneous_visits +FROM crossing +FORMAT TabSeparated" + + seam_context="$(ch_query "${seam_context_query}")" + IFS=$'\t' read -r crossing_visits per_event_homogeneous_visits click_context_homogeneous_visits <<< "${seam_context}" + + [[ "${crossing_visits}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать переходящие визиты" + (( crossing_visits > 0 )) || fail "нет визитов, переходящих через GEN_MODEL_T_END" + (( per_event_homogeneous_visits == crossing_visits )) \ + || fail "per-event фактура меняется на стыке: ${per_event_homogeneous_visits}/${crossing_visits}" + (( click_context_homogeneous_visits == crossing_visits )) \ + || fail "device/os/geo меняются на стыке: ${click_context_homogeneous_visits}/${crossing_visits}" + + echo "crossing_visits=${crossing_visits}" + echo "per_event_homogeneous_visits=${per_event_homogeneous_visits}" + echo "click_context_homogeneous_visits=${click_context_homogeneous_visits}" + + ods_context_query=" +WITH + toDateTime64('${CH_MODEL_T_END}', 6) AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end, + crossing AS ( + SELECT + click_id + FROM dds.event + WHERE event_ts >= t_end - INTERVAL 30 MINUTE + AND event_ts < t_live_end + GROUP BY click_id + HAVING min(event_ts) < t_end AND max(event_ts) >= t_end + ), + device_conflicts AS ( + SELECT + d.click_id + FROM ods.device_by_click AS d + INNER JOIN crossing AS c ON c.click_id = d.click_id + WHERE d.click_id IS NOT NULL + GROUP BY d.click_id + HAVING + uniqExact(coalesce(toString(device_type), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(device_is_mobile), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(os_name), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(os), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(os_timezone), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(user_domain_id), '__NULL__')) > 1 + ), + geo_conflicts AS ( + SELECT + g.click_id + FROM ods.geo_by_click AS g + INNER JOIN crossing AS c ON c.click_id = g.click_id + WHERE g.click_id IS NOT NULL + GROUP BY g.click_id + HAVING + uniqExact(coalesce(toString(geo_country), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(geo_timezone), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(geo_region_name), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(geo_latitude), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(geo_longitude), '__NULL__')) > 1 + OR uniqExact(coalesce(toString(ip_address), '__NULL__')) > 1 + ) +SELECT + (SELECT count() FROM crossing) AS crossing_visits, + (SELECT count() FROM device_conflicts) AS ods_device_conflicts, + (SELECT count() FROM geo_conflicts) AS ods_geo_conflicts +FORMAT TabSeparated" + + ods_context="$(ch_query "${ods_context_query}")" + IFS=$'\t' read -r ods_crossing_visits ods_device_conflicts ods_geo_conflicts <<< "${ods_context}" + + [[ "${ods_crossing_visits}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать ODS-проверку стыка" + (( ods_crossing_visits == crossing_visits )) \ + || fail "ODS и DDS нашли разное число переходящих визитов: ${ods_crossing_visits}/${crossing_visits}" + (( ods_device_conflicts == 0 )) || fail "ODS device конфликтует на стыке: ${ods_device_conflicts}" + (( ods_geo_conflicts == 0 )) || fail "ODS geo конфликтует на стыке: ${ods_geo_conflicts}" + + echo "ods_device_conflicts=${ods_device_conflicts}" + echo "ods_geo_conflicts=${ods_geo_conflicts}" +else + echo "" + echo "Проверка стыка backfill/live пропущена: CHECK_LIVE_SEAM=${CHECK_LIVE_SEAM}, live_rows=${live_rows:-0}" +fi + echo "" echo "=== Проверка ClickHouse: основные DM-витрины не пустые ==="