diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 1592bd9..c2aedb1 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -66,13 +66,17 @@ jobs: - name: Install SemaPact Databricks extra from project metadata run: uv pip install --python .venv/bin/python -e ".[databricks]" pytest - - name: Verify Databricks SDK, client factory, and existing Unity import boundaries + - name: Verify Databricks SDK read-side boundaries run: >- .venv/bin/python -c - "import inspect; from databricks.sdk import WorkspaceClient; from databricks.sdk.service.catalog import TableInfo; from semapact.importers.unity_importer import import_unity_contract; from semapact.platforms.databricks import create_databricks_workspace_client; params = inspect.signature(WorkspaceClient).parameters; assert {'host', 'token', 'profile'} <= set(params); print(WorkspaceClient.__name__, TableInfo.__name__, import_unity_contract.__name__, create_databricks_workspace_client.__name__)" + "import inspect; from databricks.sdk import WorkspaceClient; from databricks.sdk.service.catalog import TableInfo, TablesAPI; from semapact.platforms.databricks import create_databricks_workspace_client, discover_databricks_tables; client_params = inspect.signature(WorkspaceClient).parameters; list_params = inspect.signature(TablesAPI.list).parameters; assert {'host', 'token', 'profile'} <= set(client_params); assert {'catalog_name', 'schema_name'} <= set(list_params); print(WorkspaceClient.__name__, TableInfo.__name__, create_databricks_workspace_client.__name__, discover_databricks_tables.__name__)" - - name: Run Databricks platform-boundary tests with official SDK installed - run: .venv/bin/python -m pytest tests/test_databricks_client.py tests/test_observation_databricks.py + - name: Run Databricks read-side tests with official SDK installed + run: >- + .venv/bin/python -m pytest + tests/test_databricks_client.py + tests/test_databricks_discovery.py + tests/test_observation_databricks.py coverage: runs-on: ubuntu-latest diff --git a/semapact/platforms/databricks/__init__.py b/semapact/platforms/databricks/__init__.py index 2df926e..9a44ddb 100644 --- a/semapact/platforms/databricks/__init__.py +++ b/semapact/platforms/databricks/__init__.py @@ -1,5 +1,9 @@ """Databricks platform-access helpers.""" from semapact.platforms.databricks.client import create_databricks_workspace_client +from semapact.platforms.databricks.discovery import discover_databricks_tables -__all__ = ["create_databricks_workspace_client"] +__all__ = [ + "create_databricks_workspace_client", + "discover_databricks_tables", +] diff --git a/semapact/platforms/databricks/discovery.py b/semapact/platforms/databricks/discovery.py new file mode 100644 index 0000000..49de3ae --- /dev/null +++ b/semapact/platforms/databricks/discovery.py @@ -0,0 +1,69 @@ +"""Read-only Databricks table discovery. + +Discovery answers which assets exist in a Unity Catalog schema. It does not +observe table structure, create an ODCS contract, evaluate governance, or +resolve authentication. Callers supply an initialized ``WorkspaceClient``. +""" + +from __future__ import annotations + +from collections.abc import Iterable +from typing import TYPE_CHECKING, Protocol + +if TYPE_CHECKING: + from databricks.sdk import WorkspaceClient + + +class _TableInfoLike(Protocol): + full_name: str | None + name: str | None + + +class _TablesApiLike(Protocol): + def list( + self, + *, + catalog_name: str, + schema_name: str, + ) -> Iterable[_TableInfoLike]: ... + + +class _WorkspaceClientLike(Protocol): + tables: _TablesApiLike + + +def discover_databricks_tables( + *, + client: WorkspaceClient | _WorkspaceClientLike, + catalog_name: str, + schema_name: str, +) -> tuple[str, ...]: + """Return stable fully qualified names for tables visible in one UC schema.""" + catalog = _required_name(catalog_name, field="catalog_name") + schema = _required_name(schema_name, field="schema_name") + + discovered: set[str] = set() + for table in client.tables.list(catalog_name=catalog, schema_name=schema): + full_name = _text(table.full_name) + if not full_name: + name = _text(table.name) + if not name: + raise ValueError("Databricks discovery returned a table without an identity") + full_name = f"{catalog}.{schema}.{name}" + discovered.add(full_name) + + return tuple(sorted(discovered, key=lambda value: (value.casefold(), value))) + + +def _required_name(value: str, *, field: str) -> str: + cleaned = _text(value) + if not cleaned: + raise ValueError(f"{field} is required for Databricks discovery") + return cleaned + + +def _text(value: object | None) -> str | None: + if value is None: + return None + text = str(value).strip() + return text or None diff --git a/tests/test_databricks_discovery.py b/tests/test_databricks_discovery.py new file mode 100644 index 0000000..6f43d2d --- /dev/null +++ b/tests/test_databricks_discovery.py @@ -0,0 +1,120 @@ +from __future__ import annotations + +from dataclasses import dataclass + +import pytest + +from semapact.platforms.databricks.discovery import discover_databricks_tables + + +@dataclass +class _FakeTable: + full_name: str | None + name: str | None = None + + +class _FakeTablesApi: + def __init__(self, tables: list[_FakeTable]) -> None: + self._tables = tables + self.calls: list[dict[str, str]] = [] + + def list(self, *, catalog_name: str, schema_name: str) -> list[_FakeTable]: + self.calls.append( + { + "catalog_name": catalog_name, + "schema_name": schema_name, + } + ) + return self._tables + + +class _FakeWorkspaceClient: + def __init__(self, tables: list[_FakeTable]) -> None: + self.tables = _FakeTablesApi(tables) + + +def test_discover_databricks_tables_lists_without_observing_or_importing() -> None: + client = _FakeWorkspaceClient( + [ + _FakeTable("main.sales.Z_orders"), + _FakeTable("main.sales.accounts"), + _FakeTable("main.sales.accounts"), + _FakeTable(None, name="customers"), + ] + ) + + discovered = discover_databricks_tables( + client=client, + catalog_name=" main ", + schema_name=" sales ", + ) + + assert discovered == ( + "main.sales.accounts", + "main.sales.customers", + "main.sales.Z_orders", + ) + assert client.tables.calls == [ + { + "catalog_name": "main", + "schema_name": "sales", + } + ] + + +def test_discover_databricks_tables_fails_on_missing_table_identity() -> None: + client = _FakeWorkspaceClient([_FakeTable(None)]) + + with pytest.raises( + ValueError, + match="Databricks discovery returned a table without an identity", + ): + discover_databricks_tables( + client=client, + catalog_name="main", + schema_name="sales", + ) + + +@pytest.mark.parametrize( + ("catalog_name", "schema_name", "message"), + [ + ("", "sales", "catalog_name is required for Databricks discovery"), + (" ", "sales", "catalog_name is required for Databricks discovery"), + ("main", "", "schema_name is required for Databricks discovery"), + ("main", " ", "schema_name is required for Databricks discovery"), + ], +) +def test_discover_databricks_tables_validates_scope_before_sdk_call( + catalog_name: str, + schema_name: str, + message: str, +) -> None: + client = _FakeWorkspaceClient([]) + + with pytest.raises(ValueError, match=message): + discover_databricks_tables( + client=client, + catalog_name=catalog_name, + schema_name=schema_name, + ) + + assert client.tables.calls == [] + + +def test_databricks_discovery_module_has_no_contract_or_governance_dependency() -> None: + from pathlib import Path + + source = Path("semapact/platforms/databricks/discovery.py").read_text( + encoding="utf-8" + ) + + forbidden = ( + "datacontract", + "open_data_contract_standard", + "semapact.importers", + "semapact.observation", + "semapact.lifecycle", + "semapact.governance", + ) + assert not any(name in source for name in forbidden)