diff --git a/.agents/skills/semapact-system/SKILL.md b/.agents/skills/semapact-system/SKILL.md index 73af944..7da2e56 100644 --- a/.agents/skills/semapact-system/SKILL.md +++ b/.agents/skills/semapact-system/SKILL.md @@ -18,31 +18,44 @@ SemaPact is an enterprise data contract control plane and governance platform. I ## 2. Layered Architecture Boundaries (CRITICAL) ### A. Ingestion / Import Layer -- **Role:** Converts external data structures (Delta Tables, Spark DDL, Unity Catalog) into Open Data Contract Standard (ODCS) models. +- **Role:** Converts external data structures into Open Data Contract Standard (ODCS) models when a contract import is explicitly requested. - **Rules:** - Must remain strictly stateless and idempotent. - **NEVER** place merge, governance, or GitOps logic inside importers. + - Contract import is distinct from platform observation; observing a platform must not implicitly create or mutate an ODCS contract. -### B. Core Contract Model -- **Role:** Single canonical representation of the schema and metadata. +### B. Governed Contract Model +- **Role:** Single canonical representation of governed desired contract state. - **Rules:** - - The ODCS YAML/Pydantic model is the single source of truth across the architecture. No alternative models are allowed. + - The ODCS YAML/Pydantic model is the single source of truth for governed contract state. + - External platform state must not become governed truth merely because it was observed. -### C. Lifecycle Governance Layer +### C. Platform Observation Model +- **Role:** Represents point-in-time external platform state for assurance and reconciliation workflows. +- **Rules:** + - `ObservedPlatformState` is a read-side model, not an alternative canonical contract format. + - Core observation models must remain platform-neutral; provider hierarchy belongs in adapter-local mapping into a generic ordered `namespace`. + - Platform-local identity must remain distinct from ODCS contract identity. + - Provider adapters may reuse official platform SDK access, but must not route observation through ODCS import/projection. + - Observation must not invoke lifecycle merge, governance evaluation, release mutation, or platform writeback. + - Rich metadata, constraints, relationships, and lineage are evidence enrichments, not prerequisites for the minimal observed-state model. + - Converting observed state into an ODCS contract is an explicit import workflow, never an implicit observation side effect. + +### D. Lifecycle Governance Layer - **Role:** Handles breaking change checks, deprecation rules, merge policies, and version bump calculations. - **Rules:** - This is the **ONLY** place where contract lifecycle logic is allowed. - - It must remain fully decoupled from the UI and ingestion layers. + - It must remain fully decoupled from the UI, ingestion, and platform observation layers. -### D. Export Layer +### E. Export Layer - **Role:** Converts contracts to downstream assets (Great Expectations suites, Spark DDL, Graph cypher). - **Rules:** - Exporters must be read-only and **NEVER** modify the original contracts. -### E. Orchestration Layer +### F. Orchestration Layer - **Role:** Coordinates multi-step workflows (e.g. import → merge → export → PR). - **Rules:** - Coordinates execution paths but must NOT contain custom business logic. -### F. DevOps Layer +### G. DevOps Layer - **Role:** Automates PR creation, version bumps, release manifest building, and metadata auditing. diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f7af123..30abc7c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -43,6 +43,37 @@ jobs: - name: Verify package build run: uv build + databricks-extra: + runs-on: ubuntu-latest + timeout-minutes: 30 + env: + UV_NO_PROGRESS: "1" + steps: + - name: Checkout + uses: actions/checkout@v7 + + - name: Set up Python + uses: actions/setup-python@v7 + with: + python-version: "3.12" + + - name: Set up uv + uses: astral-sh/setup-uv@v7 + + - name: Create isolated environment + run: uv venv + + - name: Install SemaPact Databricks extra from project metadata + run: uv pip install --python .venv/bin/python -e ".[databricks]" pytest + + - name: Verify Databricks SDK and existing Unity import boundaries + run: >- + .venv/bin/python -c + "from databricks.sdk import WorkspaceClient; from databricks.sdk.service.catalog import TableInfo; from semapact.importers.unity_importer import import_unity_contract; print(WorkspaceClient.__name__, TableInfo.__name__, import_unity_contract.__name__)" + + - name: Run Databricks observation tests with official SDK installed + run: .venv/bin/python -m pytest tests/test_observation_databricks.py + coverage: runs-on: ubuntu-latest timeout-minutes: 30 diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 0b367d8..9ba0bdf 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -97,7 +97,36 @@ Important boundary: - `GovernanceService` creates `ChangeContext` - lifecycle/governance lower layers consume that context and must not regenerate it -### 5. Exporters +### 5. Platform Observation + +Location: + +- `semapact/observation/` + +Responsibilities: + +- represent point-in-time external platform state independently from governed contracts +- keep the core observed-state domain platform-neutral +- map provider-local identity into `platform + ordered namespace + asset` +- observe physical asset/property state without creating or mutating ODCS contracts + +Key modules: + +- `semapact/observation/models.py` +- `semapact/observation/databricks.py` + +Important boundary: + +- `ObservedPlatformState` is read-side state, not governed truth +- observed asset identity is platform-local and distinct from ODCS contract identity +- provider-specific hierarchy such as Databricks `catalog/schema` belongs in the adapter, not the core model +- Databricks observation consumes the official SDK `WorkspaceClient.tables.get(...) -> TableInfo` boundary instead of reimplementing the Unity Catalog REST transport +- observation projects `TableInfo` directly into observed state; it must not route through datacontract-cli's ODCS projection +- rich metadata, constraints, relationships, and lineage are follow-up evidence enrichments rather than prerequisites for the minimal observation model +- observation must not invoke lifecycle merge, governance evaluation, release mutation, or platform writeback +- explicit contract import remains a separate workflow + +### 6. Exporters Location: @@ -115,7 +144,7 @@ Key modules: - `semapact/quality/ge_exporter.py` - `semapact/exporters/sql_exporter.py` -### 6. Orchestration +### 7. Orchestration Location: @@ -130,7 +159,7 @@ Key module: - `semapact/orchestrator/pipeline.py` -### 7. Interfaces +### 8. Interfaces Location: @@ -148,14 +177,28 @@ Current interface: - CLI in `semapact/interfaces/cli.py` - command adapters in `semapact/interfaces/commands/` -## Current Contract Model +## Governed Contract Model + +SemaPact assumes: + +- Open Data Contract Standard (ODCS) is the single canonical representation of governed desired contract state +- `open_data_contract_standard.model.OpenDataContractStandard` is the canonical governed contract domain model + +The system may temporarily work with Python `dict` objects at contract boundaries, but contract normalization should converge back to ODCS objects or ODCS-shaped mappings. + +Platform observations are intentionally different. `ObservedPlatformState` is a non-canonical read-side model describing what an external platform reports at a point in time. It must not replace or mutate the governed ODCS contract. + +Conceptually: -SemaPact currently assumes: +```text +Approved ODCS Contract += desired governed state -- Open Data Contract Standard (ODCS) is the single canonical contract representation -- `open_data_contract_standard.model.OpenDataContractStandard` is the canonical runtime model +ObservedPlatformState += observed external platform state +``` -The system may temporarily work with Python `dict` objects at boundaries, but normalization should converge back to ODCS objects or ODCS-shaped mappings. +Converting external metadata into a new ODCS contract is an explicit import workflow. Observing platform state does not implicitly perform that conversion. ## Root Contract Governance @@ -367,10 +410,13 @@ Precedence: ## Current Design Principles -- main contract is canonical and immutable from presentation paths +- main governed contract is canonical and immutable from presentation paths +- ODCS is the canonical model for governed desired contract state +- platform observation is separate read-side state and cannot become governed truth implicitly +- core observation models are platform-neutral; provider-specific hierarchy belongs in adapters +- reuse official platform access/SDK layers where practical, while keeping ODCS import projection separate from observation - service layer is the application boundary between interfaces and system logic - lifecycle logic belongs in the lifecycle layer -- ODCS is the canonical contract model - datacontract-cli is reused where possible instead of reimplemented ## Authoritative Governance Invariants @@ -396,6 +442,10 @@ Precedence: ## Known Next Steps +- separate platform discovery, observation, and contract import application workflows, starting with Databricks Unity Catalog +- add stable observed-state fingerprints and reconciliation semantics +- add governance-relevant metadata/constraint/relationship evidence independently from the minimal observation model +- add lineage as optional runtime evidence rather than a core observation dependency - formalize draft promotion flow - continue reducing interface-specific logic that still lives near command/editor helpers -- keep converging helper logic toward ODCS model-driven behavior +- keep converging governed contract helper logic toward ODCS model-driven behavior diff --git a/pyproject.toml b/pyproject.toml index dd38b57..8437544 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -74,7 +74,7 @@ tui = [ "textual>=8.2.7", ] databricks = [ - "databricks-sql-connector>=3.0.0", + "datacontract-cli[databricks]>=0.12.0", "pyspark>=3.3.0", ] azure = [ @@ -101,7 +101,7 @@ all = [ "openai>=1.0.0", "litellm>=1.83.14", "textual>=8.2.7", - "databricks-sql-connector>=3.0.0", + "datacontract-cli[databricks]>=0.12.0", "pyspark>=3.3.0", "azure-identity>=1.16.0", "azure-storage-file-datalake>=12.16.0", diff --git a/semapact/__init__.py b/semapact/__init__.py index 8bb2424..1af74dc 100644 --- a/semapact/__init__.py +++ b/semapact/__init__.py @@ -1,57 +1,96 @@ -"""SemaPact enterprise library.""" - -from semapact.core.loader import ContractLoader, load_contract -from semapact.core.validator import ContractValidator -from semapact.devops.pr_creator import AzureDevOpsConfig, PullRequestCreator -from semapact.devops.release_workflow import ( - BatchReleaseManifestBuild, - BatchReleaseTask, - ReleasePullRequestPlan, - RepositoryContractChange, - batch_manifest_build_to_dict, - build_batch_release_manifest, - build_release_pr_plan, - create_release_pull_request, - create_release_pull_requests_from_manifest, - load_batch_release_tasks, - repository_change_to_dict, -) -from semapact.exporters.sql_exporter import ( - SparkSqlContractExporter, - export_contract_to_spark_sql, -) -from semapact.importers.delta_importer import DeltaTableImporter -from semapact.importers.sql_importer import SQLFolderImporter -from semapact.lifecycle.merge_engine import ContractMergeEngine -from semapact.lifecycle.policy import evaluate_merge_policy -from semapact.orchestrator.pipeline import ContractPipeline -from semapact.quality.ge_exporter import GreatExpectationsExporter -from semapact.quality.validation import run_contract_tests - -__all__ = [ - "ContractLoader", - "ContractValidator", - "DeltaTableImporter", - "SQLFolderImporter", - "ContractMergeEngine", - "evaluate_merge_policy", - "GreatExpectationsExporter", - "SparkSqlContractExporter", - "export_contract_to_spark_sql", - "ContractPipeline", - "PullRequestCreator", - "AzureDevOpsConfig", - "BatchReleaseManifestBuild", - "BatchReleaseTask", - "ReleasePullRequestPlan", - "RepositoryContractChange", - "batch_manifest_build_to_dict", - "build_batch_release_manifest", - "build_release_pr_plan", - "create_release_pull_request", - "create_release_pull_requests_from_manifest", - "load_batch_release_tasks", - "repository_change_to_dict", - "load_contract", - "run_contract_tests", -] +"""SemaPact enterprise library. + +The package root intentionally lazy-loads public exports so installing one +optional platform capability does not import unrelated platform dependencies. +""" + +from __future__ import annotations + +from importlib import import_module +from typing import Any + +_EXPORTS: dict[str, tuple[str, str]] = { + "ContractLoader": ("semapact.core.loader", "ContractLoader"), + "load_contract": ("semapact.core.loader", "load_contract"), + "ContractValidator": ("semapact.core.validator", "ContractValidator"), + "AzureDevOpsConfig": ("semapact.devops.pr_creator", "AzureDevOpsConfig"), + "PullRequestCreator": ("semapact.devops.pr_creator", "PullRequestCreator"), + "BatchReleaseManifestBuild": ( + "semapact.devops.release_workflow", + "BatchReleaseManifestBuild", + ), + "BatchReleaseTask": ("semapact.devops.release_workflow", "BatchReleaseTask"), + "ReleasePullRequestPlan": ( + "semapact.devops.release_workflow", + "ReleasePullRequestPlan", + ), + "RepositoryContractChange": ( + "semapact.devops.release_workflow", + "RepositoryContractChange", + ), + "batch_manifest_build_to_dict": ( + "semapact.devops.release_workflow", + "batch_manifest_build_to_dict", + ), + "build_batch_release_manifest": ( + "semapact.devops.release_workflow", + "build_batch_release_manifest", + ), + "build_release_pr_plan": ( + "semapact.devops.release_workflow", + "build_release_pr_plan", + ), + "create_release_pull_request": ( + "semapact.devops.release_workflow", + "create_release_pull_request", + ), + "create_release_pull_requests_from_manifest": ( + "semapact.devops.release_workflow", + "create_release_pull_requests_from_manifest", + ), + "load_batch_release_tasks": ( + "semapact.devops.release_workflow", + "load_batch_release_tasks", + ), + "repository_change_to_dict": ( + "semapact.devops.release_workflow", + "repository_change_to_dict", + ), + "SparkSqlContractExporter": ( + "semapact.exporters.sql_exporter", + "SparkSqlContractExporter", + ), + "export_contract_to_spark_sql": ( + "semapact.exporters.sql_exporter", + "export_contract_to_spark_sql", + ), + "DeltaTableImporter": ("semapact.importers.delta_importer", "DeltaTableImporter"), + "SQLFolderImporter": ("semapact.importers.sql_importer", "SQLFolderImporter"), + "ContractMergeEngine": ("semapact.lifecycle.merge_engine", "ContractMergeEngine"), + "evaluate_merge_policy": ("semapact.lifecycle.policy", "evaluate_merge_policy"), + "ContractPipeline": ("semapact.orchestrator.pipeline", "ContractPipeline"), + "GreatExpectationsExporter": ( + "semapact.quality.ge_exporter", + "GreatExpectationsExporter", + ), + "run_contract_tests": ("semapact.quality.validation", "run_contract_tests"), +} + +__all__ = list(_EXPORTS) + + +def __getattr__(name: str) -> Any: + """Resolve public exports only when callers actually request them.""" + target = _EXPORTS.get(name) + if target is None: + raise AttributeError(f"module {__name__!r} has no attribute {name!r}") + + module_name, attribute = target + value = getattr(import_module(module_name), attribute) + globals()[name] = value + return value + + +def __dir__() -> list[str]: + """Include lazy public exports in interactive discovery.""" + return sorted({*globals(), *__all__}) diff --git a/semapact/importers/__init__.py b/semapact/importers/__init__.py index 0cad142..c33349c 100644 --- a/semapact/importers/__init__.py +++ b/semapact/importers/__init__.py @@ -1,15 +1,29 @@ +"""SemaPact importer registrations. + +Optional importer dependencies are registered only when their provider package +is installed, so one platform extra does not require unrelated platforms. +""" + +from __future__ import annotations + +from importlib.util import find_spec + from datacontract.imports.importer_factory import importer_factory -from semapact.importers.delta_importer import DeltaTableImporter from semapact.importers.sql_importer import SQLFolderImporter -# Register SemaPact custom importers with datacontract-cli importer factory. -importer_factory.register_importer("delta", DeltaTableImporter) -importer_factory.register_importer("delta-table", DeltaTableImporter) +__all__ = ["SQLFolderImporter"] + +# SQL-folder support has no provider-specific optional dependency beyond the +# base SemaPact/datacontract-cli installation. importer_factory.register_importer("sql-folder", SQLFolderImporter) importer_factory.register_importer("delta-ddl", SQLFolderImporter) -__all__ = [ - "DeltaTableImporter", - "SQLFolderImporter", -] +# Delta support depends on the separate ``deltalake`` extra. Import and register +# it only when that optional dependency is actually available. +if find_spec("deltalake") is not None: + from semapact.importers.delta_importer import DeltaTableImporter + + importer_factory.register_importer("delta", DeltaTableImporter) + importer_factory.register_importer("delta-table", DeltaTableImporter) + __all__.append("DeltaTableImporter") diff --git a/semapact/observation/__init__.py b/semapact/observation/__init__.py new file mode 100644 index 0000000..dafdd71 --- /dev/null +++ b/semapact/observation/__init__.py @@ -0,0 +1,19 @@ +"""Platform-neutral observation domain for SemaPact read-side state.""" + +from semapact.observation.models import ( + ObservedAsset, + ObservedAssetIdentity, + ObservedPlatformState, + ObservedProperty, + ObservedPropertyIdentity, + serialize_observed_state, +) + +__all__ = [ + "ObservedAsset", + "ObservedAssetIdentity", + "ObservedPlatformState", + "ObservedProperty", + "ObservedPropertyIdentity", + "serialize_observed_state", +] diff --git a/semapact/observation/databricks.py b/semapact/observation/databricks.py new file mode 100644 index 0000000..5c959b7 --- /dev/null +++ b/semapact/observation/databricks.py @@ -0,0 +1,190 @@ +"""Databricks observation adapter backed by the official Databricks SDK. + +The adapter consumes the same ``WorkspaceClient.tables.get(...) -> TableInfo`` +boundary used by datacontract-cli, but projects that source metadata into +SemaPact's platform-neutral observation model instead of into ODCS. + +Authentication and credential resolution are caller concerns. This module +accepts an already initialized/authenticated ``WorkspaceClient`` and must not +resolve PATs, OAuth credentials, Azure identity, profiles, service principals, +or other authentication mechanisms itself. +""" + +from __future__ import annotations + +from datetime import datetime, timezone +from typing import TYPE_CHECKING, Any, Mapping, Protocol + +from semapact.observation.models import ( + ObservedAsset, + ObservedAssetIdentity, + ObservedPlatformState, + ObservedProperty, + ObservedPropertyIdentity, +) + +if TYPE_CHECKING: + from databricks.sdk import WorkspaceClient + from databricks.sdk.service.catalog import TableInfo + +DATABRICKS_PLATFORM = "databricks" + + +class _TableInfoLike(Protocol): + def as_dict(self) -> dict[str, Any]: ... + + +class _TablesApiLike(Protocol): + def get(self, full_name: str) -> _TableInfoLike: ... + + +class _WorkspaceClientLike(Protocol): + tables: _TablesApiLike + + +def observe_databricks_table( + *, + client: WorkspaceClient | _WorkspaceClientLike, + table_fqn: str, + source_identifier: str, + captured_at: datetime | None = None, +) -> ObservedPlatformState: + """Observe one Databricks table without generating or mutating an ODCS contract. + + ``client`` is expected to be an already initialized/authenticated official + ``databricks.sdk.WorkspaceClient`` in production. Authentication method and + credential resolution are intentionally outside this adapter's scope. + The client is injectable so core observation tests require no live + Databricks workspace. + """ + if not table_fqn: + raise ValueError("table_fqn is required for Databricks observation") + if not source_identifier: + raise ValueError("source_identifier is required for Databricks observation") + + observed_at = captured_at or datetime.now(timezone.utc) + _require_aware_datetime(observed_at) + + table = client.tables.get(table_fqn) + return map_databricks_table_info( + table, + source_identifier=source_identifier, + captured_at=observed_at, + table_fqn=table_fqn, + ) + + +def map_databricks_table_info( + table: TableInfo | _TableInfoLike, + *, + source_identifier: str, + captured_at: datetime, + table_fqn: str | None = None, +) -> ObservedPlatformState: + """Project an SDK ``TableInfo`` into the platform-neutral observation model.""" + _require_aware_datetime(captured_at) + if not source_identifier: + raise ValueError("source_identifier is required for Databricks observation") + + metadata = table.as_dict() + if not isinstance(metadata, Mapping): + raise TypeError("Databricks TableInfo.as_dict() must return a mapping") + + identity = _asset_identity(metadata, table_fqn=table_fqn) + asset = ObservedAsset( + identity=identity, + asset_type=_text(metadata.get("table_type")), + properties=_properties(metadata.get("columns"), identity=identity), + ) + + return ObservedPlatformState( + platform=DATABRICKS_PLATFORM, + source_identifier=source_identifier.rstrip("/"), + assets=(asset,), + captured_at=captured_at, + fingerprint=None, + ) + + +def _asset_identity( + metadata: Mapping[str, Any], *, table_fqn: str | None +) -> ObservedAssetIdentity: + catalog = _text(metadata.get("catalog_name")) + schema = _text(metadata.get("schema_name")) + asset = _text(metadata.get("name")) + + full_name = _text(metadata.get("full_name") or table_fqn) + fallback = _split_table_fqn(full_name) if full_name else None + if fallback: + catalog = catalog or fallback[0] + schema = schema or fallback[1] + asset = asset or fallback[2] + + if not catalog or not schema or not asset: + raise ValueError("Databricks TableInfo must identify catalog, schema, and table name") + + return ObservedAssetIdentity( + platform=DATABRICKS_PLATFORM, + namespace=(catalog, schema), + asset=asset, + ) + + +def _split_table_fqn(value: str) -> tuple[str, str, str] | None: + parts = tuple(part.strip() for part in value.split(".")) + if len(parts) != 3 or not all(parts): + return None + return parts + + +def _properties( + value: Any, *, identity: ObservedAssetIdentity +) -> tuple[ObservedProperty, ...]: + if not isinstance(value, list): + return () + + observed: list[tuple[int | None, ObservedProperty]] = [] + for item in value: + if not isinstance(item, Mapping): + continue + + name = _text(item.get("name")) + if not name: + continue + + nullable = item.get("nullable") if isinstance(item.get("nullable"), bool) else None + position = item.get("position") + if isinstance(position, bool) or not isinstance(position, int): + position = None + + observed.append( + ( + position, + ObservedProperty( + identity=ObservedPropertyIdentity(asset=identity, property=name), + physical_type=_text(item.get("type_text") or item.get("type_name")), + nullable=nullable, + ), + ) + ) + + observed.sort( + key=lambda item: ( + item[0] is None, + item[0] if item[0] is not None else 0, + item[1].identity.property.casefold(), + ) + ) + return tuple(item[1] for item in observed) + + +def _text(value: Any) -> str | None: + if value is None: + return None + text = str(value).strip() + return text or None + + +def _require_aware_datetime(value: datetime) -> None: + if value.tzinfo is None or value.utcoffset() is None: + raise ValueError("captured_at must be timezone-aware") diff --git a/semapact/observation/models.py b/semapact/observation/models.py new file mode 100644 index 0000000..683571b --- /dev/null +++ b/semapact/observation/models.py @@ -0,0 +1,94 @@ +"""Platform-neutral read-side models for observed external state. + +An observation describes what an external data platform reports at a point in +time. It is deliberately separate from the governed ODCS contract model and is +never governed truth by itself. +""" + +from __future__ import annotations + +import json +from datetime import datetime + +from pydantic import BaseModel, ConfigDict + + +class ObservationModel(BaseModel): + """Shared immutable base for observation models.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + +class ObservedAssetIdentity(ObservationModel): + """Platform-local identity for one observed asset. + + ``namespace`` is intentionally provider-neutral. A Databricks adapter may + populate it with ``(catalog, schema)`` while another platform may use a + different hierarchy without changing the domain model. + """ + + platform: str + namespace: tuple[str, ...] = () + asset: str + + @property + def canonical_key(self) -> tuple[str, ...]: + """Return a case-normalized identity key without provider semantics.""" + return ( + self.platform.casefold(), + *(part.casefold() for part in self.namespace), + self.asset.casefold(), + ) + + +class ObservedPropertyIdentity(ObservationModel): + """Identity for a property within an observed asset.""" + + asset: ObservedAssetIdentity + property: str + + @property + def canonical_key(self) -> tuple[str, ...]: + """Return the case-normalized property identity key.""" + return (*self.asset.canonical_key, self.property.casefold()) + + +class ObservedProperty(ObservationModel): + """Observed physical property/column state.""" + + identity: ObservedPropertyIdentity + physical_type: str | None = None + nullable: bool | None = None + + +class ObservedAsset(ObservationModel): + """Observed physical state for one external asset.""" + + identity: ObservedAssetIdentity + asset_type: str | None = None + properties: tuple[ObservedProperty, ...] = () + + +class ObservedPlatformState(ObservationModel): + """Point-in-time platform observation independent from governed ODCS state. + + ``fingerprint`` is intentionally optional. Stable normalization and hashing + belong to the dedicated observed-state fingerprint capability; this model + only reserves the field used by that later capability. + """ + + platform: str + source_identifier: str + assets: tuple[ObservedAsset, ...] + captured_at: datetime + fingerprint: str | None = None + + +def serialize_observed_state(state: ObservedPlatformState) -> str: + """Serialize observed state deterministically for machine-readable use.""" + return json.dumps( + state.model_dump(mode="json"), + ensure_ascii=False, + separators=(",", ":"), + sort_keys=True, + ) diff --git a/tests/fixtures/observation/databricks/orders_table_info.json b/tests/fixtures/observation/databricks/orders_table_info.json new file mode 100644 index 0000000..42b613b --- /dev/null +++ b/tests/fixtures/observation/databricks/orders_table_info.json @@ -0,0 +1,30 @@ +{ + "name": "orders", + "catalog_name": "main", + "schema_name": "silver", + "full_name": "main.silver.orders", + "table_type": "MANAGED", + "columns": [ + { + "name": "customer_id", + "type_name": "STRING", + "type_text": "string", + "position": 1, + "nullable": false + }, + { + "name": "note", + "type_name": "STRING", + "type_text": "string", + "position": 2, + "nullable": true + }, + { + "name": "order_id", + "type_name": "BIGINT", + "type_text": "bigint", + "position": 0, + "nullable": false + } + ] +} diff --git a/tests/test_observation_databricks.py b/tests/test_observation_databricks.py new file mode 100644 index 0000000..a4070cb --- /dev/null +++ b/tests/test_observation_databricks.py @@ -0,0 +1,180 @@ +from __future__ import annotations + +import inspect +import json +from datetime import datetime, timezone +from pathlib import Path + +import pytest +from pydantic import ValidationError + +from semapact.observation import databricks as databricks_observation +from semapact.observation.databricks import ( + map_databricks_table_info, + observe_databricks_table, +) +from semapact.observation.models import ObservedPlatformState, serialize_observed_state + +FIXTURE = ( + Path(__file__).parent + / "fixtures" + / "observation" + / "databricks" + / "orders_table_info.json" +) +CAPTURED_AT = datetime(2026, 8, 30, 3, 0, tzinfo=timezone.utc) + + +def _payload() -> dict[str, object]: + return json.loads(FIXTURE.read_text(encoding="utf-8")) + + +class _FakeTableInfo: + def __init__(self, payload: dict[str, object]) -> None: + self._payload = payload + + def as_dict(self) -> dict[str, object]: + return self._payload + + +class _FakeTablesApi: + def __init__(self, table: _FakeTableInfo) -> None: + self._table = table + self.calls: list[str] = [] + + def get(self, full_name: str) -> _FakeTableInfo: + self.calls.append(full_name) + return self._table + + +class _FakeWorkspaceClient: + def __init__(self, table: _FakeTableInfo) -> None: + self.tables = _FakeTablesApi(table) + + +def test_databricks_table_info_maps_to_platform_neutral_observation() -> None: + state = map_databricks_table_info( + _FakeTableInfo(_payload()), + source_identifier="https://adb.example/", + captured_at=CAPTURED_AT, + ) + + assert isinstance(state, ObservedPlatformState) + assert state.platform == "databricks" + assert state.source_identifier == "https://adb.example" + assert state.captured_at == CAPTURED_AT + assert state.fingerprint is None + + assert len(state.assets) == 1 + asset = state.assets[0] + assert asset.identity.platform == "databricks" + assert asset.identity.namespace == ("main", "silver") + assert asset.identity.asset == "orders" + assert asset.identity.canonical_key == ( + "databricks", + "main", + "silver", + "orders", + ) + assert asset.asset_type == "MANAGED" + + assert [item.identity.property for item in asset.properties] == [ + "order_id", + "customer_id", + "note", + ] + assert asset.properties[0].physical_type == "bigint" + assert asset.properties[0].nullable is False + assert asset.properties[1].physical_type == "string" + assert asset.properties[2].nullable is True + + +def test_observation_model_identity_does_not_encode_databricks_namespace_names() -> None: + state = map_databricks_table_info( + _FakeTableInfo(_payload()), + source_identifier="workspace-a", + captured_at=CAPTURED_AT, + ) + + identity_fields = set(type(state.assets[0].identity).model_fields) + assert identity_fields == {"platform", "namespace", "asset"} + assert "catalog" not in identity_fields + assert "schema" not in identity_fields + + +def test_serialization_is_stable_when_databricks_column_order_changes() -> None: + payload = _payload() + reordered = dict(payload) + reordered["columns"] = list(reversed(payload["columns"])) # type: ignore[index] + + left = map_databricks_table_info( + _FakeTableInfo(payload), + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + right = map_databricks_table_info( + _FakeTableInfo(reordered), + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + + assert left == right + assert serialize_observed_state(left) == serialize_observed_state(right) + + +def test_observe_databricks_table_uses_workspace_client_tables_get() -> None: + client = _FakeWorkspaceClient(_FakeTableInfo(_payload())) + + state = observe_databricks_table( + client=client, + table_fqn="main.silver.orders", + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + + assert client.tables.calls == ["main.silver.orders"] + assert state.assets[0].identity.asset == "orders" + + +def test_mapper_accepts_official_databricks_sdk_table_info_when_extra_is_installed() -> None: + catalog = pytest.importorskip("databricks.sdk.service.catalog") + table = catalog.TableInfo.from_dict(_payload()) + + state = map_databricks_table_info( + table, + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + + assert state.assets[0].identity.namespace == ("main", "silver") + assert state.assets[0].properties[0].identity.property == "order_id" + + +def test_observation_models_are_immutable() -> None: + state = map_databricks_table_info( + _FakeTableInfo(_payload()), + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + + with pytest.raises(ValidationError): + state.platform = "other-platform" # type: ignore[misc] + + +def test_databricks_observer_does_not_depend_on_contract_projection_layers() -> None: + source = inspect.getsource(databricks_observation) + + assert "open_data_contract_standard" not in source + assert "datacontract.imports" not in source + assert "semapact.importers" not in source + assert "semapact.lifecycle" not in source + assert "semapact.governance" not in source + + +def test_observation_requires_timezone_aware_capture_time() -> None: + with pytest.raises(ValueError, match="timezone-aware"): + map_databricks_table_info( + _FakeTableInfo(_payload()), + source_identifier="https://adb.example", + captured_at=datetime(2026, 8, 30, 3, 0), + )