diff --git a/.scratch/feature-data-generator/issues/04-user-population-and-returns.md b/.scratch/feature-data-generator/issues/04-user-population-and-returns.md index e5fdc0f..8347557 100644 --- a/.scratch/feature-data-generator/issues/04-user-population-and-returns.md +++ b/.scratch/feature-data-generator/issues/04-user-population-and-returns.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Популяция пользователей и возвраты @@ -17,14 +17,14 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] На длинной симуляции получается здоровая пирамида +- [x] На длинной симуляции получается здоровая пирамида `users < sessions < events`. -- [ ] Один пользователь может иметь несколько `click_id` во времени. -- [ ] Один `click_id` принадлежит ровно одному `user_domain_id`. -- [ ] Пользователь не получает новый визит раньше минимального кулдауна возврата. -- [ ] Если доступных возвращающихся пользователей нет, новый визит достаётся +- [x] Один пользователь может иметь несколько `click_id` во времени. +- [x] Один `click_id` принадлежит ровно одному `user_domain_id`. +- [x] Пользователь не получает новый визит раньше минимального кулдауна возврата. +- [x] Если доступных возвращающихся пользователей нет, новый визит достаётся новому пользователю. -- [ ] Размер активной популяции не растёт выше заданного потолка; при +- [x] Размер активной популяции не растёт выше заданного потолка; при переполнении вытесняется давно неактивный пользователь без активного визита. ## Blocked by @@ -33,3 +33,17 @@ Status: ready-for-agent ## Comments +2026-06-11: реализовано через ограниченную популяцию в `TickStreamGenerator`. +Профиль пользователя переиспользуется между визитами, активный пользователь не +получает второй визит, после завершения действует `GEN_MIN_RETURN_MINUTES`. +Если доступных возвратов нет или срабатывает `GEN_P_NEW_USER`, создаётся новый +пользователь с вытеснением давно неактивного профиля без роста выше +`GEN_POPULATION_MAX`. + +Проверка: `uv run --with-requirements generator/requirements.txt pytest +generator/tests -q` — 88 passed. В песочнице `/snap/bin/uv` падает на +`snap-confine`, поэтому полный pytest запускался вне песочницы. + +После свежего ревью исправлены две детали: UUID теперь создаются через единый +ГПСЧ генератора, а кулдаун пользователя начинается от запланированного времени +последнего события визита, а не от времени позднего тика. diff --git a/generator/README.md b/generator/README.md index 7237644..e1facbc 100644 --- a/generator/README.md +++ b/generator/README.md @@ -52,6 +52,9 @@ generator-service -> Kafka topics -> (потребители отдельно) строго растущие запланированные `event_timestamp`. - Тиковый слой хранит активные визиты между вызовами и выпускает только события, у которых наступил `event_timestamp`; завершённые визиты удаляются из памяти. +- Тиковый слой ведёт ограниченную популяцию пользователей: один + `user_domain_id` может вернуться в новом `click_id` после кулдауна, а при + переполнении вытесняется давно неактивный пользователь. - Сохраняем связи `event_id <-> location`, `click_id <-> device/geo`. ## Конфигурация (env) @@ -66,7 +69,9 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_MAX_EVENTS_PER_TICK` | Максимум событий за тик | `50` | | `GEN_MAX_SESSION_EVENTS` | Потолок длины одного визита, защита от петель | `30` | | `GEN_MAX_ACTIVE_SESSIONS` | Потолок одновременных активных визитов | `200` | -| `GEN_POPULATION_MAX` | Будущий потолок популяции пользователей; сейчас нужен для валидации активных визитов | `300` | +| `GEN_POPULATION_MAX` | Потолок активной популяции пользователей | `300` | +| `GEN_P_NEW_USER` | Вероятность отдать новый визит новому пользователю | `0.15` | +| `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | | `GEN_DATA_DIR` | Путь к JSONL файлам | `/data` | | `GEN_SEED` | Сид для воспроизводимости | — | | `GEN_ENABLED` | Включить генерацию | `true` | diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index 675e76c..3353662 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -36,6 +36,12 @@ class Config: population_max: int = field( default_factory=lambda: int(os.getenv("GEN_POPULATION_MAX", "300")) ) + p_new_user: float = field( + default_factory=lambda: float(os.getenv("GEN_P_NEW_USER", "0.15")) + ) + min_return_minutes: int = field( + default_factory=lambda: int(os.getenv("GEN_MIN_RETURN_MINUTES", "30")) + ) data_dir: Path = field( default_factory=lambda: Path(os.getenv("GEN_DATA_DIR", "/data")) ) @@ -68,6 +74,10 @@ class Config: raise ValueError("GEN_MAX_ACTIVE_SESSIONS must be >= 1") if self.population_max < 1: raise ValueError("GEN_POPULATION_MAX must be >= 1") + if not 0 <= self.p_new_user <= 1: + raise ValueError("GEN_P_NEW_USER must be between 0 and 1") + if self.min_return_minutes < 0: + raise ValueError("GEN_MIN_RETURN_MINUTES must be >= 0") if self.max_active_sessions >= self.population_max: raise ValueError("GEN_MAX_ACTIVE_SESSIONS must be < GEN_POPULATION_MAX") if not self.data_dir.exists(): diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py index 50bc276..a69383f 100644 --- a/generator/src/clickstream_generator/generation.py +++ b/generator/src/clickstream_generator/generation.py @@ -77,7 +77,7 @@ class EventGenerator: def _new_uuid(self) -> str: """Генерирует новый UUID.""" - return str(uuid.uuid4()) + return str(uuid.UUID(int=self.rng.getrandbits(128), version=4)) def _current_timestamp(self) -> str: """Возвращает текущую метку времени в формате JSONL.""" @@ -133,6 +133,7 @@ class EventGenerator: self, batch_size: int, planned_start_at: datetime | None = None, + user_profile: dict[str, dict] | None = None, ) -> dict[str, list[dict]]: """Генерирует один визит с сохранением связей.""" if not self.dictionary.browser_events: @@ -177,8 +178,16 @@ class EventGenerator: base_click_id = base_browser["click_id"] base_browser_events = [base_browser for _ in range(len(visit_path))] - base_device = self.dictionary.device_by_click_id.get(base_click_id) - base_geo = self.dictionary.geo_by_click_id.get(base_click_id) + base_device = ( + user_profile["device"] + if user_profile is not None + else self.dictionary.device_by_click_id.get(base_click_id) + ) + base_geo = ( + user_profile["geo"] + if user_profile is not None + else self.dictionary.geo_by_click_id.get(base_click_id) + ) new_click_id = self._new_uuid() planned_timestamp = planned_start_at or datetime.now(timezone.utc) if planned_timestamp.tzinfo is not None: diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py index 614138f..16a66ef 100644 --- a/generator/src/clickstream_generator/runtime.py +++ b/generator/src/clickstream_generator/runtime.py @@ -1,7 +1,7 @@ """Тиковый слой генератора с активными визитами между вызовами.""" from dataclasses import dataclass -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from weakref import WeakKeyDictionary from clickstream_generator.generation import EventGenerator @@ -10,12 +10,29 @@ from clickstream_generator.generation import EventGenerator TOPICS = ("browser_events", "location_events", "device_events", "geo_events") +@dataclass +class UserProfile: + """Постоянный профиль пользователя между визитами.""" + + user_domain_id: str + seed_click_id: str + device: dict + geo: dict + active_click_id: str | None = None + last_finished_at: datetime | None = None + + @property + def is_active(self) -> bool: + return self.active_click_id is not None + + @dataclass class ActiveVisit: """Запланированный визит, который выпускается по тикам.""" batch: dict[str, list[dict]] timestamps: list[datetime] + user: UserProfile | None = None next_index: int = 0 @property @@ -23,6 +40,80 @@ class ActiveVisit: return self.next_index >= len(self.timestamps) +class UserPopulation: + """Ограниченная популяция пользователей для тикового потока.""" + + def __init__(self, generator: EventGenerator): + self.generator = generator + self.users: list[UserProfile] = [ + self._create_user() + for _ in range(generator.config.population_max) + ] + + def choose_for_visit(self, tick_time: datetime) -> UserProfile | None: + """Выбирает пользователя без активного визита и кулдауна.""" + available = [ + user for user in self.users + if self._is_available(user, tick_time) + ] + if not available: + return self._rotate_new_user() + if self.generator.rng.random() < self.generator.config.p_new_user: + return self._rotate_new_user() or self.generator.rng.choice(available) + return self.generator.rng.choice(available) + + def start_visit(self, user: UserProfile, click_id: str) -> None: + user.active_click_id = click_id + + def finish_visit(self, user: UserProfile | None, finished_at: datetime) -> None: + if user is None: + return + user.active_click_id = None + user.last_finished_at = finished_at + + def _create_user(self) -> UserProfile: + seed_click_id = self.generator.rng.choice(self._profile_seed_click_ids()) + device = { + **self.generator.dictionary.device_by_click_id[seed_click_id], + "user_domain_id": self.generator._new_uuid(), + } + geo = self.generator.dictionary.geo_by_click_id[seed_click_id] + return UserProfile( + user_domain_id=device["user_domain_id"], + seed_click_id=seed_click_id, + device=device, + geo=geo, + ) + + def _rotate_new_user(self) -> UserProfile | None: + inactive_users = [user for user in self.users if not user.is_active] + if not inactive_users: + return None + + new_user = self._create_user() + victim = min( + inactive_users, + key=lambda user: user.last_finished_at or datetime.min, + ) + self.users[self.users.index(victim)] = new_user + return new_user + + def _is_available(self, user: UserProfile, tick_time: datetime) -> bool: + if user.is_active: + return False + if user.last_finished_at is None: + return True + cooldown = timedelta(minutes=self.generator.config.min_return_minutes) + return tick_time - user.last_finished_at >= cooldown + + def _profile_seed_click_ids(self) -> list[str]: + return [ + click_id + for click_id in self.generator.dictionary.device_by_click_id + if click_id in self.generator.dictionary.geo_by_click_id + ] + + def _empty_batch() -> dict[str, list[dict]]: return {topic: [] for topic in TOPICS} @@ -44,8 +135,17 @@ class TickStreamGenerator: def __init__(self, generator: EventGenerator): self.generator = generator self.active_visits: list[ActiveVisit] = [] + self.population = UserPopulation(generator) self._pending_event_budget = 0 + @property + def population_size(self) -> int: + return len(self.population.users) + + @property + def population_user_ids(self) -> set[str]: + return {user.user_domain_id for user in self.population.users} + def generate_tick( self, event_budget: int, @@ -74,9 +174,15 @@ class TickStreamGenerator: self._pending_event_budget > 0 and len(self.active_visits) < self.generator.config.max_active_sessions ): + user = self.population.choose_for_visit(tick_time) + if user is None: + self._pending_event_budget = 0 + return + visit_batch = self.generator.generate_batch( self.generator.config.max_session_events, planned_start_at=tick_time, + user_profile={"device": user.device, "geo": user.geo}, ) timestamps = [ _parse_timestamp(event["event_timestamp"]) @@ -85,7 +191,11 @@ class TickStreamGenerator: if not timestamps: break - self.active_visits.append(ActiveVisit(batch=visit_batch, timestamps=timestamps)) + 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) + ) self._pending_event_budget -= len(timestamps) if len(self.active_visits) >= self.generator.config.max_active_sessions: @@ -104,6 +214,10 @@ class TickStreamGenerator: visit.next_index += 1 def _drop_finished_visits(self) -> None: + for visit in self.active_visits: + if visit.is_finished: + self.population.finish_visit(visit.user, visit.timestamps[-1]) + self.active_visits = [ visit for visit in self.active_visits diff --git a/generator/tests/test_config.py b/generator/tests/test_config.py index 2201bd5..4746668 100644 --- a/generator/tests/test_config.py +++ b/generator/tests/test_config.py @@ -35,6 +35,18 @@ class TestConfigValidation: with pytest.raises(ValueError, match="GEN_MAX_ACTIVE_SESSIONS"): replace(base_config, max_active_sessions=10, population_max=10) + def test_new_user_probability_must_be_valid_share(self, base_config): + """Вероятность нового пользователя должна быть долей от 0 до 1.""" + from dataclasses import replace + with pytest.raises(ValueError, match="GEN_P_NEW_USER"): + replace(base_config, p_new_user=1.5) + + def test_min_return_minutes_must_not_be_negative(self, base_config): + """Кулдаун возврата не может быть отрицательным.""" + from dataclasses import replace + with pytest.raises(ValueError, match="GEN_MIN_RETURN_MINUTES"): + replace(base_config, min_return_minutes=-1) + def test_data_dir_must_exist(self, base_config): """data_dir должен существовать.""" from dataclasses import replace diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index ea07356..d59bf62 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -37,6 +37,41 @@ def _page_path_visits(generator, visits_count: int): ] +def _run_stream( + stream, + tick_at: datetime, + ticks_count: int, + event_budget: int, + tick_seconds: int, +): + events = [] + for tick_index in range(ticks_count): + batch = stream.generate_tick( + event_budget=event_budget, + tick_started_at=tick_at + timedelta(seconds=tick_seconds * tick_index), + ) + events.extend(batch["device_events"]) + return events + + +def _run_stream_batches( + stream, + tick_at: datetime, + ticks_count: int, + event_budget: int, + tick_seconds: int, +): + batches = [] + for tick_index in range(ticks_count): + batches.append( + stream.generate_tick( + event_budget=event_budget, + tick_started_at=tick_at + timedelta(seconds=tick_seconds * tick_index), + ) + ) + return batches + + class TestEventGeneration: """Тесты генерации событий.""" @@ -137,7 +172,7 @@ class TestEventGeneration: def test_tick_stream_keeps_active_visit_between_ticks(self, event_dictionary, base_config): """Один визит может выпускать события в нескольких последовательных тиках.""" - config = replace(base_config, tick_seconds=60, max_session_events=5) + config = replace(base_config, tick_seconds=30 * 60, max_session_events=5) generator = EventGenerator(event_dictionary, config) stream = TickStreamGenerator(generator) first_tick_at = datetime(2026, 6, 11, 12, 0) @@ -230,6 +265,228 @@ class TestEventGeneration: assert {event["click_id"] for event in final_tick["browser_events"]} == {click_id} assert later_tick["browser_events"] == [] + def test_tick_stream_reuses_users_across_visits(self, event_dictionary, base_config): + """Длинная симуляция даёт пирамиду users < sessions < events.""" + config = replace( + base_config, + tick_seconds=60, + max_session_events=3, + max_active_sessions=50, + population_max=60, + p_new_user=0, + min_return_minutes=0, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + + events = _run_stream( + stream=stream, + tick_at=datetime(2026, 6, 11, 12, 0), + ticks_count=240, + event_budget=6, + tick_seconds=config.tick_seconds, + ) + users_by_session = {} + for event in events: + users_by_session.setdefault(event["click_id"], set()).add(event["user_domain_id"]) + + unique_users = {event["user_domain_id"] for event in events} + repeated_users = [ + user_id + for user_id in unique_users + if len({ + event["click_id"] + for event in events + if event["user_domain_id"] == user_id + }) > 1 + ] + + assert len(unique_users) < len(users_by_session) < len(events) + assert repeated_users + assert all(len(user_ids) == 1 for user_ids in users_by_session.values()) + + def test_tick_stream_replays_same_flow_with_same_seed( + self, event_dictionary, base_config + ): + """Одинаковый seed даёт одинаковые решения популяции и визитов.""" + config = replace( + base_config, + tick_seconds=60, + max_session_events=3, + max_active_sessions=5, + population_max=6, + p_new_user=0, + min_return_minutes=0, + ) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_stream = TickStreamGenerator(EventGenerator(event_dictionary, config)) + second_stream = TickStreamGenerator(EventGenerator(event_dictionary, config)) + + first_events = [] + second_events = [] + for tick_index in range(10): + tick_time = tick_at + timedelta(seconds=config.tick_seconds * tick_index) + first_batch = first_stream.generate_tick( + event_budget=3, + tick_started_at=tick_time, + ) + second_batch = second_stream.generate_tick( + event_budget=3, + tick_started_at=tick_time, + ) + first_events.extend(first_batch["device_events"]) + second_events.extend(second_batch["device_events"]) + + first_decisions = [ + (event["click_id"], event["user_domain_id"]) + for event in first_events + ] + second_decisions = [ + (event["click_id"], event["user_domain_id"]) + for event in second_events + ] + + assert first_decisions == second_decisions + + def test_returning_user_waits_for_cooldown_after_visit_end( + self, event_dictionary, base_config + ): + """Пользователь не получает новый визит раньше кулдауна возврата.""" + config = replace( + base_config, + tick_seconds=60, + max_session_events=3, + max_active_sessions=5, + population_max=20, + p_new_user=0, + min_return_minutes=30, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + + batches = _run_stream_batches( + stream=stream, + tick_at=datetime(2026, 6, 11, 12, 0), + ticks_count=240, + event_budget=1, + tick_seconds=config.tick_seconds, + ) + users_by_session = {} + times_by_session = {} + for batch in batches: + for device_event in batch["device_events"]: + users_by_session.setdefault( + device_event["click_id"], + device_event["user_domain_id"], + ) + for browser_event in batch["browser_events"]: + times_by_session.setdefault(browser_event["click_id"], []).append( + datetime.fromisoformat( + browser_event["event_timestamp"].replace(" ", "T") + ) + ) + + sessions_by_user = {} + for click_id, timestamps in times_by_session.items(): + sessions_by_user.setdefault(users_by_session[click_id], []).append( + (min(timestamps), max(timestamps)) + ) + + repeated_users = [ + sessions + for sessions in sessions_by_user.values() + if len(sessions) > 1 + ] + + assert repeated_users + for sessions in repeated_users: + ordered_sessions = sorted(sessions) + for previous, current in zip(ordered_sessions, ordered_sessions[1:]): + previous_end = previous[1] + current_start = current[0] + assert current_start - previous_end >= timedelta( + minutes=config.min_return_minutes + ) + + def test_visit_cooldown_starts_from_planned_last_event_time( + self, event_dictionary, base_config + ): + """Кулдаун считается от запланированного конца визита, а не от позднего тика.""" + config = replace( + base_config, + tick_seconds=60, + max_session_events=3, + max_active_sessions=1, + population_max=2, + p_new_user=0, + min_return_minutes=30, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=1, tick_started_at=tick_at) + user_id = first_tick["device_events"][0]["user_domain_id"] + planned_visit_end = stream.active_visits[0].timestamps[-1] + + stream.generate_tick( + event_budget=0, + tick_started_at=planned_visit_end + timedelta(hours=1), + ) + user = next( + user + for user in stream.population.users + if user.user_domain_id == user_id + ) + + assert user.last_finished_at == planned_visit_end + + def test_tick_stream_creates_new_user_when_no_returning_user_is_available( + self, event_dictionary, base_config + ): + """Если все прежние пользователи в кулдауне, новый визит получает нового пользователя.""" + config = replace( + base_config, + tick_seconds=60, + max_session_events=2, + max_active_sessions=1, + population_max=2, + p_new_user=0, + min_return_minutes=240, + ) + generator = EventGenerator(event_dictionary, config) + stream = TickStreamGenerator(generator) + tick_at = datetime(2026, 6, 11, 12, 0) + + first_tick = stream.generate_tick(event_budget=5, tick_started_at=tick_at) + second_tick = stream.generate_tick( + event_budget=5, + tick_started_at=tick_at + timedelta(hours=1), + ) + third_tick = stream.generate_tick( + event_budget=5, + tick_started_at=tick_at + timedelta(hours=2), + ) + + first_two_users = { + event["user_domain_id"] + for event in first_tick["device_events"] + second_tick["device_events"] + } + later_users = { + event["user_domain_id"] + for event in third_tick["device_events"] + } + first_user = first_tick["device_events"][0]["user_domain_id"] + second_new_users = first_two_users - {first_user} + new_users = later_users - first_two_users + + assert new_users + assert stream.population_size <= config.population_max + assert first_user not in stream.population_user_ids + assert second_new_users <= stream.population_user_ids + assert new_users <= stream.population_user_ids + def test_generate_batch_creates_one_connected_visit(self, event_dictionary, base_config): """Публичный вызов генератора создаёт один связанный визит.""" generator = EventGenerator(event_dictionary, base_config)