diff --git a/.changeset/kafka-table-engine.md b/.changeset/kafka-table-engine.md new file mode 100644 index 00000000..07238fde --- /dev/null +++ b/.changeset/kafka-table-engine.md @@ -0,0 +1,13 @@ +--- +"@chkit/core": minor +"@chkit/clickhouse": patch +"@chkit/plugin-pull": minor +"chkit": minor +--- + +Support Kafka engine tables without MergeTree key clauses, with escaped literal +settings, pull round trips, and normalized drift/check comparisons. Preserve the +existing setting-string contract for other engines. Reject unsupported Kafka +changes before writing migration artifacts and document an explicit, destructive- +gated replacement workflow. Handle quoted delimiters and escaped trailing +backslashes in introspection and migration statement splitting. diff --git a/.github/workflows/kafka.yml b/.github/workflows/kafka.yml new file mode 100644 index 00000000..da91f8ed --- /dev/null +++ b/.github/workflows/kafka.yml @@ -0,0 +1,48 @@ +name: Kafka integration + +on: + pull_request: + paths: + - 'packages/**' + - 'chkit_python/**' + - 'test/kafka/**' + - '.github/workflows/kafka.yml' + - 'bun.lock' + push: + branches: [main] + +permissions: + contents: read + +jobs: + kafka: + runs-on: blacksmith-8vcpu-ubuntu-2404 + strategy: + fail-fast: false + matrix: + clickhouse: ['25.3', '26.3'] + env: + CLICKHOUSE_VERSION: ${{ matrix.clickhouse }} + steps: + - uses: actions/checkout@v4 + - uses: ./.github/actions/setup + - uses: actions/setup-python@v5 + with: + python-version: '3.12' + - name: Build TypeScript and install Python + run: | + bunx turbo run build --filter=chkit... --filter=@chkit/plugin-pull + python -m venv .venv + .venv/bin/pip install -e './chkit_python[dev]' + - name: Start Kafka and ClickHouse + run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml up -d --wait + - name: TypeScript Kafka workflow + run: bun test test/kafka/kafka.e2e.test.ts + - name: Python Kafka workflow + run: .venv/bin/python -m pytest test/kafka/test_python_e2e.py -q + - name: Service logs + if: failure() + run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml logs --tail=150 + - name: Clean up + if: always() + run: docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml down -v diff --git a/apps/docs/src/content/docs/guides/clickhouse-compatibility.md b/apps/docs/src/content/docs/guides/clickhouse-compatibility.md index 8571cbab..2fb0db9a 100644 --- a/apps/docs/src/content/docs/guides/clickhouse-compatibility.md +++ b/apps/docs/src/content/docs/guides/clickhouse-compatibility.md @@ -13,6 +13,12 @@ The continuous test suite runs against ObsessionDB, so the `SharedMergeTree`/`Sh ## Version-gated features +[Kafka engine tables](/schema/kafka/) support creation, pull, and schema drift, +validated on self-hosted ClickHouse 25.3 and 26.3. Existing queue changes require +explicit replacement; generic Kafka column/settings ALTERs are not generated. +The server must provide the Kafka engine and broker connectivity. Distributed and +other integration engines remain outside this support scope. + A few schema features depend on the ClickHouse version of your target: | Feature | Requirement | diff --git a/apps/docs/src/content/docs/schema/dsl-reference.mdx b/apps/docs/src/content/docs/schema/dsl-reference.mdx index 965eb426..510f2c21 100644 --- a/apps/docs/src/content/docs/schema/dsl-reference.mdx +++ b/apps/docs/src/content/docs/schema/dsl-reference.mdx @@ -168,6 +168,10 @@ Creates a table definition. ### Optional fields +For [Kafka tables](/schema/kafka/), omit `primaryKey` and `orderBy`. Kafka also +rejects storage clauses (`partitionBy`, `uniqueKey`, `ttl`, indexes, projections) +and column defaults. Its setting strings are escaped SQL literals. + | Field | Type | Description | |-------|------|-------------| | `partitionBy` | `string` | Partition expression, e.g. `'toYYYYMM(created_at)'` | @@ -907,6 +911,9 @@ chkit validates schema definitions and throws a `ChxValidationError` if any issu ## Structural vs. alterable properties +The table rules below apply to MergeTree-family tables. Kafka changes require an +[explicit replacement](/schema/kafka/#changing-a-queue); chkit refuses generic ALTERs. + When a property changes, chkit determines whether the table can be altered in place or must be dropped and recreated. **Structural** (drop + recreate): `engine`, `primaryKey`, `orderBy`, `partitionBy`, `uniqueKey` diff --git a/apps/docs/src/content/docs/schema/kafka.md b/apps/docs/src/content/docs/schema/kafka.md new file mode 100644 index 00000000..7a914650 --- /dev/null +++ b/apps/docs/src/content/docs/schema/kafka.md @@ -0,0 +1,155 @@ +--- +title: Kafka tables +description: Define Kafka queues and materialized-view ingestion pipelines in schema code. +sidebar: + order: 4 +--- + +Use `table({ engine: 'Kafka', ... })` to manage a Kafka queue together with its +materialized view and storage table. Kafka tables have columns and engine settings, +but no `primaryKey`, `orderBy`, partitioning, TTL, indexes, or projections. + +```ts +import { schema, table, materializedView } from '@chkit/core' + +const columns = [ + { name: 'id', type: 'String' }, + { name: 'event_time', type: 'DateTime64(3)' }, +] + +const queue = table({ + database: 'analytics', name: 'events_queue', engine: 'Kafka', columns, + settings: { + kafka_broker_list: 'kafka01:9092,kafka02:9092', + kafka_topic_list: 'events', + kafka_group_name: 'analytics_events', + kafka_format: 'JSONEachRow', + kafka_num_consumers: 1, + input_format_skip_unknown_fields: true, + }, +}) + +const events = table({ + database: 'analytics', name: 'events', engine: 'MergeTree', columns, + primaryKey: ['event_time', 'id'], orderBy: ['event_time', 'id'], + partitionBy: 'toYYYYMM(event_time)', +}) + +const consumer = materializedView({ + database: 'analytics', name: 'events_consumer', + to: { database: 'analytics', name: 'events' }, + as: 'SELECT id, event_time FROM analytics.events_queue', +}) + +export default schema(queue, events, consumer) +``` + +The Python DSL supports the same pipeline natively: + +```python +from chkit import materialized_view, schema, table + +columns = [{"name": "id", "type": "String"}, + {"name": "event_time", "type": "DateTime64(3)"}] +queue = table( + database="analytics", name="events_queue", engine="Kafka", columns=columns, + settings={ + "kafka_broker_list": "kafka01:9092,kafka02:9092", + "kafka_topic_list": "events", + "kafka_group_name": "analytics_events", + "kafka_format": "JSONEachRow", + "kafka_num_consumers": 1, + "input_format_skip_unknown_fields": True, + }, +) +events = table( + database="analytics", name="events", engine="MergeTree", columns=columns, + primary_key=["event_time", "id"], order_by=["event_time", "id"], + partition_by="toYYYYMM(event_time)", +) +consumer = materialized_view( + database="analytics", name="events_consumer", + to={"database": "analytics", "name": "events"}, + as_="SELECT id, event_time FROM analytics.events_queue", +) +definitions = schema(queue, events, consumer) +``` + +Run `chkit generate`, review the SQL, then `chkit migrate --apply`. Tables are +created before materialized views. Attaching the view starts background consumption. +The Kafka engine must be available on the target server, which must be able to +reach the brokers. This workflow is tested on self-hosted ClickHouse 25.3 and 26.3. + +## Settings and validation + +Kafka settings are literal values: strings are quoted and escaped, numbers remain +numbers, and booleans render as `1` or `0`. Pass `kafka_format: 'JSONEachRow'`, not +a string containing SQL quotes. Existing MergeTree setting strings keep their +previous raw SQL behavior; this addition does not reinterpret older snapshots. + +With `engine: 'Kafka'` or `'Kafka()'`, supply nonempty `kafka_broker_list`, +`kafka_topic_list`, `kafka_group_name`, and `kafka_format` settings. Positional +arguments such as `Kafka('broker:9092', 'topic', 'group', 'JSONEachRow')` and +`Kafka(named_collection)` are also retained on pull. The server validates those +arguments and named collections. + +Kafka columns cannot have `DEFAULT` values. Compute defaults in the consuming +materialized view instead. Format-specific settings are passed through; check +their availability on your ClickHouse version. + +`kafka_auto_offset_reset` is not a standard ClickHouse Kafka table setting. Set +`auto_offset_reset` in the server's extended Kafka configuration. See the +[ClickHouse Kafka reference](https://clickhouse.com/docs/engines/table-engines/integrations/kafka). + +## Changing a queue + +`generate` refuses changes to an existing Kafka table's columns, settings, engine, +or comment with `kafka_change_requires_replacement`. It leaves migrations and the +snapshot untouched. Kafka does not support the generic column/settings ALTER +operations used for MergeTree tables, and automatically replacing an active queue +could disrupt ingestion. + +For an intentional replacement: + +1. Remove the queue and **all consuming materialized views** from your schema. + Keep the destination storage table. Generate a migration with + `chkit generate --name stop-events-queue`. +2. Re-add the updated queue and views. Generate a second migration with + `chkit generate --name restart-events-queue`. +3. Review both migrations, the interruption window, and the consumer group and + offset behavior. Apply with `chkit migrate --apply --allow-destructive`. + +The generated drop migration removes views before synchronously dropping the +queue; the second migration creates the queue before its views. The storage table +is retained, and both generated snapshots remain consistent with the migration +history. This is an explicit replacement, not a zero-downtime operation. Existing +consumer-group offsets, broker retention, and in-flight batches determine whether +messages resume, replay, or are unavailable. chkit does not promise exactly-once +delivery or reset offsets. Coordinate other consumers and unmanaged views yourself. + +`generate --empty` still provides manual SQL, but intentionally does not update +schema snapshots. Use the two generated migrations above when replacing a managed +queue so the snapshot follows the change. + +## Pull, drift, and scope + +`chkit pull schema` emits a Kafka definition without sorting keys and decodes its +settings into literal values. `drift` and `check` compare the engine, columns, and +settings declared in the snapshot. Numeric/boolean metadata and SQL string quoting +are normalized without removing meaningful whitespace from strings. + +In Python, use `chkit pull --database analytics --out-file schema.py`, +`chkit drift --live`, and `chkit check --live`. Python's default drift/check also +inspect local schema changes against the snapshot and report Kafka changes that +require replacement, including when scoped with `--table`. + +Consumer offsets, lag, assignments, topic contents, and server configuration or +named-collection contents are runtime/external state, outside schema drift. Server +settings absent from the desired schema are not treated as drift. + +Prefer server-side configuration for credentials. Pull warns when credential +settings are returned in plaintext or redacted; `[HIDDEN]` credentials must be +resolved before generating SQL. Values evaluated from environment variables in +schema code are still written into snapshots and migration files. + +Distributed and other integration engines are outside this release's scope. diff --git a/apps/docs/src/content/docs/schema/overview.mdx b/apps/docs/src/content/docs/schema/overview.mdx index 3e38562a..a8ccf060 100644 --- a/apps/docs/src/content/docs/schema/overview.mdx +++ b/apps/docs/src/content/docs/schema/overview.mdx @@ -69,6 +69,7 @@ chkit discovers schema files using the `schema` glob in your [configuration](/co - [DSL Reference](/schema/dsl-reference/) — every function, option, and column type, in both languages. - [Refreshable Views](/schema/refreshable-views/) — using ClickHouse refreshable materialized views from chkit. +- [Kafka tables](/schema/kafka/) — queue → materialized view → storage pipelines, pull, drift, and explicit replacements. ## Related diff --git a/chkit_python/CHANGELOG.md b/chkit_python/CHANGELOG.md index ab230e06..2bbf2b55 100644 --- a/chkit_python/CHANGELOG.md +++ b/chkit_python/CHANGELOG.md @@ -2,10 +2,21 @@ ## Unreleased +### Added - Add `SkipIndexText` for full-text index generation, introspection, pull, and drift. Preserve quoted SQL literals, normalize ClickHouse’s fixed granularity, and reject malformed or unsupported metadata. Exercise adversarial round trips and actual indexed search results on ClickHouse 26.3 and 26.8. +- Native Kafka table definitions without sorting keys, escaped literal settings, + pull round trips, and normalized live drift/check comparisons. +- Validation and migration guards reject unsupported Kafka changes before + writing artifacts. Offline drift/check report changes requiring replacement; + explicit drop/create migrations preserve the existing destructive-operation gate. +- Kafka ingestion and replacement integration tests on ClickHouse 25.3 and 26.3. + +### Fixed +- Preserve quoted clause names, delimiters, whitespace, and escaped trailing + backslashes in table introspection and migration statement splitting. ## 0.2.0 — 2026-08-10 diff --git a/chkit_python/src/chkit/cli/commands/check.py b/chkit_python/src/chkit/cli/commands/check.py index 8846efc6..e7fe2aee 100644 --- a/chkit_python/src/chkit/cli/commands/check.py +++ b/chkit_python/src/chkit/cli/commands/check.py @@ -22,6 +22,7 @@ from chkit.cli.commands.drift_compare import summarize_drift_reasons from chkit.cli.commands.drift_payload import build_drift_payload from chkit.cli.commands.migrate_scope import filter_pending_by_scope +from chkit.cli.commands.snapshot_drift import plan_snapshot_drift from chkit.cli.config_loader import load_config from chkit.cli.journal_store import JournalStore from chkit.cli.migration_store import ( @@ -32,14 +33,12 @@ from chkit.cli.plugin_runtime import load_plugin_runtime from chkit.cli.schema_loader import load_schema from chkit.cli.table_scope import ( - filter_plan_by_table_scope, resolve_table_scope, table_keys_from_definitions, ) from chkit.clickhouse.client import ClickHouseClient from chkit.core.canonical import canonicalize_definitions from chkit.core.model import ChxConfigEnv -from chkit.core.planner import plan_diff from chkit.core.validate import validate_definitions from chkit.plugins import ChxOnCheckContext, ChxPlugin @@ -106,13 +105,11 @@ def run( # noqa: PLR0912, PLR0915 drift_ops: list[str] = [] drift_reason_counts: dict[str, int] = {} + replacement_drift = False if snapshot is not None: - plan = plan_diff(snapshot_defs, schema_defs) - if table_scope.enabled: - filtered = filter_plan_by_table_scope( - plan, set(table_scope.matched_tables) - ) - plan = filtered.plan + plan, replacement_issues = plan_snapshot_drift(snapshot_defs, schema_defs, table_scope) + replacement_drift = bool(replacement_issues) + issues.extend(issue.model_dump(mode="json") for issue in replacement_issues) drift_ops = [op.key for op in plan.operations] live_drifted = False @@ -157,9 +154,7 @@ def run( # noqa: PLR0912, PLR0915 failed_checks.append("pending_migrations") if fail_on_mismatch and mismatches: failed_checks.append("checksum_mismatch") - drift_fired = ( - bool(drift_ops) if not live else live_drifted or bool(drift_ops) - ) + drift_fired = bool(drift_ops) or replacement_drift or (live and live_drifted) if fail_on_drift and drift_fired: # Match the TS finding code: ``schema_drift`` (was ``drift``). failed_checks.append("schema_drift") diff --git a/chkit_python/src/chkit/cli/commands/drift.py b/chkit_python/src/chkit/cli/commands/drift.py index 28512800..8c3bdef2 100644 --- a/chkit_python/src/chkit/cli/commands/drift.py +++ b/chkit_python/src/chkit/cli/commands/drift.py @@ -1,9 +1,7 @@ -"""`chkit drift` — compare the on-disk snapshot against the current schema. +"""`chkit drift` — compare the snapshot against the current schema or live DB. -Mirrors the TypeScript ``driftCommand`` for the snapshot-vs-schema check. -The TS port also reaches into ClickHouse to compare against the live database; -that full DB-side introspection is not in this first-base Python port, so the -command focuses on the snapshot/schema diff produced by ``plan_diff``. +The default inspects local schema changes; ``--live`` compares the snapshot +against ClickHouse, matching the TypeScript command's default behavior. Output (human): @@ -31,18 +29,17 @@ import typer from chkit.cli.commands.drift_payload import build_drift_payload +from chkit.cli.commands.snapshot_drift import plan_snapshot_drift from chkit.cli.config_loader import load_config from chkit.cli.migration_store import read_snapshot from chkit.cli.schema_loader import load_schema from chkit.cli.table_scope import ( - filter_plan_by_table_scope, resolve_table_scope, table_keys_from_definitions, ) from chkit.clickhouse.client import ClickHouseClient from chkit.core.canonical import canonicalize_definitions from chkit.core.model import ChxConfigEnv -from chkit.core.planner import plan_diff def run( # noqa: PLR0912, PLR0915 @@ -148,20 +145,18 @@ def run( # noqa: PLR0912, PLR0915 typer.echo(f"- {detail.table}: {', '.join(detail.reason_codes)}") return - plan = plan_diff(snapshot_defs, schema_defs) - if table_scope.enabled: - filtered = filter_plan_by_table_scope(plan, set(table_scope.matched_tables)) - plan = filtered.plan + plan, issues = plan_snapshot_drift(snapshot_defs, schema_defs, table_scope) + drifted = bool(plan.operations or issues) snapshot_file = meta_dir / "snapshot.json" payload = { "snapshotFile": str(snapshot_file), - "drifted": bool(plan.operations), + "drifted": drifted, "operations": [op.model_dump(by_alias=True) for op in plan.operations], - "renameSuggestions": [ - s.model_dump(by_alias=True) for s in plan.rename_suggestions - ], + "renameSuggestions": [s.model_dump(by_alias=True) for s in plan.rename_suggestions], } + if issues: + payload["issues"] = [issue.model_dump(mode="json") for issue in issues] if output_json: typer.echo(json.dumps(payload, indent=2)) @@ -169,7 +164,9 @@ def run( # noqa: PLR0912, PLR0915 typer.echo(f"Snapshot file: {snapshot_file}") typer.echo(f"Expected operations: {len(plan.operations)}") - typer.echo(f"Drifted: {'yes' if plan.operations else 'no'}") + typer.echo(f"Drifted: {'yes' if drifted else 'no'}") + for issue in issues: + typer.echo(f"- [{issue.code}] {issue.message}") if plan.operations: typer.echo("") typer.echo("Operations:") diff --git a/chkit_python/src/chkit/cli/commands/drift_compare.py b/chkit_python/src/chkit/cli/commands/drift_compare.py index 99ee29d7..71869ecb 100644 --- a/chkit_python/src/chkit/cli/commands/drift_compare.py +++ b/chkit_python/src/chkit/cli/commands/drift_compare.py @@ -21,6 +21,7 @@ diff_settings, ) from chkit.clickhouse.introspect import IntrospectedTable, SchemaObjectKind +from chkit.core.kafka import is_kafka_engine, kafka_setting_fingerprint, parse_kafka_settings from chkit.core.model import ( ColumnDefinition, ProjectionDefinition, @@ -344,7 +345,18 @@ def compare_table_shape( # noqa: PLR0912, PLR0915 extra_columns = column_diff.extra changed_columns = column_diff.changed - setting_diffs = diff_settings(expected.settings or {}, actual.settings) + if is_kafka_engine(expected.engine): + expected_settings = { + key: kafka_setting_fingerprint(value) + for key, value in (expected.settings or {}).items() + } + actual_settings = { + key: kafka_setting_fingerprint(value) + for key, value in parse_kafka_settings(actual.settings).items() + } + setting_diffs = diff_settings(expected_settings, actual_settings) + else: + setting_diffs = diff_settings(expected.settings or {}, actual.settings) expected_indexes = { idx.name: _normalize_index_shape(idx) for idx in (expected.indexes or []) diff --git a/chkit_python/src/chkit/cli/commands/pull.py b/chkit_python/src/chkit/cli/commands/pull.py index c33d5deb..9f2ed103 100644 --- a/chkit_python/src/chkit/cli/commands/pull.py +++ b/chkit_python/src/chkit/cli/commands/pull.py @@ -55,6 +55,7 @@ list_schema_objects, list_table_details, ) +from chkit.core.kafka import is_kafka_engine, parse_kafka_settings from chkit.core.model import ( ChxConfigEnv, DictionaryDefinition, @@ -85,7 +86,12 @@ def _introspected_table_to_definition( if not item.columns: return None - settings = {k: _coerce_setting_value(v) for k, v in (item.settings or {}).items()} + kafka = is_kafka_engine(item.engine or "") + settings = ( + parse_kafka_settings(item.settings or {}) + if kafka + else {k: _coerce_setting_value(v) for k, v in (item.settings or {}).items()} + ) indexes: list[ SkipIndexMinmax @@ -102,8 +108,8 @@ def _introspected_table_to_definition( name=item.name, engine=item.engine or "MergeTree", columns=list(item.columns), - primary_key=_split_clause(item.primary_key) or [item.columns[0].name], - order_by=_split_clause(item.order_by) or [item.columns[0].name], + primary_key=[] if kafka else _split_clause(item.primary_key) or [item.columns[0].name], + order_by=[] if kafka else _split_clause(item.order_by) or [item.columns[0].name], unique_key=_split_clause(item.unique_key) or None, partition_by=item.partition_by or None, ttl=item.ttl or None, @@ -352,6 +358,24 @@ def _summarize_skipped_objects( ) +def _kafka_setting_warnings(definitions: Sequence[SchemaDefinition]) -> list[str]: + warnings: list[str] = [] + for definition in definitions: + if not isinstance(definition, TableDefinition) or not is_kafka_engine(definition.engine): + continue + for key, value in (definition.settings or {}).items(): + if re.search(r"password|secret|token", key, re.IGNORECASE): + label = f'Kafka table "{definition.database}.{definition.name}" setting {key}' + warnings.append( + f"{label} was redacted by ClickHouse. Restore it or use server-side " + "configuration before generating migrations." + if value == "[HIDDEN]" + else f"{label} contains a credential returned by ClickHouse and written " + "into the schema file. Prefer server-side configuration." + ) + return warnings + + def _write_schema_file(out_file: Path, content: str, *, overwrite: bool) -> None: if out_file.exists() and not overwrite: msg = ( @@ -459,6 +483,7 @@ def run( # noqa: PLR0917 ) warnings = _dictionary_password_warnings(definitions) + warnings.extend(_kafka_setting_warnings(definitions)) payload: dict[str, object] = { "command": "schema", diff --git a/chkit_python/src/chkit/cli/commands/pull_render.py b/chkit_python/src/chkit/cli/commands/pull_render.py index 857b8108..c723cca5 100644 --- a/chkit_python/src/chkit/cli/commands/pull_render.py +++ b/chkit_python/src/chkit/cli/commands/pull_render.py @@ -33,6 +33,7 @@ from collections.abc import Sequence from chkit.core.canonical import canonicalize_definitions +from chkit.core.kafka import is_kafka_engine from chkit.core.model import ( ColumnCodec, ColumnCodecSpec, @@ -182,8 +183,9 @@ def _render_table(variable_name: str, definition: TableDefinition) -> list[str]: ] lines.extend(f" {_render_column(column)}," for column in definition.columns) lines.append(" ],") - lines.append(f" primary_key={_render_string_list(definition.primary_key)},") - lines.append(f" order_by={_render_string_list(definition.order_by)},") + if not is_kafka_engine(definition.engine): + lines.append(f" primary_key={_render_string_list(definition.primary_key)},") + lines.append(f" order_by={_render_string_list(definition.order_by)},") if definition.unique_key: lines.append(f" unique_key={_render_string_list(definition.unique_key)},") if definition.partition_by: diff --git a/chkit_python/src/chkit/cli/commands/snapshot_drift.py b/chkit_python/src/chkit/cli/commands/snapshot_drift.py new file mode 100644 index 00000000..ff39a97c --- /dev/null +++ b/chkit_python/src/chkit/cli/commands/snapshot_drift.py @@ -0,0 +1,43 @@ +"""Inspect local drift without turning unsupported Kafka changes into SQL.""" + +from chkit.cli.table_scope import TableScope, filter_plan_by_table_scope +from chkit.core.model import ( + ChxValidationError, + MigrationPlan, + SchemaDefinition, + ValidationIssue, +) +from chkit.core.planner import plan_diff + + +def plan_snapshot_drift( + previous: list[SchemaDefinition], + current: list[SchemaDefinition], + scope: TableScope, +) -> tuple[MigrationPlan, list[ValidationIssue]]: + """Report blocked replacements alongside the remaining, plannable changes. + + Migration generation still fails on these issues. Read-only checks can + report them and continue inspecting other tables, including scoped checks. + """ + issues: list[ValidationIssue] = [] + while True: + try: + plan = plan_diff(previous, current) + break + except ChxValidationError as error: + if not error.issues or any( + issue.code != "kafka_change_requires_replacement" for issue in error.issues + ): + raise + blocked = {(issue.kind, issue.database, issue.name) for issue in error.issues} + issues.extend( + issue + for issue in error.issues + if not scope.enabled or f"{issue.database}.{issue.name}" in scope.matched_tables + ) + previous = [d for d in previous if (d.kind, d.database, d.name) not in blocked] + current = [d for d in current if (d.kind, d.database, d.name) not in blocked] + if scope.enabled: + plan = filter_plan_by_table_scope(plan, set(scope.matched_tables)).plan + return plan, issues diff --git a/chkit_python/src/chkit/clickhouse/create_table_parser.py b/chkit_python/src/chkit/clickhouse/create_table_parser.py index fc0b03e8..392d4dad 100644 --- a/chkit_python/src/chkit/clickhouse/create_table_parser.py +++ b/chkit_python/src/chkit/clickhouse/create_table_parser.py @@ -20,6 +20,7 @@ from chkit.core.key_clause import split_top_level_comma from chkit.core.projection import normalize_projection_index from chkit.core.sql_normalizer import normalize_sql_fragment +from chkit.core.sql_scan import find_top_level_sql_pattern __all__ = [ "ProjectionDefinitionShape", @@ -44,14 +45,15 @@ class ProjectionDefinitionShape: type: str | None = None -_SETTINGS_RE = re.compile(r"\bSETTINGS\b(.*?)(?:;|$)", re.IGNORECASE | re.DOTALL) +_SETTINGS_RE = re.compile(r"\bSETTINGS\b", re.IGNORECASE) +_SETTINGS_STOP = re.compile(r"\bCOMMENT\b|;", re.IGNORECASE) _TTL_RE = re.compile(r"\bTTL\b(.*?)(?:\bSETTINGS\b|;|$)", re.IGNORECASE | re.DOTALL) _BODY_ENGINE_RE = re.compile(r"\)\s*ENGINE\s*=", re.IGNORECASE) _ENGINE_START = re.compile(r"\bENGINE\s*=\s*", re.IGNORECASE) _ENGINE_STOP = re.compile( r"\bPRIMARY\s+KEY\b|\bORDER\s+BY\b|\bPARTITION\s+BY\b|\bUNIQUE\s+KEY\b" - r"|\bSAMPLE\s+BY\b|\bTTL\b|\bSETTINGS\b|;|$", + r"|\bSAMPLE\s+BY\b|\bTTL\b|\bSETTINGS\b|\bCOMMENT\b|;|$", re.IGNORECASE, ) @@ -98,7 +100,10 @@ class ProjectionDefinitionShape: def _parse_clause( - query: str | None, start_pattern: re.Pattern[str], stop_pattern: re.Pattern[str] + query: str | None, + start_pattern: re.Pattern[str], + stop_pattern: re.Pattern[str], + preserve_whitespace: bool = False, ) -> str | None: """Slice between ``start_pattern`` and the first ``stop_pattern`` hit.""" if not query: @@ -109,16 +114,16 @@ def _parse_clause( # swallow the real clause plus everything up to the next stop keyword # (issue #190). options = _extract_table_options(query) - start = start_pattern.search(options) + start = find_top_level_sql_pattern(options, start_pattern) if start is None: return None after = options[start.end() :] - stop = stop_pattern.search(after) + stop = find_top_level_sql_pattern(after, stop_pattern) raw = after[: stop.start()] if stop is not None else after raw = raw.strip() if not raw: return None - return normalize_sql_fragment(raw) + return raw if preserve_whitespace else normalize_sql_fragment(raw) def _find_column_list_bounds(query: str) -> tuple[int, int] | None: @@ -188,10 +193,13 @@ def parse_settings_from_create_table_query(query: str | None) -> dict[str, str]: """Extract ``SETTINGS k=v, k=v`` as a dict (last write wins).""" if not query: return {} - match = _SETTINGS_RE.search(_extract_table_options(query)) + options = _extract_table_options(query) + match = find_top_level_sql_pattern(options, _SETTINGS_RE) if match is None: return {} - raw = match.group(1).strip() + tail = options[match.end() :] + stop = find_top_level_sql_pattern(tail, _SETTINGS_STOP) + raw = (tail[: stop.start()] if stop is not None else tail).strip() if not raw: return {} out: dict[str, str] = {} @@ -220,7 +228,7 @@ def parse_ttl_from_create_table_query(query: str | None) -> str | None: def parse_engine_from_create_table_query(query: str | None) -> str | None: - return _parse_clause(query, _ENGINE_START, _ENGINE_STOP) + return _parse_clause(query, _ENGINE_START, _ENGINE_STOP, preserve_whitespace=True) def parse_primary_key_from_create_table_query(query: str | None) -> str | None: diff --git a/chkit_python/src/chkit/clickhouse/introspect.py b/chkit_python/src/chkit/clickhouse/introspect.py index 0e1e776a..7d5c6ce3 100644 --- a/chkit_python/src/chkit/clickhouse/introspect.py +++ b/chkit_python/src/chkit/clickhouse/introspect.py @@ -295,7 +295,7 @@ def build_introspected_tables( name=table_row.name, engine=parse_engine_from_create_table_query( table_row.create_table_query - ), + ) or table_row.engine, primary_key=parse_primary_key_from_create_table_query( table_row.create_table_query ), diff --git a/chkit_python/src/chkit/core/kafka.py b/chkit_python/src/chkit/core/kafka.py new file mode 100644 index 00000000..a4478184 --- /dev/null +++ b/chkit_python/src/chkit/core/kafka.py @@ -0,0 +1,79 @@ +"""Kafka engine recognition and literal settings, matching the TypeScript core.""" + +from __future__ import annotations + +import re + +from chkit.core.key_clause import split_top_level_comma + + +def is_kafka_engine(engine: str) -> bool: + return re.match(r"^Kafka\s*(?:\(|$)", engine.strip(), re.IGNORECASE) is not None + + +def render_kafka_setting(value: str | int | float | bool) -> str: + if isinstance(value, str): + escaped = value.replace("\\", "\\\\").replace("'", "''") + return f"'{escaped}'" + if isinstance(value, bool): + return "1" if value else "0" + return str(value) + + +def parse_kafka_setting(value: str) -> str | int | float | bool: + trimmed = value.strip() + if trimmed.startswith("'") and trimmed.endswith("'"): + escapes: dict[str, str] = { + "n": "\n", + "r": "\r", + "t": "\t", + "b": "\b", + "f": "\f", + "a": "\a", + "v": "\v", + "0": "\0", + "N": "", + "\\": "\\", + "'": "'", + '"': '"', + } + data = bytearray() + tokens: list[str] = re.findall(r"\\x[\da-fA-F]{2}|\\[\s\S]|''|[\s\S]", trimmed[1:-1]) + for token in tokens: + if re.fullmatch(r"\\x[\da-fA-F]{2}", token): + data.append(int(token[2:], 16)) + else: + decoded = token + if token == "''": + decoded = "'" + elif token.startswith("\\"): + decoded = escapes.get(token[1:], token) + data.extend(decoded.encode("utf-8")) + return data.decode("utf-8", errors="replace") + if trimmed.lower() in {"true", "false"}: + return trimmed.lower() == "true" + if re.fullmatch(r"-?\d+", trimmed): + return int(trimmed) + return trimmed + + +def parse_kafka_settings(settings: dict[str, str]) -> dict[str, str | int | float | bool]: + return {key: parse_kafka_setting(value) for key, value in settings.items()} + + +def kafka_setting_fingerprint(value: str | int | float | bool) -> str: + if isinstance(value, bool): + return "1" if value else "0" + if isinstance(value, float) and value.is_integer(): + return str(int(value)) + return str(value) + + +def normalize_kafka_engine(engine: str) -> str: + match = re.fullmatch(r"Kafka\s*\(([\s\S]*)\)", engine.strip(), re.IGNORECASE) + args = split_top_level_comma(match.group(1) if match else "") + normalized = [ + render_kafka_setting(parse_kafka_setting(arg)) if arg.startswith("'") else arg.strip() + for arg in args + ] + return f"Kafka({', '.join(normalized)})" diff --git a/chkit_python/src/chkit/core/key_clause.py b/chkit_python/src/chkit/core/key_clause.py index cdf8b72b..01fbff87 100644 --- a/chkit_python/src/chkit/core/key_clause.py +++ b/chkit_python/src/chkit/core/key_clause.py @@ -24,12 +24,19 @@ def split_top_level_comma(text: str) -> list[str]: current: list[str] = [] depth = 0 quote: str | None = None + skip_next = False for i, ch in enumerate(text): - prev = text[i - 1] if i > 0 else "" + if skip_next: + skip_next = False + continue + nxt = text[i + 1] if i + 1 < len(text) else "" if quote is not None: current.append(ch) - if ch == quote and prev != "\\": + if nxt and (ch == "\\" or ch == quote == nxt): + current.append(nxt) + skip_next = True + elif ch == quote: quote = None continue diff --git a/chkit_python/src/chkit/core/model.py b/chkit_python/src/chkit/core/model.py index 981272d5..b93a80bf 100644 --- a/chkit_python/src/chkit/core/model.py +++ b/chkit_python/src/chkit/core/model.py @@ -670,6 +670,11 @@ class MigrationPlan(_StrictModel): ValidationIssueCode: TypeAlias = Literal[ + "kafka_unsupported_clause", + "kafka_column_default", + "kafka_missing_setting", + "kafka_invalid_setting", + "kafka_change_requires_replacement", "duplicate_object_name", "duplicate_column_name", "duplicate_index_name", @@ -762,6 +767,11 @@ def table( ) -> TableDefinition: pk = primary_key if primary_key is not None else primaryKey ob = order_by if order_by is not None else orderBy + from chkit.core.kafka import is_kafka_engine + + if is_kafka_engine(engine): + pk = pk or [] + ob = ob or [] if pk is None or ob is None: msg = "table() requires primary_key/primaryKey and order_by/orderBy" raise ValueError(msg) diff --git a/chkit_python/src/chkit/core/planner.py b/chkit_python/src/chkit/core/planner.py index 78f4dc0c..93c69e83 100644 --- a/chkit_python/src/chkit/core/planner.py +++ b/chkit_python/src/chkit/core/planner.py @@ -7,7 +7,9 @@ from chkit.core.canonical import canonicalize_definitions, definition_key from chkit.core.diff_primitives import diff_by_name, diff_clauses, diff_settings +from chkit.core.kafka import is_kafka_engine, kafka_setting_fingerprint from chkit.core.model import ( + ChxValidationError, ColumnDefinition, ColumnRenameSuggestion, DictionaryDefinition, @@ -19,6 +21,7 @@ SchemaDefinition, SkipIndexDefinition, TableDefinition, + ValidationIssue, ViewDefinition, _RiskSummary, ) @@ -57,7 +60,8 @@ def _push_drop( type="drop_table", key=definition_key(definition), risk=risk, - sql=f"DROP TABLE IF EXISTS {definition.database}.{definition.name};", + sql=f"DROP TABLE IF EXISTS {definition.database}.{definition.name}" + + (" SYNC;" if is_kafka_engine(definition.engine) else ";"), ) ) return @@ -364,6 +368,39 @@ def _diff_materialized_view( def _diff_tables( old: TableDefinition, new: TableDefinition ) -> tuple[list[MigrationOperation], list[ColumnRenameSuggestion]]: + if is_kafka_engine(old.engine) or is_kafka_engine(new.engine): + old_settings = { + key: kafka_setting_fingerprint(value) for key, value in (old.settings or {}).items() + } + new_settings = { + key: kafka_setting_fingerprint(value) for key, value in (new.settings or {}).items() + } + old_columns = [(column.name, _column_identity(column)) for column in old.columns] + new_columns = [(column.name, _column_identity(column)) for column in new.columns] + if ( + _requires_table_recreate(old, new) + or old_columns != new_columns + or old_settings != new_settings + or (old.comment or "") != (new.comment or "") + ): + raise ChxValidationError( + [ + ValidationIssue( + code="kafka_change_requires_replacement", + kind="table", + database=new.database, + name=new.name, + message=f"Kafka table {new.database}.{new.name} requires an explicit " + "replacement; column, engine and setting ALTERs are not supported. " + "Remove the queue and its consuming materialized views from the " + "schema and generate a drop migration, then re-add the updated " + "definitions and generate a create migration. Review " + "both migrations and consumer-group/offset behavior before applying with " + "--allow-destructive.", + ) + ] + ) + return [], [] if _requires_table_recreate(old, new): return ( [ diff --git a/chkit_python/src/chkit/core/sql.py b/chkit_python/src/chkit/core/sql.py index b7accbdf..78a8863f 100644 --- a/chkit_python/src/chkit/core/sql.py +++ b/chkit_python/src/chkit/core/sql.py @@ -8,6 +8,7 @@ from pydantic import TypeAdapter from chkit.core.codec import render_codec +from chkit.core.kafka import is_kafka_engine, render_kafka_setting from chkit.core.key_clause import is_plain_column_reference, normalize_key_columns from chkit.core.model import ( ColumnDefinition, @@ -143,12 +144,13 @@ def _render_table_sql(definition: TableDefinition) -> str: clauses: list[str] = [] if definition.partition_by is not None: clauses.append(f"PARTITION BY {definition.partition_by}") - clauses.append( - f"PRIMARY KEY ({_render_key_clause_columns(definition.primary_key, column_names)})" - ) - clauses.append( - f"ORDER BY ({_render_key_clause_columns(definition.order_by, column_names)})" - ) + if not is_kafka_engine(definition.engine): + clauses.append( + f"PRIMARY KEY ({_render_key_clause_columns(definition.primary_key, column_names)})" + ) + clauses.append( + f"ORDER BY ({_render_key_clause_columns(definition.order_by, column_names)})" + ) if definition.unique_key is not None and len(definition.unique_key) > 0: clauses.append( f"UNIQUE KEY ({_render_key_clause_columns(definition.unique_key, column_names)})" @@ -156,7 +158,15 @@ def _render_table_sql(definition: TableDefinition) -> str: if definition.ttl is not None: clauses.append(f"TTL {definition.ttl}") if definition.settings is not None and len(definition.settings) > 0: - clauses.append(f"SETTINGS {_render_settings_clause(definition.settings)}") + rendered_settings = ( + ", ".join( + f"{key} = {render_kafka_setting(value)}" + for key, value in definition.settings.items() + ) + if is_kafka_engine(definition.engine) + else _render_settings_clause(definition.settings) + ) + clauses.append(f"SETTINGS {rendered_settings}") if definition.comment is not None and len(definition.comment) > 0: escaped = definition.comment.replace("'", "''") clauses.append(f"COMMENT '{escaped}'") diff --git a/chkit_python/src/chkit/core/sql_normalizer.py b/chkit_python/src/chkit/core/sql_normalizer.py index 3d0dc368..ed21e6d9 100644 --- a/chkit_python/src/chkit/core/sql_normalizer.py +++ b/chkit_python/src/chkit/core/sql_normalizer.py @@ -5,6 +5,8 @@ import re from typing import Final +from chkit.core.kafka import is_kafka_engine, normalize_kafka_engine + _WHITESPACE: Final[re.Pattern[str]] = re.compile(r"\s+") @@ -13,6 +15,8 @@ def normalize_sql_fragment(value: str) -> str: def normalize_engine(engine: str) -> str: + if is_kafka_engine(engine): + return normalize_kafka_engine(engine) normalized = engine.strip() if normalized.startswith("Shared"): normalized = normalized[len("Shared") :] diff --git a/chkit_python/src/chkit/core/sql_scan.py b/chkit_python/src/chkit/core/sql_scan.py new file mode 100644 index 00000000..f5f72172 --- /dev/null +++ b/chkit_python/src/chkit/core/sql_scan.py @@ -0,0 +1,34 @@ +"""Locate SQL clauses outside string literals, identifiers and parentheses.""" + +from __future__ import annotations + +import re + + +def find_top_level_sql_pattern(sql: str, pattern: re.Pattern[str]) -> re.Match[str] | None: + quote: str | None = None + depth = 0 + i = 0 + while i < len(sql): + char = sql[i] + if quote is not None: + if char == "\\": + i += 2 + continue + if char == quote: + if i + 1 < len(sql) and sql[i + 1] == quote: + i += 2 + continue + quote = None + elif char in {"'", '"', "`"}: + quote = char + elif char == "(": + depth += 1 + elif char == ")": + depth -= 1 + elif depth == 0: + match = pattern.match(sql, i) + if match is not None: + return match + i += 1 + return None diff --git a/chkit_python/src/chkit/core/sql_splitter.py b/chkit_python/src/chkit/core/sql_splitter.py index 33035536..0fbe8930 100644 --- a/chkit_python/src/chkit/core/sql_splitter.py +++ b/chkit_python/src/chkit/core/sql_splitter.py @@ -29,10 +29,14 @@ def _handle_in_block_comment(state: _SplitterState, ch: str, nxt: str) -> int: return 1 -def _handle_in_quote(state: _SplitterState, ch: str, prev: str) -> None: +def _handle_in_quote(state: _SplitterState, ch: str, nxt: str) -> int: state.current.append(ch) - if ch == state.quote and prev != "\\": + if nxt and (ch == "\\" or ch == state.quote == nxt): + state.current.append(nxt) + return 2 + if ch == state.quote: state.quote = None + return 1 def _flush_statement(state: _SplitterState) -> None: @@ -55,7 +59,6 @@ def split_sql_statements(text: str) -> list[str]: while i < n: ch = text[i] nxt = text[i + 1] if i + 1 < n else "" - prev = text[i - 1] if i > 0 else "" if state.in_line_comment: _handle_in_line_comment(state, ch) @@ -65,8 +68,7 @@ def split_sql_statements(text: str) -> list[str]: i += _handle_in_block_comment(state, ch, nxt) continue if state.quote is not None: - _handle_in_quote(state, ch, prev) - i += 1 + i += _handle_in_quote(state, ch, nxt) continue if ch == "-" and nxt == "-": state.current.append(ch) diff --git a/chkit_python/src/chkit/core/validate.py b/chkit_python/src/chkit/core/validate.py index a7e96374..d279a614 100644 --- a/chkit_python/src/chkit/core/validate.py +++ b/chkit_python/src/chkit/core/validate.py @@ -2,12 +2,14 @@ from __future__ import annotations +import math import re from collections.abc import Iterable from typing import Final from chkit.core.canonical import definition_key from chkit.core.codec import canonicalize_codec, is_general_codec, is_raw_codec +from chkit.core.kafka import is_kafka_engine from chkit.core.key_clause import is_plain_column_reference, normalize_key_columns from chkit.core.model import ( ChxValidationError, @@ -108,7 +110,78 @@ def _validate_indexes(definition: TableDefinition, issues: list[ValidationIssue] f'Text index "{index.name}": {exc}') +def _validate_kafka_table(definition: TableDefinition, issues: list[ValidationIssue]) -> None: + if not is_kafka_engine(definition.engine): + return + label = f"Kafka table {definition.database}.{definition.name}" + clauses = { + "primaryKey": definition.primary_key, + "orderBy": definition.order_by, + "uniqueKey": definition.unique_key, + "partitionBy": definition.partition_by, + "ttl": definition.ttl, + "indexes": definition.indexes, + "projections": definition.projections, + } + for field, value in clauses.items(): + if value: + _push( + issues, + definition, + "kafka_unsupported_clause", + f"{label} does not support {field}. Put storage clauses on the destination table.", + ) + for column in definition.columns: + if column.default is not None: + _push( + issues, + definition, + "kafka_column_default", + f'{label} column "{column.name}" cannot have a DEFAULT. ' + "Compute defaults in the materialized view.", + ) + settings = definition.settings or {} + if re.fullmatch(r"Kafka\s*(?:\(\s*\))?", definition.engine.strip(), re.IGNORECASE): + for key in ("kafka_broker_list", "kafka_topic_list", "kafka_group_name", "kafka_format"): + setting = settings.get(key) + if not isinstance(setting, str) or not setting.strip(): + _push( + issues, + definition, + "kafka_missing_setting", + f"{label} requires a nonempty {key} string " + "(or engine arguments / a named collection).", + ) + for key, setting in settings.items(): + if not re.fullmatch(r"[A-Za-z_][A-Za-z0-9_]*", key) or ( + isinstance(setting, float) and not math.isfinite(setting) + ): + _push( + issues, + definition, + "kafka_invalid_setting", + f"{label} has an invalid setting key or value: {key}.", + ) + if key == "kafka_auto_offset_reset": + _push( + issues, + definition, + "kafka_invalid_setting", + "kafka_auto_offset_reset is not a Kafka table setting on standard ClickHouse. " + "Configure auto_offset_reset in the server Kafka configuration.", + ) + if re.search(r"password|secret|token", key, re.IGNORECASE) and setting == "[HIDDEN]": + _push( + issues, + definition, + "kafka_invalid_setting", + f"Kafka setting {key} was redacted by ClickHouse. Restore the credential or " + "use server-side configuration before generating migrations.", + ) + + def _validate_table(definition: TableDefinition, issues: list[ValidationIssue]) -> None: + _validate_kafka_table(definition, issues) column_seen: set[str] = set() column_set: set[str] = set() for column in definition.columns: diff --git a/chkit_python/tests/test_kafka.py b/chkit_python/tests/test_kafka.py new file mode 100644 index 00000000..4afe5901 --- /dev/null +++ b/chkit_python/tests/test_kafka.py @@ -0,0 +1,168 @@ +"""Kafka DSL, SQL, parser, pull and migration-safety parity with TypeScript.""" + +from dataclasses import replace + +import pytest + +from chkit import materialized_view, table +from chkit.cli.commands.drift_compare import compare_table_shape +from chkit.cli.commands.pull import _introspected_table_to_definition +from chkit.cli.commands.pull_render import render_schema_file +from chkit.cli.commands.snapshot_drift import plan_snapshot_drift +from chkit.cli.table_scope import resolve_table_scope +from chkit.clickhouse.create_table_parser import ( + parse_engine_from_create_table_query, + parse_settings_from_create_table_query, +) +from chkit.clickhouse.introspect import IntrospectedTable +from chkit.core.kafka import parse_kafka_setting, render_kafka_setting +from chkit.core.model import ChxValidationError +from chkit.core.planner import plan_diff +from chkit.core.sql import to_create_sql +from chkit.core.sql_normalizer import normalize_engine +from chkit.core.sql_splitter import extract_executable_statements +from chkit.core.validate import validate_definitions + + +def queue(): + return table( + database="app", + name="queue", + engine="Kafka", + columns=[{"name": "id", "type": "String"}], + settings={ + "kafka_broker_list": "a:9092,b:9092", + "kafka_topic_list": "events", + "kafka_group_name": "consumer", + "kafka_format": "JSONEachRow", + "kafka_num_consumers": 1, + "kafka_commit_on_select": False, + }, + ) + + +def test_queue_sql_and_validation(): + sql = to_create_sql(queue()) + assert "PRIMARY KEY" not in sql + assert "ORDER BY" not in sql + assert "kafka_broker_list = 'a:9092,b:9092'" in sql + assert "kafka_commit_on_select = 0" in sql + assert validate_definitions([queue()]) == [] + bad = queue().model_copy(update={"primary_key": ["id"], "ttl": "id"}) + assert [issue.code for issue in validate_definitions([bad])] == [ + "kafka_unsupported_clause", + "kafka_unsupported_clause", + ] + assert len(validate_definitions([queue().model_copy(update={"settings": {}})])) == 4 + default = queue().columns[0].model_copy(update={"default": ""}) + assert ( + validate_definitions([queue().model_copy(update={"columns": [default]})])[0].code + == "kafka_column_default" + ) + with pytest.raises(ValueError, match="requires primary_key"): + table(database="app", name="stored", engine="MergeTree", columns=[]) + + +@pytest.mark.parametrize( + "value", ["a'b", "a\\b\\", "a; SETTINGS COMMENT", "a\nb\tc", "'quoted'", "é"] +) +def test_literal_round_trip(value): + assert parse_kafka_setting(render_kafka_setting(value)) == value + + +def test_hex_control_and_unknown_escapes(): + assert parse_kafka_setting(r"'\xC3\xA9\a\v\N\q\%'") == "é\a\v\\q\\%" + + +def test_parser_pull_and_drift(): + original = queue() + sql = to_create_sql(original) + settings = parse_settings_from_create_table_query(sql) + actual = IntrospectedTable( + database="app", + name="queue", + engine="Kafka()", + columns=original.columns, + settings=settings, + indexes=[], + projections=[], + ) + assert compare_table_shape(original, actual) is None + pulled = _introspected_table_to_definition(actual) + assert pulled is not None + assert plan_diff([original], [pulled]).operations == [] + source = render_schema_file([pulled]) + assert "primary_key=" not in source + assert "order_by=" not in source + assert "'a:9092,b:9092'" in source or '"a:9092,b:9092"' in source + changed = replace(actual, settings={**settings, "kafka_group_name": "'different'"}) + drift = compare_table_shape(original, changed) + assert drift is not None + assert drift.setting_diffs == ["kafka_group_name"] + + +def test_quoted_delimiters_and_engine_arguments(): + engine = "Kafka('host:9092', 'SETTINGS; COMMENT', 'a b', 'JSONEachRow')" + sql = f"CREATE TABLE q (id String) ENGINE = {engine} SETTINGS kafka_client_id = 'a;\\\\';" + assert parse_engine_from_create_table_query(sql) == engine + assert ( + parse_kafka_setting(parse_settings_from_create_table_query(sql)["kafka_client_id"]) + == "a;\\" + ) + assert len(extract_executable_statements(sql + "\nSELECT 1;")) == 2 + assert ( + normalize_engine("Kafka( 'b', 'a b', 'c', 'JSONEachRow' )") + == "Kafka('b', 'a b', 'c', 'JSONEachRow')" + ) + + +def test_rejects_unsupported_changes_and_drops_synchronously(): + original = queue() + for changed in [ + original.model_copy( + update={"settings": {**(original.settings or {}), "kafka_num_consumers": 2}} + ), + original.model_copy( + update={"columns": [original.columns[0].model_copy(update={"name": "renamed"})]} + ), + original.model_copy( + update={"engine": "MergeTree", "primary_key": ["id"], "order_by": ["id"]} + ), + ]: + with pytest.raises(ChxValidationError) as error: + plan_diff([original], [changed]) + assert error.value.issues[0].code == "kafka_change_requires_replacement" + mv = materialized_view( + database="app", + name="mv", + to={"database": "app", "name": "stored"}, + as_="SELECT id FROM app.queue", + ) + drops = plan_diff([original, mv], []).operations + assert [op.type for op in drops] == ["drop_materialized_view", "drop_table"] + assert drops[1].sql == "DROP TABLE IF EXISTS app.queue SYNC;" + + +def test_snapshot_checks_report_replacements_and_keep_other_scoped_changes(): + original = queue() + changed = original.model_copy(update={"comment": "changed"}) + stored = table( + database="app", + name="stored", + engine="MergeTree", + columns=list(original.columns), + primary_key=["id"], + order_by=["id"], + ) + edited = stored.model_copy(update={"settings": {"index_granularity": 4096}}) + for selector, expected_issues, expected_operations in [ + (None, 1, 1), + ("queue", 1, 0), + ("stored", 0, 1), + ("missing", 0, 0), + ]: + scope = resolve_table_scope(selector, ["app.queue", "app.stored"]) + plan, issues = plan_snapshot_drift([original, stored], [changed, edited], scope) + assert len(issues) == expected_issues + assert len(plan.operations) == expected_operations + assert all("app.queue" not in operation.key for operation in plan.operations) diff --git a/packages/cli/src/commands/drift/compare.ts b/packages/cli/src/commands/drift/compare.ts index ffa74a66..97adfc24 100644 --- a/packages/cli/src/commands/drift/compare.ts +++ b/packages/cli/src/commands/drift/compare.ts @@ -1,5 +1,8 @@ import { normalizeEngine as coreNormalizeEngine, + isKafkaEngine, + parseKafkaSettings, + kafkaSettingFingerprint, isIndexProjection, normalizeProjectionIndex, normalizeSQLFragment, @@ -279,7 +282,14 @@ export function compareTableShape(expected: TableDefinition, actual: ActualTable const extraColumns = columnDiff.extra const changedColumns = columnDiff.changed - const settingDiffs = diffSettings(expected.settings ?? {}, actual.settings) + const kafka = isKafkaEngine(expected.engine) + const expectedSettings = kafka + ? Object.fromEntries(Object.entries(expected.settings ?? {}).map(([key, value]) => [key, kafkaSettingFingerprint(value)])) + : expected.settings ?? {} + const actualSettings = kafka + ? Object.fromEntries(Object.entries(parseKafkaSettings(actual.settings)).map(([key, value]) => [key, kafkaSettingFingerprint(value)])) + : actual.settings + const settingDiffs = diffSettings(expectedSettings, actualSettings) const expectedIndexes = new Map( (expected.indexes ?? []).map((idx) => [idx.name, normalizeIndexShape(idx)]) diff --git a/packages/cli/src/test/kafka-drift.test.ts b/packages/cli/src/test/kafka-drift.test.ts new file mode 100644 index 00000000..76a22bfd --- /dev/null +++ b/packages/cli/src/test/kafka-drift.test.ts @@ -0,0 +1,27 @@ +import { expect, test } from 'bun:test' +import { table } from '@chkit/core' +import { compareTableShape } from '../commands/drift/compare.js' + +test('Kafka drift compares decoded setting literals, preserving significant string whitespace', () => { + const expected = table({ + database: 'app', + name: 'q', + engine: 'Kafka', + columns: [{ name: 'id', type: 'String' }], + settings: { kafka_client_id: 'a b', kafka_commit_on_select: false, kafka_num_consumers: 1 }, + }) + const actual = { + engine: 'Kafka()', + columns: expected.columns, + indexes: [], + projections: [], + settings: { kafka_client_id: "'a b'", kafka_commit_on_select: '0', kafka_num_consumers: '1' }, + } + expect(compareTableShape(expected, actual)).toBeNull() + expect( + compareTableShape(expected, { + ...actual, + settings: { ...actual.settings, kafka_client_id: "'a b'" }, + })?.settingDiffs, + ).toEqual(['kafka_client_id']) +}) diff --git a/packages/clickhouse/src/create-table-parser.ts b/packages/clickhouse/src/create-table-parser.ts index 09fbed99..5011c335 100644 --- a/packages/clickhouse/src/create-table-parser.ts +++ b/packages/clickhouse/src/create-table-parser.ts @@ -1,4 +1,4 @@ -import { normalizeProjectionIndex, normalizeSQLFragment, splitTopLevelComma } from '@chkit/core' +import { findTopLevelSQLPattern, normalizeProjectionIndex, normalizeSQLFragment, splitTopLevelComma } from '@chkit/core' type ProjectionDefinitionShape = | { name: string; query: string } @@ -7,7 +7,8 @@ type ProjectionDefinitionShape = function parseClauseFromCreateTableQuery( createTableQuery: string | undefined, clausePattern: RegExp, - stopPattern: RegExp + stopPattern: RegExp, + preserveWhitespace = false ): string | undefined { if (!createTableQuery) return undefined // Table-level clauses (ENGINE, ORDER BY, PRIMARY KEY, ...) only appear after @@ -15,13 +16,13 @@ function parseClauseFromCreateTableQuery( // body — e.g. the `ORDER BY` of a projection's SELECT — and swallow the real // clause plus everything up to the next stop keyword (issue #190). const options = extractTableOptions(createTableQuery) - const start = options.match(clausePattern) - if (!start || start.index === undefined) return undefined - const afterClause = options.slice(start.index + start[0].length) - const stop = afterClause.match(stopPattern) + const start = findTopLevelSQLPattern(options, clausePattern) + if (!start) return undefined + const afterClause = options.slice(start.index + start.length) + const stop = findTopLevelSQLPattern(afterClause, stopPattern) const raw = (stop ? afterClause.slice(0, stop.index) : afterClause).trim() if (!raw) return undefined - return normalizeSQLFragment(raw) + return preserveWhitespace ? raw : normalizeSQLFragment(raw) } /** @@ -88,9 +89,12 @@ function extractTableOptions(createTableQuery: string): string { export function parseSettingsFromCreateTableQuery(createTableQuery: string | undefined): Record { if (!createTableQuery) return {} - const settingsMatch = extractTableOptions(createTableQuery).match(/\bSETTINGS\b([\s\S]*?)(?:;|$)/i) - if (!settingsMatch?.[1]) return {} - const rawSettings = settingsMatch[1].trim() + const options = extractTableOptions(createTableQuery) + const start = findTopLevelSQLPattern(options, /\bSETTINGS\b/i) + if (!start) return {} + const tail = options.slice(start.index + start.length) + const stop = findTopLevelSQLPattern(tail, /\bCOMMENT\b|;/i) + const rawSettings = (stop ? tail.slice(0, stop.index) : tail).trim() if (!rawSettings) return {} const items = splitTopLevelComma(rawSettings) const out: Record = {} @@ -117,7 +121,8 @@ export function parseEngineFromCreateTableQuery(createTableQuery: string | undef return parseClauseFromCreateTableQuery( createTableQuery, /\bENGINE\s*=\s*/i, - /\bPRIMARY\s+KEY\b|\bORDER\s+BY\b|\bPARTITION\s+BY\b|\bUNIQUE\s+KEY\b|\bSAMPLE\s+BY\b|\bTTL\b|\bSETTINGS\b|;|$/i + /\bPRIMARY\s+KEY\b|\bORDER\s+BY\b|\bPARTITION\s+BY\b|\bUNIQUE\s+KEY\b|\bSAMPLE\s+BY\b|\bTTL\b|\bSETTINGS\b|\bCOMMENT\b|;|$/i, + true ) } diff --git a/packages/clickhouse/src/index.ts b/packages/clickhouse/src/index.ts index 8f5c482a..29c14dc5 100644 --- a/packages/clickhouse/src/index.ts +++ b/packages/clickhouse/src/index.ts @@ -314,7 +314,7 @@ export function buildIntrospectedTables( return { database: row.database, name: row.name, - engine: parseEngineFromCreateTableQuery(row.create_table_query), + engine: parseEngineFromCreateTableQuery(row.create_table_query) ?? row.engine, primaryKey: parsePrimaryKeyFromCreateTableQuery(row.create_table_query), orderBy: parseOrderByFromCreateTableQuery(row.create_table_query), uniqueKey: parseUniqueKeyFromCreateTableQuery(row.create_table_query), diff --git a/packages/clickhouse/src/kafka-parser.test.ts b/packages/clickhouse/src/kafka-parser.test.ts new file mode 100644 index 00000000..51fbdbef --- /dev/null +++ b/packages/clickhouse/src/kafka-parser.test.ts @@ -0,0 +1,23 @@ +import { expect, test } from 'bun:test' +import { parseKafkaSettings } from '@chkit/core' +import { + parseEngineFromCreateTableQuery, + parseSettingsFromCreateTableQuery, +} from './create-table-parser.js' + +test('parses Kafka settings containing SQL keywords, semicolons, quotes and trailing backslashes', () => { + const sql = String.raw`CREATE TABLE q (id String) ENGINE = Kafka SETTINGS kafka_broker_list = 'a:9092,b:9092', kafka_client_id = 'a; SETTINGS COMMENT ''b'' \\', kafka_num_consumers = 1 COMMENT 'queue';` + expect(parseEngineFromCreateTableQuery(sql)).toBe('Kafka') + expect(parseKafkaSettings(parseSettingsFromCreateTableQuery(sql))).toEqual({ + kafka_broker_list: 'a:9092,b:9092', + kafka_client_id: "a; SETTINGS COMMENT 'b' \\", + kafka_num_consumers: 1, + }) +}) + +test('engine arguments preserve spaces and quoted clause keywords', () => { + const engine = "Kafka('host:9092', 'SETTINGS; COMMENT', 'a b', 'JSONEachRow')" + const sql = `CREATE TABLE q (id String) ENGINE = ${engine} SETTINGS kafka_num_consumers = 1;` + expect(parseEngineFromCreateTableQuery(sql)).toBe(engine) + expect(parseSettingsFromCreateTableQuery(sql)).toEqual({ kafka_num_consumers: '1' }) +}) diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 648cf754..2af4c46e 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -1,5 +1,7 @@ export * from './flags.js' export * from './model.js' +export { isKafkaEngine, parseKafkaSettings, kafkaSettingFingerprint } from './kafka.js' +export { findTopLevelSQLPattern } from './sql-scan.js' export { SYNTHESIZED_CONFIG_PATH, isSynthesizedConfigPath } from './config-path.js' export { canonicalizeDefinition, diff --git a/packages/core/src/kafka.test.ts b/packages/core/src/kafka.test.ts new file mode 100644 index 00000000..48aa068c --- /dev/null +++ b/packages/core/src/kafka.test.ts @@ -0,0 +1,166 @@ +import { describe, expect, test } from 'bun:test' +import { + canonicalizeDefinitions, + materializedView, + normalizeEngine, + planDiff, + table, + toCreateSQL, + validateDefinitions, +} from './index.js' +import { parseKafkaSetting, renderKafkaSetting } from './kafka.js' + +const queue = () => + table({ + database: 'app', + name: 'queue', + engine: 'Kafka', + columns: [{ name: 'id', type: 'String' }], + settings: { + kafka_broker_list: 'one:9092,two:9092', + kafka_topic_list: 'events', + kafka_group_name: 'consumer', + kafka_format: 'JSONEachRow', + kafka_num_consumers: 1, + }, + }) + +describe('Kafka tables', () => { + test('decodes ClickHouse hex/control escapes and preserves unknown escapes', () => { + expect(parseKafkaSetting(String.raw`'\xC3\xA9\a\v\N\q\%'`)).toBe('é\u0007\v\\q\\%') + }) + test('renders a queue without sorting clauses and with literal settings', () => { + const sql = toCreateSQL(queue()) + expect(sql).not.toContain('PRIMARY KEY') + expect(sql).not.toContain('ORDER BY') + expect(sql).toContain("kafka_broker_list = 'one:9092,two:9092'") + expect(sql).toContain('kafka_num_consumers = 1') + expect( + toCreateSQL({ ...queue(), settings: { ...queue().settings, kafka_commit_on_select: false } }), + ).toContain('kafka_commit_on_select = 0') + }) + + test.each([ + "a'b", + 'a\\b\\', + 'a; b, SETTINGS COMMENT', + 'a\nb\tc', + "'quoted'", + ])('round-trips literal %j', (value) => + expect(parseKafkaSetting(renderKafkaSetting(value))).toBe(value)) + + test('supports positional arguments and named collections', () => { + expect( + toCreateSQL( + table({ + database: 'app', + name: 'q', + engine: "Kafka('b:9092', 'topic', 'group', 'JSONEachRow')", + columns: [{ name: 'id', type: 'String' }], + }), + ), + ).toContain('ENGINE = Kafka(') + expect( + toCreateSQL( + table({ + database: 'app', + name: 'q', + engine: 'Kafka(kafka_config)', + columns: [{ name: 'id', type: 'String' }], + }), + ), + ).toContain('ENGINE = Kafka(kafka_config)') + expect(normalizeEngine("Kafka( 'broker', 'a b', 'c\\\\', 'JSONEachRow' )")).toBe( + "Kafka('broker', 'a b', 'c\\\\', 'JSONEachRow')", + ) + }) + + test('rejects invalid Kafka definitions at runtime', () => { + expect( + validateDefinitions([ + { + ...queue(), + primaryKey: ['id'], + ttl: 'id', + projections: [{ name: 'p', query: 'SELECT id' }], + }, + ]).map((x) => x.code), + ).toEqual(['kafka_unsupported_clause', 'kafka_unsupported_clause', 'kafka_unsupported_clause']) + expect( + validateDefinitions([ + { ...queue(), columns: [{ name: 'id', type: 'String', default: '' }] }, + ])[0]?.code, + ).toBe('kafka_column_default') + expect(validateDefinitions([{ ...queue(), settings: {} }])).toHaveLength(4) + expect( + validateDefinitions([ + { ...queue(), settings: { ...queue().settings, kafka_auto_offset_reset: 'earliest' } }, + ])[0]?.code, + ).toBe('kafka_invalid_setting') + expect( + validateDefinitions([ + { ...queue(), settings: { ...queue().settings, kafka_sasl_password: '[HIDDEN]' } }, + ])[0]?.code, + ).toBe('kafka_invalid_setting') + }) + + test('preserves existing MergeTree key and raw-settings behavior', () => { + const stored = table({ + database: 'app', + name: 'stored', + engine: 'MergeTree', + columns: [{ name: 'id', type: 'String' }], + primaryKey: ['id'], + orderBy: ['id'], + settings: { storage_policy: "'default'" }, + }) + expect(toCreateSQL(stored)).toContain("SETTINGS storage_policy = 'default'") + expect(toCreateSQL(stored)).toContain('ORDER BY (`id`)') + // @ts-expect-error MergeTree still requires keys at the public API boundary. + table({ database: 'app', name: 't', engine: 'MergeTree', columns: [] }) + }) + + test('creates queues before MVs and drops MVs before synchronously dropping queues', () => { + const defs = [ + queue(), + materializedView({ + database: 'app', + name: 'mv', + to: { database: 'app', name: 'stored' }, + as: 'SELECT id FROM app.queue', + }), + ] + expect(planDiff([], defs).operations.map((x) => x.type)).toEqual([ + 'create_database', + 'create_table', + 'create_materialized_view', + ]) + const drops = planDiff(defs, []).operations + expect(drops.map((x) => x.type)).toEqual(['drop_materialized_view', 'drop_table']) + expect(drops[1]?.sql).toBe('DROP TABLE IF EXISTS app.queue SYNC;') + expect(drops[1]?.risk).toBe('danger') + }) + + test('refuses unsupported ALTER and implicit engine replacement', () => { + const original = queue() + for (const updated of [ + { ...original, settings: { ...original.settings, kafka_num_consumers: 2 } }, + { ...original, settings: { ...original.settings, kafka_group_name: 'new-group' } }, + { ...original, columns: [...original.columns, { name: 'extra', type: 'String' }] }, + { ...original, columns: [{ name: 'renamed', type: 'String' }] }, + { ...original, engine: 'MergeTree', primaryKey: ['id'], orderBy: ['id'] }, + ]) + expect(() => planDiff([original], [updated])).toThrow('Schema validation failed') + }) + + test('no-op round trips tolerate numeric setting metadata and empty key normalization', () => { + const original = queue() + const pulled = { + ...original, + engine: 'Kafka()', + settings: { ...original.settings, kafka_num_consumers: '1' }, + } + expect(planDiff([original], [pulled]).operations).toEqual([]) + expect(planDiff(canonicalizeDefinitions([original]), [original]).operations).toEqual([]) + }) +}) diff --git a/packages/core/src/kafka.ts b/packages/core/src/kafka.ts new file mode 100644 index 00000000..31496607 --- /dev/null +++ b/packages/core/src/kafka.ts @@ -0,0 +1,72 @@ +import { splitTopLevelComma } from './key-clause.js' + +export function isKafkaEngine(engine: string): boolean { + return /^Kafka\s*(?:\(|$)/i.test(engine.trim()) +} + +/** Kafka has literal settings; existing engines retain their raw SQL setting contract. */ +export function renderKafkaSetting(value: string | number | boolean): string { + if (typeof value === 'string') { + return `'${value.replace(/\\/g, '\\\\').replace(/'/g, "''")}'` + } + return typeof value === 'boolean' ? (value ? '1' : '0') : String(value) +} + +/** Decode a ClickHouse literal from SHOW CREATE, without evaluating expressions. */ +export function parseKafkaSetting(value: string): string | number | boolean { + const trimmed = value.trim() + if (trimmed.startsWith("'") && trimmed.endsWith("'")) { + const escapes: Record = { + n: '\n', + r: '\r', + t: '\t', + b: '\b', + f: '\f', + a: '\u0007', + v: '\v', + '0': '\0', + N: '', + '\\': '\\', + "'": "'", + '"': '"', + } + const bytes: number[] = [] + const encoder = new TextEncoder() + for (const token of trimmed.slice(1, -1).match(/\\x[\da-fA-F]{2}|\\[\s\S]|''|[\s\S]/gu) ?? []) { + if (/^\\x[\da-fA-F]{2}$/.test(token)) bytes.push(Number.parseInt(token.slice(2), 16)) + else { + const decoded = + token === "''" ? "'" : token.startsWith('\\') ? (escapes[token[1] ?? ''] ?? token) : token + bytes.push(...encoder.encode(decoded)) + } + } + return new TextDecoder().decode(new Uint8Array(bytes)) + } + if (/^(true|false)$/i.test(trimmed)) return trimmed.toLowerCase() === 'true' + if (/^-?\d+(?:\.\d+)?$/.test(trimmed) && Number.isSafeInteger(Number(trimmed))) + return Number(trimmed) + // Keep large integers lossless. String comparison handles numeric metadata. + return trimmed +} + +export function parseKafkaSettings( + settings: Record, +): Record { + return Object.fromEntries( + Object.entries(settings).map(([key, value]) => [key, parseKafkaSetting(value)]), + ) +} + +export function kafkaSettingFingerprint(value: string | number | boolean): string { + return typeof value === 'boolean' ? (value ? '1' : '0') : String(value) +} + +export function normalizeKafkaEngine(engine: string): string { + const args = engine.trim().match(/^Kafka\s*\(([\s\S]*)\)$/i)?.[1] ?? '' + return `Kafka(${splitTopLevelComma(args) + .map((arg) => { + const parsed = parseKafkaSetting(arg) + return arg.trim().startsWith("'") ? renderKafkaSetting(parsed) : arg.trim() + }) + .join(', ')})` +} diff --git a/packages/core/src/key-clause.ts b/packages/core/src/key-clause.ts index 7de7f3c4..495c2d78 100644 --- a/packages/core/src/key-clause.ts +++ b/packages/core/src/key-clause.ts @@ -6,11 +6,14 @@ export function splitTopLevelComma(input: string): string[] { for (let i = 0; i < input.length; i += 1) { const char = input[i] ?? '' - const prev = i > 0 ? input[i - 1] : '' if (quote) { current += char - if (char === quote && prev !== '\\') quote = null + if (char === '\\' && i + 1 < input.length) current += input[++i] + else if (char === quote) { + if (input[i + 1] === quote) current += input[++i] + else quote = null + } continue } diff --git a/packages/core/src/model-types.ts b/packages/core/src/model-types.ts index 16515d3b..c0e9f438 100644 --- a/packages/core/src/model-types.ts +++ b/packages/core/src/model-types.ts @@ -159,6 +159,20 @@ export interface TableDefinition { plugins?: TablePlugins } +/** Kafka queues do not have MergeTree sorting/storage clauses. */ +export type KafkaTableInput = Omit & { + engine: 'Kafka' | `Kafka(${string})` + primaryKey?: never + orderBy?: never + partitionBy?: never + uniqueKey?: never + ttl?: never + indexes?: never + projections?: never +} + export interface ViewDefinition { kind: 'view' database: string @@ -390,6 +404,11 @@ export interface MigrationPlan { } export type ValidationIssueCode = + | 'kafka_unsupported_clause' + | 'kafka_column_default' + | 'kafka_missing_setting' + | 'kafka_invalid_setting' + | 'kafka_change_requires_replacement' | 'duplicate_object_name' | 'duplicate_column_name' | 'duplicate_index_name' diff --git a/packages/core/src/model.ts b/packages/core/src/model.ts index 6fb4cb20..6c2b9a3f 100644 --- a/packages/core/src/model.ts +++ b/packages/core/src/model.ts @@ -6,6 +6,7 @@ import type { ChxResolvedConfig, ChxUserConfig, DictionaryDefinition, + KafkaTableInput, MaterializedViewDefinition, SchemaDefinition, TableDefinition, @@ -91,8 +92,10 @@ export function resolveConfig(config: ChxUserConfig): ChxResolvedConfig { } } -export function table(input: Omit): TableDefinition { - return { ...input, kind: 'table' } +export function table(input: KafkaTableInput): TableDefinition +export function table(input: Omit): TableDefinition +export function table(input: KafkaTableInput | Omit): TableDefinition { + return { ...input, primaryKey: input.primaryKey ?? [], orderBy: input.orderBy ?? [], kind: 'table' } } export function view(input: Omit): ViewDefinition { diff --git a/packages/core/src/planner.ts b/packages/core/src/planner.ts index be49b800..5d180f2c 100644 --- a/packages/core/src/planner.ts +++ b/packages/core/src/planner.ts @@ -30,6 +30,8 @@ import { } from './sql.js' import { textIndexFingerprint } from './text-index.js' import { assertValidDefinitions } from './validate.js' +import { ChxValidationError } from './model.js' +import { isKafkaEngine, kafkaSettingFingerprint } from './kafka.js' function createMap(definitions: SchemaDefinition[]): Map { return new Map(definitions.map((def) => [definitionKey(def), def])) @@ -45,7 +47,7 @@ function pushDropOperation( type: 'drop_table', key: definitionKey(def), risk, - sql: `DROP TABLE IF EXISTS ${def.database}.${def.name};`, + sql: `DROP TABLE IF EXISTS ${def.database}.${def.name}${isKafkaEngine(def.engine) ? ' SYNC' : ''};`, }) return } @@ -354,6 +356,21 @@ function diffDictionary( } function diffTables(oldDef: TableDefinition, newDef: TableDefinition): TableDiffResult { + if (isKafkaEngine(oldDef.engine) || isKafkaEngine(newDef.engine)) { + const settings = (def: TableDefinition) => Object.fromEntries( + Object.entries(def.settings ?? {}).map(([key, value]) => [key, kafkaSettingFingerprint(value)]) + ) + if (requiresTableRecreate(oldDef, newDef) + || JSON.stringify(oldDef.columns.map(column => [column.name, normalizeColumn(column)])) !== JSON.stringify(newDef.columns.map(column => [column.name, normalizeColumn(column)])) + || diffSettings(settings(oldDef), settings(newDef)).changes.length > 0 + || (oldDef.comment ?? '') !== (newDef.comment ?? '')) { + throw new ChxValidationError([{ + code: 'kafka_change_requires_replacement', kind: 'table', database: newDef.database, name: newDef.name, + message: `Kafka table ${newDef.database}.${newDef.name} requires an explicit replacement; column, engine and setting ALTERs are not supported. Remove the queue and its consuming materialized views from the schema and generate a drop migration, then re-add the updated definitions and generate a create migration. Review both migrations and consumer-group/offset behavior before applying with --allow-destructive.`, + }]) + } + return { operations: [], renameSuggestions: [] } + } if (requiresTableRecreate(oldDef, newDef)) { return { operations: [ diff --git a/packages/core/src/sql-normalizer.ts b/packages/core/src/sql-normalizer.ts index 9b8bc930..6a107164 100644 --- a/packages/core/src/sql-normalizer.ts +++ b/packages/core/src/sql-normalizer.ts @@ -1,8 +1,11 @@ +import { isKafkaEngine, normalizeKafkaEngine } from './kafka.js' + export function normalizeSQLFragment(value: string): string { return value.replace(/\s+/g, ' ').trim() } export function normalizeEngine(engine: string): string { + if (isKafkaEngine(engine)) return normalizeKafkaEngine(engine) let normalized = engine.trim().replace(/^Shared/, '') if (!normalized.includes('(')) { normalized += '()' diff --git a/packages/core/src/sql-scan.ts b/packages/core/src/sql-scan.ts new file mode 100644 index 00000000..c5270c53 --- /dev/null +++ b/packages/core/src/sql-scan.ts @@ -0,0 +1,39 @@ +/** Find a clause/delimiter outside quoted strings, identifiers and parentheses. */ +export function findTopLevelSQLPattern( + sql: string, + pattern: RegExp, +): { index: number; length: number } | undefined { + let quote: string | undefined + let depth = 0 + for (let i = 0; i < sql.length; i += 1) { + const char = sql[i] + if (quote) { + if (char === '\\') i += 1 + else if (char === quote) { + if (sql[i + 1] === quote) i += 1 + else quote = undefined + } + continue + } + if (char === "'" || char === '"' || char === '`') { + quote = char + continue + } + if (char === '(') { + depth += 1 + continue + } + if (char === ')') { + depth -= 1 + continue + } + if ( + depth !== 0 || + (/[A-Za-z_]/.test(char ?? '') && i > 0 && /[A-Za-z0-9_]/.test(sql[i - 1] ?? '')) + ) + continue + const match = sql.slice(i).match(pattern) + if (match?.index === 0) return { index: i, length: match[0].length } + } + return undefined +} diff --git a/packages/core/src/sql-splitter.test.ts b/packages/core/src/sql-splitter.test.ts index 0b3edea1..f1ab5898 100644 --- a/packages/core/src/sql-splitter.test.ts +++ b/packages/core/src/sql-splitter.test.ts @@ -114,3 +114,12 @@ describe('extractExecutableStatements', () => { ]) }) }) +test('escaped trailing backslash closes a setting literal before the next statement', () => { + const sql = String.raw`CREATE TABLE q (id String) ENGINE = Kafka SETTINGS kafka_client_id = 'a;\\'; +-- next operation +CREATE MATERIALIZED VIEW mv TO stored AS SELECT id FROM q;` + expect(extractExecutableStatements(sql)).toEqual([ + String.raw`CREATE TABLE q (id String) ENGINE = Kafka SETTINGS kafka_client_id = 'a;\\';`, + 'CREATE MATERIALIZED VIEW mv TO stored AS SELECT id FROM q;', + ]) +}) diff --git a/packages/core/src/sql-splitter.ts b/packages/core/src/sql-splitter.ts index d4ab2322..fc24637e 100644 --- a/packages/core/src/sql-splitter.ts +++ b/packages/core/src/sql-splitter.ts @@ -22,12 +22,17 @@ export function splitSqlStatements(sql: string): string[] { if (inSingleQuote) { current += ch + if (ch === '\\' && next) { + current += next + i += 1 + continue + } if (ch === "'" && next === "'") { current += next i += 1 continue } - if (ch === "'" && sql[i - 1] !== '\\') { + if (ch === "'") { inSingleQuote = false } continue @@ -35,7 +40,12 @@ export function splitSqlStatements(sql: string): string[] { if (inDoubleQuote) { current += ch - if (ch === '"' && sql[i - 1] !== '\\') { + if (ch === '\\' && next) { + current += next + i += 1 + continue + } + if (ch === '"') { inDoubleQuote = false } continue @@ -125,12 +135,17 @@ export function extractExecutableStatements(sql: string): string[] { if (inSingleQuote) { stripped += ch + if (ch === '\\' && next) { + stripped += next + i += 1 + continue + } if (ch === "'" && next === "'") { stripped += next i += 1 continue } - if (ch === "'" && sql[i - 1] !== '\\') { + if (ch === "'") { inSingleQuote = false } continue @@ -138,7 +153,12 @@ export function extractExecutableStatements(sql: string): string[] { if (inDoubleQuote) { stripped += ch - if (ch === '"' && sql[i - 1] !== '\\') { + if (ch === '\\' && next) { + stripped += next + i += 1 + continue + } + if (ch === '"') { inDoubleQuote = false } continue diff --git a/packages/core/src/sql.ts b/packages/core/src/sql.ts index 1247633c..8ac3d364 100644 --- a/packages/core/src/sql.ts +++ b/packages/core/src/sql.ts @@ -11,6 +11,7 @@ import type { ViewDefinition, } from './model.js' import { renderCodec } from './codec.js' +import { isKafkaEngine, renderKafkaSetting } from './kafka.js' import { isPlainColumnReference, normalizeKeyColumns } from './key-clause.js' import { renderProjectionBody } from './projection.js' import { TEXT_INDEX_GRANULARITY, renderTextIndexType } from './text-index.js' @@ -76,8 +77,10 @@ function renderTableSQL(def: TableDefinition): string { const columnNames = new Set(def.columns.map((column) => column.name)) const clauses: string[] = [] if (def.partitionBy) clauses.push(`PARTITION BY ${def.partitionBy}`) - clauses.push(`PRIMARY KEY (${renderKeyClauseColumns(def.primaryKey, columnNames)})`) - clauses.push(`ORDER BY (${renderKeyClauseColumns(def.orderBy, columnNames)})`) + if (!isKafkaEngine(def.engine)) { + clauses.push(`PRIMARY KEY (${renderKeyClauseColumns(def.primaryKey, columnNames)})`) + clauses.push(`ORDER BY (${renderKeyClauseColumns(def.orderBy, columnNames)})`) + } if (def.uniqueKey && def.uniqueKey.length > 0) { clauses.push(`UNIQUE KEY (${renderKeyClauseColumns(def.uniqueKey, columnNames)})`) } @@ -85,7 +88,7 @@ function renderTableSQL(def: TableDefinition): string { if (def.settings && Object.keys(def.settings).length > 0) { clauses.push( `SETTINGS ${Object.entries(def.settings) - .map(([k, v]) => `${k} = ${v}`) + .map(([k, v]) => `${k} = ${isKafkaEngine(def.engine) ? renderKafkaSetting(v) : v}`) .join(', ')}` ) } diff --git a/packages/core/src/validate.ts b/packages/core/src/validate.ts index c467fc56..ecc2b68b 100644 --- a/packages/core/src/validate.ts +++ b/packages/core/src/validate.ts @@ -3,6 +3,7 @@ import { definitionKey } from './canonical.js' import { canonicalizeCodec, isGeneralCodec, isRawCodec } from './codec.js' import { isPlainColumnReference, normalizeKeyColumns } from './key-clause.js' import { isIndexProjection, normalizeProjectionIndex } from './projection.js' +import { isKafkaEngine } from './kafka.js' import type { ColumnDefinition, DictionaryDefinition, @@ -77,6 +78,37 @@ function validateColumnCodec( } function validateTableDefinition(def: TableDefinition, issues: ValidationIssue[]): void { + if (isKafkaEngine(def.engine)) { + for (const field of ['primaryKey', 'orderBy', 'uniqueKey', 'partitionBy', 'ttl', 'indexes', 'projections'] as const) { + const value = def[field] + if (Array.isArray(value) ? value.length > 0 : Boolean(value)) { + pushValidationIssue(issues, def, 'kafka_unsupported_clause', `Kafka table ${def.database}.${def.name} does not support ${field}. Put storage clauses on the destination table.`) + } + } + for (const column of def.columns) { + if (column.default !== undefined) { + pushValidationIssue(issues, def, 'kafka_column_default', `Kafka table ${def.database}.${def.name} column "${column.name}" cannot have a DEFAULT. Compute defaults in the materialized view.`) + } + } + if (/^Kafka\s*(?:\(\s*\))?$/i.test(def.engine.trim())) { + for (const key of ['kafka_broker_list', 'kafka_topic_list', 'kafka_group_name', 'kafka_format']) { + if (typeof def.settings?.[key] !== 'string' || !String(def.settings[key]).trim()) { + pushValidationIssue(issues, def, 'kafka_missing_setting', `Kafka table ${def.database}.${def.name} requires a nonempty ${key} string (or engine arguments / a named collection).`) + } + } + } + for (const [key, value] of Object.entries(def.settings ?? {})) { + if (!/^[A-Za-z_][A-Za-z0-9_]*$/.test(key) || (typeof value === 'number' && !Number.isFinite(value))) { + pushValidationIssue(issues, def, 'kafka_invalid_setting', `Kafka table ${def.database}.${def.name} has an invalid setting key or value: ${key}.`) + } + if (key === 'kafka_auto_offset_reset') { + pushValidationIssue(issues, def, 'kafka_invalid_setting', 'kafka_auto_offset_reset is not a Kafka table setting on standard ClickHouse. Configure auto_offset_reset in the server Kafka configuration.') + } + if (/password|secret|token/i.test(key) && value === '[HIDDEN]') { + pushValidationIssue(issues, def, 'kafka_invalid_setting', `Kafka setting ${key} was redacted by ClickHouse. Restore the credential or use server-side configuration before generating migrations.`) + } + } + } const columnSeen = new Set() const columnSet = new Set() for (const column of def.columns) { diff --git a/packages/plugin-pull/src/index.ts b/packages/plugin-pull/src/index.ts index 5636b288..cf8379f6 100644 --- a/packages/plugin-pull/src/index.ts +++ b/packages/plugin-pull/src/index.ts @@ -16,6 +16,8 @@ import { type DictionaryDefinition, type FlagMapping, normalizeEngine, + isKafkaEngine, + parseKafkaSettings, type ResolvedChxConfig, type SafeParseable, type SchemaDefinition, @@ -327,7 +329,15 @@ async function pullSchema(input: { const content = renderSchemaFile(definitions) const tableCount = definitions.filter((definition) => definition.kind === 'table').length const skippedObjects = summarizeSkippedObjects(objects, definitions, selectedDatabases) - const warnings = dictionaryPasswordWarnings(definitions) + const warnings = [ + ...dictionaryPasswordWarnings(definitions), + ...definitions.flatMap((def) => def.kind === 'table' && isKafkaEngine(def.engine) + ? Object.entries(def.settings ?? {}).filter(([key]) => /password|secret|token/i.test(key)).map(([key, value]) => + value === '[HIDDEN]' + ? `Kafka table "${def.database}.${def.name}" setting ${key} was redacted by ClickHouse. Restore it or use server-side configuration before generating migrations.` + : `Kafka table "${def.database}.${def.name}" setting ${key} contains a credential returned by ClickHouse and written into the schema file. Prefer server-side configuration.`) + : []), + ] return { outFile, @@ -367,7 +377,7 @@ function mapIntrospectedTableToDefinition(table: IntrospectedTable): TableDefini ...(table.uniqueKey ? { uniqueKey: splitTopLevelCommaSeparated(table.uniqueKey) } : {}), ...(table.partitionBy ? { partitionBy: table.partitionBy } : {}), ...(table.ttl ? { ttl: table.ttl } : {}), - ...(Object.keys(table.settings).length > 0 ? { settings: table.settings } : {}), + ...(Object.keys(table.settings).length > 0 ? { settings: isKafkaEngine(table.engine ?? '') ? parseKafkaSettings(table.settings) : table.settings } : {}), ...(table.indexes.length > 0 ? { indexes: table.indexes } : {}), ...(table.projections.length > 0 ? { projections: table.projections } : {}), } diff --git a/packages/plugin-pull/src/render-schema.ts b/packages/plugin-pull/src/render-schema.ts index e8130ff7..3afaa063 100644 --- a/packages/plugin-pull/src/render-schema.ts +++ b/packages/plugin-pull/src/render-schema.ts @@ -1,5 +1,6 @@ import { canonicalizeDefinitions, + isKafkaEngine, isIndexProjection, isRawCodec, type ColumnCodec, @@ -79,8 +80,10 @@ function renderTableDefinition(definition: TableDefinition, variableName: string lines.push(` ${renderColumn(column)},`) } lines.push(' ],') - lines.push(` primaryKey: ${renderStringArray(definition.primaryKey)},`) - lines.push(` orderBy: ${renderStringArray(definition.orderBy)},`) + if (!isKafkaEngine(definition.engine)) { + lines.push(` primaryKey: ${renderStringArray(definition.primaryKey)},`) + lines.push(` orderBy: ${renderStringArray(definition.orderBy)},`) + } if (definition.uniqueKey && definition.uniqueKey.length > 0) { lines.push(` uniqueKey: ${renderStringArray(definition.uniqueKey)},`) } diff --git a/test/kafka/.gitignore b/test/kafka/.gitignore new file mode 100644 index 00000000..c18dd8d8 --- /dev/null +++ b/test/kafka/.gitignore @@ -0,0 +1 @@ +__pycache__/ diff --git a/test/kafka/README.md b/test/kafka/README.md new file mode 100644 index 00000000..351742f3 --- /dev/null +++ b/test/kafka/README.md @@ -0,0 +1,35 @@ +# Kafka engine integration test + +This opt-in suite runs an isolated Kafka-compatible Redpanda broker and ClickHouse. +It exercises the real CLI and actual message consumption; it fails if the services +are unavailable. It is separate from the normal test suite because managed +ClickHouse test targets do not necessarily enable the Kafka engine. + +```sh +bun install --frozen-lockfile +bunx turbo run build --filter=chkit... --filter=@chkit/plugin-pull +docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml up -d --wait +bun test test/kafka/kafka.e2e.test.ts +python3 -m venv .venv +.venv/bin/pip install -e './chkit_python[dev]' +.venv/bin/python -m pytest test/kafka/test_python_e2e.py -q +docker compose -p chkit-issue203 -f test/kafka/docker-compose.yml down -v +``` + +The default is ClickHouse 25.3 on localhost port 18203. To test another version, +set `CLICKHOUSE_VERSION=26.3` before `docker compose up`. Use `KAFKA_TEST_HTTP_PORT` +for a different exposed port and `CLICKHOUSE_URL` for the matching test URL. +If changing the Compose project name, pass `KAFKA_TEST_PROJECT` to the test too. +Use a separate project name and port when running tests in parallel. +Run `docker compose down -v` for this project before switching ClickHouse versions +so an older server does not inherit a newer server's data files. + +The test creates unique topics/databases and cleans up only its own fixtures. +It verifies SQL migration execution, Kafka → MV → MergeTree ingestion, escaped +setting literals, pull/generate round trips, clean drift/check, settings drift, +rejection of unsupported changes without advancing snapshots, and an explicit +drop/create replacement with the existing destructive-operation gate. + +The Python test uses its native CLI and schema DSL, including `drift --live` and +`check --live`. It also checks that offline drift/check report unsupported Kafka +changes without crashing. CI runs both workflows on ClickHouse 25.3 and 26.3. diff --git a/test/kafka/clickhouse.xml b/test/kafka/clickhouse.xml new file mode 100644 index 00000000..988f634e --- /dev/null +++ b/test/kafka/clickhouse.xml @@ -0,0 +1,5 @@ + + + 18203 + earliest + diff --git a/test/kafka/docker-compose.yml b/test/kafka/docker-compose.yml new file mode 100644 index 00000000..9a5201b1 --- /dev/null +++ b/test/kafka/docker-compose.yml @@ -0,0 +1,33 @@ +services: + kafka: + image: docker.redpanda.com/redpandadata/redpanda:v25.3.1 + command: + - redpanda + - start + - --mode=dev-container + - --smp=1 + - --memory=512M + - --kafka-addr=0.0.0.0:9092 + - --advertise-kafka-addr=kafka:9092 + healthcheck: + test: [CMD, rpk, cluster, info] + interval: 2s + timeout: 5s + retries: 40 + clickhouse: + image: clickhouse/clickhouse-server:${CLICKHOUSE_VERSION:-25.3} + environment: + CLICKHOUSE_PASSWORD: chkit-kafka-test + ports: + - "127.0.0.1:${KAFKA_TEST_HTTP_PORT:-18203}:18203" + volumes: + - ./clickhouse.xml:/etc/clickhouse-server/config.d/chkit-kafka.xml:ro + - ../ci/clickhouse.xml:/etc/clickhouse-server/config.d/chkit-ci.xml:ro + depends_on: + kafka: + condition: service_healthy + healthcheck: + test: [CMD, clickhouse-client, --password, chkit-kafka-test, --query, SELECT 1] + interval: 2s + timeout: 5s + retries: 40 diff --git a/test/kafka/kafka.e2e.test.ts b/test/kafka/kafka.e2e.test.ts new file mode 100644 index 00000000..98002d4a --- /dev/null +++ b/test/kafka/kafka.e2e.test.ts @@ -0,0 +1,141 @@ +import { expect, test } from 'bun:test' +import { mkdtemp, readFile, readdir, rm, symlink, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join, resolve } from 'node:path' +import { setTimeout as sleep } from 'node:timers/promises' +import { createClickHouseExecutor } from '../../packages/clickhouse/src/index.js' +import { runCli } from '../../packages/cli/src/test/e2e-testkit.js' + +const root = resolve(import.meta.dir, '../..') +const compose = [ + 'docker', + 'compose', + '-p', + process.env.KAFKA_TEST_PROJECT ?? 'chkit-issue203', + '-f', + join(import.meta.dir, 'docker-compose.yml'), +] +const url = process.env.CLICKHOUSE_URL ?? 'http://127.0.0.1:18203' +const password = 'chkit-kafka-test' + +async function broker(args: string[], stdin?: string) { + const proc = Bun.spawn([...compose, 'exec', '-T', 'kafka', 'rpk', ...args], { + stdin: stdin === undefined ? 'ignore' : new Blob([stdin]), + stdout: 'pipe', + stderr: 'pipe', + }) + const [exitCode, stderr] = await Promise.all([proc.exited, new Response(proc.stderr).text()]) + expect(exitCode, stderr).toBe(0) +} + +test('Kafka → MV → MergeTree: generate, migrate, consume, pull, drift, check and explicit replacement', async () => { + const tag = `${Date.now()}_${Math.floor(Math.random() * 100000)}` + const database = `chkit_kafka_${tag}` + const topic = `chkit_${tag}` + const journal = `_chkit_kafka_${tag}` + const dir = await mkdtemp(join(tmpdir(), 'chkit-kafka-')) + const db = createClickHouseExecutor({ url, password, username: 'default', database: 'default' }) + const clientId = "client; COMMENT 'quoted' \\" + const source = ( + includeQueue: boolean, + consumers = 1, + ) => `import { schema, table, materializedView } from '@chkit/core' +const storage = table({ database: '${database}', name: 'events', engine: 'MergeTree', columns: [{ name: 'id', type: 'UInt64' }, { name: 'body', type: 'String' }], primaryKey: ['id'], orderBy: ['id'], settings: { index_granularity: '8192' } }) +${ + includeQueue + ? `const queue = table({ database: '${database}', name: 'queue', engine: 'Kafka', columns: storage.columns, settings: { kafka_broker_list: 'kafka:9092', kafka_topic_list: '${topic}', kafka_group_name: '${topic}', kafka_format: 'JSONEachRow', kafka_client_id: ${JSON.stringify(clientId)}, kafka_num_consumers: ${consumers}, kafka_flush_interval_ms: 100, kafka_commit_on_select: false, input_format_skip_unknown_fields: true } }) +const mv = materializedView({ database: '${database}', name: 'consumer', to: { database: '${database}', name: 'events' }, as: 'SELECT id, body FROM ${database}.queue' })` + : '' +} +export default schema(storage${includeQueue ? ', queue, mv' : ''}) +` + const cli = (args: string[], success = true) => { + const result = runCli(dir, [...args, '--config', join(dir, 'clickhouse.config.ts'), '--json'], { + CHKIT_JOURNAL_TABLE: journal, + }) + if (success) expect(result.exitCode, result.stdout + result.stderr).toBe(0) + return result + } + const waitForIds = async (ids: number[]) => { + const deadline = Date.now() + 30000 + let actual: number[] = [] + do { + actual = ( + await db.query<{ id: string }>(`SELECT id FROM ${database}.events ORDER BY id`) + ).map((row) => Number(row.id)) + if (JSON.stringify(actual) === JSON.stringify(ids)) return + await sleep(200) + } while (Date.now() < deadline) + expect(actual).toEqual(ids) + } + try { + await symlink(join(root, 'node_modules'), join(dir, 'node_modules')) + await writeFile(join(dir, 'schema.ts'), source(true)) + await writeFile( + join(dir, 'clickhouse.config.ts'), + `import { pull } from '@chkit/plugin-pull' +export default { schema: './schema.ts', outDir: './chkit', plugins: [pull()], clickhouse: { url: '${url}', username: 'default', password: '${password}', database: 'default' } }`, + ) + await broker(['topic', 'create', topic, '-p', '2']) + cli(['generate', '--name', 'create', '--migration-id', '001']) + cli(['migrate', '--apply']) + await broker( + ['topic', 'produce', topic], + '{"id":1,"body":"hello","ignored":true}\n{"id":2,"body":"world"}\n', + ) + await waitForIds([1, 2]) + expect(JSON.parse(cli(['drift']).stdout).drifted).toBe(false) + cli(['check']) + + // Round-trip through the actual pull command and rerun generation. + const pulled = join(dir, 'pulled.ts') + cli(['pull', 'schema', '--database', database, '--out-file', pulled]) + const content = await readFile(pulled, 'utf8') + const queueBlock = content.slice(content.indexOf('name: "queue"')) + expect(queueBlock).not.toContain('primaryKey:') + expect(content).toContain(JSON.stringify(clientId)) + await writeFile(join(dir, 'schema.ts'), content) + const roundTrip = cli(['generate', '--dryrun']) + expect(JSON.parse(roundTrip.stdout).operationCount, roundTrip.stdout).toBe(0) + + // Settings drift is observable without trying unsupported Kafka ALTERs. + const snapshotPath = join(dir, 'chkit/meta/snapshot.json') + const snapshotText = await readFile(snapshotPath, 'utf8') + const snapshot = JSON.parse(snapshotText) + snapshot.definitions.find( + (def: { name: string }) => def.name === 'queue', + ).settings.kafka_group_name = 'different' + await writeFile(snapshotPath, JSON.stringify(snapshot)) + expect(JSON.parse(cli(['drift'], false).stdout).tableDrift[0].settingDiffs).toContain( + 'kafka_group_name', + ) + expect(cli(['check'], false).exitCode).not.toBe(0) + await writeFile(snapshotPath, snapshotText) + + // A refused change must not update either snapshot or migration files. + const beforeFiles = await readdir(join(dir, 'chkit/migrations')) + await writeFile(join(dir, 'schema.ts'), source(true, 2)) + const blocked = cli(['generate', '--name', 'unsafe'], false) + expect(blocked.exitCode).not.toBe(0) + expect(blocked.stdout).toContain('kafka_change_requires_replacement') + expect(await readFile(snapshotPath, 'utf8')).toBe(snapshotText) + expect(await readdir(join(dir, 'chkit/migrations'))).toEqual(beforeFiles) + + // Explicit two-migration replacement keeps the storage table and snapshot. + await writeFile(join(dir, 'schema.ts'), source(false)) + cli(['generate', '--name', 'stop-queue', '--migration-id', '002']) + await writeFile(join(dir, 'schema.ts'), source(true, 2)) + cli(['generate', '--name', 'restart-queue', '--migration-id', '003']) + expect(cli(['migrate', '--apply'], false).exitCode).not.toBe(0) + cli(['migrate', '--apply', '--allow-destructive']) + await broker(['topic', 'produce', topic], '{"id":3,"body":"after replacement"}\n') + await waitForIds([1, 2, 3]) + cli(['check']) + } finally { + await db.command(`DROP DATABASE IF EXISTS ${database} SYNC`) + await db.command(`DROP TABLE IF EXISTS default.${journal} SYNC`) + await db.close() + await broker(['topic', 'delete', topic]) + await rm(dir, { recursive: true, force: true }) + } +}, 120000) diff --git a/test/kafka/test_python_e2e.py b/test/kafka/test_python_e2e.py new file mode 100644 index 00000000..ea47fb16 --- /dev/null +++ b/test/kafka/test_python_e2e.py @@ -0,0 +1,161 @@ +"""Opt-in native Python CLI Kafka workflow against test/kafka/docker-compose.yml.""" + +from __future__ import annotations + +import json +import os +import subprocess +import sys +import time +from pathlib import Path +from urllib.parse import urlparse + +import clickhouse_connect + + +def test_python_kafka_pipeline(tmp_path: Path) -> None: + tag = f"{time.time_ns()}" + database, topic, journal = ( + f"chkit_py_kafka_{tag}", + f"chkit_py_{tag}", + f"_chkit_py_{tag}", + ) + url = os.environ.get("CLICKHOUSE_URL", "http://127.0.0.1:18203") + address = urlparse(url) + password = "chkit-kafka-test" + db = clickhouse_connect.get_client( + host=address.hostname, port=address.port, username="default", password=password + ) + compose = [ + "docker", + "compose", + "-p", + os.environ.get("KAFKA_TEST_PROJECT", "chkit-issue203"), + "-f", + str(Path(__file__).with_name("docker-compose.yml")), + ] + + def broker(args: list[str], stdin: str | None = None) -> None: + subprocess.run( + [*compose, "exec", "-T", "kafka", "rpk", *args], + input=stdin, + text=True, + capture_output=True, + check=True, + ) + + def source(include_queue: bool, consumers: int = 1) -> str: + return f"""from chkit import schema, table, materialized_view +storage = table(database={database!r}, name="events", engine="MergeTree", + columns=[{{"name": "id", "type": "UInt64"}}, {{"name": "body", "type": "String"}}], + primary_key=["id"], order_by=["id"], settings={{"index_granularity": 8192}}) +""" + ( + f""" +queue = table(database={database!r}, name="queue", engine="Kafka", columns=list(storage.columns), + settings={{"kafka_broker_list": "kafka:9092", "kafka_topic_list": {topic!r}, + "kafka_group_name": {topic!r}, "kafka_format": "JSONEachRow", + "kafka_client_id": {"client; COMMENT 'quoted' " + chr(92)!r}, + "kafka_num_consumers": {consumers}, "kafka_flush_interval_ms": 100, + "kafka_commit_on_select": False, "input_format_skip_unknown_fields": True}}) +mv = materialized_view(database={database!r}, name="consumer", + to={{"database": {database!r}, "name": "events"}}, as_="SELECT id, body FROM {database}.queue") +definitions = schema(storage, queue, mv) +""" + if include_queue + else "definitions = schema(storage)\n" + ) + + def cli(args: list[str], success: bool = True) -> subprocess.CompletedProcess[str]: + result = subprocess.run( + [ + str(Path(sys.executable).with_name("chkit")), + *args, + "--config", + str(tmp_path / "clickhouse.config.py"), + "--json", + ], + cwd=tmp_path, + env={**os.environ, "CHKIT_JOURNAL_TABLE": journal}, + text=True, + capture_output=True, + check=False, + ) + if success: + assert result.returncode == 0, result.stdout + result.stderr + return result + + def wait_for_ids(ids: list[int]) -> None: + deadline = time.monotonic() + 30 + while time.monotonic() < deadline: + actual = [ + row[0] + for row in db.query( + f"SELECT id FROM {database}.events ORDER BY id" + ).result_rows + ] + if actual == ids: + return + time.sleep(0.2) + assert actual == ids + + schema_file = tmp_path / "schema.py" + try: + schema_file.write_text(source(True)) + (tmp_path / "clickhouse.config.py").write_text( + f"config = {{'schema': './schema.py', 'outDir': './chkit', 'clickhouse': " + f"{{'url': {url!r}, 'username': 'default', 'password': {password!r}, 'database': 'default'}}}}" + ) + broker(["topic", "create", topic, "-p", "2"]) + cli(["generate", "--name", "create", "--migration-id", "001"]) + cli(["migrate", "--apply"]) + broker( + ["topic", "produce", topic], + '{"id":1,"body":"hello","ignored":true}\n{"id":2,"body":"world"}\n', + ) + wait_for_ids([1, 2]) + assert json.loads(cli(["drift", "--live"]).stdout)["drifted"] is False + cli(["check", "--live"]) + + pulled = tmp_path / "pulled.py" + cli(["pull", "--database", database, "--out-file", str(pulled)]) + schema_file.write_text(pulled.read_text()) + plan = cli(["generate", "--dryrun"]) + assert json.loads(plan.stdout)["operationCount"] == 0, plan.stdout + + snapshot_file = tmp_path / "chkit/meta/snapshot.json" + before = snapshot_file.read_text() + snapshot = json.loads(before) + next(d for d in snapshot["definitions"] if d["name"] == "queue")["settings"][ + "kafka_group_name" + ] = "other" + snapshot_file.write_text(json.dumps(snapshot)) + assert "kafka_group_name" in cli(["drift", "--live"]).stdout + assert "kafka_change_requires_replacement" in cli(["drift"]).stdout + for flags in ([], ["--live"]): + result = cli(["check", *flags], success=False) + assert result.returncode != 0 + assert "kafka_change_requires_replacement" in result.stdout + snapshot_file.write_text(before) + + schema_file.write_text(source(True, 2)) + migrations = sorted((tmp_path / "chkit/migrations").iterdir()) + blocked = cli(["generate", "--name", "unsafe"], success=False) + assert blocked.returncode != 0 + assert "kafka_change_requires_replacement" in blocked.stdout + blocked.stderr + assert snapshot_file.read_text() == before + assert sorted((tmp_path / "chkit/migrations").iterdir()) == migrations + + schema_file.write_text(source(False)) + cli(["generate", "--name", "stop", "--migration-id", "002"]) + schema_file.write_text(source(True, 2)) + cli(["generate", "--name", "restart", "--migration-id", "003"]) + assert cli(["migrate", "--apply"], success=False).returncode != 0 + cli(["migrate", "--apply", "--allow-destructive"]) + broker(["topic", "produce", topic], '{"id":3,"body":"after replacement"}\n') + wait_for_ids([1, 2, 3]) + cli(["check", "--live"]) + finally: + db.command(f"DROP DATABASE IF EXISTS {database} SYNC") + db.command(f"DROP TABLE IF EXISTS default.{journal} SYNC") + db.close() + broker(["topic", "delete", topic])