fix(generator): сохранена фактура визита при восстановлении

- Зачем:
  - визит после восстановления не должен менять браузер и источники перехода внутри одного click_id.
- Что:
  - добавлен base_click_id в state v3 для восстановления донора фактуры.
  - исправлено восстановление timestamp offset без потери микросекунд.
  - расширены тесты и стыковая проверка browser/source и device/os/geo.
- Проверка:
  - uv run --with-requirements generator/requirements.txt pytest generator/tests -q.
  - bash -n scripts/check_generated_analytics.sh.
  - git diff --cached --check.
This commit is contained in:
2026-07-04 21:46:23 +03:00
parent 0cbfe9b32e
commit e2d06841de
13 changed files with 464 additions and 77 deletions
+1 -1
View File
@@ -46,7 +46,7 @@ generated-history-analytics:
# Повторяемая проверка после прогона стартовой истории # Повторяемая проверка после прогона стартовой истории
generated-history-check: generated-history-check:
COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh CHECK_LIVE_SEAM="$${CHECK_LIVE_SEAM:-auto}" COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_analytics.sh
# Сгенерировать стартовую историю и сохранить портативный артефакт # Сгенерировать стартовую историю и сохранить портативный артефакт
startup-history-export: startup-history-export:
+8 -1
View File
@@ -147,7 +147,7 @@ make generator-logs
| `GEN_RUN_MODE` | Режим генератора | `live` | | `GEN_RUN_MODE` | Режим генератора | `live` |
| `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` | | `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` |
| `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | пусто | | `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | пусто |
| `GEN_STATE_ENABLED` | Сохранять state v2 между рестартами | `true` | | `GEN_STATE_ENABLED` | Сохранять state v3 между рестартами | `true` |
| `GEN_STATE_RESET` | Сбросить state при старте | `false` | | `GEN_STATE_RESET` | Сбросить state при старте | `false` |
В `docker-compose.yml` через окружение переопределяются демо-параметры модели и В `docker-compose.yml` через окружение переопределяются демо-параметры модели и
@@ -177,6 +177,13 @@ CI короткая повторная проверка после уже гот
make generated-history-check make generated-history-check
``` ```
После live-продолжения из `T_end` эта же команда автоматически включает проверку
стыка. Для принудительной проверки:
```bash
CHECK_LIVE_SEAM=1 GEN_LIVE_CHECK_MINUTES=10 make generated-history-check
```
По умолчанию команда использует быстрый профиль `ci`: 6 часов модельного По умолчанию команда использует быстрый профиль `ci`: 6 часов модельного
времени. Историю на 2 суток с суточной волной можно получить одной командой. времени. Историю на 2 суток с суточной волной можно получить одной командой.
В live-продолжении `daily-wave` идёт с ×60 и тикает раз в секунду, поэтому В live-продолжении `daily-wave` идёт с ×60 и тикает раз в секунду, поэтому
@@ -174,6 +174,11 @@ state-записей остаётся прежним контрактом воз
- `wall_timestamp` — настенная UTC-метка, когда это состояние было сохранено; - `wall_timestamp` — настенная UTC-метка, когда это состояние было сохранено;
- `model_time_speed`, `model_timezone` и `model_t0`. - `model_time_speed`, `model_timezone` и `model_t0`.
State v3 также хранит для каждого активного визита `base_click_id` — click_id
донора браузерной и source-фактуры из статического сида. При восстановлении
активный визит берёт per-event поля от этого донора; неизвестный донор считается
битым state, а не поводом выбрать запасную фактуру.
При восстановлении живого режима после сбоя модельная точка считается так: При восстановлении живого режима после сбоя модельная точка считается так:
```text ```text
+2 -2
View File
@@ -3,7 +3,7 @@
> **Статус (2026-06-11):** исторический дефект старой плоской генерации закрыт > **Статус (2026-06-11):** исторический дефект старой плоской генерации закрыт
> для режима `steady-stream`. Генератор строит визиты с общим `click_id`, > для режима `steady-stream`. Генератор строит визиты с общим `click_id`,
> монотонным временем событий, путём по страницам воронки, популяцией > монотонным временем событий, путём по страницам воронки, популяцией
> возвращающихся пользователей и состоянием v2 для активных визитов. > возвращающихся пользователей и состоянием для активных визитов.
> >
> Эта заметка больше не является предупреждением «генератор концептуально > Эта заметка больше не является предупреждением «генератор концептуально
> сломан». Она оставлена как учебный разбор старого дефекта и как место для > сломан». Она оставлена как учебный разбор старого дефекта и как место для
@@ -122,7 +122,7 @@ device_event = {**base_device, "click_id": new_click_id} # user_domain_i
времени. времени.
3. **Время событий** внутри визита строго растёт и не прилипает к одному 3. **Время событий** внутри визита строго растёт и не прилипает к одному
`now()` для всего батча. `now()` для всего батча.
4. **Состояние v2** сохраняет популяцию, активные визиты, накопленный бюджет 4. **Состояние** сохраняет популяцию, активные визиты, накопленный бюджет
рождения визитов, номер тика и состояние ГПСЧ. рождения визитов, номер тика и состояние ГПСЧ.
После этой переделки поток на длинном окне и при штатных параметрах даёт После этой переделки поток на длинном окне и при штатных параметрах даёт
+5 -4
View File
@@ -30,7 +30,7 @@ generator-service -> Kafka topics -> (потребители отдельно)
| `src/clickstream_generator/intensity.py` | расчёт событийного бюджета тика | | `src/clickstream_generator/intensity.py` | расчёт событийного бюджета тика |
| `src/clickstream_generator/runtime.py` | тиковый слой: активные визиты и выпуск созревших событий | | `src/clickstream_generator/runtime.py` | тиковый слой: активные визиты и выпуск созревших событий |
| `src/clickstream_generator/kafka_io.py` | Kafka publisher, история batch, Kafka-state и служебные топики | | `src/clickstream_generator/kafka_io.py` | Kafka publisher, история batch, Kafka-state и служебные топики |
| `src/clickstream_generator/state.py` | сериализуемое состояние генератора v2 | | `src/clickstream_generator/state.py` | сериализуемое состояние генератора v3 |
| `src/clickstream_generator/metrics.py` | Prometheus-метрики | | `src/clickstream_generator/metrics.py` | Prometheus-метрики |
| `src/clickstream_generator/service.py` | основной цикл сервиса | | `src/clickstream_generator/service.py` | основной цикл сервиса |
| `generator.py` | запуск сервиса и совместимый фасад | | `generator.py` | запуск сервиса и совместимый фасад |
@@ -102,7 +102,7 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator
реальному часу запуска процесса. реальному часу запуска процесса.
В режиме `backfill` генератор без сна проходит от `GEN_MODEL_T0` до В режиме `backfill` генератор без сна проходит от `GEN_MODEL_T0` до
`GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v2 `GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v3
на `T_end` в `generator_state` и пишет manifest в compact-topic на `T_end` в `generator_state` и пишет manifest в compact-topic
`generator_startup_history_manifest`. Live-запуск с теми же настройками `generator_startup_history_manifest`. Live-запуск с теми же настройками
использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты. использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты.
@@ -286,11 +286,12 @@ docker compose exec kafka /opt/kafka/bin/kafka-console-consumer.sh \
### Как работает ### Как работает
1. После каждого успешного тика состояние v2 сохраняется в `generator_state`. 1. После каждого успешного тика состояние v3 сохраняется в `generator_state`.
2. При старте генератор читает последнее состояние из топика. 2. При старте генератор читает последнее состояние из топика.
3. Если состояние найдено, сервис восстанавливает номер тика, состояние ГПСЧ, 3. Если состояние найдено, сервис восстанавливает номер тика, состояние ГПСЧ,
популяцию пользователей, накопленный бюджет рождения визитов и активные популяцию пользователей, накопленный бюджет рождения визитов и активные
визиты. визиты. Для активного визита state хранит `base_click_id` — донора браузерной
и source-фактуры из статического сида.
4. Если состояния нет или оно невалидно, генератор начинает с чистого листа. 4. Если состояния нет или оно невалидно, генератор начинает с чистого листа.
Активный визит после простоя до 30 минут продолжается со своими исходными Активный визит после простоя до 30 минут продолжается со своими исходными
@@ -142,13 +142,27 @@ class EventGenerator:
user_profile: dict[str, dict] | None = None, user_profile: dict[str, dict] | None = None,
) -> dict[str, list[dict]]: ) -> dict[str, list[dict]]:
"""Генерирует один визит с сохранением связей.""" """Генерирует один визит с сохранением связей."""
batch, _base_click_id = self.generate_visit_batch(
batch_size,
planned_start_at=planned_start_at,
user_profile=user_profile,
)
return batch
def generate_visit_batch(
self,
batch_size: int,
planned_start_at: datetime | None = None,
user_profile: dict[str, dict] | None = None,
) -> tuple[dict[str, list[dict]], str | None]:
"""Генерирует визит и возвращает click_id донора фактуры."""
if not self.dictionary.browser_events: if not self.dictionary.browser_events:
return { return {
"browser_events": [], "browser_events": [],
"location_events": [], "location_events": [],
"device_events": [], "device_events": [],
"geo_events": [], "geo_events": [],
} }, None
batch = { batch = {
"browser_events": [], "browser_events": [],
@@ -158,7 +172,7 @@ class EventGenerator:
} }
if batch_size <= 0: if batch_size <= 0:
return batch return batch, None
max_visit_events = min(batch_size, self.config.max_session_events) max_visit_events = min(batch_size, self.config.max_session_events)
min_visit_events = min(2, max_visit_events) min_visit_events = min(2, max_visit_events)
@@ -181,8 +195,18 @@ class EventGenerator:
base_browser_events = self.dictionary.browser_by_click_id[base_click_id][:len(visit_path)] base_browser_events = self.dictionary.browser_by_click_id[base_click_id][:len(visit_path)]
else: else:
base_browser = self.rng.choice(self.dictionary.browser_events) base_browser = self.rng.choice(self.dictionary.browser_events)
# Запасная ветка тоже восстановима: state хранит донора, а индекс идёт по кругу.
base_click_id = base_browser["click_id"] base_click_id = base_browser["click_id"]
base_browser_events = [base_browser for _ in range(len(visit_path))] source_events = self.dictionary.browser_by_click_id[base_click_id]
if not all(
event["event_id"] in self.dictionary.location_by_event_id
for event in source_events
):
raise ValueError(f"Fixture base_click_id has incomplete locations: {base_click_id}")
base_browser_events = [
source_events[event_index % len(source_events)]
for event_index in range(len(visit_path))
]
base_device = ( base_device = (
user_profile["device"] user_profile["device"]
@@ -229,4 +253,4 @@ class EventGenerator:
planned_timestamp += timedelta(seconds=self._visit_pause_seconds()) planned_timestamp += timedelta(seconds=self._visit_pause_seconds())
return batch return batch, base_click_id
+30 -12
View File
@@ -35,6 +35,7 @@ class ActiveVisit:
batch: dict[str, list[dict]] batch: dict[str, list[dict]]
timestamps: list[datetime] timestamps: list[datetime]
base_click_id: str
user: UserProfile | None = None user: UserProfile | None = None
next_index: int = 0 next_index: int = 0
@@ -130,7 +131,12 @@ def _datetime_to_state(value: datetime | None) -> str | None:
def _timestamp_to_state_offset(started_at: datetime, timestamp: datetime) -> int: def _timestamp_to_state_offset(started_at: datetime, timestamp: datetime) -> int:
return int((timestamp - started_at).total_seconds() * 1_000_000) delta = timestamp - started_at
return (
delta.days * 86_400_000_000
+ delta.seconds * 1_000_000
+ delta.microseconds
)
def _format_event_timestamp(timestamp: datetime) -> str: def _format_event_timestamp(timestamp: datetime) -> str:
@@ -217,6 +223,7 @@ class TickStreamGenerator:
return { return {
"user_domain_id": visit.user.user_domain_id if visit.user else None, "user_domain_id": visit.user.user_domain_id if visit.user else None,
"click_id": browser_events[0]["click_id"], "click_id": browser_events[0]["click_id"],
"base_click_id": visit.base_click_id,
"next_index": visit.next_index, "next_index": visit.next_index,
"started_at": started_at.isoformat(), "started_at": started_at.isoformat(),
"offsets_us": [ "offsets_us": [
@@ -235,7 +242,7 @@ class TickStreamGenerator:
resume_model_at: datetime | None = None, resume_model_at: datetime | None = None,
restarted_at: datetime | None = None, restarted_at: datetime | None = None,
) -> None: ) -> None:
"""Восстанавливает популяцию и активные визиты из state v2.""" """Восстанавливает популяцию и активные визиты из state."""
users = [ users = [
self._user_from_state(item) self._user_from_state(item)
for item in state.population for item in state.population
@@ -302,6 +309,7 @@ class TickStreamGenerator:
] ]
batch = self._compact_visit_batch( batch = self._compact_visit_batch(
click_id=item["click_id"], click_id=item["click_id"],
base_click_id=item["base_click_id"],
user=user, user=user,
timestamps=timestamps, timestamps=timestamps,
page_url_paths=item["page_url_paths"], page_url_paths=item["page_url_paths"],
@@ -309,6 +317,7 @@ class TickStreamGenerator:
return ActiveVisit( return ActiveVisit(
batch=batch, batch=batch,
timestamps=timestamps, timestamps=timestamps,
base_click_id=item["base_click_id"],
user=user, user=user,
next_index=item["next_index"], next_index=item["next_index"],
) )
@@ -316,24 +325,26 @@ class TickStreamGenerator:
def _compact_visit_batch( def _compact_visit_batch(
self, self,
click_id: str, click_id: str,
base_click_id: str,
user: UserProfile, user: UserProfile,
timestamps: list[datetime], timestamps: list[datetime],
page_url_paths: list[str], page_url_paths: list[str],
) -> dict[str, list[dict]]: ) -> dict[str, list[dict]]:
batch = _empty_batch() batch = _empty_batch()
browser_templates = self.generator.dictionary.browser_by_click_id.get( if base_click_id not in self.generator.dictionary.browser_by_click_id:
user.seed_click_id, raise ValueError(f"Unknown fixture base_click_id: {base_click_id}")
self.generator.dictionary.browser_events, browser_templates = self.generator.dictionary.browser_by_click_id[base_click_id]
) if not browser_templates:
raise ValueError(f"Fixture base_click_id has no browser events: {base_click_id}")
for event_index, (timestamp, page_url_path) in enumerate( for event_index, (timestamp, page_url_path) in enumerate(
zip(timestamps, page_url_paths) zip(timestamps, page_url_paths)
): ):
browser_template = browser_templates[event_index % len(browser_templates)] browser_template = browser_templates[event_index % len(browser_templates)]
location_template = self.generator.dictionary.location_by_event_id.get( source_event_id = browser_template["event_id"]
browser_template["event_id"], if source_event_id not in self.generator.dictionary.location_by_event_id:
self.generator.dictionary.location_events[0], raise ValueError(f"Unknown fixture location event_id: {source_event_id}")
) location_template = self.generator.dictionary.location_by_event_id[source_event_id]
event_id = _stable_event_id(click_id, event_index) event_id = _stable_event_id(click_id, event_index)
batch["browser_events"].append( batch["browser_events"].append(
{ {
@@ -428,7 +439,7 @@ class TickStreamGenerator:
self._pending_visit_births = 0.0 self._pending_visit_births = 0.0
return return
visit_batch = self.generator.generate_batch( visit_batch, base_click_id = self.generator.generate_visit_batch(
self.generator.config.max_session_events, self.generator.config.max_session_events,
planned_start_at=tick_time, planned_start_at=tick_time,
user_profile={"device": user.device, "geo": user.geo}, user_profile={"device": user.device, "geo": user.geo},
@@ -439,11 +450,18 @@ class TickStreamGenerator:
] ]
if not timestamps: if not timestamps:
break break
if base_click_id is None:
raise ValueError("Generated active visit is missing fixture base_click_id")
click_id = visit_batch["browser_events"][0]["click_id"] click_id = visit_batch["browser_events"][0]["click_id"]
self.population.start_visit(user, click_id) self.population.start_visit(user, click_id)
self.active_visits.append( self.active_visits.append(
ActiveVisit(batch=visit_batch, timestamps=timestamps, user=user) ActiveVisit(
batch=visit_batch,
timestamps=timestamps,
base_click_id=base_click_id,
user=user,
)
) )
self._pending_visit_births -= 1.0 self._pending_visit_births -= 1.0
@@ -205,7 +205,7 @@ class GeneratorService:
state, state,
wall_now_utc: datetime, wall_now_utc: datetime,
) -> datetime: ) -> datetime:
"""Считает модельную точку live-восстановления по state v2.""" """Считает модельную точку live-восстановления по state."""
wall_now_utc = self._as_aware_utc(wall_now_utc) wall_now_utc = self._as_aware_utc(wall_now_utc)
wall_saved_at = self._as_aware_utc(state.wall_timestamp) wall_saved_at = self._as_aware_utc(state.wall_timestamp)
idle_seconds = max(0.0, (wall_now_utc - wall_saved_at).total_seconds()) idle_seconds = max(0.0, (wall_now_utc - wall_saved_at).total_seconds())
+7 -3
View File
@@ -10,7 +10,7 @@ from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
logger = logging.getLogger("generator") logger = logging.getLogger("generator")
STATE_VERSION = "2.0" STATE_VERSION = "3.0"
def _nested_list_to_tuple(obj): def _nested_list_to_tuple(obj):
@@ -74,7 +74,7 @@ def _validate_resume_fields(data: dict) -> None:
raise ValueError("gen_seed must be an integer or null") raise ValueError("gen_seed must be an integer or null")
def _validate_v2_payload(data: dict) -> None: def _validate_v3_payload(data: dict) -> None:
_validate_resume_fields(data) _validate_resume_fields(data)
population = data.get("population") population = data.get("population")
active_visits = data.get("active_visits") active_visits = data.get("active_visits")
@@ -123,6 +123,7 @@ def _validate_v2_payload(data: dict) -> None:
( (
"user_domain_id", "user_domain_id",
"click_id", "click_id",
"base_click_id",
"next_index", "next_index",
"started_at", "started_at",
"offsets_us", "offsets_us",
@@ -132,6 +133,7 @@ def _validate_v2_payload(data: dict) -> None:
) )
user_domain_id = visit["user_domain_id"] user_domain_id = visit["user_domain_id"]
click_id = visit["click_id"] click_id = visit["click_id"]
base_click_id = visit["base_click_id"]
offsets = visit["offsets_us"] offsets = visit["offsets_us"]
page_url_paths = visit["page_url_paths"] page_url_paths = visit["page_url_paths"]
next_index = visit["next_index"] next_index = visit["next_index"]
@@ -139,6 +141,8 @@ def _validate_v2_payload(data: dict) -> None:
raise ValueError(f"active_visits[{index}].user_domain_id is unknown") raise ValueError(f"active_visits[{index}].user_domain_id is unknown")
if not isinstance(click_id, str) or not click_id: if not isinstance(click_id, str) or not click_id:
raise ValueError(f"active_visits[{index}].click_id must be a string") raise ValueError(f"active_visits[{index}].click_id must be a string")
if not isinstance(base_click_id, str) or not base_click_id:
raise ValueError(f"active_visits[{index}].base_click_id must be a string")
if not isinstance(offsets, list) or not offsets: if not isinstance(offsets, list) or not offsets:
raise ValueError(f"active_visits[{index}].offsets_us must be a non-empty list") raise ValueError(f"active_visits[{index}].offsets_us must be a non-empty list")
if not all(isinstance(offset, int) and offset >= 0 for offset in offsets): if not all(isinstance(offset, int) and offset >= 0 for offset in offsets):
@@ -235,7 +239,7 @@ class GeneratorState:
version, version,
) )
raise ValueError(f"unsupported state version: {version}") raise ValueError(f"unsupported state version: {version}")
_validate_v2_payload(data) _validate_v3_payload(data)
model_timestamp = _parse_aware_utc( model_timestamp = _parse_aware_utc(
data["model_timestamp"], data["model_timestamp"],
"model_timestamp", "model_timestamp",
+116 -1
View File
@@ -18,6 +18,7 @@ from generator import (
generate_tick_batch, generate_tick_batch,
hour_factor, hour_factor,
) )
from clickstream_generator.runtime import _timestamp_to_state_offset
ALLOWED_PAGE_PATHS = { ALLOWED_PAGE_PATHS = {
@@ -37,6 +38,24 @@ def _parse_event_timestamps(batch):
] ]
def _without_fields(batch, excluded_fields):
return {
topic: [
{
key: value
for key, value in event.items()
if key not in excluded_fields
}
for event in events
]
for topic, events in batch.items()
}
def _without_event_ids(batch):
return _without_fields(batch, {"event_id"})
def _page_path_visits(generator, visits_count: int): def _page_path_visits(generator, visits_count: int):
return [ return [
[ [
@@ -297,7 +316,7 @@ class TestEventGeneration:
def test_tick_stream_state_stays_compact_at_active_session_limit( def test_tick_stream_state_stays_compact_at_active_session_limit(
self, event_dictionary, base_config self, event_dictionary, base_config
): ):
"""State v2 не хранит полные события активных визитов.""" """State не хранит полные события активных визитов."""
config = replace( config = replace(
base_config, base_config,
max_session_events=30, max_session_events=30,
@@ -386,6 +405,102 @@ class TestEventGeneration:
for timestamp in resumed_timestamps for timestamp in resumed_timestamps
) )
def test_tick_stream_restored_visit_keeps_per_event_fixture(
self, event_dictionary, base_config
):
"""Восстановленный визит продолжает ту же фактуру событий."""
config = replace(base_config, max_session_events=8)
tick_at = datetime(2026, 6, 11, 12, 0)
resume_at = tick_at + timedelta(hours=2)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
stream.generate_tick(event_budget=10, tick_started_at=tick_at)
state = stream.to_state(
tick=1,
rng_state=generator.rng.getstate(),
last_batch_id="batch-1",
last_timestamp=tick_at,
)
restored_state = type(state).from_dict(json.loads(json.dumps(state.to_dict())))
restored_generator = EventGenerator(event_dictionary, config)
restored_stream = TickStreamGenerator(restored_generator)
restored_stream.restore_state(restored_state, resume_model_at=tick_at)
uninterrupted = stream.generate_tick(event_budget=0, tick_started_at=resume_at)
restored = restored_stream.generate_tick(event_budget=0, tick_started_at=resume_at)
assert restored["browser_events"]
assert _without_event_ids(restored) == _without_event_ids(uninterrupted)
def test_tick_stream_restored_visit_keeps_cyclic_fallback_fixture(
self, event_dictionary, base_config
):
"""Короткий донор фактуры после восстановления повторяется так же."""
config = replace(base_config, max_session_events=30)
tick_at = datetime(2026, 6, 11, 12, 0)
resume_at = tick_at + timedelta(hours=12)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
stream.generate_tick(event_budget=100, tick_started_at=tick_at)
state = stream.to_state(
tick=1,
rng_state=generator.rng.getstate(),
last_batch_id="batch-1",
last_timestamp=tick_at,
)
assert any(
len(event_dictionary.browser_by_click_id[visit["base_click_id"]])
< len(visit["offsets_us"])
for visit in state.active_visits
)
restored_state = type(state).from_dict(json.loads(json.dumps(state.to_dict())))
restored_generator = EventGenerator(event_dictionary, config)
restored_stream = TickStreamGenerator(restored_generator)
restored_stream.restore_state(restored_state, resume_model_at=tick_at)
uninterrupted = stream.generate_tick(event_budget=0, tick_started_at=resume_at)
restored = restored_stream.generate_tick(event_budget=0, tick_started_at=resume_at)
assert restored["browser_events"]
assert _without_event_ids(restored) == _without_event_ids(uninterrupted)
def test_state_offset_preserves_microseconds_without_float_rounding(self):
"""Смещение state не теряет микросекунду на float-округлении."""
started_at = datetime(2026, 6, 11, 12, 0, 0)
timestamp = datetime(2026, 6, 11, 12, 8, 34, 130779)
assert _timestamp_to_state_offset(started_at, timestamp) == 514_130_779
def test_tick_stream_restore_rejects_unknown_fixture_donor(
self, event_dictionary, base_config
):
"""Неизвестный донор фактуры в state даёт ошибку восстановления."""
config = replace(base_config, max_session_events=5)
tick_at = datetime(2026, 6, 11, 12, 0)
generator = EventGenerator(event_dictionary, config)
stream = TickStreamGenerator(generator)
stream.generate_tick(event_budget=10, tick_started_at=tick_at)
state = stream.to_state(
tick=1,
rng_state=generator.rng.getstate(),
last_batch_id="batch-1",
last_timestamp=tick_at,
)
state_data = json.loads(json.dumps(state.to_dict()))
state_data["active_visits"][0]["base_click_id"] = "missing-click-id"
restored_state = type(state).from_dict(state_data)
restored_generator = EventGenerator(event_dictionary, config)
restored_stream = TickStreamGenerator(restored_generator)
with pytest.raises(ValueError, match="Unknown fixture base_click_id"):
restored_stream.restore_state(restored_state, resume_model_at=tick_at)
def test_tick_stream_restored_after_long_idle_closes_overdue_visit_without_replay( def test_tick_stream_restored_after_long_idle_closes_overdue_visit_without_replay(
self, event_dictionary, base_config self, event_dictionary, base_config
): ):
+10 -10
View File
@@ -926,11 +926,11 @@ class TestGeneratorServiceBackfill:
} }
class TestGeneratorServiceStateV2: class TestGeneratorServiceState:
"""Тесты подключения state v2 к сервисному запуску.""" """Тесты подключения state к сервисному запуску."""
def test_start_restores_tick_stream_state_v2(self, base_config, event_dictionary): def test_start_restores_tick_stream_state(self, base_config, event_dictionary):
"""Сервис восстанавливает популяцию и активные визиты из state v2.""" """Сервис восстанавливает популяцию и активные визиты из state."""
source_generator = EventGenerator(event_dictionary, base_config) source_generator = EventGenerator(event_dictionary, base_config)
source_stream = TickStreamGenerator(source_generator) source_stream = TickStreamGenerator(source_generator)
tick_at = datetime.now(timezone.utc).replace(tzinfo=None) tick_at = datetime.now(timezone.utc).replace(tzinfo=None)
@@ -1073,8 +1073,8 @@ class TestGeneratorServiceStateV2:
assert service._model_time == model_t_end assert service._model_time == model_t_end
assert service.stream.active_visit_count == source_stream.active_visit_count assert service.stream.active_visit_count == source_stream.active_visit_count
def test_save_state_writes_tick_stream_state_v2(self, base_config): def test_save_state_writes_tick_stream_state(self, base_config):
"""Сервис сохраняет v2-снимок тикового слоя.""" """Сервис сохраняет снимок тикового слоя."""
service = GeneratorService(base_config) service = GeneratorService(base_config)
service.state_manager = MagicMock() service.state_manager = MagicMock()
tick_at = datetime.now(timezone.utc).replace(tzinfo=None) tick_at = datetime.now(timezone.utc).replace(tzinfo=None)
@@ -1085,7 +1085,7 @@ class TestGeneratorServiceStateV2:
service._save_state("batch-1") service._save_state("batch-1")
saved_state = service.state_manager.save.call_args.args[0] saved_state = service.state_manager.save.call_args.args[0]
assert saved_state.version == "2.0" assert saved_state.version == "3.0"
assert saved_state.model_timestamp == tick_at.replace(tzinfo=timezone.utc) assert saved_state.model_timestamp == tick_at.replace(tzinfo=timezone.utc)
assert saved_state.wall_timestamp.tzinfo is not None assert saved_state.wall_timestamp.tzinfo is not None
assert saved_state.model_time_speed == base_config.model_time_speed assert saved_state.model_time_speed == base_config.model_time_speed
@@ -1162,13 +1162,13 @@ class TestGeneratorServiceStateV2:
state_manager.load.assert_not_called() state_manager.load.assert_not_called()
assert service._tick == 0 assert service._tick == 0
def test_invalid_restored_v2_state_starts_fresh(self, base_config, caplog): def test_invalid_restored_state_starts_fresh(self, base_config, caplog):
"""Сервис не падает, если v2 state ссылается на неизвестный профиль.""" """Сервис не падает, если state ссылается на неизвестный профиль."""
state_manager = MagicMock() state_manager = MagicMock()
state_manager.load.return_value = GeneratorState( state_manager.load.return_value = GeneratorState(
tick=9, tick=9,
rng_state=random.Random(42).getstate(), rng_state=random.Random(42).getstate(),
last_batch_id="bad-v2", last_batch_id="bad-state",
last_timestamp=datetime.now(timezone.utc), last_timestamp=datetime.now(timezone.utc),
model_timestamp=base_config.model_t0, model_timestamp=base_config.model_t0,
wall_timestamp=datetime.now(timezone.utc), wall_timestamp=datetime.now(timezone.utc),
+67 -38
View File
@@ -18,8 +18,8 @@ def _make_valid_rng_state(seed: int = 42):
return rng.getstate() return rng.getstate()
def _make_valid_v2_state_data() -> dict: def _make_valid_state_data() -> dict:
"""Создаёт минимальный валидный state v2 для тестов загрузки.""" """Создаёт минимальный валидный state для тестов загрузки."""
return { return {
"tick": 42, "tick": 42,
"rng_state": list(_make_valid_rng_state(42)), "rng_state": list(_make_valid_rng_state(42)),
@@ -31,7 +31,7 @@ def _make_valid_v2_state_data() -> dict:
"model_timezone": "UTC", "model_timezone": "UTC",
"model_t0": "2026-01-01T00:00:00+00:00", "model_t0": "2026-01-01T00:00:00+00:00",
"gen_seed": 42, "gen_seed": 42,
"version": "2.0", "version": "3.0",
"population": [ "population": [
{ {
"user_domain_id": "user-1", "user_domain_id": "user-1",
@@ -44,6 +44,7 @@ def _make_valid_v2_state_data() -> dict:
{ {
"user_domain_id": "user-1", "user_domain_id": "user-1",
"click_id": "visit-1", "click_id": "visit-1",
"base_click_id": "seed-1",
"next_index": 1, "next_index": 1,
"started_at": "2026-06-11T12:00:00", "started_at": "2026-06-11T12:00:00",
"offsets_us": [0, 60_000_000], "offsets_us": [0, 60_000_000],
@@ -66,7 +67,7 @@ def _minimal_population() -> list[dict]:
def _with_resume_fields(data: dict) -> dict: def _with_resume_fields(data: dict) -> dict:
"""Добавляет обязательные поля state v2, не связанные с проверяемой ошибкой.""" """Добавляет обязательные поля state, не связанные с проверяемой ошибкой."""
return { return {
**data, **data,
"model_timestamp": "2026-01-01T10:00:00+00:00", "model_timestamp": "2026-01-01T10:00:00+00:00",
@@ -97,7 +98,7 @@ class TestGeneratorState:
model_timezone="UTC", model_timezone="UTC",
model_t0=now, model_t0=now,
gen_seed=42, gen_seed=42,
version="2.0", version="3.0",
) )
assert state.tick == 42 assert state.tick == 42
@@ -110,7 +111,7 @@ class TestGeneratorState:
assert state.model_timezone == "UTC" assert state.model_timezone == "UTC"
assert state.model_t0 == now assert state.model_t0 == now
assert state.gen_seed == 42 assert state.gen_seed == 42
assert state.version == "2.0" assert state.version == "3.0"
def test_default_version(self): def test_default_version(self):
"""Новые состояния по умолчанию пишутся в версии 2.""" """Новые состояния по умолчанию пишутся в версии 2."""
@@ -124,7 +125,7 @@ class TestGeneratorState:
last_timestamp=now, last_timestamp=now,
) )
assert state.version == "2.0" assert state.version == "3.0"
def test_to_dict_serialization(self): def test_to_dict_serialization(self):
"""Сериализация в словарь (JSON-safe, без pickle).""" """Сериализация в словарь (JSON-safe, без pickle)."""
@@ -150,7 +151,7 @@ class TestGeneratorState:
assert data["model_timezone"] == "UTC" assert data["model_timezone"] == "UTC"
assert data["model_t0"] == now.isoformat() assert data["model_t0"] == now.isoformat()
assert data["gen_seed"] is None assert data["gen_seed"] is None
assert data["version"] == "2.0" assert data["version"] == "3.0"
# Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle) # Проверяем что rng_state сериализован как tuple (JSON-safe, без pickle)
assert "rng_state" in data assert "rng_state" in data
@@ -225,8 +226,8 @@ class TestGeneratorState:
assert next_values == values_after assert next_values == values_after
def test_version_2_roundtrip_keeps_population_and_active_visits(self): def test_state_roundtrip_keeps_population_and_active_visits(self):
"""State v2 хранит популяцию и активные визиты в JSON.""" """State хранит популяцию и активные визиты в JSON."""
rng_state = _make_valid_rng_state(42) rng_state = _make_valid_rng_state(42)
state = GeneratorState( state = GeneratorState(
tick=7, tick=7,
@@ -239,7 +240,7 @@ class TestGeneratorState:
model_timezone="Europe/Moscow", model_timezone="Europe/Moscow",
model_t0=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), model_t0=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc),
gen_seed=42, gen_seed=42,
version="2.0", version="3.0",
population=[ population=[
{ {
"user_domain_id": "user-1", "user_domain_id": "user-1",
@@ -252,6 +253,7 @@ class TestGeneratorState:
{ {
"user_domain_id": "user-1", "user_domain_id": "user-1",
"click_id": "visit-1", "click_id": "visit-1",
"base_click_id": "seed-1",
"next_index": 1, "next_index": 1,
"started_at": "2026-06-11T12:00:00", "started_at": "2026-06-11T12:00:00",
"offsets_us": [0, 60_000_000], "offsets_us": [0, 60_000_000],
@@ -263,7 +265,7 @@ class TestGeneratorState:
restored = GeneratorState.from_dict(json.loads(json.dumps(state.to_dict()))) restored = GeneratorState.from_dict(json.loads(json.dumps(state.to_dict())))
assert restored.version == "2.0" assert restored.version == "3.0"
assert restored.tick == state.tick assert restored.tick == state.tick
assert restored.rng_state == rng_state assert restored.rng_state == rng_state
assert restored.population == state.population assert restored.population == state.population
@@ -281,7 +283,7 @@ class TestGeneratorStateValidation:
"""Тесты валидации состояния и graceful degradation.""" """Тесты валидации состояния и graceful degradation."""
def test_from_dict_missing_version_raises(self): def test_from_dict_missing_version_raises(self):
"""from_dict выбрасывает исключение при отсутствии версии v2.""" """from_dict выбрасывает исключение при отсутствии версии."""
data = { data = {
"tick": 42, "tick": 42,
"last_batch_id": "test", "last_batch_id": "test",
@@ -298,7 +300,7 @@ class TestGeneratorStateValidation:
"rng_state": "not_a_tuple", "rng_state": "not_a_tuple",
"last_batch_id": "test", "last_batch_id": "test",
"last_timestamp": "2024-01-01T00:00:00+00:00", "last_timestamp": "2024-01-01T00:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
}) })
@@ -313,7 +315,7 @@ class TestGeneratorStateValidation:
"rng_state": [1], # Слишком короткий "rng_state": [1], # Слишком короткий
"last_batch_id": "test", "last_batch_id": "test",
"last_timestamp": "2024-01-01T00:00:00+00:00", "last_timestamp": "2024-01-01T00:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
}) })
@@ -328,7 +330,7 @@ class TestGeneratorStateValidation:
"rng_state": [999, [1, 2, 3], None], # Невалидный state "rng_state": [999, [1, 2, 3], None], # Невалидный state
"last_batch_id": "test", "last_batch_id": "test",
"last_timestamp": "2024-01-01T00:00:00+00:00", "last_timestamp": "2024-01-01T00:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
}) })
@@ -343,7 +345,7 @@ class TestGeneratorStateValidation:
"rng_state": "invalid", "rng_state": "invalid",
"last_batch_id": "test", "last_batch_id": "test",
"last_timestamp": "2024-01-01T00:00:00+00:00", "last_timestamp": "2024-01-01T00:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
}) })
@@ -365,7 +367,7 @@ class TestGeneratorStateValidation:
"model_timezone": "UTC", "model_timezone": "UTC",
"model_t0": "2026-01-01T00:00:00+00:00", "model_t0": "2026-01-01T00:00:00+00:00",
"gen_seed": 42, "gen_seed": 42,
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
} }
@@ -378,7 +380,7 @@ class TestGeneratorStateValidation:
def test_from_dict_safe_returns_none_on_invalid_gen_seed(self): def test_from_dict_safe_returns_none_on_invalid_gen_seed(self):
"""gen_seed в JSON state должен быть числом или null.""" """gen_seed в JSON state должен быть числом или null."""
data = _make_valid_v2_state_data() data = _make_valid_state_data()
data["gen_seed"] = "42" data["gen_seed"] = "42"
result = GeneratorState.from_dict_safe(data) result = GeneratorState.from_dict_safe(data)
@@ -387,7 +389,7 @@ class TestGeneratorStateValidation:
def test_from_dict_safe_returns_none_on_bool_model_time_speed(self): def test_from_dict_safe_returns_none_on_bool_model_time_speed(self):
"""model_time_speed не принимает bool как числовую скорость.""" """model_time_speed не принимает bool как числовую скорость."""
data = _make_valid_v2_state_data() data = _make_valid_state_data()
data["model_time_speed"] = True data["model_time_speed"] = True
result = GeneratorState.from_dict_safe(data) result = GeneratorState.from_dict_safe(data)
@@ -395,14 +397,14 @@ class TestGeneratorStateValidation:
assert result is None assert result is None
def test_from_dict_safe_returns_none_without_model_resume_fields(self): def test_from_dict_safe_returns_none_without_model_resume_fields(self):
"""State v2 без связки модельного и настенного времени несовместим.""" """State без связки модельного и настенного времени несовместим."""
rng = random.Random(42) rng = random.Random(42)
data = { data = {
"tick": 42, "tick": 42,
"rng_state": list(rng.getstate()), "rng_state": list(rng.getstate()),
"last_batch_id": "test", "last_batch_id": "test",
"last_timestamp": "2024-01-01T00:00:00+00:00", "last_timestamp": "2024-01-01T00:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": _minimal_population(), "population": _minimal_population(),
"active_visits": [], "active_visits": [],
} }
@@ -497,7 +499,7 @@ class TestKafkaStateManager:
rng_state=_make_valid_rng_state(100), rng_state=_make_valid_rng_state(100),
last_batch_id="xyz789", last_batch_id="xyz789",
last_timestamp=now, last_timestamp=now,
version="2.0", version="3.0",
population=_minimal_population(), population=_minimal_population(),
) )
@@ -540,8 +542,8 @@ class TestKafkaStateManager:
# Должно вернуть None из-за невалидного state # Должно вернуть None из-за невалидного state
assert result is None assert result is None
def test_load_invalid_v2_nested_state_returns_none(self, caplog): def test_load_invalid_nested_state_returns_none(self, caplog):
"""Битое state v2 с валидным rng_state даёт чистый старт.""" """Битое state с валидным rng_state даёт чистый старт."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
@@ -555,12 +557,13 @@ class TestKafkaStateManager:
"rng_state": list(_make_valid_rng_state(42)), "rng_state": list(_make_valid_rng_state(42)),
"last_batch_id": "bad-v2", "last_batch_id": "bad-v2",
"last_timestamp": "2026-06-11T12:00:00+00:00", "last_timestamp": "2026-06-11T12:00:00+00:00",
"version": "2.0", "version": "3.0",
"population": [{"user_domain_id": "user-1"}], "population": [{"user_domain_id": "user-1"}],
"active_visits": [ "active_visits": [
{ {
"user_domain_id": "user-1", "user_domain_id": "user-1",
"click_id": "visit-1", "click_id": "visit-1",
"base_click_id": "seed-1",
"next_index": 1, "next_index": 1,
"started_at": "2026-06-11T12:00:00", "started_at": "2026-06-11T12:00:00",
"offsets_us": [0], "offsets_us": [0],
@@ -579,15 +582,15 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_empty_population_v2_returns_none(self, caplog): def test_load_empty_population_returns_none(self, caplog):
"""Пустая популяция в state v2 не восстанавливается.""" """Пустая популяция в state не восстанавливается."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_v2_state_data() bad_state = _make_valid_state_data()
bad_state["population"] = [] bad_state["population"] = []
bad_state["active_visits"] = [] bad_state["active_visits"] = []
@@ -606,15 +609,15 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_bad_pending_births_v2_returns_none(self, caplog): def test_load_bad_pending_births_returns_none(self, caplog):
"""Нечисловой pending_visit_births в state v2 не восстанавливается.""" """Нечисловой pending_visit_births в state не восстанавливается."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_v2_state_data() bad_state = _make_valid_state_data()
bad_state["pending_visit_births"] = "bad" bad_state["pending_visit_births"] = "bad"
mock_message = MagicMock() mock_message = MagicMock()
@@ -632,7 +635,33 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_active_visit_with_unknown_user_v2_returns_none(self, caplog): def test_load_active_visit_without_base_click_id_returns_none(self, caplog):
"""Активный визит без донора фактуры не восстанавливается."""
with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class:
mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_state_data()
del bad_state["active_visits"][0]["base_click_id"]
mock_message = MagicMock()
mock_message.key = b"default"
mock_message.value = bad_state
mock_consumer = MagicMock()
mock_consumer.__iter__ = MagicMock(return_value=iter([mock_message]))
mock_consumer_class.return_value = mock_consumer
manager = KafkaStateManager("kafka:29092")
with caplog.at_level(logging.WARNING, logger="generator"):
result = manager.load()
assert result is None
assert "base_click_id" in caplog.text
def test_load_active_visit_with_unknown_user_returns_none(self, caplog):
"""Активный визит должен ссылаться на пользователя из популяции.""" """Активный визит должен ссылаться на пользователя из популяции."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
@@ -640,7 +669,7 @@ class TestKafkaStateManager:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_v2_state_data() bad_state = _make_valid_state_data()
bad_state["active_visits"][0]["user_domain_id"] = "missing-user" bad_state["active_visits"][0]["user_domain_id"] = "missing-user"
mock_message = MagicMock() mock_message = MagicMock()
@@ -658,7 +687,7 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_active_visit_with_conflicting_click_id_v2_returns_none(self, caplog): def test_load_active_visit_with_conflicting_click_id_returns_none(self, caplog):
"""active_click_id пользователя не должен противоречить визиту.""" """active_click_id пользователя не должен противоречить визиту."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
@@ -666,7 +695,7 @@ class TestKafkaStateManager:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_v2_state_data() bad_state = _make_valid_state_data()
bad_state["population"][0]["active_click_id"] = "other-visit" bad_state["population"][0]["active_click_id"] = "other-visit"
mock_message = MagicMock() mock_message = MagicMock()
@@ -684,7 +713,7 @@ class TestKafkaStateManager:
assert result is None assert result is None
assert "Invalid state" in caplog.text assert "Invalid state" in caplog.text
def test_load_population_ghost_active_click_id_v2_returns_none(self, caplog): def test_load_population_ghost_active_click_id_returns_none(self, caplog):
"""active_click_id пользователя должен иметь соответствующий активный визит.""" """active_click_id пользователя должен иметь соответствующий активный визит."""
with patch("generator._import_kafka") as mock_import, \ with patch("generator._import_kafka") as mock_import, \
patch("kafka.KafkaConsumer") as mock_consumer_class: patch("kafka.KafkaConsumer") as mock_consumer_class:
@@ -692,7 +721,7 @@ class TestKafkaStateManager:
mock_producer_class = MagicMock() mock_producer_class = MagicMock()
mock_import.return_value = (mock_producer_class, None) mock_import.return_value = (mock_producer_class, None)
bad_state = _make_valid_v2_state_data() bad_state = _make_valid_state_data()
bad_state["population"][0]["active_click_id"] = "ghost" bad_state["population"][0]["active_click_id"] = "ghost"
bad_state["active_visits"] = [] bad_state["active_visits"] = []
+184
View File
@@ -12,6 +12,8 @@ CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}"
GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}" GEN_MODEL_T0="${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00}"
GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}" GEN_MODEL_T_END="${GEN_MODEL_T_END:-2026-01-01T06:00:00+00:00}"
REQUIRE_SUPERSET="${REQUIRE_SUPERSET:-1}" REQUIRE_SUPERSET="${REQUIRE_SUPERSET:-1}"
CHECK_LIVE_SEAM="${CHECK_LIVE_SEAM:-auto}"
GEN_LIVE_CHECK_MINUTES="${GEN_LIVE_CHECK_MINUTES:-10}"
fail() { fail() {
echo "Ошибка: $*" >&2 echo "Ошибка: $*" >&2
@@ -46,6 +48,13 @@ pg_query() {
CH_MODEL_T0="$(clickhouse_datetime_literal "${GEN_MODEL_T0}")" CH_MODEL_T0="$(clickhouse_datetime_literal "${GEN_MODEL_T0}")"
CH_MODEL_T_END="$(clickhouse_datetime_literal "${GEN_MODEL_T_END}")" CH_MODEL_T_END="$(clickhouse_datetime_literal "${GEN_MODEL_T_END}")"
[[ "${GEN_LIVE_CHECK_MINUTES}" =~ ^[0-9]+$ ]] \
|| fail "GEN_LIVE_CHECK_MINUTES должен быть целым числом минут"
case "${CHECK_LIVE_SEAM}" in
0|1|auto) ;;
*) fail "CHECK_LIVE_SEAM должен быть 0, 1 или auto" ;;
esac
echo "=== Проверка ClickHouse: данные генерации в DM ===" echo "=== Проверка ClickHouse: данные генерации в DM ==="
stats_query=" stats_query="
@@ -261,6 +270,181 @@ echo "confirmation=${contains_confirmation}"
echo "monotonic_ok=${contains_monotonic_ok}" echo "monotonic_ok=${contains_monotonic_ok}"
echo "confirmation_share=${contains_confirmation_share}" echo "confirmation_share=${contains_confirmation_share}"
should_check_live_seam="${CHECK_LIVE_SEAM}"
if [[ "${should_check_live_seam}" == "auto" ]]; then
live_rows_query="
WITH
toDateTime64('${CH_MODEL_T_END}', 6) AS t_end,
t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end
SELECT count()
FROM dds.event
WHERE event_ts >= t_end AND event_ts < t_live_end
FORMAT TabSeparated"
live_rows="$(ch_query "${live_rows_query}")"
[[ "${live_rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать live-строки DDS после GEN_MODEL_T_END"
if (( live_rows > 0 )); then
should_check_live_seam="1"
else
should_check_live_seam="0"
fi
fi
if [[ "${should_check_live_seam}" == "1" ]]; then
echo ""
echo "=== Проверка ClickHouse: стык backfill/live в DDS ==="
seam_duplicates_query="
WITH
toDateTime64('${CH_MODEL_T0}', 6) AS t0,
toDateTime64('${CH_MODEL_T_END}', 6) AS t_end,
t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end
SELECT
count() AS events,
uniqExact(event_id) AS unique_events,
events - unique_events AS duplicate_events
FROM dds.event
WHERE event_ts >= t0 AND event_ts < t_live_end
FORMAT TabSeparated"
seam_duplicates="$(ch_query "${seam_duplicates_query}")"
IFS=$'\t' read -r seam_events seam_unique_events seam_duplicate_events <<< "${seam_duplicates}"
[[ "${seam_events}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать события DDS на стыке"
(( seam_events > 0 )) || fail "нет событий DDS в диапазоне проверки стыка"
(( seam_duplicate_events == 0 )) || fail "на стыке найдены дубли event_id: ${seam_duplicate_events}"
echo "events=${seam_events}"
echo "unique_events=${seam_unique_events}"
echo "duplicate_events=${seam_duplicate_events}"
seam_context_query="
WITH
toDateTime64('${CH_MODEL_T_END}', 6) AS t_end,
t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end,
crossing AS (
SELECT
e.click_id,
min(e.event_ts) AS first_ts,
max(e.event_ts) AS last_ts,
groupUniqArray(coalesce(toString(e.browser_name), '__NULL__')) AS browser_names,
groupUniqArray(coalesce(toString(e.browser_language), '__NULL__')) AS browser_languages,
groupUniqArray(coalesce(toString(e.browser_user_agent), '__NULL__')) AS browser_user_agents,
groupUniqArray(coalesce(toString(e.referer_url), '__NULL__')) AS referer_urls,
groupUniqArray(coalesce(toString(e.referer_medium), '__NULL__')) AS referer_mediums,
groupUniqArray(coalesce(toString(e.utm_medium), '__NULL__')) AS utm_mediums,
groupUniqArray(coalesce(toString(e.utm_source), '__NULL__')) AS utm_sources,
groupUniqArray(coalesce(toString(e.utm_content), '__NULL__')) AS utm_contents,
groupUniqArray(coalesce(toString(e.utm_campaign), '__NULL__')) AS utm_campaigns,
groupUniqArray(coalesce(toString(c.device_type), '__NULL__')) AS device_types,
groupUniqArray(coalesce(toString(c.os_name), '__NULL__')) AS os_names,
groupUniqArray(coalesce(toString(c.geo_country), '__NULL__')) AS geo_countries
FROM dds.event AS e
LEFT JOIN dds.click AS c ON c.click_id = e.click_id
WHERE e.event_ts >= t_end - INTERVAL 30 MINUTE
AND e.event_ts < t_live_end
GROUP BY e.click_id
HAVING first_ts < t_end AND last_ts >= t_end
)
SELECT
count() AS crossing_visits,
countIf(
length(browser_names) = 1
AND length(browser_languages) = 1
AND length(browser_user_agents) = 1
AND length(referer_urls) = 1
AND length(referer_mediums) = 1
AND length(utm_mediums) = 1
AND length(utm_sources) = 1
AND length(utm_contents) = 1
AND length(utm_campaigns) = 1
) AS per_event_homogeneous_visits,
countIf(
length(device_types) = 1
AND length(os_names) = 1
AND length(geo_countries) = 1
) AS click_context_homogeneous_visits
FROM crossing
FORMAT TabSeparated"
seam_context="$(ch_query "${seam_context_query}")"
IFS=$'\t' read -r crossing_visits per_event_homogeneous_visits click_context_homogeneous_visits <<< "${seam_context}"
[[ "${crossing_visits}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать переходящие визиты"
(( crossing_visits > 0 )) || fail "нет визитов, переходящих через GEN_MODEL_T_END"
(( per_event_homogeneous_visits == crossing_visits )) \
|| fail "per-event фактура меняется на стыке: ${per_event_homogeneous_visits}/${crossing_visits}"
(( click_context_homogeneous_visits == crossing_visits )) \
|| fail "device/os/geo меняются на стыке: ${click_context_homogeneous_visits}/${crossing_visits}"
echo "crossing_visits=${crossing_visits}"
echo "per_event_homogeneous_visits=${per_event_homogeneous_visits}"
echo "click_context_homogeneous_visits=${click_context_homogeneous_visits}"
ods_context_query="
WITH
toDateTime64('${CH_MODEL_T_END}', 6) AS t_end,
t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end,
crossing AS (
SELECT
click_id
FROM dds.event
WHERE event_ts >= t_end - INTERVAL 30 MINUTE
AND event_ts < t_live_end
GROUP BY click_id
HAVING min(event_ts) < t_end AND max(event_ts) >= t_end
),
device_conflicts AS (
SELECT
d.click_id
FROM ods.device_by_click AS d
INNER JOIN crossing AS c ON c.click_id = d.click_id
WHERE d.click_id IS NOT NULL
GROUP BY d.click_id
HAVING
uniqExact(coalesce(toString(device_type), '__NULL__')) > 1
OR uniqExact(coalesce(toString(device_is_mobile), '__NULL__')) > 1
OR uniqExact(coalesce(toString(os_name), '__NULL__')) > 1
OR uniqExact(coalesce(toString(os), '__NULL__')) > 1
OR uniqExact(coalesce(toString(os_timezone), '__NULL__')) > 1
OR uniqExact(coalesce(toString(user_domain_id), '__NULL__')) > 1
),
geo_conflicts AS (
SELECT
g.click_id
FROM ods.geo_by_click AS g
INNER JOIN crossing AS c ON c.click_id = g.click_id
WHERE g.click_id IS NOT NULL
GROUP BY g.click_id
HAVING
uniqExact(coalesce(toString(geo_country), '__NULL__')) > 1
OR uniqExact(coalesce(toString(geo_timezone), '__NULL__')) > 1
OR uniqExact(coalesce(toString(geo_region_name), '__NULL__')) > 1
OR uniqExact(coalesce(toString(geo_latitude), '__NULL__')) > 1
OR uniqExact(coalesce(toString(geo_longitude), '__NULL__')) > 1
OR uniqExact(coalesce(toString(ip_address), '__NULL__')) > 1
)
SELECT
(SELECT count() FROM crossing) AS crossing_visits,
(SELECT count() FROM device_conflicts) AS ods_device_conflicts,
(SELECT count() FROM geo_conflicts) AS ods_geo_conflicts
FORMAT TabSeparated"
ods_context="$(ch_query "${ods_context_query}")"
IFS=$'\t' read -r ods_crossing_visits ods_device_conflicts ods_geo_conflicts <<< "${ods_context}"
[[ "${ods_crossing_visits}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать ODS-проверку стыка"
(( ods_crossing_visits == crossing_visits )) \
|| fail "ODS и DDS нашли разное число переходящих визитов: ${ods_crossing_visits}/${crossing_visits}"
(( ods_device_conflicts == 0 )) || fail "ODS device конфликтует на стыке: ${ods_device_conflicts}"
(( ods_geo_conflicts == 0 )) || fail "ODS geo конфликтует на стыке: ${ods_geo_conflicts}"
echo "ods_device_conflicts=${ods_device_conflicts}"
echo "ods_geo_conflicts=${ods_geo_conflicts}"
else
echo ""
echo "Проверка стыка backfill/live пропущена: CHECK_LIVE_SEAM=${CHECK_LIVE_SEAM}, live_rows=${live_rows:-0}"
fi
echo "" echo ""
echo "=== Проверка ClickHouse: основные DM-витрины не пустые ===" echo "=== Проверка ClickHouse: основные DM-витрины не пустые ==="