diff --git a/AGENTS.md b/AGENTS.md index 36c3748..a986010 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -8,6 +8,25 @@ Стенд строится по спеке [«Боевой реализм стенда (v2)»](docs/specs/2026-07-30-stand-v2-realism.md). +Отсюда следует, кому что поручать. Читатель здесь — не побочный потребитель, +а тот, ради кого стенд существует: код, комментарии и документы и есть +продукт. Поэтому всё, что человек будет читать и разбирать, пишет модель с +чувством меры — по умолчанию Opus. Кодексу остаётся техническая работа без +вкусовых решений: миграции, разбор журналов и объёмного вывода, обход +репозитория, повторяющиеся правки; там его больши́е лимиты работают на нас. +Ревью идёт наоборот: разбор реализации на уместность ведёт Opus, а Кодекс +даёт независимое второе мнение другой родословной — линии ревью всегда +разных родословных, кто бы ни был исполнителем. + +Разрез проходит не по задаче целиком, а по подзадачам, и намечается при +планировании: вкусовое — Опусу, руки — Кодексу. Сначала идёт вкусовая часть: +она задаёт форму, имена и границы, под которые механическая потом +подстраивается. Писатель в дереве в каждый момент один — либо +последовательно, либо по непересекающимся файлам. + +Экономить лимиты на том, что читает человек, — ложная экономия: переписывать +выйдет дороже, чем сразу написать хорошо. + ## Язык Пиши на ясном русском языке. Иностранные слова оставляй только там, где у @@ -48,8 +67,7 @@ Airflow) и названия из кода. Если для понятия ес Трекер — Gitea на `git.dementev.space`, работа через CLI `tea` (логин по умолчанию настроен, репозиторий определяется по git remote). Команды и -подводные камни описаны в -[доке предшественника](https://git.dementev.space/ddmitry/clickstream-ch-kafka-superset-demo/src/branch/main/docs/agents/issue-tracker.md). +подводные камни — в [`docs/agents/issue-tracker.md`](docs/agents/issue-tracker.md). - Спека фичи — файл в `docs/specs/`, источник истины, версионируется с кодом. - Корневой issue фичи — тонкий: ссылка на спеку и чек-лист дочерних issues @@ -60,8 +78,33 @@ Airflow) и названия из кода. Если для понятия ес - Метки триажа — пять ролей: `needs-triage`, `needs-info`, `ready-for-agent`, `ready-for-human`, `wontfix`. Карта и её тикеты — метки `wayfinder:*`. +## Agent skills + +### Issue tracker + +Задачи — в Gitea на `git.dementev.space`, все операции через CLI `tea`. +См. [`docs/agents/issue-tracker.md`](docs/agents/issue-tracker.md). + +### Triage labels + +Пять канонических меток триажа без переименований, уже заведены в трекере. +См. [`docs/agents/triage-labels.md`](docs/agents/triage-labels.md). + +### Domain docs + +Один контекст: `CONTEXT.md` в корне и `docs/adr/`. +См. [`docs/agents/domain.md`](docs/agents/domain.md). + ## Структура - Новые документы — в `docs/` или в профильных подпапках, не в корне. - Состав доков v2 определяется по ходу этапов, набор предшественника не копируется (спека, раздел 12). +- Имена файлов в `docs/specs/` и `docs/research/` — `ГГГГ-ММ-ДД-краткое-имя.md`: + дата создания документа и слаг строчными латинскими буквами через дефис + (`2026-07-30-stand-v2-realism.md`). Дата фиксирует, когда документ появился, + и при правках не меняется: файлы сортируются по времени, а история живёт + в git. +- Имена файлов в `docs/adr/` — `NNNN-краткое-имя.md`: сквозной номер из четырёх + цифр и слаг (`0001-stand-services.md`). Решения нумеруются подряд, дата + в имени не нужна. diff --git a/README.md b/README.md index 485bca2..87610c4 100644 --- a/README.md +++ b/README.md @@ -38,29 +38,55 @@ make smoke все локальные данные стенда, а следующий `make up` создаст их с новыми значениями. -`make up` собирает локальный образ Superset, поднимает весь стенд и ждёт -здорового состояния долгоживущих контейнеров. Одноразовые `airflow-init` и +`make up` собирает локальные образы Airflow и Superset, поднимает весь стенд и +ждёт здорового состояния долгоживущих контейнеров. В образ Airflow добавлены +закреплённые клиенты ClickHouse и Kafka. Одноразовые `airflow-init` и `superset-init` завершаются с кодом 0. Первый обновляет схему Airflow, подготавливает администратора и подключение к `clickhouse-01`. Второй обновляет Superset, создаёт администратора и импортирует подключение к `clickhouse-02`. `make smoke` проверяет согласованность `.env.example` с Compose, зависимости машины, здоровье контейнеров, Kafka через порт машины, три цели Prometheus, -источник Grafana, компоненты Airflow, ручной запуск примера DAG, метаданные и -подключение Superset. В конце проверка ждёт 20 секунд покоя, печатает общую -память контейнеров и падает при превышении 3,4 ГБ. Временный топик Kafka и -проверочный запуск DAG удаляются. +источник Grafana, компоненты Airflow, ручной запуск пробников `test_clickhouse` +и `test_kafka`, метаданные и подключение Superset. Первый пробник создаёт +таблицы на обеих нодах и читает через `Distributed` на ноде 2 строку из +локальной таблицы ноды 1. Второй пишет в Kafka и читает свой маркер. В конце +проверка ждёт 20 секунд покоя, печатает общую память контейнеров и падает при +превышении 3,4 ГБ. Временный топик проверки с машины и запуски DAG удаляются; +постоянный топик пробника сохраняется, а старые записи чистит Kafka. `make smoke-cluster` запускает отдельную глубокую проверку ClickHouse: описание кластера, макросы, связь с keeper, `ReplicatedMergeTree`, `Distributed`, очередь распределённых DDL и очистку временных таблиц. -`make config-test` проверяет Compose, синтаксис Bash и Python и пробельные -ошибки в diff без запуска стенда. +`make config-test` проверяет Compose, синтаксис Bash и Python, малые проверки +логики пробников и пробельные ошибки в diff без запуска стенда. `make smoke-guards` сначала проверяет аварийную семантику кластерной проверки, -а затем останавливает Prometheus и убеждается, что общая проверка называет его -и завершается с ошибкой. В конце стенд восстанавливается. +а затем удаляет служебную таблицу пробника только на второй ноде и +останавливает Prometheus с Kafka. Общая проверка должна назвать Prometheus и +оба пробника, после чего завершиться с ошибкой. В конце стенд восстанавливается. + +### Какую проверку когда запускать + +Проверки выстроены лесенкой: чем дороже прогон, тем больше связей он трогает. + +- `make config-test` — секунды, стенд поднимать не нужно. Видит только то, что + есть в файлах, и о работоспособности не говорит ничего. Дёшево настолько, что + можно гонять перед каждым коммитом. +- `make smoke` — минута-две на поднятом стенде. Дороже, но проверяет связи + между службами, а не отдельные файлы: это интеграционная проверка. +- `make smoke-cluster` — около минуты. Одна связь, зато до дна: межнодовое + устройство ClickHouse. +- `make smoke-guards` — около пяти минут, и отвечает на другой вопрос. Не + «работает ли стенд», а «умеют ли проверки падать»: она намеренно ломает стенд + и смотрит, покраснеет ли `make smoke` и назовёт ли виновника, потом чинит и + убеждается, что стенд снова зелёный. Отсюда и три прогона `make smoke` + внутри — до поломки, во время неё и после починки. + +Обычный рабочий цикл — `make config-test` и `make smoke`. `make smoke-guards` +нужна тому, кто правит сами проверки или пробники: без неё легко завести +проверку, которая зелена всегда. Остановить контейнеры без удаления данных можно командой `make down`. Для полного сброса с удалением всех именованных томов используйте `make clean`. @@ -134,6 +160,10 @@ Airflow закреплён на 3.3.0. Состав обязательных п [архитектурой Airflow 3.3](https://airflow.apache.org/docs/apache-airflow/3.3.0/core-concepts/overview.html), [публичным интерфейсом](https://airflow.apache.org/docs/apache-airflow/3.3.0/public-airflow-interface.html) и [описанием здоровья](https://airflow.apache.org/docs/apache-airflow/3.3.0/administration-and-deployment/logging-monitoring/check-health.html). +Для пробников проверены публичный `Connection.get` из `airflow.sdk` и клиент +`clickhouse-connect`. Официальный провайдер Kafka сам использует +`confluent-kafka`; отдельное подключение и обёртки провайдера здесь не нужны, +поэтому прямой клиент оставляет образ и пример короче. Superset закреплён на 6.1.0; драйвер `clickhouse-connect`, форма `clickhousedb://` и драйвер Postgres сверены с [документацией подключений Superset](https://superset.apache.org/user-docs/6.1.0/databases/) diff --git a/compose.yaml b/compose.yaml index b62c31a..f98f036 100644 --- a/compose.yaml +++ b/compose.yaml @@ -24,7 +24,11 @@ x-clickhouse-common: &clickhouse-common start_period: 10s x-airflow-common: &airflow-common - image: ${AIRFLOW_IMAGE:-apache/airflow:3.3.0} + image: clickstream-airflow:local + build: + context: ./infra/airflow + args: + AIRFLOW_BASE_IMAGE: ${AIRFLOW_IMAGE:-apache/airflow:3.3.0} environment: AIRFLOW__CORE__EXECUTOR: LocalExecutor AIRFLOW__CORE__PARALLELISM: "4" diff --git a/dags/example_clickstream_hello.py b/dags/example_clickstream_hello.py deleted file mode 100644 index caa932e..0000000 --- a/dags/example_clickstream_hello.py +++ /dev/null @@ -1,34 +0,0 @@ -"""Пример независимого DAG для проверки Airflow 3.""" - -from __future__ import annotations - -import datetime - -from airflow.sdk import dag, task - - -@dag( - dag_id="example_clickstream_hello", - schedule=None, - start_date=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc), - catchup=False, - tags=["пример"], - doc_md=__doc__, -) -def example_clickstream_hello(): - @task - def extract() -> dict[str, int]: - return {"clicks": 3, "views": 10} - - @task - def calculate_ctr(counters: dict[str, int]) -> float: - return round(counters["clicks"] / counters["views"], 3) - - @task - def show_result(ctr: float) -> None: - print(f"CTR учебного примера: {ctr}") - - show_result(calculate_ctr(extract())) - - -example_clickstream_hello() diff --git a/dags/test_clickhouse.py b/dags/test_clickhouse.py new file mode 100644 index 0000000..fa34817 --- /dev/null +++ b/dags/test_clickhouse.py @@ -0,0 +1,215 @@ +"""Сквозная проверка подключения Airflow к кластеру ClickHouse.""" + +from __future__ import annotations + +import datetime +import uuid + +from airflow.sdk import Connection, dag, get_current_context, task + +CLUSTER = "clickstream_cluster" +LOCAL_TABLE = "airflow_probe_local" +DISTRIBUTED_TABLE = "airflow_probe_distributed" + + +def _table_engines(client, *, query_node_2: bool) -> list[tuple[str, str]]: + source = ( + "remote('clickhouse-02:9000', system.tables)" + if query_node_2 + else "system.tables" + ) + result = client.query( + f""" + SELECT name, engine + FROM {source} + WHERE database = 'default' + AND name IN ('{LOCAL_TABLE}', '{DISTRIBUTED_TABLE}') + ORDER BY name + """ + ) + return result.result_rows + + +def _drop_tables(client) -> None: + client.command( + f"DROP TABLE IF EXISTS default.{DISTRIBUTED_TABLE} " + f"ON CLUSTER {CLUSTER} SYNC" + ) + client.command( + f"DROP TABLE IF EXISTS default.{LOCAL_TABLE} ON CLUSTER {CLUSTER} SYNC" + ) + + +def _assert_tables_absent(client) -> None: + for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): + remaining = _table_engines(client, query_node_2=query_node_2) + if remaining: + raise RuntimeError( + f"служебные таблицы остались на {node_name}: {remaining}" + ) + + +def _assert_marker_path( + *, + local_rows: list[tuple[str, str]], + distributed_rows: list[tuple[int, str, str]], + node_2_hostname: str, + marker: str, +) -> None: + if len(local_rows) != 1 or local_rows[0][1] != marker: + raise RuntimeError(f"маркер не найден в локальной таблице: {marker}") + node_1_hostname = local_rows[0][0] + if node_1_hostname == node_2_hostname: + raise RuntimeError("запись и чтение маркера должны выполняться с разных нод") + expected_rows = [(1, node_1_hostname, marker)] + if distributed_rows != expected_rows: + raise RuntimeError( + "нода 2 не прочитала маркер первого шарда через Distributed: " + f"{marker}, получено {distributed_rows}" + ) + + +def _cleanup_clickhouse_client(client, original_error: BaseException | None) -> None: + cleanup_errors: list[Exception] = [] + try: + _drop_tables(client) + _assert_tables_absent(client) + except Exception as error: + cleanup_errors.append(error) + try: + client.close() + except Exception as error: + cleanup_errors.append(error) + + if original_error is not None: + for error in cleanup_errors: + original_error.add_note( + f"Дополнительная ошибка очистки ClickHouse: {error}" + ) + return + if len(cleanup_errors) == 1: + raise cleanup_errors[0] + if cleanup_errors: + raise ExceptionGroup("ошибки очистки ClickHouse", cleanup_errors) + + +@dag( + dag_id="test_clickhouse", + schedule=None, + start_date=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc), + catchup=False, + tags=["проверка"], + doc_md=__doc__, +) +def test_clickhouse(): + @task + def check_cluster_path() -> None: + import clickhouse_connect + + connection = Connection.get("clickhouse_default") + client = clickhouse_connect.get_client( + host=connection.host, + port=connection.port, + username=connection.login or "default", + password=connection.password or "", + database=connection.schema or "default", + connect_timeout=5, + send_receive_timeout=30, + ) + marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}" + expected_tables = [ + (DISTRIBUTED_TABLE, "Distributed"), + (LOCAL_TABLE, "ReplicatedMergeTree"), + ] + probe_error = None + + try: + _drop_tables(client) + _assert_tables_absent(client) + client.command( + f""" + CREATE TABLE default.{LOCAL_TABLE} ON CLUSTER {CLUSTER} + ( + marker String + ) + ENGINE = ReplicatedMergeTree( + '/clickhouse/tables/{{shard}}/{LOCAL_TABLE}', + '{{replica}}' + ) + ORDER BY marker + """ + ) + client.command( + f""" + CREATE TABLE default.{DISTRIBUTED_TABLE} ON CLUSTER {CLUSTER} + AS default.{LOCAL_TABLE} + ENGINE = Distributed( + '{CLUSTER}', + 'default', + '{LOCAL_TABLE}', + cityHash64(marker) + ) + """ + ) + + for query_node_2, node_name in ((False, "ноде 1"), (True, "ноде 2")): + actual_tables = _table_engines( + client, + query_node_2=query_node_2, + ) + if actual_tables != expected_tables: + raise RuntimeError( + f"неверный набор таблиц на {node_name}: {actual_tables}" + ) + + client.insert( + f"default.{LOCAL_TABLE}", + [[marker]], + column_names=["marker"], + ) + local_rows = client.query( + f""" + SELECT hostName(), marker + FROM default.{LOCAL_TABLE} + WHERE marker = {{marker:String}} + """, + parameters={"marker": marker}, + ).result_rows + node_2_rows = client.query( + """ + SELECT hostName() + FROM remote('clickhouse-02:9000', system.one) + """ + ).result_rows + if len(node_2_rows) != 1: + raise RuntimeError( + f"не удалось определить имя ноды 2: {node_2_rows}" + ) + distributed_rows = client.query( + f""" + SELECT _shard_num, hostName(), marker + FROM remote( + 'clickhouse-02:9000', + 'default', + '{DISTRIBUTED_TABLE}' + ) + WHERE marker = {{marker:String}} + """, + parameters={"marker": marker}, + ).result_rows + _assert_marker_path( + local_rows=local_rows, + distributed_rows=distributed_rows, + node_2_hostname=node_2_rows[0][0], + marker=marker, + ) + except BaseException as error: + probe_error = error + raise + finally: + _cleanup_clickhouse_client(client, probe_error) + + check_cluster_path() + + +test_clickhouse() diff --git a/dags/test_kafka.py b/dags/test_kafka.py new file mode 100644 index 0000000..dec32c9 --- /dev/null +++ b/dags/test_kafka.py @@ -0,0 +1,98 @@ +"""Проверка записи и чтения сообщения Airflow через Kafka.""" + +from __future__ import annotations + +import datetime +import time +import uuid + +from airflow.sdk import dag, get_current_context, task + +BROKER = "kafka:9092" +TOPIC = "airflow_integration_probe" + + +@dag( + dag_id="test_kafka", + schedule=None, + start_date=datetime.datetime(2026, 1, 1, tzinfo=datetime.timezone.utc), + catchup=False, + tags=["проверка"], + doc_md=__doc__, +) +def test_kafka(): + @task + def check_round_trip() -> None: + from confluent_kafka import Consumer, Producer, TopicPartition + + marker = f"{get_current_context()['run_id']}:{uuid.uuid4()}" + marker_bytes = marker.encode() + group_id = f"airflow-probe-{uuid.uuid4()}" + producer = None + consumer = None + + try: + # Стенд опирается на автосоздание: срок хранения чистит записи, а + # не топик; в бою автосоздание обычно выключают. + delivery_errors: list[str] = [] + delivered_offsets: list[tuple[int, int]] = [] + + def on_delivery(error, message) -> None: + if error is not None: + delivery_errors.append(str(error)) + else: + delivered_offsets.append((message.partition(), message.offset())) + + producer = Producer( + { + "bootstrap.servers": BROKER, + "client.id": "airflow-test-kafka", + "message.timeout.ms": 10_000, + "socket.timeout.ms": 5_000, + } + ) + producer.produce( + TOPIC, + key=marker_bytes, + value=marker_bytes, + on_delivery=on_delivery, + ) + undelivered = producer.flush(10) + if undelivered or delivery_errors or len(delivered_offsets) != 1: + raise RuntimeError( + "Kafka не подтвердила запись маркера: " + f"не доставлено {undelivered}, ошибки {delivery_errors}, " + f"смещения {delivered_offsets}" + ) + + consumer = Consumer( + { + "bootstrap.servers": BROKER, + "group.id": group_id, + "enable.auto.commit": False, + "session.timeout.ms": 6_000, + "socket.timeout.ms": 5_000, + } + ) + partition, offset = delivered_offsets[0] + consumer.assign([TopicPartition(TOPIC, partition, offset)]) + read_deadline = time.monotonic() + 20 + while time.monotonic() < read_deadline: + message = consumer.poll(0.5) + if message is None: + continue + if message.error(): + raise RuntimeError(f"Kafka вернула ошибку чтения: {message.error()}") + if message.value() == marker_bytes: + return + raise RuntimeError(f"Kafka не вернула свой маркер за 20 секунд: {marker}") + finally: + if producer is not None: + producer.flush(1) + if consumer is not None: + consumer.close() + + check_round_trip() + + +test_kafka() diff --git a/docs/adr/0001-stand-services.md b/docs/adr/0001-stand-services.md index f645e04..49309fc 100644 --- a/docs/adr/0001-stand-services.md +++ b/docs/adr/0001-stand-services.md @@ -7,7 +7,7 @@ Стенд использует Airflow 3.3.0 с LocalExecutor. Airflow разделён на одноразовый `airflow-init` и три долгоживущих процесса: API, планировщик и обработчик DAG. Triggerer не запускается: в каркасе нет отложенных задач. -Пример DAG написан через публичный `airflow.sdk`. +Два интеграционных пробника написаны через публичный `airflow.sdk`. Для входа выбран SimpleAuthManager. У него нет команды создания пользователя, поэтому `airflow-init` записывает пароль администратора из окружения в @@ -25,9 +25,16 @@ JSON-файл отдельного тома `airflow_auth`. Новый имен Airflow получает подготовленное общее подключение к `clickhouse-01`, а Superset — подключение `clickhousedb://` к `clickhouse-02`. Провайдер -ClickHouse для Airflow пока не нужен. Разные ноды создают учебную ловушку: -забытый `ON CLUSTER` проявится в Superset, даже если операция Airflow на первой -ноде прошла успешно. +ClickHouse для Airflow не нужен: локальный образ содержит прямой клиент. +Пробник создаёт служебные таблицы `ON CLUSTER`, пишет на первой ноде и через +`remote` читает `Distributed` на второй. Разные ноды сохраняют учебную +ловушку: забытый `ON CLUSTER` проявится в Superset, даже если операция Airflow +на первой ноде прошла успешно. + +Клиент Kafka также добавлен прямо в образ Airflow. Официальный провайдер +использует тот же `confluent-kafka`, а пробнику не нужны его подключение, +операторы и обёртки. Топик пробника постоянный, сообщение каждого запуска +отличается уникальным маркером. Superset собирается от `apache/superset:6.1.0`. В образ добавлены `clickhouse-connect` для ClickHouse и `psycopg2-binary` для метаданных @@ -61,6 +68,11 @@ Prometheus читает встроенные точки метрик двух с Оттуда взяты `airflow.sdk`, обязательный отдельный обработчик DAG, возможность не запускать triggerer и проверка конкретных компонентов здоровья. +Для интеграционных пробников проверены `Connection.get` в публичном +`airflow.sdk`, запросы через `clickhouse-connect` и состав официального +провайдера Kafka. Выбран прямой `confluent-kafka`: провайдер строит свои +подключения и обёртки поверх него, которые двум коротким пробникам не нужны. + Драйверы и строки подключения проверены по документации Superset 6.1.0: [подключения к базам](https://superset.apache.org/user-docs/6.1.0/databases/), [ClickHouse](https://superset.apache.org/user-docs/databases/supported/clickhouse/) diff --git a/docs/agents/domain.md b/docs/agents/domain.md new file mode 100644 index 0000000..fe22ec3 --- /dev/null +++ b/docs/agents/domain.md @@ -0,0 +1,54 @@ +# Доки предметной области + +Как скиллам читать документацию репозитория, прежде чем лезть в код. + +## Что прочитать до разбора кода + +- **`CONTEXT.md`** в корне — словарь понятий проекта. +- **`docs/adr/`** — решения. Читать те ADR, что касаются области, в которой + сейчас работаешь. +- **`docs/specs/`** — спеки фич; для этого репозитория спека и есть источник + истины по тому, что строим (см. `docs/agents/issue-tracker.md`). + +Если файла нет — **просто идти дальше молча**. Не сообщать об отсутствии и не +предлагать создать заранее. Скилл `/domain-modeling` (через `/grill-with-docs` +и `/improve-codebase-architecture`) заводит их лениво, когда термин или решение +действительно понадобилось зафиксировать. + +## Разметка: один контекст + +Репозиторий одноконтекстный — один `CONTEXT.md` в корне и один `docs/adr/`: + +``` +/ +├── CONTEXT.md ← пока не создан, появится по надобности +├── docs/ +│ ├── adr/ ← 0001-stand-services.md +│ ├── specs/ +│ └── research/ +├── dags/ +├── infra/ +└── scripts/ +``` + +Карта контекстов (`CONTEXT-MAP.md` в корне и по `CONTEXT.md` на контекст) — +для больших многопакетных репозиториев; здесь она не нужна. Если репозиторий +когда-нибудь разъедется на несколько контекстов, разметку менять здесь. + +## Пользоваться словарём + +Если в выводе называешь понятие предметной области (заголовок issue, гипотеза, +имя теста, предложение по рефакторингу) — бери термин ровно в том виде, как он +записан в `CONTEXT.md`. Не уходить в синонимы, от которых словарь отказался. + +Понятия нет в словаре — это сигнал: либо выдумываешь язык, которого в проекте +нет (стоит передумать), либо нашёл настоящий пробел (отметить для +`/domain-modeling`). + +## Спорить с ADR вслух + +Если то, что ты предлагаешь, противоречит принятому ADR — сказать об этом прямо, +а не переехать решение молча: + +> _Противоречит ADR-0001 (состав сервисов стенда) — но открыть заново стоит, +> потому что…_ diff --git a/docs/agents/issue-tracker.md b/docs/agents/issue-tracker.md new file mode 100644 index 0000000..5296c87 --- /dev/null +++ b/docs/agents/issue-tracker.md @@ -0,0 +1,132 @@ +# Issue tracker: Gitea + +Задачи этого репозитория живут в Gitea на `git.dementev.space` +(`ddmitry/clickstream-data-platform`, это remote `origin`). Все операции — +через CLI [`tea`](https://gitea.com/gitea/tea), официальный клиент Gitea; по +устройству он близок к `gh` и `glab`. Логин и репозиторий `tea` определяет сам +по git remote в текущем каталоге. + +## Перед первым запуском + +- **Бинарник.** Скачивается с `https://dl.gitea.com/tea/<версия>/` (файл + `tea-<версия>-linux-amd64` и `.sha256` рядом), кладётся в `~/.local/bin/tea`. + Проверка: `tea --version`. +- **Вход.** `tea logins add --name git.dementev.space --url + https://git.dementev.space`, токен передаётся переменной + `GITEA_SERVER_TOKEN` (не аргументом командной строки — он попадёт в историю + оболочки). Логин уже добавлен и назначен по умолчанию, так что `tea` работает + из любого каталога. +- **Скоупы токена:** `read:user` (без него `tea` откажется добавлять логин), + `write:issue`, `write:repository`. Токен выпускается в UI: Settings → + Applications. Нехватка скоупа выглядит не как «нет прав», а как невнятная + ошибка или пустой ответ. +- **Прокси.** Домен `dementev.space` должен быть в `NO_PROXY`, иначе запросы + уходят в прокси и виснут. В обычной оболочке это делает `proxy-client` из + `~/dotfiles`. Если переменная не подхватилась, короткий разовый префикс: + `NO_PROXY='*' tea ...`. + +## Команды + +- **Создать issue:** `tea issues create --title "..." --description "..."`. + Многострочное тело удобнее собрать heredoc'ом в переменную и подставить + как `--description "$BODY"`. +- **Прочитать issue:** `tea issues <номер> --comments`. +- **Список:** `tea issues list --state open --output json --fields + index,title,labels,assignees`. Фильтры: `--labels`, `--assignee`, + `--keyword`. +- **Комментарий:** `tea comments add <номер> -d "..."`. +- **Метки:** `tea issues edit <номер> --add-labels "..."` / `--remove-labels + "..."`. Список меток репозитория — `tea labels list`, создать новую — + `tea labels create --name "..." --color "..."`. +- **Закрыть:** `tea issues close <номер>`. Комментария при закрытии команда не + принимает — сначала `tea comments add`, потом `close`. +- **Взять в работу:** `tea issues edit <номер> --add-assignees ddmitry`. + Сокращения вида `@me` в `tea` нет, имя пишется целиком. +- **Чего нет в CLI** — через `tea api `: команда ходит в REST API Gitea + уже с сохранённым токеном, например + `tea api repos/ddmitry/clickstream-data-platform/issues/18`. + +## Слияние PR не закрывает issue + +Gitea понимает только английские ключевые слова автозакрытия: `closes`, +`fixes`, `resolves` (`Closes #13`). Русское «Закрывает #13» в теле PR — обычная +ссылка, issue останется открытым. Наступали не раз: последний случай — слияние +PR #17, где issue #13 пришлось закрывать руками. + +Порядок после слияния: закрыть задачу (`tea issues close <номер>`) и тикнуть +её пункт в чек-листе родительского issue этапа — Gitea чек-листы сама не +обновляет. + +## Спека — источник истины + +- Спецификация фичи — файл в `docs/specs/`, версионируется с кодом. +- Корневой issue фичи — **тонкий**: ссылка на спеку + чек-лист дочерних issues + (`- [ ] #NN`). Содержание спеки в issue не дублируется — истина одна, в git. +- Дочерние issues — полноценные самодостаточные постановки: цель, критерии + приёмки чекбоксами, границы («что трогать нельзя»), «сначала прочитать», + команды проверки. +- Итоговые резолюции и решения — в спеку или ADR тем же PR; issue — рабочая + переписка, она не обязана переживать фичу. + +## Когда скилл говорит «опубликовать в issue tracker» + +Создать issue в Gitea: `tea issues create ...`. + +## Когда скилл говорит «достать тикет» + +`tea issues <номер> --comments`. + +## PR как поверхность триажа + +**Нет** — одиночный учебный репозиторий, внешних PR не ждём. (Если включить — +`/triage` начнёт гонять PR через те же метки и состояния командами +`tea pulls ...`.) + +## Wayfinding-операции + +Используются `/wayfinder`. Карта — один issue, тикеты — дочерние issues. + +- **Карта**: issue с меткой `wayfinder:map` (Notes / Decisions-so-far / Fog + в теле). +- **Дочерний тикет**: вложенных issues в Gitea нет, поэтому связь держится + двумя ссылками — пункт списка `- [ ] #NN` в теле карты и строка + `Part of #<карта>` в начале тела тикета. Метки: `wayfinder:<тип>` + (`research` / `prototype` / `grilling` / `task`). +- **Блокировки**: нативные зависимости Gitea — единственный источник истины, + текстовых строк `Blocked by:` в телах тикетов нет. В CLI их команд нет, + работаем через `tea api` (`{owner}` и `{repo}` подставляются из текущего + репозитория): + - добавить блокер: `tea api repos/{owner}/{repo}/issues//dependencies + -F index=<блокер> -f owner=ddmitry -f repo=clickstream-data-platform` + — поля `owner` и `repo` обязательны, без них API отвечает + «repository does not exist»; + - кто блокирует тикет: `GET .../issues//dependencies`; + - кого блокирует тикет: `GET .../issues//blocks`; + - снять блокировку: тот же путь методом `DELETE` с тем же телом. + + Тикет разблокирован, когда у всех блокеров `state == "closed"`. +- **Фронтир**: открытые дети карты минус заблокированные и назначенные; первый + в порядке карты. Блокеры проверяются запросом `dependencies` по каждому + кандидату. +- **Взять в работу**: `tea issues edit --add-assignees ddmitry` — первая + запись за сессию. +- **Закрыть**: `tea comments add -d "<ответ>"`, затем `tea issues close + `, затем указатель на контекст (суть + ссылка) в Decisions-so-far карты. + +## Что проверено и когда + +2026-07-31: Gitea 1.27.0 (`tea api version`), `tea` 0.15.0. Живыми запросами по +этому репозиторию проверены `tea labels list` и `tea issues list`. Остальной +набор команд, флагов и приёмы с `tea api` перенесены из доки +предшественника — там они снимались с `tea <команда> --help` установленного +бинарника и проверялись живыми запросами 2026-07-29. При обновлении `tea` стоит +перечитать `--help`: набор флагов между версиями менялся. + +## Архив + +- Трекер предшественника — репозиторий + [`ddmitry/clickstream-ch-kafka-superset-demo`](https://git.dementev.space/ddmitry/clickstream-ch-kafka-superset-demo). + Задачи оттуда сюда не переносились: платформа v2 начата с чистого трекера. +- До 2026-07-26 задачи предшественника жили в GitHub Issues + (`dementev-dev/…`); аккаунт заблокирован, номера воссозданы в Gitea один + в один. diff --git a/docs/agents/triage-labels.md b/docs/agents/triage-labels.md new file mode 100644 index 0000000..2af78cd --- /dev/null +++ b/docs/agents/triage-labels.md @@ -0,0 +1,22 @@ +# Метки триажа + +Скиллы говорят о пяти канонических ролях триажа. Эта таблица переводит роли в +метки, которые реально заведены в трекере репозитория. + +| Роль в скиллах | Метка у нас | Значение | +| ----------------- | ----------------- | ------------------------------------------- | +| `needs-triage` | `needs-triage` | Задачу надо оценить, решение не принято | +| `needs-info` | `needs-info` | Ждём уточнений от автора | +| `ready-for-agent` | `ready-for-agent` | Постановка полная, можно отдавать агенту | +| `ready-for-human` | `ready-for-human` | Нужен человек | +| `wontfix` | `wontfix` | Делать не будем | + +Имена совпадают с каноническими: переименований нет. Когда скилл говорит про +роль («поставь метку готовности для агента»), берём строку из правой колонки. + +Все пять меток уже созданы в Gitea (проверено `tea labels list` 2026-07-31) — +заводить их заново не нужно. Кроме них в репозитории живут метки карты +`wayfinder:*` (`map`, `research`, `prototype`, `grilling`, `task`); к триажу они +отношения не имеют, ими размечает `/wayfinder`. + +Если вокабуляр меток поменяется — править правую колонку здесь, а не в скиллах. diff --git a/infra/airflow/Dockerfile b/infra/airflow/Dockerfile new file mode 100644 index 0000000..a7cee69 --- /dev/null +++ b/infra/airflow/Dockerfile @@ -0,0 +1,13 @@ +ARG AIRFLOW_BASE_IMAGE +FROM ${AIRFLOW_BASE_IMAGE} + +USER root + +# Официальный провайдер Kafka использует тот же confluent-kafka. Пробникам +# не нужны его подключения и обёртки, поэтому в образ добавлен сам клиент. +RUN uv pip install --python /home/airflow/.local/bin/python --no-cache \ + "clickhouse-connect==1.6.0" \ + "confluent-kafka==2.15.0" \ + && python -c "import clickhouse_connect, confluent_kafka" + +USER airflow diff --git a/scripts/clickhouse-smoke.sh b/scripts/clickhouse-smoke.sh index 1baa7de..15432b9 100755 --- a/scripts/clickhouse-smoke.sh +++ b/scripts/clickhouse-smoke.sh @@ -10,10 +10,14 @@ compose() { docker compose --project-directory "$ROOT_DIR" "$@" } +# Ввод закрыт намеренно: у запроса INSERT clickhouse-client дочитывает данные +# из стандартного ввода и ждёт его конца. Если проверку запустили не из +# терминала, а из фонового процесса с открытым вводом, конца не наступает +# никогда, и команда висит без сообщений. query() { local service="$1" local sql="$2" - compose exec -T "$service" clickhouse-client --query "$sql" + compose exec -T "$service" clickhouse-client --query "$sql" 0) ' >/dev/null <<<"$config_json" +grep -qx 'ARG AIRFLOW_BASE_IMAGE' "$ROOT_DIR/infra/airflow/Dockerfile" +jq -e ' + .services["airflow-init"].image == "clickstream-airflow:local" and + (.services["airflow-init"].build.args.AIRFLOW_BASE_IMAGE | length > 0) and + all( + .services[]; + ((.environment // {}) | has("_PIP_ADDITIONAL_REQUIREMENTS") | not) + ) +' >/dev/null <<<"$config_json" jq -e ' .services["airflow-init"] as $service | ($service.environment.AIRFLOW__CORE__SIMPLE_AUTH_MANAGER_PASSWORDS_FILE == @@ -39,6 +48,20 @@ if [[ "${#shell_files[@]}" -eq 0 ]]; then fi bash -n "${shell_files[@]}" PYTHONPYCACHEPREFIX="$CACHE_DIR" uv run --no-project python -m compileall -q "$ROOT_DIR/dags" +unit_status=0 +unit_output="$( + PYTHONPYCACHEPREFIX="$CACHE_DIR" \ + uv run --no-project python "$ROOT_DIR/tests/dag-probes-unit.py" 2>&1 +)" || unit_status=$? +printf '%s\n' "$unit_output" +if [[ "$unit_status" -ne 0 ]]; then + exit "$unit_status" +fi +if grep -Eq '^(Ran [0-9]+ tests|OK|FAILED)' <<<"$unit_output" || + ! grep -Eq '^ИТОГ: пройдено [0-9]+, ошибок 0$' <<<"$unit_output"; then + printf 'ОШИБКА: малые проверки пробников вывели итог не на русском языке.\n' >&2 + exit 1 +fi git -C "$ROOT_DIR" diff --check printf 'ЗЕЛЁНО: Compose, Bash, Python и пробельные ошибки diff проверены.\n' diff --git a/scripts/stand-smoke.sh b/scripts/stand-smoke.sh index ae6d0c1..b9e862c 100755 --- a/scripts/stand-smoke.sh +++ b/scripts/stand-smoke.sh @@ -17,6 +17,7 @@ kafka_cleanup_image='' kafka_cleanup_port='' kafka_cleanup_topic='' airflow_cleanup_port='' +airflow_cleanup_dag_id='' airflow_cleanup_original_paused='' airflow_cleanup_run_id='' airflow_cleanup_token='' @@ -232,7 +233,7 @@ cleanup_airflow_run() { encoded_run_id="$(jq -rn --arg value "$airflow_cleanup_run_id" '$value | @uri')" if curl -sf --max-time 10 -X DELETE \ -H "Authorization: Bearer ${airflow_cleanup_token}" \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello/dagRuns/${encoded_run_id}" \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${airflow_cleanup_dag_id}/dagRuns/${encoded_run_id}" \ >/dev/null 2>&1; then airflow_cleanup_run_id='' return 0 @@ -250,7 +251,7 @@ cleanup_airflow_pause() { -H "Authorization: Bearer ${airflow_cleanup_token}" \ -H 'Content-Type: application/json' \ -d "$payload" \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello" \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${airflow_cleanup_dag_id}" \ >/dev/null 2>&1; then airflow_cleanup_original_paused='' return 0 @@ -383,18 +384,84 @@ check_grafana_datasource() { fi } +run_airflow_probe() { + local dag + local dag_id="$1" + local encoded_run_id + local response + local run_state='' + local unpaused + local -i attempt + + dag="$(curl -sf --max-time 10 \ + -H "Authorization: Bearer ${airflow_cleanup_token}" \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${dag_id}" \ + 2>/dev/null || true)" + if ! jq -e --arg dag_id "$dag_id" '.dag_id == $dag_id' \ + >/dev/null 2>&1 <<<"$dag"; then + fail "пробник ${dag_id} не виден через API Airflow" + return 1 + fi + + airflow_cleanup_dag_id="$dag_id" + airflow_cleanup_original_paused="$(jq -r '.is_paused | tostring' <<<"$dag")" + unpaused="$(curl -sf --max-time 10 -X PATCH \ + -H "Authorization: Bearer ${airflow_cleanup_token}" \ + -H 'Content-Type: application/json' \ + -d '{"is_paused":false}' \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${dag_id}" \ + 2>/dev/null || true)" + if ! jq -e '.is_paused == false' >/dev/null 2>&1 <<<"$unpaused"; then + fail "Airflow не смог включить пробник ${dag_id} перед ручным запуском" + cleanup_airflow_pause || true + return 1 + fi + + response="$(curl -sf --max-time 10 -X POST \ + -H "Authorization: Bearer ${airflow_cleanup_token}" \ + -H 'Content-Type: application/json' \ + -d '{"logical_date":null}' \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${dag_id}/dagRuns" \ + 2>/dev/null || true)" + airflow_cleanup_run_id="$(jq -r '.dag_run_id // empty' <<<"$response")" + if [[ -z "$airflow_cleanup_run_id" ]]; then + fail "Airflow не создал ручной запуск пробника ${dag_id}" + cleanup_airflow_pause || true + return 1 + fi + + encoded_run_id="$(jq -rn --arg value "$airflow_cleanup_run_id" '$value | @uri')" + for ((attempt = 1; attempt <= 60; attempt++)); do + response="$(curl -sf --max-time 10 \ + -H "Authorization: Bearer ${airflow_cleanup_token}" \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${dag_id}/dagRuns/${encoded_run_id}" \ + 2>/dev/null || true)" + run_state="$(jq -r '.state // empty' <<<"$response")" + [[ "$run_state" == 'success' || "$run_state" == 'failed' ]] && break + sleep 2 + done + + if [[ "$run_state" == 'success' ]] && \ + cleanup_airflow_run && cleanup_airflow_pause; then + pass "пробник ${dag_id} завершился успешно, запуск удалён, исходная пауза восстановлена" + return 0 + fi + + fail "пробник ${dag_id} не завершился чисто: состояние ${run_state:-неизвестно}" + cleanup_airflow_run || true + cleanup_airflow_pause || true + return 1 +} + check_airflow() { local config local connection local dag - local encoded_run_id + local dag_id local health local password local response - local run_state='' - local unpaused local user - local -i attempt config="$(compose config --format json 2>/dev/null || true)" user="$(jq -r '.services["airflow-apiserver"].environment.AIRFLOW_ADMIN_USER // empty' <<<"$config")" @@ -419,29 +486,22 @@ check_airflow() { '{username: $username, password: $password}')" \ "http://127.0.0.1:${airflow_cleanup_port}/auth/token" 2>/dev/null || true)" airflow_cleanup_token="$(jq -r '.access_token // empty' <<<"$response")" - dag="$(curl -sf --max-time 10 \ - -H "Authorization: Bearer ${airflow_cleanup_token}" \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello" \ - 2>/dev/null || true)" - if [[ -n "$airflow_cleanup_token" ]] && \ - jq -e '.dag_id == "example_clickstream_hello"' >/dev/null 2>&1 <<<"$dag"; then - airflow_cleanup_original_paused="$(jq -r '.is_paused | tostring' <<<"$dag")" - pass 'учётные данные администратора Airflow принимаются, пример DAG виден через API' - else - fail 'Airflow не принял учётные данные администратора или не показал пример DAG' - return - fi - - unpaused="$(curl -sf --max-time 10 -X PATCH \ - -H "Authorization: Bearer ${airflow_cleanup_token}" \ - -H 'Content-Type: application/json' \ - -d '{"is_paused":false}' \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello" \ - 2>/dev/null || true)" - if ! jq -e '.is_paused == false' >/dev/null 2>&1 <<<"$unpaused"; then - fail 'Airflow не смог включить пример DAG перед ручным запуском' + if [[ -z "$airflow_cleanup_token" ]]; then + fail 'Airflow не принял учётные данные администратора' return fi + for dag_id in test_clickhouse test_kafka; do + dag="$(curl -sf --max-time 10 \ + -H "Authorization: Bearer ${airflow_cleanup_token}" \ + "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/${dag_id}" \ + 2>/dev/null || true)" + if ! jq -e --arg dag_id "$dag_id" '.dag_id == $dag_id' \ + >/dev/null 2>&1 <<<"$dag"; then + fail "Airflow не показал пробник ${dag_id}" + return + fi + done + pass 'учётные данные администратора Airflow принимаются, оба пробника видны через API' connection="$(curl -sf --max-time 10 \ -H "Authorization: Bearer ${airflow_cleanup_token}" \ @@ -463,34 +523,8 @@ check_airflow() { fail 'подключение Airflow не указывает на clickhouse-01 или нода недоступна из контейнера' fi - response="$(curl -sf --max-time 10 -X POST \ - -H "Authorization: Bearer ${airflow_cleanup_token}" \ - -H 'Content-Type: application/json' \ - -d '{"logical_date":null}' \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello/dagRuns" \ - 2>/dev/null || true)" - airflow_cleanup_run_id="$(jq -r '.dag_run_id // empty' <<<"$response")" - if [[ -z "$airflow_cleanup_run_id" ]]; then - fail 'Airflow не создал ручной запуск примера DAG' - return - fi - - encoded_run_id="$(jq -rn --arg value "$airflow_cleanup_run_id" '$value | @uri')" - for ((attempt = 1; attempt <= 60; attempt++)); do - response="$(curl -sf --max-time 10 \ - -H "Authorization: Bearer ${airflow_cleanup_token}" \ - "http://127.0.0.1:${airflow_cleanup_port}/api/v2/dags/example_clickstream_hello/dagRuns/${encoded_run_id}" \ - 2>/dev/null || true)" - run_state="$(jq -r '.state // empty' <<<"$response")" - [[ "$run_state" == 'success' || "$run_state" == 'failed' ]] && break - sleep 2 - done - - if [[ "$run_state" == 'success' ]] && cleanup_airflow_run && cleanup_airflow_pause; then - pass 'ручной запуск примера DAG завершился успешно, удалён, исходная пауза восстановлена' - else - fail "ручной запуск примера DAG не завершился чисто: состояние ${run_state:-неизвестно}" - fi + run_airflow_probe test_clickhouse || true + run_airflow_probe test_kafka || true } check_superset() { diff --git a/tests/dag-probes-unit.py b/tests/dag-probes-unit.py new file mode 100644 index 0000000..78995ed --- /dev/null +++ b/tests/dag-probes-unit.py @@ -0,0 +1,184 @@ +"""Малые проверки логики пробника ClickHouse без запуска Airflow.""" + +from __future__ import annotations + +import importlib.util +import sys +import types +import unittest +from pathlib import Path + + +def load_clickhouse_dag_module(): + airflow_module = types.ModuleType("airflow") + sdk_module = types.ModuleType("airflow.sdk") + + def dag(**_kwargs): + def decorate(function): + return function + + return decorate + + def task(function): + def declare_task(*_args, **_kwargs): + return None + + return declare_task + + sdk_module.Connection = object + sdk_module.dag = dag + sdk_module.get_current_context = lambda: {} + sdk_module.task = task + airflow_module.sdk = sdk_module + sys.modules["airflow"] = airflow_module + sys.modules["airflow.sdk"] = sdk_module + + dag_path = Path(__file__).resolve().parents[1] / "dags" / "test_clickhouse.py" + spec = importlib.util.spec_from_file_location("test_clickhouse_dag", dag_path) + if spec is None or spec.loader is None: + raise RuntimeError("не удалось загрузить модуль пробника ClickHouse") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def load_kafka_dag_task(): + airflow_module = types.ModuleType("airflow") + sdk_module = types.ModuleType("airflow.sdk") + captured_tasks = {} + + def dag(**_kwargs): + def decorate(function): + return function + + return decorate + + def task(function): + captured_tasks[function.__name__] = function + + def declare_task(*_args, **_kwargs): + return None + + return declare_task + + sdk_module.dag = dag + sdk_module.get_current_context = lambda: {"run_id": "unit-test"} + sdk_module.task = task + airflow_module.sdk = sdk_module + sys.modules["airflow"] = airflow_module + sys.modules["airflow.sdk"] = sdk_module + + dag_path = Path(__file__).resolve().parents[1] / "dags" / "test_kafka.py" + spec = importlib.util.spec_from_file_location("test_kafka_dag", dag_path) + if spec is None or spec.loader is None: + raise RuntimeError("не удалось загрузить модуль пробника Kafka") + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return captured_tasks["check_round_trip"] + + +class FailingCleanupClient: + def __init__(self) -> None: + self.closed = False + + def command(self, _sql: str) -> None: + raise RuntimeError("очистка недоступна") + + def close(self) -> None: + self.closed = True + + +class ClickHouseProbeTests(unittest.TestCase): + @classmethod + def setUpClass(cls) -> None: + cls.module = load_clickhouse_dag_module() + + def test_marker_path_requires_different_nodes_and_first_shard(self) -> None: + marker = "свой-маркер" + self.module._assert_marker_path( + local_rows=[("clickhouse-01-host", marker)], + distributed_rows=[(1, "clickhouse-01-host", marker)], + node_2_hostname="clickhouse-02-host", + marker=marker, + ) + + with self.assertRaisesRegex(RuntimeError, "разных нод"): + self.module._assert_marker_path( + local_rows=[("clickhouse-02-host", marker)], + distributed_rows=[(2, "clickhouse-02-host", marker)], + node_2_hostname="clickhouse-02-host", + marker=marker, + ) + with self.assertRaisesRegex(RuntimeError, "первого шарда"): + self.module._assert_marker_path( + local_rows=[("clickhouse-01-host", marker)], + distributed_rows=[(2, "clickhouse-01-host", marker)], + node_2_hostname="clickhouse-02-host", + marker=marker, + ) + + def test_cleanup_keeps_original_error_as_primary(self) -> None: + client = FailingCleanupClient() + original_error = RuntimeError("маркер не найден") + + self.module._cleanup_clickhouse_client(client, original_error) + + self.assertTrue(client.closed) + self.assertIn("очистка недоступна", "\n".join(original_error.__notes__)) + + +class KafkaProbeTests(unittest.TestCase): + def test_producer_flushes_when_consumer_creation_fails(self) -> None: + kafka_module = types.ModuleType("confluent_kafka") + flush_timeouts = [] + + class Message: + def partition(self) -> int: + return 0 + + def offset(self) -> int: + return 1 + + class Producer: + def __init__(self, _config) -> None: + pass + + def produce(self, _topic, **kwargs) -> None: + kwargs["on_delivery"](None, Message()) + + def flush(self, timeout: int) -> int: + flush_timeouts.append(timeout) + return 0 + + class Consumer: + def __init__(self, _config) -> None: + raise RuntimeError("чтение недоступно") + + kafka_module.Consumer = Consumer + kafka_module.Producer = Producer + kafka_module.TopicPartition = object + sys.modules["confluent_kafka"] = kafka_module + check_round_trip = load_kafka_dag_task() + + with self.assertRaisesRegex(RuntimeError, "чтение недоступно"): + check_round_trip() + + self.assertEqual(flush_timeouts, [10, 1]) + + +def run_tests() -> int: + suite = unittest.defaultTestLoader.loadTestsFromModule(sys.modules[__name__]) + result = unittest.TestResult() + suite.run(result) + + problems = result.failures + result.errors + for test, details in problems: + print(f"ОШИБКА: {test.id()}", file=sys.stderr) + print(details, file=sys.stderr) + passed = result.testsRun - len(problems) - len(result.skipped) + print(f"ИТОГ: пройдено {passed}, ошибок {len(problems)}") + return 0 if result.wasSuccessful() else 1 + + +if __name__ == "__main__": + raise SystemExit(run_tests()) diff --git a/tests/docs-guards.sh b/tests/docs-guards.sh index 4b154d1..91d1aad 100755 --- a/tests/docs-guards.sh +++ b/tests/docs-guards.sh @@ -1,38 +1,117 @@ #!/usr/bin/env bash +# Сторож документации. +# +# README устаревает молча: порт поменяли в compose.yaml, а в описании остался +# старый — и это выясняется через месяц, когда кто-то по нему подключается. +# Здесь собраны утверждения README, которые дёшево проверить текстом и дорого +# обнаружить сломанными. Стенд поднимать не нужно; запускается в составе +# `make config-test`. +# +# Чего сторож НЕ делает: он не проверяет, что README понятен или полон. Только +# то, что перечисленные ниже факты не разошлись с кодом. +# +# Добавляя проверку, формулируй утверждение так, как оно должно читаться в +# отчёте: строка печатается и при успехе, и при провале, поэтому по красной +# строке сразу видно, что именно перестало быть правдой. + set -euo pipefail readonly ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" readonly README="$ROOT_DIR/README.md" +readonly SMOKE="$ROOT_DIR/scripts/stand-smoke.sh" +passed=0 -grep -Eq 'нода 1.*28123.*29000' "$README" -grep -Eq 'нода 2.*28124.*29001' "$README" -printf 'ЗЕЛЁНО: README перечисляет HTTP- и нативные порты обеих нод.\n' - -grep -Eq 'После первого запуска.*`make clean`' "$README" -printf 'ЗЕЛЁНО: README объясняет сброс томов после смены исходных учётных данных.\n' - -grep -Fxq -- '- нода 2 — `http://127.0.0.1:28124`, нативный порт `29001`;' "$README" -printf 'ЗЕЛЁНО: список портов остаётся единым списком.\n' - -awk ' - previous == "После изменения `infra/clickhouse/config.d/prometheus.xml` выполните" && - $0 == "`docker compose restart clickhouse-01 clickhouse-02`: обычный `make up` не" { - found = 1 - } - {previous = $0} - END {exit !found} -' "$README" -printf 'ЗЕЛЁНО: README требует перезапуск ClickHouse после изменения настройки метрик.\n' - -grep_status=0 -grep -q '3,4 GB' "$ROOT_DIR/scripts/stand-smoke.sh" || grep_status=$? -if [[ "$grep_status" -eq 0 ]]; then - printf 'ОШИБКА: отчёт проверки использует латинское обозначение GB.\n' >&2 +fail() { + printf 'ОШИБКА: %s\n' "$1" >&2 exit 1 -elif [[ "$grep_status" -ne 1 ]]; then - printf 'ОШИБКА: не удалось проверить обозначение единицы памяти.\n' >&2 - exit 1 -fi -printf 'ЗЕЛЁНО: отчёт проверки использует русское обозначение ГБ.\n' +} -printf 'ИТОГ: пройдено 5, ошибок 0\n' +# check «утверждение» команда... — запускает команду и считает результат. +# Успех: «ЗЕЛЁНО: утверждение». Провал: «ОШИБКА: не подтвердилось: утверждение» +# и выход с кодом 1. Раньше проверки падали через `set -e` молча: код возврата +# был единственным следом, и какая именно проверка не прошла — не сообщалось. +check() { + local claim="$1" + shift + if "$@"; then + passed=$((passed + 1)) + printf 'ЗЕЛЁНО: %s\n' "$claim" + else + fail "не подтвердилось: $claim" + fi +} + +# Порты обеих нод. Ломается ровно тогда, когда порт поменяли в compose.yaml и +# забыли документацию — самая частая причина расхождения. +ports_documented() { + grep -Eq 'нода 1.*28123.*29000' "$README" && + grep -Eq 'нода 2.*28124.*29001' "$README" +} + +# Совет про сброс томов. Пароли Postgres и Grafana применяются при создании +# тома: без `make clean` смена значений в .env ничего не даёт, и человек +# полчаса ищет, почему его не пускает. +clean_advice_present() { + grep -Eq 'После первого запуска.*`make clean`' "$README" +} + +# Список портов остаётся единым списком. Совпадение точное намеренно: проверка +# стережёт не сам факт (он проверен выше), а то, что строку не выдернули из +# списка в отдельный абзац при правке соседнего текста. +ports_stay_one_list() { + grep -Fxq -- '- нода 2 — `http://127.0.0.1:28124`, нативный порт `29001`;' "$README" +} + +# Состав `make config-test`. Разбор пробников добавлен в него отдельной +# проверкой, и README должен называть её: иначе читатель считает, что дешёвая +# ступень трогает только Compose, и гоняет полный стенд ради того, что видно +# без него. Проверка стоит на двух соседних строках — см. ниже про перенос. +probe_checks_documented() { + awk ' + previous == "`make config-test` проверяет Compose, синтаксис Bash и Python, малые проверки" && + $0 == "логики пробников и пробельные ошибки в diff без запуска стенда." { + found = 1 + } + {previous = $0} + END {exit !found} + ' "$README" +} + +# Урок из предшественника: `make up` не трогает уже созданные контейнеры, и +# после правки настройки метрик серверы молча работают со старой +# конфигурацией. README обязан требовать явный перезапуск. Проверка сверяет +# две соседние строки целиком, а не подстроку: так фразу нельзя незаметно +# разорвать переносом или переписать наполовину. +restart_lesson_present() { + awk ' + previous == "После изменения `infra/clickhouse/config.d/prometheus.xml` выполните" && + $0 == "`docker compose restart clickhouse-01 clickhouse-02`: обычный `make up` не" { + found = 1 + } + {previous = $0} + END {exit !found} + ' "$README" +} + +# Единица памяти в отчёте проверки — русская «ГБ», а не латинская «GB» +# (контракт языка из AGENTS.md). Здесь успех — это отсутствие образца, поэтому +# код возврата grep разбирается вручную: 1 — не нашли, и это хорошо; 0 — нашли +# латинское; больше 1 — сам grep не отработал, и молчать об этом нельзя. +smoke_uses_russian_unit() { + local status=0 + grep -q '3,4 GB' "$SMOKE" || status=$? + case "$status" in + 1) return 0 ;; + 0) return 1 ;; + *) fail "не удалось проверить обозначение единицы памяти в $SMOKE" ;; + esac +} + +check 'README перечисляет HTTP- и нативные порты обеих нод' ports_documented +check 'README объясняет сброс томов после смены исходных учётных данных' clean_advice_present +check 'список портов остаётся единым списком' ports_stay_one_list +check 'README перечисляет малые проверки пробников в составе config-test' probe_checks_documented +check 'README требует перезапуск ClickHouse после изменения настройки метрик' restart_lesson_present +check 'отчёт проверки использует русское обозначение ГБ' smoke_uses_russian_unit + +printf 'ИТОГ: пройдено %d, ошибок 0\n' "$passed" diff --git a/tests/stand-smoke-guards.sh b/tests/stand-smoke-guards.sh index 54c54cd..b91b3a2 100755 --- a/tests/stand-smoke-guards.sh +++ b/tests/stand-smoke-guards.sh @@ -4,13 +4,100 @@ set -uo pipefail readonly ROOT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd)" read -r -a COMPOSE_CMD <<<"${COMPOSE_BIN:-docker compose}" readonly LOG_FILE="$(mktemp)" +readonly CLICKHOUSE_LOCAL_TABLE="airflow_probe_local" +readonly CLICKHOUSE_DISTRIBUTED_TABLE="airflow_probe_distributed" restored=0 passed=0 +red_smoke_pid='' compose() { "${COMPOSE_CMD[@]}" --project-directory "$ROOT_DIR" "$@" } +published_port() { + local binding + local service="$1" + local container_port="$2" + binding="$(compose port "$service" "$container_port" 2>/dev/null || true)" + printf '%s\n' "${binding##*:}" +} + +query_node_2() { + local sql="$1" + local port + port="$(published_port clickhouse-02 8123)" + curl -sf --max-time 2 --data-binary "$sql" "http://127.0.0.1:${port}/" +} + +break_clickhouse_probe() { + local deadline + local table_state + + deadline=$((EPOCHSECONDS + 60)) + printf 'Ожидание служебных таблиц на ноде 2, предел 60 секунд...\n' + while [[ "$EPOCHSECONDS" -lt "$deadline" ]]; do + table_state="$(query_node_2 " + SELECT + countIf(name = '${CLICKHOUSE_LOCAL_TABLE}'), + countIf(name = '${CLICKHOUSE_DISTRIBUTED_TABLE}') + FROM system.tables + WHERE database = 'default' + AND name IN ( + '${CLICKHOUSE_LOCAL_TABLE}', + '${CLICKHOUSE_DISTRIBUTED_TABLE}' + ) + FORMAT TSV + " 2>/dev/null || true)" + if [[ "$table_state" == $'1\t1' ]]; then + query_node_2 \ + "DROP TABLE IF EXISTS default.${CLICKHOUSE_LOCAL_TABLE} SYNC" \ + >/dev/null + table_state="$(query_node_2 " + SELECT + countIf(name = '${CLICKHOUSE_LOCAL_TABLE}'), + countIf(name = '${CLICKHOUSE_DISTRIBUTED_TABLE}') + FROM system.tables + WHERE database = 'default' + AND name IN ( + '${CLICKHOUSE_LOCAL_TABLE}', + '${CLICKHOUSE_DISTRIBUTED_TABLE}' + ) + FORMAT TSV + " 2>/dev/null || true)" + if [[ "$table_state" == $'0\t1' ]]; then + printf 'Служебная локальная таблица удалена только на ноде 2.\n' + return + fi + return 1 + fi + sleep 0.05 + done + return 1 +} + +clickhouse_break_is_reported() { + compose exec -T airflow-scheduler bash -ceu ' + latest="$( + find /opt/airflow/logs/dag_id=test_clickhouse \ + -type f -name "attempt=1.log" -printf "%T@ %p\n" | + sort -nr | + head -n 1 + )" + latest="${latest#* }" + test -n "$latest" + grep -Eq \ + "неверный набор таблиц на ноде 2|Unknown table expression identifier '\''default.airflow_probe_local'\''" \ + "$latest" + ' +} + +probe_failure_is_reported() { + local dag_id="$1" + grep -Eq \ + "ОШИБКА: (пробник ${dag_id}|Airflow не показал пробник ${dag_id})" \ + "$LOG_FILE" +} + restore_stand() { make --no-print-directory -C "$ROOT_DIR" COMPOSE="${COMPOSE_BIN:-docker compose}" up >/dev/null } @@ -18,6 +105,10 @@ restore_stand() { on_exit() { local status=$? trap - EXIT INT TERM + if [[ -n "$red_smoke_pid" ]]; then + kill -TERM -- "-$red_smoke_pid" >/dev/null 2>&1 || true + wait "$red_smoke_pid" >/dev/null 2>&1 || true + fi rm -f "$LOG_FILE" if [[ "$restored" -eq 0 ]] && ! restore_stand; then printf 'ОШИБКА: не удалось восстановить стенд после проверки.\n' >&2 @@ -38,14 +129,45 @@ else fi compose stop prometheus >/dev/null -make --no-print-directory -C "$ROOT_DIR" COMPOSE="${COMPOSE_BIN:-docker compose}" smoke >"$LOG_FILE" 2>&1 +setsid make --no-print-directory -C "$ROOT_DIR" \ + COMPOSE="${COMPOSE_BIN:-docker compose}" smoke >"$LOG_FILE" 2>&1 & +red_smoke_pid=$! +if ! break_clickhouse_probe; then + printf 'ОШИБКА: не удалось удалить служебную таблицу только на ноде 2.\n' >&2 + exit 1 +fi +compose stop kafka >/dev/null +wait "$red_smoke_pid" status=$? +red_smoke_pid='' -if [[ "$status" -ne 0 ]] && grep -q 'ОШИБКА: сервис prometheus' "$LOG_FILE"; then +red_path_ok=1 +if [[ "$status" -eq 0 ]]; then + printf 'ОШИБКА: make smoke остался зелёным после внесённых поломок.\n' >&2 + red_path_ok=0 +fi +if ! grep -q 'ОШИБКА: сервис prometheus' "$LOG_FILE"; then + printf 'ОШИБКА: make smoke не назвал остановленный prometheus.\n' >&2 + red_path_ok=0 +fi +if ! probe_failure_is_reported test_clickhouse; then + printf 'ОШИБКА: make smoke не назвал пробник test_clickhouse.\n' >&2 + red_path_ok=0 +fi +if ! probe_failure_is_reported test_kafka; then + printf 'ОШИБКА: make smoke не назвал пробник test_kafka.\n' >&2 + red_path_ok=0 +fi +if ! clickhouse_break_is_reported; then + printf 'ОШИБКА: журнал test_clickhouse не связал отказ с удалённой таблицей на ноде 2.\n' >&2 + red_path_ok=0 +fi + +if [[ "$red_path_ok" -eq 1 ]]; then passed=$((passed + 1)) - printf 'ЗЕЛЁНО: остановленный prometheus делает make smoke красным и назван в отчёте.\n' + printf 'ЗЕЛЁНО: удаление таблицы на ноде 2 и остановка prometheus с kafka делают make smoke красным; оба пробника названы в отчёте.\n' else - printf 'ОШИБКА: make smoke не обнаружил остановленный prometheus; код=%s.\n' "$status" >&2 + printf 'ОШИБКА: проверка красного пути завершилась с кодом make smoke %s.\n' "$status" >&2 exit 1 fi