From 7906e8bd88c8adfabfa7e79a97d4bd40c4f3ba45 Mon Sep 17 00:00:00 2001 From: Syed Annas Date: Sun, 4 Oct 2026 19:36:02 +0100 Subject: [PATCH 1/2] fix: generate a valid duckdb gateway config for DuckLake pipelines MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A `ducklake` entry in config.yaml fell through the generic credentials path and produced an invalid duckdb gateway config, which surfaced as a ConfigError from connection.py with a message identical to an unrelated failure — so the real diagnostic was lost. DuckLake connections are a duckdb gateway with the lake attached as a catalog, so build that block explicitly. Unsupported catalog drivers are cut to a clear ClickException rather than propagating the generic error; postgres/mysql/MotherDuck catalogs are not yet mapped and are tracked on the issue. Tests: 3 new (duckdb catalog, sqlite catalog, unsupported driver) plus the existing dlt pipeline tests. Verified to fail on unpatched source and pass on this commit. Signed-off-by: Syed Annas --- sqlmesh/integrations/dlt.py | 50 ++++++++++++++++----- tests/cli/test_cli.py | 86 +++++++++++++++++++++++++++++++++++++ 2 files changed, 126 insertions(+), 10 deletions(-) diff --git a/sqlmesh/integrations/dlt.py b/sqlmesh/integrations/dlt.py index d9cced8deb..52b9fdfb7d 100644 --- a/sqlmesh/integrations/dlt.py +++ b/sqlmesh/integrations/dlt.py @@ -64,16 +64,19 @@ def generate_dlt_models_and_settings( connection_config = None else: client = pipeline.destination_client() - config = client.config - credentials = config.credentials - configs = { - key: value - for key in dir(credentials) - if not key.startswith("_") - and not callable(value := getattr(credentials, key)) - and value is not None - } - connection_config = format_config(configs, db_type) + if db_type == "ducklake": + connection_config = format_ducklake_config(client.config) + else: + config = client.config + credentials = config.credentials + configs = { + key: value + for key in dir(credentials) + if not key.startswith("_") + and not callable(value := getattr(credentials, key)) + and value is not None + } + connection_config = format_config(configs, db_type) dlt_tables = { name: table @@ -209,6 +212,33 @@ def generate_incremental_model( """ +def format_ducklake_config(client_config: t.Any) -> str: + """Generate a duckdb-gateway connection block with the DuckLake attached as catalog.""" + creds = client_config.credentials + catalog = creds.catalog + drivername = getattr(catalog, "drivername", "") or "" + if drivername not in ("duckdb", "sqlite"): + raise click.ClickException( + f"Unsupported DuckLake catalog '{drivername}'. SQLMesh dlt init currently supports " + "file-backed catalogs (duckdb, sqlite); postgres/mysql/MotherDuck catalogs are not " + "yet mapped. Tracked in SQLMesh/sqlmesh#5914." + ) + alias = creds.ducklake_name or "ducklake" + lines = [ + " type: duckdb", + " catalogs:", + f" {alias}:", + " type: ducklake", + f" path: {catalog.database}", + f" data_path: {creds.storage_url}", + ] + metadata_schema = creds.metadata_schema or alias + lines.append(f" metadata_schema: {metadata_schema}") + if getattr(client_config, "override_data_path", False): + lines.append(" override_data_path: true") + return "\n".join(lines) + + def format_config(configs: t.Dict[str, str], db_type: str) -> str: """Generate a string for the gateway connection config.""" config = { diff --git a/tests/cli/test_cli.py b/tests/cli/test_cli.py index f1540727b1..d43db7217a 100644 --- a/tests/cli/test_cli.py +++ b/tests/cli/test_cli.py @@ -1457,6 +1457,92 @@ def test_dlt_pipeline(runner, tmp_path): remove(dataset_path) +def _stub_ducklake_dlt(monkeypatch, drivername="sqlite"): + """Install a fake `dlt` module exposing a ducklake pipeline. No dlt install needed.""" + import sys + import types + + catalog = types.SimpleNamespace(drivername=drivername, database="/tmp/x/mre_ducklake.sqlite") + credentials = types.SimpleNamespace( + ducklake_name="mre_ducklake", + metadata_schema=None, + catalog=catalog, + storage_url="/tmp/x/mre_ducklake.files", + ) + client_config = types.SimpleNamespace(credentials=credentials, override_data_path=False) + pipeline = types.SimpleNamespace( + destination=types.SimpleNamespace(to_name=lambda dest: "ducklake"), + default_schema=types.SimpleNamespace( + tables={}, _dlt_tables_prefix="_dlt", loads_table_name="_dlt_loads" + ), + dataset_name="mre", + ) + pipeline._get_load_storage = lambda: types.SimpleNamespace(list_loaded_packages=lambda: []) + pipeline.destination_client = lambda: types.SimpleNamespace(config=client_config) + + dlt_fake = types.ModuleType("dlt") + dlt_fake.attach = lambda pipeline_name, pipelines_dir="": pipeline + + schema_utils = types.ModuleType("dlt.common.schema.utils") + schema_utils.has_table_seen_data = lambda table: True + schema_utils.is_complete_column = lambda col: True + + pipeline_exceptions = types.ModuleType("dlt.pipeline.exceptions") + pipeline_exceptions.CannotRestorePipelineException = type( + "CannotRestorePipelineException", (Exception,), {} + ) + + for name, module in { + "dlt": dlt_fake, + "dlt.common": types.ModuleType("dlt.common"), + "dlt.common.schema": types.ModuleType("dlt.common.schema"), + "dlt.common.schema.utils": schema_utils, + "dlt.pipeline": types.ModuleType("dlt.pipeline"), + "dlt.pipeline.exceptions": pipeline_exceptions, + }.items(): + monkeypatch.setitem(sys.modules, name, module) + + +@pytest.mark.parametrize( + "drivername", ["sqlite", "duckdb"], ids=["sqlite-catalog", "duckdb-catalog"] +) +def test_dlt_ducklake_pipeline(monkeypatch, drivername): + import yaml + + from sqlmesh.core.config.connection import DuckDBConnectionConfig, parse_connection_config + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, drivername=drivername) + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + + assert "type: duckdb" in connection_config + assert "type: ducklake" in connection_config + assert "path: /tmp/x/mre_ducklake.sqlite" in connection_config + assert "data_path: /tmp/x/mre_ducklake.files" in connection_config + + # Round-trip: the exact ConfigError from #5914 no longer fires + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + config = parse_connection_config(parsed) + assert isinstance(config, DuckDBConnectionConfig) + attach_sql = next(iter(config.catalogs.values())).to_sql("mre_ducklake") + assert attach_sql == ( + "ATTACH IF NOT EXISTS 'ducklake:/tmp/x/mre_ducklake.sqlite' AS mre_ducklake " + "(DATA_PATH '/tmp/x/mre_ducklake.files', METADATA_SCHEMA 'mre_ducklake')" + ) + + +def test_dlt_ducklake_unsupported_catalog(monkeypatch): + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, drivername="postgres") + + with pytest.raises(ClickException, match="Unsupported DuckLake catalog"): + dlt_module.generate_dlt_models_and_settings(pipeline_name="mre_ducklake", dialect="duckdb") + + @time_machine.travel(FREEZE_TIME) def test_environments(runner, tmp_path): create_example_project(tmp_path) From 7518b13f8de63549f44abb36aa9fd179d99a6c8a Mon Sep 17 00:00:00 2001 From: Syed Annas Date: Sun, 4 Oct 2026 19:50:43 +0100 Subject: [PATCH 2/2] fix: quote ducklake YAML scalars, harden tests and typing Signed-off-by: Syed Annas --- sqlmesh/integrations/dlt.py | 50 +++++++++-- tests/cli/test_cli.py | 173 +++++++++++++++++++++++++++++++++--- 2 files changed, 205 insertions(+), 18 deletions(-) diff --git a/sqlmesh/integrations/dlt.py b/sqlmesh/integrations/dlt.py index 52b9fdfb7d..c321e3e182 100644 --- a/sqlmesh/integrations/dlt.py +++ b/sqlmesh/integrations/dlt.py @@ -1,3 +1,6 @@ +from __future__ import annotations + +import json import typing as t import click from datetime import datetime, timedelta, timezone @@ -8,6 +11,10 @@ from sqlmesh.utils.date import yesterday_ds +if t.TYPE_CHECKING: + from dlt.destinations.impl.ducklake.configuration import DuckLakeClientConfiguration + + def generate_dlt_models_and_settings( pipeline_name: str, dialect: str, @@ -65,7 +72,11 @@ def generate_dlt_models_and_settings( else: client = pipeline.destination_client() if db_type == "ducklake": - connection_config = format_ducklake_config(client.config) + # Cast: reachable only for ducklake pipelines, so client.config is the + # DuckLake client configuration at runtime (statically the base type). + connection_config = format_ducklake_config( + t.cast("DuckLakeClientConfiguration", client.config) + ) else: config = client.config credentials = config.credentials @@ -212,7 +223,25 @@ def generate_incremental_model( """ -def format_ducklake_config(client_config: t.Any) -> str: +def _yaml_inline(value: str) -> str: + """Emit one YAML scalar inline, quoting only when plain would not round-trip. + + Ordinary locators (alphanumerics plus / _ . -) are returned unchanged so the + generated config stays byte-identical; anything else (leading quote, newline, + ': ', ' #', spaces, etc.) is double-quoted via JSON (single line, valid YAML) + so yaml.safe_load round-trips instead of raising ScannerError. No PyYAML + dependency: json double-quotes are valid YAML double-quotes. + """ + if ( + value + and (value[0].isalnum() or value[0] in "/_") + and all(ch.isalnum() or ch in "_./-" for ch in value) + ): + return value + return json.dumps(value) + + +def format_ducklake_config(client_config: DuckLakeClientConfiguration) -> str: """Generate a duckdb-gateway connection block with the DuckLake attached as catalog.""" creds = client_config.credentials catalog = creds.catalog @@ -224,16 +253,18 @@ def format_ducklake_config(client_config: t.Any) -> str: "yet mapped. Tracked in SQLMesh/sqlmesh#5914." ) alias = creds.ducklake_name or "ducklake" + catalog_database = str(catalog.database or "") + storage_url = str(creds.storage_url or "") lines = [ " type: duckdb", " catalogs:", - f" {alias}:", + f" {_yaml_inline(alias)}:", " type: ducklake", - f" path: {catalog.database}", - f" data_path: {creds.storage_url}", + f" path: {_yaml_inline(catalog_database)}", + f" data_path: {_yaml_inline(storage_url)}", ] metadata_schema = creds.metadata_schema or alias - lines.append(f" metadata_schema: {metadata_schema}") + lines.append(f" metadata_schema: {_yaml_inline(str(metadata_schema))}") if getattr(client_config, "override_data_path", False): lines.append(" override_data_path: true") return "\n".join(lines) @@ -241,6 +272,13 @@ def format_ducklake_config(client_config: t.Any) -> str: def format_config(configs: t.Dict[str, str], db_type: str) -> str: """Generate a string for the gateway connection config.""" + # NOTE (SQLMesh#5914 scope cut): only the `ducklake` destination is mapped + # (see format_ducklake_config). Any other unrecognised dlt `db_type` + # (e.g. weaviate, pandas, qdrant, typos) still falls through to + # parse_connection_config below and surfaces as + # ConfigError("Unknown connection type ''."). That is a known + # limitation, not a regression introduced here; #5914 reports only the + # ducklake destination ("When using dlt with a `ducklake` destination ..."). config = { "type": db_type, } diff --git a/tests/cli/test_cli.py b/tests/cli/test_cli.py index d43db7217a..3588d0a38a 100644 --- a/tests/cli/test_cli.py +++ b/tests/cli/test_cli.py @@ -1457,19 +1457,29 @@ def test_dlt_pipeline(runner, tmp_path): remove(dataset_path) -def _stub_ducklake_dlt(monkeypatch, drivername="sqlite"): +def _stub_ducklake_dlt( + monkeypatch, + drivername="sqlite", + ducklake_name="mre_ducklake", + metadata_schema=None, + override_data_path=False, + database="/tmp/x/mre_ducklake.sqlite", + storage_url="/tmp/x/mre_ducklake.files", +): """Install a fake `dlt` module exposing a ducklake pipeline. No dlt install needed.""" import sys import types - catalog = types.SimpleNamespace(drivername=drivername, database="/tmp/x/mre_ducklake.sqlite") + catalog = types.SimpleNamespace(drivername=drivername, database=database) credentials = types.SimpleNamespace( - ducklake_name="mre_ducklake", - metadata_schema=None, + ducklake_name=ducklake_name, + metadata_schema=metadata_schema, catalog=catalog, - storage_url="/tmp/x/mre_ducklake.files", + storage_url=storage_url, + ) + client_config = types.SimpleNamespace( + credentials=credentials, override_data_path=override_data_path ) - client_config = types.SimpleNamespace(credentials=credentials, override_data_path=False) pipeline = types.SimpleNamespace( destination=types.SimpleNamespace(to_name=lambda dest: "ducklake"), default_schema=types.SimpleNamespace( @@ -1518,13 +1528,30 @@ def test_dlt_ducklake_pipeline(monkeypatch, drivername): pipeline_name="mre_ducklake", dialect="duckdb" ) - assert "type: duckdb" in connection_config - assert "type: ducklake" in connection_config - assert "path: /tmp/x/mre_ducklake.sqlite" in connection_config - assert "data_path: /tmp/x/mre_ducklake.files" in connection_config + # Byte-identical for ordinary paths (no quoting, key order and indent unchanged) + assert connection_config == ( + " type: duckdb\n" + " catalogs:\n" + " mre_ducklake:\n" + " type: ducklake\n" + " path: /tmp/x/mre_ducklake.sqlite\n" + " data_path: /tmp/x/mre_ducklake.files\n" + " metadata_schema: mre_ducklake" + ) - # Round-trip: the exact ConfigError from #5914 no longer fires + # Structural parse instead of substring checks: malformed YAML cannot pass parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + assert parsed["type"] == "duckdb" + assert set(parsed["catalogs"]) == {"mre_ducklake"} + lake = parsed["catalogs"]["mre_ducklake"] + assert lake == { + "type": "ducklake", + "path": "/tmp/x/mre_ducklake.sqlite", + "data_path": "/tmp/x/mre_ducklake.files", + "metadata_schema": "mre_ducklake", + } + + # Round-trip: the exact ConfigError from #5914 no longer fires config = parse_connection_config(parsed) assert isinstance(config, DuckDBConnectionConfig) attach_sql = next(iter(config.catalogs.values())).to_sql("mre_ducklake") @@ -1534,14 +1561,136 @@ def test_dlt_ducklake_pipeline(monkeypatch, drivername): ) +def test_dlt_ducklake_explicit_metadata_schema(monkeypatch): + import yaml + + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, metadata_schema="custom_meta") + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + lake = parsed["catalogs"]["mre_ducklake"] + assert lake["metadata_schema"] == "custom_meta" + assert lake["path"] == "/tmp/x/mre_ducklake.sqlite" + + +def test_dlt_ducklake_override_data_path(monkeypatch): + import yaml + + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, override_data_path=True) + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + lake = parsed["catalogs"]["mre_ducklake"] + assert lake["override_data_path"] is True + + +def test_dlt_ducklake_custom_name(monkeypatch): + import yaml + + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, ducklake_name="my_lake") + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + assert set(parsed["catalogs"]) == {"my_lake"} + lake = parsed["catalogs"]["my_lake"] + assert lake["metadata_schema"] == "my_lake" + + +def test_dlt_ducklake_yaml_inline_helper(): + from sqlmesh.integrations.dlt import _yaml_inline + + # Ordinary values stay byte-identical (no quotes) + assert _yaml_inline("/tmp/x/mre_ducklake.sqlite") == "/tmp/x/mre_ducklake.sqlite" + assert _yaml_inline("mre_ducklake") == "mre_ducklake" + # Pathological values are quoted single-line and round-trip + import yaml + + for pathological in ( + "'/tmp/quote/mre_ducklake.sqlite", + "/tmp/new\nline/mre.sqlite", + "a: b # c", + " leading-space", + ): + emitted = _yaml_inline(pathological) + assert "\n" not in emitted + doc = f"connection:\n path: {emitted}\n" + assert yaml.safe_load(doc)["connection"]["path"] == pathological + + +@pytest.mark.parametrize( + "pathological", + ["'/tmp/quote/mre_ducklake.sqlite", "/tmp/new\nline/mre.sqlite"], + ids=["leading-quote", "newline"], +) +def test_dlt_ducklake_pathological_paths_round_trip(monkeypatch, pathological): + import yaml + + from sqlmesh.core.config.connection import DuckDBConnectionConfig, parse_connection_config + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch, database=pathological) + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + # Must not raise ScannerError; values must round-trip exactly + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + assert parsed["catalogs"]["mre_ducklake"]["path"] == pathological + config = parse_connection_config(parsed) + assert isinstance(config, DuckDBConnectionConfig) + + +def test_dlt_ducklake_block_coexists_with_second_catalog(monkeypatch): + import yaml + + from sqlmesh.core.config.connection import DuckDBConnectionConfig, parse_connection_config + from sqlmesh.integrations import dlt as dlt_module + + _stub_ducklake_dlt(monkeypatch) + + _, connection_config, _ = dlt_module.generate_dlt_models_and_settings( + pipeline_name="mre_ducklake", dialect="duckdb" + ) + parsed = yaml.safe_load("connection:\n" + connection_config)["connection"] + parsed["catalogs"]["other"] = {"type": "ducklake", "path": "/tmp/x/other.sqlite"} + config = parse_connection_config(parsed) + assert isinstance(config, DuckDBConnectionConfig) + assert set(config.catalogs) == {"mre_ducklake", "other"} + + def test_dlt_ducklake_unsupported_catalog(monkeypatch): from sqlmesh.integrations import dlt as dlt_module _stub_ducklake_dlt(monkeypatch, drivername="postgres") - with pytest.raises(ClickException, match="Unsupported DuckLake catalog"): + called = {} + + orig = dlt_module.format_ducklake_config + + def _spy(client_config): + called["branch"] = True + return orig(client_config) + + monkeypatch.setattr(dlt_module, "format_ducklake_config", _spy) + + with pytest.raises(ClickException, match="Unsupported DuckLake catalog 'postgres'") as exc_info: dlt_module.generate_dlt_models_and_settings(pipeline_name="mre_ducklake", dialect="duckdb") + assert called.get("branch") is True + assert "postgres" in str(exc_info.value) + @time_machine.travel(FREEZE_TIME) def test_environments(runner, tmp_path):