- Зачем:
- режим кормления стенда порциями: менти триггерит «следующий день»,
видит полный цикл DWH за один шаг (задача 13, вариант 2 — генерация
от слепка T_end).
- Что:
- новый ограниченный режим next-day: восстановление мира из state,
генерация ровно [T_end, T_end+24h), публикация данные -> state ->
манифест (манифест — точка фиксации, автоотката нет).
- операция next-day в DAG generator_control: своя предпроверка границы
вместо clean-guard, идемпотентность через параметр expected_t_end.
- цепочка границ — накопительное поле boundaries в манифесте, старый
формат читается как [T0, T_end]; импорт не изменён.
- новая проверка цепочки (make generated-history-chain-check): непарные
счётчики и однородность по каждой границе, явный статус нулевого
стыка, хвост за границей по всем четырём топикам, литералы в UTC
с микросекундами.
- документация OPERATIONS.md: глагол, предпроверка, восстановление
после сбоя, ограничение retention; в задаче 13 — решения двух слепых
ревью постановки и кода с аргументами отклонений.
- Проверка:
- make test: 204 теста генератора + 31 контракт корня, зелёные.
- make generated-history-chain-check: зелёный, 2 внутренние границы,
непарные счётчики нулевые; учебный цикл: DM 322 -> 10026 -> 19196
за два next-day подряд.
- make generated-history-runtime-check (регрессия задачи 20): зелёный.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
1493 lines
60 KiB
Python
1493 lines
60 KiB
Python
"""
|
||
Тесты для 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
|
||
|
||
|
||
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."""
|
||
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 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)
|
||
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.assert_called_once()
|
||
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_state_then_cumulative_manifest(
|
||
self, base_config
|
||
):
|
||
"""Next-day пишет [T_end, T_end+24h), затем state и manifest."""
|
||
from clickstream_generator.startup_history_artifact import (
|
||
StartupHistoryArtifactBuilder,
|
||
build_manifest,
|
||
)
|
||
|
||
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)
|
||
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.side_effect = (
|
||
lambda _manifest: order.append("manifest")
|
||
)
|
||
|
||
class DataReader:
|
||
def load(self):
|
||
return {
|
||
topic: old_batch[topic] + published[topic]
|
||
for topic in old_batch
|
||
}
|
||
|
||
service.data_reader = DataReader()
|
||
service.restore_from_startup_history(state, model_t_end=current_t_end)
|
||
|
||
service._run_next_day(manifest)
|
||
|
||
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", "state", "manifest"]
|
||
|
||
first_day_events = len(browser_events)
|
||
first_day_checksum = saved_manifest["topics"]["browser_events"][
|
||
"checksum_sha256"
|
||
]
|
||
service._run_next_day(saved_manifest)
|
||
|
||
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(DataReader().load())
|
||
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 order == [
|
||
"data",
|
||
"state",
|
||
"manifest",
|
||
"data",
|
||
"state",
|
||
"manifest",
|
||
]
|
||
|
||
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()],
|
||
}
|
||
|
||
with pytest.raises(RuntimeError, match="Публикация next-day не удалась"):
|
||
service._run_next_day(manifest)
|
||
|
||
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
|