""" Тесты для GeneratorService и интеграционных сценариев. """ import logging import random import hashlib import json import signal from dataclasses import replace from datetime import datetime, timedelta, timezone from time import sleep as real_sleep from unittest.mock import MagicMock, patch import pytest from generator import ( Config, EventDictionary, EventGenerator, EXPECTED_VISIT_EVENTS, GeneratorService, GeneratorState, KafkaBatchHistory, TickStreamGenerator, main, ) from clickstream_generator.state import UnsupportedStateVersionError from clickstream_generator.service import IncompatibleStateError class TestBatchRecordWithDictConversion: """Тесты конвертации BatchRecord в dict.""" def test_dict_contains_all_batch_info(self): """Словарь содержит всю информацию о батче.""" from datetime import datetime, timezone from generator import BatchRecord started = datetime(2024, 6, 15, 12, 0, 0, tzinfo=timezone.utc) finished = datetime(2024, 6, 15, 12, 0, 5, tzinfo=timezone.utc) record = BatchRecord( batch_id="batch_001", started_at=started, finished_at=finished, sent_total=400, sent_browser=100, sent_location=100, sent_device=100, sent_geo=100, status="success", error_message=None, ) data = record.to_dict() # Проверяем структуру JSON assert isinstance(data, dict) assert data["batch_id"] == "batch_001" assert data["sent_total"] == 400 assert data["status"] == "success" assert "started_at" in data assert "finished_at" in data # Проверяем что можно сериализовать в JSON import json json_str = json.dumps(data) assert isinstance(json_str, str) # Проверяем что можно десериализовать restored = json.loads(json_str) assert restored["batch_id"] == "batch_001" class TestGeneratorServiceInit: """Тесты инициализации GeneratorService.""" def test_service_initializes_dictionary(self, base_config): """Service загружает словарь при инициализации.""" service = GeneratorService(base_config) assert service.dictionary is not None assert len(service.dictionary.browser_events) == 1000 assert service.generator is not None assert service.config == base_config def test_service_history_is_none_before_start(self, base_config): """История None до вызова start.""" service = GeneratorService(base_config) # История и publisher инициализируются в start() assert service.history is None assert service.publisher is None class TestGeneratorServiceDisabled: """Тесты отключенного генератора.""" def test_disabled_generator_logs_warning(self, base_config, caplog): """Отключенный генератор логирует warning.""" from dataclasses import replace import logging disabled_config = replace(base_config, enabled=False) service = GeneratorService(disabled_config) with caplog.at_level(logging.WARNING): service.start() assert "disabled" in caplog.text.lower() or "GEN_ENABLED" in caplog.text class TestGeneratorServiceSteadyStream: """Проверки сервисного тика без настоящей Kafka.""" def test_service_live_tick_uses_model_time_for_events_and_day_factor(self, base_config): """Живой тик пишет события от T0 и не зависит от реального часа запуска.""" 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, ) night_wall_run = self._run_service_ticks( config, wall_now=datetime(2026, 6, 14, 3, 0, tzinfo=timezone.utc), ticks_count=2, ) day_wall_run = self._run_service_ticks( config, wall_now=datetime(2026, 6, 14, 11, 0, tzinfo=timezone.utc), ticks_count=2, ) 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"] for event in night_wall_run["browser_events"] } assert timestamps == { "2026-01-01 10:00:00.000000", "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 test_successful_live_tick_does_not_log_per_tick_info(self, base_config, caplog): """Успешный live-тик не пишет подробный пер-тиковый журнал на INFO.""" config = replace( base_config, tick_seconds=1, lambda_base_per_min=60, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=1, max_active_sessions=10, population_max=11, ) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) service.history = MagicMock() service._running = True def stop_after_first_tick(_sleep_seconds): service._running = False with caplog.at_level(logging.INFO, logger="generator"), \ patch.object( service._shutdown_event, "wait", side_effect=stop_after_first_tick, ): service._main_loop() assert "=== Tick" not in caplog.text assert "Generating with event budget" not in caplog.text assert "completed:" not in caplog.text assert "sent" not in caplog.text def test_sigterm_during_publish_finishes_current_batch(self, base_config): """SIGTERM посреди публикации не обрывает связанные топики batch.""" 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=3, metrics_enabled=False, state_enabled=False, ) service = GeneratorService(config) publisher = MagicMock() history = MagicMock() published_topics = [] def publish_and_request_stop(topic, events): published_topics.append(topic) if len(published_topics) == 1: service.request_stop(signal.SIGTERM, None) return len(events), 0 publisher.publish.side_effect = publish_and_request_stop with patch( "clickstream_generator.service.start_http_server", side_effect=AssertionError("тест не должен открывать порт метрик"), ), \ patch("clickstream_generator.service.ensure_topics"), \ patch( "clickstream_generator.service.KafkaPublisher", return_value=publisher, ), \ patch( "clickstream_generator.service.KafkaBatchHistory", return_value=history, ), \ patch.object( service.generator, "_calculate_events_count", return_value=int(EXPECTED_VISIT_EVENTS), ): service.start() assert published_topics == [ "browser_events", "location_events", "device_events", "geo_events", ] publisher.flush.assert_called_once_with() history.add.assert_called_once() history.flush.assert_called_once_with() publisher.close.assert_called_once_with() history.close.assert_called_once_with() assert service._tick == 1 def test_main_connects_sigterm_to_graceful_stop(self, base_config): """Точка входа передаёт SIGTERM сервису как штатный запрос остановки.""" service = MagicMock() with patch("clickstream_generator.service.Config", return_value=base_config), \ patch("clickstream_generator.service.GeneratorService", return_value=service), \ patch("clickstream_generator.service.signal.signal") as register_signal: main() register_signal.assert_called_once_with(signal.SIGTERM, service.request_stop) service.start.assert_called_once_with() def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int): service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) 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 def now(cls, tz=None): if tz is None: return wall_now.replace(tzinfo=None) return wall_now.astimezone(tz) def stop_after_tick(_sleep_seconds): nonlocal sleep_calls sleep_calls += 1 if sleep_calls >= ticks_count: service._running = False 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.object( service._shutdown_event, "wait", side_effect=stop_after_tick, ): service._main_loop() published = {} 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): """Сервисные тики публикуют несколько связанных событий одного визита.""" config = replace(base_config, tick_seconds=1, max_session_events=3) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) service.history = MagicMock() service._running = True sleep_calls = 0 def stop_after_second_tick(sleep_seconds): nonlocal sleep_calls sleep_calls += 1 if sleep_calls == 1: real_sleep(sleep_seconds) else: service._running = False with patch.object( service.generator, "_calculate_events_count", side_effect=[int(EXPECTED_VISIT_EVENTS), 0], ), patch.object( service.generator, "_visit_pause_seconds", return_value=0.05, ), patch.object(service._shutdown_event, "wait") as sleep_mock: sleep_mock.side_effect = stop_after_second_tick service._main_loop() published = {} for call in service.publisher.publish.call_args_list: topic, events = call.args published.setdefault(topic, []).extend(events) browser_events = published["browser_events"] location_events = published["location_events"] device_events = published["device_events"] geo_events = published["geo_events"] assert set(published) == { "browser_events", "location_events", "device_events", "geo_events", } assert len(browser_events) >= 2 assert len({event["click_id"] for event in browser_events}) == 1 assert len({event["event_id"] for event in browser_events}) == len(browser_events) assert {event["event_id"] for event in location_events} == { event["event_id"] for event in browser_events } assert {event["click_id"] for event in device_events} == { browser_events[0]["click_id"] } assert {event["click_id"] for event in geo_events} == { browser_events[0]["click_id"] } history_records = [ call.args[0] for call in service.history.add.call_args_list ] assert [record.status for record in history_records] == ["success", "success"] assert sum(record.sent_browser for record in history_records) == len(browser_events) assert sum(record.sent_location for record in history_records) == len(location_events) assert sum(record.sent_device for record in history_records) == len(device_events) assert sum(record.sent_geo for record in history_records) == len(geo_events) class TestGeneratorServiceBackfill: """Проверки режима промотки стартовой истории.""" def test_backfill_publishes_half_open_history_state_and_manifest( self, base_config ): """Backfill пишет [T0, T_end), state на T_end и повторяемый manifest.""" from clickstream_generator.startup_history_artifact import ( cumulative_counter_reference, ) model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=3) config = replace( base_config, run_mode="backfill", model_t0=model_t0, model_t_end=model_t_end, tick_seconds=60, lambda_base_per_min=600, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, state_enabled=True, ) first = self._run_backfill(config) second = self._run_backfill(config) browser_events = first["published"]["browser_events"] timestamps = [ datetime.fromisoformat(event["event_timestamp"].replace(" ", "T")) for event in browser_events ] saved_state = first["state_manager"].save.call_args.args[0] manifest = first["manifest_manager"].save.call_args.args[0] assert browser_events assert min(timestamps) >= model_t0.replace(tzinfo=None) assert max(timestamps) < model_t_end.replace(tzinfo=None) assert saved_state.model_timestamp == model_t_end assert saved_state.last_timestamp == model_t_end assert manifest["run_mode"] == "backfill" assert manifest["model_t0"] == model_t0.isoformat() assert manifest["model_t_end"] == model_t_end.isoformat() assert manifest["state"]["last_batch_id"] == saved_state.last_batch_id assert saved_state.cumulative_manifest_counters == ( cumulative_counter_reference( manifest["cumulative_manifest_counters"] ) ) assert manifest["topics"]["browser_events"]["rows"] == len(browser_events) for topic in ( "browser_events", "location_events", "device_events", "geo_events", ): assert manifest["topics"][topic]["min_event_timestamp"] is not None assert manifest["topics"][topic]["max_event_timestamp"] is not None assert manifest["totals"]["events"] == len(browser_events) assert manifest["totals"]["visits"] == len({ event["click_id"] for event in browser_events }) assert manifest["totals"]["users"] == len({ event["user_domain_id"] for event in first["published"]["device_events"] }) assert first["digest"] == second["digest"] assert first["manifest_digest"] == second["manifest_digest"] def test_backfill_writes_portable_artifact_when_requested( self, base_config, tmp_path ): """Backfill пишет переносимый файл с событиями, state и manifest.""" from clickstream_generator.startup_history_artifact import ( load_startup_history_artifact, validate_startup_history_artifact, ) model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=1) artifact_path = tmp_path / "startup-history.json" config = replace( base_config, run_mode="backfill", model_t0=model_t0, model_t_end=model_t_end, tick_seconds=60, lambda_base_per_min=600, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, state_enabled=True, startup_history_artifact=artifact_path, ) self._run_backfill(config) artifact = load_startup_history_artifact(artifact_path) assert "cumulative_manifest_counters" not in artifact["state"] assert "cumulative_manifest_counters" not in artifact["manifest"] state, manifest, topics, _ = validate_startup_history_artifact( artifact, expected_config=config, ) assert topics["browser_events"] assert state.last_batch_id == manifest["state"]["last_batch_id"] assert manifest["totals"]["events"] == len(topics["browser_events"]) def test_live_start_uses_startup_manifest_without_wall_delta( self, base_config, event_dictionary ): """Live-запуск из backfill-state стартует с T_end без wall-дельты.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=5) config = replace( base_config, model_t0=model_t0, tick_seconds=60, model_time_speed=3600, ) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="startup-history-state", last_timestamp=model_t_end, model_timestamp=model_t_end, wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) manifest = { "run_mode": "backfill", "gen_seed": config.seed, "model_t0": model_t0.isoformat(), "model_t_end": model_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), "state_version": state.version, "state": {"last_batch_id": state.last_batch_id}, } state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = manifest 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(config) service.start() assert service._model_time == model_t_end assert service._tick == state.tick def test_live_start_rejects_startup_manifest_with_different_config_t_end( self, base_config, event_dictionary ): """Manifest от другого T_end не считается стартовой историей запуска.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) manifest_t_end = model_t0 + timedelta(minutes=5) config_t_end = model_t0 + timedelta(minutes=10) wall_saved_at = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) wall_restarted_at = wall_saved_at + timedelta(seconds=30) config = replace( base_config, model_t0=model_t0, model_t_end=config_t_end, tick_seconds=60, model_time_speed=10, ) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="startup-history-state", last_timestamp=manifest_t_end, model_timestamp=manifest_t_end, wall_timestamp=wall_saved_at, model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) manifest = { "run_mode": "backfill", "gen_seed": config.seed, "model_t0": model_t0.isoformat(), "model_t_end": manifest_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), "state_version": state.version, "state": {"last_batch_id": state.last_batch_id}, } state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = manifest class FrozenDateTime(datetime): @classmethod def now(cls, tz=None): if tz is None: return wall_restarted_at.replace(tzinfo=None) return wall_restarted_at.astimezone(tz) 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch("clickstream_generator.service.datetime", FrozenDateTime), \ patch.object( GeneratorService, "restore_from_startup_history", ) as restore_from_startup_history, \ patch.object( GeneratorService, "_restore_live_state", ) as restore_live_state, \ patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(config) with pytest.raises(ValueError, match="GEN_MODEL_T_END.*GEN_STATE_RESET=true"): service.start() restore_from_startup_history.assert_not_called() restore_live_state.assert_not_called() def test_live_start_fails_loudly_on_readable_incompatible_state( self, base_config, event_dictionary ): """Читаемый state другого мира при продолжении даёт жёсткий отказ.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) config = replace( base_config, model_t0=model_t0, tick_seconds=60, model_time_speed=10, state_reset=False, ) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="live-state", last_timestamp=model_t0, model_timestamp=model_t0, wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed + 1, ) state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = None 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(config) with pytest.raises(ValueError, match="GEN_SEED.*GEN_STATE_RESET=true"): service.start() def test_live_start_rejects_orphan_startup_state_without_live_restore( self, base_config, event_dictionary, caplog ): """Orphan startup-history state не восстанавливается как live-state.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=5) config = replace(base_config, model_t0=model_t0, model_time_speed=10) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="startup-history-orphan", last_timestamp=model_t_end, model_timestamp=model_t_end, wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = None 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch.object( GeneratorService, "restore_from_startup_history", ) as restore_from_startup_history, \ patch.object( GeneratorService, "_restore_live_state", ) as restore_live_state, \ patch.object(GeneratorService, "_main_loop", return_value=None), \ caplog.at_level(logging.WARNING, logger="generator"): service = GeneratorService(config) with pytest.raises( ValueError, match="generator_startup_history_manifest.*GEN_STATE_RESET=true", ): service.start() restore_from_startup_history.assert_not_called() restore_live_state.assert_not_called() def test_startup_history_state_checks_state_fields( self, base_config, event_dictionary ): """Startup-history state сверяется с manifest/config по полям state.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=5) config = replace(base_config, model_t0=model_t0, model_time_speed=10) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="startup-history-state", last_timestamp=model_t_end, model_timestamp=model_t_end, wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed + 1, ) manifest = { "run_mode": "backfill", "gen_seed": config.seed, "model_t0": model_t0.isoformat(), "model_t_end": model_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), "state_version": state.version, "state": {"last_batch_id": state.last_batch_id}, } service = GeneratorService(config) assert not service._is_startup_history_state(state, manifest) def test_startup_history_state_checks_state_version( self, base_config, event_dictionary ): """Startup-history state сверяется с версией state из manifest.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=5) config = replace(base_config, model_t0=model_t0, model_time_speed=10) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="startup-history-state", last_timestamp=model_t_end, model_timestamp=model_t_end, wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) manifest = { "run_mode": "backfill", "gen_seed": config.seed, "model_t0": model_t0.isoformat(), "model_t_end": model_t_end.isoformat(), "model_timezone": config.model_timezone, "generation_settings": GeneratorService(config)._generation_settings(), "state_version": "1.0", "state": {"last_batch_id": state.last_batch_id}, } service = GeneratorService(config) assert not service._is_startup_history_state(state, manifest) def test_backfill_publish_error_does_not_save_state_or_manifest( self, base_config ): """Backfill не создаёт валидный артефакт при ошибке публикации.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) config = replace( base_config, run_mode="backfill", model_t0=model_t0, model_t_end=model_t0 + timedelta(minutes=1), tick_seconds=60, lambda_base_per_min=600, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, state_enabled=True, ) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 1) if topic == "location_events" else (len(events), 0) ) service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() with pytest.raises(RuntimeError, match="Backfill publish failed"): service._run_backfill() service.state_manager.save.assert_not_called() service.state_manager.flush.assert_not_called() service.manifest_manager.save.assert_not_called() service.manifest_manager.flush.assert_not_called() def test_backfill_manifest_save_error_does_not_save_state(self, base_config): """Если manifest не записан, state стартовой истории не сохраняется.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) config = replace( base_config, run_mode="backfill", model_t0=model_t0, model_t_end=model_t0 + timedelta(minutes=1), tick_seconds=60, lambda_base_per_min=600, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, state_enabled=True, ) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) service.publisher.flush.return_value = None service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() service.manifest_manager.save.side_effect = RuntimeError("manifest down") with pytest.raises(RuntimeError, match="manifest down"): service._run_backfill() service.manifest_manager.save.assert_called_once() service.state_manager.save.assert_not_called() service.state_manager.flush.assert_not_called() def test_backfill_manifest_flush_error_does_not_save_state(self, base_config): """Если manifest не сброшен в Kafka, state стартовой истории не сохраняется.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) config = replace( base_config, run_mode="backfill", model_t0=model_t0, model_t_end=model_t0 + timedelta(minutes=1), tick_seconds=60, lambda_base_per_min=600, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, state_enabled=True, ) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) service.publisher.flush.return_value = None service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() service.manifest_manager.flush.side_effect = RuntimeError("flush down") with pytest.raises(RuntimeError, match="flush down"): service._run_backfill() service.manifest_manager.save_counter_chunk.assert_called() service.manifest_manager.save.assert_not_called() service.manifest_manager.flush.assert_called_once() service.state_manager.save.assert_not_called() service.state_manager.flush.assert_not_called() def _run_backfill(self, config): service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = ( lambda topic, events: (len(events), 0) ) service.publisher.flush.return_value = None service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() service._run_backfill() published = {} for call in service.publisher.publish.call_args_list: topic, events = call.args published.setdefault(topic, []).extend(events) digest = hashlib.sha256( json.dumps(published, sort_keys=True, default=str).encode() ).hexdigest() manifest = service.manifest_manager.save.call_args.args[0] stable_manifest = { key: value for key, value in manifest.items() if key != "generated_at" } manifest_digest = hashlib.sha256( json.dumps(stable_manifest, sort_keys=True, default=str).encode() ).hexdigest() return { "published": published, "state_manager": service.state_manager, "manifest_manager": service.manifest_manager, "digest": digest, "manifest_digest": manifest_digest, } class TestGeneratorServiceNextDay: """Проверки ограниченной доливки следующего модельного дня.""" def test_next_day_publishes_24h_then_chunks_manifest_and_state( self, base_config ): """Next-day пишет день, фрагменты множеств, manifest и затем state.""" from clickstream_generator.startup_history_artifact import ( StartupHistoryArtifactBuilder, build_manifest, cumulative_counter_reference, ) model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) current_t_end = model_t0 + timedelta(hours=1) config = replace( base_config, run_mode="next-day", model_t0=model_t0, model_t_end=current_t_end, tick_seconds=3600, model_time_speed=1, lambda_base_per_min=60, jitter_pct=0, min_events_per_tick=1, max_events_per_tick=1000, max_session_events=5, max_active_sessions=250, population_max=251, ) service = GeneratorService(config) state = service.stream.to_state( tick=1, rng_state=service.generator.rng.getstate(), last_batch_id="startup-history-initial", last_timestamp=current_t_end, model_timestamp=current_t_end, wall_timestamp=current_t_end, model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) old_batch = { "browser_events": [{ "event_id": "old-event", "click_id": "old-click", "event_timestamp": "2026-01-01 00:30:00.000000", }], "location_events": [{"event_id": "old-event"}], "device_events": [{ "click_id": "old-click", "user_domain_id": "old-user", }], "geo_events": [{"click_id": "old-click"}], } initial_builder = StartupHistoryArtifactBuilder() initial_builder.add_batch(old_batch) state.cumulative_manifest_counters = cumulative_counter_reference( initial_builder.counters.to_state() ) manifest = build_manifest(config, initial_builder.counters, state) published = {topic: [] for topic in old_batch} service.publisher = MagicMock() def publish(topic, events): published[topic].extend(events) return len(events), 0 service.publisher.publish.side_effect = publish service.history = MagicMock() order = [] service.publisher.flush.side_effect = lambda: order.append("data") service.state_manager = MagicMock() service.state_manager.save.side_effect = lambda _state: order.append("state") service.manifest_manager = MagicMock() service.manifest_manager.save_counter_chunk.side_effect = ( lambda _chunk_sha256, _chunk: order.append("counter_chunk") ) service.manifest_manager.save.side_effect = ( lambda _manifest: order.append("manifest") ) class DataReader: def load(self): raise AssertionError("next-day не должен перечитывать Kafka") service.data_reader = DataReader() service.restore_from_startup_history(state, model_t_end=current_t_end) service._run_next_day(manifest, state) browser_events = published["browser_events"] timestamps = [ datetime.fromisoformat(event["event_timestamp"].replace(" ", "T")) for event in browser_events ] target_t_end = current_t_end + timedelta(hours=24) saved_state = service.state_manager.save.call_args.args[0] saved_manifest = service.manifest_manager.save.call_args.args[0] assert browser_events assert min(timestamps) >= current_t_end.replace(tzinfo=None) assert max(timestamps) < target_t_end.replace(tzinfo=None) assert saved_state.model_timestamp == target_t_end assert saved_manifest["model_t_end"] == target_t_end.isoformat() assert saved_manifest["boundaries"] == [ model_t0.isoformat(), current_t_end.isoformat(), target_t_end.isoformat(), ] assert saved_manifest["totals"]["events"] == len(browser_events) + 1 assert order == ["data", "counter_chunk", "manifest", "state"] first_day_events = len(browser_events) first_day_checksum = saved_manifest["topics"]["browser_events"][ "checksum_sha256" ] service._run_next_day(saved_manifest, saved_state) second_state = service.state_manager.save.call_args.args[0] second_manifest = service.manifest_manager.save.call_args.args[0] expected_builder = StartupHistoryArtifactBuilder() expected_builder.add_batch({ topic: old_batch[topic] + published[topic] for topic in old_batch }) assert second_state.model_timestamp == target_t_end + timedelta(hours=24) assert second_manifest["boundaries"] == [ model_t0.isoformat(), current_t_end.isoformat(), target_t_end.isoformat(), (target_t_end + timedelta(hours=24)).isoformat(), ] assert second_manifest["totals"]["events"] == len(browser_events) + 1 assert len(browser_events) > first_day_events assert ( second_manifest["topics"]["browser_events"]["checksum_sha256"] != first_day_checksum ) assert second_manifest["topics"] == expected_builder.counters.to_manifest_topics() assert second_manifest["totals"] == expected_builder.counters.to_manifest_totals() assert second_state.cumulative_manifest_counters == ( cumulative_counter_reference( second_manifest["cumulative_manifest_counters"] ) ) assert order == [ "data", "counter_chunk", "manifest", "state", "data", "counter_chunk", "manifest", "state", ] def test_next_day_requires_cumulative_counters_from_import(self, base_config): """Старый локальный state требует повторного импорта, а не чтения Kafka.""" model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) config = replace( base_config, run_mode="next-day", model_t0=model_t0, model_t_end=model_t0 + timedelta(hours=1), ) service = GeneratorService(config) state = service.stream.to_state( tick=1, rng_state=service.generator.rng.getstate(), last_batch_id="startup-history-old", last_timestamp=config.model_t_end, model_timestamp=config.model_t_end, model_t0=config.model_t0, ) service.publisher = MagicMock() service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() service._model_time = config.model_t_end manifest = { "model_t0": model_t0.isoformat(), "model_t_end": config.model_t_end.isoformat(), "boundaries": [model_t0.isoformat(), config.model_t_end.isoformat()], } with pytest.raises(IncompatibleStateError, match="повторно запустите import"): service._run_next_day(manifest, state) service.publisher.publish.assert_not_called() def test_next_day_publish_error_does_not_move_state_or_manifest(self, base_config): """Ошибка data-топика оставляет обе точки фиксации без изменений.""" model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) config = replace( base_config, run_mode="next-day", model_t0=model_t0, model_t_end=model_t0 + timedelta(hours=1), tick_seconds=3600, model_time_speed=1, max_active_sessions=250, population_max=251, ) service = GeneratorService(config) service.publisher = MagicMock() service.publisher.publish.side_effect = lambda topic, events: ( (len(events), 1) if topic == "location_events" else (len(events), 0) ) service.history = MagicMock() service.state_manager = MagicMock() service.manifest_manager = MagicMock() service._model_time = config.model_t_end manifest = { "model_t0": model_t0.isoformat(), "model_t_end": config.model_t_end.isoformat(), "boundaries": [model_t0.isoformat(), config.model_t_end.isoformat()], } from clickstream_generator.startup_history_artifact import ( ManifestCounters, cumulative_counter_reference, ) state = service.stream.to_state( tick=1, rng_state=service.generator.rng.getstate(), last_batch_id="startup-history-test", last_timestamp=config.model_t_end, model_timestamp=config.model_t_end, model_t0=config.model_t0, ) empty_counters = ManifestCounters() manifest["cumulative_manifest_counters"] = empty_counters.to_state() manifest["topics"] = empty_counters.to_manifest_topics() manifest["totals"] = empty_counters.to_manifest_totals() state.cumulative_manifest_counters = cumulative_counter_reference( manifest["cumulative_manifest_counters"] ) with pytest.raises(RuntimeError, match="Публикация next-day не удалась"): service._run_next_day(manifest, state) service.state_manager.save.assert_not_called() service.manifest_manager.save.assert_not_called() class TestGeneratorServiceState: """Тесты подключения state к сервисному запуску.""" def test_start_restores_tick_stream_state(self, base_config, event_dictionary): """Сервис восстанавливает популяцию и активные визиты из state.""" source_generator = EventGenerator(event_dictionary, base_config) source_stream = TickStreamGenerator(source_generator) tick_at = datetime.now(timezone.utc).replace(tzinfo=None) 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, model_timestamp=tick_at.replace(tzinfo=timezone.utc), wall_timestamp=tick_at.replace(tzinfo=timezone.utc), model_time_speed=base_config.model_time_speed, model_timezone=base_config.model_timezone, model_t0=base_config.model_t0, gen_seed=base_config.seed, ) state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = None 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_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_start_restores_model_time_from_state_wall_delta( self, base_config, event_dictionary ): """Live-восстановление считает точку модели из сохранённой wall-метки.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) wall_saved_at = datetime(2026, 6, 14, 12, 0, tzinfo=timezone.utc) wall_restarted_at = wall_saved_at + timedelta(seconds=30) config = replace( base_config, model_t0=model_t0, model_time_speed=10, tick_seconds=60, ) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) state = source_stream.to_state( tick=3, rng_state=source_generator.rng.getstate(), last_batch_id="batch-3", last_timestamp=model_t0, model_timestamp=model_t0, wall_timestamp=wall_saved_at, model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = None class FrozenDateTime(datetime): @classmethod def now(cls, tz=None): if tz is None: return wall_restarted_at.replace(tzinfo=None) return wall_restarted_at.astimezone(tz) 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch("clickstream_generator.service.datetime", FrozenDateTime), \ patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(config) service.start() assert service._tick == 3 assert service._model_time == model_t0 + timedelta(seconds=300) assert service.stream.active_visit_count == source_stream.active_visit_count def test_restore_from_startup_history_uses_passed_model_point_without_wall_delta( self, base_config, event_dictionary ): """Стартовая история продолжает с T_end, а не с wall-простоя.""" model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=5) old_wall_saved_at = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) config = replace( base_config, model_t0=model_t0, model_time_speed=3600, tick_seconds=60, max_session_events=5, ) source_generator = EventGenerator(event_dictionary, config) source_stream = TickStreamGenerator(source_generator) source_stream.generate_tick(event_budget=10, tick_started_at=model_t_end) state = source_stream.to_state( tick=99, rng_state=source_generator.rng.getstate(), last_batch_id="history-end", last_timestamp=model_t_end, model_timestamp=model_t_end, wall_timestamp=old_wall_saved_at, model_time_speed=config.model_time_speed, model_timezone=config.model_timezone, model_t0=config.model_t0, gen_seed=config.seed, ) service = GeneratorService(config) service.restore_from_startup_history(state, model_t_end=model_t_end) assert service._tick == 99 assert service._model_time == model_t_end assert service.stream.active_visit_count == source_stream.active_visit_count def test_save_state_writes_tick_stream_state(self, base_config): """Сервис сохраняет снимок тикового слоя.""" 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._model_time = tick_at.replace(tzinfo=timezone.utc) service._save_state("batch-1") saved_state = service.state_manager.save.call_args.args[0] assert saved_state.version == "3.0" assert saved_state.model_timestamp == tick_at.replace(tzinfo=timezone.utc) assert saved_state.wall_timestamp.tzinfo is not None assert saved_state.model_time_speed == base_config.model_time_speed assert saved_state.model_timezone == base_config.model_timezone assert saved_state.model_t0 == base_config.model_t0 assert saved_state.gen_seed == base_config.seed assert saved_state.population assert saved_state.active_visits service.state_manager.flush.assert_called_once() def test_incompatible_seed_state_fails_loudly(self, base_config, event_dictionary, caplog): """State от другого GEN_SEED даёт жёсткий отказ при продолжении.""" source_config = replace(base_config, seed=7) source_generator = EventGenerator(event_dictionary, source_config) source_stream = TickStreamGenerator(source_generator) tick_at = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) source_stream.generate_tick(event_budget=10, tick_started_at=tick_at) state = source_stream.to_state( tick=5, rng_state=source_generator.rng.getstate(), last_batch_id="other-seed", last_timestamp=tick_at, model_timestamp=tick_at, wall_timestamp=datetime(2026, 6, 14, 12, 0, tzinfo=timezone.utc), model_time_speed=base_config.model_time_speed, model_timezone=base_config.model_timezone, model_t0=base_config.model_t0, gen_seed=source_config.seed, ) state_manager = MagicMock() state_manager.load.return_value = state manifest_manager = MagicMock() manifest_manager.load.return_value = None 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_manager, ), \ patch.object(GeneratorService, "_main_loop", return_value=None), \ caplog.at_level(logging.WARNING, logger="generator"): service = GeneratorService(base_config) with pytest.raises(ValueError, match="GEN_SEED.*GEN_STATE_RESET=true"): service.start() assert "Readable generator state is incompatible" in caplog.text 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_known_old_state_version_fails_instead_of_starting_new_world( self, base_config, caplog ): """GEN_STATE_RESET=false не скрывает старый state за чистым стартом.""" state_manager = MagicMock() state_manager.load.side_effect = UnsupportedStateVersionError("2.0", "3.0") 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.ERROR, logger="generator"): service = GeneratorService(base_config) with pytest.raises(ValueError, match="state version 2.0.*GEN_STATE_RESET=true"): service.start() assert service._tick == 0 assert "Readable generator state uses an old format" in caplog.text def test_invalid_restored_state_starts_fresh(self, base_config, caplog): """Сервис не падает, если state ссылается на неизвестный профиль.""" state_manager = MagicMock() state_manager.load.return_value = GeneratorState( tick=9, rng_state=random.Random(42).getstate(), last_batch_id="bad-state", last_timestamp=datetime.now(timezone.utc), model_timestamp=base_config.model_t0, wall_timestamp=datetime.now(timezone.utc), model_time_speed=base_config.model_time_speed, model_timezone=base_config.model_timezone, model_t0=base_config.model_t0, gen_seed=base_config.seed, population=[ { "user_domain_id": "user-unknown", "seed_click_id": "missing-click-id", } ], active_visits=[], ) manifest_manager = MagicMock() manifest_manager.load.return_value = None 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( "clickstream_generator.service.KafkaStartupHistoryManifest", return_value=manifest_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