diff --git a/sqlmesh/integrations/dlt.py b/sqlmesh/integrations/dlt.py index d9cced8deb..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, @@ -64,16 +71,23 @@ 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": + # 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 + 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,8 +223,62 @@ def generate_incremental_model( """ +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 + 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" + catalog_database = str(catalog.database or "") + storage_url = str(creds.storage_url or "") + lines = [ + " type: duckdb", + " catalogs:", + f" {_yaml_inline(alias)}:", + " type: ducklake", + 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: {_yaml_inline(str(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.""" + # 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 f1540727b1..3588d0a38a 100644 --- a/tests/cli/test_cli.py +++ b/tests/cli/test_cli.py @@ -1457,6 +1457,241 @@ def test_dlt_pipeline(runner, tmp_path): remove(dataset_path) +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=database) + credentials = types.SimpleNamespace( + ducklake_name=ducklake_name, + metadata_schema=metadata_schema, + catalog=catalog, + storage_url=storage_url, + ) + client_config = types.SimpleNamespace( + credentials=credentials, override_data_path=override_data_path + ) + 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" + ) + + # 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" + ) + + # 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") + 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_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") + + 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): create_example_project(tmp_path)