From 50b4f9a15bb16cc58dea5588a6911ac4dc983910 Mon Sep 17 00:00:00 2001 From: mday-io Date: Wed, 9 Sep 2026 14:11:08 -0400 Subject: [PATCH 1/2] fix(clickhouse): strip virtual catalogs from executed queries and inserts Signed-off-by: mday-io --- sqlmesh/core/engine_adapter/clickhouse.py | 10 +++ tests/core/engine_adapter/test_clickhouse.py | 69 ++++++++++++++++++++ 2 files changed, 79 insertions(+) diff --git a/sqlmesh/core/engine_adapter/clickhouse.py b/sqlmesh/core/engine_adapter/clickhouse.py index 92e3c33e6f..23dfb91ad2 100644 --- a/sqlmesh/core/engine_adapter/clickhouse.py +++ b/sqlmesh/core/engine_adapter/clickhouse.py @@ -61,6 +61,16 @@ def inject_virtual_catalog(self, gateway: str) -> None: configured = self._extra_config.get("virtual_catalog") self._default_catalog = f"__{gateway}__" if configured is None else configured + def _to_sql(self, expression: exp.Expr, quote: bool = True, **kwargs: t.Any) -> str: + """Render queries without the synthetic catalog unsupported by ClickHouse.""" + virtual_catalog = self._default_catalog or self._extra_config.get("virtual_catalog") + if virtual_catalog and isinstance(expression, (exp.Query, exp.Insert)): + expression = expression.copy() + for reference in expression.find_all(exp.Table, exp.Column): + if reference.text("catalog") == virtual_catalog: + reference.set("catalog", None) + return super()._to_sql(expression, quote=quote, **kwargs) + @property def engine_run_mode(self) -> EngineRunMode: if self._extra_config.get("cloud_mode"): diff --git a/tests/core/engine_adapter/test_clickhouse.py b/tests/core/engine_adapter/test_clickhouse.py index a3dfe0fdda..6ed552075f 100644 --- a/tests/core/engine_adapter/test_clickhouse.py +++ b/tests/core/engine_adapter/test_clickhouse.py @@ -1596,6 +1596,75 @@ def test_virtual_catalog_stripped_in_alter_table(make_mocked_engine_adapter: t.C assert "ALTER TABLE" in sql_calls[0] +@pytest.mark.parametrize( + "query_sql, expected_sql", + [ + ( + 'INSERT INTO __ch_gw__.mydb.target ("id") ' + "SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source", + 'INSERT INTO "mydb"."target" ("id") SELECT "mydb"."source"."id" FROM "mydb"."source"', + ), + ( + "SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source", + 'SELECT "mydb"."source"."id" FROM "mydb"."source"', + ), + ], +) +def test_virtual_catalog_stripped_from_execute_queries( + make_mocked_engine_adapter: t.Callable, query_sql: str, expected_sql: str +): + adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter) + adapter.inject_virtual_catalog("ch_gw") + query = parse_one(query_sql, dialect="clickhouse") + original_sql = query.sql(dialect="clickhouse") + + adapter.execute(query) + + assert query.sql(dialect="clickhouse") == original_sql + assert to_sql_calls(adapter) == [expected_sql] + + +def test_virtual_catalog_execute_preserves_unconfigured_catalog_and_literals( + make_mocked_engine_adapter: t.Callable, +): + adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter) + query = parse_one( + "SELECT other_catalog.mydb.source.id, '__ch_gw__.literal' FROM other_catalog.mydb.source", + dialect="clickhouse", + ) + + adapter.execute(query) + + assert to_sql_calls(adapter) == [ + 'SELECT "other_catalog"."mydb"."source"."id", \'__ch_gw__.literal\' ' + 'FROM "other_catalog"."mydb"."source"' + ] + + +def test_virtual_catalog_execute_uses_configured_catalog_fallback( + make_mocked_engine_adapter: t.Callable, +): + adapter = make_mocked_engine_adapter( + ClickhouseEngineAdapter, virtual_catalog="configured_catalog" + ) + query = parse_one( + "SELECT configured_catalog.mydb.source.id, other_catalog.otherdb.source.id, " + "'configured_catalog.literal' FROM configured_catalog.mydb.source " + "JOIN other_catalog.otherdb.source ON configured_catalog.mydb.source.id = " + "other_catalog.otherdb.source.id", + dialect="clickhouse", + ) + + adapter.execute(query) + + assert to_sql_calls(adapter) == [ + 'SELECT "mydb"."source"."id", "other_catalog"."otherdb"."source"."id", ' + '\'configured_catalog.literal\' FROM "mydb"."source" JOIN ' + '"other_catalog"."otherdb"."source" ON "mydb"."source"."id" = ' + '"other_catalog"."otherdb"."source"."id"' + ] + + def test_virtual_catalog_stripped_from_create_view_source( make_mocked_engine_adapter: t.Callable, ): From 74f820e474402f750018898cc39d73e8ee3bac2b Mon Sep 17 00:00:00 2001 From: mday-io Date: Fri, 25 Sep 2026 17:35:01 +0000 Subject: [PATCH 2/2] fix(clickhouse): strip virtual catalog from all executed statements Strip the injected virtual catalog from every rendered expression rather than only queries and inserts, so CTAS and DELETE statements no longer leak it. Drop the fallback to the configured virtual_catalog when none was injected, matching the other catalog strip sites, and only copy the expression when a reference actually needs rewriting. Consolidate the tests: the unconfigured-catalog test never enabled the feature, and the fallback test only covered the removed fallback. Signed-off-by: mday-io --- sqlmesh/core/engine_adapter/clickhouse.py | 17 +++--- tests/core/engine_adapter/test_clickhouse.py | 54 ++++++++------------ 2 files changed, 33 insertions(+), 38 deletions(-) diff --git a/sqlmesh/core/engine_adapter/clickhouse.py b/sqlmesh/core/engine_adapter/clickhouse.py index 23dfb91ad2..6bd14a54cc 100644 --- a/sqlmesh/core/engine_adapter/clickhouse.py +++ b/sqlmesh/core/engine_adapter/clickhouse.py @@ -62,15 +62,20 @@ def inject_virtual_catalog(self, gateway: str) -> None: self._default_catalog = f"__{gateway}__" if configured is None else configured def _to_sql(self, expression: exp.Expr, quote: bool = True, **kwargs: t.Any) -> str: - """Render queries without the synthetic catalog unsupported by ClickHouse.""" - virtual_catalog = self._default_catalog or self._extra_config.get("virtual_catalog") - if virtual_catalog and isinstance(expression, (exp.Query, exp.Insert)): + """Render SQL without the virtual catalog, which ClickHouse does not support.""" + if self._default_catalog and any(self._virtual_catalog_references(expression)): expression = expression.copy() - for reference in expression.find_all(exp.Table, exp.Column): - if reference.text("catalog") == virtual_catalog: - reference.set("catalog", None) + for reference in list(self._virtual_catalog_references(expression)): + reference.set("catalog", None) return super()._to_sql(expression, quote=quote, **kwargs) + def _virtual_catalog_references(self, expression: exp.Expr) -> t.Iterator[exp.Expr]: + return ( + reference + for reference in expression.find_all(exp.Table, exp.Column) + if reference.text("catalog") == self._default_catalog + ) + @property def engine_run_mode(self) -> EngineRunMode: if self._extra_config.get("cloud_mode"): diff --git a/tests/core/engine_adapter/test_clickhouse.py b/tests/core/engine_adapter/test_clickhouse.py index 6ed552075f..5af04304be 100644 --- a/tests/core/engine_adapter/test_clickhouse.py +++ b/tests/core/engine_adapter/test_clickhouse.py @@ -1608,6 +1608,14 @@ def test_virtual_catalog_stripped_in_alter_table(make_mocked_engine_adapter: t.C "SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source", 'SELECT "mydb"."source"."id" FROM "mydb"."source"', ), + ( + "SELECT __ch_gw__.mydb.source.id, '__ch_gw__.literal' FROM __ch_gw__.mydb.source " + "JOIN other_catalog.otherdb.source ON __ch_gw__.mydb.source.id = " + "other_catalog.otherdb.source.id", + 'SELECT "mydb"."source"."id", \'__ch_gw__.literal\' FROM "mydb"."source" JOIN ' + '"other_catalog"."otherdb"."source" ON "mydb"."source"."id" = ' + '"other_catalog"."otherdb"."source"."id"', + ), ], ) def test_virtual_catalog_stripped_from_execute_queries( @@ -1624,44 +1632,26 @@ def test_virtual_catalog_stripped_from_execute_queries( assert to_sql_calls(adapter) == [expected_sql] -def test_virtual_catalog_execute_preserves_unconfigured_catalog_and_literals( - make_mocked_engine_adapter: t.Callable, -): +def test_virtual_catalog_stripped_from_ctas_and_delete(make_mocked_engine_adapter: t.Callable): adapter = make_mocked_engine_adapter(ClickhouseEngineAdapter) - query = parse_one( - "SELECT other_catalog.mydb.source.id, '__ch_gw__.literal' FROM other_catalog.mydb.source", - dialect="clickhouse", - ) - - adapter.execute(query) - - assert to_sql_calls(adapter) == [ - 'SELECT "other_catalog"."mydb"."source"."id", \'__ch_gw__.literal\' ' - 'FROM "other_catalog"."mydb"."source"' - ] - + adapter.inject_virtual_catalog("ch_gw") -def test_virtual_catalog_execute_uses_configured_catalog_fallback( - make_mocked_engine_adapter: t.Callable, -): - adapter = make_mocked_engine_adapter( - ClickhouseEngineAdapter, virtual_catalog="configured_catalog" + adapter.ctas( + "__ch_gw__.mydb.target", + parse_one("SELECT __ch_gw__.mydb.source.id FROM __ch_gw__.mydb.source"), + {"id": exp.DataType.build("Int32")}, ) - query = parse_one( - "SELECT configured_catalog.mydb.source.id, other_catalog.otherdb.source.id, " - "'configured_catalog.literal' FROM configured_catalog.mydb.source " - "JOIN other_catalog.otherdb.source ON configured_catalog.mydb.source.id = " - "other_catalog.otherdb.source.id", - dialect="clickhouse", + adapter.delete_from( + "__ch_gw__.mydb.target", + "__ch_gw__.mydb.target.id IN (SELECT id FROM __ch_gw__.mydb.source)", ) - adapter.execute(query) - assert to_sql_calls(adapter) == [ - 'SELECT "mydb"."source"."id", "other_catalog"."otherdb"."source"."id", ' - '\'configured_catalog.literal\' FROM "mydb"."source" JOIN ' - '"other_catalog"."otherdb"."source" ON "mydb"."source"."id" = ' - '"other_catalog"."otherdb"."source"."id"' + 'CREATE TABLE IF NOT EXISTS "mydb"."target" ENGINE=MergeTree ORDER BY () AS ' + 'SELECT CAST("id" AS Nullable(Int32)) AS "id" FROM ' + '(SELECT "mydb"."source"."id" FROM "mydb"."source") AS "_subquery"', + 'DELETE FROM "mydb"."target" WHERE "mydb"."target"."id" IN ' + '(SELECT "id" FROM "mydb"."source")', ]