From dab120a07e3ad958b5c0c7b8daa4bb9659b7ff30 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 12 Jul 2026 19:33:28 +0300 Subject: [PATCH] =?UTF-8?q?fix(generator):=20=D1=83=D1=81=D1=82=D1=80?= =?UTF-8?q?=D0=B0=D0=BD=D1=91=D0=BD=20=D1=84=D0=BB=D0=B0=D0=BA=D0=B8=20run?= =?UTF-8?q?time-=D0=B3=D0=B5=D0=B9=D1=82=D0=B0=20=D1=81=D1=82=D1=8B=D0=BA?= =?UTF-8?q?=D0=B0=20backfill/live?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - гейт generated-history-runtime-check падал через раз: docker compose stop убивал генератор SIGKILL'ом посреди batch (SIGTERM не ловился), непарные строки маскировались под «смену фактуры» (задача 20). - Что: - генератор грациозно завершается по SIGTERM: текущий batch дописывается во все топики с flush и записью history; compose даёт минуту grace. - runtime-check ждёт пересекающий визит в STG (предусловие проверки), seam-check получил precheck непарных live-строк; в STG-запросах закреплён 'UTC' против сдвига наивных меток в поясе сервера. - контрактные тесты усилены, задача 20 закрыта, блокер задачи 13 снят. - Проверка: - тесты: 192 passed (generator), 21 passed (контракты); - стенд: 3 подряд зелёных make generated-history-runtime-check; красный сценарий (искажение referer_url пересекающего визита) валит гейт прежним сообщением при нулевых непарных счётчиках. --- .../13-backfill-top-up-from-snapshot.md | 5 +- .../issues/20-flaky-runtime-seam-check.md | 54 ++++++-- docker-compose.yml | 2 + docs/OPERATIONS.md | 9 +- .../src/clickstream_generator/service.py | 17 ++- generator/tests/test_service.py | 86 +++++++++++- scripts/check_generated_analytics.sh | 108 +++++++++++++++ .../run_generated_history_runtime_check.sh | 52 ++++++- ...test_generated_analytics_check_contract.py | 130 +++++++++++++++++- 9 files changed, 433 insertions(+), 30 deletions(-) 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 3269b52..09dff48 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 @@ -81,8 +81,9 @@ Status: ready-for-agent ## Blocked by -- `20-flaky-runtime-seam-check.md` — доверенный стабильный гейт стыка нужен - до навешивания на него цепочки границ; фикс 20 идёт первым. +- Нет. Последний блокер снят 2026-07-12: `20-flaky-runtime-seam-check.md` + закрыта, гейт стыка стабилен (3 подряд зелёных прогона, красный сценарий + ловится). - Исторические блокеры закрыты: `15-world-boundary-after-cross-review.md` (граница миров) и `17-trusted-checks-startup-history-superset.md` (доверенные проверки) — done. diff --git a/.scratch/generator-model-time-startup-history/issues/20-flaky-runtime-seam-check.md b/.scratch/generator-model-time-startup-history/issues/20-flaky-runtime-seam-check.md index df2f312..a9d22fd 100644 --- a/.scratch/generator-model-time-startup-history/issues/20-flaky-runtime-seam-check.md +++ b/.scratch/generator-model-time-startup-history/issues/20-flaky-runtime-seam-check.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: done # Runtime-проверка стыка нестабильна: фактура меняется через раз @@ -99,24 +99,54 @@ location/device/geo -> DDS строит `dds.event` через `LEFT JOIN locati сужение окна `GEN_LIVE_CHECK_MINUTES` и прочие способы снизить вероятность — гонку они не убирают и фиксом не считаются. +## Принятый остаточный риск (решение координатора, 2026-07-12) + +Ревью отметило: `stop_grace_period: 1m` не ограничивает публикацию по +времени — при зависании Kafka дольше минуты SIGKILL всё ещё оборвёт batch. +Риск принят: штатный тик публикуется за секунды (минута — многократный +запас), а ограничивать публикацию таймером значило бы рвать batch уже по +построению. Ключевое отличие от исходного флаки: такой обрыв больше не +маскируется под «смену фактуры» — precheck непарных строк назовёт его +явно, гейт упадёт громко и честно. + ## Acceptance criteria - [x] Причина расхождения 8/19 найдена и названа (код, не догадка) — см. «Диагноз» выше. -- [ ] Генератор корректно завершается по SIGTERM: текущий batch дописывается +- [x] Генератор корректно завершается по SIGTERM: текущий batch дописывается во все четыре топика целиком (flush + запись history), потом процесс выходит; runtime-check дожидается фактической остановки контейнера - (направление фикса, пункт 1). -- [ ] Seam-check различает «непарные live-строки» и «смена фактуры»: + (направление фикса, пункт 1). Плюс `stop_grace_period: 1m` в compose. +- [x] Seam-check различает «непарные live-строки» и «смена фактуры»: precheck называет реальную причину (пункт 2). -- [ ] Красный сценарий по-прежнему ловится: настоящая смена per-event - фактуры внутри `click_id` (инъекция в тесте или контролируемое искажение - данных) валит гейт с прежним сообщением — precheck и фикс не сделали - проверку мягче. -- [ ] Стабильность обоснована структурно (гонка снята по построению: - завершение только на границе batch), а не статистикой прогонов; - `make generated-history-runtime-check` — N подряд зелёных (N >= 3) - как дымовая проверка поверх этого довода, зафиксировано в задаче. +- [x] Красный сценарий по-прежнему ловится: контролируемое искажение + `referer_url` у события пересекающего визита в `dds.event` уронило гейт + с прежним сообщением «per-event фактура меняется на стыке: 18/19» при + нулевых счётчиках непарных строк (стендовая приёмка 2026-07-12). +- [x] Стабильность обоснована структурно (SIGTERM ставит флаг, тик + дописывает все четыре топика + flush + history и выходит на границе + batch — тест `test_sigterm_during_publish_finishes_current_batch`); + дымовая проверка поверх довода: 3 подряд зелёных + `make generated-history-runtime-check` (2026-07-12, прогоны 19:25, + 19:28, 19:30). + +## Находки стендовой приёмки (2026-07-12, исправлены в этой же задаче) + +Ревью по чтению кода их поймать не могло — вскрылись только прогонами: + +1. **Гейт жил на побочном эффекте бага.** Пересекающие визиты для проверки + стыка появлялись только потому, что генератор игнорировал SIGTERM и + дописывал ~10 секунд данных до SIGKILL. После фикса live-хвост стал + коротким и пересекающих визитов могло не быть вовсе. Фикс: шаг 6 + runtime-check ждёт не «любую новую STG-строку», а появления + пересекающего визита в STG (детерминированное предусловие проверки, + то же окно, что у seam-SQL). +2. **Часовой пояс в STG-запросах.** Сырые `event_timestamp` наивные и + означают UTC, сервер ClickHouse — Europe/Moscow: сравнение с границей, + заданной с `+00:00`, уезжало на 3 часа, precheck был зелёным вакуумно + (не видел ни одной live-строки). Фикс: явный `'UTC'` в + `parseDateTime64BestEffort*` с обеих сторон сравнения (4 места), + закреплено контрактными тестами. ## Blocked by diff --git a/docker-compose.yml b/docker-compose.yml index 770bb92..a12caff 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -332,6 +332,8 @@ services: depends_on: - kafka restart: unless-stopped + # SIGTERM останавливает генератор на границе пакета; минута нужна на сброс буфера Kafka. + stop_grace_period: 1m healthcheck: test: ["CMD", "python", "-c", "import sys; sys.exit(0)"] interval: 30s diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 4f0384d..8ce706d 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -220,9 +220,12 @@ CHECK_LIVE_SEAM=1 GEN_LIVE_CHECK_MINUTES=10 make generated-history-check make generated-history-runtime-check ``` -По умолчанию он берёт профиль `daily-wave`, но сжимает историю до `1h`, ждёт -новые STG-строки от live до 25 секунд, делает второй batch и проверяет стык с -`CHECK_LIVE_SEAM=1`. +По умолчанию он берёт профиль `daily-wave`, но сжимает историю до `1h` и до +25 секунд ждёт в STG визит с browser-событиями по обе стороны границы. Затем +он делает второй batch и проверяет стык с `CHECK_LIVE_SEAM=1`. +При остановке generator получает SIGTERM и завершает текущий пакет. Параметр +`stop_grace_period: 1m` даёт время дописать четыре топика, сбросить буферы и +записать историю пакета до принудительной остановки контейнера. Если машина медленная, можно увеличить только ожидания: ```bash diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index 3d7a7f4..e30d5eb 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -1,7 +1,9 @@ """Основной сервисный цикл генератора.""" import logging +import signal import sys +import threading import time import uuid from datetime import datetime, timedelta, timezone @@ -54,6 +56,8 @@ class GeneratorService: self.state_manager: KafkaStateManager | None = None self.manifest_manager: KafkaStartupHistoryManifest | None = None self._running = False + self._stop_requested = False + self._shutdown_event = threading.Event() self._tick = 0 self._model_time = config.model_t0 @@ -168,7 +172,7 @@ class GeneratorService: else: logger.info("State management disabled") - self._running = True + self._running = not self._stop_requested try: self._main_loop() @@ -181,6 +185,7 @@ class GeneratorService: """Останавливает сервис.""" logger.info("Stopping generator service...") self._running = False + self._shutdown_event.set() if self.publisher: self.publisher.close() if self.history: @@ -190,6 +195,13 @@ class GeneratorService: if self.manifest_manager: self.manifest_manager.close() + def request_stop(self, signum, _frame) -> None: + """Просит завершить сервис после текущего batch.""" + logger.info("Received shutdown signal %s", signum) + self._stop_requested = True + self._running = False + self._shutdown_event.set() + def restore_from_startup_history( self, state, @@ -642,7 +654,7 @@ class GeneratorService: sleep_time = max(0, self.config.tick_seconds - elapsed) if sleep_time > 0: logger.debug(f"Sleeping for {sleep_time:.1f}s until next tick") - time.sleep(sleep_time) + self._shutdown_event.wait(sleep_time) def _advance_model_time(self) -> None: """Сдвигает модельное время после успешного live-тика.""" @@ -656,6 +668,7 @@ def main(): try: config = Config() service = GeneratorService(config) + signal.signal(signal.SIGTERM, service.request_stop) service.start() except Exception as e: logger.exception(f"Fatal error: {e}") diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index eccfdca..c391838 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -6,6 +6,7 @@ import logging import random import hashlib import json +import signal from dataclasses import replace from datetime import datetime, timedelta, timezone from time import sleep as real_sleep @@ -21,6 +22,7 @@ from generator import ( GeneratorState, KafkaBatchHistory, TickStreamGenerator, + main, ) from clickstream_generator.state import UnsupportedStateVersionError @@ -210,8 +212,9 @@ class TestGeneratorServiceSteadyStream: service._running = False with caplog.at_level(logging.INFO, logger="generator"), \ - patch( - "clickstream_generator.service.time.sleep", + patch.object( + service._shutdown_event, + "wait", side_effect=stop_after_first_tick, ): service._main_loop() @@ -221,6 +224,77 @@ class TestGeneratorServiceSteadyStream: assert "completed:" not in caplog.text assert "sent" not in caplog.text + def test_sigterm_during_publish_finishes_current_batch(self, base_config): + """SIGTERM посреди публикации не обрывает связанные топики batch.""" + config = replace( + base_config, + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=3, + metrics_enabled=False, + state_enabled=False, + ) + service = GeneratorService(config) + publisher = MagicMock() + history = MagicMock() + published_topics = [] + + def publish_and_request_stop(topic, events): + published_topics.append(topic) + if len(published_topics) == 1: + service.request_stop(signal.SIGTERM, None) + return len(events), 0 + + publisher.publish.side_effect = publish_and_request_stop + + with patch( + "clickstream_generator.service.start_http_server", + side_effect=AssertionError("тест не должен открывать порт метрик"), + ), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch( + "clickstream_generator.service.KafkaPublisher", + return_value=publisher, + ), \ + patch( + "clickstream_generator.service.KafkaBatchHistory", + return_value=history, + ), \ + patch.object( + service.generator, + "_calculate_events_count", + return_value=int(EXPECTED_VISIT_EVENTS), + ): + service.start() + + assert published_topics == [ + "browser_events", + "location_events", + "device_events", + "geo_events", + ] + publisher.flush.assert_called_once_with() + history.add.assert_called_once() + history.flush.assert_called_once_with() + publisher.close.assert_called_once_with() + history.close.assert_called_once_with() + assert service._tick == 1 + + def test_main_connects_sigterm_to_graceful_stop(self, base_config): + """Точка входа передаёт SIGTERM сервису как штатный запрос остановки.""" + service = MagicMock() + + with patch("clickstream_generator.service.Config", return_value=base_config), \ + patch("clickstream_generator.service.GeneratorService", return_value=service), \ + patch("clickstream_generator.service.signal.signal") as register_signal: + main() + + register_signal.assert_called_once_with(signal.SIGTERM, service.request_stop) + service.start.assert_called_once_with() + def _run_service_ticks(self, config, wall_now: datetime, ticks_count: int): service = GeneratorService(config) service.publisher = MagicMock() @@ -256,7 +330,11 @@ class TestGeneratorServiceSteadyStream: "_calculate_events_count", side_effect=calculate_events_count, ), \ - patch("clickstream_generator.service.time.sleep", side_effect=stop_after_tick): + patch.object( + service._shutdown_event, + "wait", + side_effect=stop_after_tick, + ): service._main_loop() published = {} @@ -294,7 +372,7 @@ class TestGeneratorServiceSteadyStream: service.generator, "_visit_pause_seconds", return_value=0.05, - ), patch("clickstream_generator.service.time.sleep") as sleep_mock: + ), patch.object(service._shutdown_event, "wait") as sleep_mock: sleep_mock.side_effect = stop_after_second_tick service._main_loop() diff --git a/scripts/check_generated_analytics.sh b/scripts/check_generated_analytics.sh index 72cd11c..3d70729 100644 --- a/scripts/check_generated_analytics.sh +++ b/scripts/check_generated_analytics.sh @@ -335,6 +335,114 @@ if [[ "${should_check_live_seam}" == "1" ]]; then echo "" echo "=== Проверка ClickHouse: стык backfill/live в DDS ===" + live_pairing_query=" +WITH + parseDateTime64BestEffort('${manifest_model_t_end}', 6, 'UTC') AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end, + live_browser AS ( + SELECT + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id + FROM stg.browser_raw + WHERE parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + ) >= t_end + AND parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + ) < t_live_end + AND event_id IS NOT NULL + AND click_id IS NOT NULL + GROUP BY event_id, click_id + ), + live_click_ids AS ( + SELECT click_id + FROM live_browser + GROUP BY click_id + ), + browser_counts AS ( + SELECT click_id, count() AS browser_rows + FROM + ( + SELECT + toUUIDOrNull(JSONExtractString(raw, 'event_id')) AS event_id, + toUUIDOrNull(JSONExtractString(raw, 'click_id')) AS click_id + FROM stg.browser_raw + WHERE event_id IS NOT NULL AND click_id IS NOT NULL + GROUP BY event_id, click_id + ) + WHERE click_id IN (SELECT click_id FROM live_click_ids) + 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 IS NOT NULL + AND click_id IN (SELECT click_id FROM live_click_ids) + 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 IS NOT NULL + AND click_id IN (SELECT click_id FROM live_click_ids) + GROUP BY click_id + ) +SELECT + (SELECT count() FROM live_browser) AS live_browser_rows, + ( + SELECT count() + FROM live_browser AS b + LEFT JOIN location_ids AS l ON l.event_id = b.event_id + WHERE l.event_id IS NULL + ) AS unpaired_location_rows, + ( + SELECT ifNull(sum(if( + b.browser_rows > ifNull(d.device_rows, 0), + 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(if( + b.browser_rows > ifNull(g.geo_rows, 0), + 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" + + live_pairing="$(ch_query "${live_pairing_query}")" + IFS=$'\t' read -r live_browser_rows unpaired_location_rows unpaired_device_rows unpaired_geo_rows <<< "${live_pairing}" + + [[ "${live_browser_rows}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_location_rows}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_device_rows}" =~ ^[0-9]+$ ]] \ + && [[ "${unpaired_geo_rows}" =~ ^[0-9]+$ ]] \ + || fail "не удалось прочитать пары live-строк" + if (( unpaired_location_rows > 0 || unpaired_device_rows > 0 || unpaired_geo_rows > 0 )); then + fail "непарные live-строки: location=${unpaired_location_rows}, device=${unpaired_device_rows}, geo=${unpaired_geo_rows} из browser=${live_browser_rows}" + fi + + echo "live_browser_rows=${live_browser_rows}" + echo "unpaired_location_rows=${unpaired_location_rows}" + echo "unpaired_device_rows=${unpaired_device_rows}" + echo "unpaired_geo_rows=${unpaired_geo_rows}" + seam_duplicates_query=" WITH toDateTime64('${CH_MODEL_T0}', 6) AS t0, diff --git a/scripts/run_generated_history_runtime_check.sh b/scripts/run_generated_history_runtime_check.sh index 9060277..32ab7d7 100644 --- a/scripts/run_generated_history_runtime_check.sh +++ b/scripts/run_generated_history_runtime_check.sh @@ -57,9 +57,47 @@ SELECT FORMAT TabSeparated" } +stg_crossing_visits() { + ${COMPOSE_BIN} exec -T clickhouse clickhouse-client \ + --user=default \ + --password=123456 \ + --query " +WITH + parseDateTime64BestEffort('${GEN_MODEL_T_END}', 6, 'UTC') AS t_end, + t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE AS t_live_end, + browser AS ( + SELECT + JSONExtractString(raw, 'click_id') AS click_id, + parseDateTime64BestEffortOrNull( + JSONExtractString(raw, 'event_timestamp'), 6, 'UTC' + ) AS event_ts + FROM stg.browser_raw + WHERE click_id != '' + AND event_ts >= t_end - INTERVAL 30 MINUTE + AND event_ts < t_live_end + ) +SELECT count() +FROM +( + SELECT + click_id, + min(event_ts) AS first_ts, + max(event_ts) AS last_ts + FROM browser + GROUP BY click_id + HAVING first_ts < t_end AND last_ts >= t_end +) +FORMAT TabSeparated" +} + cleanup_live_generator() { ${COMPOSE_BIN} stop generator >/dev/null 2>&1 || true } + +stop_live_generator() { + # docker compose stop возвращается только после фактической остановки контейнера. + ${COMPOSE_BIN} stop generator +} trap cleanup_live_generator EXIT echo "=== Быстрая runtime-проверка startup-history/live seam ===" @@ -103,7 +141,7 @@ sleep "${WAIT_STG_SECONDS}" bash "${SCRIPT_DIR}/run_batch.sh" stg_rows_before_live="$(stg_total_rows)" -echo "Шаг 6: live-продолжение, ждём новые STG-строки до ${LIVE_SECONDS} сек." +echo "Шаг 6: live-продолжение, ждём переходящий визит в STG до ${LIVE_SECONDS} сек." GEN_RUN_MODE=live \ GEN_STATE_RESET=false \ GEN_LAUNCH_PROFILE="${GEN_LAUNCH_PROFILE}" \ @@ -120,19 +158,23 @@ GEN_MAX_EVENTS_PER_TICK="${GEN_MAX_EVENTS_PER_TICK}" \ ${COMPOSE_BIN} up -d --build generator deadline=$((SECONDS + LIVE_SECONDS)) while true; do - stg_rows_after_live="$(stg_total_rows)" - if (( stg_rows_after_live > stg_rows_before_live )); then + crossing_visits="$(stg_crossing_visits)" + [[ "${crossing_visits}" =~ ^[0-9]+$ ]] \ + || { echo "Ошибка: не удалось прочитать переходящие визиты из STG." >&2; exit 1; } + if (( crossing_visits > 0 )); then + stg_rows_after_live="$(stg_total_rows)" echo "live_stg_rows_before=${stg_rows_before_live}" echo "live_stg_rows_after=${stg_rows_after_live}" + echo "live_stg_crossing_visits=${crossing_visits}" break fi if (( SECONDS >= deadline )); then - echo "Ошибка: live-продолжение не записало новые STG-строки за ${LIVE_SECONDS} сек." >&2 + echo "Ошибка: live-продолжение не создало переходящий визит в STG за ${LIVE_SECONDS} сек." >&2 exit 1 fi sleep 1 done -cleanup_live_generator +stop_live_generator echo "Шаг 7: второй batch после live" sleep "${WAIT_STG_SECONDS}" diff --git a/tests/test_generated_analytics_check_contract.py b/tests/test_generated_analytics_check_contract.py index ad53ed7..96f022c 100644 --- a/tests/test_generated_analytics_check_contract.py +++ b/tests/test_generated_analytics_check_contract.py @@ -3,9 +3,23 @@ import os import subprocess import textwrap +import pytest + 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" + + +def _required_runtime_line(script, command): + matches = [ + line_number + for line_number, line in enumerate(script.splitlines()) + if line == command + ] + assert len(matches) == 1, f"{command} должна быть отдельной командой ровно один раз" + return matches[0] def _fake_compose( @@ -13,6 +27,8 @@ def _fake_compose( *, manifest_profile="ci", live_rows="3", + live_pairing="3\t0\t0\t0", + seam_context="19\t19\t19", views_mode="ok", ): fake = tmp_path / "docker-compose" @@ -53,8 +69,10 @@ def _fake_compose( printf '%s\\n' '{live_rows}' elif [[ "$query" == *"uniqExact(event_id)"* ]]; then printf '%s\\n' '16200\t16200\t0' + elif [[ "$query" == *"unpaired_location_rows"* ]]; then + printf '%s\\n' '{live_pairing}' elif [[ "$query" == *"per_event_homogeneous_visits"* ]]; then - printf '%s\\n' '19\t19\t19' + printf '%s\\n' '{seam_context}' elif [[ "$query" == *"ods_device_rows"* ]]; then printf '%s\\n' '19\t19\t19\t0\t0' elif [[ "$query" == *"SELECT source, rows"* ]]; then @@ -165,6 +183,71 @@ def test_generated_history_check_requires_live_rows_when_seam_is_required(tmp_pa assert "live-продолжение не записало строки" in result.stderr +def test_generated_history_check_names_unpaired_live_rows(tmp_path): + """Неполный live-batch получает отдельную ошибку до проверки фактуры.""" + fake_compose = _fake_compose( + tmp_path, + live_pairing="3\t1\t0\t0", + seam_context="19\t18\t19", + ) + + result = _run_generated_history_check( + tmp_path, + fake_compose=fake_compose, + PROFILE="ci", + CHECK_LIVE_SEAM="1", + ) + + assert result.returncode != 0 + assert "непарные live-строки: location=1, device=0, geo=0 из browser=3" in result.stderr + assert "per-event фактура меняется на стыке" not in result.stderr + + +def test_live_pairing_precheck_counts_device_and_geo_messages_in_stg(): + """Старая ODS-строка click_id не скрывает пропуск live device/geo.""" + script = CHECK_SCRIPT.read_text(encoding="utf-8") + live_pairing_query = script.split('live_pairing_query="', 1)[1].split( + 'FORMAT TabSeparated"', + 1, + )[0] + + assert "FROM stg.browser_raw" in live_pairing_query + assert "FROM stg.location_raw" in live_pairing_query + assert "FROM stg.device_raw" in live_pairing_query + assert "FROM stg.geo_raw" in live_pairing_query + assert "browser_rows - ifNull(d.device_rows, 0)" in live_pairing_query + assert "browser_rows - ifNull(g.geo_rows, 0)" in live_pairing_query + assert "FROM ods.device_by_click" not in live_pairing_query + assert "FROM ods.geo_by_click" not in live_pairing_query + assert ( + "parseDateTime64BestEffort('${manifest_model_t_end}', 6, 'UTC') AS t_end" + in live_pairing_query + ) + assert live_pairing_query.count( + "JSONExtractString(raw, 'event_timestamp'), 6, 'UTC'" + ) == 2 + + +def test_generated_history_check_still_rejects_real_per_event_change(tmp_path): + """Полные пары не скрывают настоящую смену per-event фактуры.""" + fake_compose = _fake_compose( + tmp_path, + live_pairing="3\t0\t0\t0", + seam_context="19\t18\t19", + ) + + result = _run_generated_history_check( + tmp_path, + fake_compose=fake_compose, + PROFILE="ci", + CHECK_LIVE_SEAM="1", + ) + + assert result.returncode != 0 + assert "per-event фактура меняется на стыке: 18/19" in result.stderr + assert "непарные live-строки" not in result.stderr + + def test_generated_history_check_fails_when_dm_views_query_fails(tmp_path): """Ошибка запроса DM-витрин не превращается в пустой успешный цикл.""" fake_compose = _fake_compose(tmp_path, views_mode="fail") @@ -201,14 +284,57 @@ def test_runtime_gate_has_bounded_daily_wave_live_seam_path(): assert 'PROFILE="${PROFILE:-daily-wave}"' in script assert 'GEN_HISTORY_DURATION="${GEN_HISTORY_DURATION:-1h}"' in script assert 'LIVE_SECONDS="${LIVE_SECONDS:-25}"' in script - assert "stg_rows_after_live > stg_rows_before_live" in script + assert "stg_crossing_visits()" in script + assert "crossing_visits=\"$(stg_crossing_visits)\"" in script + assert "crossing_visits > 0" in script + assert "t_end - INTERVAL 30 MINUTE" in script + assert "t_end + INTERVAL ${GEN_LIVE_CHECK_MINUTES} MINUTE" in script + assert "HAVING first_ts < t_end AND last_ts >= t_end" in script + assert ( + "parseDateTime64BestEffort('${GEN_MODEL_T_END}', 6, 'UTC') AS t_end" + in script + ) + assert "JSONExtractString(raw, 'event_timestamp'), 6, 'UTC'" in script + assert "stg_rows_after_live > stg_rows_before_live" not in script + assert "не создало переходящий визит" in script assert "GEN_RUN_MODE=live" in script assert "GEN_STATE_RESET=false" in script assert "up -d --build generator" in script + stop_call_line = _required_runtime_line(script, "stop_live_generator") + step_7_line = _required_runtime_line( + script, + 'echo "Шаг 7: второй batch после live"', + ) + assert stop_call_line < step_7_line assert "CHECK_LIVE_SEAM=1" in script assert "REQUIRE_SUPERSET=0" in script +def test_runtime_gate_contract_rejects_removed_stop_call(): + """Определение функции не скрывает удалённый вызов перед вторым batch.""" + script = ( + REPO_ROOT / "scripts" / "run_generated_history_runtime_check.sh" + ).read_text(encoding="utf-8") + script_without_call = script.replace("\nstop_live_generator\n", "\n", 1) + + assert "stop_live_generator" in script_without_call + with pytest.raises(AssertionError, match="отдельной командой"): + _required_runtime_line(script_without_call, "stop_live_generator") + + +def test_generator_shutdown_grace_covers_slow_batch_and_is_documented(): + """Compose даёт текущему batch минуту на завершение после SIGTERM.""" + compose = COMPOSE_FILE.read_text(encoding="utf-8") + operations = OPERATIONS_DOC.read_text(encoding="utf-8") + generator_service = compose.split("\n generator:\n", 1)[1].split( + "\n kafka-exporter:\n", + 1, + )[0] + + assert "stop_grace_period: 1m" in generator_service + assert "stop_grace_period: 1m" in operations + + def test_generated_history_check_rejects_empty_ods_context(): """ODS-блок стыка падает, если в ODS нет строк для переходящих визитов.""" script = (REPO_ROOT / "scripts" / "check_generated_analytics.sh").read_text(