Files
clickstream-ch-kafka-supers…/generator/tests/test_generation.py
T
ddadmin fcb06da16e feat(generator): добавлено восстановление state по модельному времени
- Зачем:
  - рестарт генератора должен продолжать поток от модельной точки без дублей и смешивания state разных настроек.
- Что:
  - state v2 хранит модельную и настенную метки, скорость, timezone, T0 и seed.
  - live-восстановление считает модельную точку по настенной дельте и проверяет совместимость config.
  - добавлен путь восстановления от T_end и тесты короткого и долгого простоя.
- Проверка:
  - make generator-test.
  - ClickHouse-сценарии короткого и долгого восстановления state.
  - reviewer gate issue 04 пройден после исправления совместимости state.
2026-06-14 18:47:34 +03:00

1225 lines
49 KiB
Python

"""
Тесты генерации событий.
"""
import json
import random
import uuid
from dataclasses import replace
from datetime import datetime, timedelta, timezone
import pytest
from generator import (
EventGenerator,
EventDictionary,
EXPECTED_VISIT_EVENTS,
TickStreamGenerator,
calculate_events_count,
generate_tick_batch,
hour_factor,
)
ALLOWED_PAGE_PATHS = {
"/home",
"/product_a",
"/product_b",
"/cart",
"/payment",
"/confirmation",
}
def _parse_event_timestamps(batch):
return [
datetime.fromisoformat(event["event_timestamp"].replace(" ", "T"))
for event in batch["browser_events"]
]
def _page_path_visits(generator, visits_count: int):
return [
[
event["page_url_path"]
for event in generator.generate_batch(30)["location_events"]
]
for _ in range(visits_count)
]
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:
"""Тесты генерации событий."""
def test_event_dictionary_loads(self, event_dictionary):
"""Словарь событий загружается корректно."""
assert len(event_dictionary.browser_events) == 1000
assert len(event_dictionary.location_events) == 1000
assert len(event_dictionary.device_events) == 1000
assert len(event_dictionary.geo_events) == 1000
def test_dictionary_consistency(self, event_dictionary):
"""Все записи в словаре имеют корректные связи."""
browser_event_ids = {e["event_id"] for e in event_dictionary.browser_events}
browser_click_ids = {e["click_id"] for e in event_dictionary.browser_events}
location_orphaned = sum(
1 for loc in event_dictionary.location_events
if loc["event_id"] not in browser_event_ids
)
device_orphaned = sum(
1 for dev in event_dictionary.device_events
if dev["click_id"] not in browser_click_ids
)
geo_orphaned = sum(
1 for geo in event_dictionary.geo_events
if geo["click_id"] not in browser_click_ids
)
assert location_orphaned == 0, f"Found {location_orphaned} orphaned location records"
assert device_orphaned == 0, f"Found {device_orphaned} orphaned device records"
assert geo_orphaned == 0, f"Found {geo_orphaned} orphaned geo records"
def test_generate_batch_structure(self, event_dictionary, base_config):
"""Батч имеет правильную структуру."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(10)
assert "browser_events" in batch
assert "location_events" in batch
assert "device_events" in batch
assert "geo_events" in batch
def test_generate_batch_respects_requested_visit_budget(self, event_dictionary, base_config):
"""Визит не превышает запрошенный бюджет событий."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(10)
browser_count = len(batch["browser_events"])
assert 1 <= browser_count <= 10
assert len(batch["location_events"]) == browser_count
assert len(batch["device_events"]) == browser_count
assert len(batch["geo_events"]) == browser_count
def test_generate_tick_batch_uses_requested_event_budget(self, event_dictionary, base_config):
"""Тиковый батч трактует входное число как событийный бюджет."""
generator = EventGenerator(event_dictionary, base_config)
batch = generate_tick_batch(generator, 20)
browser_count = len(batch["browser_events"])
assert 1 <= browser_count <= 20
assert len(batch["location_events"]) == browser_count
assert len(batch["device_events"]) == browser_count
assert len(batch["geo_events"]) == browser_count
assert len({event["click_id"] for event in batch["browser_events"]}) == browser_count
def test_tick_stream_keeps_long_window_intensity_near_event_budget(
self, event_dictionary, base_config
):
"""Длинное окно держит среднюю интенсивность около событийного бюджета."""
config = replace(
base_config,
tick_seconds=5,
jitter_pct=0,
min_events_per_tick=17,
max_events_per_tick=17,
max_active_sessions=10_000,
population_max=10_001,
)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
ticks_count = 12 * 60
total_events = 0
for tick_index in range(ticks_count):
batch = stream.generate_tick(
event_budget=17,
tick_started_at=tick_at + timedelta(seconds=5 * tick_index),
)
total_events += len(batch["browser_events"])
events_per_minute = total_events / (ticks_count * config.tick_seconds / 60)
expected_events_per_minute = 17 * 60 / config.tick_seconds
assert events_per_minute == pytest.approx(expected_events_per_minute, rel=0.10)
def test_tick_stream_births_visits_from_average_visit_length_budget(
self, event_dictionary, base_config
):
"""Новый визит рождается после накопления бюджета средней длины визита."""
config = replace(
base_config,
tick_seconds=60,
max_active_sessions=100,
population_max=101,
)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
partial_budget_tick = stream.generate_tick(
event_budget=9,
tick_started_at=tick_at,
)
completed_budget_tick = stream.generate_tick(
event_budget=1,
tick_started_at=tick_at + timedelta(seconds=config.tick_seconds),
)
assert partial_budget_tick["browser_events"] == []
assert len(completed_budget_tick["browser_events"]) == 1
def test_tick_stream_keeps_active_visit_between_ticks(self, event_dictionary, base_config):
"""Один визит может выпускать события в нескольких последовательных тиках."""
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)
first_tick = stream.generate_tick(event_budget=10, tick_started_at=first_tick_at)
second_tick = stream.generate_tick(
event_budget=0,
tick_started_at=first_tick_at + timedelta(seconds=config.tick_seconds),
)
first_click_id = first_tick["browser_events"][0]["click_id"]
second_click_ids = {
event["click_id"]
for event in second_tick["browser_events"]
}
second_timestamps = _parse_event_timestamps(second_tick)
assert first_click_id in second_click_ids
assert all(
timestamp < first_tick_at + timedelta(seconds=config.tick_seconds)
for timestamp in second_timestamps
)
def test_tick_stream_releases_only_matured_events(self, event_dictionary, base_config):
"""Тик не выпускает будущие события активного визита."""
generator = EventGenerator(event_dictionary, base_config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
first_tick = stream.generate_tick(event_budget=10, tick_started_at=tick_at)
same_time_tick = stream.generate_tick(event_budget=0, tick_started_at=tick_at)
assert len(first_tick["browser_events"]) == 1
assert same_time_tick["browser_events"] == []
assert same_time_tick["location_events"] == []
assert same_time_tick["device_events"] == []
assert same_time_tick["geo_events"] == []
def test_tick_stream_does_not_reemit_finished_visit(self, event_dictionary, base_config):
"""Завершённый визит больше не выпускает события в следующих тиках."""
config = replace(base_config, max_session_events=2)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
first_tick = stream.generate_tick(event_budget=10, tick_started_at=tick_at)
final_tick = stream.generate_tick(
event_budget=0,
tick_started_at=tick_at + timedelta(hours=1),
)
later_tick = stream.generate_tick(
event_budget=0,
tick_started_at=tick_at + timedelta(hours=2),
)
click_id = first_tick["browser_events"][0]["click_id"]
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]
original_visit = stream.active_visits[0]
started_at = datetime.fromisoformat(visit_state["started_at"])
next_index = visit_state["next_index"]
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 {event["click_id"] for event in resumed["browser_events"]} == {
visit_state["click_id"]
}
assert [
event["page_url_path"]
for event in resumed["location_events"]
] == visit_state["page_url_paths"][next_index:next_index + len(resumed_timestamps)]
assert resumed["device_events"] == original_visit.batch["device_events"][
next_index:next_index + len(resumed_timestamps)
]
assert resumed["geo_events"] == original_visit.batch["geo_events"][
next_index:next_index + len(resumed_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
):
"""При заполненном потолке новые рождения пропускаются без накопления бюджета."""
config = replace(
base_config,
max_session_events=2,
max_active_sessions=1,
population_max=2,
)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
first_tick = stream.generate_tick(event_budget=10, tick_started_at=tick_at)
blocked_tick = stream.generate_tick(event_budget=10, tick_started_at=tick_at)
final_tick = stream.generate_tick(
event_budget=0,
tick_started_at=tick_at + timedelta(hours=1),
)
later_tick = stream.generate_tick(
event_budget=0,
tick_started_at=tick_at + timedelta(hours=2),
)
click_id = first_tick["browser_events"][0]["click_id"]
assert len(first_tick["browser_events"]) == 1
assert blocked_tick["browser_events"] == []
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_new_user_share_matches_config_after_warmup(
self, event_dictionary, base_config
):
"""Доля новых пользователей в потоке близка к GEN_P_NEW_USER."""
config = replace(
base_config,
tick_seconds=60,
jitter_pct=0,
min_events_per_tick=30,
max_events_per_tick=30,
lambda_base_per_min=30,
max_active_sessions=200,
population_max=300,
p_new_user=0.15,
min_return_minutes=30,
)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
tick_at = datetime(2026, 6, 11, 12, 0)
warmup_ticks = 240
measure_ticks = 480
seen_users = set()
seen_sessions = set()
total_sessions = 0
new_user_sessions = 0
for tick_index in range(warmup_ticks + measure_ticks):
batch = stream.generate_tick(
event_budget=30,
tick_started_at=tick_at + timedelta(seconds=60 * tick_index),
)
users_by_session = {
event["click_id"]: event["user_domain_id"]
for event in batch["device_events"]
}
for click_id, user_id in users_by_session.items():
if click_id in seen_sessions:
continue
seen_sessions.add(click_id)
if tick_index >= warmup_ticks:
total_sessions += 1
if user_id not in seen_users:
new_user_sessions += 1
seen_users.add(user_id)
new_user_share = new_user_sessions / total_sessions
assert new_user_share == pytest.approx(config.p_new_user, rel=0.35)
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_returning_user_pause_average_matches_model_scale(
self, event_dictionary, base_config
):
"""Средняя пауза между визитами одного пользователя имеет часовой масштаб."""
config = replace(
base_config,
tick_seconds=60,
jitter_pct=0,
min_events_per_tick=30,
max_events_per_tick=30,
lambda_base_per_min=30,
max_active_sessions=200,
population_max=300,
p_new_user=0.15,
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=16 * 60,
event_budget=30,
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))
)
pauses_minutes = []
visit_durations_minutes = []
for sessions in sessions_by_user.values():
ordered_sessions = sorted(sessions)
visit_durations_minutes.extend(
(finished_at - started_at).total_seconds() / 60
for started_at, finished_at in ordered_sessions
)
for previous, current in zip(ordered_sessions, ordered_sessions[1:]):
pauses_minutes.append((current[0] - previous[1]).total_seconds() / 60)
mean_pause_minutes = sum(pauses_minutes) / len(pauses_minutes)
mean_visit_duration_minutes = (
sum(visit_durations_minutes) / len(visit_durations_minutes)
)
expected_pause_minutes = (
config.population_max
/ (
config.lambda_base_per_min
/ EXPECTED_VISIT_EVENTS
* (1 - config.p_new_user)
)
- mean_visit_duration_minutes
)
assert min(pauses_minutes) >= config.min_return_minutes
assert mean_pause_minutes == pytest.approx(expected_pause_minutes, rel=0.30)
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=10, 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=10, tick_started_at=tick_at)
second_tick = stream.generate_tick(
event_budget=10,
tick_started_at=tick_at + timedelta(hours=1),
)
third_tick = stream.generate_tick(
event_budget=10,
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)
batch = generator.generate_batch(3)
original_event_ids = {e["event_id"] for e in event_dictionary.browser_events}
original_click_ids = {e["click_id"] for e in event_dictionary.browser_events}
browser_events = batch["browser_events"]
location_events = batch["location_events"]
device_events = batch["device_events"]
geo_events = batch["geo_events"]
click_ids = {event["click_id"] for event in browser_events}
event_ids = [event["event_id"] for event in browser_events]
assert len(browser_events) > 1
assert len(click_ids) == 1
click_id = next(iter(click_ids))
assert click_id not in original_click_ids
uuid.UUID(click_id)
assert len(set(event_ids)) == len(event_ids)
assert all(event_id not in original_event_ids for event_id in event_ids)
for event_id in event_ids:
uuid.UUID(event_id)
assert {event["event_id"] for event in location_events} == set(event_ids)
assert {event["click_id"] for event in device_events} == {click_id}
assert {event["click_id"] for event in geo_events} == {click_id}
device_context = [{k: v for k, v in event.items() if k != "click_id"} for event in device_events]
geo_context = [{k: v for k, v in event.items() if k != "click_id"} for event in geo_events]
assert len({event["user_domain_id"] for event in device_events}) == 1
assert all(context == device_context[0] for context in device_context)
assert all(context == geo_context[0] for context in geo_context)
def test_generate_batch_creates_multi_event_visit_when_budget_allows(
self, event_dictionary, base_config
):
"""Минимальный связанный визит не схлопывается в одно событие."""
config = replace(base_config, seed=2)
generator = EventGenerator(event_dictionary, config)
batch = generator.generate_batch(3)
assert len(batch["browser_events"]) >= 2
def test_event_ids_are_new_uuids(self, event_dictionary, base_config):
"""event_id и click_id — новые UUID, не из оригинальных данных."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(1)
original_event_ids = {e["event_id"] for e in event_dictionary.browser_events}
original_click_ids = {e["click_id"] for e in event_dictionary.browser_events}
browser_event = batch["browser_events"][0]
assert browser_event["event_id"] not in original_event_ids
assert browser_event["click_id"] not in original_click_ids
# Проверяем, что это валидный UUID
uuid.UUID(browser_event["event_id"])
uuid.UUID(browser_event["click_id"])
def test_event_timestamp_format(self, event_dictionary, base_config):
"""event_timestamp имеет правильный формат."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(1)
browser_event = batch["browser_events"][0]
timestamp = browser_event["event_timestamp"]
# Должен парситься как datetime
dt = datetime.fromisoformat(timestamp.replace(" ", "T"))
assert dt.year >= 2024
def test_visit_event_timestamps_strictly_increase(self, event_dictionary, base_config):
"""Время событий внутри одного визита строго возрастает."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(8)
timestamps = _parse_event_timestamps(batch)
assert all(
previous < current
for previous, current in zip(timestamps, timestamps[1:])
)
def test_visit_event_timestamps_use_planned_user_pauses(self, event_dictionary, base_config):
"""Метки времени визита разделены пользовательскими паузами, а не временем цикла."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(8)
timestamps = _parse_event_timestamps(batch)
pauses_seconds = [
(current - previous).total_seconds()
for previous, current in zip(timestamps, timestamps[1:])
]
assert min(pauses_seconds) >= 1.0
def test_visit_path_uses_known_funnel_pages(self, event_dictionary, base_config):
"""Путь визита состоит из страниц воронки."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(30)
page_paths = [event["page_url_path"] for event in batch["location_events"]]
assert page_paths
assert set(page_paths) <= ALLOWED_PAGE_PATHS
def test_visit_starts_from_calibrated_start_distribution(self, event_dictionary, base_config):
"""Визиты стартуют не только с /home, а по стартовому распределению."""
generator = EventGenerator(event_dictionary, base_config)
visits = _page_path_visits(generator, 1000)
first_pages = [visit[0] for visit in visits]
home_share = first_pages.count("/home") / len(first_pages)
product_entry_share = (
first_pages.count("/product_a") + first_pages.count("/product_b")
) / len(first_pages)
assert 0.54 <= home_share <= 0.64
assert product_entry_share >= 0.25
def test_visit_length_is_capped_by_session_limit(self, event_dictionary, base_config):
"""Длина визита ограничена потолком, который защищает от петель."""
config = replace(base_config, max_session_events=4)
generator = EventGenerator(event_dictionary, config)
visits = _page_path_visits(generator, 200)
assert max(len(visit) for visit in visits) <= 4
assert any(len(visit) == 4 for visit in visits)
def test_visit_length_distribution_matches_seed_scale(
self, event_dictionary, base_config
):
"""Длина визита сопоставима с сидом: медиана и среднее около 10."""
generator = EventGenerator(event_dictionary, base_config)
visits = _page_path_visits(generator, 2000)
visit_lengths = sorted(len(visit) for visit in visits)
median_length = visit_lengths[len(visit_lengths) // 2]
mean_length = sum(visit_lengths) / len(visit_lengths)
assert 8 <= median_length <= 12
assert mean_length == pytest.approx(EXPECTED_VISIT_EVENTS, rel=0.10)
assert max(visit_lengths) <= base_config.max_session_events
def test_visit_pauses_stay_below_session_timeout_scale(self, event_dictionary, base_config):
"""Паузы внутри визита остаются меньше 30 минут, p95 — единицы минут."""
generator = EventGenerator(event_dictionary, base_config)
pauses_seconds = []
for _ in range(500):
timestamps = _parse_event_timestamps(generator.generate_batch(30))
pauses_seconds.extend(
(current - previous).total_seconds()
for previous, current in zip(timestamps, timestamps[1:])
)
pauses_seconds.sort()
p95 = pauses_seconds[int(len(pauses_seconds) * 0.95)]
assert max(pauses_seconds) < 30 * 60
assert p95 < 5 * 60
def test_confirmation_share_matches_seed_scale(self, event_dictionary, base_config):
"""Около четверти визитов доходят до /confirmation."""
generator = EventGenerator(event_dictionary, base_config)
visits = _page_path_visits(generator, 1000)
confirmation_share = (
sum("/confirmation" in visit for visit in visits) / len(visits)
)
assert 0.20 <= confirmation_share <= 0.30
def test_visit_funnel_monotonically_fades(self, event_dictionary, base_config):
"""Воронка по визитам монотонно затухает от главной до подтверждения."""
generator = EventGenerator(event_dictionary, base_config)
visits = _page_path_visits(generator, 2000)
home_visits = sum("/home" in visit for visit in visits)
product_visits = sum(
"/product_a" in visit or "/product_b" in visit
for visit in visits
)
cart_visits = sum("/cart" in visit for visit in visits)
payment_visits = sum("/payment" in visit for visit in visits)
confirmation_visits = sum("/confirmation" in visit for visit in visits)
assert home_visits >= product_visits >= cart_visits
assert cart_visits >= payment_visits >= confirmation_visits
def test_visit_can_continue_after_confirmation(self, event_dictionary, base_config):
"""/confirmation не обязан быть последним событием визита."""
generator = EventGenerator(event_dictionary, base_config)
visits = _page_path_visits(generator, 1000)
assert any(
page_path == "/confirmation" and index < len(visit) - 1
for visit in visits
for index, page_path in enumerate(visit)
)
def test_links_consistency(self, event_dictionary, base_config):
"""Связи между событиями сохраняются."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(5)
for i, browser in enumerate(batch["browser_events"]):
event_id = browser["event_id"]
click_id = browser["click_id"]
# Location должен иметь тот же event_id
assert batch["location_events"][i]["event_id"] == event_id
# Device и Geo должны иметь тот же click_id
assert batch["device_events"][i]["click_id"] == click_id
assert batch["geo_events"][i]["click_id"] == click_id
def test_required_fields_present(self, event_dictionary, base_config):
"""Все обязательные поля присутствуют в событиях."""
generator = EventGenerator(event_dictionary, base_config)
batch = generator.generate_batch(1)
browser = batch["browser_events"][0]
required_fields = [
"event_id", "event_timestamp", "event_type", "click_id",
"browser_name", "browser_user_agent", "browser_language"
]
for field in required_fields:
assert field in browser, f"Missing field: {field}"
class TestPoissonDistribution:
"""Тесты статистической модели."""
def test_hour_factor_uses_model_timezone(self):
"""Дневной коэффициент считается по заданному часовому поясу модели."""
assert hour_factor(
datetime(2026, 1, 1, 2, 30, tzinfo=timezone.utc),
"Europe/Moscow",
) == 0.7
assert hour_factor(
datetime(2026, 1, 1, 6, 30, tzinfo=timezone.utc),
"Europe/Moscow",
) == 1.2
def test_event_budget_mean_follows_lambda_and_hour_factor(self, base_config):
"""Средний событийный бюджет следует λ и часовому коэффициенту."""
config = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=30,
jitter_pct=0,
min_events_per_tick=1,
max_events_per_tick=100,
model_t0=datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
rng = random.Random(config.seed)
samples = [calculate_events_count(config, rng) for _ in range(1000)]
mean_budget = sum(samples) / len(samples)
assert mean_budget == pytest.approx(30 * 1.2, rel=0.15)
def test_default_tick_budget_floor_does_not_outgrow_target_lambda(self, base_config):
"""Дефолтная нижняя граница бюджета не разгоняет lambda=30 на тике 5 секунд."""
config = replace(
base_config,
tick_seconds=5,
lambda_base_per_min=30,
jitter_pct=0,
max_events_per_tick=50,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
rng = random.Random(config.seed)
samples = [calculate_events_count(config, rng) for _ in range(1000)]
events_per_minute = sum(samples) / len(samples) * 60 / config.tick_seconds
assert events_per_minute == pytest.approx(config.lambda_base_per_min, rel=0.20)
def test_event_budget_uses_model_tick_duration(self, base_config):
"""При ускорении событийный бюджет растёт по модельной длительности тика."""
base = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=30,
jitter_pct=0,
min_events_per_tick=1,
max_events_per_tick=10_000,
model_time_speed=1,
)
accelerated = replace(base, model_time_speed=10)
model_tick_at = datetime(2026, 1, 1, 7, 0)
normal_rng = random.Random(base.seed)
accelerated_rng = random.Random(accelerated.seed)
normal_samples = [
calculate_events_count(base, normal_rng, now=model_tick_at)
for _ in range(300)
]
accelerated_samples = [
calculate_events_count(accelerated, accelerated_rng, now=model_tick_at)
for _ in range(300)
]
normal_mean = sum(normal_samples) / len(normal_samples)
accelerated_mean = sum(accelerated_samples) / len(accelerated_samples)
assert accelerated_mean / normal_mean == pytest.approx(10, rel=0.15)
def test_large_event_budget_does_not_stick_on_knuth_underflow(self, base_config):
"""Для λ > 1000 средний бюджет растёт вместе с целевой интенсивностью."""
config = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=1200,
jitter_pct=0,
min_events_per_tick=0,
max_events_per_tick=10_000,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
rng = random.Random(config.seed)
samples = [calculate_events_count(config, rng) for _ in range(500)]
mean_budget = sum(samples) / len(samples)
assert mean_budget == pytest.approx(1200, rel=0.05)
assert mean_budget > 1000
def test_event_generator_hour_factor_defaults_to_model_t0(
self, event_dictionary, base_config
):
"""Wrapper без аргумента берёт модельную точку, а не настенный час."""
config = replace(
base_config,
model_t0=datetime(2026, 1, 1, 6, 30, tzinfo=timezone.utc),
model_timezone="Europe/Moscow",
)
generator = EventGenerator(event_dictionary, config)
assert generator._hour_factor() == 1.2
def test_calculate_events_respects_bounds(self, event_dictionary, base_config):
"""Расчет количества событий уважает границы."""
generator = EventGenerator(event_dictionary, base_config)
samples = [generator._calculate_events_count() for _ in range(100)]
assert all(s >= base_config.min_events_per_tick for s in samples)
assert all(s <= base_config.max_events_per_tick for s in samples)
def test_jitter_increases_variance(self, event_dictionary, base_config):
"""Jitter увеличивает дисперсию."""
config_with_jitter = replace(
base_config,
tick_seconds=60,
lambda_base_per_min=30,
jitter_pct=50,
min_events_per_tick=1,
max_events_per_tick=100,
model_t0=datetime(2026, 1, 1, 7, 0, tzinfo=timezone.utc),
model_timezone="UTC",
)
config_without_jitter = replace(config_with_jitter, jitter_pct=0)
gen_with = EventGenerator(event_dictionary, config_with_jitter)
gen_without = EventGenerator(event_dictionary, config_without_jitter)
samples_with = [gen_with._calculate_events_count() for _ in range(200)]
samples_without = [gen_without._calculate_events_count() for _ in range(200)]
mean_with = sum(samples_with) / len(samples_with)
mean_without = sum(samples_without) / len(samples_without)
var_with = sum((x - mean_with) ** 2 for x in samples_with) / len(samples_with)
var_without = sum((x - mean_without) ** 2 for x in samples_without) / len(samples_without)
assert var_with > var_without, \
f"Jitter should increase variance: {var_with} vs {var_without}"
def test_mean_is_reasonable(self, event_dictionary, base_config):
"""Среднее значение в разумных пределах."""
config = replace(
base_config,
jitter_pct=0,
min_events_per_tick=1,
)
generator = EventGenerator(event_dictionary, config)
samples = [generator._calculate_events_count() for _ in range(500)]
mean = sum(samples) / len(samples)
# Ожидаем: lambda_base * tick_seconds / 60 * hour_factor
# hour_factor обычно 0.7-1.2
expected_base = config.lambda_base_per_min * config.tick_seconds / 60.0
# Допустимое отклонение до 50%
assert mean > expected_base * 0.5, f"Mean {mean} too low (expected ~{expected_base})"
assert mean < expected_base * 1.5, f"Mean {mean} too high (expected ~{expected_base})"
class TestEmptyData:
"""Тесты обработки пустых данных."""
def test_empty_jsonl_raises_error(self, empty_temp_dir):
"""Пустые JSONL файлы вызывают ValueError при загрузке."""
with pytest.raises(ValueError, match="browser_events.jsonl is empty"):
EventDictionary.load(empty_temp_dir)