diff --git a/.changeset/text-index-round-trips.md b/.changeset/text-index-round-trips.md new file mode 100644 index 00000000..6af47325 --- /dev/null +++ b/.changeset/text-index-round-trips.md @@ -0,0 +1,14 @@ +--- +"@chkit/core": minor +"@chkit/clickhouse": patch +"@chkit/plugin-pull": minor +"chkit": minor +--- + +Support ClickHouse `text` indexes in schemas, migrations, pull, and drift. Text +indexes require a tokenizer and support preprocessing, postprocessing, phrase +search, and dictionary/posting-list options when supported by the server. +Granularity is automatic. Preserve whitespace and escapes inside SQL literals, +compare parameter order and SQL formatting consistently, and reject unsupported +or malformed metadata instead of silently losing settings. Includes Python parity +and live adversarial round-trip tests on ClickHouse 26.3 and 26.8. diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index eaa6541a..dca86cbb 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -77,6 +77,48 @@ jobs: - name: Check packed tarball deps run: bun run check:packed-deps + text-index: + name: Text indexes (ClickHouse ${{ matrix.clickhouse }}) + runs-on: blacksmith-4vcpu-ubuntu-2404 + strategy: + fail-fast: false + matrix: + clickhouse: ['26.3', '26.8'] + services: + clickhouse: + image: clickhouse/clickhouse-server:${{ matrix.clickhouse }} + env: + CLICKHOUSE_PASSWORD: chkit-ci + CLICKHOUSE_DEFAULT_ACCESS_MANAGEMENT: 1 + ports: + - 8123:8123 + options: >- + --ulimit nofile=262144:262144 + --health-cmd "clickhouse-client --password chkit-ci --query 'SELECT 1'" + --health-interval 5s + --health-timeout 5s + --health-retries 20 + env: + CLICKHOUSE_URL: http://127.0.0.1:8123 + CLICKHOUSE_USER: default + CLICKHOUSE_PASSWORD: chkit-ci + CLICKHOUSE_DB: default + steps: + - uses: actions/checkout@v4 + - uses: ./.github/actions/setup + - uses: actions/setup-python@v5 + with: + python-version: '3.12' + - name: Build test dependencies + run: bunx turbo run build --filter=chkit... --filter=@chkit/plugin-pull... + - name: Test TypeScript text indexes + run: bun test packages/core/src/text-index.test.ts packages/cli/src/test/text-index.e2e.test.ts + - name: Install Python test dependencies + run: python -m pip install -e './chkit_python[dev]' + - name: Test Python text indexes + working-directory: chkit_python + run: python -m pytest tests/test_text_index.py tests/test_text_index_e2e.py -q + obsessiondb: runs-on: blacksmith-8vcpu-ubuntu-2404 needs: verify @@ -116,7 +158,7 @@ jobs: deploy: runs-on: blacksmith-4vcpu-ubuntu-2404 - needs: [verify, obsessiondb] + needs: [verify, obsessiondb, text-index] if: github.event_name == 'push' && github.ref == 'refs/heads/main' steps: - name: Checkout diff --git a/apps/docs/src/content/docs/schema/dsl-reference.mdx b/apps/docs/src/content/docs/schema/dsl-reference.mdx index 8236164c..965eb426 100644 --- a/apps/docs/src/content/docs/schema/dsl-reference.mdx +++ b/apps/docs/src/content/docs/schema/dsl-reference.mdx @@ -358,8 +358,8 @@ Each entry in the `indexes` array is a `SkipIndexDefinition`. The shared base fi |-------|------|-------------| | `name` | `string` | Index name | | `expression` | `string` | Indexed expression | -| `type` | `'minmax' \| 'set' \| 'bloom_filter' \| 'tokenbf_v1' \| 'ngrambf_v1'` | Index type | -| `granularity` | `number` | Index granularity | +| `type` | `'minmax' \| 'set' \| 'bloom_filter' \| 'tokenbf_v1' \| 'ngrambf_v1' \| 'text'` | Index type | +| `granularity` | `number` | Required for other indexes; optional and ignored for `text`, which always uses `100000000` | Type-specific fields: @@ -371,6 +371,58 @@ Type-specific fields: | `tokenbf_v1` | `sizeBytes`, `hashFunctions`, `randomSeed` (all `number`) | — | Maps to `tokenbf_v1(size_bytes, n_hash, seed)` | | `ngrambf_v1` | `ngramSize`, `sizeBytes`, `hashFunctions`, `randomSeed` (all `number`) | — | Maps to `ngrambf_v1(n, size_bytes, n_hash, seed)` | +### Full-text indexes + +Use `type: 'text'` on ClickHouse 26.2 or newer. `tokenizer` is a required SQL +expression, such as `splitByNonAlpha`, `ngrams(3)`, or `splitByString([' ', ';'])`. +Quoted whitespace, Unicode, and escaped characters retain their meaning through +generation, pull, and drift checks. Granularity is automatic: ClickHouse indexes +an entire part and ignores any supplied granularity. + +| Field | Type | Meaning | +|-------|------|---------| +| `tokenizer` | `string` | Required SQL tokenizer | +| `preprocessor` | `string` | Optional SQL expression applied before tokenization | +| `postprocessor` | `string` | Optional SQL expression applied to each token; requires server support | +| `supportPhraseSearch` | `boolean` | Store token positions; requires server support and the table setting `allow_experimental_text_index_phrase_search: 1` | +| `dictionaryBlockSize` | `number` | Positive integer dictionary block size | +| `dictionaryBlockFrontcodingCompression` | `boolean` | Enable or disable dictionary front coding | +| `postingListBlockSize` | `number` | Positive integer posting-list block size | +| `postingListCodec` | `'none' \| 'bitpacking'` | Posting-list compression | + +```ts +indexes: [{ + name: 'idx_body', + expression: 'body', + type: 'text', + tokenizer: "splitByString([' ', ';'])", + preprocessor: 'lower(body)', +}] +``` + +Python accepts the same dictionary fields, or +`SkipIndexText(name="idx_body", expression="body", tokenizer="splitByNonAlpha")`. +Snake-case names such as `posting_list_codec` are also accepted. + +Basic text indexes and tuning options are tested against ClickHouse 26.3 and 26.8. +The newer `postprocessor` and `supportPhraseSearch` options are exercised on 26.8; +26.2 availability of the index does not imply availability of every later option. +ClickHouse remains responsible for validating tokenizer/function availability and +server-specific parameter limits. Pull fails with an explicit error for unknown +text-index parameters instead of silently discarding them. + +Avoid column names that are also SQL literals (`true`, `false`, `inf`, `infinity`, +or `nan`) in text-index expressions on older servers. ClickHouse 26.3 can remove +their required identifier quotes from index metadata, preventing a lossless pull. +chkit treats meaningful quote differences as drift; it does not assume a column +reference and a literal are equivalent. Ordinary identifier quoting and switching +between backticks and double quotes do not require an index rebuild. + +Adding or changing an index does not automatically index historical parts. Run +`ALTER TABLE database.table MATERIALIZE INDEX idx_body` when historical data must +be indexed; materialization consumes database resources. Existing rows remain +queryable before materialization. + ```ts @@ -405,7 +457,7 @@ Type-specific fields: }, ] ``` - Model classes are importable when dicts feel too loose: `SkipIndexMinmax`, `SkipIndexSet`, `SkipIndexBloomFilter`, `SkipIndexTokenBF`, `SkipIndexNgramBF`. + Model classes are importable when dicts feel too loose: `SkipIndexMinmax`, `SkipIndexSet`, `SkipIndexBloomFilter`, `SkipIndexTokenBF`, `SkipIndexNgramBF`, `SkipIndexText`. diff --git a/chkit_python/CHANGELOG.md b/chkit_python/CHANGELOG.md index 81cbe519..ab230e06 100644 --- a/chkit_python/CHANGELOG.md +++ b/chkit_python/CHANGELOG.md @@ -1,5 +1,12 @@ # Changelog +## Unreleased + +- 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. + ## 0.2.0 — 2026-08-10 **Full parity with the TypeScript chkit.** Every remaining gap is closed; diff --git a/chkit_python/pyproject.toml b/chkit_python/pyproject.toml index 89f2953f..0c46b11c 100644 --- a/chkit_python/pyproject.toml +++ b/chkit_python/pyproject.toml @@ -124,6 +124,9 @@ ignore = [ # dispatch and per-operation-type branches are clearer as flat code than as # data-driven dispatch tables. PERF401 prefers comprehensions over append loops # but the imperative style is explicit and intentional here. +# Text-index lexical scanning uses byte values and explicit parser state. +"src/chkit/core/text_index_sql.py" = ["PLR0912", "PLR0915", "PLR2004"] +"src/chkit/core/text_index.py" = ["PLR0912", "PLR2004"] "src/chkit/core/codec.py" = ["PLR0911", "PLR0912", "PLR2004", "E501"] "src/chkit/core/planner.py" = ["PLR0911", "PLR0912", "PERF401"] "src/chkit/core/model.py" = ["E501", "PLC0415"] diff --git a/chkit_python/src/chkit/__init__.py b/chkit_python/src/chkit/__init__.py index c4c43c21..fe1fe444 100644 --- a/chkit_python/src/chkit/__init__.py +++ b/chkit_python/src/chkit/__init__.py @@ -47,6 +47,7 @@ SkipIndexMinmax, SkipIndexNgramBF, SkipIndexSet, + SkipIndexText, SkipIndexTokenBF, ) @@ -75,6 +76,7 @@ "SkipIndexMinmax", "SkipIndexNgramBF", "SkipIndexSet", + "SkipIndexText", "SkipIndexTokenBF", "TableDefinition", "TableRef", diff --git a/chkit_python/src/chkit/cli/commands/drift_compare.py b/chkit_python/src/chkit/cli/commands/drift_compare.py index 2eacc443..99ee29d7 100644 --- a/chkit_python/src/chkit/cli/commands/drift_compare.py +++ b/chkit_python/src/chkit/cli/commands/drift_compare.py @@ -29,6 +29,7 @@ ) from chkit.core.projection import is_index_projection, normalize_projection_index from chkit.core.sql_normalizer import normalize_engine, normalize_sql_fragment +from chkit.core.text_index import render_text_index_type, text_index_fingerprint _MIN_QUOTED_LEN = 2 @@ -244,6 +245,8 @@ def _normalize_default_value(value: str) -> str: def _render_index_type_fingerprint(index: SkipIndexDefinition) -> str: + if index.type == "text": + return render_text_index_type(index) if index.type == "minmax": return "minmax" if index.type == "set": @@ -285,6 +288,8 @@ def _strip_enclosing_parens(value: str) -> str: def _normalize_index_shape(index: SkipIndexDefinition) -> str: + if index.type == "text": + return text_index_fingerprint(index) return "|".join( [ f"expr={_strip_enclosing_parens(normalize_sql_fragment(index.expression))}", diff --git a/chkit_python/src/chkit/cli/commands/pull.py b/chkit_python/src/chkit/cli/commands/pull.py index 0e3ccf1b..c33d5deb 100644 --- a/chkit_python/src/chkit/cli/commands/pull.py +++ b/chkit_python/src/chkit/cli/commands/pull.py @@ -65,6 +65,7 @@ SkipIndexMinmax, SkipIndexNgramBF, SkipIndexSet, + SkipIndexText, SkipIndexTokenBF, TableDefinition, TableRef, @@ -92,6 +93,7 @@ def _introspected_table_to_definition( | SkipIndexBloomFilter | SkipIndexTokenBF | SkipIndexNgramBF + | SkipIndexText | dict[str, object] ] = list(item.indexes) diff --git a/chkit_python/src/chkit/cli/commands/pull_render.py b/chkit_python/src/chkit/cli/commands/pull_render.py index 3df505e2..857b8108 100644 --- a/chkit_python/src/chkit/cli/commands/pull_render.py +++ b/chkit_python/src/chkit/cli/commands/pull_render.py @@ -99,6 +99,7 @@ def render_schema_file( # noqa: PLR0912, PLR0915 "bloom_filter": "SkipIndexBloomFilter", "tokenbf_v1": "SkipIndexTokenBF", "ngrambf_v1": "SkipIndexNgramBF", + "text": "SkipIndexText", }[idx.type] ) @@ -245,13 +246,23 @@ def _render_index(index: SkipIndexDefinition) -> str: parts.append(f"size_bytes={index.size_bytes}") parts.append(f"hash_functions={index.hash_functions}") parts.append(f"random_seed={index.random_seed}") - parts.append(f"granularity={index.granularity}") + elif index.type == "text": + for field in ("tokenizer", "preprocessor", "postprocessor", "support_phrase_search", + "dictionary_block_size", "dictionary_block_frontcoding_compression", + "posting_list_block_size", "posting_list_codec"): + value = getattr(index, field) + if value is not None: + rendered = _render_string(value) if isinstance(value, str) else str(value) + parts.append(f"{field}={rendered}") + if index.type != "text": + parts.append(f"granularity={index.granularity}") type_class = { "minmax": "SkipIndexMinmax", "set": "SkipIndexSet", "bloom_filter": "SkipIndexBloomFilter", "tokenbf_v1": "SkipIndexTokenBF", "ngrambf_v1": "SkipIndexNgramBF", + "text": "SkipIndexText", }[index.type] return f"{type_class}({', '.join(parts)})" diff --git a/chkit_python/src/chkit/clickhouse/introspect.py b/chkit_python/src/chkit/clickhouse/introspect.py index 1d05d8aa..0e1e776a 100644 --- a/chkit_python/src/chkit/clickhouse/introspect.py +++ b/chkit_python/src/chkit/clickhouse/introspect.py @@ -45,9 +45,12 @@ SkipIndexMinmax, SkipIndexNgramBF, SkipIndexSet, + SkipIndexText, SkipIndexTokenBF, ) from chkit.core.sql_normalizer import normalize_sql_fragment +from chkit.core.text_index import parse_text_index_params +from chkit.core.text_index_sql import normalize_text_index_sql SchemaObjectKind: TypeAlias = Literal["table", "view", "materialized_view", "dictionary"] @@ -122,7 +125,7 @@ class IntrospectedTable: _NULLABLE_RE = re.compile(r"^Nullable\((.+)\)$") -_INDEX_TYPE_RE = re.compile(r"^(\w+)\((.+)\)$") +_INDEX_TYPE_RE = re.compile(r"^(\w+)\((.*)\)$", re.DOTALL) def infer_schema_kind_from_engine(engine: str) -> SchemaObjectKind | None: @@ -212,6 +215,12 @@ def normalize_index_from_system_row(row: SystemSkippingIndexRow) -> SkipIndexDef base_name = match.group(1) if match is not None else row.type args_str = match.group(2) if match is not None else None + if base_name == "text": + base_payload["expression"] = normalize_text_index_sql(row.expr) + return SkipIndexText.model_validate( + {**base_payload, **parse_text_index_params(args_str or "")} + ) + if base_name == "minmax": return SkipIndexMinmax(**base_payload) diff --git a/chkit_python/src/chkit/core/canonical.py b/chkit_python/src/chkit/core/canonical.py index 6cab2693..2fe9bf9d 100644 --- a/chkit_python/src/chkit/core/canonical.py +++ b/chkit_python/src/chkit/core/canonical.py @@ -26,6 +26,7 @@ ) from chkit.core.projection import canonicalize_projection from chkit.core.sql_normalizer import normalize_engine, normalize_sql_fragment +from chkit.core.text_index import canonicalize_text_index _T = TypeVar("_T") @@ -63,6 +64,8 @@ def _canonicalize_column(column: ColumnDefinition) -> ColumnDefinition: def _canonicalize_index(index: SkipIndexDefinition) -> SkipIndexDefinition: + if index.type == "text": + return canonicalize_text_index(index) return index.model_copy(update={"expression": normalize_sql_fragment(index.expression)}) diff --git a/chkit_python/src/chkit/core/model.py b/chkit_python/src/chkit/core/model.py index 9214324e..981272d5 100644 --- a/chkit_python/src/chkit/core/model.py +++ b/chkit_python/src/chkit/core/model.py @@ -209,12 +209,33 @@ class SkipIndexNgramBF(_SkipIndexBase): ) +class SkipIndexText(_SkipIndexBase): + """Full-text index (26.2+); ClickHouse fixes granularity per part.""" + + type: Literal["text"] = "text" + granularity: int = 100_000_000 + tokenizer: str + preprocessor: str | None = None + postprocessor: str | None = None + support_phrase_search: bool | None = Field(default=None, alias="supportPhraseSearch") + dictionary_block_size: int | None = Field(default=None, alias="dictionaryBlockSize") + dictionary_block_frontcoding_compression: bool | None = Field( + default=None, alias="dictionaryBlockFrontcodingCompression" + ) + posting_list_block_size: int | None = Field(default=None, alias="postingListBlockSize") + posting_list_codec: Literal["none", "bitpacking"] | None = Field(default=None, alias="postingListCodec") + + model_config = ConfigDict(frozen=True, extra="forbid", strict=True, + validate_assignment=True, populate_by_name=True) + + SkipIndexDefinition: TypeAlias = Annotated[ SkipIndexMinmax | SkipIndexSet | SkipIndexBloomFilter | SkipIndexTokenBF - | SkipIndexNgramBF, + | SkipIndexNgramBF + | SkipIndexText, Field(discriminator="type"), ] @@ -652,6 +673,7 @@ class MigrationPlan(_StrictModel): "duplicate_object_name", "duplicate_column_name", "duplicate_index_name", + "text_index_invalid_parameters", "duplicate_projection_name", "projection_ambiguous_kind", "projection_empty_index", diff --git a/chkit_python/src/chkit/core/planner.py b/chkit_python/src/chkit/core/planner.py index 532df8af..78f4dc0c 100644 --- a/chkit_python/src/chkit/core/planner.py +++ b/chkit_python/src/chkit/core/planner.py @@ -38,6 +38,7 @@ render_dictionary_sql, to_create_sql, ) +from chkit.core.text_index import text_index_fingerprint from chkit.core.validate import assert_valid_definitions @@ -232,6 +233,8 @@ def _index_identity(index: SkipIndexDefinition) -> str: def _indexes_equal(left: SkipIndexDefinition, right: SkipIndexDefinition) -> bool: + if left.type == "text" and right.type == "text": + return text_index_fingerprint(left) == text_index_fingerprint(right) return _index_identity(left) == _index_identity(right) diff --git a/chkit_python/src/chkit/core/sql.py b/chkit_python/src/chkit/core/sql.py index 1f739e1d..b7accbdf 100644 --- a/chkit_python/src/chkit/core/sql.py +++ b/chkit_python/src/chkit/core/sql.py @@ -27,6 +27,7 @@ ViewDefinition, ) from chkit.core.projection import render_projection_body +from chkit.core.text_index import TEXT_INDEX_GRANULARITY, render_text_index_type from chkit.core.validate import assert_valid_definitions _COLUMN_ADAPTER: TypeAdapter[ColumnDefinition] = TypeAdapter(ColumnDefinition) @@ -91,14 +92,15 @@ def _render_key_clause_columns(columns: list[str], column_names: set[str]) -> st def _render_index_type(idx: SkipIndexDefinition) -> str: + if idx.type == "text": + return render_text_index_type(idx) if isinstance(idx, SkipIndexMinmax): return "minmax" if isinstance(idx, SkipIndexSet): return f"set({idx.max_rows})" if isinstance(idx, SkipIndexBloomFilter): - if idx.false_positive_rate is not None: - return f"bloom_filter({idx.false_positive_rate})" - return "bloom_filter" + rate = idx.false_positive_rate + return "bloom_filter" if rate is None else f"bloom_filter({rate})" if isinstance(idx, SkipIndexTokenBF): return f"tokenbf_v1({idx.size_bytes}, {idx.hash_functions}, {idx.random_seed})" # SkipIndexNgramBF is the only remaining variant in the discriminated union. @@ -126,7 +128,8 @@ def _render_projection(p: ProjectionDefinition) -> str: def _render_index_line(idx: SkipIndexDefinition) -> str: return ( f"INDEX `{idx.name}` ({idx.expression}) " - f"TYPE {_render_index_type(idx)} GRANULARITY {idx.granularity}" + f"TYPE {_render_index_type(idx)} GRANULARITY " + f"{TEXT_INDEX_GRANULARITY if idx.type == 'text' else idx.granularity}" ) @@ -344,7 +347,8 @@ def render_alter_add_index(definition: TableDefinition, index: IndexInput) -> st return ( f"ALTER TABLE {definition.database}.{definition.name} " f"ADD INDEX IF NOT EXISTS `{normalized.name}` ({normalized.expression}) " - f"TYPE {_render_index_type(normalized)} GRANULARITY {normalized.granularity};" + f"TYPE {_render_index_type(normalized)} GRANULARITY " + f"{TEXT_INDEX_GRANULARITY if normalized.type == 'text' else normalized.granularity};" ) diff --git a/chkit_python/src/chkit/core/text_index.py b/chkit_python/src/chkit/core/text_index.py new file mode 100644 index 00000000..05f4810d --- /dev/null +++ b/chkit_python/src/chkit/core/text_index.py @@ -0,0 +1,126 @@ +"""Text index normalization and SQL round trips, in parity with TypeScript.""" + +from __future__ import annotations + +import json +import re +from typing import TYPE_CHECKING + +from chkit.core.text_index_sql import ( + format_text_sql, + normalize_text_index_sql, + text_expression_fingerprint, + text_sql_fingerprint, + text_sql_tokens, +) + +if TYPE_CHECKING: + from chkit.core.model import SkipIndexText + +TEXT_INDEX_GRANULARITY = 100_000_000 +PARAMETERS = ( + ("tokenizer", "tokenizer", "sql"), + ("preprocessor", "preprocessor", "sql"), + ("postprocessor", "postprocessor", "sql"), + ("support_phrase_search", "support_phrase_search", "boolean"), + ("dictionary_block_size", "dictionary_block_size", "number"), + ( + "dictionary_block_frontcoding_compression", + "dictionary_block_frontcoding_compression", + "boolean", + ), + ("posting_list_block_size", "posting_list_block_size", "number"), + ("posting_list_codec", "posting_list_codec", "codec"), +) + + +def render_text_index_type(index: SkipIndexText) -> str: + if not text_sql_tokens(index.tokenizer): + raise ValueError("A non-empty tokenizer is required") + parts = [] + for field, key, kind in PARAMETERS: + value = getattr(index, field) + if value is None: + continue + if kind == "sql": + if not isinstance(value, str) or not text_sql_tokens(value): + raise ValueError(f"{field} must be non-empty SQL") + sql = normalize_text_index_sql(value) + elif kind == "boolean": + if not isinstance(value, bool): + raise ValueError(f"{field} must be a boolean") + sql = "1" if value else "0" + elif kind == "number": + if type(value) is not int or not 0 < value <= 2**53 - 1: + raise ValueError(f"{field} must be a positive safe integer") + sql = str(value) + else: + if value not in ("none", "bitpacking"): + raise ValueError(f"{field} must be none or bitpacking") + sql = f"'{value}'" + parts.append(f"{key} = {sql}") + return f"text({', '.join(parts)})" + + +def parse_text_index_params(args: str) -> dict[str, str | int | bool]: + tokens = text_sql_tokens(args) + groups: list[list[str]] = [[]] + depth = 0 + for token in tokens: + if token in ("(", "[", "{"): + depth += 1 + if token in (")", "]", "}"): + depth -= 1 + if token == "," and depth == 0: + groups.append([]) + else: + groups[-1].append(token) + params: dict[str, str | int | bool] = {} + for group in groups: + parameter = next((p for p in PARAMETERS if group and p[1] == group[0]), None) + if parameter is None or len(group) < 3 or group[1] != "=": + raise ValueError( + f"Invalid or unsupported text index parameter: {format_text_sql(group)}" + ) + field, key, kind = parameter + if field in params: + raise ValueError(f"Duplicate text index parameter: {key}") + raw = format_text_sql(group[2:]) + if kind == "boolean": + if raw.lower() not in ("0", "1", "true", "false"): + raise ValueError(f"Invalid boolean for {key}: {raw}") + params[field] = raw.lower() in ("1", "true") + elif kind == "number": + if not re.fullmatch(r"\d+", raw) or not 0 < int(raw) <= 2**53 - 1: + raise ValueError(f"Invalid integer for {key}: {raw}") + params[field] = int(raw) + elif kind == "codec": + if raw not in ("'none'", "'bitpacking'"): + raise ValueError(f"Invalid codec: {raw}") + params[field] = raw[1:-1] + else: + params[field] = raw + if not params.get("tokenizer"): + raise ValueError("A non-empty tokenizer is required") + return params + + +def canonicalize_text_index(index: SkipIndexText) -> SkipIndexText: + sql = render_text_index_type(index) + return index.model_copy( + update={ + "expression": text_expression_fingerprint(index.expression), + "granularity": TEXT_INDEX_GRANULARITY, + **parse_text_index_params(sql[5:-1]), + } + ) + + +def text_index_fingerprint(index: SkipIndexText) -> str: + canonical = canonicalize_text_index(index) + values = canonical.model_dump(exclude_none=True) + for field in ("expression", "tokenizer", "preprocessor", "postprocessor"): + value = getattr(canonical, field) + if value is not None: + values[field] = text_sql_fingerprint(value) + return json.dumps(values, sort_keys=True) diff --git a/chkit_python/src/chkit/core/text_index_sql.py b/chkit_python/src/chkit/core/text_index_sql.py new file mode 100644 index 00000000..051ed0cb --- /dev/null +++ b/chkit_python/src/chkit/core/text_index_sql.py @@ -0,0 +1,178 @@ +"""Small text-index SQL lexer; mirrors core/src/text-index-sql.ts.""" + +from __future__ import annotations + +import re + +_ESCAPES = {"0": 0, "a": 7, "b": 8, "t": 9, "n": 10, "v": 11, "f": 12, "r": 13, "e": 27} +_WORD = re.compile(r"(?:\d+(?:\.\d*)?(?:[eE][+-]?\d+)?|[\w$]+|->|<=|>=|!=|<>|\|\||::|==)") +_CLOSE = {")": "(", "]": "[", "}": "{"} +# Mirror ClickHouse's writeProbablyQuotedStringImpl (src/IO/WriteHelpers.cpp), +# plus NULL, which its isValidIdentifier helper excludes separately. +_QUOTED_IDENTIFIERS = { + "null", + "true", + "false", + "inf", + "infinity", + "nan", + "distinct", + "all", + "some", + "table", + "select", + "from", + "top", + "values", +} + + +def _string_literal(body: str) -> str: + data = bytearray() + i = 0 + while i < len(body): + char = body[i] + if body[i : i + 2] == "''": + data.append(39) + i += 2 + elif char == "\\": + next_char = body[i + 1] + hex_value = body[i + 2 : i + 4] + if next_char == "x" and re.fullmatch(r"[\da-fA-F]{2}", hex_value): + data.append(int(hex_value, 16)) + i += 4 + else: + if next_char == "x": + raise ValueError("Invalid hexadecimal SQL escape") + if next_char == "N": + i += 2 + continue + if next_char not in _ESCAPES and next_char not in ("'", '"', "`", "\\", "/", "="): + data.append(92) + if next_char in _ESCAPES: + data.append(_ESCAPES[next_char]) + else: + data.extend(next_char.encode()) + i += 2 + else: + data.extend(char.encode()) + i += 1 + parts = [] + for byte in data: + if byte == 39: + parts.append("\\'") + elif byte == 92: + parts.append("\\\\") + elif 32 <= byte < 127: + parts.append(chr(byte)) + else: + parts.append(f"\\x{byte:02x}") + return "'" + "".join(parts) + "'" + + +def text_sql_tokens(sql: str) -> list[str]: + tokens: list[str] = [] + brackets: list[str] = [] + i = 0 + while i < len(sql): + char = sql[i] + if char.isspace(): + i += 1 + continue + if sql.startswith("--", i): + end = sql.find("\n", i + 2) + i = len(sql) if end < 0 else end + 1 + continue + if sql.startswith("/*", i): + end = sql.find("*/", i + 2) + if end < 0: + raise ValueError("Unterminated SQL comment") + i = end + 2 + continue + if char in ("'", '"', "`"): + start = i + i += 1 + closed = False + while i < len(sql): + if sql[i] == "\\": + i += 2 + continue + current = sql[i] + i += 1 + if current == char: + if i < len(sql) and sql[i] == char: + i += 1 + continue + closed = True + break + if not closed: + raise ValueError("Unterminated SQL quote") + raw = sql[start:i] + tokens.append(_string_literal(raw[1:-1]) if char == "'" else raw) + continue + if char == ";": + raise ValueError("Expected a SQL expression, not a statement") + if char in "([{": + brackets.append(char) + if char in _CLOSE and (not brackets or brackets.pop() != _CLOSE[char]): + raise ValueError("Unbalanced SQL brackets") + match = _WORD.match(sql, i) + token = match.group() if match else char + tokens.append(token) + i += len(token) + if brackets: + raise ValueError("Unbalanced SQL brackets") + return tokens + + +def format_text_sql(tokens: list[str]) -> str: + parts = [] + for i, token in enumerate(tokens): + previous = tokens[i - 1] if i else "" + tight = ( + i == 0 + or token in (")", "]", "}", ",", ".") + or previous in ("(", "[", "{", ".") + or (token == "(" and re.fullmatch(r"[^\W\d]\w*", previous)) + ) + parts.append(("" if tight else " ") + token) + return "".join(parts) + + +def normalize_text_index_sql(sql: str) -> str: + return format_text_sql(text_sql_tokens(sql)) + + +def text_expression_fingerprint(sql: str) -> str: + tokens = text_sql_tokens(sql) + while tokens and tokens[0] == "(" and tokens[-1] == ")": + depth = 0 + wraps = True + for i, token in enumerate(tokens): + if token == "(": + depth += 1 + if token == ")": + depth -= 1 + if depth == 0 and i < len(tokens) - 1: + wraps = False + break + if not wraps: + break + tokens = tokens[1:-1] + return format_text_sql(tokens) + + +def text_sql_fingerprint(sql: str) -> list[str]: + """Ignore redundant simple identifier quotes only for comparisons.""" + result = [] + for token in text_sql_tokens(sql): + match = re.fullmatch(r'([`"])([A-Za-z_][A-Za-z0-9_]*)\1', token) + if match: + identifier = match[2] + # Normalize quote style without erasing meaningful quotes or case. + result.append( + f"`{identifier}`" if identifier.lower() in _QUOTED_IDENTIFIERS else identifier + ) + else: + result.append(token) + return result diff --git a/chkit_python/src/chkit/core/validate.py b/chkit_python/src/chkit/core/validate.py index c6ba8012..a7e96374 100644 --- a/chkit_python/src/chkit/core/validate.py +++ b/chkit_python/src/chkit/core/validate.py @@ -20,6 +20,7 @@ ValidationIssueCode, ) from chkit.core.projection import is_index_projection, normalize_projection_index +from chkit.core.text_index import render_text_index_type def _push( @@ -86,6 +87,27 @@ def _validate_column_codec( ) +def _validate_indexes(definition: TableDefinition, issues: list[ValidationIssue]) -> None: + index_seen: set[str] = set() + for index in definition.indexes or []: + if index.name in index_seen: + _push( + issues, + definition, + "duplicate_index_name", + f'Table {definition.database}.{definition.name} ' + f'has duplicate index name "{index.name}"', + ) + continue + index_seen.add(index.name) + if index.type == "text": + try: + render_text_index_type(index) + except ValueError as exc: + _push(issues, definition, "text_index_invalid_parameters", + f'Text index "{index.name}": {exc}') + + def _validate_table(definition: TableDefinition, issues: list[ValidationIssue]) -> None: column_seen: set[str] = set() column_set: set[str] = set() @@ -103,18 +125,7 @@ def _validate_table(definition: TableDefinition, issues: list[ValidationIssue]) column_set.add(column.name) _validate_column_codec(definition, column, issues) - index_seen: set[str] = set() - for index in definition.indexes or []: - if index.name in index_seen: - _push( - issues, - definition, - "duplicate_index_name", - f'Table {definition.database}.{definition.name} ' - f'has duplicate index name "{index.name}"', - ) - continue - index_seen.add(index.name) + _validate_indexes(definition, issues) projection_seen: set[str] = set() for projection in definition.projections or []: diff --git a/chkit_python/tests/test_text_index.py b/chkit_python/tests/test_text_index.py new file mode 100644 index 00000000..efe22c37 --- /dev/null +++ b/chkit_python/tests/test_text_index.py @@ -0,0 +1,176 @@ +"""Shared adversarial corpus exercises both Python and TypeScript implementations.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from chkit import SkipIndexText, table +from chkit.cli.commands.pull_render import render_schema_file +from chkit.core.canonical import canonicalize_definitions +from chkit.core.planner import plan_diff +from chkit.core.sql import to_create_sql +from chkit.core.text_index import ( + canonicalize_text_index, + parse_text_index_params, + render_text_index_type, +) +from chkit.core.text_index_sql import normalize_text_index_sql +from chkit.core.validate import validate_definitions + +CASES = json.loads( + (Path(__file__).resolve().parents[2] / "test/fixtures/text-index.json").read_text() +) +IDENTIFIERS = json.loads( + (Path(__file__).resolve().parents[2] / "test/fixtures/text-index-identifiers.json").read_text() +) + + +def docs(index, name="docs", database="app"): + return table( + database=database, + name=name, + engine="MergeTree()", + columns=[{"name": "id", "type": "UInt64"}, {"name": "body", "type": "String"}], + primary_key=["id"], + order_by=["id"], + indexes=[index], + ) + + +@pytest.mark.parametrize("case", CASES, ids=lambda case: case["name"]) +def test_text_round_trip(case): + params = {key: value for key, value in case.items() if key != "name"} + index = SkipIndexText.model_validate({"name": "idx", "expression": "body", **params}) + sql = render_text_index_type(index) + parsed = index.model_copy(update=parse_text_index_params(sql[5:-1])) + assert render_text_index_type(parsed) == sql + assert canonicalize_text_index(canonicalize_text_index(index)) == canonicalize_text_index(index) + namespace = {} + exec(compile(render_schema_file([docs(index)]), "schema.py", "exec"), namespace) + assert canonicalize_definitions(namespace["definitions"]) == canonicalize_definitions( + [docs(index)] + ) + + +def test_literal_whitespace_changes_require_migrations(): + index = SkipIndexText( + name="idx", + expression="body", + tokenizer="splitByString([' '])", + preprocessor="replaceAll(body, ' ', ' ')", + ) + original = docs(index) + assert "splitByString([' '])" in to_create_sql(original) + for change in ( + {"tokenizer": "splitByString([' '])"}, + {"preprocessor": "replaceAll(body, ' ', ' ')"}, + ): + changed = docs(index.model_copy(update=change)) + assert [op.type for op in plan_diff([original], [changed]).operations] == [ + "alter_table_drop_index", + "alter_table_add_index", + ] + + +@pytest.mark.parametrize("granularity", [1, 64, 100000000]) +def test_granularity_does_not_change_schema(granularity): + index = SkipIndexText(name="idx", expression="body", tokenizer="splitByNonAlpha") + changed = docs(index.model_copy(update={"granularity": granularity})) + assert "GRANULARITY 100000000" in to_create_sql(changed) + assert plan_diff([docs(index)], [changed]).operations == [] + + +def test_equivalent_escapes(): + assert normalize_text_index_sql("splitByString(['\t', 'é', ''''])") == normalize_text_index_sql( + r"splitByString(['\x09', '\xc3\xa9', '\''])" + ) + + +@pytest.mark.parametrize( + "args", + [ + "", + "tokenizer =", + "tokenizer = ngrams(3", + "tokenizer = 'oops", + "tokenizer = ngrams(3],)", + "tokenizer = splitByNonAlpha,", + "tokenizer = splitByNonAlpha, tokenizer = ngrams(2)", + "tokenizer = splitByNonAlpha, unknown = 1", + "tokenizer = splitByNonAlpha, support_phrase_search = 2", + "tokenizer = splitByNonAlpha, dictionary_block_size = 0", + "tokenizer = splitByNonAlpha, dictionary_block_size = 1.5", + "tokenizer = splitByNonAlpha, dictionary_block_size = 9007199254740992", + "tokenizer = splitByNonAlpha, posting_list_codec = 'bogus'", + "tokenizer = splitByNonAlpha; SELECT 1", + "tokenizer = /* unclosed", + ], +) +def test_rejects_lossy_or_malformed_metadata(args): + with pytest.raises(ValueError, match=r".+"): + parse_text_index_params(args) + + +@pytest.mark.parametrize( + "params", + [ + {"tokenizer": ""}, + {"tokenizer": "/*empty*/"}, + {"dictionary_block_size": -1}, + {"posting_list_block_size": 0}, + {"preprocessor": ""}, + ], +) +def test_validation_reports_invalid_parameters(params): + index = SkipIndexText.model_validate( + {"name": "idx", "expression": "body", "tokenizer": "splitByNonAlpha", **params} + ) + assert "text_index_invalid_parameters" in [ + issue.code for issue in validate_definitions([docs(index)]) + ] + + +def test_all_options_and_parameter_order(): + index = SkipIndexText( + name="idx", + expression="concat(body, ') (')", + tokenizer="splitByNonAlpha", + postprocessor="lower(body)", + support_phrase_search=False, + dictionary_block_frontcoding_compression=False, + dictionary_block_size=512, + posting_list_block_size=1024, + posting_list_codec="bitpacking", + ) + sql = render_text_index_type(index) + assert "support_phrase_search = 0" in sql + pulled = index.model_copy(update=parse_text_index_params(sql[5:-1])) + assert pulled == index + + +def test_identifier_quotes_do_not_rebuild_or_conflate_string_literals(): + index = SkipIndexText( + name="idx", expression="body", tokenizer="splitByNonAlpha", preprocessor="lower(`body`)" + ) + unquoted = index.model_copy(update={"preprocessor": "lower(body)"}) + literal = index.model_copy(update={"preprocessor": "lower('body')"}) + assert plan_diff([docs(index)], [docs(unquoted)]).operations == [] + assert len(plan_diff([docs(index)], [docs(literal)]).operations) == 2 + + +@pytest.mark.parametrize("word", IDENTIFIERS) +def test_meaningful_identifier_quotes(word): + index = SkipIndexText(name="idx", expression="body", tokenizer="splitByNonAlpha") + for name in (word, word.upper(), word.capitalize()): + for field in ("expression", "preprocessor", "postprocessor"): + quoted = docs(index.model_copy(update={field: f"toString(`{name}`)"})) + double_quoted = docs(index.model_copy(update={field: f'toString("{name}")'})) + unquoted = docs(index.model_copy(update={field: f"toString({name})"})) + assert plan_diff([quoted], [double_quoted]).operations == [] + assert [op.type for op in plan_diff([quoted], [unquoted]).operations] == [ + "alter_table_drop_index", + "alter_table_add_index", + ] diff --git a/chkit_python/tests/test_text_index_e2e.py b/chkit_python/tests/test_text_index_e2e.py new file mode 100644 index 00000000..e88798ea --- /dev/null +++ b/chkit_python/tests/test_text_index_e2e.py @@ -0,0 +1,250 @@ +"""Execute real CREATE/ALTER/pull/drift and compare indexed search results.""" + +from __future__ import annotations + +import json +from pathlib import Path + +import pytest + +from chkit import SkipIndexText +from chkit.cli.commands.drift_compare import compare_table_shape +from chkit.cli.commands.pull_render import render_schema_file +from chkit.clickhouse.client import ClickHouseClient +from chkit.clickhouse.introspect import list_table_details +from chkit.core.model import ChxResolvedClickHouseConfig +from chkit.core.planner import plan_diff +from chkit.core.sql import to_create_sql +from chkit.core.text_index import text_index_fingerprint +from chkit.core.text_index_sql import normalize_text_index_sql +from tests.e2e_testkit import create_prefix, get_required_env +from tests.test_text_index import docs + +CASES = json.loads( + (Path(__file__).resolve().parents[2] / "test/fixtures/text-index.json").read_text() +) + + +@pytest.fixture +def text_client(): + env = get_required_env() + with ClickHouseClient.connect( + ChxResolvedClickHouseConfig( + url=env.clickhouse_url, + username=env.clickhouse_user, + password=env.clickhouse_password, + database=env.clickhouse_database, + secure=env.clickhouse_url.startswith("https:"), + ) + ) as client: + yield client, env.clickhouse_database + + +@pytest.mark.parametrize("case", CASES, ids=lambda case: case["name"]) +def test_live_text_index_round_trip(text_client, case): + client, database = text_client + name = create_prefix("py_text") + "docs" + params = {key: value for key, value in case.items() if key != "name"} + index = SkipIndexText.model_validate( + {"name": "idx", "expression": "body", "granularity": 1, **params} + ) + definition = docs(index, name=name, database=database) + full_name = f"{database}.{name}" + clone_name = f"{database}.{name}_clone" + try: + pre = f", preprocessor = {index.preprocessor}" if index.preprocessor else "" + client.execute( + f"CREATE TABLE {full_name} (id UInt64, body String, INDEX idx ({index.expression}) " + f"TYPE text(tokenizer = {index.tokenizer}{pre}) GRANULARITY 1) ENGINE=MergeTree ORDER BY id" + ) + client.execute( + f"INSERT INTO {full_name} VALUES (1, 'alpha beta gamma'), (2, 'alpha beta gamma'), (3, 'alpha'), (4, 'é東京😀'), (5, 'a,b=c')" + ) + actual = next(item for item in list_table_details(client, [database]) if item.name == name) + assert compare_table_shape(definition, actual) is None + assert actual.indexes[0].granularity == 100000000 + namespace = {} + exec( + compile( + render_schema_file([definition.model_copy(update={"indexes": actual.indexes})]), + "schema.py", + "exec", + ), + namespace, + ) + pulled = namespace["definitions"][0] + assert plan_diff([definition], [pulled]).operations == [] + client.execute(to_create_sql(pulled.model_copy(update={"name": name + "_clone"}))) + client.execute(f"INSERT INTO {clone_name} SELECT * FROM {full_name}") + predicate = f"WHERE hasAllTokens({index.expression}, ['alpha']) ORDER BY id" + original = client.query(f"SELECT id FROM {full_name} {predicate}").rows + assert client.query(f"SELECT id FROM {clone_name} {predicate}").rows == original + if case["name"] == "two spaces": + assert [int(row["id"]) for row in original] == [1, 3] + finally: + client.execute(f"DROP TABLE IF EXISTS {full_name} SYNC") + client.execute(f"DROP TABLE IF EXISTS {clone_name} SYNC") + + +def test_live_add_and_change_text_index(text_client): + client, database = text_client + name = create_prefix("py_text_alter") + "docs" + full_name = f"{database}.{name}" + index = SkipIndexText( + name="idx", expression="body", tokenizer="splitByString([' '])", granularity=1 + ) + with_index = docs(index, name=name, database=database) + without_index = with_index.model_copy(update={"indexes": []}) + changed = with_index.model_copy( + update={"indexes": [index.model_copy(update={"tokenizer": "splitByString([' '])"})]} + ) + try: + client.execute(to_create_sql(without_index)) + for before, after in ((without_index, with_index), (with_index, changed)): + plan = plan_diff([before], [after]) + assert plan.operations + for op in plan.operations: + client.execute(op.sql) + actual = next( + item for item in list_table_details(client, [database]) if item.name == name + ) + assert compare_table_shape(after, actual) is None + assert "index_mismatch" in compare_table_shape(with_index, actual).reason_codes + finally: + client.execute(f"DROP TABLE IF EXISTS {full_name} SYNC") + + +def test_tuning_newer_options_and_materializing_existing_rows(text_client): + client, database = text_client + version = str(client.query("SELECT version() AS version").rows[0]["version"]) + newer_options = tuple(int(part) for part in version.split(".")[:2]) >= (26, 8) + name = create_prefix("py_text_options") + "docs" + full_name = f"{database}.{name}" + index = SkipIndexText( + name="idx", + expression="body", + tokenizer="splitByNonAlpha", + dictionary_block_size=512, + dictionary_block_frontcoding_compression=False, + posting_list_block_size=1024, + posting_list_codec="bitpacking", + postprocessor="lower(body)" if newer_options else None, + support_phrase_search=True if newer_options else None, + ) + definition = docs(index, name=name, database=database) + if newer_options: + definition = definition.model_copy( + update={"settings": {"allow_experimental_text_index_phrase_search": 1}} + ) + without_index = definition.model_copy(update={"indexes": []}) + try: + client.execute(to_create_sql(without_index)) + client.execute(f"INSERT INTO {full_name} VALUES (1, 'hello world'), (2, 'goodbye world')") + for op in plan_diff([without_index], [definition]).operations: + client.execute(op.sql) + client.execute(f"ALTER TABLE {full_name} MATERIALIZE INDEX idx SETTINGS mutations_sync = 2") + actual = next(item for item in list_table_details(client, [database]) if item.name == name) + assert compare_table_shape(definition, actual) is None + namespace = {} + exec( + compile( + render_schema_file([definition.model_copy(update={"indexes": actual.indexes})]), + "schema.py", + "exec", + ), + namespace, + ) + assert plan_diff([definition], namespace["definitions"]).operations == [] + fn = "hasPhrase(body, 'hello world')" if newer_options else "hasAllTokens(body, ['hello'])" + assert [ + int(row["id"]) for row in client.query(f"SELECT id FROM {full_name} WHERE {fn}").rows + ] == [1] + finally: + client.execute(f"DROP TABLE IF EXISTS {full_name} SYNC") + + +def test_normalization_preserves_every_printable_clickhouse_escape(text_client): + client, _ = text_client + for code in range(32, 127): + sql = "'\\" + chr(code) + "'" + if code == 120: + with pytest.raises(ValueError, match="Invalid hexadecimal"): + normalize_text_index_sql(sql) + continue + rows = client.query( + f"SELECT hex({sql}) AS original, hex({normalize_text_index_sql(sql)}) AS normalized" + ).rows + assert rows[0]["normalized"] == rows[0]["original"], repr(sql) + + +def test_quoted_literal_names_remain_distinct_from_constants(text_client): + client, _ = text_client + for word in ("null", "true", "false", "inf", "infinity", "nan"): + for name in (word, word.upper(), word.capitalize()): + quoted, unquoted = f"toString(`{name}`)", f"toString({name})" + rows = client.query( + f"SELECT {quoted} AS quoted_value, {unquoted} AS literal_value " + f"FROM (SELECT 'sentinel' AS `{name}`)" + ).rows + assert rows[0]["quoted_value"] == "sentinel" + assert rows[0]["literal_value"] != "sentinel" + index = SkipIndexText(name="idx", expression=quoted, tokenizer="splitByNonAlpha") + assert text_index_fingerprint(index) != text_index_fingerprint( + index.model_copy(update={"expression": unquoted}) + ) + + +def test_quoted_null_column_round_trips_and_literal_change_migrates(text_client): + client, database = text_client + name = create_prefix("py_text_keyword") + "docs" + full_name = f"{database}.{name}" + index = SkipIndexText( + name="idx", + tokenizer="splitByNonAlpha", + expression="concat(body, ifNull(\"NULL\", 'missing'))", + ) + base = docs(index, name=name, database=database) + definition = type(base).model_validate( + {**base.model_dump(), "columns": [*base.columns, {"name": "NULL", "type": "String"}]} + ) + changed_index = index.model_copy(update={"expression": "concat(body, ifNull(NULL, 'missing'))"}) + changed = definition.model_copy(update={"indexes": [changed_index]}) + + def get_actual(): + return next(item for item in list_table_details(client, [database]) if item.name == name) + + def search(expression): + return client.query( + f"SELECT id FROM {full_name} WHERE hasAllTokens({expression}, ['alpha'])" + ).rows + + try: + client.execute(to_create_sql(definition)) + client.execute(f"INSERT INTO {full_name} VALUES (1, 'doc ', 'alpha')") + actual = get_actual() + assert compare_table_shape(definition, actual) is None + assert "index_mismatch" in compare_table_shape(changed, actual).reason_codes + namespace = {} + exec( + compile( + render_schema_file([definition.model_copy(update={"indexes": actual.indexes})]), + "schema.py", + "exec", + ), + namespace, + ) + pulled = namespace["definitions"][0] + assert plan_diff([definition], [pulled]).operations == [] + assert [int(row["id"]) for row in search(index.expression)] == [1] + plan = plan_diff([pulled], [changed]) + assert [op.type for op in plan.operations] == [ + "alter_table_drop_index", + "alter_table_add_index", + ] + for op in plan.operations: + client.execute(op.sql) + client.execute(f"ALTER TABLE {full_name} MATERIALIZE INDEX idx SETTINGS mutations_sync = 2") + assert compare_table_shape(changed, get_actual()) is None + assert search(changed_index.expression) == [] + finally: + client.execute(f"DROP TABLE IF EXISTS {full_name} SYNC") diff --git a/packages/cli/src/commands/drift/compare.ts b/packages/cli/src/commands/drift/compare.ts index 009fae9a..ffa74a66 100644 --- a/packages/cli/src/commands/drift/compare.ts +++ b/packages/cli/src/commands/drift/compare.ts @@ -6,6 +6,7 @@ import { type ColumnDefinition, type ProjectionDefinition, type SkipIndexDefinition, + textIndexFingerprint, type TableDefinition, } from '@chkit/core' import { diffByName, diffNamedShapeMaps, diffSettings } from './diff.js' @@ -205,6 +206,8 @@ function normalizeColumnShape(column: ColumnDefinition): string { function renderIndexTypeFingerprint(index: SkipIndexDefinition): string { switch (index.type) { + case 'text': + return textIndexFingerprint(index) case 'minmax': return 'minmax' case 'set': @@ -235,6 +238,7 @@ function stripEnclosingParens(value: string): string { } function normalizeIndexShape(index: SkipIndexDefinition): string { + if (index.type === 'text') return textIndexFingerprint(index) return [ `expr=${stripEnclosingParens(normalizeSQLFragment(index.expression))}`, `type=${renderIndexTypeFingerprint(index)}`, diff --git a/packages/cli/src/test/text-index.e2e.test.ts b/packages/cli/src/test/text-index.e2e.test.ts new file mode 100644 index 00000000..f5571f8d --- /dev/null +++ b/packages/cli/src/test/text-index.e2e.test.ts @@ -0,0 +1,341 @@ +import { describe, expect, test } from 'bun:test' +import { mkdtemp, rm, writeFile } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { + canonicalizeDefinitions, + planDiff, + table, + textIndexFingerprint, + toCreateSQL, + type TableDefinition, + type TextSkipIndex, +} from '@chkit/core' +import { renderSchemaFile } from '../../../plugin-pull/src/render-schema.js' +import fixtures from '../../../../test/fixtures/text-index.json' +import { compareTableShape } from '../commands/drift/compare.js' +import { + CORE_ENTRY, + createJournalTableName, + createLiveExecutor, + createPrefix, + formatTestDiagnostic, + getRequiredEnv, + runCli, + runCliWithRetry, + waitForTable, +} from './e2e-testkit.js' + +const env = getRequiredEnv() +const docs = (name: string, index: TextSkipIndex): TableDefinition => + table({ + database: env.clickhouseDatabase, + name, + engine: 'MergeTree()', + columns: [ + { name: 'id', type: 'UInt64' }, + { name: 'body', type: 'String' }, + ], + primaryKey: ['id'], + orderBy: ['id'], + indexes: [index], + }) + +async function loadPulled(definition: TableDefinition, dir: string): Promise { + const path = join(dir, `${definition.name}.ts`) + await writeFile( + path, + renderSchemaFile([definition]).replace("'@chkit/core'", JSON.stringify(CORE_ENTRY)), + ) + return (await import(path)).default[0] +} + +describe('text index live round trips', () => { + for (const { name: label, ...params } of fixtures) { + test(label, async () => { + const executor = createLiveExecutor(env) + const name = `${createPrefix('text')}docs` + const index: TextSkipIndex = { + name: 'idx', + expression: 'body', + type: 'text', + granularity: 1, + ...params, + } + const definition = docs(name, index) + const fullName = `${definition.database}.${name}` + const cloneName = `${name}_clone` + const dir = await mkdtemp(join(tmpdir(), 'chkit-text-pull-')) + try { + // Start from SQL exactly as a user wrote it, independently of our renderer. + await executor.command( + `CREATE TABLE ${fullName} (id UInt64, body String, INDEX idx (${index.expression}) TYPE text(tokenizer = ${index.tokenizer}${index.preprocessor ? `, preprocessor = ${index.preprocessor}` : ''}) GRANULARITY 1) ENGINE=MergeTree ORDER BY id`, + ) + await waitForTable(executor, definition.database, name) + await executor.command( + `INSERT INTO ${fullName} VALUES (1, 'alpha beta gamma'), (2, 'alpha beta gamma'), (3, 'alpha'), (4, 'é東京😀'), (5, 'a,b=c')`, + ) + const actual = (await executor.listTableDetails([definition.database])).find( + (item) => item.name === name, + ) + if (!actual) throw new Error('Missing test table') + expect(compareTableShape(definition, actual)).toBeNull() + expect(actual.indexes[0]?.granularity).toBe(100000000) + const pulled = await loadPulled({ ...definition, indexes: actual.indexes }, dir) + expect(planDiff([definition], [pulled]).operations).toEqual([]) + await executor.command(toCreateSQL({ ...pulled, name: cloneName })) + await executor.command( + `INSERT INTO ${definition.database}.${cloneName} SELECT * FROM ${fullName}`, + ) + const query = (target: string) => + executor.query<{ id: number }>( + `SELECT id FROM ${target} WHERE hasAllTokens(${index.expression}, ['alpha']) ORDER BY id`, + ) + const original = await query(fullName) + expect(await query(`${definition.database}.${cloneName}`)).toEqual(original) + if (label === 'two spaces') expect(original.map((row) => Number(row.id))).toEqual([1, 3]) + } finally { + await executor.command(`DROP TABLE IF EXISTS ${fullName} SYNC`) + await executor.command(`DROP TABLE IF EXISTS ${definition.database}.${cloneName} SYNC`) + await executor.close() + await rm(dir, { recursive: true, force: true }) + } + }, 60000) + } + + test('generate, migrate, drift, change tokenizer, and migrate again', async () => { + const executor = createLiveExecutor(env) + const name = `${createPrefix('text_cli')}docs` + const definition = docs(name, { + name: 'idx', + type: 'text', + expression: 'body', + tokenizer: "splitByString([' '])", + granularity: 1, + }) + const dir = await mkdtemp(join(tmpdir(), 'chkit-text-cli-')) + const schemaPath = join(dir, 'schema.ts') + const configPath = join(dir, 'clickhouse.config.ts') + const journal = createJournalTableName('text_index') + const extraEnv = { CHKIT_JOURNAL_TABLE: journal } + const command = (args: string[]) => { + const result = runCli(dir, [...args, '--config', configPath, '--json'], extraEnv) + if (result.exitCode !== 0) throw new Error(formatTestDiagnostic(args.join(' '), result)) + return JSON.parse(result.stdout) + } + try { + await writeFile( + configPath, + `export default ${JSON.stringify({ schema: schemaPath, outDir: join(dir, 'chkit'), migrationsDir: join(dir, 'chkit/migrations'), metaDir: join(dir, 'chkit/meta'), clickhouse: { url: env.clickhouseUrl, username: env.clickhouseUser, password: env.clickhousePassword, database: env.clickhouseDatabase } })}`, + ) + const writeSchema = async (value: TableDefinition) => + writeFile( + schemaPath, + renderSchemaFile([value]).replace("'@chkit/core'", JSON.stringify(CORE_ENTRY)), + ) + const migrate = async () => { + const result = await runCliWithRetry( + dir, + ['migrate', '--execute', '--config', configPath, '--json'], + { extraEnv }, + ) + if (result.exitCode !== 0) throw new Error(formatTestDiagnostic('migrate', result)) + } + await writeSchema(definition) + command(['generate']) + await migrate() + await waitForTable(executor, definition.database, name) + expect(command(['drift', '--table', `${definition.database}.${name}`]).drifted).toBe(false) + const changed = { + ...definition, + indexes: [ + { + ...definition.indexes?.[0], + tokenizer: "splitByString([' '])", + } as TextSkipIndex, + ], + } + expect(planDiff([definition], [changed]).operations).toHaveLength(2) + await writeSchema(changed) + command(['generate']) + await migrate() + expect(command(['drift', '--table', `${definition.database}.${name}`]).drifted).toBe(false) + const actual = (await executor.listTableDetails([definition.database])).find( + (item) => item.name === name, + ) + if (!actual) throw new Error('Missing test table') + expect(compareTableShape(definition, actual)?.reasonCodes).toContain('index_mismatch') + expect(compareTableShape(changed, actual)).toBeNull() + expect(canonicalizeDefinitions([changed])).not.toEqual(canonicalizeDefinitions([definition])) + } finally { + await executor.command(`DROP TABLE IF EXISTS ${definition.database}.${name} SYNC`) + await executor.command(`DROP TABLE IF EXISTS ${definition.database}.${journal} SYNC`) + await executor.close() + await rm(dir, { recursive: true, force: true }) + } + }, 120000) +}) + +test('text index tuning, newer options, and materializing existing rows', async () => { + const executor = createLiveExecutor(env) + const name = `${createPrefix('text_options')}docs` + const [{ version }] = await executor.query<{ version: string }>('SELECT version() AS version') + const [major, minor] = version.split('.').map(Number) + const newerOptions = major > 26 || (major === 26 && minor >= 8) + const index: TextSkipIndex = { + name: 'idx', + expression: 'body', + type: 'text', + tokenizer: 'splitByNonAlpha', + dictionaryBlockSize: 512, + dictionaryBlockFrontcodingCompression: false, + postingListBlockSize: 1024, + postingListCodec: 'bitpacking', + ...(newerOptions ? { postprocessor: 'lower(body)', supportPhraseSearch: true } : {}), + } + const definition = { + ...docs(name, index), + ...(newerOptions ? { settings: { allow_experimental_text_index_phrase_search: 1 } } : {}), + } + const fullName = `${definition.database}.${name}` + const dir = await mkdtemp(join(tmpdir(), 'chkit-text-options-')) + try { + await executor.command(toCreateSQL({ ...definition, indexes: [] })) + await executor.command( + `INSERT INTO ${fullName} VALUES (1, 'hello world'), (2, 'goodbye world')`, + ) + for (const op of planDiff([{ ...definition, indexes: [] }], [definition]).operations) + await executor.command(op.sql) + await executor.command( + `ALTER TABLE ${fullName} MATERIALIZE INDEX idx SETTINGS mutations_sync = 2`, + ) + const actual = (await executor.listTableDetails([definition.database])).find( + (item) => item.name === name, + ) + if (!actual) throw new Error('Missing test table') + expect(compareTableShape(definition, actual)).toBeNull() + const pulled = await loadPulled({ ...definition, indexes: actual.indexes }, dir) + expect(planDiff([definition], [pulled]).operations).toEqual([]) + const fn = newerOptions ? "hasPhrase(body, 'hello world')" : "hasAllTokens(body, ['hello'])" + expect( + (await executor.query<{ id: number }>(`SELECT id FROM ${fullName} WHERE ${fn}`)).map((row) => + Number(row.id), + ), + ).toEqual([1]) + } finally { + await executor.command(`DROP TABLE IF EXISTS ${fullName} SYNC`) + await executor.close() + await rm(dir, { recursive: true, force: true }) + } +}, 60000) + +test('normalization preserves every printable ClickHouse string escape', async () => { + const { normalizeTextIndexSQL } = await import('@chkit/core') + const executor = createLiveExecutor(env) + try { + for (let code = 32; code < 127; code++) { + const sql = `'\\${String.fromCharCode(code)}'` + if (code === 120) { + expect(() => normalizeTextIndexSQL(sql)).toThrow('Invalid hexadecimal') + continue + } + const rows = await executor.query<{ + original: string + normalized: string + }>(`SELECT hex(${sql}) AS original, hex(${normalizeTextIndexSQL(sql)}) AS normalized`) + expect(rows[0]?.normalized, `escape ${JSON.stringify(sql)}`).toBe(rows[0]?.original) + } + } finally { + await executor.close() + } +}) + +test('quoted literal names remain distinct from constants in ClickHouse and planning', async () => { + const executor = createLiveExecutor(env) + try { + for (const word of ['null', 'true', 'false', 'inf', 'infinity', 'nan']) { + for (const name of [word, word.toUpperCase(), word[0]?.toUpperCase() + word.slice(1)]) { + const quoted = `toString(\`${name}\`)` + const unquoted = `toString(${name})` + const rows = await executor.query<{ + quoted_value: string + literal_value: string | null + }>( + `SELECT ${quoted} AS quoted_value, ${unquoted} AS literal_value FROM (SELECT 'sentinel' AS \`${name}\`)`, + ) + expect(rows[0]?.quoted_value).toBe('sentinel') + expect(rows[0]?.literal_value).not.toBe('sentinel') + const index: TextSkipIndex = { + name: 'idx', + type: 'text', + expression: quoted, + tokenizer: 'splitByNonAlpha', + } + expect(textIndexFingerprint(index)).not.toBe( + textIndexFingerprint({ ...index, expression: unquoted }), + ) + } + } + } finally { + await executor.close() + } +}) + +test('quoted NULL column round-trips and changing it to a literal migrates the index', async () => { + const executor = createLiveExecutor(env) + const name = `${createPrefix('text_keyword')}docs` + const index: TextSkipIndex = { + name: 'idx', + type: 'text', + tokenizer: 'splitByNonAlpha', + expression: 'concat(body, ifNull("NULL", \'missing\'))', + } + const base = docs(name, index) + const definition = { + ...base, + columns: [...base.columns, { name: 'NULL', type: 'String' }], + } + const changedIndex = { + ...index, + expression: "concat(body, ifNull(NULL, 'missing'))", + } + const changed = { ...definition, indexes: [changedIndex] } + const fullName = `${definition.database}.${name}` + const dir = await mkdtemp(join(tmpdir(), 'chkit-text-keyword-')) + try { + await executor.command(toCreateSQL(definition)) + await executor.command(`INSERT INTO ${fullName} VALUES (1, 'doc ', 'alpha')`) + const getActual = async () => { + const actual = (await executor.listTableDetails([definition.database])).find( + (item) => item.name === name, + ) + if (!actual) throw new Error('Missing test table') + return actual + } + const actual = await getActual() + expect(compareTableShape(definition, actual)).toBeNull() + expect(compareTableShape(changed, actual)?.reasonCodes).toContain('index_mismatch') + const pulled = await loadPulled({ ...definition, indexes: actual.indexes }, dir) + expect(planDiff([definition], [pulled]).operations).toEqual([]) + const search = (expression: string) => + executor.query<{ id: number }>( + `SELECT id FROM ${fullName} WHERE hasAllTokens(${expression}, ['alpha'])`, + ) + expect((await search(index.expression)).map((row) => Number(row.id))).toEqual([1]) + const plan = planDiff([pulled], [changed]) + expect(plan.operations.map((op) => op.type)).toEqual([ + 'alter_table_drop_index', + 'alter_table_add_index', + ]) + for (const op of plan.operations) await executor.command(op.sql) + await executor.command( + `ALTER TABLE ${fullName} MATERIALIZE INDEX idx SETTINGS mutations_sync = 2`, + ) + expect(compareTableShape(changed, await getActual())).toBeNull() + expect(await search(changedIndex.expression)).toEqual([]) + } finally { + await executor.command(`DROP TABLE IF EXISTS ${fullName} SYNC`) + await executor.close() + await rm(dir, { recursive: true, force: true }) + } +}, 60000) diff --git a/packages/clickhouse/src/index.ts b/packages/clickhouse/src/index.ts index bc3376fd..8f5c482a 100644 --- a/packages/clickhouse/src/index.ts +++ b/packages/clickhouse/src/index.ts @@ -6,6 +6,8 @@ import { type ProjectionDefinition, parseCodec, type SkipIndexDefinition, + parseTextIndexParams, + normalizeTextIndexSQL, } from '@chkit/core' import { type ClickHouseSettings, ClickHouseLogLevel, createClient } from '@clickhouse/client' import { getLogger } from '@logtape/logtape' @@ -205,6 +207,7 @@ export function normalizeColumnFromSystemRow( } type ParsedIndexShape = + | ({ type: 'text' } & ReturnType) | { type: 'minmax' } | { type: 'set'; maxRows: number } | { type: 'bloom_filter'; falsePositiveRate?: number } @@ -231,8 +234,9 @@ function splitArgs(args: string | undefined): number[] { } function parseIndexType(value: string): ParsedIndexShape { - const match = value.match(/^(\w+)\((.+)\)$/) + const match = value.match(/^(\w+)\((.*)\)$/s) const baseName = match?.[1] ?? value + if (baseName === 'text') return { type: 'text', ...parseTextIndexParams(match?.[2] ?? '') } const args = splitArgs(match?.[2]) switch (baseName) { @@ -268,7 +272,7 @@ export function normalizeIndexFromSystemRow( const parsed = parseIndexType(row.type) return { name: row.name, - expression: normalizeSQLFragment(row.expr), + expression: parsed.type === 'text' ? normalizeTextIndexSQL(row.expr) : normalizeSQLFragment(row.expr), granularity: row.granularity, ...parsed, } diff --git a/packages/core/src/canonical.ts b/packages/core/src/canonical.ts index ae077f55..46313b9d 100644 --- a/packages/core/src/canonical.ts +++ b/packages/core/src/canonical.ts @@ -13,6 +13,7 @@ import { normalizeKeyColumns } from './key-clause.js' import { isSchemaDefinition } from './model.js' import { canonicalizeCodec } from './codec.js' import { canonicalizeProjection } from './projection.js' +import { canonicalizeTextIndex } from './text-index.js' import { normalizeEngine, normalizeSQLFragment } from './sql-normalizer.js' function sortByName(items: T[]): T[] { @@ -38,6 +39,7 @@ function canonicalizeColumn(column: ColumnDefinition): ColumnDefinition { } function canonicalizeIndex(index: SkipIndexDefinition): SkipIndexDefinition { + if (index.type === 'text') return canonicalizeTextIndex(index) return { ...index, expression: normalizeSQLFragment(index.expression), diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 1743f344..648cf754 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -8,6 +8,14 @@ export { } from './canonical.js' export { planDiff } from './planner.js' export { createSnapshot } from './snapshot.js' +export { + TEXT_INDEX_GRANULARITY, + canonicalizeTextIndex, + normalizeTextIndexSQL, + parseTextIndexParams, + renderTextIndexType, + textIndexFingerprint, +} from './text-index.js' export { splitTopLevelComma } from './key-clause.js' export { isIndexProjection, normalizeProjectionIndex } from './projection.js' export { normalizeEngine, normalizeSQLFragment } from './sql-normalizer.js' diff --git a/packages/core/src/model-types.ts b/packages/core/src/model-types.ts index c4a29ffe..16515d3b 100644 --- a/packages/core/src/model-types.ts +++ b/packages/core/src/model-types.ts @@ -68,6 +68,22 @@ interface SkipIndexBase { granularity: number } +/** Full-text index. ClickHouse ignores granularity and uses one index per part. */ +export interface TextSkipIndex extends Omit { + type: 'text' + /** SQL tokenizer, e.g. splitByNonAlpha or splitByString([' ', ';']). */ + tokenizer: string + granularity?: number + preprocessor?: string + /** Requires a ClickHouse version that supports postprocessing. */ + postprocessor?: string + supportPhraseSearch?: boolean + dictionaryBlockSize?: number + dictionaryBlockFrontcodingCompression?: boolean + postingListBlockSize?: number + postingListCodec?: 'none' | 'bitpacking' +} + /** * Skip index with structured, discriminated args per type. Arg signatures * come from ClickHouse MergeTree docs: @@ -76,11 +92,12 @@ interface SkipIndexBase { * - `bloom_filter([false_positive_rate])` — optional float, default 0.025 * - `tokenbf_v1(size_bytes, n_hash, seed)` — 3 required ints * - `ngrambf_v1(n, size_bytes, n_hash, seed)` — 4 required ints + * - `text(tokenizer = ..., ...)` — named parameters, automatic granularity * * ClickHouse 26+ requires `set(0)` not bare `set`; `maxRows` is required * so this is encoded naturally. */ -export type SkipIndexDefinition = SkipIndexBase & +export type SkipIndexDefinition = TextSkipIndex | SkipIndexBase & ( | { type: 'minmax' } | { type: 'set'; maxRows: number } @@ -376,6 +393,7 @@ export type ValidationIssueCode = | 'duplicate_object_name' | 'duplicate_column_name' | 'duplicate_index_name' + | 'text_index_invalid_parameters' | 'duplicate_projection_name' | 'projection_ambiguous_kind' | 'projection_empty_index' diff --git a/packages/core/src/planner.ts b/packages/core/src/planner.ts index 5689bd8f..be49b800 100644 --- a/packages/core/src/planner.ts +++ b/packages/core/src/planner.ts @@ -28,6 +28,7 @@ import { renderDictionarySQL, toCreateSQL, } from './sql.js' +import { textIndexFingerprint } from './text-index.js' import { assertValidDefinitions } from './validate.js' function createMap(definitions: SchemaDefinition[]): Map { @@ -414,7 +415,9 @@ function diffTables(oldDef: TableDefinition, newDef: TableDefinition): TableDiff oldDef.indexes ?? [], newDef.indexes ?? [], (index) => index.name, - (left, right) => JSON.stringify(left) === JSON.stringify(right) + (left, right) => left.type === 'text' && right.type === 'text' + ? textIndexFingerprint(left) === textIndexFingerprint(right) + : JSON.stringify(left) === JSON.stringify(right) ) for (const index of indexDiff.added) { ops.push( { diff --git a/packages/core/src/sql.ts b/packages/core/src/sql.ts index fe5e9fca..1247633c 100644 --- a/packages/core/src/sql.ts +++ b/packages/core/src/sql.ts @@ -13,6 +13,7 @@ import type { import { renderCodec } from './codec.js' import { isPlainColumnReference, normalizeKeyColumns } from './key-clause.js' import { renderProjectionBody } from './projection.js' +import { TEXT_INDEX_GRANULARITY, renderTextIndexType } from './text-index.js' import { assertValidDefinitions } from './validate.js' function renderDefault(value: string | number | boolean): string { @@ -44,6 +45,8 @@ function renderKeyClauseColumns(columns: string[], columnNames: Set): st function renderIndexType(idx: SkipIndexDefinition): string { switch (idx.type) { + case 'text': + return renderTextIndexType(idx) case 'minmax': return 'minmax' case 'set': @@ -63,7 +66,7 @@ function renderTableSQL(def: TableDefinition): string { const columns = def.columns.map(renderColumn) const indexes = (def.indexes ?? []).map( (idx) => - `INDEX \`${idx.name}\` (${idx.expression}) TYPE ${renderIndexType(idx)} GRANULARITY ${idx.granularity}` + `INDEX \`${idx.name}\` (${idx.expression}) TYPE ${renderIndexType(idx)} GRANULARITY ${idx.type === 'text' ? TEXT_INDEX_GRANULARITY : idx.granularity}` ) const projections = (def.projections ?? []).map( (projection) => `PROJECTION \`${projection.name}\` ${renderProjectionBody(projection)}` @@ -218,7 +221,7 @@ export function renderAlterRemoveCodec(def: TableDefinition, columnName: string) } export function renderAlterAddIndex(def: TableDefinition, index: SkipIndexDefinition): string { - return `ALTER TABLE ${def.database}.${def.name} ADD INDEX IF NOT EXISTS \`${index.name}\` (${index.expression}) TYPE ${renderIndexType(index)} GRANULARITY ${index.granularity};` + return `ALTER TABLE ${def.database}.${def.name} ADD INDEX IF NOT EXISTS \`${index.name}\` (${index.expression}) TYPE ${renderIndexType(index)} GRANULARITY ${index.type === 'text' ? TEXT_INDEX_GRANULARITY : index.granularity};` } export function renderAlterDropIndex(def: TableDefinition, indexName: string): string { diff --git a/packages/core/src/text-index-sql.ts b/packages/core/src/text-index-sql.ts new file mode 100644 index 00000000..f6f8cf3f --- /dev/null +++ b/packages/core/src/text-index-sql.ts @@ -0,0 +1,182 @@ +// A small lexer for text-index SQL fragments, not a SQL grammar. Keep literals +// opaque to whitespace/parenthesis handling and canonicalize their byte escapes. +const ESCAPES: Record = { + '0': 0, + a: 7, + b: 8, + t: 9, + n: 10, + v: 11, + f: 12, + r: 13, + e: 27, +} +const UTF8 = new TextEncoder() + +// Literal names and ambiguous keywords must remain identifiers when compared. +// Match ClickHouse's writeProbablyQuotedStringImpl (src/IO/WriteHelpers.cpp), +// including NULL, which its isValidIdentifier helper excludes separately. +const QUOTED_IDENTIFIERS = new Set([ + 'null', + 'true', + 'false', + 'inf', + 'infinity', + 'nan', + 'distinct', + 'all', + 'some', + 'table', + 'select', + 'from', + 'top', + 'values', +]) + +function stringLiteral(body: string): string { + const bytes: number[] = [] + for (let i = 0; i < body.length; ) { + const char = body[i] ?? '' + if (char === "'" && body[i + 1] === "'") { + bytes.push(39) + i += 2 + } else if (char === '\\') { + const next = String.fromCodePoint(body.codePointAt(i + 1) ?? 0) + const hex = body.slice(i + 2, i + 4) + if (next === 'x' && /^[\da-f]{2}$/i.test(hex)) { + bytes.push(Number.parseInt(hex, 16)) + i += 4 + } else { + // Unknown escapes (including regex \s and LIKE \%) keep their backslash. + if (next === 'x') throw new Error('Invalid hexadecimal SQL escape') + if (next === 'N') { + i += 2 + continue + } + if (ESCAPES[next] === undefined && !["'", '"', '`', '\\', '/', '='].includes(next)) + bytes.push(92) + bytes.push(...(ESCAPES[next] !== undefined ? [ESCAPES[next] ?? 0] : UTF8.encode(next))) + i += 1 + next.length + } + } else { + const point = String.fromCodePoint(body.codePointAt(i) ?? 0) + bytes.push(...UTF8.encode(point)) + i += point.length + } + } + return ( + "'" + + bytes + .map((byte) => { + if (byte === 39) return "\\'" + if (byte === 92) return '\\\\' + if (byte >= 32 && byte < 127) return String.fromCharCode(byte) + return `\\x${byte.toString(16).padStart(2, '0')}` + }) + .join('') + + "'" + ) +} + +export function textSQLTokens(sql: string): string[] { + const tokens: string[] = [] + const brackets: string[] = [] + for (let i = 0; i < sql.length; ) { + const char = sql[i] ?? '' + if (/\s/.test(char)) { + i++ + continue + } + if (sql.startsWith('--', i)) { + const end = sql.indexOf('\n', i + 2) + i = end < 0 ? sql.length : end + 1 + continue + } + if (sql.startsWith('/*', i)) { + const end = sql.indexOf('*/', i + 2) + if (end < 0) throw new Error('Unterminated SQL comment') + i = end + 2 + continue + } + if (char === "'" || char === '"' || char === '`') { + const start = i++ + let closed = false + while (i < sql.length) { + if (sql[i] === '\\') { + i += 2 + continue + } + if (sql[i++] === char) { + if (sql[i] === char) { + i++ + continue + } + closed = true + break + } + } + if (!closed) throw new Error('Unterminated SQL quote') + const raw = sql.slice(start, i) + tokens.push(char === "'" ? stringLiteral(raw.slice(1, -1)) : raw) + continue + } + if (char === ';') throw new Error('Expected a SQL expression, not a statement') + if ('([{'.includes(char)) brackets.push(char) + if (')]}'.includes(char) && brackets.pop() !== { ')': '(', ']': '[', '}': '{' }[char]) { + throw new Error('Unbalanced SQL brackets') + } + const word = sql + .slice(i) + .match(/^(?:\d+(?:\.\d*)?(?:[eE][+-]?\d+)?|[\p{L}\p{N}_$]+|->|<=|>=|!=|<>|\|\||::|==)/u)?.[0] + tokens.push(word ?? char) + i += word?.length ?? 1 + } + if (brackets.length) throw new Error('Unbalanced SQL brackets') + return tokens +} + +export function formatTextSQL(tokens: string[]): string { + let sql = '' + for (const [i, token] of tokens.entries()) { + const previous = tokens[i - 1] + const tight = + i === 0 || + [')', ']', '}', ',', '.'].includes(token) || + ['(', '[', '{', '.'].includes(previous ?? '') || + (token === '(' && /^[\p{L}_][\p{L}\p{N}_]*$/u.test(previous ?? '')) + sql += (tight ? '' : ' ') + token + } + return sql +} + +export function normalizeTextIndexSQL(sql: string): string { + return formatTextSQL(textSQLTokens(sql)) +} + +export function textExpressionFingerprint(sql: string): string { + let tokens = textSQLTokens(sql) + // ClickHouse may add or remove the outer INDEX expression parentheses. + while (tokens[0] === '(' && tokens.at(-1) === ')') { + let depth = 0 + const wraps = tokens.every((token, i) => { + if (token === '(') depth++ + if (token === ')') depth-- + return depth !== 0 || i === tokens.length - 1 + }) + if (!wraps) break + tokens = tokens.slice(1, -1) + } + return formatTextSQL(tokens) +} + +// Compare redundant identifier quotes without removing them from generated SQL. +export function textSQLFingerprint(sql: string): string { + return JSON.stringify( + textSQLTokens(sql).map((token) => { + const identifier = token.match(/^([`"])([A-Za-z_][A-Za-z0-9_]*)\1$/)?.[2] + if (identifier === undefined) return token + // Backticks and double quotes name the same identifier; preserve its case. + return QUOTED_IDENTIFIERS.has(identifier.toLowerCase()) ? `\`${identifier}\`` : identifier + }), + ) +} diff --git a/packages/core/src/text-index.test.ts b/packages/core/src/text-index.test.ts new file mode 100644 index 00000000..dd357b81 --- /dev/null +++ b/packages/core/src/text-index.test.ts @@ -0,0 +1,212 @@ +import { describe, expect, test } from 'bun:test' +import fixtures from '../../../test/fixtures/text-index.json' +import identifiers from '../../../test/fixtures/text-index-identifiers.json' +import { + canonicalizeDefinitions, + planDiff, + table, + toCreateSQL, + validateDefinitions, + type TextSkipIndex, +} from './index.js' +import { + canonicalizeTextIndex, + normalizeTextIndexSQL, + parseTextIndexParams, + renderTextIndexType, + textIndexFingerprint, +} from './text-index.js' + +const index = (params: Partial = {}): TextSkipIndex => ({ + name: 'idx', + expression: 'body', + type: 'text', + tokenizer: 'splitByNonAlpha', + ...params, +}) +const docs = (idx: TextSkipIndex) => + table({ + database: 'app', + name: 'docs', + columns: [ + { name: 'id', type: 'UInt64' }, + { name: 'body', type: 'String' }, + ], + engine: 'MergeTree()', + primaryKey: ['id'], + orderBy: ['id'], + indexes: [idx], + }) + +describe('text indexes', () => { + for (const { name, ...params } of fixtures) { + test(`lossless canonical round trip: ${name}`, () => { + const idx = index(params) + const sql = renderTextIndexType(idx) + expect(renderTextIndexType(parseTextIndexParams(sql.slice(5, -1)))).toBe(sql) + expect(canonicalizeTextIndex(canonicalizeTextIndex(idx))).toEqual(canonicalizeTextIndex(idx)) + }) + } + + test('preserves meaningful whitespace and detects separator/preprocessor changes', () => { + const original = docs( + index({ + tokenizer: "splitByString([' '])", + preprocessor: "replaceAll(body, ' ', ' ')", + }), + ) + expect(toCreateSQL(original)).toContain("splitByString([' '])") + for (const change of [ + { tokenizer: "splitByString([' '])" }, + { preprocessor: "replaceAll(body, ' ', ' ')" }, + ]) { + const changed = docs({ + ...original.indexes?.[0], + ...change, + } as TextSkipIndex) + expect(planDiff([original], [changed]).operations.map((op) => op.type)).toEqual([ + 'alter_table_drop_index', + 'alter_table_add_index', + ]) + } + }) + + test('granularity is automatic and never causes a migration', () => { + const original = docs(index()) + for (const granularity of [1, 64, 100000000]) { + const changed = docs(index({ granularity })) + expect(toCreateSQL(changed)).toContain('GRANULARITY 100000000') + expect(planDiff([original], [changed]).operations).toEqual([]) + expect(canonicalizeDefinitions([original])).toEqual(canonicalizeDefinitions([changed])) + } + }) + + test('fingerprint ignores formatting, parameter order and outer expression parentheses', () => { + const original = index({ + expression: "concat(body, ') (')", + tokenizer: "splitByString([' ', ','])", + dictionaryBlockSize: 512, + dictionaryBlockFrontcodingCompression: false, + postingListCodec: 'none', + }) + const parsed = parseTextIndexParams( + "posting_list_codec = 'none', dictionary_block_frontcoding_compression = FALSE, tokenizer = splitByString ( [ ' ' , ',' ] ), dictionary_block_size = 512", + ) + expect( + textIndexFingerprint({ + ...original, + ...parsed, + expression: "((concat(body, ') (')))", + }), + ).toBe(textIndexFingerprint(original)) + expect(textIndexFingerprint({ ...original, expression: "concat(body, ') (')" })).not.toBe( + textIndexFingerprint(original), + ) + }) + + test('equivalent string escape spellings compare equally', () => { + expect(normalizeTextIndexSQL("splitByString(['\t', 'é', ''''])")).toBe( + normalizeTextIndexSQL(String.raw`splitByString(['\x09', '\xc3\xa9', '\''])`), + ) + }) + + for (const args of [ + '', + 'tokenizer =', + 'tokenizer = ngrams(3', + "tokenizer = 'oops", + 'tokenizer = ngrams(3],)', + 'tokenizer = splitByNonAlpha,', + 'tokenizer = splitByNonAlpha, tokenizer = ngrams(2)', + 'tokenizer = splitByNonAlpha, unknown = 1', + 'tokenizer = splitByNonAlpha, support_phrase_search = 2', + 'tokenizer = splitByNonAlpha, dictionary_block_size = 0', + 'tokenizer = splitByNonAlpha, dictionary_block_size = 1.5', + 'tokenizer = splitByNonAlpha, dictionary_block_size = 9007199254740992', + "tokenizer = splitByNonAlpha, posting_list_codec = 'bogus'", + 'tokenizer = splitByNonAlpha; SELECT 1', + 'tokenizer = /* unclosed', + ]) { + test(`rejects lossy or malformed metadata: ${args}`, () => + expect(() => parseTextIndexParams(args)).toThrow()) + } + + for (const params of [ + { tokenizer: '' }, + { tokenizer: '/*empty*/' }, + { dictionaryBlockSize: -1 }, + { postingListBlockSize: Number.NaN }, + { preprocessor: '' }, + { supportPhraseSearch: 'false' }, + { postingListCodec: "none'" }, + ]) { + test(`validation reports invalid parameters: ${JSON.stringify(params)}`, () => { + expect( + validateDefinitions([docs(index(params as Partial))]).map( + (issue) => issue.code, + ), + ).toContain('text_index_invalid_parameters') + }) + } + + test('retains false options and all tuning fields', () => { + const idx = index({ + postprocessor: 'lower(body)', + supportPhraseSearch: false, + dictionaryBlockFrontcodingCompression: false, + dictionaryBlockSize: 512, + postingListBlockSize: 1024, + postingListCodec: 'bitpacking', + }) + expect(canonicalizeTextIndex(idx)).toMatchObject(idx) + expect(renderTextIndexType(idx)).toContain('support_phrase_search = 0') + }) +}) + +test('identifier quoting does not rebuild text indexes or conflate string literals', () => { + const quoted = docs(index({ preprocessor: 'lower(`body`)' })) + const unquoted = docs(index({ preprocessor: 'lower(body)' })) + expect(planDiff([quoted], [unquoted]).operations).toEqual([]) + expect( + planDiff([quoted], [docs(index({ preprocessor: "lower('body')" }))]).operations, + ).toHaveLength(2) +}) + +test('text index SQL generation and planning work without the Node Buffer global', () => { + const descriptor = Object.getOwnPropertyDescriptor(globalThis, 'Buffer') + try { + Reflect.deleteProperty(globalThis, 'Buffer') + // Include both ordinary and escaped Unicode, as well as equivalent byte escapes. + const original = docs(index({ tokenizer: "splitByString(['é', '😀', '\\é', '\\😀'])" })) + const equivalent = docs( + index({ + tokenizer: String.raw`splitByString(['\xc3\xa9', '\xf0\x9f\x98\x80', '\\\xc3\xa9', '\\\xf0\x9f\x98\x80'])`, + }), + ) + expect(toCreateSQL(original)).toBe(toCreateSQL(equivalent)) + expect(canonicalizeDefinitions([original])).toEqual(canonicalizeDefinitions([equivalent])) + expect(planDiff([original], [equivalent]).operations).toEqual([]) + expect( + planDiff([original], [docs(index({ tokenizer: "splitByString([' '])" }))]).operations, + ).toHaveLength(2) + } finally { + if (descriptor) Object.defineProperty(globalThis, 'Buffer', descriptor) + } +}) + +for (const word of identifiers) { + test(`preserves meaningful identifier quotes: ${word}`, () => { + for (const name of [word, word.toUpperCase(), word[0]?.toUpperCase() + word.slice(1)]) { + for (const field of ['expression', 'preprocessor', 'postprocessor'] as const) { + const quoted = docs(index({ [field]: `toString(\`${name}\`)` })) + const doubleQuoted = docs(index({ [field]: `toString("${name}")` })) + const unquoted = docs(index({ [field]: `toString(${name})` })) + expect(planDiff([quoted], [doubleQuoted]).operations).toEqual([]) + expect(planDiff([quoted], [unquoted]).operations.map((op) => op.type)).toEqual([ + 'alter_table_drop_index', + 'alter_table_add_index', + ]) + } + } + }) +} diff --git a/packages/core/src/text-index.ts b/packages/core/src/text-index.ts new file mode 100644 index 00000000..b8fde7cf --- /dev/null +++ b/packages/core/src/text-index.ts @@ -0,0 +1,119 @@ +import type { TextSkipIndex } from './model-types.js' +import { + formatTextSQL, + normalizeTextIndexSQL, + textExpressionFingerprint, + textSQLTokens, + textSQLFingerprint, +} from './text-index-sql.js' + +export { normalizeTextIndexSQL } from './text-index-sql.js' +export const TEXT_INDEX_GRANULARITY = 100_000_000 + +type Params = Omit +const PARAMETERS = [ + ['tokenizer', 'tokenizer', 'sql'], + ['preprocessor', 'preprocessor', 'sql'], + ['postprocessor', 'postprocessor', 'sql'], + ['supportPhraseSearch', 'support_phrase_search', 'boolean'], + ['dictionaryBlockSize', 'dictionary_block_size', 'number'], + ['dictionaryBlockFrontcodingCompression', 'dictionary_block_frontcoding_compression', 'boolean'], + ['postingListBlockSize', 'posting_list_block_size', 'number'], + ['postingListCodec', 'posting_list_codec', 'codec'], +] as const + +export function renderTextIndexType(index: Params): string { + if (typeof index.tokenizer !== 'string' || !textSQLTokens(index.tokenizer).length) { + throw new Error('A non-empty tokenizer is required') + } + const parts: string[] = [] + for (const [field, key, kind] of PARAMETERS) { + const value = index[field] + if (value === undefined) continue + let sql: string + if (kind === 'sql') { + if (typeof value !== 'string' || !textSQLTokens(value).length) + throw new Error(`${field} must be non-empty SQL`) + sql = normalizeTextIndexSQL(value) + } else if (kind === 'boolean') { + if (typeof value !== 'boolean') throw new Error(`${field} must be a boolean`) + sql = value ? '1' : '0' + } else if (kind === 'number') { + if (typeof value !== 'number' || !Number.isSafeInteger(value) || value <= 0) + throw new Error(`${field} must be a positive safe integer`) + sql = String(value) + } else { + if (value !== 'none' && value !== 'bitpacking') + throw new Error(`${field} must be none or bitpacking`) + sql = `'${value}'` + } + parts.push(`${key} = ${sql}`) + } + return `text(${parts.join(', ')})` +} + +/** Fail explicitly on unknown or malformed parameters rather than losing them on pull. */ +export function parseTextIndexParams(args: string): Params { + const tokens = textSQLTokens(args) + let group: string[] = [] + const groups: string[][] = [group] + let depth = 0 + for (const token of tokens) { + if (['(', '[', '{'].includes(token)) depth++ + if ([')', ']', '}'].includes(token)) depth-- + if (token === ',' && depth === 0) { + group = [] + groups.push(group) + } else group.push(token) + } + const params: Record = {} + for (const group of groups) { + const [key, equals, ...value] = group + const parameter = PARAMETERS.find((item) => item[1] === key) + if (!parameter || equals !== '=' || value.length === 0) + throw new Error(`Invalid or unsupported text index parameter: ${formatTextSQL(group)}`) + const [field, , kind] = parameter + if (field in params) throw new Error(`Duplicate text index parameter: ${key}`) + const raw = formatTextSQL(value) + if (kind === 'boolean') { + if (!/^(0|1|true|false)$/i.test(raw)) throw new Error(`Invalid boolean for ${key}: ${raw}`) + params[field] = raw === '1' || raw.toLowerCase() === 'true' + } else if (kind === 'number') { + if (!/^\d+$/.test(raw)) throw new Error(`Invalid integer for ${key}: ${raw}`) + params[field] = Number(raw) + } else if (kind === 'codec') { + params[field] = raw.slice(1, -1) + if (raw !== "'none'" && raw !== "'bitpacking'") throw new Error(`Invalid codec: ${raw}`) + } else params[field] = raw + } + const result = params as Params + renderTextIndexType(result) + return result +} + +export function canonicalizeTextIndex(index: TextSkipIndex): TextSkipIndex { + // Parse/render fixes key order as well as SQL formatting for migration snapshots. + const type = renderTextIndexType(index) + return { + name: index.name, + expression: textExpressionFingerprint(index.expression), + type: 'text', + granularity: TEXT_INDEX_GRANULARITY, + ...parseTextIndexParams(type.slice(5, -1)), + } +} + +export function textIndexFingerprint(index: TextSkipIndex): string { + const canonical = canonicalizeTextIndex(index) + return JSON.stringify({ + ...canonical, + expression: textSQLFingerprint(canonical.expression), + tokenizer: textSQLFingerprint(canonical.tokenizer), + preprocessor: + canonical.preprocessor === undefined ? undefined : textSQLFingerprint(canonical.preprocessor), + postprocessor: + canonical.postprocessor === undefined + ? undefined + : textSQLFingerprint(canonical.postprocessor), + }) +} diff --git a/packages/core/src/validate.ts b/packages/core/src/validate.ts index 89602348..c467fc56 100644 --- a/packages/core/src/validate.ts +++ b/packages/core/src/validate.ts @@ -1,3 +1,4 @@ +import { renderTextIndexType } from './text-index.js' import { definitionKey } from './canonical.js' import { canonicalizeCodec, isGeneralCodec, isRawCodec } from './codec.js' import { isPlainColumnReference, normalizeKeyColumns } from './key-clause.js' @@ -105,6 +106,14 @@ function validateTableDefinition(def: TableDefinition, issues: ValidationIssue[] continue } indexSeen.add(index.name) + if (index.type === 'text') { + try { + renderTextIndexType(index) + } catch (error) { + pushValidationIssue(issues, def, 'text_index_invalid_parameters', + `Text index "${index.name}": ${error instanceof Error ? error.message : String(error)}`) + } + } } const projectionSeen = new Set() diff --git a/packages/plugin-pull/src/render-schema.ts b/packages/plugin-pull/src/render-schema.ts index 486c066e..e8130ff7 100644 --- a/packages/plugin-pull/src/render-schema.ts +++ b/packages/plugin-pull/src/render-schema.ts @@ -212,6 +212,13 @@ function renderIndex(index: SkipIndexDefinition): string { `type: ${renderString(index.type)}`, ] switch (index.type) { + case 'text': + for (const field of ['tokenizer', 'preprocessor', 'postprocessor', 'supportPhraseSearch', + 'dictionaryBlockSize', 'dictionaryBlockFrontcodingCompression', 'postingListBlockSize', 'postingListCodec'] as const) { + const value = index[field] + if (value !== undefined) parts.push(`${field}: ${typeof value === 'string' ? renderString(value) : value}`) + } + break case 'minmax': break case 'set': @@ -238,7 +245,7 @@ function renderIndex(index: SkipIndexDefinition): string { ) break } - parts.push(`granularity: ${index.granularity}`) + if (index.type !== 'text') parts.push(`granularity: ${index.granularity}`) return `{ ${parts.join(', ')} }` } diff --git a/test/fixtures/text-index-identifiers.json b/test/fixtures/text-index-identifiers.json new file mode 100644 index 00000000..ad437187 --- /dev/null +++ b/test/fixtures/text-index-identifiers.json @@ -0,0 +1,16 @@ +[ + "null", + "true", + "false", + "inf", + "infinity", + "nan", + "distinct", + "all", + "some", + "table", + "select", + "from", + "top", + "values" +] diff --git a/test/fixtures/text-index.json b/test/fixtures/text-index.json new file mode 100644 index 00000000..5c89a168 --- /dev/null +++ b/test/fixtures/text-index.json @@ -0,0 +1,100 @@ +[ + { + "name": "default", + "tokenizer": "splitByNonAlpha" + }, + { + "name": "quoted name", + "tokenizer": "'splitByNonAlpha'" + }, + { + "name": "ngrams", + "tokenizer": "ngrams ( 3 )" + }, + { + "name": "sparse grams", + "tokenizer": "sparseGrams(3, 5, 4)" + }, + { + "name": "two spaces", + "tokenizer": "splitByString([' '])" + }, + { + "name": "actual tab and newline", + "tokenizer": "splitByString(['\t', '\n'])" + }, + { + "name": "escaped tab and newline", + "tokenizer": "splitByString(['\\t', '\\n'])" + }, + { + "name": "commas and equals", + "tokenizer": "splitByString([',', '=', '), tokenizer = bogus'])" + }, + { + "name": "backslash at end", + "tokenizer": "splitByString(['\\\\', ','])" + }, + { + "name": "escaped quote", + "tokenizer": "splitByString(['\\'', ';'])" + }, + { + "name": "doubled quote", + "tokenizer": "splitByString(['''', ';'])" + }, + { + "name": "unicode", + "tokenizer": "splitByString(['é', '東京', '😀'])" + }, + { + "name": "hex escapes", + "tokenizer": "splitByString(['\\x20\\x20', '\\xc3\\xa9'])" + }, + { + "name": "comment text in strings", + "tokenizer": "splitByString(['--', '/*', '*/', ';'])" + }, + { + "name": "SQL comments", + "tokenizer": "splitByString(/* commas , = */ [' ', -- keep line break\n ';'])" + }, + { + "name": "preprocessor whitespace", + "tokenizer": "splitByNonAlpha", + "preprocessor": "replaceAll(body, ' ', ' ')" + }, + { + "name": "preprocessor escaped literal", + "tokenizer": "splitByNonAlpha", + "preprocessor": "replaceAll(body, '\\\\', ' ')" + }, + { + "name": "parenthesis in expression literal", + "expression": "concat(body, ') (')", + "tokenizer": "splitByNonAlpha" + }, + { + "name": "unknown escapes and regex", + "tokenizer": "splitByNonAlpha", + "preprocessor": "replaceRegexpAll(body, '\\s+', ' ')" + }, + { + "name": "escaped unicode", + "tokenizer": "splitByString(['\\é', '\\😀'])" + }, + { + "name": "escape and empty sequence", + "tokenizer": "splitByString(['\\e', 'a\\Nb'])" + }, + { + "name": "quoted column expression", + "expression": "lower(`body`)", + "tokenizer": "splitByNonAlpha" + }, + { + "name": "quoted column preprocessor", + "tokenizer": "splitByNonAlpha", + "preprocessor": "lower(\"body\")" + } +]