diff --git a/.scratch/generator-model-time-startup-history/issues/13-backfill-top-up-from-snapshot.md b/.scratch/generator-model-time-startup-history/issues/13-backfill-top-up-from-snapshot.md index 13a040a..ea7e6c1 100644 --- a/.scratch/generator-model-time-startup-history/issues/13-backfill-top-up-from-snapshot.md +++ b/.scratch/generator-model-time-startup-history/issues/13-backfill-top-up-from-snapshot.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: done # Глагол next-day: собрать следующий модельный день @@ -111,26 +111,26 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Глагол `next-day` в `generator_control`: батч ровно за +- [x] Глагол `next-day` в `generator_control`: батч ровно за `[T_end, T_end + 24h)` от `T_end` манифеста; после прогона манифест содержит новую границу в `boundaries` и новый `T_end`. -- [ ] Предпроверка границы: на пустом стенде (нет манифеста) next-day +- [x] Предпроверка границы: на пустом стенде (нет манифеста) next-day громко падает; при работающем live-генераторе — громко падает; `clean`-guard и существующие глаголы `backfill|import|check` не изменены. -- [ ] Идемпотентность: запуск с `expected_t_end`, не равным `T_end` +- [x] Идемпотентность: запуск с `expected_t_end`, не равным `T_end` манифеста, — громкий отказ с обеими границами в тексте; повторный запуск того же дня после успеха — тот же отказ, дублей и второго мира нет. -- [ ] Проверка цепочки (новая make-цель): проходит по всем границам из +- [x] Проверка цепочки (новая make-цель): проходит по всем границам из `boundaries`; после N прогонов next-day в цепочке N новых границ; непарные счётчики нулевые, per-event поля визитов через каждую границу однородны (браузер, referer, utm). -- [ ] Учебный цикл целиком и повторяемо: next-day -> прогон `etl_pipeline` +- [x] Учебный цикл целиком и повторяемо: next-day -> прогон `etl_pipeline` -> в DM-витринах появились строки ровно за новый день (счётчики до/после, прирост только в диапазоне нового дня); проверено на двух next-day подряд. -- [ ] Регрессий нет: тесты генератора и контракты корня зелёные +- [x] Регрессий нет: тесты генератора и контракты корня зелёные (`make test`), сценарий `make generated-history-runtime-check` задачи 20 остаётся зелёным. -- [ ] Документация: OPERATIONS.md — глагол, предпроверка, поведение при +- [x] Документация: OPERATIONS.md — глагол, предпроверка, поведение при сбое до фиксации манифеста (восстановление переимпортом артефакта). ## Границы (что не трогать) @@ -158,3 +158,51 @@ state и закрыл тихие fallback'и. Прежняя редакция э «Доливка стартовой истории кусочком от слепка». Проверенные 2026-07-07 факты (import «всё или ничего», донор задачи 09) подтверждены обоими ревью 2026-07-12 по коду ветки. + +## Решения по слепым ревью реализации (2026-07-12) + +### Полная перечитка Kafka и конечный срок хранения — отклонено в задаче 13 + +Накопительные `visits` и `users` требуют точного объединения идентификаторов, +а текущая контрольная сумма — SHA-256 последовательности событий. Из одних +итоговых чисел и готовой контрольной суммы нельзя точно добавить новый день. Для +инкрементального расчёта пришлось бы хранить множества идентификаторов и новое +состояние hash в manifest. Это изменило бы формат артефакта сверх разрешённого +поля `boundaries` и нарушило бы решение 6. + +Поэтому в задаче 13 остаётся точная пересборка по доступной истории Kafka. +Её проверяет +`test_next_day_publishes_24h_then_state_then_cumulative_manifest`: два дня, +точные накопительные счётчики и суммы. Риск принят только для ручного учебного +цикла. Ограничение по времени, памяти и сроку хранения записано в OPERATIONS.md. +Задача про расписание из решения 8 до включения обязана выбрать бессрочное +хранение data-топиков или новый согласованный формат накопительного состояния. + +### Device/geo min/max timestamps после пересборки — отклонено + +Эти два диапазона являются диагностическими полями статистики топика и не +участвуют в `compare_clickhouse_stats_to_manifest`, проверке цепочки или +критериях приёмки. Контракт задачи требует накопительные строки и контрольные +суммы; они пересчитываются точно и проверяются тестом двух `next-day`. +Выравнивание +device/geo min/max с потактовым backfill потребовало бы сохранять связь каждой +повторной строки с исходным browser-событием, которой в этих сообщениях нет. +Менять формат ради неиспользуемых диагностических полей в задаче 13 не следует. + +### Нулевой стык — допустим с явным статусом + +Диагноз текущего стенда подтвердил естественный ночной провал, а не потерю +визитов при восстановлении. На границе `2026-01-02T01:00:00+00:00` DDS и STG +дали `crossing_visits=0`. За предыдущий час завершилось 27 визитов, за +следующий началось 24. Последний визит закончился в `00:56:57.766327`, новый +начался в `01:00:00`; пауза составила 183 секунды. State на границе содержал +`active_visits=0`, поэтому restore не мог отбросить активный визит. Для +контроля: на предыдущей границе state содержал один активный визит, и DDS/STG +нашли один переходящий визит. + +Поэтому нулевой стык проходит с явным статусом «однородность неприменима». +Пустая порция, хвост за последней границей и непарные строки по-прежнему дают +отказ. Остаточный слепой участок принят: однородность нельзя проверить без +переходящего визита. Restore закреплён тестом двух последовательных `next-day` +и детерминизмом задачи 09, а live-гейт задачи 20 отдельно проверяет стык, где +переходящий визит гарантирован сценарием. diff --git a/Makefile b/Makefile index 01e0384..c5db886 100644 --- a/Makefile +++ b/Makefile @@ -1,5 +1,5 @@ .PHONY: up down clean ddl data transform logs \ - generated-history-analytics generated-history-check generated-history-runtime-check \ + generated-history-analytics generated-history-check generated-history-runtime-check generated-history-chain-check \ startup-history-export startup-history-import startup-history-check \ reload-monitoring recover-monitoring \ superset-init superset-dashboard superset-ui superset-restart \ @@ -57,6 +57,10 @@ generated-history-check: generated-history-runtime-check: COMPOSE_BIN="$(COMPOSE)" bash ./scripts/run_generated_history_runtime_check.sh +# Проверить все завершённые стыки в накопительной цепочке модельных границ +generated-history-chain-check: + COMPOSE_BIN="$(COMPOSE)" bash ./scripts/check_generated_history_chain.sh + # Сгенерировать стартовую историю и сохранить портативный артефакт startup-history-export: COMPOSE_BIN="$(COMPOSE)" bash ./scripts/export_startup_history_artifact.sh diff --git a/airflow/dags/generator_control_dag.py b/airflow/dags/generator_control_dag.py index 519f3d5..38bc76f 100644 --- a/airflow/dags/generator_control_dag.py +++ b/airflow/dags/generator_control_dag.py @@ -23,13 +23,19 @@ from airflow_clickhouse_plugin.hooks.clickhouse import ClickHouseHook from clickstream_generator.airflow_control import ( KAFKA_BOOTSTRAP_SERVERS, assert_clickhouse_matches_manifest, + assert_expected_t_end, assert_live_generator_not_running, + assert_live_generator_not_running_for_next_day, + assert_next_day_snapshot, assert_stand_clean, build_control_env, + build_next_day_env, default_artifact_path, load_manifest_from_kafka, + load_state_from_kafka, run_backfill, run_import, + run_next_day, target_dag_trigger_error, validate_import_artifact, ) @@ -101,6 +107,8 @@ def choose_operation(**context) -> str: return "check_etl_not_paused_before_backfill" if operation == "import": return "check_etl_not_paused_before_import" + if operation == "next-day": + return "check_etl_not_paused_before_next_day" if operation == "check": return "check_only" raise ValueError(f"Неизвестная операция: {operation}") @@ -143,6 +151,31 @@ def run_import_task(**context) -> None: run_import(env, artifact_path) +def precheck_next_day(**context) -> None: + """Проверяет точку продолжения и готовит неизменные настройки мира.""" + assert_live_generator_not_running_for_next_day() + manifest = load_manifest_from_kafka() + state = load_state_from_kafka() + assert_next_day_snapshot(manifest, state) + assert_expected_t_end( + str(_param(context, "expected_t_end") or "").strip(), + str(manifest["model_t_end"]), + ) + context["ti"].xcom_push( + key="generator_env", + value=build_next_day_env(manifest), + ) + + +def run_next_day_task(**context) -> None: + """Генерирует следующий модельный день от проверенного state.""" + env = context["ti"].xcom_pull( + task_ids="precheck_next_day", + key="generator_env", + ) + run_next_day(env) + + def check_manifest_task(**context) -> None: """Сверяет витрину ClickHouse с manifest стартовой истории.""" manifest = load_manifest_from_kafka() @@ -164,6 +197,8 @@ def assert_target_dag_not_paused(dag_id: str, session=None) -> None: raise AirflowException(error) +# Context7, Airflow 2.10.5: Param поддерживает enum, а max_active_runs ограничивает +# число одновременных DAG run. Поэтому форму и блокировку next-day держим в DAG. with DAG( dag_id="generator_control", description="Пульт стартовой истории генератора", @@ -178,11 +213,11 @@ with DAG( "operation": Param( "backfill", type="string", - enum=["backfill", "import", "check"], + enum=["backfill", "import", "next-day", "check"], title="Операция", description=( - "Что сделать: создать историю, импортировать артефакт " - "или проверить витрины." + "Что сделать: создать историю, импортировать артефакт, " + "добавить следующий день или проверить витрины." ), ), "profile": Param( @@ -219,6 +254,15 @@ with DAG( "Import: что читать; пусто — путь по умолчанию в data." ), ), + "expected_t_end": Param( + "", + type="string", + title="Ожидаемая граница next-day", + description=( + "Необязательный model_t_end до запуска. Защищает от " + "повторной доливки того же дня." + ), + ), }, ) as dag: route = BranchPythonOperator( @@ -244,6 +288,15 @@ with DAG( python_callable=run_import_task, ) + precheck_next_day_task = PythonOperator( + task_id="precheck_next_day", + python_callable=precheck_next_day, + ) + next_day_task = PythonOperator( + task_id="run_next_day", + python_callable=run_next_day_task, + ) + check_etl_not_paused_before_backfill = PythonOperator( task_id="check_etl_not_paused_before_backfill", python_callable=assert_target_dag_not_paused, @@ -254,6 +307,11 @@ with DAG( python_callable=assert_target_dag_not_paused, op_kwargs={"dag_id": ETL_DAG_ID}, ) + check_etl_not_paused_before_next_day = PythonOperator( + task_id="check_etl_not_paused_before_next_day", + python_callable=assert_target_dag_not_paused, + op_kwargs={"dag_id": ETL_DAG_ID}, + ) trigger_etl = TriggerDagRunOperator( task_id="trigger_etl", @@ -284,11 +342,14 @@ with DAG( route >> [ check_etl_not_paused_before_backfill, check_etl_not_paused_before_import, + check_etl_not_paused_before_next_day, check_only, ] check_etl_not_paused_before_backfill >> precheck_backfill_task check_etl_not_paused_before_import >> precheck_import_task + check_etl_not_paused_before_next_day >> precheck_next_day_task precheck_backfill_task >> backfill_task >> trigger_etl precheck_import_task >> import_task >> trigger_etl + precheck_next_day_task >> next_day_task >> trigger_etl trigger_etl >> check_after_etl >> done check_only >> done diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 8ce706d..10e16f6 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -14,6 +14,8 @@ - `CHECK_LIVE_SEAM=0 make generated-history-check` (повторяемая проверка ClickHouse и Superset после прогона только стартовой истории) - `make generated-history-runtime-check` (короткая проверка стыка backfill/live) +- `make generated-history-chain-check` (проверка завершённых стыков между + порциями истории) - `make ddl` (применяет SQL из `sql/ddl/00_databases.sql` и `sql/ddl/*/*.sql` в ClickHouse) - `make data` (архивный путь: заливает `data/*.jsonl` в Kafka; не основной источник аналитики) - `make transform` (запускает batch-процесс ODS -> DDS -> DM) @@ -59,19 +61,55 @@ volumes или live-генератором. Для стыка backfill/live от дождаться `success`; - `import` — импортировать портативный артефакт, затем запустить `etl_pipeline` и дождаться `success`; + - `next-day` — восстановить мир из state, добавить 24 модельных часа, + затем запустить `etl_pipeline` и дождаться `success`; - `check` — сверить ClickHouse с manifest из Kafka. - Параметры: - - `operation` (`backfill` / `import` / `check`); + - `operation` (`backfill` / `import` / `next-day` / `check`); - `profile` — список берётся из `PROFILES` генератора; - `duration` — `6h`, `2d` и т.п.; пусто означает длительность профиля; - `seed`, `model_time_speed` — необязательные переопределения мира; - - `artifact_path` — для `backfill` путь сохранения, для `import` путь чтения. + - `artifact_path` — для `backfill` путь сохранения, для `import` путь чтения; + - `expected_t_end` — необязательная ожидаемая граница перед `next-day`. + При расхождении запуск показывает ожидаемое и фактическое значения. Backfill/import требуют чистый стенд: пустые data-топики Kafka и пустые `stg.*_raw`. При отказе очистите стенд через `make clean`. Операции `continue` в DAG нет: live-генератор — долгоживущий сервис, его запускают с консоли через `make generator-continue`. +`next-day` работает на непустом стенде и не использует проверку чистоты. +Перед записью пульт требует manifest, state ровно на его `T_end` и остановленный +live-генератор. Настройки мира берутся из manifest; поля `profile`, `duration`, +`seed` и `model_time_speed` формы для этой операции не применяются. Один запуск +добавляет полуоткрытый диапазон `[T_end, T_end + 24h)` в UTC. Новая граница +появляется в `boundaries`; старый manifest без поля читается как `[T0, T_end]`. +Расписание остаётся выключенным (`schedule=None`), а `max_active_runs=1` не даёт +двум доливкам выполняться параллельно. + +После двух доливок проверьте завершённые стыки: + +```bash +make generated-history-chain-check +``` + +Проверка проходит по внутренним границам `boundaries`, ищет непарные строки и +смену browser/referer/utm внутри переходящих визитов. Она не меняет +`make generated-history-runtime-check` для стыка backfill/live. + +Точка фиксации `next-day` — новый manifest. Порядок записи: data-топики, state, +manifest. Автоматического отката нет. Если запуск упал до публикации manifest, +не повторяйте доливку поверх возможного хвоста. Очистите стенд и переимпортируйте +последний исправный портативный артефакт, затем повторите `next-day`. + +Текущая версия пересчитывает накопительные счётчики и контрольные суммы по всей +доступной истории data-топиков Kafka. Поэтому время выполнения и расход памяти +каждой доливки растут вместе с историей. Стандартный срок хранения Kafka тоже +ограничивает долгую работу стенда. Для ручного учебного цикла запускайте +соседние дни без долгих пауз. Перед включением расписания отдельная задача +должна выбрать одно из решений: бессрочное хранение data-топиков или новый +формат накопительного состояния manifest. + Перед запуском `etl_pipeline` пульт проверяет, что DAG не стоит на паузе. Если стоит, задача падает сразу с подсказкой снять паузу в UI или командой: @@ -165,7 +203,7 @@ make generator-logs | `GEN_MODEL_T_END` | Правая граница стартовой истории для `backfill` | пусто | | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | -| `GEN_RUN_MODE` | Режим генератора | `live` | +| `GEN_RUN_MODE` | Режим: `live`, `backfill` или `next-day` | `live` | | `GEN_LAUNCH_PROFILE` | Имя профиля запуска для логов | `ci` | | `GEN_STARTUP_HISTORY_ARTIFACT` | JSON-файл для экспорта стартовой истории в режиме `backfill` | пусто | | `GEN_STATE_ENABLED` | Сохранять state v3 между рестартами | `true` | @@ -254,6 +292,10 @@ GEN_HISTORY_DURATION=2d make generated-history-analytics контрольными числами. При live-запуске с теми же настройками генератор видит, что state совпадает с manifest, и стартует ровно с `T_end` без настенной дельты. +Новый backfill записывает `boundaries=[T0, T_end]`. Каждый успешный `next-day` +добавляет одну границу, сдвигает `model_t_end` и пересчитывает накопительные +счётчики и контрольные суммы по всей истории. + Если историю нужно сохранить в файл и восстановить на чистом стенде без новой генерации, используйте [runbook стартовой истории](./runbooks/startup-history.md). diff --git a/generator/src/clickstream_generator/airflow_control.py b/generator/src/clickstream_generator/airflow_control.py index d127976..eacc21a 100644 --- a/generator/src/clickstream_generator/airflow_control.py +++ b/generator/src/clickstream_generator/airflow_control.py @@ -6,6 +6,7 @@ import os import urllib.error import urllib.request from contextlib import contextmanager +from datetime import datetime, timezone from pathlib import Path from typing import Iterator @@ -23,7 +24,9 @@ from clickstream_generator.startup_history_artifact import ( compare_clickhouse_stats_to_manifest, import_startup_history_artifact, load_startup_history_artifact, + manifest_boundaries, validate_startup_history_artifact, + validate_manifest_state, ) from clickstream_generator.stand_clean import ( STG_EMPTY_SQL, @@ -48,6 +51,20 @@ WORLD_OVERRIDE_KEYS = { "GEN_MAX_EVENTS_PER_TICK", } +NEXT_DAY_SETTING_ENV_KEYS = { + "tick_seconds": "GEN_TICK_SECONDS", + "lambda_base_per_min": "GEN_LAMBDA_BASE_PER_MIN", + "jitter_pct": "GEN_JITTER_PCT", + "min_events_per_tick": "GEN_MIN_EVENTS_PER_TICK", + "max_events_per_tick": "GEN_MAX_EVENTS_PER_TICK", + "max_session_events": "GEN_MAX_SESSION_EVENTS", + "max_active_sessions": "GEN_MAX_ACTIVE_SESSIONS", + "population_max": "GEN_POPULATION_MAX", + "p_new_user": "GEN_P_NEW_USER", + "min_return_minutes": "GEN_MIN_RETURN_MINUTES", + "model_time_speed": "GEN_MODEL_TIME_SPEED", +} + CLICKHOUSE_STATS_SQL = """ WITH toDateTime64('{model_t0}', 6) AS t0, @@ -96,6 +113,73 @@ def build_control_env( return env +def build_next_day_env(manifest: dict) -> dict[str, str]: + """Восстанавливает окружение следующего дня из manifest мира.""" + manifest_boundaries(manifest) + settings = manifest.get("generation_settings") or {} + missing = [key for key in NEXT_DAY_SETTING_ENV_KEYS if key not in settings] + if missing: + raise ValueError( + "в manifest не хватает generation_settings: " + + ", ".join(sorted(missing)) + ) + overrides = { + "GEN_SEED": str(manifest["gen_seed"]), + "GEN_MODEL_T0": str(manifest["model_t0"]), + "GEN_MODEL_T_END": str(manifest["model_t_end"]), + "GEN_MODEL_TIMEZONE": str(manifest["model_timezone"]), + } + env = build_launch_env( + "next-day", + profile_name=str(manifest.get("launch_profile") or "ci"), + overrides=overrides, + ) + for setting, env_key in NEXT_DAY_SETTING_ENV_KEYS.items(): + env[env_key] = str(settings[setting]) + env["KAFKA_BOOTSTRAP_SERVERS"] = KAFKA_BOOTSTRAP_SERVERS + env["GEN_DATA_DIR"] = AIRFLOW_DATA_DIR + env["GEN_METRICS_ENABLED"] = "false" + return env + + +def assert_expected_t_end(expected_t_end: str, actual_t_end: str) -> None: + """Отвергает повторный запуск от уже сдвинутой границы.""" + if not expected_t_end: + return + try: + expected = _parse_utc_timestamp(expected_t_end) + actual = _parse_utc_timestamp(actual_t_end) + except ValueError as exc: + raise RuntimeError(f"Некорректная граница next-day: {exc}") from exc + if expected != actual: + raise RuntimeError( + "Граница next-day изменилась: " + f"ожидалась {expected_t_end}, фактическая {actual_t_end}." + ) + + +def assert_next_day_snapshot(manifest: dict | None, state) -> None: + """Проверяет наличие и согласованность точки продолжения next-day.""" + if not manifest: + raise RuntimeError("Manifest стартовой истории не найден в Kafka.") + if state is None: + raise RuntimeError("State генератора не найден в Kafka.") + try: + manifest_boundaries(manifest) + validate_manifest_state(manifest, state) + except (KeyError, TypeError, ValueError) as exc: + raise RuntimeError( + f"state не согласован с T_end manifest: {exc}" + ) from exc + + +def _parse_utc_timestamp(value: str): + timestamp = datetime.fromisoformat(value.replace("Z", "+00:00")) + if timestamp.tzinfo is None: + raise ValueError("время должно содержать часовой пояс") + return timestamp.astimezone(timezone.utc) + + def assert_stand_clean(kafka_bootstrap_servers: str, clickhouse_hook) -> None: """Проверяет, что backfill/import не смешает миры.""" assert_kafka_data_topics_empty(kafka_bootstrap_servers) @@ -111,14 +195,31 @@ def assert_stg_tables_empty(clickhouse_hook) -> None: def assert_live_generator_not_running(url: str = GENERATOR_METRICS_URL) -> None: """Мягко предупреждает о live-сервисе по HTTP-метрикам.""" + _assert_metrics_not_responding( + url, + "Live-генератор отвечает на metrics-порту. Остановите его с консоли " + "перед backfill/import, чтобы не смешать миры.", + ) + + +def assert_live_generator_not_running_for_next_day( + url: str = GENERATOR_METRICS_URL, +) -> None: + """Падает, если next-day может смешаться с работающим live-потоком.""" + _assert_metrics_not_responding( + url, + "Live-генератор отвечает на metrics-порту. Остановите его с консоли " + "перед next-day, чтобы не смешать дни.", + ) + + +def _assert_metrics_not_responding(url: str, error_message: str) -> None: + """Проверяет, что HTTP-метрики live-генератора недоступны.""" try: urllib.request.urlopen(url, timeout=2).close() except (urllib.error.URLError, TimeoutError, OSError): return - raise RuntimeError( - "Live-генератор отвечает на metrics-порту. Остановите его с консоли " - "перед backfill/import, чтобы не смешать миры." - ) + raise RuntimeError(error_message) def run_backfill(env: dict[str, str]) -> None: @@ -128,6 +229,13 @@ def run_backfill(env: dict[str, str]) -> None: GeneratorService(config).start() +def run_next_day(env: dict[str, str]) -> None: + """Запускает ограниченную доливку следующего дня в Airflow worker.""" + with _patched_environ(env): + config = Config() + GeneratorService(config).start() + + def run_import(env: dict[str, str], artifact_path: str) -> dict: """Импортирует портативный артефакт в Kafka.""" with _patched_environ(env): @@ -175,6 +283,17 @@ def load_manifest_from_kafka( return manifest +def load_state_from_kafka( + kafka_bootstrap_servers: str = KAFKA_BOOTSTRAP_SERVERS, +): + """Читает state генератора из compact-топика.""" + manager = KafkaStateManager(kafka_bootstrap_servers) + try: + return manager.load() + finally: + manager.close() + + def assert_clickhouse_matches_manifest(manifest: dict, clickhouse_hook) -> None: """Сверяет контрольные числа ClickHouse с manifest.""" model_t0 = _clickhouse_datetime_literal(str(manifest["model_t0"])) diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index d7ac0dc..e870370 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -137,8 +137,8 @@ class Config: raise ValueError(f"Unknown GEN_MODEL_TIMEZONE: {self.model_timezone}") from e if self.model_time_speed <= 0: raise ValueError("GEN_MODEL_TIME_SPEED must be > 0") - if self.run_mode not in {"live", "backfill"}: - raise ValueError("GEN_RUN_MODE must be live or backfill") + if self.run_mode not in {"live", "backfill", "next-day"}: + raise ValueError("GEN_RUN_MODE должен быть live, backfill или next-day") if self.model_t_end is not None: object.__setattr__( self, @@ -150,6 +150,11 @@ class Config: raise ValueError("GEN_MODEL_T_END is required for backfill") if self.model_t_end <= self.model_t0: raise ValueError("GEN_MODEL_T_END must be after GEN_MODEL_T0") + if self.run_mode == "next-day": + if self.model_t_end is None: + raise ValueError("для next-day требуется GEN_MODEL_T_END") + if self.model_t_end <= self.model_t0: + raise ValueError("GEN_MODEL_T_END должен быть позже GEN_MODEL_T0") def _parse_optional_model_timestamp(value: str | None) -> datetime | None: diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py index 4da2ca6..de77cb6 100644 --- a/generator/src/clickstream_generator/kafka_io.py +++ b/generator/src/clickstream_generator/kafka_io.py @@ -13,6 +13,7 @@ from clickstream_generator.state import UnsupportedStateVersionError logger = logging.getLogger("generator") +DATA_TOPICS = ("browser_events", "location_events", "device_events", "geo_events") _kafka_imported = False KafkaProducer = None @@ -280,6 +281,63 @@ class KafkaStartupHistoryManifest: return None +class KafkaDataTopicReader: + """Читает устойчивый снимок data-топиков для накопительного manifest.""" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + + def load(self) -> dict[str, list[dict]]: + """Возвращает все события data-топиков после остановки live-потока.""" + from kafka import KafkaConsumer, TopicPartition + + consumer = KafkaConsumer( + bootstrap_servers=self.bootstrap_servers, + enable_auto_commit=False, + value_deserializer=lambda value: json.loads(value.decode("utf-8")), + ) + topics = {topic: [] for topic in DATA_TOPICS} + try: + partitions = [] + for topic in DATA_TOPICS: + topic_partitions = consumer.partitions_for_topic(topic) or set() + partitions.extend( + TopicPartition(topic, partition) + for partition in topic_partitions + ) + if not partitions: + return topics + + consumer.assign(partitions) + consumer.seek_to_beginning(*partitions) + end_offsets = consumer.end_offsets(partitions) + empty_polls = 0 + while any( + consumer.position(partition) < end_offsets[partition] + for partition in partitions + ): + records = consumer.poll(timeout_ms=1000) + if not records: + empty_polls += 1 + if empty_polls >= 5: + raise RuntimeError( + "не удалось дочитать data-топики до зафиксированных " + "конечных смещений" + ) + continue + empty_polls = 0 + for partition, messages in records.items(): + for message in messages: + if ( + message.offset < end_offsets[partition] + and message.value is not None + ): + topics[message.topic].append(message.value) + finally: + consumer.close() + return topics + + def ensure_topics(bootstrap_servers: str) -> None: """Создаёт служебные топики, если их ещё нет.""" from kafka import KafkaAdminClient diff --git a/generator/src/clickstream_generator/launch.py b/generator/src/clickstream_generator/launch.py index 7aa5719..7659a9a 100644 --- a/generator/src/clickstream_generator/launch.py +++ b/generator/src/clickstream_generator/launch.py @@ -93,8 +93,8 @@ def build_launch_env( overrides: dict[str, str] | None = None, ) -> dict[str, str]: """Возвращает старые env-переменные для нового глагола запуска.""" - if verb not in {"backfill", "continue", "reset"}: - raise ValueError("verb must be backfill, continue or reset") + if verb not in {"backfill", "continue", "next-day", "reset"}: + raise ValueError("глагол должен быть backfill, continue, next-day или reset") if profile_name not in PROFILES: known = ", ".join(sorted(PROFILES)) raise ValueError(f"Unknown launch profile {profile_name!r}; known: {known}") @@ -124,6 +124,13 @@ def build_launch_env( env["GEN_STATE_RESET"] = "false" env["GEN_HISTORY_DURATION"] = selected_duration env["GEN_MODEL_T_END"] = _model_t_end(env["GEN_MODEL_T0"], selected_duration) + elif verb == "next-day": + current_t_end = overrides.get("GEN_MODEL_T_END") + if not current_t_end: + raise ValueError("для next-day требуется GEN_MODEL_T_END") + env["GEN_RUN_MODE"] = "next-day" + env["GEN_STATE_RESET"] = "false" + env["GEN_MODEL_T_END"] = current_t_end else: env["GEN_RUN_MODE"] = "live" env["GEN_STATE_RESET"] = "true" @@ -148,7 +155,10 @@ def _model_t_end(model_t0: str, duration: str) -> str: def _parse_args(argv: list[str]) -> argparse.Namespace: parser = argparse.ArgumentParser(description="Глаголы запуска генератора") - parser.add_argument("verb", choices=("backfill", "continue", "reset")) + parser.add_argument( + "verb", + choices=("backfill", "continue", "next-day", "reset"), + ) parser.add_argument( "--profile", default=os.getenv("PROFILE") or os.getenv("GEN_LAUNCH_PROFILE") or "ci", diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index e30d5eb..2755897 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -16,6 +16,7 @@ from clickstream_generator.generation import EventGenerator from clickstream_generator.kafka_io import ( BatchRecord, KafkaBatchHistory, + KafkaDataTopicReader, KafkaPublisher, KafkaStateManager, KafkaStartupHistoryManifest, @@ -29,8 +30,10 @@ from clickstream_generator.metrics import ( from clickstream_generator.runtime import TickStreamGenerator from clickstream_generator.startup_history_artifact import ( StartupHistoryArtifactBuilder, + ManifestCounters, build_manifest, generation_settings_from_config, + manifest_boundaries, write_startup_history_artifact, ) from clickstream_generator.state import UnsupportedStateVersionError @@ -55,6 +58,7 @@ 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() @@ -89,6 +93,8 @@ class GeneratorService: if self.config.run_mode == "backfill" and not self.config.state_enabled: raise ValueError("GEN_STATE_ENABLED must be true for backfill") + if self.config.run_mode == "next-day" and not self.config.state_enabled: + raise ValueError("для next-day GEN_STATE_ENABLED должен быть true") if self.config.state_enabled: self.state_manager = KafkaStateManager(self.config.kafka_bootstrap_servers) @@ -102,6 +108,38 @@ class GeneratorService: self.stop() return + if self.config.run_mode == "next-day": + self.manifest_manager = KafkaStartupHistoryManifest( + self.config.kafka_bootstrap_servers + ) + restored_state = self.state_manager.load() + manifest = self.manifest_manager.load() + if restored_state is None: + raise IncompatibleStateError( + "state генератора для next-day не найден" + ) + mismatches = self._startup_history_mismatch_fields( + restored_state, + manifest, + ) + if mismatches: + raise IncompatibleStateError( + "слепок next-day не согласован: " + ", ".join(mismatches) + ) + model_t_end = self._as_aware_utc( + datetime.fromisoformat(manifest["model_t_end"]) + ) + self.restore_from_startup_history( + restored_state, + model_t_end=model_t_end, + ) + self.data_reader = KafkaDataTopicReader( + self.config.kafka_bootstrap_servers + ) + self._run_next_day(manifest) + self.stop() + return + if not self.config.state_reset: try: restored_state = self.state_manager.load() @@ -481,6 +519,117 @@ class GeneratorService: f"status={status}, sent_counts={sent_counts}" ) + def _run_next_day(self, manifest: dict) -> None: + """Доливает ровно 24 модельных часа от границы manifest.""" + if not self.publisher: + raise RuntimeError("publisher не инициализирован") + if not self.history: + raise RuntimeError("история batch не инициализирована") + if not self.state_manager or not self.manifest_manager: + raise RuntimeError("менеджеры state и manifest не инициализированы") + + current_t_end = self._as_aware_utc( + datetime.fromisoformat(manifest["model_t_end"]) + ) + if self._model_time != current_t_end: + raise IncompatibleStateError( + "восстановленное время next-day не совпадает с T_end manifest" + ) + target_t_end = current_t_end + timedelta(hours=24) + logger.info( + "Запуск next-day от %s до %s", + current_t_end.isoformat(), + target_t_end.isoformat(), + ) + + while self._model_time < target_t_end: + self._tick += 1 + batch_id = f"next-day-{self._tick:08d}" + model_time = self._model_time + started_at = datetime.now(timezone.utc) + events_count = self.generator._calculate_events_count(now=model_time) + batch = self.stream.generate_tick( + events_count, + tick_started_at=model_time, + ) + total_sent, sent_counts, status = self._publish_batch(batch) + self._raise_on_next_day_publish_error(batch_id, status, sent_counts) + self._write_batch_history( + batch_id=batch_id, + started_at=started_at, + sent_counts=sent_counts, + total_sent=total_sent, + status=status, + error_message=None, + ) + self._advance_model_time() + + final_batch = self.stream.drain_until( + target_t_end, + include_boundary=False, + ) + if any(final_batch.values()): + batch_id = f"next-day-{self._tick + 1:08d}-final" + 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) + self._write_batch_history( + batch_id=batch_id, + started_at=started_at, + sent_counts=sent_counts, + total_sent=total_sent, + status=status, + error_message=None, + ) + + 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, + rng_state=self.generator.rng.getstate(), + last_batch_id=state_batch_id, + last_timestamp=target_t_end, + model_timestamp=target_t_end, + wall_timestamp=datetime.now(timezone.utc), + model_time_speed=self.config.model_time_speed, + model_timezone=self.config.model_timezone, + model_t0=self.config.model_t0, + gen_seed=self.config.seed, + ) + boundaries = manifest_boundaries(manifest) + [target_t_end.isoformat()] + updated_manifest = build_manifest( + self.config, + counters, + state, + model_t_end=target_t_end, + boundaries=boundaries, + ) + + self.state_manager.save(state) + self.state_manager.flush() + self.manifest_manager.save(updated_manifest) + self.manifest_manager.flush() + + def _raise_on_next_day_publish_error( + self, + batch_id: str, + status: str, + sent_counts: dict[str, dict[str, int]], + ) -> None: + """Останавливает next-day до записи state и manifest при ошибке Kafka.""" + if status == "success": + return + raise RuntimeError( + f"Публикация next-day не удалась для batch {batch_id}: " + f"status={status}, sent_counts={sent_counts}" + ) + def _publish_batch( self, batch: dict[str, list[dict]], diff --git a/generator/src/clickstream_generator/startup_history_artifact.py b/generator/src/clickstream_generator/startup_history_artifact.py index db68a56..e77f9e6 100644 --- a/generator/src/clickstream_generator/startup_history_artifact.py +++ b/generator/src/clickstream_generator/startup_history_artifact.py @@ -287,16 +287,30 @@ def generation_settings_from_config(config: Config) -> dict: } -def build_manifest(config: Config, counters: ManifestCounters, state: GeneratorState) -> dict: +def build_manifest( + config: Config, + counters: ManifestCounters, + state: GeneratorState, + *, + model_t_end: datetime | None = None, + boundaries: list[str] | None = None, +) -> dict: """Строит manifest стартовой истории.""" - if config.model_t_end is None: + selected_t_end = model_t_end or config.model_t_end + if selected_t_end is None: raise ValueError("GEN_MODEL_T_END is required for startup history manifest") - return { + selected_boundaries = ( + boundaries + if boundaries is not None + else [config.model_t0.isoformat(), selected_t_end.isoformat()] + ) + manifest = { "manifest_version": "1.0", "generated_at": datetime.now(timezone.utc).isoformat(), "gen_seed": config.seed, "model_t0": config.model_t0.isoformat(), - "model_t_end": config.model_t_end.isoformat(), + "model_t_end": selected_t_end.isoformat(), + "boundaries": list(selected_boundaries), "model_timezone": config.model_timezone, "run_mode": "backfill", "launch_profile": config.launch_profile, @@ -311,6 +325,28 @@ def build_manifest(config: Config, counters: ManifestCounters, state: GeneratorS "topics": counters.to_manifest_topics(), "totals": counters.to_manifest_totals(), } + manifest_boundaries(manifest) + return manifest + + +def manifest_boundaries(manifest: dict) -> list[str]: + """Возвращает проверенную цепочку границ, включая старый формат.""" + model_t0 = _parse_manifest_timestamp(manifest.get("model_t0")) + model_t_end = _parse_manifest_timestamp(manifest.get("model_t_end")) + raw_boundaries = manifest.get("boundaries") + if raw_boundaries is None: + return [manifest["model_t0"], manifest["model_t_end"]] + if not isinstance(raw_boundaries, list) or len(raw_boundaries) < 2: + raise ValueError("manifest boundaries должен содержать T0 и T_end") + + parsed = [_parse_manifest_timestamp(value) for value in raw_boundaries] + if parsed[0] != model_t0: + raise ValueError("manifest boundaries не начинается с model_t0") + if parsed[-1] != model_t_end: + raise ValueError("manifest boundaries не заканчивается на model_t_end") + if any(left >= right for left, right in zip(parsed, parsed[1:])): + raise ValueError("manifest boundaries должен строго возрастать") + return list(raw_boundaries) def write_startup_history_artifact(path: str | Path, artifact: dict) -> None: @@ -359,6 +395,7 @@ def validate_startup_history_artifact( raise ValueError(f"artifact raw topic {topic} must be a list") state = GeneratorState.from_dict(state_payload) + manifest_boundaries(manifest) _validate_manifest_state(manifest, state) _validate_manifest_topics(manifest, topics) _validate_raw_topics(topics, raw_topics) @@ -487,6 +524,11 @@ def _validate_manifest_state(manifest: dict, state: GeneratorState) -> None: raise ValueError("startup history artifact mismatch: " + ", ".join(mismatches)) +def validate_manifest_state(manifest: dict, state: GeneratorState) -> None: + """Проверяет согласованность manifest и state без данных артефакта.""" + _validate_manifest_state(manifest, state) + + def _validate_manifest_topics(manifest: dict, topics: dict[str, list[dict]]) -> None: counters = ManifestCounters() counters.add_batch({topic: topics[topic] for topic in TOPICS}) @@ -652,6 +694,34 @@ def _run_kafka_manifest_summary(args: argparse.Namespace) -> int: return 0 +def _run_kafka_boundaries(args: argparse.Namespace) -> int: + config = Config() + manifest_manager = KafkaStartupHistoryManifest(config.kafka_bootstrap_servers) + try: + manifest = manifest_manager.load() + finally: + manifest_manager.close() + if not manifest: + raise RuntimeError("manifest стартовой истории не найден в Kafka") + for boundary in manifest_boundaries(manifest): + print(boundary) + return 0 + + +def _run_kafka_topic_rows(args: argparse.Namespace) -> int: + config = Config() + manifest_manager = KafkaStartupHistoryManifest(config.kafka_bootstrap_servers) + try: + manifest = manifest_manager.load() + finally: + manifest_manager.close() + if not manifest: + raise RuntimeError("manifest стартовой истории не найден в Kafka") + topics = manifest.get("topics") or {} + print("\t".join(str(topics[topic]["rows"]) for topic in TOPICS)) + return 0 + + def main(argv: list[str] | None = None) -> int: parser = argparse.ArgumentParser(description="Startup history artifact tools") subparsers = parser.add_subparsers(dest="command", required=True) @@ -671,6 +741,12 @@ def main(argv: list[str] | None = None) -> int: kafka_summary_parser = subparsers.add_parser("kafka-manifest-summary") kafka_summary_parser.set_defaults(func=_run_kafka_manifest_summary) + kafka_boundaries_parser = subparsers.add_parser("kafka-boundaries") + kafka_boundaries_parser.set_defaults(func=_run_kafka_boundaries) + + kafka_topic_rows_parser = subparsers.add_parser("kafka-topic-rows") + kafka_topic_rows_parser.set_defaults(func=_run_kafka_topic_rows) + args = parser.parse_args(argv) return args.func(args) diff --git a/generator/tests/test_airflow_control.py b/generator/tests/test_airflow_control.py index 994496d..931f926 100644 --- a/generator/tests/test_airflow_control.py +++ b/generator/tests/test_airflow_control.py @@ -50,6 +50,94 @@ def test_import_env_requires_artifact_path_and_uses_backfill_contract(): assert env["GEN_HISTORY_DURATION"] == "2d" +def test_next_day_env_uses_world_settings_and_current_manifest_boundary(): + """Next-day восстанавливает настройки мира из manifest, а не из формы DAG.""" + from clickstream_generator.airflow_control import build_next_day_env + + manifest = { + "gen_seed": 42, + "model_t0": "2026-01-01T00:00:00+00:00", + "model_t_end": "2026-01-03T00:00:00+00:00", + "model_timezone": "UTC", + "launch_profile": "daily-wave", + "generation_settings": { + "tick_seconds": 1, + "lambda_base_per_min": 60, + "jitter_pct": 0, + "min_events_per_tick": 1, + "max_events_per_tick": 1000, + "max_session_events": 30, + "max_active_sessions": 200, + "population_max": 300, + "p_new_user": 0.15, + "min_return_minutes": 30, + "model_time_speed": 60, + }, + } + + env = build_next_day_env(manifest) + + assert env["GEN_RUN_MODE"] == "next-day" + assert env["GEN_STATE_RESET"] == "false" + assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00" + assert env["GEN_MODEL_TIME_SPEED"] == "60" + assert env["GEN_MAX_SESSION_EVENTS"] == "30" + assert env["GEN_DATA_DIR"] == "/opt/airflow/data" + + +def test_next_day_env_map_covers_every_generation_setting(base_config): + """Новая настройка генерации не может потеряться между manifest и env.""" + from clickstream_generator.airflow_control import NEXT_DAY_SETTING_ENV_KEYS + from clickstream_generator.startup_history_artifact import ( + generation_settings_from_config, + ) + + assert set(NEXT_DAY_SETTING_ENV_KEYS) == set( + generation_settings_from_config(base_config) + ) + + +def test_next_day_precheck_rejects_stale_expected_boundary_with_both_values(): + """Повторный запуск с прежней границей падает до записи данных.""" + from clickstream_generator.airflow_control import assert_expected_t_end + + with pytest.raises(RuntimeError) as exc_info: + assert_expected_t_end( + "2026-01-03T00:00:00+00:00", + "2026-01-04T00:00:00+00:00", + ) + + message = str(exc_info.value) + assert "2026-01-03T00:00:00+00:00" in message + assert "2026-01-04T00:00:00+00:00" in message + + +def test_next_day_precheck_requires_manifest_and_matching_state(): + """Next-day громко отвергает пустой стенд и state не на T_end.""" + from clickstream_generator.airflow_control import assert_next_day_snapshot + from test_startup_history_artifact import _state + + state = _state() + with pytest.raises(RuntimeError, match="Manifest"): + assert_next_day_snapshot(None, state) + + manifest = { + "run_mode": "backfill", + "model_t0": state.model_t0.isoformat(), + "model_t_end": "2026-01-01T02:00:00+00:00", + "gen_seed": state.gen_seed, + "model_timezone": state.model_timezone, + "state_version": state.version, + "state": { + "last_batch_id": state.last_batch_id, + "model_timestamp": "2026-01-01T02:00:00+00:00", + }, + "generation_settings": {"model_time_speed": state.model_time_speed}, + } + with pytest.raises(RuntimeError, match="state.*T_end"): + assert_next_day_snapshot(manifest, state) + + def test_target_dag_trigger_error_covers_missing_paused_and_ready_states(): """Проверка зависимого DAG различает три реальные ветки.""" from clickstream_generator.airflow_control import target_dag_trigger_error diff --git a/generator/tests/test_generator_control_dag_contract.py b/generator/tests/test_generator_control_dag_contract.py index f746785..916c4b4 100644 --- a/generator/tests/test_generator_control_dag_contract.py +++ b/generator/tests/test_generator_control_dag_contract.py @@ -30,14 +30,15 @@ def test_trigger_form_has_expected_param_enums(): """Форма запуска ограничивает операции и профили.""" text = DAG_PATH.read_text(encoding="utf-8") - assert 'enum=["backfill", "import", "check"]' in text + assert 'enum=["backfill", "import", "next-day", "check"]' in text assert "enum=sorted(PROFILES)" in text assert '"duration": Param(' in text assert '"artifact_path": Param(' in text + assert '"expected_t_end": Param(' in text def test_dag_branches_and_waits_for_etl_completion(): - """Backfill/import запускают ETL и ждут его завершения перед check.""" + """Backfill/import/next-day запускают ETL и ждут завершения перед check.""" text = DAG_PATH.read_text(encoding="utf-8") assert "BranchPythonOperator" in text @@ -48,6 +49,23 @@ def test_dag_branches_and_waits_for_etl_completion(): assert 'failed_states=["failed"]' in text +def test_next_day_has_own_boundary_precheck_and_serial_execution(): + """Next-day не использует clean-guard и не допускает параллельных запусков.""" + text = DAG_PATH.read_text(encoding="utf-8") + + assert "schedule=None" in text + assert "max_active_runs=1" in text + assert 'return "check_etl_not_paused_before_next_day"' in text + assert "assert_next_day_snapshot" in text + assert "assert_expected_t_end" in text + next_day_precheck = text.split("def precheck_next_day", maxsplit=1)[1].split( + "\ndef ", maxsplit=1 + )[0] + assert "assert_stand_clean" not in next_day_precheck + assert "assert_live_generator_not_running" in next_day_precheck + assert "run_next_day" in text + + def test_generator_control_prechecks_etl_dag_not_paused_before_waiting(): """Пульт проверяет паузу etl_pipeline до долгого ожидания.""" text = DAG_PATH.read_text(encoding="utf-8") @@ -62,8 +80,10 @@ def test_generator_control_prechecks_etl_dag_not_paused_before_waiting(): assert "fail_when_dag_is_paused" in text assert 'return "check_etl_not_paused_before_backfill"' in text assert 'return "check_etl_not_paused_before_import"' in text + assert 'return "check_etl_not_paused_before_next_day"' in text assert "check_etl_not_paused_before_backfill >> precheck_backfill_task" in text assert "check_etl_not_paused_before_import >> precheck_import_task" in text + assert "check_etl_not_paused_before_next_day >> precheck_next_day_task" in text assert text.index("check_etl_not_paused_before_backfill >> precheck_backfill_task") < text.index( "precheck_backfill_task >> backfill_task" ) diff --git a/generator/tests/test_launch.py b/generator/tests/test_launch.py index cf33b7a..fa7660e 100644 --- a/generator/tests/test_launch.py +++ b/generator/tests/test_launch.py @@ -43,7 +43,11 @@ def test_continue_and_reset_verbs_map_to_existing_low_level_flags(): from clickstream_generator.launch import build_launch_env continue_env = build_launch_env("continue", profile_name="ci") - reset_env = build_launch_env("reset", profile_name="ci") + reset_env = build_launch_env( + "reset", + profile_name="ci", + overrides={"GEN_MODEL_T_END": "2026-01-09T00:00:00+00:00"}, + ) assert continue_env["GEN_RUN_MODE"] == "live" assert continue_env["GEN_STATE_RESET"] == "false" @@ -53,6 +57,24 @@ def test_continue_and_reset_verbs_map_to_existing_low_level_flags(): assert "GEN_MODEL_T_END" not in reset_env +def test_next_day_verb_restores_snapshot_without_reset(): + """Глагол next-day включает отдельный ограниченный режим от текущей границы.""" + from clickstream_generator.launch import build_launch_env + + env = build_launch_env( + "next-day", + profile_name="ci", + overrides={ + "GEN_MODEL_T0": "2026-01-01T00:00:00+00:00", + "GEN_MODEL_T_END": "2026-01-03T00:00:00+00:00", + }, + ) + + assert env["GEN_RUN_MODE"] == "next-day" + assert env["GEN_STATE_RESET"] == "false" + assert env["GEN_MODEL_T_END"] == "2026-01-03T00:00:00+00:00" + + def test_ci_profile_keeps_previous_default_backfill_window(): """CI-профиль сохраняет прежний 6-часовой проверочный запуск.""" from clickstream_generator.launch import build_launch_env diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index c391838..4eda6f4 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -1005,6 +1005,185 @@ class TestGeneratorServiceBackfill: } +class TestGeneratorServiceNextDay: + """Проверки ограниченной доливки следующего модельного дня.""" + + def test_next_day_publishes_24h_then_state_then_cumulative_manifest( + self, base_config + ): + """Next-day пишет [T_end, T_end+24h), затем state и manifest.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + ) + + model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) + current_t_end = model_t0 + timedelta(hours=1) + config = replace( + base_config, + run_mode="next-day", + model_t0=model_t0, + model_t_end=current_t_end, + tick_seconds=3600, + model_time_speed=1, + lambda_base_per_min=60, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + ) + service = GeneratorService(config) + state = service.stream.to_state( + tick=1, + rng_state=service.generator.rng.getstate(), + last_batch_id="startup-history-initial", + last_timestamp=current_t_end, + model_timestamp=current_t_end, + wall_timestamp=current_t_end, + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed, + ) + old_batch = { + "browser_events": [{ + "event_id": "old-event", + "click_id": "old-click", + "event_timestamp": "2026-01-01 00:30:00.000000", + }], + "location_events": [{"event_id": "old-event"}], + "device_events": [{ + "click_id": "old-click", + "user_domain_id": "old-user", + }], + "geo_events": [{"click_id": "old-click"}], + } + initial_builder = StartupHistoryArtifactBuilder() + initial_builder.add_batch(old_batch) + manifest = build_manifest(config, initial_builder.counters, state) + + published = {topic: [] for topic in old_batch} + service.publisher = MagicMock() + + def publish(topic, events): + published[topic].extend(events) + return len(events), 0 + + service.publisher.publish.side_effect = publish + service.history = MagicMock() + order = [] + service.publisher.flush.side_effect = lambda: order.append("data") + service.state_manager = MagicMock() + service.state_manager.save.side_effect = lambda _state: order.append("state") + service.manifest_manager = MagicMock() + service.manifest_manager.save.side_effect = ( + lambda _manifest: order.append("manifest") + ) + + class DataReader: + def load(self): + return { + topic: old_batch[topic] + published[topic] + for topic in old_batch + } + + service.data_reader = DataReader() + service.restore_from_startup_history(state, model_t_end=current_t_end) + + service._run_next_day(manifest) + + browser_events = published["browser_events"] + timestamps = [ + datetime.fromisoformat(event["event_timestamp"].replace(" ", "T")) + for event in browser_events + ] + target_t_end = current_t_end + timedelta(hours=24) + saved_state = service.state_manager.save.call_args.args[0] + saved_manifest = service.manifest_manager.save.call_args.args[0] + + assert browser_events + assert min(timestamps) >= current_t_end.replace(tzinfo=None) + assert max(timestamps) < target_t_end.replace(tzinfo=None) + assert saved_state.model_timestamp == target_t_end + assert saved_manifest["model_t_end"] == target_t_end.isoformat() + assert saved_manifest["boundaries"] == [ + model_t0.isoformat(), + current_t_end.isoformat(), + target_t_end.isoformat(), + ] + assert saved_manifest["totals"]["events"] == len(browser_events) + 1 + assert order == ["data", "state", "manifest"] + + first_day_events = len(browser_events) + first_day_checksum = saved_manifest["topics"]["browser_events"][ + "checksum_sha256" + ] + service._run_next_day(saved_manifest) + + second_state = service.state_manager.save.call_args.args[0] + second_manifest = service.manifest_manager.save.call_args.args[0] + expected_builder = StartupHistoryArtifactBuilder() + expected_builder.add_batch(DataReader().load()) + assert second_state.model_timestamp == target_t_end + timedelta(hours=24) + assert second_manifest["boundaries"] == [ + model_t0.isoformat(), + current_t_end.isoformat(), + target_t_end.isoformat(), + (target_t_end + timedelta(hours=24)).isoformat(), + ] + assert second_manifest["totals"]["events"] == len(browser_events) + 1 + assert len(browser_events) > first_day_events + assert ( + second_manifest["topics"]["browser_events"]["checksum_sha256"] + != first_day_checksum + ) + assert second_manifest["topics"] == expected_builder.counters.to_manifest_topics() + assert order == [ + "data", + "state", + "manifest", + "data", + "state", + "manifest", + ] + + def test_next_day_publish_error_does_not_move_state_or_manifest(self, base_config): + """Ошибка data-топика оставляет обе точки фиксации без изменений.""" + model_t0 = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) + config = replace( + base_config, + run_mode="next-day", + model_t0=model_t0, + model_t_end=model_t0 + timedelta(hours=1), + tick_seconds=3600, + model_time_speed=1, + max_active_sessions=250, + population_max=251, + ) + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = lambda topic, events: ( + (len(events), 1) if topic == "location_events" else (len(events), 0) + ) + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + service._model_time = config.model_t_end + manifest = { + "model_t0": model_t0.isoformat(), + "model_t_end": config.model_t_end.isoformat(), + "boundaries": [model_t0.isoformat(), config.model_t_end.isoformat()], + } + + with pytest.raises(RuntimeError, match="Публикация next-day не удалась"): + service._run_next_day(manifest) + + service.state_manager.save.assert_not_called() + service.manifest_manager.save.assert_not_called() + + class TestGeneratorServiceState: """Тесты подключения state к сервисному запуску.""" diff --git a/generator/tests/test_startup_history_artifact.py b/generator/tests/test_startup_history_artifact.py index c966c96..c90d511 100644 --- a/generator/tests/test_startup_history_artifact.py +++ b/generator/tests/test_startup_history_artifact.py @@ -167,6 +167,92 @@ def test_manifest_and_artifact_show_launch_profile(base_config): path.unlink(missing_ok=True) +def test_manifest_records_initial_boundary_chain(base_config): + """Новый manifest явно хранит начало и правую границу истории.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + manifest_boundaries, + ) + + state = _state() + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + manifest = build_manifest( + config=replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ), + counters=builder.counters, + state=state, + ) + + assert manifest["boundaries"] == [ + "2026-01-01T00:00:00+00:00", + "2026-01-01T01:00:00+00:00", + ] + assert manifest_boundaries(manifest) == manifest["boundaries"] + + +def test_legacy_manifest_without_boundaries_gets_endpoint_chain(): + """Старый manifest без boundaries читается как пара [T0, T_end].""" + from clickstream_generator.startup_history_artifact import manifest_boundaries + + manifest = { + "model_t0": "2026-01-01T00:00:00+00:00", + "model_t_end": "2026-01-03T00:00:00+00:00", + } + + assert manifest_boundaries(manifest) == [ + "2026-01-01T00:00:00+00:00", + "2026-01-03T00:00:00+00:00", + ] + + +def test_manifest_rejects_boundary_chain_with_wrong_endpoint(): + """Цепочка не может расходиться с текущим model_t_end.""" + from clickstream_generator.startup_history_artifact import manifest_boundaries + + manifest = { + "model_t0": "2026-01-01T00:00:00+00:00", + "model_t_end": "2026-01-03T00:00:00+00:00", + "boundaries": [ + "2026-01-01T00:00:00+00:00", + "2026-01-02T00:00:00+00:00", + ], + } + + with pytest.raises(ValueError, match="model_t_end"): + manifest_boundaries(manifest) + + +def test_build_manifest_rejects_explicit_empty_boundaries(base_config): + """Явно пустая цепочка не подменяется границами по умолчанию.""" + from clickstream_generator.startup_history_artifact import ( + StartupHistoryArtifactBuilder, + build_manifest, + ) + + state = _state() + builder = StartupHistoryArtifactBuilder() + builder.add_batch(_batch()) + config = replace( + base_config, + model_t0=state.model_t0, + model_t_end=state.model_timestamp, + model_time_speed=1, + model_timezone="UTC", + seed=42, + ) + + with pytest.raises(ValueError, match="T0 и T_end"): + build_manifest(config, builder.counters, state, boundaries=[]) + + def test_artifact_validation_rejects_mismatched_manifest(base_config): """Валидация отвергает артефакт, где manifest не совпадает с событиями.""" from clickstream_generator.startup_history_artifact import ( diff --git a/scripts/check_generated_history_chain.sh b/scripts/check_generated_history_chain.sh new file mode 100644 index 0000000..ee3af68 --- /dev/null +++ b/scripts/check_generated_history_chain.sh @@ -0,0 +1,275 @@ +#!/usr/bin/env bash +# +# Проверяет завершённые стыки между порциями модельной истории. + +set -euo pipefail + +COMPOSE_BIN="${COMPOSE_BIN:-docker compose}" +CLICKHOUSE_SERVICE="${CLICKHOUSE_SERVICE:-clickhouse}" +CLICKHOUSE_USER="${CLICKHOUSE_USER:-default}" +CLICKHOUSE_PASSWORD="${CLICKHOUSE_PASSWORD:-123456}" + +fail() { + echo "Ошибка: $*" >&2 + exit 1 +} + +clickhouse_datetime_literal() { + local value="$1" + local normalized + if ! normalized="$(date --utc --date="${value}" '+%Y-%m-%d %H:%M:%S.%6N')"; then + fail "не удалось привести границу к UTC: ${value}" + fi + echo "${normalized}" +} + +ch_query() { + ${COMPOSE_BIN} exec -T "${CLICKHOUSE_SERVICE}" clickhouse-client \ + --user="${CLICKHOUSE_USER}" \ + --password="${CLICKHOUSE_PASSWORD}" \ + --query "$1" +} + +mapfile -t boundaries < <( + ${COMPOSE_BIN} run --rm --no-deps generator \ + python -m clickstream_generator.startup_history_artifact kafka-boundaries +) + +boundary_count="${#boundaries[@]}" +(( boundary_count >= 2 )) || fail "manifest boundaries должен содержать T0 и T_end" + +manifest_topic_rows="$(${COMPOSE_BIN} run --rm --no-deps generator \ + python -m clickstream_generator.startup_history_artifact kafka-topic-rows)" +IFS=$'\t' read -r expected_browser_rows expected_location_rows \ + expected_device_rows expected_geo_rows <<< "${manifest_topic_rows}" +for rows in \ + "${expected_browser_rows}" \ + "${expected_location_rows}" \ + "${expected_device_rows}" \ + "${expected_geo_rows}"; do + [[ "${rows}" =~ ^[0-9]+$ ]] || fail "manifest содержит неверные счётчики топиков" +done + +echo "=== Проверка порций истории ===" +for ((index = 0; index < boundary_count - 1; index++)); do + left="$(clickhouse_datetime_literal "${boundaries[index]}")" + right="$(clickhouse_datetime_literal "${boundaries[index + 1]}")" + rows="$(ch_query " +WITH + toDateTime64('${left}', 6, 'UTC') AS segment_start, + toDateTime64('${right}', 6, 'UTC') AS segment_end +SELECT count() +FROM dm.v_events_enriched +WHERE parseDateTime64BestEffortOrNull(toString(event_ts), 6, 'UTC') >= segment_start + AND parseDateTime64BestEffortOrNull(toString(event_ts), 6, 'UTC') < segment_end +FORMAT TabSeparated")" + [[ "${rows}" =~ ^[0-9]+$ ]] || fail "не удалось прочитать порцию ${boundaries[index]}" + (( rows > 0 )) || fail "порция [${boundaries[index]}, ${boundaries[index + 1]}) пуста" + echo "segment_${index}_rows=${rows}" +done + +actual_topic_rows="$(ch_query " +SELECT + (SELECT count() FROM stg.browser_raw) AS actual_topic_rows, + (SELECT count() FROM stg.location_raw) AS actual_location_rows, + (SELECT count() FROM stg.device_raw) AS actual_device_rows, + (SELECT count() FROM stg.geo_raw) AS actual_geo_rows +FORMAT TabSeparated")" +IFS=$'\t' read -r actual_browser_rows actual_location_rows \ + actual_device_rows actual_geo_rows <<< "${actual_topic_rows}" +topic_names=(browser_events location_events device_events geo_events) +expected_rows=( + "${expected_browser_rows}" + "${expected_location_rows}" + "${expected_device_rows}" + "${expected_geo_rows}" +) +actual_rows=( + "${actual_browser_rows}" + "${actual_location_rows}" + "${actual_device_rows}" + "${actual_geo_rows}" +) +for ((index = 0; index < ${#topic_names[@]}; index++)); do + [[ "${actual_rows[index]}" =~ ^[0-9]+$ ]] \ + || fail "не удалось прочитать число строк ${topic_names[index]} в STG" + if (( actual_rows[index] > expected_rows[index] )); then + fail "найден хвост ${topic_names[index]} после manifest: STG=${actual_rows[index]}, manifest=${expected_rows[index]}" + fi + if (( actual_rows[index] < expected_rows[index] )); then + fail "в STG не хватает строк ${topic_names[index]}: STG=${actual_rows[index]}, manifest=${expected_rows[index]}" + fi +done + +last_boundary="$(clickhouse_datetime_literal "${boundaries[boundary_count - 1]}")" +data_tail_rows="$(ch_query " +SELECT count() AS data_tail_rows +FROM stg.browser_raw +WHERE parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + ) >= toDateTime64('${last_boundary}', 6, 'UTC') +FORMAT TabSeparated")" +[[ "${data_tail_rows}" =~ ^[0-9]+$ ]] \ + || fail "не удалось проверить хвост после ${boundaries[boundary_count - 1]}" +(( data_tail_rows == 0 )) \ + || fail "найден хвост данных после последней границы manifest: ${data_tail_rows}" + +if (( boundary_count == 2 )); then + echo "Завершённых внутренних стыков пока нет." + exit 0 +fi + +echo "=== Проверка внутренних стыков boundaries ===" +for ((index = 1; index < boundary_count - 1; index++)); do + boundary="$(clickhouse_datetime_literal "${boundaries[index]}")" + + pairing="$(ch_query " +WITH + toDateTime64('${boundary}', 6, 'UTC') AS boundary, + browser_rows AS ( + SELECT + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + ) AS event_ts + FROM stg.browser_raw + WHERE event_id IS NOT NULL AND click_id IS NOT NULL AND event_ts IS NOT NULL + GROUP BY event_id, click_id, event_ts + ), + crossing AS ( + SELECT click_id + FROM browser_rows + GROUP BY click_id + HAVING min(event_ts) < boundary AND max(event_ts) >= boundary + ), + browser_counts AS ( + SELECT click_id, count() AS browser_rows + FROM browser_rows + WHERE click_id IN (SELECT click_id FROM crossing) + GROUP BY click_id + ), + location_ids AS ( + SELECT toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id + FROM stg.location_raw + WHERE event_id IS NOT NULL + GROUP BY event_id + ), + device_counts AS ( + SELECT + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + count() AS device_rows + FROM stg.device_raw + WHERE click_id IN (SELECT click_id FROM crossing) + GROUP BY click_id + ), + geo_counts AS ( + SELECT + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id, + count() AS geo_rows + FROM stg.geo_raw + WHERE click_id IN (SELECT click_id FROM crossing) + GROUP BY click_id + ) +SELECT + (SELECT count() FROM crossing) AS crossing_visits, + ( + SELECT count() + FROM browser_rows AS b + LEFT JOIN location_ids AS l ON l.event_id = b.event_id + WHERE b.click_id IN (SELECT click_id FROM crossing) + AND l.event_id IS NULL + ) AS unpaired_location_rows, + ( + SELECT ifNull(sum(greatest(b.browser_rows - ifNull(d.device_rows, 0), 0)), 0) + FROM browser_counts AS b + LEFT JOIN device_counts AS d ON d.click_id = b.click_id + ) AS unpaired_device_rows, + ( + SELECT ifNull(sum(greatest(b.browser_rows - ifNull(g.geo_rows, 0), 0)), 0) + FROM browser_counts AS b + LEFT JOIN geo_counts AS g ON g.click_id = b.click_id + ) AS unpaired_geo_rows +SETTINGS join_use_nulls = 1 +FORMAT TabSeparated")" + + IFS=$'\t' read -r crossing_visits unpaired_location_rows unpaired_device_rows unpaired_geo_rows <<< "${pairing}" + [[ "${crossing_visits}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_location_rows}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_device_rows}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_geo_rows}" =~ ^[0-9]+$ ]] \ + || fail "не удалось прочитать пары на границе ${boundaries[index]}" + if (( unpaired_location_rows > 0 || unpaired_device_rows > 0 || unpaired_geo_rows > 0 )); then + fail "непарные строки на границе ${boundaries[index]}: location=${unpaired_location_rows}, device=${unpaired_device_rows}, geo=${unpaired_geo_rows}" + fi + + context="$(ch_query " +WITH + toDateTime64('${boundary}', 6, 'UTC') AS boundary, + crossing AS ( + SELECT click_id + FROM + ( + SELECT + click_id, + parseDateTime64BestEffortOrNull( + toString(event_ts), 6, 'UTC' + ) AS model_event_ts + FROM dds.event + WHERE event_ts IS NOT NULL + ) + GROUP BY click_id + HAVING min(model_event_ts) < boundary AND max(model_event_ts) >= boundary + ), + facts AS ( + SELECT + e.click_id, + 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 + FROM dds.event AS e + INNER JOIN crossing AS c ON c.click_id = e.click_id + GROUP BY e.click_id + ) +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 +FROM facts +FORMAT TabSeparated")" + + IFS=$'\t' read -r context_crossing_visits per_event_homogeneous_visits <<< "${context}" + [[ "${context_crossing_visits}" =~ ^[0-9]+$ ]] \ + && [[ "${per_event_homogeneous_visits}" =~ ^[0-9]+$ ]] \ + || fail "не удалось прочитать фактуру на границе ${boundaries[index]}" + (( context_crossing_visits == crossing_visits )) \ + || fail "STG и DDS нашли разное число переходящих визитов на ${boundaries[index]}" + if (( crossing_visits == 0 )); then + echo "boundary ${index}/${boundary_count}: crossing_visits=0 — однородность неприменима" + else + (( per_event_homogeneous_visits == crossing_visits )) \ + || fail "browser/referer/utm меняются на границе ${boundaries[index]}: ${per_event_homogeneous_visits}/${crossing_visits}" + echo "boundary ${index}/${boundary_count}: crossing_visits=${crossing_visits} — однородность подтверждена" + fi + + echo "boundary_${index}=${boundaries[index]}" + echo "boundary_${index}_crossing_visits=${crossing_visits}" + echo "boundary_${index}_unpaired_location_rows=${unpaired_location_rows}" + echo "boundary_${index}_unpaired_device_rows=${unpaired_device_rows}" + echo "boundary_${index}_unpaired_geo_rows=${unpaired_geo_rows}" +done diff --git a/tests/test_generated_analytics_check_contract.py b/tests/test_generated_analytics_check_contract.py index 96f022c..0ecc4c8 100644 --- a/tests/test_generated_analytics_check_contract.py +++ b/tests/test_generated_analytics_check_contract.py @@ -1,5 +1,6 @@ from pathlib import Path import os +import re import subprocess import textwrap @@ -10,6 +11,7 @@ REPO_ROOT = Path(__file__).resolve().parents[1] CHECK_SCRIPT = REPO_ROOT / "scripts" / "check_generated_analytics.sh" COMPOSE_FILE = REPO_ROOT / "docker-compose.yml" OPERATIONS_DOC = REPO_ROOT / "docs" / "OPERATIONS.md" +CHAIN_CHECK_SCRIPT = REPO_ROOT / "scripts" / "check_generated_history_chain.sh" def _required_runtime_line(script, command): @@ -345,3 +347,229 @@ def test_generated_history_check_rejects_empty_ods_context(): assert "ods_geo_rows" in script assert "ODS device пустой на стыке" in script assert "ODS geo пустой на стыке" in script + + +def test_boundary_chain_has_separate_batch_check_and_keeps_live_gate_unchanged(): + """Цепочка boundaries проверяется отдельной целью, не режимом live-гейта.""" + chain = CHAIN_CHECK_SCRIPT.read_text(encoding="utf-8") + live = CHECK_SCRIPT.read_text(encoding="utf-8") + makefile = (REPO_ROOT / "Makefile").read_text(encoding="utf-8") + + assert "generated-history-chain-check" in makefile + assert "check_generated_history_chain.sh" in makefile + assert "kafka-boundaries" in chain + assert "boundaries" in chain + assert "unpaired_location_rows" in chain + assert "unpaired_device_rows" in chain + assert "unpaired_geo_rows" in chain + assert "browser_names" in chain + assert "referer_urls" in chain + assert "utm_sources" in chain + assert "per_event_homogeneous_visits" in chain + assert "toDateTime64('${boundary}', 6, 'UTC') AS boundary" in chain + assert "index < boundary_count - 1" in chain + assert "kafka-boundaries" not in live + + +def _fake_chain_compose(tmp_path): + fake = tmp_path / "chain-compose" + fake.write_text( + textwrap.dedent( + """\ + #!/usr/bin/env bash + set -euo pipefail + + args="$*" + query="" + previous="" + for arg in "$@"; do + if [[ "$previous" == "--query" ]]; then + query="$arg" + break + fi + previous="$arg" + done + + if [[ "$args" == *"kafka-boundaries"* ]]; then + printf '%b' "${CHAIN_BOUNDARIES}" + exit 0 + fi + + if [[ "$args" == *"kafka-topic-rows"* ]]; then + printf '%s\n' "${CHAIN_MANIFEST_TOPIC_ROWS:-10 10 10 10}" + exit 0 + fi + + if [[ "$args" == *"clickhouse-client"* ]]; then + printf '%s\n' "$query" >> "${CHAIN_QUERY_LOG}" + if [[ "$query" == *"data_tail_rows"* ]]; then + printf '%s\n' "${CHAIN_TAIL_ROWS:-0}" + elif [[ "$query" == *"actual_topic_rows"* ]]; then + printf '%s\n' "${CHAIN_ACTUAL_TOPIC_ROWS:-10 10 10 10}" + elif [[ "$query" == *"FROM dm.v_events_enriched"* ]]; then + printf '%s\n' "${CHAIN_SEGMENT_ROWS:-1}" + elif [[ "$query" == *"unpaired_location_rows"* ]]; then + printf '%s\n' "${CHAIN_PAIRING:-1 0 0 0}" + elif [[ "$query" == *"per_event_homogeneous_visits"* ]]; then + printf '%s\n' "${CHAIN_CONTEXT:-1 1}" + else + printf '%s\n' '1' + fi + exit 0 + fi + + echo "unexpected compose call: $args" >&2 + exit 92 + """ + ), + encoding="utf-8", + ) + fake.chmod(0o755) + return fake + + +def _run_chain_check(tmp_path, **overrides): + query_log = tmp_path / "queries.log" + env = os.environ.copy() + env.update( + { + "COMPOSE_BIN": str(_fake_chain_compose(tmp_path)), + "CHAIN_BOUNDARIES": ( + "2026-01-01T00:00:00+00:00\\n" + "2026-01-02T00:00:00+00:00\\n" + ), + "CHAIN_QUERY_LOG": str(query_log), + } + ) + env.update(overrides) + result = subprocess.run( + ["bash", str(CHAIN_CHECK_SCRIPT)], + cwd=REPO_ROOT, + env=env, + text=True, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + timeout=30, + check=False, + ) + return result, query_log.read_text(encoding="utf-8") + + +def test_boundary_chain_rejects_data_tail_after_manifest_end(tmp_path): + """Сообщения после последней границы означают незавершённый next-day.""" + result, queries = _run_chain_check(tmp_path, CHAIN_TAIL_ROWS="1") + + assert result.returncode != 0 + assert "хвост" in result.stderr + assert "FROM stg.browser_raw" in queries + assert ">= toDateTime64('2026-01-02 00:00:00.000000', 6, 'UTC')" in queries + + +@pytest.mark.parametrize( + ("actual_rows", "topic"), + [ + ("10\t11\t10\t10", "location_events"), + ("10\t10\t11\t10", "device_events"), + ("10\t10\t10\t11", "geo_events"), + ], +) +def test_boundary_chain_rejects_non_browser_topic_tail( + tmp_path, + actual_rows, + topic, +): + """Хвост отдельной Kafka-темы обнаруживается без browser-события.""" + result, queries = _run_chain_check( + tmp_path, + CHAIN_ACTUAL_TOPIC_ROWS=actual_rows, + ) + + assert result.returncode != 0 + assert "хвост" in result.stderr + assert topic in result.stderr + assert "FROM stg.location_raw" in queries + assert "FROM stg.device_raw" in queries + assert "FROM stg.geo_raw" in queries + + +def test_boundary_chain_reports_boundary_without_crossing_visits(tmp_path): + """Естественный нулевой стык проходит с явным статусом без однородности.""" + result, _ = _run_chain_check( + tmp_path, + CHAIN_BOUNDARIES=( + "2026-01-01T00:00:00+00:00\\n" + "2026-01-02T00:00:00+00:00\\n" + "2026-01-03T00:00:00+00:00\\n" + ), + CHAIN_PAIRING="0\t0\t0\t0", + CHAIN_CONTEXT="0\t0", + ) + + assert result.returncode == 0, result.stderr + assert ( + "boundary 1/3: crossing_visits=0 — однородность неприменима" + in result.stdout + ) + + +def test_boundary_chain_still_rejects_unpaired_rows_without_crossing_visits( + tmp_path, +): + """Нулевой стык не отключает обязательную проверку пар сообщений.""" + result, _ = _run_chain_check( + tmp_path, + CHAIN_BOUNDARIES=( + "2026-01-01T00:00:00+00:00\\n" + "2026-01-02T00:00:00+00:00\\n" + "2026-01-03T00:00:00+00:00\\n" + ), + CHAIN_PAIRING="0\t1\t0\t0", + CHAIN_CONTEXT="0\t0", + ) + + assert result.returncode != 0 + assert "непарные строки" in result.stderr + + +def test_boundary_chain_still_rejects_empty_segment(tmp_path): + """Отсутствие переходящих визитов не разрешает пустую порцию дня.""" + result, _ = _run_chain_check(tmp_path, CHAIN_SEGMENT_ROWS="0") + + assert result.returncode != 0 + assert "порция" in result.stderr + assert "пуста" in result.stderr + + +def test_boundary_chain_normalizes_offsets_and_uses_explicit_utc(tmp_path): + """Граница с +03:00 проверяет тот же UTC-интервал, а не местные часы.""" + result, queries = _run_chain_check( + tmp_path, + CHAIN_BOUNDARIES=( + "2026-01-01T03:00:00+03:00\\n" + "2026-01-02T03:00:00+03:00\\n" + "2026-01-03T03:00:00+03:00\\n" + ), + ) + + assert result.returncode == 0, result.stderr + assert "toDateTime64('2026-01-01 00:00:00.000000', 6, 'UTC')" in queries + assert "toDateTime64('2026-01-02 00:00:00.000000', 6, 'UTC')" in queries + assert "toDateTime64('2026-01-03 00:00:00.000000', 6, 'UTC')" in queries + literals = re.findall(r"toDateTime64\([^)]*\)", queries) + assert len(literals) == 7 + assert all(", 6, 'UTC')" in literal for literal in literals) + + +def test_boundary_chain_preserves_fractional_boundary_seconds(tmp_path): + """Приведение к UTC сохраняет микросекунды полуоткрытой границы.""" + result, queries = _run_chain_check( + tmp_path, + CHAIN_BOUNDARIES=( + "2026-01-01T00:00:00.500000+00:00\\n" + "2026-01-02T00:00:00.500000+00:00\\n" + ), + ) + + assert result.returncode == 0, result.stderr + assert "toDateTime64('2026-01-01 00:00:00.500000', 6, 'UTC')" in queries + assert "toDateTime64('2026-01-02 00:00:00.500000', 6, 'UTC')" in queries