diff --git a/.scratch/generator-model-time-startup-history/issues/03-model-speed-and-day-factor.md b/.scratch/generator-model-time-startup-history/issues/03-model-speed-and-day-factor.md index 84a2c6a..f9a576e 100644 --- a/.scratch/generator-model-time-startup-history/issues/03-model-speed-and-day-factor.md +++ b/.scratch/generator-model-time-startup-history/issues/03-model-speed-and-day-factor.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # ×K и дневной коэффициент по модельному времени @@ -19,21 +19,68 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] На ×1 поведение остаётся совместимым с обычным живым режимом. -- [ ] На ×K модельные `event_timestamp` за короткий реальный прогон покрывают +- [x] На ×1 поведение остаётся совместимым с обычным живым режимом. +- [x] На ×K модельные `event_timestamp` за короткий реальный прогон покрывают примерно `K` раз большую модельную длительность. -- [ ] Событийный бюджет считается по модельной длительности тика, а не по +- [x] Событийный бюджет считается по модельной длительности тика, а не по реальному времени сна процесса. -- [ ] Дневной коэффициент меняется при переходе модельного времени через +- [x] Дневной коэффициент меняется при переходе модельного времени через дневные/ночные часы в зафиксированном часовом поясе, независимо от реального часа запуска. -- [ ] ClickHouse-проверка показывает повторяемые контрольные числа при тех же +- [x] ClickHouse-проверка показывает повторяемые контрольные числа при тех же `GEN_SEED`, `T0`, скорости и настройках. -- [ ] Worker даёт промежуточный статус, если статистический прогон или стендовая +- [x] Worker даёт промежуточный статус, если статистический прогон или стендовая проверка занимает заметное время. -- [ ] В `generator/README.md` или `docs/OPERATIONS.md` кратко описано, что ×K +- [x] В `generator/README.md` или `docs/OPERATIONS.md` кратко описано, что ×K ускоряет именно модельное время стенда. +## Решение + +`calculate_events_count()` теперь считает λ тика по модельной длительности: +`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. При `GEN_MODEL_TIME_SPEED=1` +формула остаётся прежней. + +Для больших λ старый Knuth-алгоритм заменён нормальным приближением, чтобы +ускорение ×K не упиралось в underflow `exp(-λ)` и событийный бюджет продолжал +расти вместе с модельной длительностью. + +Добавлены проверки: + +- средний бюджет ×10 примерно в 10 раз больше бюджета ×1; +- большой λ не залипает на старом underflow-ограничении; +- `hour_factor()` берёт час из заданного `GEN_MODEL_TIMEZONE`; +- сервисный live-тик при ×10 сдвигает `event_timestamp` с `10:00` на `10:10`; +- старый live-сценарий ×1 остаётся зелёным. + +Документация в `generator/README.md` и `docs/OPERATIONS.md` уточняет, что ×K +ускоряет модельную длительность тика и событийный бюджет. + +## Проверка + +- `make generator-test` — `120 passed`. +- ClickHouse, быстрый прогон ×60 с `GEN_STATE_RESET=true`, + `GEN_MODEL_T0=2026-01-01T05:58:00+00:00`, `GEN_TICK_SECONDS=1`: за короткий + реальный прогон `event_ts` покрыл `2026-01-01 05:58:00` … + `2026-01-01 06:09:00`; после 06:00 UTC поминутные числа выросли с `40/42` до + примерно `58..63`, что показывает смену ночного коэффициента. +- ClickHouse, повторяемость: два чистых запуска с одинаковыми + `GEN_SEED=4242`, `GEN_MODEL_T0=2026-01-01T10:00:00+00:00`, + `GEN_MODEL_TIME_SPEED=60`, `GEN_TICK_SECONDS=60`, + `GEN_STATE_RESET=true` дали одинаковые контрольные числа: + `events=73`, `unique_events=73`, `unique_clicks=73`, + `min_event_ts=max_event_ts=2026-01-01 10:00:00`. + +## Риски и границы + +- State v2 resume, backfill, manifest и startup history не менялись. Известные + риски по ним остаются задачами 04/05. +- Сервис всё ещё пишет операционные метки state и истории пачек по настенному + времени; это соответствует текущему контракту issue 03. +- Нормальное приближение для больших λ не является точным Poisson, но для + событийного бюджета учебного демо достаточно сохраняет масштаб, среднее и + дисперсию. Точную статистическую выборку можно выделить в отдельную задачу, + если она станет учебной целью. + ## Blocked by - `.scratch/generator-model-time-startup-history/issues/02-model-time-to-clickhouse.md` diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 6dba349..2cb4684 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -123,7 +123,9 @@ GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator Для повторяемой проверки используйте чистый стенд и явный сброс состояния генератора. `event_timestamp` в событиях — модельное время от `GEN_MODEL_T0`, а -не настенное время запуска процесса. +не настенное время запуска процесса. При `GEN_MODEL_TIME_SPEED=K` один тик +покрывает `GEN_TICK_SECONDS * K` модельных секунд, и событийный бюджет считается +по этой модельной длительности. ```bash make clean diff --git a/generator/README.md b/generator/README.md index e7a9a19..17d61c7 100644 --- a/generator/README.md +++ b/generator/README.md @@ -92,8 +92,10 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator В живом режиме `event_timestamp` берётся из модельного времени: первый чистый тик стартует от `GEN_MODEL_T0`, дальше модельная точка сдвигается на -`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Дневной коэффициент считается по -`GEN_MODEL_TIMEZONE`, а не по реальному часу запуска процесса. +`GEN_TICK_SECONDS * GEN_MODEL_TIME_SPEED`. Событийный бюджет тика считается по +этой же модельной длительности, поэтому ×K даёт больше событий за короткий +реальный прогон. Дневной коэффициент считается по `GEN_MODEL_TIMEZONE`, а не по +реальному часу запуска процесса. Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены безопасными внутренними значениями `kafka:29092` и `/data`. diff --git a/generator/src/clickstream_generator/generation.py b/generator/src/clickstream_generator/generation.py index 0f4ab69..1db9454 100644 --- a/generator/src/clickstream_generator/generation.py +++ b/generator/src/clickstream_generator/generation.py @@ -129,7 +129,7 @@ class EventGenerator: def _hour_factor(self, now: datetime | None = None) -> float: """Совместимый wrapper над расчётом часового коэффициента.""" - return hour_factor(now, self.config.model_timezone) + return hour_factor(now or self.config.model_t0, self.config.model_timezone) def _calculate_events_count(self, now: datetime | None = None) -> int: """Совместимый wrapper над расчётом событийного бюджета.""" diff --git a/generator/src/clickstream_generator/intensity.py b/generator/src/clickstream_generator/intensity.py index 9e3c921..e718ad0 100644 --- a/generator/src/clickstream_generator/intensity.py +++ b/generator/src/clickstream_generator/intensity.py @@ -2,17 +2,22 @@ import math import random -from datetime import datetime, timezone +from datetime import datetime from zoneinfo import ZoneInfo from clickstream_generator.config import Config +POISSON_KNUTH_MAX_LAMBDA = 100.0 + + def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> float: """Возвращает коэффициент интенсивности в зависимости от часа дня.""" - current = now or datetime.now(timezone.utc) + current = now + if current is None: + raise ValueError("now is required for hour_factor") if current.tzinfo is None: - current = current.replace(tzinfo=timezone.utc) + current = current.replace(tzinfo=ZoneInfo("UTC")) current = current.astimezone(ZoneInfo(model_timezone)) hour = current.hour if 9 <= hour <= 18: @@ -22,18 +27,10 @@ def hour_factor(now: datetime | None = None, model_timezone: str = "UTC") -> flo return 1.0 -def calculate_events_count( - config: Config, - rng: random.Random, - now: datetime | None = None, -) -> int: - """Вычисляет количество событий для текущего тика (Poisson + jitter).""" - if now is None: - factor = hour_factor() - else: - factor = hour_factor(now, config.model_timezone) - lambda_minute = config.lambda_base_per_min * factor - lambda_tick = lambda_minute * (config.tick_seconds / 60.0) +def _sample_poisson(lambda_tick: float, rng: random.Random) -> int: + """Разыгрывает Poisson без underflow на больших λ.""" + if lambda_tick >= POISSON_KNUTH_MAX_LAMBDA: + return max(0, round(rng.gauss(lambda_tick, math.sqrt(lambda_tick)))) count = 0 threshold = math.exp(-lambda_tick) @@ -41,7 +38,22 @@ def calculate_events_count( while product > threshold: product *= rng.random() count += 1 - count -= 1 + return count - 1 + + +def calculate_events_count( + config: Config, + rng: random.Random, + now: datetime | None = None, +) -> int: + """Вычисляет количество событий для текущего тика (Poisson + jitter).""" + model_time = now or config.model_t0 + factor = hour_factor(model_time, config.model_timezone) + lambda_minute = config.lambda_base_per_min * factor + model_tick_seconds = config.tick_seconds * config.model_time_speed + lambda_tick = lambda_minute * (model_tick_seconds / 60.0) + + count = _sample_poisson(lambda_tick, rng) if config.jitter_pct > 0: jitter_factor = 1.0 + rng.uniform( diff --git a/generator/tests/test_generation.py b/generator/tests/test_generation.py index 2971c57..a20f2d7 100644 --- a/generator/tests/test_generation.py +++ b/generator/tests/test_generation.py @@ -6,7 +6,7 @@ import json import random import uuid from dataclasses import replace -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone import pytest from generator import ( @@ -16,6 +16,7 @@ from generator import ( TickStreamGenerator, calculate_events_count, generate_tick_batch, + hour_factor, ) @@ -1030,14 +1031,19 @@ class TestEventGeneration: class TestPoissonDistribution: """Тесты статистической модели.""" - def test_event_budget_mean_follows_lambda_and_hour_factor( - self, base_config, monkeypatch - ): + 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): """Средний событийный бюджет следует λ и часовому коэффициенту.""" - monkeypatch.setattr( - "clickstream_generator.intensity.hour_factor", - lambda: 1.2, - ) config = replace( base_config, tick_seconds=60, @@ -1045,6 +1051,8 @@ class TestPoissonDistribution: 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) @@ -1053,20 +1061,16 @@ class TestPoissonDistribution: 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, monkeypatch - ): + def test_default_tick_budget_floor_does_not_outgrow_target_lambda(self, base_config): """Дефолтная нижняя граница бюджета не разгоняет lambda=30 на тике 5 секунд.""" - monkeypatch.setattr( - "clickstream_generator.intensity.hour_factor", - lambda: 1.0, - ) 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) @@ -1075,6 +1079,69 @@ class TestPoissonDistribution: 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) @@ -1084,10 +1151,21 @@ class TestPoissonDistribution: 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, config_no_jitter): + def test_jitter_increases_variance(self, event_dictionary, base_config): """Jitter увеличивает дисперсию.""" - gen_with = EventGenerator(event_dictionary, base_config) - gen_without = EventGenerator(event_dictionary, config_no_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)] diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index a3ee7ee..510dc66 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -135,6 +135,11 @@ class TestGeneratorServiceSteadyStream: ) assert len(night_wall_run["browser_events"]) == len(day_wall_run["browser_events"]) + assert night_wall_run["budget_model_times"] == [ + datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc), + datetime(2026, 1, 1, 10, 1, tzinfo=timezone.utc), + ] + assert day_wall_run["budget_model_times"] == night_wall_run["budget_model_times"] assert night_wall_run["browser_events"] timestamps = { event["event_timestamp"] @@ -145,6 +150,38 @@ class TestGeneratorServiceSteadyStream: "2026-01-01 10:01:00.000000", } + def test_service_live_tick_advances_event_timestamps_by_model_speed(self, base_config): + """При ×K сервис сдвигает события на ускоренный модельный шаг.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + config = replace( + base_config, + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=1, + max_active_sessions=250, + population_max=251, + model_t0=model_t0, + model_time_speed=10, + ) + + published = self._run_service_ticks( + config, + wall_now=datetime(2026, 6, 14, 11, 0, tzinfo=timezone.utc), + ticks_count=2, + ) + + timestamps = { + event["event_timestamp"] + for event in published["browser_events"] + } + assert timestamps == { + "2026-01-01 10:00:00.000000", + "2026-01-01 10:10:00.000000", + } + def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int): service = GeneratorService(config) service.publisher = MagicMock() @@ -154,6 +191,8 @@ class TestGeneratorServiceSteadyStream: service.history = MagicMock() service._running = True sleep_calls = 0 + budget_model_times = [] + original_calculate_events_count = service.generator._calculate_events_count class FrozenDateTime(datetime): @classmethod @@ -168,7 +207,16 @@ class TestGeneratorServiceSteadyStream: if sleep_calls >= ticks_count: service._running = False - with patch("clickstream_generator.intensity.datetime", FrozenDateTime), \ + def calculate_events_count(now=None): + budget_model_times.append(now) + return original_calculate_events_count(now=now) + + with patch("clickstream_generator.service.datetime", FrozenDateTime), \ + patch.object( + service.generator, + "_calculate_events_count", + side_effect=calculate_events_count, + ), \ patch("clickstream_generator.service.time.sleep", side_effect=stop_after_tick): service._main_loop() @@ -176,6 +224,7 @@ class TestGeneratorServiceSteadyStream: for call in service.publisher.publish.call_args_list: topic, events = call.args published.setdefault(topic, []).extend(events) + published["budget_model_times"] = budget_model_times return published def test_service_ticks_publish_connected_multi_event_visit(self, base_config):