From 640050e40958ede554ca27e5da2f8ce5f167f609 Mon Sep 17 00:00:00 2001 From: Dmitry Dementiev Date: Thu, 11 Jun 2026 18:01:57 +0300 Subject: [PATCH] =?UTF-8?q?feat(generator):=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=BE=20=D1=81=D0=BE=D1=81=D1=82=D0=BE?= =?UTF-8?q?=D1=8F=D0=BD=D0=B8=D0=B5=20=D0=B2=D0=B5=D1=80=D1=81=D0=B8=D0=B8?= =?UTF-8?q?=202=20=D0=B4=D0=BB=D1=8F=20=D1=80=D0=B5=D1=81=D1=82=D0=B0?= =?UTF-8?q?=D1=80=D1=82=D0=BE=D0=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - генератор должен переживать рестарт без потери популяции пользователей и коротко прерванных активных визитов. - Что: - добавлен компактный state v2 для популяции, активных визитов и остатка бюджета рождений. - сервис генератора переведён на единый тиковый поток с сохранением и восстановлением состояния. - добавлена безопасная деградация для старого state v1 и битого state v2. - покрыты короткий и долгий простой, reset состояния и валидация вложенного state. - Проверка: - uv run --with-requirements generator/requirements.txt pytest generator/tests -q. - git diff --check. --- .../issues/06-state-v2-and-restart.md | 27 +- .../src/clickstream_generator/runtime.py | 218 ++++++++++++ .../src/clickstream_generator/service.py | 32 +- generator/src/clickstream_generator/state.py | 130 ++++++- generator/tests/test_generation.py | 142 ++++++++ generator/tests/test_service.py | 120 ++++++- generator/tests/test_state.py | 330 +++++++++++++++++- 7 files changed, 964 insertions(+), 35 deletions(-) diff --git a/.scratch/feature-data-generator/issues/06-state-v2-and-restart.md b/.scratch/feature-data-generator/issues/06-state-v2-and-restart.md index 4b412c6..084d5cc 100644 --- a/.scratch/feature-data-generator/issues/06-state-v2-and-restart.md +++ b/.scratch/feature-data-generator/issues/06-state-v2-and-restart.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Состояние версии 2 и рестарт @@ -17,16 +17,16 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Состояние версии 2 сериализуется в JSON и восстанавливает тик, ГПСЧ, +- [x] Состояние версии 2 сериализуется в JSON и восстанавливает тик, ГПСЧ, популяцию пользователей и активные визиты. -- [ ] Старое состояние версии 1 или битое состояние не валит генератор: +- [x] Старое состояние версии 1 или битое состояние не валит генератор: фиксируется предупреждение, генератор начинает с чистого листа. -- [ ] После простоя не больше 30 минут активный визит продолжается, а созревшие +- [x] После простоя не больше 30 минут активный визит продолжается, а созревшие события досылаются со своими исходными запланированными метками времени. -- [ ] После долгого простоя просроченный активный визит закрывается без досылки +- [x] После долгого простоя просроченный активный визит закрывается без досылки остатка. -- [ ] Популяция пользователей переживает простой любой длины. -- [ ] `GEN_STATE_RESET=true` явно сбрасывает состояние, как и раньше. +- [x] Популяция пользователей переживает простой любой длины. +- [x] `GEN_STATE_RESET=true` явно сбрасывает состояние, как и раньше. ## Blocked by @@ -40,3 +40,16 @@ Status: ready-for-agent Преждевременный кодовый задел state v2 был откатан, чтобы он не попал в коммиты следующих задач. Возвращаться к этой задаче нужно после реализации зависимостей и заново проверять все acceptance criteria на актуальной модели генератора. + +2026-06-11: реализовано состояние версии 2 для `TickStreamGenerator`: в снимок +попадают тик, ГПСЧ, компактная популяция пользователей, компактные активные +визиты и остаток бюджета рождений визитов. `GeneratorService` теперь владеет +одним тиковым потоком, сохраняет его снимок и восстанавливает его при старте. +Короткий простой продолжает активные визиты и досылает созревшие события с +исходными метками времени; долгий простой закрывает просроченные визиты без +досылки остатка; популяция сохраняется. + +Старое state v1 и битое v2-state деградируют в чистый старт с предупреждением; +`GEN_STATE_RESET=true` по-прежнему пропускает загрузку состояния. Проверка: +`uv run --with-requirements generator/requirements.txt pytest generator/tests -q` +— 112 passed. diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py index 0f77c10..1f0fe25 100644 --- a/generator/src/clickstream_generator/runtime.py +++ b/generator/src/clickstream_generator/runtime.py @@ -1,13 +1,16 @@ """Тиковый слой генератора с активными визитами между вызовами.""" +import uuid from dataclasses import dataclass from datetime import datetime, timedelta, timezone from weakref import WeakKeyDictionary from clickstream_generator.generation import EXPECTED_VISIT_EVENTS, EventGenerator +from clickstream_generator.state import GeneratorState TOPICS = ("browser_events", "location_events", "device_events", "geo_events") +RESTART_VISIT_GRACE = timedelta(minutes=30) @dataclass @@ -122,6 +125,22 @@ def _parse_timestamp(value: str) -> datetime: return datetime.fromisoformat(value.replace(" ", "T")) +def _datetime_to_state(value: datetime | None) -> str | None: + return value.isoformat() if value is not None else None + + +def _timestamp_to_state_offset(started_at: datetime, timestamp: datetime) -> int: + return int((timestamp - started_at).total_seconds() * 1_000_000) + + +def _format_event_timestamp(timestamp: datetime) -> str: + return timestamp.strftime("%Y-%m-%d %H:%M:%S.%f") + + +def _stable_event_id(click_id: str, event_index: int) -> str: + return str(uuid.uuid5(uuid.NAMESPACE_URL, f"{click_id}:{event_index}")) + + def _normalize_tick_time(tick_started_at: datetime | None) -> datetime: tick_time = tick_started_at or datetime.now(timezone.utc) if tick_time.tzinfo is not None: @@ -146,6 +165,205 @@ class TickStreamGenerator: def population_user_ids(self) -> set[str]: return {user.user_domain_id for user in self.population.users} + @property + def active_visit_count(self) -> int: + return len(self.active_visits) + + def to_state( + self, + tick: int, + rng_state: tuple, + last_batch_id: str, + last_timestamp: datetime, + ) -> GeneratorState: + """Возвращает JSON-сериализуемый снимок тикового слоя.""" + return GeneratorState( + tick=tick, + rng_state=rng_state, + last_batch_id=last_batch_id, + last_timestamp=last_timestamp, + population=[ + { + "user_domain_id": user.user_domain_id, + "seed_click_id": user.seed_click_id, + "active_click_id": user.active_click_id, + "last_finished_at": _datetime_to_state(user.last_finished_at), + } + for user in self.population.users + ], + active_visits=[ + self._visit_to_state(visit) + for visit in self.active_visits + ], + pending_visit_births=self._pending_visit_births, + ) + + def _visit_to_state(self, visit: ActiveVisit) -> dict: + started_at = visit.timestamps[0] + browser_events = visit.batch["browser_events"] + location_events = visit.batch["location_events"] + return { + "user_domain_id": visit.user.user_domain_id if visit.user else None, + "click_id": browser_events[0]["click_id"], + "next_index": visit.next_index, + "started_at": started_at.isoformat(), + "offsets_us": [ + _timestamp_to_state_offset(started_at, timestamp) + for timestamp in visit.timestamps + ], + "page_url_paths": [ + event["page_url_path"] + for event in location_events + ], + } + + def restore_state( + self, + state: GeneratorState, + restarted_at: datetime | None = None, + ) -> None: + """Восстанавливает популяцию и активные визиты из state v2.""" + users = [ + self._user_from_state(item) + for item in state.population + ] + users_by_id = {user.user_domain_id: user for user in users} + + self.population.users = users + self.active_visits = [] + restarted_time = ( + _normalize_tick_time(restarted_at) + if restarted_at is not None + else None + ) + for item in state.active_visits: + visit = self._visit_from_state(item, users_by_id) + if self._is_overdue_after_restart(visit, restarted_time): + self.population.finish_visit( + visit.user, + self._last_released_at(visit, state.last_timestamp), + ) + continue + self.active_visits.append(visit) + self._pending_visit_births = state.pending_visit_births + + def _user_from_state(self, item: dict) -> UserProfile: + seed_click_id = item["seed_click_id"] + if seed_click_id not in self.generator.dictionary.device_by_click_id: + raise ValueError(f"Unknown user seed_click_id: {seed_click_id}") + if seed_click_id not in self.generator.dictionary.geo_by_click_id: + raise ValueError(f"Unknown user geo seed_click_id: {seed_click_id}") + + user_domain_id = item["user_domain_id"] + device = { + **self.generator.dictionary.device_by_click_id[seed_click_id], + "user_domain_id": user_domain_id, + } + geo = self.generator.dictionary.geo_by_click_id[seed_click_id] + return UserProfile( + user_domain_id=user_domain_id, + seed_click_id=seed_click_id, + device=device, + geo=geo, + active_click_id=item.get("active_click_id"), + last_finished_at=( + _parse_timestamp(item["last_finished_at"]) + if item.get("last_finished_at") + else None + ), + ) + + def _visit_from_state( + self, + item: dict, + users_by_id: dict[str, UserProfile], + ) -> ActiveVisit: + user = users_by_id.get(item["user_domain_id"]) + if user is None: + raise ValueError(f"Unknown active visit user: {item['user_domain_id']}") + + started_at = _parse_timestamp(item["started_at"]) + timestamps = [ + started_at + timedelta(microseconds=offset_us) + for offset_us in item["offsets_us"] + ] + batch = self._compact_visit_batch( + click_id=item["click_id"], + user=user, + timestamps=timestamps, + page_url_paths=item["page_url_paths"], + ) + return ActiveVisit( + batch=batch, + timestamps=timestamps, + user=user, + next_index=item["next_index"], + ) + + def _compact_visit_batch( + self, + 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, + ) + + 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], + ) + event_id = _stable_event_id(click_id, event_index) + batch["browser_events"].append( + { + **browser_template, + "event_id": event_id, + "click_id": click_id, + "event_timestamp": _format_event_timestamp(timestamp), + } + ) + batch["location_events"].append( + { + **location_template, + "event_id": event_id, + "page_url": f"http://www.dummywebsite.com{page_url_path}", + "page_url_path": page_url_path, + } + ) + batch["device_events"].append({**user.device, "click_id": click_id}) + batch["geo_events"].append({**user.geo, "click_id": click_id}) + + return batch + + def _last_released_at( + self, + visit: ActiveVisit, + fallback: datetime, + ) -> datetime: + if visit.next_index > 0: + return visit.timestamps[visit.next_index - 1] + if fallback.tzinfo is not None: + return fallback.astimezone(timezone.utc).replace(tzinfo=None) + return fallback + + def _is_overdue_after_restart( + self, + visit: ActiveVisit, + restarted_time: datetime | None, + ) -> bool: + if restarted_time is None or visit.is_finished: + return False + next_timestamp = visit.timestamps[visit.next_index] + return restarted_time - next_timestamp > RESTART_VISIT_GRACE + def generate_tick( self, event_budget: int, diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 2029160..5946f77 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -23,8 +23,7 @@ from clickstream_generator.metrics import ( METRICS_LAST_SUCCESS, METRICS_TICK_DURATION, ) -from clickstream_generator.runtime import generate_tick_batch -from clickstream_generator.state import GeneratorState +from clickstream_generator.runtime import TickStreamGenerator logger = logging.getLogger("generator") @@ -37,6 +36,7 @@ class GeneratorService: self.config = config self.dictionary = EventDictionary.load(config.data_dir) self.generator = EventGenerator(self.dictionary, config) + self.stream = TickStreamGenerator(self.generator) self.publisher: KafkaPublisher | None = None self.history: KafkaBatchHistory | None = None self.state_manager: KafkaStateManager | None = None @@ -71,12 +71,24 @@ class GeneratorService: if not self.config.state_reset: restored_state = self.state_manager.load() if restored_state: - self._tick = restored_state.tick - self.generator.rng.setstate(restored_state.rng_state) - logger.info( - f"Restored state: continuing from tick {self._tick}, " - f"last_batch_id={restored_state.last_batch_id}" - ) + try: + self.generator.rng.setstate(restored_state.rng_state) + self.stream.restore_state( + restored_state, + restarted_at=datetime.now(timezone.utc), + ) + self._tick = restored_state.tick + logger.info( + f"Restored state: continuing from tick {self._tick}, " + f"last_batch_id={restored_state.last_batch_id}" + ) + except Exception as e: + logger.warning( + f"State data was invalid, starting fresh: {e}" + ) + self.generator = EventGenerator(self.dictionary, self.config) + self.stream = TickStreamGenerator(self.generator) + self._tick = 0 else: logger.info("State reset requested, starting fresh") else: @@ -108,7 +120,7 @@ class GeneratorService: return try: - state = GeneratorState( + state = self.stream.to_state( tick=self._tick, rng_state=self.generator.rng.getstate(), last_batch_id=batch_id, @@ -136,7 +148,7 @@ class GeneratorService: logger.info(f"Generating with event budget ~{events_count}") gen_start = time.time() - batch = generate_tick_batch(self.generator, events_count) + batch = self.stream.generate_tick(events_count) gen_duration = time.time() - gen_start pub_start = time.time() diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py index e25ebe3..15ca3cb 100644 --- a/generator/src/clickstream_generator/state.py +++ b/generator/src/clickstream_generator/state.py @@ -2,13 +2,16 @@ import logging import random -from dataclasses import dataclass +from dataclasses import dataclass, field from datetime import datetime logger = logging.getLogger("generator") +STATE_VERSION = "2.0" + + def _nested_list_to_tuple(obj): """Рекурсивно преобразует list в tuple для восстановления RNG state.""" if isinstance(obj, list): @@ -16,6 +19,107 @@ def _nested_list_to_tuple(obj): return obj +def _require_keys(item: dict, keys: tuple[str, ...], label: str) -> None: + missing = [key for key in keys if key not in item] + if missing: + raise ValueError(f"{label} missing fields: {', '.join(missing)}") + + +def _validate_v2_payload(data: dict) -> None: + population = data.get("population") + active_visits = data.get("active_visits") + pending_visit_births = data.get("pending_visit_births", 0.0) + + if not isinstance(population, list): + raise ValueError("population must be a list") + if not population: + raise ValueError("population must be non-empty") + if not isinstance(active_visits, list): + raise ValueError("active_visits must be a list") + if not isinstance(pending_visit_births, int | float): + raise ValueError("pending_visit_births must be a number") + if not 0 <= pending_visit_births < 1_000_000: + raise ValueError("pending_visit_births is out of range") + + users_by_id = {} + for index, user in enumerate(population): + if not isinstance(user, dict): + raise ValueError(f"population[{index}] must be an object") + _require_keys(user, ("user_domain_id", "seed_click_id"), f"population[{index}]") + user_domain_id = user["user_domain_id"] + seed_click_id = user["seed_click_id"] + active_click_id = user.get("active_click_id") + last_finished_at = user.get("last_finished_at") + if not isinstance(user_domain_id, str) or not user_domain_id: + raise ValueError(f"population[{index}].user_domain_id must be a string") + if not isinstance(seed_click_id, str) or not seed_click_id: + raise ValueError(f"population[{index}].seed_click_id must be a string") + if active_click_id is not None and not isinstance(active_click_id, str): + raise ValueError(f"population[{index}].active_click_id must be a string or null") + if last_finished_at is not None: + if not isinstance(last_finished_at, str): + raise ValueError(f"population[{index}].last_finished_at must be a string or null") + datetime.fromisoformat(last_finished_at) + if user_domain_id in users_by_id: + raise ValueError(f"duplicate population user_domain_id: {user_domain_id}") + users_by_id[user_domain_id] = user + + active_visit_pairs = set() + for index, visit in enumerate(active_visits): + if not isinstance(visit, dict): + raise ValueError(f"active_visits[{index}] must be an object") + _require_keys( + visit, + ( + "user_domain_id", + "click_id", + "next_index", + "started_at", + "offsets_us", + "page_url_paths", + ), + f"active_visits[{index}]", + ) + user_domain_id = visit["user_domain_id"] + click_id = visit["click_id"] + offsets = visit["offsets_us"] + page_url_paths = visit["page_url_paths"] + next_index = visit["next_index"] + if not isinstance(user_domain_id, str) or user_domain_id not in users_by_id: + 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(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): + raise ValueError(f"active_visits[{index}].offsets_us must contain non-negative integers") + if offsets != sorted(offsets): + raise ValueError(f"active_visits[{index}].offsets_us must be sorted") + if not isinstance(page_url_paths, list) or len(page_url_paths) != len(offsets): + raise ValueError( + f"active_visits[{index}].page_url_paths must match offsets_us length" + ) + if not all(isinstance(path, str) and path.startswith("/") for path in page_url_paths): + raise ValueError(f"active_visits[{index}].page_url_paths must contain paths") + if not isinstance(next_index, int) or not 0 <= next_index <= len(offsets): + raise ValueError(f"active_visits[{index}].next_index is out of range") + datetime.fromisoformat(visit["started_at"]) + + user_active_click_id = users_by_id[user_domain_id].get("active_click_id") + if user_active_click_id != click_id: + raise ValueError( + f"population active_click_id conflicts with active_visits[{index}]" + ) + active_visit_pairs.add((user_domain_id, click_id)) + + for user_domain_id, user in users_by_id.items(): + active_click_id = user.get("active_click_id") + if active_click_id and (user_domain_id, active_click_id) not in active_visit_pairs: + raise ValueError( + f"population user {user_domain_id} has active_click_id without active visit" + ) + + @dataclass class GeneratorState: """Состояние генератора для восстановления после рестарта.""" @@ -24,7 +128,10 @@ class GeneratorState: rng_state: tuple last_batch_id: str last_timestamp: datetime - version: str = "1.0" + version: str = STATE_VERSION + population: list[dict] = field(default_factory=list) + active_visits: list[dict] = field(default_factory=list) + pending_visit_births: float = 0.0 def to_dict(self) -> dict: """Конвертирует в словарь для JSON-сериализации.""" @@ -34,12 +141,26 @@ class GeneratorState: "last_batch_id": self.last_batch_id, "last_timestamp": self.last_timestamp.isoformat(), "version": self.version, + "population": self.population, + "active_visits": self.active_visits, + "pending_visit_births": self.pending_visit_births, } @classmethod def from_dict(cls, data: dict) -> "GeneratorState": """Создаёт состояние из словаря.""" try: + if not isinstance(data, dict): + raise ValueError("state must be an object") + version = data.get("version", "1.0") + if version != STATE_VERSION: + logger.warning( + "Unsupported generator state version %s, will start fresh", + version, + ) + raise ValueError(f"unsupported state version: {version}") + _validate_v2_payload(data) + rng_state_raw = data.get("rng_state") if not rng_state_raw: logger.warning("State missing rng_state field") @@ -65,7 +186,10 @@ class GeneratorState: last_timestamp=datetime.fromisoformat( data.get("last_timestamp", "1970-01-01T00:00:00+00:00") ), - version=data.get("version", "1.0"), + version=version, + population=data.get("population", []), + active_visits=data.get("active_visits", []), + pending_visit_births=data.get("pending_visit_births", 0.0), ) except Exception as e: logger.warning(f"Invalid state format, will start fresh: {e}") diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index fcdb4e4..2971c57 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -2,6 +2,7 @@ Тесты генерации событий. """ +import json import random import uuid from dataclasses import replace @@ -268,6 +269,147 @@ class TestEventGeneration: assert {event["click_id"] for event in final_tick["browser_events"]} == {click_id} assert later_tick["browser_events"] == [] + def test_tick_stream_state_roundtrip_keeps_population_and_active_visits( + self, event_dictionary, base_config + ): + """Снимок тикового слоя восстанавливает популяцию и активные визиты.""" + generator = EventGenerator(event_dictionary, base_config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + 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, base_config) + restored_stream = TickStreamGenerator(restored_generator) + restored_stream.restore_state(restored_state, restarted_at=tick_at) + + assert restored_stream.population_user_ids == stream.population_user_ids + assert restored_stream.active_visit_count == stream.active_visit_count + + def test_tick_stream_state_stays_compact_at_active_session_limit( + self, event_dictionary, base_config + ): + """State v2 не хранит полные события активных визитов.""" + config = replace( + base_config, + max_session_events=30, + max_active_sessions=200, + population_max=300, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + stream.generate_tick( + event_budget=int(EXPECTED_VISIT_EVENTS * config.max_active_sessions), + 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_bytes = len(json.dumps(state.to_dict())) + + assert stream.active_visit_count == config.max_active_sessions + assert state_bytes < 250_000 + + def test_tick_stream_restored_after_short_idle_releases_due_original_timestamps( + self, event_dictionary, base_config + ): + """После короткого простоя активный визит продолжается со старыми метками.""" + config = replace(base_config, max_session_events=5) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + 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, + ) + visit_state = state.active_visits[0] + started_at = datetime.fromisoformat(visit_state["started_at"]) + planned_timestamps = { + (started_at + timedelta(microseconds=offset_us)).strftime( + "%Y-%m-%d %H:%M:%S.%f" + ) + for offset_us in visit_state["offsets_us"] + } + + restored_generator = EventGenerator(event_dictionary, config) + restored_stream = TickStreamGenerator(restored_generator) + restored_stream.restore_state( + type(state).from_dict(json.loads(json.dumps(state.to_dict()))), + restarted_at=tick_at + timedelta(minutes=29), + ) + + resumed = restored_stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(minutes=29), + ) + resumed_timestamps = [ + event["event_timestamp"] + for event in resumed["browser_events"] + ] + + assert resumed_timestamps + assert set(resumed_timestamps).issubset(planned_timestamps) + assert all( + datetime.fromisoformat(timestamp.replace(" ", "T")) <= tick_at + timedelta(minutes=29) + for timestamp in resumed_timestamps + ) + + def test_tick_stream_restored_after_long_idle_closes_overdue_visit_without_replay( + self, event_dictionary, base_config + ): + """После долгого простоя просроченный визит закрывается без досылки.""" + config = replace(base_config, max_session_events=5) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + stream.generate_tick(event_budget=10, tick_started_at=tick_at) + population_ids = stream.population_user_ids + state = stream.to_state( + tick=1, + rng_state=generator.rng.getstate(), + last_batch_id="batch-1", + last_timestamp=tick_at, + ) + + restored_generator = EventGenerator(event_dictionary, config) + restored_stream = TickStreamGenerator(restored_generator) + restored_stream.restore_state( + type(state).from_dict(json.loads(json.dumps(state.to_dict()))), + restarted_at=tick_at + timedelta(hours=1), + ) + + resumed = restored_stream.generate_tick( + event_budget=0, + tick_started_at=tick_at + timedelta(hours=1), + ) + restored_user = next( + user for user in restored_stream.population.users + if user.user_domain_id == state.active_visits[0]["user_domain_id"] + ) + last_sent_at = datetime.fromisoformat(state.active_visits[0]["started_at"]) + + assert resumed["browser_events"] == [] + assert restored_stream.active_visit_count == 0 + assert restored_stream.population_user_ids == population_ids + assert restored_user.last_finished_at == last_sent_at + def test_tick_stream_drops_births_when_active_limit_is_reached( self, event_dictionary, base_config ): diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 95c92da..66a43a4 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -2,11 +2,21 @@ Тесты для GeneratorService и интеграционных сценариев. """ +import logging +import random +from dataclasses import replace +from datetime import datetime, timezone from unittest.mock import MagicMock, patch import pytest from generator import ( - Config, EventDictionary, GeneratorService, KafkaBatchHistory + Config, + EventDictionary, + EventGenerator, + GeneratorService, + GeneratorState, + KafkaBatchHistory, + TickStreamGenerator, ) @@ -90,3 +100,111 @@ class TestGeneratorServiceDisabled: service.start() assert "disabled" in caplog.text.lower() or "GEN_ENABLED" in caplog.text + + +class TestGeneratorServiceStateV2: + """Тесты подключения state v2 к сервисному запуску.""" + + def test_start_restores_tick_stream_state_v2(self, base_config, event_dictionary): + """Сервис восстанавливает популяцию и активные визиты из state v2.""" + source_generator = EventGenerator(event_dictionary, base_config) + source_stream = TickStreamGenerator(source_generator) + tick_at = datetime.now(timezone.utc).replace(tzinfo=None) + source_stream.generate_tick(event_budget=10, tick_started_at=tick_at) + state = source_stream.to_state( + tick=3, + rng_state=source_generator.rng.getstate(), + last_batch_id="batch-3", + last_timestamp=tick_at, + ) + + state_manager = MagicMock() + state_manager.load.return_value = state + + 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.object(GeneratorService, "_main_loop", return_value=None): + + service = GeneratorService(base_config) + service.start() + + assert service._tick == 3 + assert service.stream.population_user_ids == source_stream.population_user_ids + assert service.stream.active_visit_count == source_stream.active_visit_count + + def test_save_state_writes_tick_stream_state_v2(self, base_config): + """Сервис сохраняет v2-снимок тикового слоя.""" + service = GeneratorService(base_config) + service.state_manager = MagicMock() + tick_at = datetime.now(timezone.utc).replace(tzinfo=None) + service.stream.generate_tick(event_budget=10, tick_started_at=tick_at) + service._tick = 1 + + service._save_state("batch-1") + + saved_state = service.state_manager.save.call_args.args[0] + assert saved_state.version == "2.0" + assert saved_state.population + assert saved_state.active_visits + service.state_manager.flush.assert_called_once() + + def test_state_reset_skips_loading_saved_state(self, base_config): + """GEN_STATE_RESET=true запускает сервис с чистого состояния.""" + reset_config = replace(base_config, state_reset=True) + state_manager = MagicMock() + + 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.object(GeneratorService, "_main_loop", return_value=None): + + service = GeneratorService(reset_config) + service.start() + + state_manager.load.assert_not_called() + assert service._tick == 0 + + def test_invalid_restored_v2_state_starts_fresh(self, base_config, caplog): + """Сервис не падает, если v2 state ссылается на неизвестный профиль.""" + state_manager = MagicMock() + state_manager.load.return_value = GeneratorState( + tick=9, + rng_state=random.Random(42).getstate(), + last_batch_id="bad-v2", + last_timestamp=datetime.now(timezone.utc), + population=[ + { + "user_domain_id": "user-unknown", + "seed_click_id": "missing-click-id", + } + ], + active_visits=[], + ) + + 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.object(GeneratorService, "_main_loop", return_value=None), \ + caplog.at_level(logging.WARNING, logger="generator"): + + service = GeneratorService(base_config) + service.start() + + assert service._tick == 0 + assert "State data was invalid" in caplog.text diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index da55499..7c52c64 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -2,6 +2,7 @@ Тесты сохранения и восстановления состояния генератора. """ import json +import logging import random from datetime import datetime, timezone from unittest.mock import MagicMock, patch @@ -17,6 +18,47 @@ def _make_valid_rng_state(seed: int = 42): return rng.getstate() +def _make_valid_v2_state_data() -> dict: + """Создаёт минимальный валидный state v2 для тестов загрузки.""" + return { + "tick": 42, + "rng_state": list(_make_valid_rng_state(42)), + "last_batch_id": "v2", + "last_timestamp": "2026-06-11T12:00:00+00:00", + "version": "2.0", + "population": [ + { + "user_domain_id": "user-1", + "seed_click_id": "seed-1", + "active_click_id": "visit-1", + "last_finished_at": None, + } + ], + "active_visits": [ + { + "user_domain_id": "user-1", + "click_id": "visit-1", + "next_index": 1, + "started_at": "2026-06-11T12:00:00", + "offsets_us": [0, 60_000_000], + "page_url_paths": ["/home", "/cart"], + } + ], + "pending_visit_births": 0.5, + } + + +def _minimal_population() -> list[dict]: + return [ + { + "user_domain_id": "user-1", + "seed_click_id": "seed-1", + "active_click_id": None, + "last_finished_at": None, + } + ] + + class TestGeneratorState: """Тесты структуры состояния генератора.""" @@ -30,17 +72,17 @@ class TestGeneratorState: rng_state=rng_state, last_batch_id="abc123", last_timestamp=now, - version="1.0", + version="2.0", ) assert state.tick == 42 assert state.rng_state == rng_state assert state.last_batch_id == "abc123" assert state.last_timestamp == now - assert state.version == "1.0" + assert state.version == "2.0" def test_default_version(self): - """Версия по умолчанию.""" + """Новые состояния по умолчанию пишутся в версии 2.""" now = datetime.now(timezone.utc) rng_state = _make_valid_rng_state(42) @@ -51,7 +93,7 @@ class TestGeneratorState: last_timestamp=now, ) - assert state.version == "1.0" + assert state.version == "2.0" def test_to_dict_serialization(self): """Сериализация в словарь (JSON-safe, без pickle).""" @@ -63,6 +105,7 @@ class TestGeneratorState: rng_state=rng_state, last_batch_id="abc123", last_timestamp=now, + population=_minimal_population(), ) data = state.to_dict() @@ -70,7 +113,7 @@ class TestGeneratorState: assert data["tick"] == 42 assert data["last_batch_id"] == "abc123" assert data["last_timestamp"] == now.isoformat() - assert data["version"] == "1.0" + assert data["version"] == "2.0" # Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle) assert "rng_state" in data @@ -93,6 +136,7 @@ class TestGeneratorState: rng_state=rng_state, last_batch_id="abc123", last_timestamp=now, + population=_minimal_population(), ) # Сериализуем и десериализуем @@ -117,6 +161,7 @@ class TestGeneratorState: rng_state=rng.getstate(), last_batch_id="test", last_timestamp=datetime.now(timezone.utc), + population=_minimal_population(), ) # Десериализуем @@ -137,19 +182,58 @@ class TestGeneratorState: assert next_values == values_after + def test_version_2_roundtrip_keeps_population_and_active_visits(self): + """State v2 хранит популяцию и активные визиты в JSON.""" + rng_state = _make_valid_rng_state(42) + state = GeneratorState( + tick=7, + rng_state=rng_state, + last_batch_id="batch-7", + last_timestamp=datetime(2026, 6, 11, 12, 0, tzinfo=timezone.utc), + version="2.0", + population=[ + { + "user_domain_id": "user-1", + "seed_click_id": "seed-1", + "active_click_id": "visit-1", + "last_finished_at": "2026-06-11T11:30:00", + } + ], + active_visits=[ + { + "user_domain_id": "user-1", + "click_id": "visit-1", + "next_index": 1, + "started_at": "2026-06-11T12:00:00", + "offsets_us": [0, 60_000_000], + "page_url_paths": ["/home", "/cart"], + } + ], + pending_visit_births=0.5, + ) + + restored = GeneratorState.from_dict(json.loads(json.dumps(state.to_dict()))) + + assert restored.version == "2.0" + assert restored.tick == state.tick + assert restored.rng_state == rng_state + assert restored.population == state.population + assert restored.active_visits == state.active_visits + assert restored.pending_visit_births == 0.5 + class TestGeneratorStateValidation: """Тесты валидации состояния и graceful degradation.""" - def test_from_dict_missing_rng_state_raises(self): - """from_dict выбрасывает исключение при отсутствии rng_state.""" + def test_from_dict_missing_version_raises(self): + """from_dict выбрасывает исключение при отсутствии версии v2.""" data = { "tick": 42, "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", } - with pytest.raises(ValueError, match="rng_state"): + with pytest.raises(ValueError, match="version"): GeneratorState.from_dict(data) def test_from_dict_invalid_rng_state_raises(self): @@ -159,6 +243,9 @@ class TestGeneratorStateValidation: "rng_state": "not_a_tuple", "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", + "version": "2.0", + "population": _minimal_population(), + "active_visits": [], } with pytest.raises(ValueError): @@ -171,6 +258,9 @@ class TestGeneratorStateValidation: "rng_state": [1], # Слишком короткий "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", + "version": "2.0", + "population": _minimal_population(), + "active_visits": [], } with pytest.raises(ValueError): @@ -183,6 +273,9 @@ 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", + "population": [], + "active_visits": [], } with pytest.raises(ValueError): @@ -195,6 +288,9 @@ class TestGeneratorStateValidation: "rng_state": "invalid", "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", + "version": "2.0", + "population": [], + "active_visits": [], } result = GeneratorState.from_dict_safe(data) @@ -208,23 +304,24 @@ class TestGeneratorStateValidation: "rng_state": list(rng.getstate()), # JSON сериализует tuple как list "last_batch_id": "test", "last_timestamp": "2024-01-01T00:00:00+00:00", + "version": "2.0", + "population": _minimal_population(), + "active_visits": [], } result = GeneratorState.from_dict_safe(data) assert result is not None assert result.tick == 42 - def test_from_dict_uses_defaults_for_missing_fields(self): - """from_dict использует defaults для отсутствующих полей.""" + def test_from_dict_rejects_old_state_without_version(self): + """from_dict не восстанавливает старое state v1 без версии.""" rng = random.Random(42) data = { "rng_state": list(rng.getstate()), } - result = GeneratorState.from_dict(data) - assert result.tick == 0 - assert result.last_batch_id == "" - assert result.version == "1.0" + with pytest.raises(ValueError, match="version"): + GeneratorState.from_dict(data) class TestKafkaStateManager: @@ -302,6 +399,8 @@ class TestKafkaStateManager: rng_state=_make_valid_rng_state(100), last_batch_id="xyz789", last_timestamp=now, + version="2.0", + population=_minimal_population(), ) # Мокаем consumer с сообщением @@ -343,6 +442,208 @@ class TestKafkaStateManager: # Должно вернуть None из-за невалидного state assert result is None + def test_load_invalid_v2_nested_state_returns_none(self, caplog): + """Битое state v2 с валидным rng_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) + + mock_message = MagicMock() + mock_message.key = b"default" + mock_message.value = { + "tick": 42, + "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", + "population": [{"user_domain_id": "user-1"}], + "active_visits": [ + { + "user_domain_id": "user-1", + "click_id": "visit-1", + "next_index": 1, + "started_at": "2026-06-11T12:00:00", + "offsets_us": [0], + } + ], + } + + 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 "Invalid state" in caplog.text + + def test_load_empty_population_v2_returns_none(self, caplog): + """Пустая популяция в state v2 не восстанавливается.""" + 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["population"] = [] + bad_state["active_visits"] = [] + + 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 "Invalid state" in caplog.text + + def test_load_bad_pending_births_v2_returns_none(self, caplog): + """Нечисловой pending_visit_births в state v2 не восстанавливается.""" + 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["pending_visit_births"] = "bad" + + 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 "Invalid state" in caplog.text + + def test_load_active_visit_with_unknown_user_v2_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_v2_state_data() + bad_state["active_visits"][0]["user_domain_id"] = "missing-user" + + 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 "Invalid state" in caplog.text + + def test_load_active_visit_with_conflicting_click_id_v2_returns_none(self, caplog): + """active_click_id пользователя не должен противоречить визиту.""" + 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["population"][0]["active_click_id"] = "other-visit" + + 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 "Invalid state" in caplog.text + + def test_load_population_ghost_active_click_id_v2_returns_none(self, caplog): + """active_click_id пользователя должен иметь соответствующий активный визит.""" + 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["population"][0]["active_click_id"] = "ghost" + bad_state["active_visits"] = [] + + 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 "Invalid state" in caplog.text + + def test_load_version_1_state_returns_none_with_warning(self, caplog): + """Старое state v1 не восстанавливается и даёт чистый старт.""" + 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) + + old_state = GeneratorState( + tick=100, + rng_state=_make_valid_rng_state(100), + last_batch_id="old", + last_timestamp=datetime.now(timezone.utc), + version="1.0", + ) + + mock_message = MagicMock() + mock_message.key = b"default" + mock_message.value = old_state.to_dict() + + 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 "version" in caplog.text + def test_load_ignores_wrong_key(self): """Загрузка игнорирует сообщения с другим ключом.""" with patch("generator._import_kafka") as mock_import, \ @@ -435,6 +736,7 @@ class TestJsonSafeState: rng_state=rng.getstate(), last_batch_id="test123", last_timestamp=datetime.now(timezone.utc), + population=_minimal_population(), ) # Сериализуем через JSON (как в Kafka)