From d5408f28e9085e360c13260b7d0e7c1f5d874c93 Mon Sep 17 00:00:00 2001 From: Dmitry Dementev Date: Sun, 14 Jun 2026 19:45:58 +0300 Subject: [PATCH] =?UTF-8?q?feat(generator):=20=D0=B4=D0=BE=D0=B1=D0=B0?= =?UTF-8?q?=D0=B2=D0=BB=D0=B5=D0=BD=D0=B0=20=D1=81=D1=82=D0=B0=D1=80=D1=82?= =?UTF-8?q?=D0=BE=D0=B2=D0=B0=D1=8F=20=D0=B8=D1=81=D1=82=D0=BE=D1=80=D0=B8?= =?UTF-8?q?=D1=8F=20=D1=87=D0=B5=D1=80=D0=B5=D0=B7=20backfill?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Зачем: - стенду нужна повторяемая история с живым продолжением от модельной границы без дублей и разрыва визитов. - Что: - добавлен backfill-режим с `GEN_MODEL_T_END`, manifest и state на `T_end`. - live-запуск восстанавливается из manifest без настенной дельты и проверяет совместимость state. - добавлены SQL-проверки формы данных, повторяемости и стыка backfill с live. - Проверка: - make generator-test. - два чистых ClickHouse-прогона backfill дали одинаковые manifest checksums и digest. - reviewer gate issue 05 пройден после исправлений state/manifest. --- ...-startup-history-backfill-to-clickhouse.md | 194 ++++++- docker-compose.yml | 1 + docs/OPERATIONS.md | 282 ++++++++++ ...enerator-model-time-and-startup-history.md | 15 +- generator/README.md | 10 + generator/generator.py | 2 + generator/src/clickstream_generator/config.py | 23 + .../src/clickstream_generator/kafka_io.py | 90 +++- .../src/clickstream_generator/runtime.py | 28 +- .../src/clickstream_generator/service.py | 402 ++++++++++++++- generator/tests/test_service.py | 483 ++++++++++++++++++ 11 files changed, 1509 insertions(+), 21 deletions(-) diff --git a/.scratch/generator-model-time-startup-history/issues/05-startup-history-backfill-to-clickhouse.md b/.scratch/generator-model-time-startup-history/issues/05-startup-history-backfill-to-clickhouse.md index e7d01da..b5cf5fa 100644 --- a/.scratch/generator-model-time-startup-history/issues/05-startup-history-backfill-to-clickhouse.md +++ b/.scratch/generator-model-time-startup-history/issues/05-startup-history-backfill-to-clickhouse.md @@ -1,4 +1,4 @@ -Status: ready-for-agent +Status: ready-for-human # Стартовая история до ClickHouse @@ -19,30 +19,200 @@ Status: ready-for-agent ## Acceptance criteria -- [ ] Промотка прошлого создаёт события за `[T0, T_end)`, слепок состояния и +- [x] Промотка прошлого создаёт события за `[T0, T_end)`, слепок состояния и манифест с контрольными данными. -- [ ] При одинаковых `GEN_SEED`, `T0`, `T_end` и настройках артефакт промотки +- [x] При одинаковых `GEN_SEED`, `T0`, `T_end` и настройках артефакт промотки прошлого повторяем точно; проверки в ClickHouse следуют правилам допуска из контракта задачи 1. -- [ ] История загружается в ClickHouse штатной или явно описанной командой. -- [ ] SQL-проверка показывает здоровую пирамиду: пользователей меньше, чем +- [x] История загружается в ClickHouse штатной или явно описанной командой. +- [x] SQL-проверка показывает здоровую пирамиду: пользователей меньше, чем визитов, визитов меньше, чем событий. -- [ ] SQL-проверка показывает возвраты: у части пользователей больше одного +- [x] SQL-проверка показывает возвраты: у части пользователей больше одного визита. -- [ ] SQL-проверка длины визита проверяет форму, а не только среднее: долю +- [x] SQL-проверка длины визита проверяет форму, а не только среднее: долю коротких визитов, медиану и долю срезов о потолок. -- [ ] SQL-проверка воронки `/home -> товары -> /cart -> /payment -> +- [x] SQL-проверка воронки `/home -> товары -> /cart -> /payment -> /confirmation` монотонно убывает, а доля дошедших до `/confirmation` в согласованном коридоре. -- [ ] Живое продолжение после `T_end` не создаёт дублей на границе и не выглядит +- [x] Живое продолжение после `T_end` не создаёт дублей на границе и не выглядит как независимый второй мир. -- [ ] Визиты, которые переходят через `T_end`, остаются однородными: контекст +- [x] Визиты, которые переходят через `T_end`, остаются однородными: контекст визита не меняется на стыке стартовой истории и живого продолжения. -- [ ] Подготовлены данные, команды и SQL-проверки, достаточные для внешнего +- [x] Подготовлены данные, команды и SQL-проверки, достаточные для внешнего review gate по распределениям и двум путям генерации из `PRD.md`. -- [ ] Worker даёт промежуточный статус, если промотка, загрузка в ClickHouse или +- [x] Worker даёт промежуточный статус, если промотка, загрузка в ClickHouse или распределительные проверки занимают заметное время. ## Blocked by - `.scratch/generator-model-time-startup-history/issues/04-state-v2-model-resume.md` + +## Фактический прогон worker-а + +Дата проверки: 2026-06-14. + +Настройки стенда: + +- `GEN_SEED=4242` +- `GEN_MODEL_T0=2026-01-01T00:00:00+00:00` +- `GEN_MODEL_T_END=2026-01-01T06:00:00+00:00` +- `GEN_MODEL_TIMEZONE=UTC` +- `GEN_MODEL_TIME_SPEED=1` +- `GEN_TICK_SECONDS=60` +- `GEN_LAMBDA_BASE_PER_MIN=120` +- `GEN_JITTER_PCT=0` +- `GEN_MIN_EVENTS_PER_TICK=1` +- `GEN_MAX_EVENTS_PER_TICK=1000` + +Команды проверки: + +```bash +make clean +docker compose up -d clickhouse kafka +make ddl +docker compose build generator + +docker compose run --rm --no-deps \ + -e GEN_RUN_MODE=backfill \ + -e GEN_STATE_RESET=true \ + -e GEN_SEED=4242 \ + -e GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ + -e GEN_MODEL_T_END=2026-01-01T06:00:00+00:00 \ + -e GEN_MODEL_TIMEZONE=UTC \ + -e GEN_MODEL_TIME_SPEED=1 \ + -e GEN_TICK_SECONDS=60 \ + -e GEN_LAMBDA_BASE_PER_MIN=120 \ + -e GEN_JITTER_PCT=0 \ + -e GEN_MIN_EVENTS_PER_TICK=1 \ + -e GEN_MAX_EVENTS_PER_TICK=1000 \ + generator + +sleep 10 +bash scripts/run_batch.sh +``` + +Повторяемость проверена двумя чистыми прогонами с полным сбросом +ClickHouse/Kafka/state через `make clean`. В обоих прогонах manifest и ClickHouse +дали одинаковые контрольные числа: + +- manifest `browser_events.checksum_sha256`: + `3b1946bd7fd3d669e2461fcbae88d8aed6980bb5f126779fb4e7480ce15a3269` +- manifest `location_events.checksum_sha256`: + `370816c8d0b0311592282b9b75372e18862db9ef4984f7032885e935f9436522` +- manifest `device_events.checksum_sha256`: + `37e157ac7d5a3203c1c27b2a11abc1b9cb9fc7f74f79c77645faa09412e08b33` +- manifest `geo_events.checksum_sha256`: + `ec0e36e4613e05b9a5e72490c7565cf7e5f767d39a7a58b383e56e2eb191fc20` +- ClickHouse digest: + `CC1A73E65E897D1F1FC982CFA4237A07` + +Backfill в ClickHouse: + +- `events=31825` +- `unique_events=31825` +- `visits=3020` +- `users=1048` +- `min_event_ts=2026-01-01 00:00:00.000000` +- `max_event_ts=2026-01-01 05:59:59.521599` +- `pyramid_ok=1` +- `half_open_ok=1` + +Проверка возвратов: + +- `users=1048` +- `returning_users=708` +- `returning_share=0.6755725190839694` +- `max_visits_per_user=11` + +Форма длины визита: + +- `visits=3020` +- `short_visit_share=0.16490066225165562` +- `median_events_per_visit=8` +- `avg_events_per_visit=10.538079470198676` +- `capped_visit_share=0.0619205298013245` +- `median_duration_sec=202` +- `p95_duration_sec=823` +- `max_events_per_visit=30` + +Воронка: + +- строгая упорядоченная: + `home=2743`, `products=1626`, `cart=995`, `payment=673`, + `confirmation=394`, `monotonic_ok=1`, + `ordered_confirmation_share=0.14363835216915785`; +- калибровочная по наличию страницы в визите: + `home=2743`, `products=2706`, `cart=2003`, `payment=1368`, + `confirmation=824`, `monotonic_ok=1`, + `confirmation_share_all_visits=0.2728476821192053`, + `confirmation_share_from_home=0.3004010207801677`. + +Live-продолжение: + +```bash +docker compose run -d --name issue05-live --no-deps \ + -e GEN_RUN_MODE=live \ + -e GEN_STATE_RESET=false \ + -e GEN_SEED=4242 \ + -e GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ + -e GEN_MODEL_T_END=2026-01-01T06:00:00+00:00 \ + -e GEN_MODEL_TIMEZONE=UTC \ + -e GEN_MODEL_TIME_SPEED=1 \ + -e GEN_TICK_SECONDS=60 \ + -e GEN_LAMBDA_BASE_PER_MIN=120 \ + -e GEN_JITTER_PCT=0 \ + -e GEN_MIN_EVENTS_PER_TICK=1 \ + -e GEN_MAX_EVENTS_PER_TICK=1000 \ + generator +``` + +Лог live-запуска подтвердил восстановление из стартовой истории: +`model_time=2026-01-01T06:00:00+00:00`, затем прошли тики `361`-`364`. + +После live и повторного `bash scripts/run_batch.sh`: + +- `users=1079` +- `visits=3069` +- `events=32145` +- `min_event_ts=2026-01-01 00:00:00.000000` +- `max_event_ts=2026-01-01 06:03:00.000000` +- `history_events=31825` +- `live_events=320` +- `duplicate_events=0` +- `boundary_events=13` + +`boundary_events=13` — это live-события ровно на `T_end`; они не входили в +backfill `[T0, T_end)`. + +Визиты через `T_end`: + +- `crossing_visits=36` +- `visits_with_both_sides=36` +- `homogeneous_visits=36` +- `context_ok=1` +- `min_before_events=2` +- `min_after_events=1` +- `max_after_events=10` + +## Риски для review gate + +- Manifest, state и события связаны метаданными (`GEN_SEED`, `T0`, `T_end`, + настройки генерации, `last_batch_id`) и checksum событий по топикам. Полной + криптографической связки `state + events + manifest` пока нет. +- Для воронки есть две SQL-формы. Строгая ordered-форма проверяет порядок + `/home -> товары -> /cart -> /payment -> /confirmation`. Калибровочная + contains-форма проверяет, была ли страница в визите. Для review gate + `confirmation_share` сравнивается с коридором мат-модели по contains-форме. + +## Review gate + +- Саморевью worker-а нашло блокирующий риск: backfill не должен сохранять + state/manifest при частичной публикации. Исправлено fail-fast поведением и + тестом. +- Первый reviewer gate нашёл два блокера: startup-history state мог остаться без + manifest и восстановиться как live-state, а startup-history restore не сверял + поля самого state. Оба пункта исправлены до коммита. +- Повторный reviewer gate: `gate pass`, новых блокирующих находок нет. +- Проверки после исправлений: `make generator-test` — 134 passed. +- Остаточный риск процесса: reviewer был со свежим контекстом, но той же + родословной, а не внешней моделью. diff --git a/docker-compose.yml b/docker-compose.yml index 9570ef7..4e4c98b 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -303,6 +303,7 @@ services: GEN_P_NEW_USER: ${GEN_P_NEW_USER:-0.15} GEN_MIN_RETURN_MINUTES: ${GEN_MIN_RETURN_MINUTES:-30} GEN_MODEL_T0: ${GEN_MODEL_T0:-2026-01-01T00:00:00+00:00} + GEN_MODEL_T_END: ${GEN_MODEL_T_END:-} GEN_MODEL_TIMEZONE: ${GEN_MODEL_TIMEZONE:-UTC} GEN_MODEL_TIME_SPEED: ${GEN_MODEL_TIME_SPEED:-1} GEN_RUN_MODE: ${GEN_RUN_MODE:-live} diff --git a/docs/OPERATIONS.md b/docs/OPERATIONS.md index 2cb4684..dd7c1bf 100644 --- a/docs/OPERATIONS.md +++ b/docs/OPERATIONS.md @@ -103,6 +103,7 @@ make generator-logs | `GEN_P_NEW_USER` | Доля визитов новых пользователей | `0.15` | | `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | | `GEN_MODEL_T0` | Стартовая модельная точка, ISO 8601 с часовым поясом | `2026-01-01T00:00:00+00:00` | +| `GEN_MODEL_T_END` | Правая граница стартовой истории для `backfill` | пусто | | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | | `GEN_RUN_MODE` | Режим генератора | `live` | @@ -119,6 +120,287 @@ GEN_STATE_RESET=true GEN_LAMBDA_BASE_PER_MIN=60 docker compose up -d generator Контейнерные `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены внутренними значениями `kafka:29092` и `/data`. +### Стартовая история через backfill + +`GEN_RUN_MODE=backfill` быстро проматывает модельное прошлое от `GEN_MODEL_T0` +до `GEN_MODEL_T_END` без сна. В Kafka попадают события только за полуоткрытый +отрезок `[T0, T_end)`. В compact-topic `generator_state` сохраняется state на +`T_end`, а в `generator_startup_history_manifest` — manifest с настройками и +контрольными числами. При live-запуске с теми же настройками генератор видит, +что state совпадает с manifest, и стартует ровно с `T_end` без настенной дельты. + +Для чистого повтора проще всего пересоздать volumes. Это сбрасывает ClickHouse, +Kafka-топики данных и compact-topic state. + +```bash +make clean +docker compose up -d clickhouse kafka +make ddl + +GEN_RUN_MODE=backfill \ +GEN_STATE_RESET=true \ +GEN_SEED=4242 \ +GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ +GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 \ +GEN_MODEL_TIMEZONE=UTC \ +GEN_MODEL_TIME_SPEED=1 \ +GEN_TICK_SECONDS=60 \ +GEN_LAMBDA_BASE_PER_MIN=60 \ +GEN_JITTER_PCT=0 \ +docker compose run --rm --no-deps generator + +# Materialized View слоя STG читает Kafka сама; даём ей коротко догнать. +sleep 10 +bash scripts/run_batch.sh +``` + +Manifest можно посмотреть так: + +```bash +docker compose exec -T kafka /opt/kafka/bin/kafka-console-consumer.sh \ + --bootstrap-server kafka:29092 \ + --topic generator_startup_history_manifest \ + --from-beginning \ + --property print.key=true \ + --timeout-ms 5000 +``` + +Live-продолжение стартует с этого state. Используйте те же `GEN_SEED`, `T0`, +`T_end`, часовой пояс и настройки генерации. `GEN_STATE_RESET=false` важен: иначе +слепок стартовой истории будет проигнорирован. + +```bash +docker compose run -d --name startup-history-live --no-deps \ + -e GEN_RUN_MODE=live \ + -e GEN_STATE_RESET=false \ + -e GEN_SEED=4242 \ + -e GEN_MODEL_T0=2026-01-01T00:00:00+00:00 \ + -e GEN_MODEL_T_END=2026-01-02T00:00:00+00:00 \ + -e GEN_MODEL_TIMEZONE=UTC \ + -e GEN_MODEL_TIME_SPEED=1 \ + -e GEN_TICK_SECONDS=60 \ + -e GEN_LAMBDA_BASE_PER_MIN=60 \ + -e GEN_JITTER_PCT=0 \ + generator + +sleep 130 +docker stop startup-history-live +docker rm startup-history-live +sleep 10 +bash scripts/run_batch.sh +``` + +Базовая сверка формы данных после backfill: + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end +SELECT + uniqExact(user_domain_id) AS users, + uniqExact(click_id) AS visits, + count() AS events, + users < visits AND visits < events AS pyramid_ok, + min(event_ts) AS min_event_ts, + max(event_ts) AS max_event_ts, + min_event_ts >= t0 AND max_event_ts < t_end AS half_open_ok +FROM dm.v_events_enriched +WHERE event_ts >= t0 AND event_ts < t_end; +``` + +Повторяемость чистого прогона удобно сверять коротким digest по ключевым полям +ClickHouse. Запускайте запрос после `bash scripts/run_batch.sh`; при одинаковых +`GEN_SEED`, `T0`, `T_end` и настройках значение должно повторяться. + +```bash +docker compose exec -T clickhouse clickhouse-client \ + --user=default \ + --password=123456 \ + --query " +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end +SELECT hex(sipHash128(groupArray(tuple( + event_id, + click_id, + user_domain_id, + event_ts, + page_url_path +)))) AS digest +FROM ( + SELECT + event_id, + click_id, + user_domain_id, + event_ts, + page_url_path + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + ORDER BY event_id + )" +``` + +Возвраты пользователей: + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + users AS ( + SELECT user_domain_id, uniqExact(click_id) AS visits + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + AND user_domain_id IS NOT NULL + GROUP BY user_domain_id + ) +SELECT + count() AS users, + countIf(visits > 1) AS returning_users, + returning_users / users AS returning_share +FROM users; +``` + +Форма длины визита: проверяем не только среднее, а долю коротких визитов, +медиану и долю визитов, срезанных потолком `GEN_MAX_SESSION_EVENTS`. + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + 30 AS max_session_events, + sessions AS ( + SELECT + click_id, + count() AS events_count, + dateDiff('second', min(event_ts), max(event_ts)) AS duration_sec + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + count() AS visits, + countIf(events_count <= 2) / visits AS short_visit_share, + quantileExact(0.5)(events_count) AS median_events_per_visit, + avg(events_count) AS avg_events_per_visit, + countIf(events_count = max_session_events) / visits AS capped_visit_share, + quantileExact(0.5)(duration_sec) AS median_duration_sec +FROM sessions; +``` + +Воронка должна монотонно убывать, а доля дошедших до `/confirmation` должна быть +в согласованном коридоре для текущих настроек генератора. Для review gate +`confirmation_share` сравнивается по калибровочной форме «страница была в +визите». Строгий SQL ниже проверяет отдельное свойство: упорядоченный путь +`/home -> товары -> /cart -> /payment -> /confirmation`. + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + sessions AS ( + SELECT + click_id, + minIf(event_ts, page_url_path = '/home') AS home_ts, + minIf(event_ts, page_url_path IN ('/product_a', '/product_b')) AS product_ts, + minIf(event_ts, page_url_path = '/cart') AS cart_ts, + minIf(event_ts, page_url_path = '/payment') AS payment_ts, + minIf(event_ts, page_url_path = '/confirmation') AS confirmation_ts + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + countIf(home_ts IS NOT NULL) AS home, + countIf(home_ts IS NOT NULL AND product_ts > home_ts) AS products, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts) AS cart, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts AND payment_ts > cart_ts) AS payment, + countIf(home_ts IS NOT NULL AND product_ts > home_ts AND cart_ts > product_ts AND payment_ts > cart_ts AND confirmation_ts > payment_ts) AS confirmation, + products <= home AND cart <= products AND payment <= cart AND confirmation <= payment AS monotonic_ok, + confirmation / home AS confirmation_share +FROM sessions; +``` + +Калибровочная форма воронки проверяет, что страница была в визите, без строгого +порядка событий. Именно эту форму используем для сравнения `confirmation_share` +в review gate. + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + sessions AS ( + SELECT + click_id, + countIf(page_url_path = '/home') > 0 AS has_home, + countIf(page_url_path IN ('/product_a', '/product_b')) > 0 AS has_product, + countIf(page_url_path = '/cart') > 0 AS has_cart, + countIf(page_url_path = '/payment') > 0 AS has_payment, + countIf(page_url_path = '/confirmation') > 0 AS has_confirmation + FROM dm.v_events_enriched + WHERE event_ts >= t0 AND event_ts < t_end + GROUP BY click_id + ) +SELECT + countIf(has_home) AS home, + countIf(has_home AND has_product) AS products, + countIf(has_home AND has_product AND has_cart) AS cart, + countIf(has_home AND has_product AND has_cart AND has_payment) AS payment, + countIf(has_home AND has_product AND has_cart AND has_payment AND has_confirmation) AS confirmation, + products <= home AND cart <= products AND payment <= cart AND confirmation <= payment AS monotonic_ok, + confirmation / home AS confirmation_share +FROM sessions; +``` + +Стык backfill + live проверяется после короткого live-продолжения: + +```sql +WITH + toDateTime64('2026-01-01 00:00:00', 6) AS t0, + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + toDateTime64('2026-01-02 00:10:00', 6) AS t_live_end +SELECT + count() AS events, + uniqExact(event_id) AS unique_events, + events - unique_events AS duplicate_events, + countIf(event_ts = t_end) AS boundary_events, + min(event_ts) AS min_event_ts, + max(event_ts) AS max_event_ts +FROM dm.v_events_enriched +WHERE event_ts >= t0 AND event_ts < t_live_end; +``` + +Однородность визитов, переходящих через `T_end`: + +```sql +WITH + toDateTime64('2026-01-02 00:00:00', 6) AS t_end, + crossing AS ( + SELECT + click_id, + min(event_ts) AS first_ts, + max(event_ts) AS last_ts, + groupUniqArray(user_domain_id) AS users, + groupUniqArray(device_type) AS devices, + groupUniqArray(os_name) AS os_names, + groupUniqArray(geo_country) AS countries + FROM dm.v_events_enriched + WHERE event_ts >= t_end - INTERVAL 30 MINUTE + AND event_ts < t_end + INTERVAL 30 MINUTE + GROUP BY click_id + HAVING first_ts < t_end AND last_ts >= t_end + ) +SELECT + count() AS crossing_visits, + countIf( + length(users) = 1 + AND length(devices) = 1 + AND length(os_names) = 1 + AND length(countries) = 1 + ) AS homogeneous_visits, + crossing_visits = homogeneous_visits AS context_ok +FROM crossing; +``` + ### Проверка модельного времени в ClickHouse Для повторяемой проверки используйте чистый стенд и явный сброс состояния diff --git a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md index ba860e4..5d89088 100644 --- a/docs/specs/2026-06-14-generator-model-time-and-startup-history.md +++ b/docs/specs/2026-06-14-generator-model-time-and-startup-history.md @@ -165,22 +165,29 @@ resume_model_at = #### Манифест стартовой истории Стартовая история состоит из трёх частей: события, слепок состояния и манифест. -Манифест хранится как JSON и минимум содержит: +Манифест хранится как JSON в Kafka compact-topic +`generator_startup_history_manifest`, ключ `default`. Слепок состояния хранится +в `generator_state`, ключ `default`. Манифест минимум содержит: - `manifest_version`; - `generated_at` — настенная UTC-метка создания артефакта; - `gen_seed`; - `model_t0`, `model_t_end`, `model_timezone`; - `run_mode = "backfill"`; -- настройки генерации, влияющие на поток; -- `state_version` и ссылку на файл или запись слепка состояния; +- настройки генерации, влияющие на поток, в `generation_settings`; +- `state_version` и ссылку на запись слепка состояния; - контрольные числа по каждому топику: количество строк, минимум и максимум `event_timestamp`, контрольная сумма; - итоговые контрольные числа для проверки в ClickHouse: события, визиты, пользователи и диапазон модельного времени. Нельзя смешивать события, state и манифест от разных `GEN_SEED`, `T0`, `T_end` -или настроек генерации. Такое смешивание считается ошибкой запуска. +или настроек генерации. Live-запуск использует manifest как стартовую историю +только если `state.last_batch_id`, `state.model_timestamp`, `GEN_SEED`, +`GEN_MODEL_T0`, `GEN_MODEL_T_END`, `GEN_MODEL_TIMEZONE`, `GEN_MODEL_TIME_SPEED` +и `generation_settings` совпадают. Startup-history state без подходящего +manifest считается несовместимым и ведёт к чистому старту, а не к восстановлению +по правилу live-сбоя. #### Повторяемая проверка в ClickHouse diff --git a/generator/README.md b/generator/README.md index 17d61c7..0c862e5 100644 --- a/generator/README.md +++ b/generator/README.md @@ -73,6 +73,7 @@ generator-service -> Kafka topics -> (потребители отдельно) | `GEN_P_NEW_USER` | Вероятность отдать новый визит новому пользователю | `0.15` | | `GEN_MIN_RETURN_MINUTES` | Минимальная пауза перед возвратом пользователя | `30` | | `GEN_MODEL_T0` | Стартовая модельная точка, ISO 8601 с часовым поясом | `2026-01-01T00:00:00+00:00` | +| `GEN_MODEL_T_END` | Правая граница стартовой истории для `backfill` | — | | `GEN_MODEL_TIMEZONE` | Часовой пояс модельных часов для дневного коэффициента | `UTC` | | `GEN_MODEL_TIME_SPEED` | Сколько модельных секунд проходит за одну настенную секунду | `1` | | `GEN_RUN_MODE` | Режим генератора | `live` | @@ -97,6 +98,15 @@ GEN_LAMBDA_BASE_PER_MIN=60 GEN_POPULATION_MAX=500 docker compose up -d generator реальный прогон. Дневной коэффициент считается по `GEN_MODEL_TIMEZONE`, а не по реальному часу запуска процесса. +В режиме `backfill` генератор без сна проходит от `GEN_MODEL_T0` до +`GEN_MODEL_T_END`, публикует события только за `[T0, T_end)`, сохраняет state v2 +на `T_end` в `generator_state` и пишет manifest в compact-topic +`generator_startup_history_manifest`. Live-запуск с теми же настройками +использует этот manifest, чтобы продолжить ровно с `T_end` без настенной дельты. +Для одноразового backfill-запуска через compose используйте `docker compose run +--rm generator`, а не `docker compose up generator`: у штатного сервиса включён +restart policy. + Контейнерные значения `KAFKA_BOOTSTRAP_SERVERS` и `GEN_DATA_DIR` в compose оставлены безопасными внутренними значениями `kafka:29092` и `/data`. diff --git a/generator/generator.py b/generator/generator.py index ad460a1..2952f7c 100644 --- a/generator/generator.py +++ b/generator/generator.py @@ -24,6 +24,7 @@ from clickstream_generator.kafka_io import ( KafkaBatchHistory, KafkaPublisher, KafkaStateManager, + KafkaStartupHistoryManifest, _import_kafka, _with_retry, ensure_topics, @@ -56,6 +57,7 @@ __all__ = [ "KafkaBatchHistory", "KafkaPublisher", "KafkaStateManager", + "KafkaStartupHistoryManifest", "METRICS_ERRORS_TOTAL", "METRICS_EVENTS_TOTAL", "METRICS_LAST_SUCCESS", diff --git a/generator/src/clickstream_generator/config.py b/generator/src/clickstream_generator/config.py index d53c983..aaf33b0 100644 --- a/generator/src/clickstream_generator/config.py +++ b/generator/src/clickstream_generator/config.py @@ -78,6 +78,11 @@ class Config: os.getenv("GEN_MODEL_T0", "2026-01-01T00:00:00+00:00") ) ) + model_t_end: datetime | None = field( + default_factory=lambda: _parse_optional_model_timestamp( + os.getenv("GEN_MODEL_T_END") + ) + ) model_timezone: str = field( default_factory=lambda: os.getenv("GEN_MODEL_TIMEZONE", "UTC") ) @@ -122,3 +127,21 @@ class Config: 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.model_t_end is not None: + object.__setattr__( + self, + "model_t_end", + self.model_t_end.astimezone(timezone.utc), + ) + if self.run_mode == "backfill": + if self.model_t_end is None: + 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") + + +def _parse_optional_model_timestamp(value: str | None) -> datetime | None: + """Разбирает необязательную ISO-метку модельного времени.""" + if not value: + return None + return _parse_model_timestamp(value) diff --git a/generator/src/clickstream_generator/kafka_io.py b/generator/src/clickstream_generator/kafka_io.py index 103cb20..1774211 100644 --- a/generator/src/clickstream_generator/kafka_io.py +++ b/generator/src/clickstream_generator/kafka_io.py @@ -181,6 +181,84 @@ class KafkaStateManager: return None +class KafkaStartupHistoryManifest: + """Хранение манифеста стартовой истории в Kafka compact topic.""" + + MANIFEST_TOPIC = "generator_startup_history_manifest" + MANIFEST_KEY = "default" + + def __init__(self, bootstrap_servers: str): + self.bootstrap_servers = bootstrap_servers + KafkaProducerCls, _ = _kafka_importer()() + + logger.info( + f"Connecting to Kafka for startup history manifest at {self.bootstrap_servers}" + ) + self.producer = KafkaProducerCls( + bootstrap_servers=self.bootstrap_servers, + value_serializer=lambda v: json.dumps(v).encode("utf-8"), + key_serializer=lambda k: k.encode("utf-8") if k else None, + retries=3, + retry_backoff_ms=1000, + ) + logger.info("Connected to Kafka for startup history manifest successfully") + + def save(self, manifest: dict) -> None: + """Сохраняет манифест стартовой истории.""" + def _do_send(): + self.producer.send( + self.MANIFEST_TOPIC, + key=self.MANIFEST_KEY, + value=manifest, + ) + + _retry(_do_send, max_retries=3, base_delay=0.5) + + def flush(self) -> None: + """Сбрасывает буфер с retry.""" + def _do_flush(): + self.producer.flush() + + _retry(_do_flush, max_retries=3, base_delay=0.5) + + def close(self) -> None: + """Закрывает соединение.""" + try: + self.producer.close() + except Exception as e: + logger.debug(f"Error closing manifest producer (ignored): {e}") + + def load(self) -> dict | None: + """Загружает последний манифест стартовой истории.""" + from kafka import KafkaConsumer + + logger.info(f"Loading startup history manifest from topic {self.MANIFEST_TOPIC}") + + def _do_load(): + consumer = KafkaConsumer( + self.MANIFEST_TOPIC, + bootstrap_servers=self.bootstrap_servers, + auto_offset_reset="earliest", + enable_auto_commit=False, + consumer_timeout_ms=5000, + value_deserializer=lambda v: json.loads(v.decode("utf-8")), + ) + + last_manifest = None + for message in consumer: + if message.key and message.key.decode("utf-8") == self.MANIFEST_KEY: + last_manifest = message.value + + consumer.close() + return last_manifest + + try: + return _retry(_do_load, max_retries=3, base_delay=0.5) + except Exception as e: + logger.warning(f"Failed to load startup history manifest: {e}") + return None + + def ensure_topics(bootstrap_servers: str) -> None: """Создаёт служебные топики, если их ещё нет.""" from kafka import KafkaAdminClient @@ -205,8 +283,18 @@ def ensure_topics(bootstrap_servers: str) -> None: "delete.retention.ms": "100", }, ) + manifest_topic = NewTopic( + name=KafkaStartupHistoryManifest.MANIFEST_TOPIC, + num_partitions=1, + replication_factor=1, + topic_configs={ + "cleanup.policy": "compact", + "min.cleanable.dirty.ratio": "0.1", + "delete.retention.ms": "100", + }, + ) - for topic in [history_topic, state_topic]: + for topic in [history_topic, state_topic, manifest_topic]: try: admin_client.create_topics([topic]) logger.info(f"Created topic: {topic.name}") diff --git a/generator/src/clickstream_generator/runtime.py b/generator/src/clickstream_generator/runtime.py index b11bab8..3072c14 100644 --- a/generator/src/clickstream_generator/runtime.py +++ b/generator/src/clickstream_generator/runtime.py @@ -394,6 +394,24 @@ class TickStreamGenerator: return tick_batch + def drain_until( + self, + cutoff_at: datetime, + include_boundary: bool = False, + ) -> dict[str, list[dict]]: + """Выпускает созревшие события без рождения новых визитов.""" + cutoff_time = _normalize_tick_time(cutoff_at) + tick_batch = _empty_batch() + + self._release_due_events( + cutoff_time, + tick_batch, + include_boundary=include_boundary, + ) + self._drop_finished_visits() + + return tick_batch + def _birth_visits(self, event_budget: int, tick_time: datetime) -> None: if len(self.active_visits) >= self.generator.config.max_active_sessions: self._pending_visit_births = 0.0 @@ -436,9 +454,17 @@ class TickStreamGenerator: self, tick_time: datetime, tick_batch: dict[str, list[dict]], + include_boundary: bool = True, ) -> None: for visit in self.active_visits: - while not visit.is_finished and visit.timestamps[visit.next_index] <= tick_time: + while not visit.is_finished: + next_timestamp = visit.timestamps[visit.next_index] + if include_boundary: + is_due = next_timestamp <= tick_time + else: + is_due = next_timestamp < tick_time + if not is_due: + break event_index = visit.next_index for topic in TOPICS: tick_batch[topic].append(visit.batch[topic][event_index]) diff --git a/generator/src/clickstream_generator/service.py b/generator/src/clickstream_generator/service.py index c7d8fe4..270e34f 100644 --- a/generator/src/clickstream_generator/service.py +++ b/generator/src/clickstream_generator/service.py @@ -1,5 +1,7 @@ """Основной сервисный цикл генератора.""" +import hashlib +import json import logging import sys import time @@ -16,6 +18,7 @@ from clickstream_generator.kafka_io import ( KafkaBatchHistory, KafkaPublisher, KafkaStateManager, + KafkaStartupHistoryManifest, ensure_topics, ) from clickstream_generator.metrics import ( @@ -29,6 +32,105 @@ from clickstream_generator.runtime import TickStreamGenerator logger = logging.getLogger("generator") +class _TopicManifestStats: + """Накопительные счётчики одного Kafka-топика для manifest.""" + + def __init__(self): + self.rows = 0 + self.min_event_timestamp: str | None = None + self.max_event_timestamp: str | None = None + self._checksum = hashlib.sha256() + + @property + def checksum(self) -> str: + return self._checksum.hexdigest() + + def add(self, event: dict, event_timestamp: str | None = None) -> None: + self.rows += 1 + self._checksum.update( + json.dumps(event, sort_keys=True, ensure_ascii=True).encode("utf-8") + ) + timestamp = event_timestamp if event_timestamp is not None else event.get("event_timestamp") + self.add_timestamp(timestamp) + + def add_timestamp(self, timestamp: str | None) -> None: + if timestamp is None: + return + if self.min_event_timestamp is None or timestamp < self.min_event_timestamp: + self.min_event_timestamp = timestamp + if self.max_event_timestamp is None or timestamp > self.max_event_timestamp: + self.max_event_timestamp = timestamp + + def to_dict(self) -> dict: + return { + "rows": self.rows, + "min_event_timestamp": self.min_event_timestamp, + "max_event_timestamp": self.max_event_timestamp, + "checksum_sha256": self.checksum, + } + + +class _ManifestCounters: + """Счётчики стартовой истории для manifest.""" + + def __init__(self): + self.topic_stats = { + topic: _TopicManifestStats() + for topic in ("browser_events", "location_events", "device_events", "geo_events") + } + self.click_ids: set[str] = set() + self.user_ids: set[str] = set() + + def add_batch(self, batch: dict[str, list[dict]]) -> None: + browser_events = batch.get("browser_events", []) + event_timestamps = { + event["event_id"]: event.get("event_timestamp") + for event in browser_events + if event.get("event_id") and event.get("event_timestamp") + } + click_timestamps: dict[str, list[str]] = {} + for event in browser_events: + click_id = event.get("click_id") + timestamp = event.get("event_timestamp") + if click_id and timestamp: + click_timestamps.setdefault(click_id, []).append(timestamp) + + for topic, events in batch.items(): + stats = self.topic_stats[topic] + for event in events: + event_timestamp = event.get("event_timestamp") + if event_timestamp is None and topic == "location_events": + event_timestamp = event_timestamps.get(event.get("event_id")) + if event_timestamp is None and topic in ("device_events", "geo_events"): + timestamps = click_timestamps.get(event.get("click_id")) + if timestamps: + event_timestamp = min(timestamps) + stats.add_timestamp(max(timestamps)) + stats.add(event, event_timestamp=event_timestamp) + click_id = event.get("click_id") + if topic == "browser_events" and click_id: + self.click_ids.add(click_id) + user_id = event.get("user_domain_id") + if topic == "device_events" and user_id: + self.user_ids.add(user_id) + + def to_manifest_topics(self) -> dict: + return { + topic: stats.to_dict() + for topic, stats in self.topic_stats.items() + } + + def to_manifest_totals(self) -> dict: + browser_stats = self.topic_stats["browser_events"] + return { + "events": browser_stats.rows, + "visits": len(self.click_ids), + "users": len(self.user_ids), + "min_event_timestamp": browser_stats.min_event_timestamp, + "max_event_timestamp": browser_stats.max_event_timestamp, + } + + class GeneratorService: """Основной сервис генератора.""" @@ -40,6 +142,7 @@ class GeneratorService: self.publisher: KafkaPublisher | None = None self.history: KafkaBatchHistory | None = None self.state_manager: KafkaStateManager | None = None + self.manifest_manager: KafkaStartupHistoryManifest | None = None self._running = False self._tick = 0 self._model_time = config.model_t0 @@ -66,17 +169,46 @@ class GeneratorService: self.publisher = KafkaPublisher(self.config.kafka_bootstrap_servers) self.history = KafkaBatchHistory(self.config.kafka_bootstrap_servers) + 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.state_enabled: self.state_manager = KafkaStateManager(self.config.kafka_bootstrap_servers) + if self.config.run_mode == "backfill": + self.manifest_manager = KafkaStartupHistoryManifest( + self.config.kafka_bootstrap_servers + ) + logger.info("Backfill mode starts from a fresh generator state") + self._run_backfill() + self.stop() + return + if not self.config.state_reset: restored_state = self.state_manager.load() if restored_state: try: - self._restore_live_state( - restored_state, - wall_now_utc=datetime.now(timezone.utc), + self.manifest_manager = KafkaStartupHistoryManifest( + self.config.kafka_bootstrap_servers ) + manifest = self.manifest_manager.load() + if self._is_startup_history_state(restored_state, manifest): + 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, + ) + elif self._is_startup_history_marker(restored_state): + raise ValueError( + "startup-history state without matching manifest" + ) + else: + self._restore_live_state( + restored_state, + wall_now_utc=datetime.now(timezone.utc), + ) logger.info( f"Restored state: continuing from tick {self._tick}, " f"model_time={self._model_time.isoformat()}, " @@ -113,6 +245,8 @@ class GeneratorService: self.history.close() if self.state_manager: self.state_manager.close() + if self.manifest_manager: + self.manifest_manager.close() def restore_from_startup_history( self, @@ -169,6 +303,54 @@ class GeneratorService: "state config mismatch: " + ", ".join(mismatches) ) + def _is_startup_history_state(self, state, manifest: dict | None) -> bool: + """Проверяет, что state совпадает со слепком стартовой истории.""" + if not manifest or manifest.get("run_mode") != "backfill": + return False + + expected_state = manifest.get("state") or {} + if expected_state.get("last_batch_id") != state.last_batch_id: + return False + + try: + model_t_end = self._as_aware_utc( + datetime.fromisoformat(manifest["model_t_end"]) + ) + model_t0 = self._as_aware_utc( + datetime.fromisoformat(manifest["model_t0"]) + ) + except (KeyError, TypeError, ValueError): + return False + + if self._as_aware_utc(state.model_timestamp) != model_t_end: + return False + if self._as_aware_utc(state.model_t0) != model_t0: + return False + if state.gen_seed != manifest.get("gen_seed"): + return False + if state.model_timezone != manifest.get("model_timezone"): + return False + if ( + self.config.model_t_end is not None + and self.config.model_t_end != model_t_end + ): + return False + if model_t0 != self.config.model_t0: + return False + if manifest.get("gen_seed") != self.config.seed: + return False + if manifest.get("model_timezone") != self.config.model_timezone: + return False + + settings = manifest.get("generation_settings") or {} + if abs(state.model_time_speed - settings.get("model_time_speed", -1)) > 1e-9: + return False + return self._generation_settings() == settings + + def _is_startup_history_marker(self, state) -> bool: + """Отличает state стартовой истории от обычного live-state.""" + return str(state.last_batch_id).startswith("startup-history-") + @staticmethod def _as_aware_utc(value: datetime) -> datetime: if value.tzinfo is None: @@ -204,6 +386,220 @@ class GeneratorService: logger.warning(f"Failed to save state: {e}") METRICS_ERRORS_TOTAL.labels(topic="state").inc() + def _run_backfill(self) -> None: + """Проматывает стартовую историю без сна до GEN_MODEL_T_END.""" + if self.config.model_t_end is None: + raise ValueError("GEN_MODEL_T_END is required for backfill") + if not self.publisher: + raise RuntimeError("publisher is not initialized") + if not self.history: + raise RuntimeError("history is not initialized") + + logger.info( + "Running backfill from %s to %s", + self.config.model_t0.isoformat(), + self.config.model_t_end.isoformat(), + ) + counters = _ManifestCounters() + + while self._model_time < self.config.model_t_end: + self._tick += 1 + batch_id = f"backfill-{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_backfill_publish_error(batch_id, status, sent_counts) + counters.add_batch(batch) + 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 if status == "success" else "Backfill publish error", + ) + self._advance_model_time() + + final_batch = self.stream.drain_until( + self.config.model_t_end, + include_boundary=False, + ) + final_sent = 0 + final_status = "success" + if any(final_batch.values()): + batch_id = f"backfill-{self._tick + 1:08d}-final" + started_at = datetime.now(timezone.utc) + final_sent, sent_counts, final_status = self._publish_batch(final_batch) + self._raise_on_backfill_publish_error( + batch_id, + final_status, + sent_counts, + ) + counters.add_batch(final_batch) + self._write_batch_history( + batch_id=batch_id, + started_at=started_at, + sent_counts=sent_counts, + total_sent=final_sent, + status=final_status, + error_message=None if final_status == "success" else "Backfill publish error", + ) + + if self.publisher: + self.publisher.flush() + + self._model_time = self.config.model_t_end + 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=self.config.model_t_end, + model_timestamp=self.config.model_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, + ) + + manifest = self._build_startup_history_manifest( + counters=counters, + state=state, + ) + if self.manifest_manager: + self.manifest_manager.save(manifest) + self.manifest_manager.flush() + + if self.state_manager and self.config.state_enabled: + self.state_manager.save(state) + self.state_manager.flush() + + logger.info( + "Backfill completed: events=%s, visits=%s, users=%s, final_sent=%s", + manifest["totals"]["events"], + manifest["totals"]["visits"], + manifest["totals"]["users"], + final_sent, + ) + + def _raise_on_backfill_publish_error( + self, + batch_id: str, + status: str, + sent_counts: dict[str, dict[str, int]], + ) -> None: + """Останавливает backfill до записи state/manifest при ошибке Kafka.""" + if status == "success": + return + + raise RuntimeError( + f"Backfill publish failed for batch {batch_id}: " + f"status={status}, sent_counts={sent_counts}" + ) + + def _publish_batch( + self, + batch: dict[str, list[dict]], + ) -> tuple[int, dict[str, dict[str, int]], str]: + """Публикует batch и возвращает счётчики отправки.""" + total_sent = 0 + total_errors = 0 + sent_counts = {} + + for topic, events in batch.items(): + if events: + sent, errors = self.publisher.publish(topic, events) + sent_counts[topic] = {"sent": sent, "errors": errors} + total_sent += sent + total_errors += errors + + if total_errors == 0: + status = "success" + elif total_sent > 0: + status = "partial" + else: + status = "error" + + return total_sent, sent_counts, status + + def _write_batch_history( + self, + batch_id: str, + started_at: datetime, + sent_counts: dict[str, dict[str, int]], + total_sent: int, + status: str, + error_message: str | None, + ) -> None: + """Пишет служебную историю batch.""" + if not self.history: + return + try: + record = BatchRecord( + batch_id=batch_id, + started_at=started_at, + finished_at=datetime.now(timezone.utc), + sent_total=total_sent, + sent_browser=sent_counts.get("browser_events", {}).get("sent", 0), + sent_location=sent_counts.get("location_events", {}).get("sent", 0), + sent_device=sent_counts.get("device_events", {}).get("sent", 0), + sent_geo=sent_counts.get("geo_events", {}).get("sent", 0), + status=status, + error_message=error_message, + ) + self.history.add(record) + self.history.flush() + except Exception as hist_err: + logger.warning(f"Failed to write batch history: {hist_err}") + METRICS_ERRORS_TOTAL.labels(topic="history").inc() + + def _startup_state_batch_id(self, counters) -> str: + digest = counters.topic_stats["browser_events"].checksum[:12] + return f"startup-history-{digest}" + + def _build_startup_history_manifest(self, counters, state) -> dict: + return { + "manifest_version": "1.0", + "generated_at": datetime.now(timezone.utc).isoformat(), + "gen_seed": self.config.seed, + "model_t0": self.config.model_t0.isoformat(), + "model_t_end": self.config.model_t_end.isoformat(), + "model_timezone": self.config.model_timezone, + "run_mode": "backfill", + "generation_settings": self._generation_settings(), + "state_version": state.version, + "state": { + "topic": KafkaStateManager.STATE_TOPIC, + "key": KafkaStateManager.STATE_KEY, + "last_batch_id": state.last_batch_id, + "model_timestamp": state.model_timestamp.isoformat(), + }, + "topics": counters.to_manifest_topics(), + "totals": counters.to_manifest_totals(), + } + + def _generation_settings(self) -> dict: + return { + "tick_seconds": self.config.tick_seconds, + "lambda_base_per_min": self.config.lambda_base_per_min, + "jitter_pct": self.config.jitter_pct, + "min_events_per_tick": self.config.min_events_per_tick, + "max_events_per_tick": self.config.max_events_per_tick, + "max_session_events": self.config.max_session_events, + "max_active_sessions": self.config.max_active_sessions, + "population_max": self.config.population_max, + "p_new_user": self.config.p_new_user, + "min_return_minutes": self.config.min_return_minutes, + "model_time_speed": self.config.model_time_speed, + } + def _main_loop(self): """Основной цикл тиков.""" while self._running: diff --git a/generator/tests/test_service.py b/generator/tests/test_service.py index 9df2eae..c03e515 100644 --- a/generator/tests/test_service.py +++ b/generator/tests/test_service.py @@ -4,6 +4,8 @@ import logging import random +import hashlib +import json from dataclasses import replace from datetime import datetime, timedelta, timezone from time import sleep as real_sleep @@ -301,6 +303,463 @@ class TestGeneratorServiceSteadyStream: assert sum(record.sent_geo for record in history_records) == len(geo_events) +class TestGeneratorServiceBackfill: + """Проверки режима промотки стартовой истории.""" + + def test_backfill_publishes_half_open_history_state_and_manifest( + self, base_config + ): + """Backfill пишет [T0, T_end), state на T_end и повторяемый manifest.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=3) + config = replace( + base_config, + run_mode="backfill", + model_t0=model_t0, + model_t_end=model_t_end, + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + state_enabled=True, + ) + + first = self._run_backfill(config) + second = self._run_backfill(config) + + browser_events = first["published"]["browser_events"] + timestamps = [ + datetime.fromisoformat(event["event_timestamp"].replace(" ", "T")) + for event in browser_events + ] + saved_state = first["state_manager"].save.call_args.args[0] + manifest = first["manifest_manager"].save.call_args.args[0] + + assert browser_events + assert min(timestamps) >= model_t0.replace(tzinfo=None) + assert max(timestamps) < model_t_end.replace(tzinfo=None) + assert saved_state.model_timestamp == model_t_end + assert saved_state.last_timestamp == model_t_end + assert manifest["run_mode"] == "backfill" + assert manifest["model_t0"] == model_t0.isoformat() + assert manifest["model_t_end"] == model_t_end.isoformat() + assert manifest["state"]["last_batch_id"] == saved_state.last_batch_id + assert manifest["topics"]["browser_events"]["rows"] == len(browser_events) + for topic in ( + "browser_events", + "location_events", + "device_events", + "geo_events", + ): + assert manifest["topics"][topic]["min_event_timestamp"] is not None + assert manifest["topics"][topic]["max_event_timestamp"] is not None + assert manifest["totals"]["events"] == len(browser_events) + assert manifest["totals"]["visits"] == len({ + event["click_id"] + for event in browser_events + }) + assert manifest["totals"]["users"] == len({ + event["user_domain_id"] + for event in first["published"]["device_events"] + }) + assert first["digest"] == second["digest"] + assert first["manifest_digest"] == second["manifest_digest"] + + def test_live_start_uses_startup_manifest_without_wall_delta( + self, base_config, event_dictionary + ): + """Live-запуск из backfill-state стартует с T_end без wall-дельты.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=5) + config = replace( + base_config, + model_t0=model_t0, + tick_seconds=60, + model_time_speed=3600, + ) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="startup-history-state", + last_timestamp=model_t_end, + model_timestamp=model_t_end, + wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed, + ) + manifest = { + "run_mode": "backfill", + "gen_seed": config.seed, + "model_t0": model_t0.isoformat(), + "model_t_end": model_t_end.isoformat(), + "model_timezone": config.model_timezone, + "generation_settings": GeneratorService(config)._generation_settings(), + "state": {"last_batch_id": state.last_batch_id}, + } + state_manager = MagicMock() + state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = manifest + + with patch("clickstream_generator.service.start_http_server"), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch("clickstream_generator.service.KafkaPublisher"), \ + patch("clickstream_generator.service.KafkaBatchHistory"), \ + patch( + "clickstream_generator.service.KafkaStateManager", + return_value=state_manager, + ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ + patch.object(GeneratorService, "_main_loop", return_value=None): + + service = GeneratorService(config) + service.start() + + assert service._model_time == model_t_end + assert service._tick == state.tick + + def test_live_start_rejects_startup_manifest_with_different_config_t_end( + self, base_config, event_dictionary + ): + """Manifest от другого T_end не считается стартовой историей запуска.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + manifest_t_end = model_t0 + timedelta(minutes=5) + config_t_end = model_t0 + timedelta(minutes=10) + wall_saved_at = datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc) + wall_restarted_at = wall_saved_at + timedelta(seconds=30) + config = replace( + base_config, + model_t0=model_t0, + model_t_end=config_t_end, + tick_seconds=60, + model_time_speed=10, + ) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="startup-history-state", + last_timestamp=manifest_t_end, + model_timestamp=manifest_t_end, + wall_timestamp=wall_saved_at, + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed, + ) + manifest = { + "run_mode": "backfill", + "gen_seed": config.seed, + "model_t0": model_t0.isoformat(), + "model_t_end": manifest_t_end.isoformat(), + "model_timezone": config.model_timezone, + "generation_settings": GeneratorService(config)._generation_settings(), + "state": {"last_batch_id": state.last_batch_id}, + } + state_manager = MagicMock() + state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = manifest + + class FrozenDateTime(datetime): + @classmethod + def now(cls, tz=None): + if tz is None: + return wall_restarted_at.replace(tzinfo=None) + return wall_restarted_at.astimezone(tz) + + with patch("clickstream_generator.service.start_http_server"), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch("clickstream_generator.service.KafkaPublisher"), \ + patch("clickstream_generator.service.KafkaBatchHistory"), \ + patch( + "clickstream_generator.service.KafkaStateManager", + return_value=state_manager, + ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ + patch("clickstream_generator.service.datetime", FrozenDateTime), \ + patch.object( + GeneratorService, + "restore_from_startup_history", + ) as restore_from_startup_history, \ + patch.object( + GeneratorService, + "_restore_live_state", + ) as restore_live_state, \ + patch.object(GeneratorService, "_main_loop", return_value=None): + + service = GeneratorService(config) + service.start() + + restore_from_startup_history.assert_not_called() + restore_live_state.assert_not_called() + assert service._tick == 0 + assert service._model_time == config.model_t0 + + def test_live_start_rejects_orphan_startup_state_without_live_restore( + self, base_config, event_dictionary, caplog + ): + """Orphan startup-history state не восстанавливается как live-state.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=5) + config = replace(base_config, model_t0=model_t0, model_time_speed=10) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="startup-history-orphan", + last_timestamp=model_t_end, + model_timestamp=model_t_end, + wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed, + ) + state_manager = MagicMock() + state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = None + + with patch("clickstream_generator.service.start_http_server"), \ + patch("clickstream_generator.service.ensure_topics"), \ + patch("clickstream_generator.service.KafkaPublisher"), \ + patch("clickstream_generator.service.KafkaBatchHistory"), \ + patch( + "clickstream_generator.service.KafkaStateManager", + return_value=state_manager, + ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ + patch.object( + GeneratorService, + "restore_from_startup_history", + ) as restore_from_startup_history, \ + patch.object( + GeneratorService, + "_restore_live_state", + ) as restore_live_state, \ + patch.object(GeneratorService, "_main_loop", return_value=None), \ + caplog.at_level(logging.WARNING, logger="generator"): + + service = GeneratorService(config) + service.start() + + restore_from_startup_history.assert_not_called() + restore_live_state.assert_not_called() + assert service._tick == 0 + assert service._model_time == model_t0 + assert "startup-history state without matching manifest" in caplog.text + + def test_startup_history_state_checks_state_fields( + self, base_config, event_dictionary + ): + """Startup-history state сверяется с manifest/config по полям state.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + model_t_end = model_t0 + timedelta(minutes=5) + config = replace(base_config, model_t0=model_t0, model_time_speed=10) + source_generator = EventGenerator(event_dictionary, config) + source_stream = TickStreamGenerator(source_generator) + source_stream.generate_tick(event_budget=10, tick_started_at=model_t0) + state = source_stream.to_state( + tick=5, + rng_state=source_generator.rng.getstate(), + last_batch_id="startup-history-state", + last_timestamp=model_t_end, + model_timestamp=model_t_end, + wall_timestamp=datetime(2026, 1, 1, 0, 0, tzinfo=timezone.utc), + model_time_speed=config.model_time_speed, + model_timezone=config.model_timezone, + model_t0=config.model_t0, + gen_seed=config.seed + 1, + ) + manifest = { + "run_mode": "backfill", + "gen_seed": config.seed, + "model_t0": model_t0.isoformat(), + "model_t_end": model_t_end.isoformat(), + "model_timezone": config.model_timezone, + "generation_settings": GeneratorService(config)._generation_settings(), + "state": {"last_batch_id": state.last_batch_id}, + } + + service = GeneratorService(config) + + assert not service._is_startup_history_state(state, manifest) + + def test_backfill_publish_error_does_not_save_state_or_manifest( + self, base_config + ): + """Backfill не создаёт валидный артефакт при ошибке публикации.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + config = replace( + base_config, + run_mode="backfill", + model_t0=model_t0, + model_t_end=model_t0 + timedelta(minutes=1), + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + state_enabled=True, + ) + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = ( + lambda topic, events: (len(events), 1) + if topic == "location_events" + else (len(events), 0) + ) + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + + with pytest.raises(RuntimeError, match="Backfill publish failed"): + service._run_backfill() + + service.state_manager.save.assert_not_called() + service.state_manager.flush.assert_not_called() + service.manifest_manager.save.assert_not_called() + service.manifest_manager.flush.assert_not_called() + + def test_backfill_manifest_save_error_does_not_save_state(self, base_config): + """Если manifest не записан, state стартовой истории не сохраняется.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + config = replace( + base_config, + run_mode="backfill", + model_t0=model_t0, + model_t_end=model_t0 + timedelta(minutes=1), + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + state_enabled=True, + ) + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = ( + lambda topic, events: (len(events), 0) + ) + service.publisher.flush.return_value = None + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + service.manifest_manager.save.side_effect = RuntimeError("manifest down") + + with pytest.raises(RuntimeError, match="manifest down"): + service._run_backfill() + + service.manifest_manager.save.assert_called_once() + service.state_manager.save.assert_not_called() + service.state_manager.flush.assert_not_called() + + def test_backfill_manifest_flush_error_does_not_save_state(self, base_config): + """Если manifest не сброшен в Kafka, state стартовой истории не сохраняется.""" + model_t0 = datetime(2026, 1, 1, 10, 0, tzinfo=timezone.utc) + config = replace( + base_config, + run_mode="backfill", + model_t0=model_t0, + model_t_end=model_t0 + timedelta(minutes=1), + tick_seconds=60, + lambda_base_per_min=600, + jitter_pct=0, + min_events_per_tick=1, + max_events_per_tick=1000, + max_session_events=5, + max_active_sessions=250, + population_max=251, + state_enabled=True, + ) + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = ( + lambda topic, events: (len(events), 0) + ) + service.publisher.flush.return_value = None + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + service.manifest_manager.flush.side_effect = RuntimeError("flush down") + + with pytest.raises(RuntimeError, match="flush down"): + service._run_backfill() + + service.manifest_manager.save.assert_called_once() + service.manifest_manager.flush.assert_called_once() + service.state_manager.save.assert_not_called() + service.state_manager.flush.assert_not_called() + + def _run_backfill(self, config): + service = GeneratorService(config) + service.publisher = MagicMock() + service.publisher.publish.side_effect = ( + lambda topic, events: (len(events), 0) + ) + service.publisher.flush.return_value = None + service.history = MagicMock() + service.state_manager = MagicMock() + service.manifest_manager = MagicMock() + + service._run_backfill() + + published = {} + for call in service.publisher.publish.call_args_list: + topic, events = call.args + published.setdefault(topic, []).extend(events) + + digest = hashlib.sha256( + json.dumps(published, sort_keys=True, default=str).encode() + ).hexdigest() + manifest = service.manifest_manager.save.call_args.args[0] + stable_manifest = { + key: value + for key, value in manifest.items() + if key != "generated_at" + } + manifest_digest = hashlib.sha256( + json.dumps(stable_manifest, sort_keys=True, default=str).encode() + ).hexdigest() + + return { + "published": published, + "state_manager": service.state_manager, + "manifest_manager": service.manifest_manager, + "digest": digest, + "manifest_digest": manifest_digest, + } + + class TestGeneratorServiceStateV2: """Тесты подключения state v2 к сервисному запуску.""" @@ -325,6 +784,8 @@ class TestGeneratorServiceStateV2: state_manager = MagicMock() state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = None with patch("clickstream_generator.service.start_http_server"), \ patch("clickstream_generator.service.ensure_topics"), \ @@ -334,6 +795,10 @@ class TestGeneratorServiceStateV2: "clickstream_generator.service.KafkaStateManager", return_value=state_manager, ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ patch.object(GeneratorService, "_main_loop", return_value=None): service = GeneratorService(base_config) @@ -373,6 +838,8 @@ class TestGeneratorServiceStateV2: ) state_manager = MagicMock() state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = None class FrozenDateTime(datetime): @classmethod @@ -389,6 +856,10 @@ class TestGeneratorServiceStateV2: "clickstream_generator.service.KafkaStateManager", return_value=state_manager, ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ patch("clickstream_generator.service.datetime", FrozenDateTime), \ patch.object(GeneratorService, "_main_loop", return_value=None): @@ -480,6 +951,8 @@ class TestGeneratorServiceStateV2: ) state_manager = MagicMock() state_manager.load.return_value = state + manifest_manager = MagicMock() + manifest_manager.load.return_value = None with patch("clickstream_generator.service.start_http_server"), \ patch("clickstream_generator.service.ensure_topics"), \ @@ -489,6 +962,10 @@ class TestGeneratorServiceStateV2: "clickstream_generator.service.KafkaStateManager", return_value=state_manager, ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ patch.object(GeneratorService, "_main_loop", return_value=None), \ caplog.at_level(logging.WARNING, logger="generator"): @@ -542,6 +1019,8 @@ class TestGeneratorServiceStateV2: ], active_visits=[], ) + manifest_manager = MagicMock() + manifest_manager.load.return_value = None with patch("clickstream_generator.service.start_http_server"), \ patch("clickstream_generator.service.ensure_topics"), \ @@ -551,6 +1030,10 @@ class TestGeneratorServiceStateV2: "clickstream_generator.service.KafkaStateManager", return_value=state_manager, ), \ + patch( + "clickstream_generator.service.KafkaStartupHistoryManifest", + return_value=manifest_manager, + ), \ patch.object(GeneratorService, "_main_loop", return_value=None), \ caplog.at_level(logging.WARNING, logger="generator"):