Status: ready-for-agent # Airflow-DAG — пульт управления генератором ## Parent `.scratch/generator-model-time-startup-history/PRD.md` ## Why Даже с runbook и глаголами управление генератором остаётся консольным. Идея (пользователь, 2026-06-14; приоритет поднят 2026-07-04): параметризованный DAG в Airflow как основной человеческий интерфейс стенда — форма в веб-UI, где выбираются операция, профиль, длительность. Плюсы: валидация настроек против манифеста ещё до запуска, наглядный статус, понятные ошибки. Учебный бонус: такой DAG сам по себе учебный материал — живой пример параметризованной оркестрации (тема урока 4), что ближе к цели курса, чем устройство генератора. ## Решения (дооформление 2026-07-04) - **Без Docker-доступа из Airflow.** Docker socket в контейнеры Airflow не прокидываем. Конечные операции — это обычный Python-код генератора, которому нужны только настройки и Kafka; таски DAG выполняют его сами, в worker'е Airflow. Урок из проекта `airflow-greenplum-solution` (пользователь): DAG говорит с управляемой системой по сетевому протоколу, жизненный цикл контейнеров — только хост/Makefile. - **Границы пульта.** `make up`, `make clean` (полный сброс) и старт/стоп live-сервиса — консоль; операции `continue` в пульте нет намеренно. Backfill и import работают только на чистом стенде (пустые топики данных); грязный стенд лечится `make clean` с консоли — пульт при отказе прямо подсказывает это. - **Генератор не автостартует.** Сейчас `make up` поднимает live-генератор вместе со стендом (`docker-compose.yml`, `restart: unless-stopped`) — с пультом это неверно: любой backfill был бы сразу отбит предпроверкой, а выключить live из UI нельзя. Live становится явным действием (`make generator-continue`); механизм — например, compose-профиль для сервиса `generator`, выбрать при реализации и поправить make-цели. - **Три операции, а не четыре.** `backfill` | `import` | `check`. Отдельного `export` нет: артефакт в коде пишется только по ходу backfill (`StartupHistoryArtifactBuilder` внутри `_run_backfill`), «выгрузки задним числом» не существует — поэтому у backfill есть поле `artifact_path` («сохранить артефакт», пусто — не сохранять). - **Один DAG** с параметром «операция», а не несколько DAG по операциям. - **После backfill/import DAG сам запускает ETL** (`TriggerDagRunOperator` на существующий ETL-DAG) и **дожидается его завершения** перед check: `wait_for_completion` в Airflow 2.10.5 по умолчанию `False` (сверено по Context7, ревью Codex 2026-07-04) — простой trigger запустил бы check раньше конца ETL. Менти видит цепочку «мир → пайплайн → витрины». ## What to build - **Код генератора доступен таскам Airflow.** Варианты: смонтировать `generator/src` томом (как уже смонтирован `./sql`) и добавить в `PYTHONPATH`, либо ставить пакет в `Dockerfile.airflow`. Монтировать нужно во все Airflow-сервисы (scheduler, webserver, worker): список профилей в форме читается из `PROFILES` при разборе DAG. Критерий выбора — правка генератора не должна требовать лишних пересборок. - **Выравнивание зависимостей:** версии сейчас расходятся — kafka-python 2.0.5 у генератора против 2.0.6 в `airflow/requirements.txt`; привести к одной. `prometheus-client` в образе Airflow нет — либо добавить, либо (лучше) точка входа backfill не поднимает HTTP-сервер метрик (`service.py:65-66` вызывается безусловно — понадобится небольшая склейка; это единственное место, где честно появляется новый код, зафиксировать его в PR). - **DAG `generator_control`** в `airflow/dags/`: `schedule=None`, `max_active_runs=1`, форма запуска на `Param`: - `operation`: `backfill` | `import` | `check` (выпадающий список); - `profile`: из `PROFILES` (`launch.py`), список брать динамически; - `duration`: строка вида `6h`/`2d`, пусто — из профиля; - переопределения мира для backfill (минимум скорость и seed) — через механизм `overrides` в `build_launch_env`; - `artifact_path`: для backfill — куда сохранить артефакт (пусто — не сохранять), для import — что импортировать; дефолт в общем томе `./data`. - **Ветвление по операции** — `BranchPythonOperator`, в стиле существующих DAG стенда. - **Каждая операция начинается с валидации, падение — до любых записей:** - backfill и import: топики данных пусты (переиспользовать `KafkaTopicInspector.assert_data_topics_empty`) **и STG-таблицы ClickHouse пусты** — батч-ETL пересобирает ODS из всего STG (`sql/ods/20_stg_to_ods.sql`), поэтому пустых топиков мало: старый мир в STG смешался бы с новым. При отказе в логе — подсказка про `make clean`; - import дополнительно: артефакт читается, манифест совместим с выбранными настройками (готовые проверки `startup_history_artifact.py`), в лог — перечень разошедшихся полей; - вспомогательно: проверка «live не работает» по метрикам `generator:9109` внутри сети compose (см. Notes об ограничениях). - **Тела операций — вызовы существующего кода:** `build_launch_env` + прогон генерации до `T_end` (backfill), функции `startup_history_artifact` (import), сверка контрольных чисел с манифестом (check). Источник манифеста для check — compact-топик (он есть и после backfill, и после import); артефакт — запасной вариант. - `retries=0` у содержательных тасков: оператор должен сразу видеть ошибку (проверенный приём из greenplum-проекта). - **Документация:** в runbook `docs/runbooks/startup-history.md` — раздел «Пульт в Airflow» как основной путь и явные границы пульта (что остаётся консолью и почему, включая «continue в пульте нет намеренно»); тонкость нестандартного мира: `make generator-continue` для мира с переопределёнными настройками требует тех же настроек, громкий отказ подскажет разошедшиеся поля. Ссылки из `docs/OPERATIONS.md` и `README.md`; правки compose и Makefile описать в том же PR (правило репозитория). ## Acceptance criteria - [ ] После `make up` (генератор не автостартует) стартовую историю выбранного профиля/длительности можно создать из веб-UI Airflow, не открывая консоль и не выставляя переменных окружения; следом ETL запускается из той же цепочки и `check` зелёный. - [ ] Импорт артефакта из веб-UI: несовместимый артефакт отклоняется **до** записи в топики, в логе таска — перечень разошедшихся полей. - [ ] Backfill/import на непустом стенде отклоняются предпроверкой с подсказкой про `make clean`; миры не смешиваются. - [ ] Артефакт, сохранённый при backfill из веб-UI, пригоден для консольного импорта (формат один и тот же), и наоборот. - [ ] Операция check сверяет контрольные числа ClickHouse с манифестом и падает при расхождении. - [ ] Docker недоступен из контейнеров Airflow (socket не монтируется) — это граница решения, а не упущение. - [ ] Консольные глаголы работают как раньше; `scripts/*` и DAG сходятся в одном `launch.py`, дублирования логики запуска нет. - [ ] `make up` больше не запускает live; `make generator-continue` запускает его явно; существующие сценарии (`generated-history-analytics`, CI) не сломаны. - [ ] Runbook и `docs/OPERATIONS.md` описывают пульт как основной путь и его границы. - [ ] DAG остаётся читаемым менти: витрина параметризованной оркестрации, сложность живёт в коде генератора. ## Notes - Перед реализацией сверить API формы (`Param`, enum, описания полей) и `TriggerDagRunOperator` для Airflow 2.10 через MCP Context7 — правило репозитория. - `Config` генератора читает env процесса при создании: в таске собирать окружение через `build_launch_env` и применять к процессу таска (executor — LocalExecutor, таск живёт в своём процессе). Не забыть `GEN_DATA_DIR`: `build_launch_env` его не задаёт, дефолт `/data`, а в контейнерах Airflow данные смонтированы в `/opt/airflow/data` — без явной установки `Config()` упадёт, словари `*.jsonl` не найдутся. - Права на файлы: артефакт из worker'а пишется uid'ом Airflow (50000); в контейнере генератора `./data` смонтирован **read-only**, консольные скрипты монтируют каталог артефактов отдельно. Чтение работает везде, перезапись чужого файла — не всегда; одну строку об этом — в runbook. - Проверка `generator:9109` — вспомогательная эвристика, не защита: ошибка DNS означает «сервис не поднят» (не ошибку таска); при крашлупе после громкого отказа порт мигает (HTTP-сервер стартует до валидации state); эфемерные контейнеры `docker compose run` под этим именем не видны. Основная защита от смешивания — пустота топиков данных. - Backfill — блокирующий таск на минуты; при `max_active_runs=1` для ручного стенда это нормально. - Грабли greenplum-проекта, применимые здесь: DAG создаются на паузе — в инструкции приёмки не забыть unpause перед trigger; автоматические прогоны гонять через REST API, а не CLI внутри контейнера. - Возможное развитие (не в скоупе): операция «очистить данные» из пульта (прецедент удаления/пересоздания топиков из DAG уже есть в `airflow/dags/utils/kafka_helpers.py`) — сняла бы требование `make clean` перед новым миром. - Артефакты `daily-wave`, сделанные до задачи 14, протухнут после неё (антисмешивание отработает громко) — при пересечении работ это ожидаемо. - Рекомендуемый режим ревью по coordinator-loop: обычный — state и сериализацию задача не трогает, код генератора переиспользуется как есть. ## Зависимости (проверенные предпосылки) - `07-startup-history-portable-artifact-and-usage-docs.md` — сделана. - `11-generator-launch-verbs-and-profiles.md` — сделана. - Мягкая связь с `14-fast-teaching-profile.md`: быстрый профиль появится в выпадашке сам, если список профилей брать из `PROFILES` динамически; блокером не является.