From 103ac021c8e40d7c6a1cb76a091bbafd126a10ed Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Wed, 22 Jul 2026 23:47:02 +0300 Subject: [PATCH] =?UTF-8?q?feat(generator):=20=D0=B8=D0=BD=D0=BA=D1=80?= =?UTF-8?q?=D0=B5=D0=BC=D0=B5=D0=BD=D1=82=D0=B0=D0=BB=D1=8C=D0=BD=D1=8B?= =?UTF-8?q?=D0=B5=20=D1=81=D1=87=D1=91=D1=82=D1=87=D0=B8=D0=BA=D0=B8=20man?= =?UTF-8?q?ifest=20=D0=B1=D0=B5=D0=B7=20=D0=BF=D0=B5=D1=80=D0=B5=D1=87?= =?UTF-8?q?=D0=B8=D1=82=D0=BA=D0=B8=20Kafka?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - world_next_day перечитывал всю историю топиков Kafka ради накопительных счётчиков — время прогона росло с возрастом мира (issue #5, находка F9). - Что: - счётчики засеваются при import из уже прочитанного артефакта и при backfill из потока; next-day продвигает их только событиями нового дня, полного чтения Kafka больше нет; - катящаяся контрольная сумма — сумма SHA-256 событий по модулю 2^256 (инкремент равен полному пересчёту), старый формат артефакта принимается без изменений; - точные множества click_id/user_domain_id вынесены из manifest в цепочку контент-адресуемых фрагментов (<=10 000 ID, SHA-256-цепочка, отдельный топик counter_chunks) — потолок сообщения Kafka не грозит, предел 900 000 байт проверяется явно с понятной ошибкой; - порядок записи всюду: фрагменты -> manifest -> state; старое локальное состояние отклоняется с подсказкой перезапустить import; - документация manifest/state обновлена (ARCHITECTURE, OPERATIONS, runbook startup-history). - Проверка: - make test (216+31) и make lint зелёные; - живая приёмка на чистом стенде: import 235 с; три прогона world_next_day — 716/718/716 с (плоское время, O(нового дня)); мир 3->6 дней, 561 942 события; make generated-history-chain-check — все порции и стыки однородны; - тест равенства инкремента и полного пересчёта: test_incremental_counters_equal_full_recompute. Co-Authored-By: Claude Fable 5 --- docs/ARCHITECTURE.md | 7 +- docs/OPERATIONS.md | 21 +- docs/runbooks/startup-history.md | 18 +- .../clickstream_generator/airflow_control.py | 35 ++ .../src/clickstream_generator/kafka_io.py | 57 ++- .../src/clickstream_generator/service.py | 77 +++- .../startup_history_artifact.py | 328 +++++++++++++++++- generator/src/clickstream_generator/state.py | 11 +- generator/tests/test_airflow_control.py | 5 + generator/tests/test_service.py | 111 +++++- .../tests/test_startup_history_artifact.py | 210 +++++++++++ generator/tests/test_state.py | 19 + 12 files changed, 845 insertions(+), 54 deletions(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index cbb0f12..c5c7b77 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -670,7 +670,12 @@ INSERT INTO dm.daily_traffic SELECT * FROM dm.v_daily_traffic; **DAG `world_next_day`** без параметров пакетно добавляет следующий модельный день, запускает `etl_pipeline` с `full_refresh` и сверяет manifest. У него задано -расписание каждые 30 минут, но DAG создаётся на паузе и не выполняет пропущенные интервалы. +расписание каждые 30 минут, но DAG создаётся на паузе и не выполняет пропущенные +интервалы. Счётчики manifest продолжаются из compact-топиков state/manifest и +обновляются событиями нового дня; старые data-топики Kafka не перечитываются. +Точные множества идентификаторов лежат отдельными неизменяемыми фрагментами, а +основная запись manifest хранит только числа и ссылку на их цепочку. Полная +пересборка ETL при этом сохраняется намеренно. **Подключение к ClickHouse:** - Connection: `clickhouse_default` diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 8cfa870..73d261f 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -82,7 +82,6 @@ Backfill/import требуют чистый стенд: пустые data-топ - Один запуск восстанавливает мир из state, добавляет 24 модельных часа, запускает `etl_pipeline` с полной пересборкой и сверяет витрины с manifest. - Расписание задано каждые 30 минут, но DAG создаётся на паузе; `catchup=False`. - Не включайте расписание до внедрения накопительных счётчиков manifest. - `max_active_runs=1` не даёт двум доливкам выполняться параллельно. `world_next_day` работает на непустом стенде и не использует проверку чистоты. @@ -107,13 +106,15 @@ data-топики, state, manifest. Автоматического отката не повторяйте доливку поверх возможного хвоста. Очистите стенд и переимпортируйте последний исправный портативный артефакт, затем повторите запуск `world_next_day`. -Текущая версия пересчитывает накопительные счётчики и контрольные суммы по всей -доступной истории data-топиков Kafka. Поэтому время выполнения и расход памяти -каждой доливки растут вместе с историей. Стандартный срок хранения Kafka тоже -ограничивает долгую работу стенда. Для ручного учебного цикла запускайте -соседние дни без долгих пауз. Перед включением расписания отдельная задача -должна выбрать одно из решений: бессрочное хранение data-топиков или новый -формат накопительного состояния manifest. +`next-day` не перечитывает старые data-топики Kafka. Основная запись manifest +хранит суммы, катящуюся контрольную сумму и SHA-256-ссылку на цепочку точных +множеств. Сами `click_id` и `user_domain_id` разбиты на небольшие неизменяемые +фрагменты в отдельном compact-топике. Новый день записывает только новый +фрагмент, поэтому размер одной записи не растёт с возрастом мира, а загрузка +основного manifest не сканирует всю цепочку. +State ссылается на основную запись manifest, чтобы две точки продолжения нельзя +было случайно смешать. Если локальный state создан старым кодом и раздела в нём +нет, повторно запустите `world_init` с операцией `import` на чистом стенде. Перед запуском `etl_pipeline` пульт проверяет, что DAG не стоит на паузе. Если стоит, задача падает сразу с подсказкой снять паузу в UI или командой: @@ -294,8 +295,8 @@ GEN_HISTORY_DURATION=2d make generated-history-analytics что state совпадает с manifest, и стартует ровно с `T_end` без настенной дельты. Новый backfill записывает `boundaries=[T0, T_end]`. Каждый успешный `next-day` -добавляет одну границу, сдвигает `model_t_end` и пересчитывает накопительные -счётчики и контрольные суммы по всей истории. +добавляет одну границу, сдвигает `model_t_end` и дополняет накопительные +счётчики только событиями нового дня. Полной перечитки Kafka в этом пути нет. Если историю нужно сохранить в файл и восстановить на чистом стенде без новой генерации, используйте [runbook стартовой истории](./runbooks/startup-history.md). diff --git a/docs/runbooks/startup-history.md b/docs/runbooks/startup-history.md index 9cb20e9..4b97d80 100644 --- a/docs/runbooks/startup-history.md +++ b/docs/runbooks/startup-history.md @@ -78,7 +78,8 @@ make startup-history-import 5. Когда нужен ещё один модельный день, запустите `world_next_day` с пустой формой. `world_next_day` имеет расписание каждые 30 минут, но по умолчанию стоит на паузе. -Не включайте расписание до внедрения накопительных счётчиков manifest. +Один запуск добавляет один модельный день; `max_active_runs=1` не допускает +параллельных доливок. ## Операции сопровождающего @@ -157,6 +158,21 @@ CHECK_LIVE_SEAM=0 make generated-history-check compact-топики в Kafka. ClickHouse читает данные через свои Kafka-таблицы и Materialized View, затем batch строит ODS, DDS и DM. +Во время импорта генератор за один проход по событиям создаёт локальное +накопительное состояние manifest. Эталонный файл не меняется: в нём этого +служебного раздела нет. Основная запись manifest в Kafka хранит суммы, катящуюся +контрольную сумму и ссылку на цепочку точных множеств. `click_id` и +`user_domain_id` лежат небольшими неизменяемыми фрагментами в отдельном +compact-топике. State хранит SHA-256-ссылку на основную запись. +`backfill` создаёт такое состояние сразу по ходу генерации. + +`world_next_day` восстанавливает эти числа и добавляет только события нового +дня, не читая старую историю Kafka и не перезаписывая прежние фрагменты. +Контрольная сумма складывает 256-битные +SHA-256-отпечатки записей по модулю 2^256, поэтому результат не зависит от +границ пакетов и совпадает с полным пересчётом. Если старый локальный state не +содержит ссылку на накопительные числа, очистите стенд и повторите `import`. + Импорт запускайте с тем же профилем, на котором создан артефакт. Для служебного артефакта `ci` профиль нужно задать явно: diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index 68d0171..42b32c9 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -21,8 +21,11 @@ from clickstream_generator.service import GeneratorService from clickstream_generator.startup_history_artifact import ( KafkaRawPublisher, KafkaTopicInspector, + ManifestCounters, compare_clickhouse_stats_to_manifest, + cumulative_counter_reference, import_startup_history_artifact, + known_counter_ids_from_state, load_startup_history_artifact, manifest_boundaries, validate_startup_history_artifact, @@ -171,6 +174,38 @@ def assert_next_day_snapshot(manifest: dict | None, state) -> None: raise RuntimeError( f"state не согласован с T_end manifest: {exc}" ) from exc + counter_reference = state.cumulative_manifest_counters + counter_state = manifest.get("cumulative_manifest_counters") + if counter_reference is None or counter_state is None: + raise RuntimeError( + "В state нет накопительных счётчиков manifest; " + "повторно запустите import эталонного мира." + ) + if cumulative_counter_reference(counter_state) != counter_reference: + raise RuntimeError( + "Накопительные счётчики state и manifest не совпадают; " + "повторно запустите import эталонного мира." + ) + known_click_ids, known_user_ids = known_counter_ids_from_state(state) + try: + counters = ManifestCounters.from_state( + counter_state, + known_click_ids=known_click_ids, + known_user_ids=known_user_ids, + ) + except ValueError as exc: + raise RuntimeError( + "Накопительные счётчики manifest повреждены; " + "повторно запустите import эталонного мира." + ) from exc + if ( + counters.to_manifest_topics() != manifest.get("topics") + or counters.to_manifest_totals() != manifest.get("totals") + ): + raise RuntimeError( + "Накопительные счётчики не совпадают с числами manifest; " + "повторно запустите import эталонного мира." + ) def _parse_utc_timestamp(value: str): diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py index de77cb6..1507e09 100644 --- a/generator/src/clickstream_generator/kafka_io.py +++ b/generator/src/clickstream_generator/kafka_io.py @@ -1,5 +1,6 @@ """Kafka-интеграция генератора.""" +import hashlib import json import logging import sys @@ -14,6 +15,7 @@ from clickstream_generator.state import UnsupportedStateVersionError logger = logging.getLogger("generator") DATA_TOPICS = ("browser_events", "location_events", "device_events", "geo_events") +MAX_COMPACT_MESSAGE_BYTES = 900_000 _kafka_imported = False KafkaProducer = None @@ -207,7 +209,9 @@ class KafkaStartupHistoryManifest: """Хранение манифеста стартовой истории в Kafka compact topic.""" MANIFEST_TOPIC = "generator_startup_history_manifest" + COUNTER_TOPIC = "generator_startup_history_counter_chunks" MANIFEST_KEY = "default" + COUNTER_CHUNK_KEY_PREFIX = "id-set:" def __init__(self, bootstrap_servers: str): self.bootstrap_servers = bootstrap_servers @@ -227,6 +231,8 @@ class KafkaStartupHistoryManifest: def save(self, manifest: dict) -> None: """Сохраняет манифест стартовой истории.""" + self._assert_message_size(manifest, "manifest") + def _do_send(): self.producer.send( self.MANIFEST_TOPIC, @@ -236,6 +242,40 @@ class KafkaStartupHistoryManifest: _retry(_do_send, max_retries=3, base_delay=0.5) + def save_counter_chunk(self, chunk_sha256: str, chunk: dict) -> None: + """Сохраняет неизменяемый фрагмент точных множеств идентификаторов.""" + if not isinstance(chunk_sha256, str) or len(chunk_sha256) != 64: + raise ValueError("хеш фрагмента множеств идентификаторов неверен") + actual_sha256 = hashlib.sha256( + json.dumps( + chunk, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=True, + ).encode("utf-8") + ).hexdigest() + if actual_sha256 != chunk_sha256: + raise ValueError("хеш фрагмента множеств идентификаторов не совпадает") + self._assert_message_size(chunk, "фрагмент множеств идентификаторов") + + def _do_send(): + self.producer.send( + self.COUNTER_TOPIC, + key=self.COUNTER_CHUNK_KEY_PREFIX + chunk_sha256, + value=chunk, + ) + + _retry(_do_send, max_retries=3, base_delay=0.5) + + @staticmethod + def _assert_message_size(value: dict, label: str) -> None: + size = len(json.dumps(value, ensure_ascii=True).encode("utf-8")) + if size > MAX_COMPACT_MESSAGE_BYTES: + raise ValueError( + f"{label} занимает {size} байт и превышает безопасный предел " + f"{MAX_COMPACT_MESSAGE_BYTES} байт одного сообщения Kafka" + ) + def flush(self) -> None: """Сбрасывает буфер с retry.""" def _do_flush(): @@ -372,8 +412,23 @@ def ensure_topics(bootstrap_servers: str) -> None: "delete.retention.ms": "100", }, ) + counter_topic = NewTopic( + name=KafkaStartupHistoryManifest.COUNTER_TOPIC, + num_partitions=1, + replication_factor=1, + topic_configs={ + "cleanup.policy": "compact", + "min.cleanable.dirty.ratio": "0.1", + "delete.retention.ms": "100", + }, + ) - for topic in [history_topic, state_topic, manifest_topic]: + for topic in [ + history_topic, + state_topic, + manifest_topic, + counter_topic, + ]: try: admin_client.create_topics([topic]) logger.info(f"Created topic: {topic.name}") diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 2755897..6bb19b1 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -16,7 +16,6 @@ from clickstream_generator.generation import EventGenerator from clickstream_generator.kafka_io import ( BatchRecord, KafkaBatchHistory, - KafkaDataTopicReader, KafkaPublisher, KafkaStateManager, KafkaStartupHistoryManifest, @@ -32,7 +31,9 @@ from clickstream_generator.startup_history_artifact import ( StartupHistoryArtifactBuilder, ManifestCounters, build_manifest, + cumulative_counter_reference, generation_settings_from_config, + known_counter_ids_from_state, manifest_boundaries, write_startup_history_artifact, ) @@ -58,7 +59,6 @@ class GeneratorService: self.history: KafkaBatchHistory | None = None self.state_manager: KafkaStateManager | None = None self.manifest_manager: KafkaStartupHistoryManifest | None = None - self.data_reader: KafkaDataTopicReader | None = None self._running = False self._stop_requested = False self._shutdown_event = threading.Event() @@ -133,10 +133,7 @@ class GeneratorService: restored_state, model_t_end=model_t_end, ) - self.data_reader = KafkaDataTopicReader( - self.config.kafka_bootstrap_servers - ) - self._run_next_day(manifest) + self._run_next_day(manifest, restored_state) self.stop() return @@ -471,13 +468,24 @@ class GeneratorService: model_t0=self.config.model_t0, gen_seed=self.config.seed, ) + cumulative_state, counter_chunks = ( + artifact_builder.counters.to_state_with_chunks() + ) + state.cumulative_manifest_counters = cumulative_counter_reference( + cumulative_state + ) manifest = build_manifest( config=self.config, counters=artifact_builder.counters, state=state, + cumulative_state=cumulative_state, ) if self.manifest_manager: + for chunk_sha256, chunk in counter_chunks: + self.manifest_manager.save_counter_chunk(chunk_sha256, chunk) + if counter_chunks: + self.manifest_manager.flush() self.manifest_manager.save(manifest) self.manifest_manager.flush() @@ -519,7 +527,7 @@ class GeneratorService: f"status={status}, sent_counts={sent_counts}" ) - def _run_next_day(self, manifest: dict) -> None: + def _run_next_day(self, manifest: dict, restored_state) -> None: """Доливает ровно 24 модельных часа от границы manifest.""" if not self.publisher: raise RuntimeError("publisher не инициализирован") @@ -528,6 +536,41 @@ class GeneratorService: if not self.state_manager or not self.manifest_manager: raise RuntimeError("менеджеры state и manifest не инициализированы") + counter_reference = restored_state.cumulative_manifest_counters + counter_state = manifest.get("cumulative_manifest_counters") + if counter_reference is None or counter_state is None: + raise IncompatibleStateError( + "в state нет накопительных счётчиков manifest; " + "повторно запустите import эталонного мира" + ) + if cumulative_counter_reference(counter_state) != counter_reference: + raise IncompatibleStateError( + "накопительные счётчики state и manifest не совпадают; " + "повторно запустите import эталонного мира" + ) + known_click_ids, known_user_ids = known_counter_ids_from_state( + restored_state + ) + try: + counters = ManifestCounters.from_state( + counter_state, + known_click_ids=known_click_ids, + known_user_ids=known_user_ids, + ) + except ValueError as exc: + raise IncompatibleStateError( + "накопительные счётчики manifest повреждены; " + "повторно запустите import эталонного мира" + ) from exc + if ( + counters.to_manifest_topics() != manifest.get("topics") + or counters.to_manifest_totals() != manifest.get("totals") + ): + raise IncompatibleStateError( + "накопительные счётчики не совпадают с числами manifest; " + "повторно запустите import эталонного мира" + ) + current_t_end = self._as_aware_utc( datetime.fromisoformat(manifest["model_t_end"]) ) @@ -554,6 +597,7 @@ class GeneratorService: ) total_sent, sent_counts, status = self._publish_batch(batch) self._raise_on_next_day_publish_error(batch_id, status, sent_counts) + counters.add_batch(batch) self._write_batch_history( batch_id=batch_id, started_at=started_at, @@ -573,6 +617,7 @@ class GeneratorService: started_at = datetime.now(timezone.utc) total_sent, sent_counts, status = self._publish_batch(final_batch) self._raise_on_next_day_publish_error(batch_id, status, sent_counts) + counters.add_batch(final_batch) self._write_batch_history( batch_id=batch_id, started_at=started_at, @@ -584,11 +629,6 @@ class GeneratorService: self.publisher.flush() self._model_time = target_t_end - reader = self.data_reader or KafkaDataTopicReader( - self.config.kafka_bootstrap_servers - ) - counters = ManifestCounters() - counters.add_batch(reader.load()) state_batch_id = self._startup_state_batch_id(counters) state = self.stream.to_state( tick=self._tick, @@ -602,6 +642,10 @@ class GeneratorService: model_t0=self.config.model_t0, gen_seed=self.config.seed, ) + cumulative_state, counter_chunks = counters.to_state_with_chunks() + state.cumulative_manifest_counters = cumulative_counter_reference( + cumulative_state + ) boundaries = manifest_boundaries(manifest) + [target_t_end.isoformat()] updated_manifest = build_manifest( self.config, @@ -609,12 +653,17 @@ class GeneratorService: state, model_t_end=target_t_end, boundaries=boundaries, + cumulative_state=cumulative_state, ) - self.state_manager.save(state) - self.state_manager.flush() + for chunk_sha256, chunk in counter_chunks: + self.manifest_manager.save_counter_chunk(chunk_sha256, chunk) + if counter_chunks: + self.manifest_manager.flush() self.manifest_manager.save(updated_manifest) self.manifest_manager.flush() + self.state_manager.save(state) + self.state_manager.flush() def _raise_on_next_day_publish_error( self, diff --git a/generator/src/clickstream_generator/startup_history_artifact.py b/generator/src/clickstream_generator/startup_history_artifact.py index 3ae2483..dba85d4 100644 --- a/generator/src/clickstream_generator/startup_history_artifact.py +++ b/generator/src/clickstream_generator/startup_history_artifact.py @@ -3,10 +3,12 @@ from __future__ import annotations import argparse +import base64 import hashlib import json import lzma import logging +import zlib from copy import deepcopy from datetime import datetime, timezone from pathlib import Path @@ -28,6 +30,11 @@ logger = logging.getLogger("generator") ARTIFACT_KIND = "clickstream_generator_startup_history" ARTIFACT_VERSION = "1.0" TOPICS = ("browser_events", "location_events", "device_events", "geo_events") +COUNTER_STATE_VERSION = "2.0" +CHECKSUM_MODULUS = 1 << 256 +ID_SET_CHUNK_VERSION = "1.0" +ID_SET_CHUNK_SIZE = 10_000 +ID_SET_STORAGE = "manifest-topic-sha256-chain-v1" IMPORT_TOPICS = ( "browser_events", "location_events", @@ -35,6 +42,7 @@ IMPORT_TOPICS = ( "geo_events", KafkaStateManager.STATE_TOPIC, KafkaStartupHistoryManifest.MANIFEST_TOPIC, + KafkaStartupHistoryManifest.COUNTER_TOPIC, ) @@ -45,17 +53,20 @@ class _TopicManifestStats: self.rows = 0 self.min_event_timestamp: str | None = None self.max_event_timestamp: str | None = None - self._checksum = hashlib.sha256() + self._checksum_sum = 0 @property def checksum(self) -> str: - return self._checksum.hexdigest() + return f"{self._checksum_sum:064x}" def add(self, event: dict, event_timestamp: str | None = None) -> None: self.rows += 1 - self._checksum.update( + event_digest = hashlib.sha256( json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") - ) + ).digest() + self._checksum_sum = ( + self._checksum_sum + int.from_bytes(event_digest, "big") + ) % CHECKSUM_MODULUS timestamp = ( event_timestamp if event_timestamp is not None @@ -79,6 +90,55 @@ class _TopicManifestStats: "checksum_sha256": self.checksum, } + @classmethod + def from_state(cls, payload: dict) -> "_TopicManifestStats": + """Восстанавливает накопительные числа одного топика.""" + if not isinstance(payload, dict): + raise ValueError("счётчики топика должны быть объектом") + rows = payload.get("rows") + checksum = payload.get("checksum_sha256") + if isinstance(rows, bool) or not isinstance(rows, int) or rows < 0: + raise ValueError("rows в накопительных счётчиках должен быть целым") + if not isinstance(checksum, str) or len(checksum) != 64: + raise ValueError("checksum_sha256 в накопительных счётчиках неверен") + try: + checksum_sum = int(checksum, 16) + except ValueError as exc: + raise ValueError( + "checksum_sha256 в накопительных счётчиках неверен" + ) from exc + + stats = cls() + stats.rows = rows + stats.min_event_timestamp = payload.get("min_event_timestamp") + stats.max_event_timestamp = payload.get("max_event_timestamp") + stats._checksum_sum = checksum_sum + return stats + + +class _LegacyTopicManifestStats(_TopicManifestStats): + """Прежняя контрольная сумма для проверки неизменённого артефакта.""" + + def __init__(self): + super().__init__() + self._legacy_checksum = hashlib.sha256() + + @property + def checksum(self) -> str: + return self._legacy_checksum.hexdigest() + + def add(self, event: dict, event_timestamp: str | None = None) -> None: + self.rows += 1 + self._legacy_checksum.update( + json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") + ) + timestamp = ( + event_timestamp + if event_timestamp is not None + else event.get("event_timestamp") + ) + self.add_timestamp(timestamp) + class ManifestCounters: """Счётчики стартовой истории для manifest.""" @@ -87,6 +147,11 @@ class ManifestCounters: self.topic_stats = {topic: _TopicManifestStats() for topic in TOPICS} self.click_ids: set[str] = set() self.user_ids: set[str] = set() + self._new_click_ids: set[str] = set() + self._new_user_ids: set[str] = set() + self._visit_count = 0 + self._user_count = 0 + self._previous_id_set_chain = _empty_id_set_chain() def add_batch(self, batch: dict[str, list[dict]]) -> None: browser_events = batch.get("browser_events", []) @@ -117,11 +182,23 @@ class ManifestCounters: stats.add_timestamp(max(timestamps)) stats.add(event, event_timestamp=event_timestamp) click_id = event.get("click_id") - if topic == "browser_events" and click_id: + if ( + topic == "browser_events" + and click_id + and click_id not in self.click_ids + ): self.click_ids.add(click_id) + self._new_click_ids.add(click_id) + self._visit_count += 1 user_id = event.get("user_domain_id") - if topic == "device_events" and user_id: + if ( + topic == "device_events" + and user_id + and user_id not in self.user_ids + ): self.user_ids.add(user_id) + self._new_user_ids.add(user_id) + self._user_count += 1 def to_manifest_topics(self) -> dict: return {topic: stats.to_dict() for topic, stats in self.topic_stats.items()} @@ -130,12 +207,214 @@ class ManifestCounters: browser_stats = self.topic_stats["browser_events"] return { "events": browser_stats.rows, - "visits": len(self.click_ids), - "users": len(self.user_ids), + "visits": self._visit_count, + "users": self._user_count, "min_event_timestamp": browser_stats.min_event_timestamp, "max_event_timestamp": browser_stats.max_event_timestamp, } + def to_state(self) -> dict: + """Возвращает JSON-сериализуемое накопительное состояние.""" + state, _chunks = self.to_state_with_chunks() + return state + + def to_state_with_chunks(self) -> tuple[dict, list[tuple[str, dict]]]: + """Возвращает компактное состояние и новые фрагменты точных множеств.""" + id_set_chain, chunks = _extend_id_set_chain( + self._previous_id_set_chain, + self._new_click_ids, + self._new_user_ids, + ) + state = { + "version": COUNTER_STATE_VERSION, + "checksum_algorithm": "sha256-sum-v1", + "topics": self.to_manifest_topics(), + "totals": { + "visits": self._visit_count, + "users": self._user_count, + }, + "id_sets": id_set_chain, + } + return state, chunks + + @classmethod + def from_state( + cls, + payload: dict, + *, + known_click_ids: set[str] | None = None, + known_user_ids: set[str] | None = None, + ) -> "ManifestCounters": + """Продолжает счётчики из state без чтения старых событий.""" + if not isinstance(payload, dict): + raise ValueError("накопительные счётчики должны быть объектом") + if payload.get("version") != COUNTER_STATE_VERSION: + raise ValueError("версия накопительных счётчиков не поддерживается") + if payload.get("checksum_algorithm") != "sha256-sum-v1": + raise ValueError("алгоритм накопительной контрольной суммы не поддерживается") + topic_payloads = payload.get("topics") + if not isinstance(topic_payloads, dict) or set(topic_payloads) != set(TOPICS): + raise ValueError("в накопительных счётчиках нет всех топиков") + totals = payload.get("totals") + if not isinstance(totals, dict): + raise ValueError("в накопительных счётчиках нет итогов") + visits = _non_negative_int(totals.get("visits"), "visits") + users = _non_negative_int(totals.get("users"), "users") + id_sets = payload.get("id_sets") + _validate_id_set_chain(id_sets, visits=visits, users=users) + + counters = cls() + counters.topic_stats = { + topic: _TopicManifestStats.from_state(topic_payloads[topic]) + for topic in TOPICS + } + counters.click_ids = set(known_click_ids or ()) + counters.user_ids = set(known_user_ids or ()) + counters._new_click_ids = set() + counters._new_user_ids = set() + counters._visit_count = visits + counters._user_count = users + counters._previous_id_set_chain = dict(id_sets) + return counters + + +def cumulative_counter_reference(counter_state: dict) -> dict: + """Строит компактную ссылку state на полный набор из manifest.""" + canonical = json.dumps( + counter_state, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=True, + ).encode("utf-8") + return { + "version": COUNTER_STATE_VERSION, + "manifest_sha256": hashlib.sha256(canonical).hexdigest(), + } + + +def known_counter_ids_from_state(state: GeneratorState) -> tuple[set[str], set[str]]: + """Возвращает идентификаторы, которые уже успели попасть в data-топики.""" + click_ids = {visit["click_id"] for visit in state.active_visits} + user_ids = { + user["user_domain_id"] + for user in state.population + if user.get("active_click_id") is not None + or user.get("last_finished_at") is not None + } + return click_ids, user_ids + + +def _empty_id_set_chain() -> dict: + return { + "storage": ID_SET_STORAGE, + "latest_sha256": None, + "chunks": 0, + "click_ids": 0, + "user_ids": 0, + } + + +def _non_negative_int(value: Any, label: str) -> int: + if isinstance(value, bool) or not isinstance(value, int) or value < 0: + raise ValueError(f"{label} в накопительных счётчиках должен быть целым") + return value + + +def _validate_id_set_chain(payload: Any, *, visits: int, users: int) -> None: + if not isinstance(payload, dict): + raise ValueError("в накопительных счётчиках нет множеств идентификаторов") + if payload.get("storage") != ID_SET_STORAGE: + raise ValueError("хранилище множеств идентификаторов не поддерживается") + chunks = _non_negative_int(payload.get("chunks"), "chunks") + click_ids = _non_negative_int(payload.get("click_ids"), "click_ids") + user_ids = _non_negative_int(payload.get("user_ids"), "user_ids") + latest = payload.get("latest_sha256") + if latest is not None: + if not isinstance(latest, str) or len(latest) != 64: + raise ValueError("ссылка на фрагмент множеств идентификаторов неверна") + try: + int(latest, 16) + except ValueError as exc: + raise ValueError( + "ссылка на фрагмент множеств идентификаторов неверна" + ) from exc + if (chunks == 0) != (latest is None): + raise ValueError("цепочка множеств идентификаторов повреждена") + if click_ids != visits or user_ids != users: + raise ValueError("размеры множеств идентификаторов не совпадают с итогами") + + +def _extend_id_set_chain( + previous: dict, + click_ids: set[str], + user_ids: set[str], +) -> tuple[dict, list[tuple[str, dict]]]: + _validate_id_set_chain( + previous, + visits=_non_negative_int(previous.get("click_ids"), "click_ids"), + users=_non_negative_int(previous.get("user_ids"), "user_ids"), + ) + tagged_ids = [("click_id", value) for value in sorted(click_ids)] + tagged_ids.extend(("user_id", value) for value in sorted(user_ids)) + previous_sha256 = previous["latest_sha256"] + chunk_count = previous["chunks"] + chunks: list[tuple[str, dict]] = [] + for offset in range(0, len(tagged_ids), ID_SET_CHUNK_SIZE): + part = tagged_ids[offset : offset + ID_SET_CHUNK_SIZE] + part_click_ids = { + value for kind, value in part if kind == "click_id" + } + part_user_ids = { + value for kind, value in part if kind == "user_id" + } + chunk = { + "version": ID_SET_CHUNK_VERSION, + "encoding": "json-zlib-base64-v1", + "sequence": chunk_count + 1, + "previous_sha256": previous_sha256, + "click_ids": _encode_id_set(part_click_ids), + "user_ids": _encode_id_set(part_user_ids), + } + digest = hashlib.sha256(_canonical_json(chunk)).hexdigest() + chunks.append((digest, chunk)) + previous_sha256 = digest + chunk_count += 1 + + return ( + { + "storage": ID_SET_STORAGE, + "latest_sha256": previous_sha256, + "chunks": chunk_count, + "click_ids": previous["click_ids"] + len(click_ids), + "user_ids": previous["user_ids"] + len(user_ids), + }, + chunks, + ) + + +def _canonical_json(value: dict) -> bytes: + return json.dumps( + value, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=True, + ).encode("utf-8") + + +def _encode_id_set(values: set[str]) -> str: + raw = json.dumps(sorted(values), separators=(",", ":")).encode("utf-8") + return base64.b64encode(zlib.compress(raw, level=9)).decode("ascii") + + +class _LegacyManifestCounters(ManifestCounters): + """Полный пересчёт старой суммы только при проверке артефакта.""" + + def __init__(self): + super().__init__() + self.topic_stats = { + topic: _LegacyTopicManifestStats() for topic in TOPICS + } + class StartupHistoryArtifactBuilder: """Собирает сообщения топиков для портативного файла.""" @@ -158,12 +437,16 @@ class StartupHistoryArtifactBuilder: ) def to_artifact(self, manifest: dict, state: GeneratorState) -> dict: + artifact_manifest = deepcopy(manifest) + artifact_manifest.pop("cumulative_manifest_counters", None) + artifact_state = state.to_dict() + artifact_state.pop("cumulative_manifest_counters", None) return { "artifact_version": ARTIFACT_VERSION, "kind": ARTIFACT_KIND, "created_at": datetime.now(timezone.utc).isoformat(), - "manifest": deepcopy(manifest), - "state": state.to_dict(), + "manifest": artifact_manifest, + "state": artifact_state, "topics": deepcopy(self.topics), "raw_topics": deepcopy(self.raw_topics), } @@ -295,6 +578,7 @@ def build_manifest( *, model_t_end: datetime | None = None, boundaries: list[str] | None = None, + cumulative_state: dict | None = None, ) -> dict: """Строит manifest стартовой истории.""" selected_t_end = model_t_end or config.model_t_end @@ -325,6 +609,7 @@ def build_manifest( }, "topics": counters.to_manifest_topics(), "totals": counters.to_manifest_totals(), + "cumulative_manifest_counters": cumulative_state or counters.to_state(), } manifest_boundaries(manifest) return manifest @@ -437,6 +722,16 @@ def import_startup_history_artifact( artifact, expected_config=expected_config, ) + counters = ManifestCounters() + counters.add_batch(topics) + cumulative_state, counter_chunks = counters.to_state_with_chunks() + state.cumulative_manifest_counters = cumulative_counter_reference( + cumulative_state + ) + manifest = deepcopy(manifest) + manifest["topics"] = counters.to_manifest_topics() + manifest["totals"] = counters.to_manifest_totals() + manifest["cumulative_manifest_counters"] = cumulative_state if topic_inspector is not None: topic_inspector.assert_data_topics_empty() snapshot = ( @@ -471,6 +766,10 @@ def import_startup_history_artifact( sent_by_topic[topic] = sent publisher.flush() + for chunk_sha256, chunk in counter_chunks: + manifest_manager.save_counter_chunk(chunk_sha256, chunk) + if counter_chunks: + manifest_manager.flush() manifest_manager.save(manifest) manifest_manager.flush() state_manager.save(state) @@ -548,10 +847,15 @@ def _validate_manifest_topics(manifest: dict, topics: dict[str, list[dict]]) -> counters.add_batch({topic: topics[topic] for topic in TOPICS}) expected_topics = manifest.get("topics") or {} expected_totals = manifest.get("totals") or {} - if counters.to_manifest_topics() != expected_topics: - raise ValueError("startup history artifact mismatch: topic checksums") if counters.to_manifest_totals() != expected_totals: raise ValueError("startup history artifact mismatch: totals") + if counters.to_manifest_topics() == expected_topics: + return + + legacy_counters = _LegacyManifestCounters() + legacy_counters.add_batch({topic: topics[topic] for topic in TOPICS}) + if legacy_counters.to_manifest_topics() != expected_topics: + raise ValueError("startup history artifact mismatch: topic checksums") def _validate_raw_topics( diff --git a/generator/src/clickstream_generator/state.py b/generator/src/clickstream_generator/state.py index 4d70463..9d45a67 100644 --- a/generator/src/clickstream_generator/state.py +++ b/generator/src/clickstream_generator/state.py @@ -204,6 +204,7 @@ class GeneratorState: population: list[dict] = field(default_factory=list) active_visits: list[dict] = field(default_factory=list) pending_visit_births: float = 0.0 + cumulative_manifest_counters: dict | None = None def __post_init__(self) -> None: if self.model_timestamp is None: @@ -221,7 +222,7 @@ class GeneratorState: def to_dict(self) -> dict: """Конвертирует в словарь для JSON-сериализации.""" - return { + payload = { "tick": self.tick, "rng_state": self.rng_state, "last_batch_id": self.last_batch_id, @@ -237,6 +238,11 @@ class GeneratorState: "active_visits": self.active_visits, "pending_visit_births": self.pending_visit_births, } + if self.cumulative_manifest_counters is not None: + payload["cumulative_manifest_counters"] = ( + self.cumulative_manifest_counters + ) + return payload @classmethod def from_dict(cls, data: dict) -> "GeneratorState": @@ -290,6 +296,9 @@ class GeneratorState: population=data.get("population", []), active_visits=data.get("active_visits", []), pending_visit_births=data.get("pending_visit_births", 0.0), + cumulative_manifest_counters=data.get( + "cumulative_manifest_counters" + ), ) except UnsupportedStateVersionError: raise diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py index a4806b9..03452ad 100644 --- a/generator/tests/test_airflow_control.py +++ b/generator/tests/test_airflow_control.py @@ -149,6 +149,11 @@ def test_next_day_precheck_requires_manifest_and_matching_state(): with pytest.raises(RuntimeError, match="state.*T_end"): assert_next_day_snapshot(manifest, state) + manifest["model_t_end"] = state.model_timestamp.isoformat() + manifest["state"]["model_timestamp"] = state.model_timestamp.isoformat() + with pytest.raises(RuntimeError, match="повторно запустите import"): + assert_next_day_snapshot(manifest, state) + def test_target_dag_trigger_error_covers_missing_paused_and_ready_states(): """Проверка зависимого DAG различает три реальные ветки.""" diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 4eda6f4..3668f0f 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -25,6 +25,7 @@ from generator import ( main, ) from clickstream_generator.state import UnsupportedStateVersionError +from clickstream_generator.service import IncompatibleStateError class TestBatchRecordWithDictConversion: @@ -425,6 +426,10 @@ class TestGeneratorServiceBackfill: self, base_config ): """Backfill пишет [T0, T_end), state на T_end и повторяемый manifest.""" + from clickstream_generator.startup_history_artifact import ( + cumulative_counter_reference, + ) + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) model_t_end = model_t0 + timedelta(minutes=3) config = replace( @@ -463,6 +468,11 @@ class TestGeneratorServiceBackfill: assert manifest["model_t0"] == model_t0.isoformat() assert manifest["model_t_end"] == model_t_end.isoformat() assert manifest["state"]["last_batch_id"] == saved_state.last_batch_id + assert saved_state.cumulative_manifest_counters == ( + cumulative_counter_reference( + manifest["cumulative_manifest_counters"] + ) + ) assert manifest["topics"]["browser_events"]["rows"] == len(browser_events) for topic in ( "browser_events", @@ -516,6 +526,8 @@ class TestGeneratorServiceBackfill: self._run_backfill(config) artifact = load_startup_history_artifact(artifact_path) + assert "cumulative_manifest_counters" not in artifact["state"] + assert "cumulative_manifest_counters" not in artifact["manifest"] state, manifest, topics, _ = validate_startup_history_artifact( artifact, expected_config=config, @@ -960,7 +972,8 @@ class TestGeneratorServiceBackfill: with pytest.raises(RuntimeError, match="flush down"): service._run_backfill() - service.manifest_manager.save.assert_called_once() + service.manifest_manager.save_counter_chunk.assert_called() + service.manifest_manager.save.assert_not_called() service.manifest_manager.flush.assert_called_once() service.state_manager.save.assert_not_called() service.state_manager.flush.assert_not_called() @@ -1008,13 +1021,14 @@ class TestGeneratorServiceBackfill: class TestGeneratorServiceNextDay: """Проверки ограниченной доливки следующего модельного дня.""" - def test_next_day_publishes_24h_then_state_then_cumulative_manifest( + def test_next_day_publishes_24h_then_chunks_manifest_and_state( self, base_config ): - """Next-day пишет [T_end, T_end+24h), затем state и manifest.""" + """Next-day пишет день, фрагменты множеств, manifest и затем state.""" from clickstream_generator.startup_history_artifact import ( StartupHistoryArtifactBuilder, build_manifest, + cumulative_counter_reference, ) model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) @@ -1062,6 +1076,9 @@ class TestGeneratorServiceNextDay: } initial_builder = StartupHistoryArtifactBuilder() initial_builder.add_batch(old_batch) + state.cumulative_manifest_counters = cumulative_counter_reference( + initial_builder.counters.to_state() + ) manifest = build_manifest(config, initial_builder.counters, state) published = {topic: [] for topic in old_batch} @@ -1078,21 +1095,21 @@ class TestGeneratorServiceNextDay: service.state_manager = MagicMock() service.state_manager.save.side_effect = lambda _state: order.append("state") service.manifest_manager = MagicMock() + service.manifest_manager.save_counter_chunk.side_effect = ( + lambda _chunk_sha256, _chunk: order.append("counter_chunk") + ) service.manifest_manager.save.side_effect = ( lambda _manifest: order.append("manifest") ) class DataReader: def load(self): - return { - topic: old_batch[topic] + published[topic] - for topic in old_batch - } + raise AssertionError("next-day не должен перечитывать Kafka") service.data_reader = DataReader() service.restore_from_startup_history(state, model_t_end=current_t_end) - service._run_next_day(manifest) + service._run_next_day(manifest, state) browser_events = published["browser_events"] timestamps = [ @@ -1114,18 +1131,21 @@ class TestGeneratorServiceNextDay: target_t_end.isoformat(), ] assert saved_manifest["totals"]["events"] == len(browser_events) + 1 - assert order == ["data", "state", "manifest"] + assert order == ["data", "counter_chunk", "manifest", "state"] first_day_events = len(browser_events) first_day_checksum = saved_manifest["topics"]["browser_events"][ "checksum_sha256" ] - service._run_next_day(saved_manifest) + service._run_next_day(saved_manifest, saved_state) second_state = service.state_manager.save.call_args.args[0] second_manifest = service.manifest_manager.save.call_args.args[0] expected_builder = StartupHistoryArtifactBuilder() - expected_builder.add_batch(DataReader().load()) + expected_builder.add_batch({ + topic: old_batch[topic] + published[topic] + for topic in old_batch + }) assert second_state.model_timestamp == target_t_end + timedelta(hours=24) assert second_manifest["boundaries"] == [ model_t0.isoformat(), @@ -1140,15 +1160,57 @@ class TestGeneratorServiceNextDay: != first_day_checksum ) assert second_manifest["topics"] == expected_builder.counters.to_manifest_topics() + assert second_manifest["totals"] == expected_builder.counters.to_manifest_totals() + assert second_state.cumulative_manifest_counters == ( + cumulative_counter_reference( + second_manifest["cumulative_manifest_counters"] + ) + ) assert order == [ "data", - "state", + "counter_chunk", "manifest", + "state", "data", - "state", + "counter_chunk", "manifest", + "state", ] + def test_next_day_requires_cumulative_counters_from_import(self, base_config): + """Старый локальный state требует повторного импорта, а не чтения Kafka.""" + model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) + config = replace( + base_config, + run_mode="next-day", + model_t0=model_t0, + model_t_end=model_t0 + timedelta(hours=1), + ) + service = GeneratorService(config) + state = service.stream.to_state( + tick=1, + rng_state=service.generator.rng.getstate(), + last_batch_id="startup-history-old", + last_timestamp=config.model_t_end, + model_timestamp=config.model_t_end, + model_t0=config.model_t0, + ) + service.publisher = MagicMock() + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + service._model_time = config.model_t_end + manifest = { + "model_t0": model_t0.isoformat(), + "model_t_end": config.model_t_end.isoformat(), + "boundaries": [model_t0.isoformat(), config.model_t_end.isoformat()], + } + + with pytest.raises(IncompatibleStateError, match="повторно запустите import"): + service._run_next_day(manifest, state) + + service.publisher.publish.assert_not_called() + def test_next_day_publish_error_does_not_move_state_or_manifest(self, base_config): """Ошибка data-топика оставляет обе точки фиксации без изменений.""" model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) @@ -1177,8 +1239,29 @@ class TestGeneratorServiceNextDay: "boundaries": [model_t0.isoformat(), config.model_t_end.isoformat()], } + from clickstream_generator.startup_history_artifact import ( + ManifestCounters, + cumulative_counter_reference, + ) + + state = service.stream.to_state( + tick=1, + rng_state=service.generator.rng.getstate(), + last_batch_id="startup-history-test", + last_timestamp=config.model_t_end, + model_timestamp=config.model_t_end, + model_t0=config.model_t0, + ) + empty_counters = ManifestCounters() + manifest["cumulative_manifest_counters"] = empty_counters.to_state() + manifest["topics"] = empty_counters.to_manifest_topics() + manifest["totals"] = empty_counters.to_manifest_totals() + state.cumulative_manifest_counters = cumulative_counter_reference( + manifest["cumulative_manifest_counters"] + ) + with pytest.raises(RuntimeError, match="Публикация next-day не удалась"): - service._run_next_day(manifest) + service._run_next_day(manifest, state) service.state_manager.save.assert_not_called() service.manifest_manager.save.assert_not_called() diff --git a/generator/tests/test_startup_history_artifact.py b/generator/tests/test_startup_history_artifact.py index d70c2d6..821f3be 100644 --- a/generator/tests/test_startup_history_artifact.py +++ b/generator/tests/test_startup_history_artifact.py @@ -4,6 +4,7 @@ from dataclasses import replace from datetime import datetime, timezone +import hashlib import json import lzma @@ -110,6 +111,144 @@ def _batch() -> dict[str, list[dict]]: } +def test_incremental_counters_equal_full_recompute(): + """Накопление по частям даёт те же числа, что полный пересчёт мира.""" + from clickstream_generator.startup_history_artifact import ManifestCounters + + first_day = _batch() + second_day = { + "browser_events": [ + { + "event_id": "event-2", + "click_id": "click-1", + "user_domain_id": "user-1", + "event_timestamp": "2026-01-02 00:00:00.000000", + }, + { + "event_id": "event-3", + "click_id": "click-2", + "user_domain_id": "user-2", + "event_timestamp": "2026-01-02 00:01:00.000000", + }, + ], + "location_events": [ + {"event_id": "event-2", "page_url_path": "/cart"}, + {"event_id": "event-3", "page_url_path": "/home"}, + ], + "device_events": [ + {"click_id": "click-1", "user_domain_id": "user-1"}, + {"click_id": "click-2", "user_domain_id": "user-2"}, + ], + "geo_events": [ + {"click_id": "click-1", "country": "RU"}, + {"click_id": "click-2", "country": "KZ"}, + ], + } + + incremental = ManifestCounters() + incremental.add_batch(first_day) + incremental = ManifestCounters.from_state( + incremental.to_state(), + known_click_ids={"click-1"}, + known_user_ids={"user-1"}, + ) + incremental.add_batch(second_day) + + full = ManifestCounters() + full.add_batch({ + topic: first_day[topic] + second_day[topic] + for topic in first_day + }) + + assert incremental.to_manifest_topics() == full.to_manifest_topics() + assert incremental.to_manifest_totals() == full.to_manifest_totals() + + +def test_exact_id_sets_are_split_into_bounded_hash_chain(monkeypatch): + """Manifest ссылается на цепочку малых точных фрагментов, а не растёт сам.""" + from clickstream_generator import startup_history_artifact as module + + monkeypatch.setattr(module, "ID_SET_CHUNK_SIZE", 2) + counters = module.ManifestCounters() + batch = _batch() + batch["browser_events"].extend([ + { + "event_id": "event-2", + "click_id": "click-2", + "event_timestamp": "2026-01-01 00:01:00.000000", + }, + { + "event_id": "event-3", + "click_id": "click-3", + "event_timestamp": "2026-01-01 00:02:00.000000", + }, + ]) + counters.add_batch(batch) + + state, chunks = counters.to_state_with_chunks() + + assert len(chunks) == 2 + assert state["totals"] == {"visits": 3, "users": 1} + assert state["id_sets"] == { + "storage": "manifest-topic-sha256-chain-v1", + "latest_sha256": chunks[-1][0], + "chunks": 2, + "click_ids": 3, + "user_ids": 1, + } + assert chunks[0][1]["previous_sha256"] is None + assert chunks[1][1]["previous_sha256"] == chunks[0][0] + assert "click-1" not in json.dumps(state) + + restored = module.ManifestCounters.from_state( + state, + known_click_ids={"click-3"}, + known_user_ids={"user-1"}, + ) + restored.add_batch({ + "browser_events": [ + { + "event_id": "event-4", + "click_id": "click-3", + "event_timestamp": "2026-01-02 00:00:00.000000", + }, + { + "event_id": "event-5", + "click_id": "click-4", + "event_timestamp": "2026-01-02 00:01:00.000000", + }, + ], + "location_events": [], + "device_events": [{"click_id": "click-4", "user_domain_id": "user-2"}], + "geo_events": [], + }) + next_state, next_chunks = restored.to_state_with_chunks() + + assert next_state["totals"] == {"visits": 4, "users": 2} + assert len(next_chunks) == 1 + assert next_chunks[0][1]["previous_sha256"] == chunks[-1][0] + assert next_state["id_sets"]["chunks"] == 3 + + +def test_counter_chunk_has_explicit_kafka_size_and_hash_guards(): + """Большой или неверно адресованный фрагмент падает до отправки в Kafka.""" + from clickstream_generator.kafka_io import ( + MAX_COMPACT_MESSAGE_BYTES, + KafkaStartupHistoryManifest, + ) + + manager = object.__new__(KafkaStartupHistoryManifest) + with pytest.raises(ValueError, match="безопасный предел"): + manager._assert_message_size( + {"payload": "x" * MAX_COMPACT_MESSAGE_BYTES}, + "фрагмент", + ) + with pytest.raises(ValueError, match="хеш.*не совпадает"): + manager.save_counter_chunk("0" * 64, {"version": "1.0"}) + + assert manager.COUNTER_TOPIC != manager.MANIFEST_TOPIC + + def test_artifact_roundtrip_keeps_events_state_and_manifest(base_config): """Артефакт хранит события, state и manifest одним проверяемым набором.""" from clickstream_generator.startup_history_artifact import ( @@ -173,6 +312,43 @@ def test_artifact_roundtrip_keeps_events_state_and_manifest(base_config): path.unlink(missing_ok=True) +def test_legacy_artifact_checksum_remains_valid(base_config): + """Эталонный артефакт со старой суммой не требует пересборки.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + validate_startup_history_artifact, + ) + + state = _state() + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + manifest = build_manifest(config, builder.counters, state) + for topic, events in builder.topics.items(): + checksum = hashlib.sha256() + for event in events: + checksum.update( + json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") + ) + manifest["topics"][topic]["checksum_sha256"] = checksum.hexdigest() + + artifact = builder.to_artifact(manifest=manifest, state=state) + + _, loaded_manifest, _, _ = validate_startup_history_artifact( + artifact, + expected_config=config, + ) + assert loaded_manifest["topics"] == artifact["manifest"]["topics"] + + def test_manifest_and_artifact_show_launch_profile(base_config): """Manifest и файл артефакта показывают выбранный профиль запуска.""" from clickstream_generator.startup_history_artifact import ( @@ -372,6 +548,8 @@ def test_raw_topic_value_can_keep_original_json_bytes(base_config): manifest=build_manifest(config=config, counters=builder.counters, state=state), state=state, ) + assert "cumulative_manifest_counters" not in artifact["state"] + assert "cumulative_manifest_counters" not in artifact["manifest"] raw_value = ( '{"event_timestamp":"2026-01-01 00:00:00.000000",' '"user_domain_id":"user-1","click_id":"click-1","event_id":"event-1"}' @@ -391,9 +569,15 @@ def test_raw_topic_value_can_keep_original_json_bytes(base_config): return None class CompactWriter: + def __init__(self): + self.chunks = [] + def save(self, value): self.saved = value + def save_counter_chunk(self, chunk_sha256, chunk): + self.chunks.append((chunk_sha256, chunk)) + def flush(self): return None @@ -609,8 +793,10 @@ def test_import_rolls_back_kafka_topics_after_partial_publish(base_config): def test_import_replays_events_and_compact_topics(base_config): """Импорт пишет события в Kafka и служебные compact-топики, не трогая ClickHouse.""" from clickstream_generator.startup_history_artifact import ( + ManifestCounters, StartupHistoryArtifactBuilder, build_manifest, + cumulative_counter_reference, import_startup_history_artifact, ) @@ -646,10 +832,14 @@ def test_import_replays_events_and_compact_topics(base_config): def __init__(self): self.saved = None self.flushed = False + self.chunks = [] def save(self, value): self.saved = value + def save_counter_chunk(self, chunk_sha256, chunk): + self.chunks.append((chunk_sha256, chunk)) + def flush(self): self.flushed = True @@ -678,7 +868,27 @@ def test_import_replays_events_and_compact_topics(base_config): ) assert result["events"] == 4 assert state_manager.saved.last_batch_id == state.last_batch_id + assert state_manager.saved.cumulative_manifest_counters is not None assert manifest_manager.saved["state"]["last_batch_id"] == state.last_batch_id + assert state_manager.saved.cumulative_manifest_counters == ( + cumulative_counter_reference( + manifest_manager.saved["cumulative_manifest_counters"] + ) + ) + seeded = ManifestCounters.from_state( + manifest_manager.saved["cumulative_manifest_counters"], + known_click_ids={"click-1"}, + known_user_ids={"user-1"}, + ) + assert seeded.click_ids == {"click-1"} + assert seeded.user_ids == {"user-1"} + assert len(manifest_manager.chunks) == 1 + assert ( + manifest_manager.saved["cumulative_manifest_counters"]["id_sets"][ + "latest_sha256" + ] + == manifest_manager.chunks[0][0] + ) assert publisher.flushed assert state_manager.flushed assert manifest_manager.flushed diff --git a/generator/tests/test_state.py b/generator/tests/test_state.py index e3b2228..964d243 100644 --- a/generator/tests/test_state.py +++ b/generator/tests/test_state.py @@ -227,6 +227,25 @@ class TestGeneratorState: assert next_values == values_after + def test_state_roundtrip_keeps_cumulative_manifest_counters(self): + """State сохраняет накопительные числа для следующего дня.""" + counter_reference = { + "version": "1.0", + "manifest_sha256": "a" * 64, + } + state = GeneratorState( + tick=1, + rng_state=_make_valid_rng_state(), + last_batch_id="startup-history-test", + last_timestamp=datetime.now(timezone.utc), + population=_minimal_population(), + cumulative_manifest_counters=counter_reference, + ) + + restored = GeneratorState.from_dict(json.loads(json.dumps(state.to_dict()))) + + assert restored.cumulative_manifest_counters == counter_reference + def test_state_roundtrip_keeps_population_and_active_visits(self): """State хранит популяцию и активные визиты в JSON.""" rng_state = _make_valid_rng_state(42)