Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
88 changes: 78 additions & 10 deletions sqlmesh/integrations/dlt.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,6 @@
from __future__ import annotations

import json
import typing as t
import click
from datetime import datetime, timedelta, timezone
Expand All @@ -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,
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 '<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,
}
Expand Down
235 changes: 235 additions & 0 deletions tests/cli/test_cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down