From 9402e93588cbbcb3d61704f79c0bfa188851dc8a Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:04:28 +1000 Subject: [PATCH 01/28] feat(runtime): add observed runtime state domain models --- semapact/runtime/models.py | 146 +++++++++++++++++++++++++++++++++++++ 1 file changed, 146 insertions(+) create mode 100644 semapact/runtime/models.py diff --git a/semapact/runtime/models.py b/semapact/runtime/models.py new file mode 100644 index 0000000..1e77e49 --- /dev/null +++ b/semapact/runtime/models.py @@ -0,0 +1,146 @@ +"""Read-side domain models for observed runtime state. + +Observed runtime state is deliberately separate from the governed ODCS contract model. +It describes what a platform reports now; it is never governed truth by itself. +""" + +from __future__ import annotations + +import json +from datetime import datetime + +from pydantic import BaseModel, ConfigDict + + +class RuntimeModel(BaseModel): + """Shared immutable base for runtime observation models.""" + + model_config = ConfigDict(frozen=True, extra="forbid") + + +class ObservedAssetIdentity(RuntimeModel): + """Platform-local identity for one observed runtime asset.""" + + platform: str + catalog: str + schema: str + asset: str + + @property + def canonical_key(self) -> tuple[str, str, str, str]: + """Return the case-normalized runtime identity key.""" + return ( + self.platform.casefold(), + self.catalog.casefold(), + self.schema.casefold(), + self.asset.casefold(), + ) + + +class ObservedPropertyIdentity(RuntimeModel): + """Identity for a property within an observed runtime asset.""" + + asset: ObservedAssetIdentity + property: str + + @property + def canonical_key(self) -> tuple[str, str, str, str, str]: + """Return the case-normalized runtime property identity key.""" + return (*self.asset.canonical_key, self.property.casefold()) + + +class RuntimeMetadata(RuntimeModel): + """Canonical preservation of platform metadata not modeled directly.""" + + key: str + value_json: str + + +class RuntimeTag(RuntimeModel): + """Observed platform tag.""" + + key: str + value: str + + +class ObservedProperty(RuntimeModel): + """Observed runtime property/column state.""" + + identity: ObservedPropertyIdentity + logical_type: str | None = None + physical_type: str | None = None + nullable: bool | None = None + required: bool | None = None + position: int | None = None + comment: str | None = None + tags: tuple[RuntimeTag, ...] = () + metadata: tuple[RuntimeMetadata, ...] = () + + +class ObservedConstraint(RuntimeModel): + """Observed non-relational table constraint.""" + + constraint_type: str + name: str | None = None + columns: tuple[str, ...] = () + expression: str | None = None + metadata: tuple[RuntimeMetadata, ...] = () + + +class ObservedRelationship(RuntimeModel): + """Observed relationship from one runtime asset to another.""" + + relationship_type: str + source_columns: tuple[str, ...] = () + target_reference: str | None = None + target_columns: tuple[str, ...] = () + constraint_name: str | None = None + metadata: tuple[RuntimeMetadata, ...] = () + + +class RuntimeEvidence(RuntimeModel): + """Reference to runtime evidence reported by the source platform.""" + + kind: str + reference: str + metadata: tuple[RuntimeMetadata, ...] = () + + +class ObservedAsset(RuntimeModel): + """Observed state for one runtime asset.""" + + identity: ObservedAssetIdentity + asset_type: str | None = None + data_source_format: str | None = None + owner: str | None = None + comment: str | None = None + properties: tuple[ObservedProperty, ...] = () + constraints: tuple[ObservedConstraint, ...] = () + relationships: tuple[ObservedRelationship, ...] = () + tags: tuple[RuntimeTag, ...] = () + metadata: tuple[RuntimeMetadata, ...] = () + + +class ObservedContractState(RuntimeModel): + """Point-in-time runtime observation independent from governed ODCS state. + + ``fingerprint`` is intentionally optional. M1 fingerprint normalization and + hashing are a separate capability; this model only reserves the contract. + """ + + platform: str + source_identifier: str + assets: tuple[ObservedAsset, ...] + captured_at: datetime + fingerprint: str | None = None + evidence: tuple[RuntimeEvidence, ...] = () + + +def serialize_observed_state(state: ObservedContractState) -> str: + """Serialize observed state deterministically for machine-readable use.""" + return json.dumps( + state.model_dump(mode="json"), + ensure_ascii=False, + separators=(",", ":"), + sort_keys=True, + ) From 7c08734e83bf35400804bac1863fcb5c6356718a Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:05:18 +1000 Subject: [PATCH 02/28] feat(runtime): observe Unity Catalog table state --- semapact/runtime/unity.py | 579 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 579 insertions(+) create mode 100644 semapact/runtime/unity.py diff --git a/semapact/runtime/unity.py b/semapact/runtime/unity.py new file mode 100644 index 0000000..65aa95a --- /dev/null +++ b/semapact/runtime/unity.py @@ -0,0 +1,579 @@ +"""Unity Catalog read-side observation without ODCS contract generation.""" + +from __future__ import annotations + +import json +from datetime import datetime, timezone +from typing import Any, Callable, Mapping, Sequence +from urllib.error import HTTPError, URLError +from urllib.parse import quote +from urllib.request import Request, urlopen + +from semapact.exceptions import StorageError +from semapact.runtime.models import ( + ObservedAsset, + ObservedAssetIdentity, + ObservedConstraint, + ObservedContractState, + ObservedProperty, + ObservedPropertyIdentity, + ObservedRelationship, + RuntimeEvidence, + RuntimeMetadata, + RuntimeTag, +) + +UNITY_CATALOG_PLATFORM = "databricks-unity-catalog" +UnityMetadataFetcher = Callable[[str, str, str], Mapping[str, Any]] + + +def observe_unity_table( + *, + table_fqn: str, + workspace_url: str, + token: str, + captured_at: datetime | None = None, + fetcher: UnityMetadataFetcher | None = None, +) -> ObservedContractState: + """Observe one Unity Catalog table without creating or mutating an ODCS contract.""" + if not workspace_url: + raise ValueError("workspace_url is required for Unity Catalog observation") + if not token: + raise ValueError("token is required for Unity Catalog observation") + if not table_fqn: + raise ValueError("table_fqn is required for Unity Catalog observation") + + observed_at = captured_at or datetime.now(timezone.utc) + if observed_at.tzinfo is None or observed_at.utcoffset() is None: + raise ValueError("captured_at must be timezone-aware") + + metadata_fetcher = fetcher or fetch_unity_table_metadata + payload = metadata_fetcher(workspace_url, token, table_fqn) + return map_unity_table_metadata( + payload, + source_identifier=workspace_url.rstrip("/"), + table_fqn=table_fqn, + captured_at=observed_at, + ) + + +def fetch_unity_table_metadata( + workspace_url: str, + token: str, + table_fqn: str, +) -> Mapping[str, Any]: + """Fetch the Unity Catalog TableInfo payload for one exact table.""" + if not workspace_url or not token: + raise ValueError("workspace_url and token are required for Unity observation") + + endpoint = ( + f"{workspace_url.rstrip('/')}/api/2.1/unity-catalog/tables/" + f"{quote(table_fqn, safe='')}" + ) + request = Request( + endpoint, + headers={ + "Authorization": f"Bearer {token}", + "Accept": "application/json", + }, + method="GET", + ) + + try: + with urlopen(request, timeout=8) as response: + payload = response.read().decode("utf-8") + except HTTPError as exc: + raise StorageError( + f"Unity table observation request failed: HTTP {exc.code}" + ) from exc + except URLError as exc: + raise StorageError( + f"Unity table observation request failed: {exc.reason}" + ) from exc + + parsed = json.loads(payload) + if not isinstance(parsed, dict): + raise StorageError("Unity table observation response is not a JSON object") + return parsed + + +def map_unity_table_metadata( + metadata: Mapping[str, Any], + *, + source_identifier: str, + captured_at: datetime, + table_fqn: str | None = None, +) -> ObservedContractState: + """Map a Unity TableInfo-style payload into immutable observed runtime state.""" + if captured_at.tzinfo is None or captured_at.utcoffset() is None: + raise ValueError("captured_at must be timezone-aware") + + identity = _asset_identity(metadata, table_fqn=table_fqn) + properties = _properties(metadata, identity=identity) + constraints, relationships = _constraints_and_relationships(metadata) + evidence = _runtime_evidence(metadata) + + asset = ObservedAsset( + identity=identity, + asset_type=_text(metadata.get("table_type") or metadata.get("tableType")), + data_source_format=_text( + metadata.get("data_source_format") or metadata.get("dataSourceFormat") + ), + owner=_text(metadata.get("owner")), + comment=_text(metadata.get("comment")), + properties=properties, + constraints=constraints, + relationships=relationships, + tags=_tags(metadata.get("tags")), + metadata=_metadata_entries( + metadata, + excluded={ + "catalog_name", + "catalogName", + "schema_name", + "schemaName", + "name", + "full_name", + "fullName", + "table_type", + "tableType", + "data_source_format", + "dataSourceFormat", + "owner", + "comment", + "columns", + "table_constraints", + "tableConstraints", + "constraints", + "foreign_keys", + "foreignKeys", + "tags", + "view_dependencies", + "viewDependencies", + }, + ), + ) + + return ObservedContractState( + platform=UNITY_CATALOG_PLATFORM, + source_identifier=source_identifier.rstrip("/"), + assets=(asset,), + captured_at=captured_at, + fingerprint=None, + evidence=evidence, + ) + + +def _asset_identity( + metadata: Mapping[str, Any], *, table_fqn: str | None +) -> ObservedAssetIdentity: + catalog = _text(metadata.get("catalog_name") or metadata.get("catalogName")) + schema = _text(metadata.get("schema_name") or metadata.get("schemaName")) + asset = _text(metadata.get("name")) + + full_name = _text( + metadata.get("full_name") or metadata.get("fullName") or table_fqn + ) + fallback = _split_table_fqn(full_name) if full_name else None + if fallback is not None: + 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( + "Unity metadata must identify catalog, schema, and table name" + ) + + return ObservedAssetIdentity( + platform=UNITY_CATALOG_PLATFORM, + catalog=catalog, + schema=schema, + asset=asset, + ) + + +def _split_table_fqn(value: str) -> tuple[str, str, str] | None: + parts = [part.strip() for part in value.split(".")] + if len(parts) != 3 or not all(parts): + return None + return parts[0], parts[1], parts[2] + + +def _properties( + metadata: Mapping[str, Any], *, identity: ObservedAssetIdentity +) -> tuple[ObservedProperty, ...]: + raw_columns = metadata.get("columns") + if not isinstance(raw_columns, list): + return () + + properties: list[ObservedProperty] = [] + for item in raw_columns: + if not isinstance(item, Mapping): + continue + name = _text(item.get("name")) + if not name: + continue + + nullable = _bool_or_none(item.get("nullable")) + explicit_required = _bool_or_none(item.get("required")) + required = explicit_required if explicit_required is not None else ( + None if nullable is None else not nullable + ) + + properties.append( + ObservedProperty( + identity=ObservedPropertyIdentity(asset=identity, property=name), + logical_type=_text( + item.get("logical_type") or item.get("logicalType") + ), + physical_type=_text( + item.get("type_text") + or item.get("typeText") + or item.get("type_name") + or item.get("typeName") + or item.get("type") + ), + nullable=nullable, + required=required, + position=_int_or_none(item.get("position")), + comment=_text(item.get("comment")), + tags=_tags(item.get("tags")), + metadata=_metadata_entries( + item, + excluded={ + "name", + "logical_type", + "logicalType", + "type_text", + "typeText", + "type_name", + "typeName", + "type", + "nullable", + "required", + "position", + "comment", + "tags", + }, + ), + ) + ) + + properties.sort(key=_property_sort_key) + return tuple(properties) + + +def _property_sort_key(item: ObservedProperty) -> tuple[int, int, str]: + if item.position is None: + return (1, 0, item.identity.property.casefold()) + return (0, item.position, item.identity.property.casefold()) + + +def _constraints_and_relationships( + metadata: Mapping[str, Any], +) -> tuple[tuple[ObservedConstraint, ...], tuple[ObservedRelationship, ...]]: + constraints: list[ObservedConstraint] = [] + relationships: list[ObservedRelationship] = [] + + for item in _constraint_items(metadata): + primary = _mapping(item.get("primary_key_constraint") or item.get("primaryKeyConstraint")) + foreign = _mapping(item.get("foreign_key_constraint") or item.get("foreignKeyConstraint")) + + if primary is not None: + constraints.append( + ObservedConstraint( + constraint_type="PRIMARY_KEY", + name=_text(primary.get("name")), + columns=_string_tuple( + primary.get("child_columns") or primary.get("childColumns") + ), + metadata=_metadata_entries( + primary, + excluded={"name", "child_columns", "childColumns"}, + ), + ) + ) + continue + + if foreign is not None: + relationships.append( + ObservedRelationship( + relationship_type="FOREIGN_KEY", + source_columns=_string_tuple( + foreign.get("child_columns") or foreign.get("childColumns") + ), + target_reference=_text( + foreign.get("parent_table") or foreign.get("parentTable") + ), + target_columns=_string_tuple( + foreign.get("parent_columns") or foreign.get("parentColumns") + ), + constraint_name=_text(foreign.get("name")), + metadata=_metadata_entries( + foreign, + excluded={ + "name", + "child_columns", + "childColumns", + "parent_table", + "parentTable", + "parent_columns", + "parentColumns", + }, + ), + ) + ) + continue + + constraint_type = _text( + item.get("constraint_type") + or item.get("constraintType") + or item.get("type") + or item.get("kind") + ) + normalized_type = (constraint_type or "UNKNOWN").upper().replace(" ", "_") + source_columns = _string_tuple( + item.get("columns") + or item.get("column_names") + or item.get("columnNames") + or item.get("from_columns") + or item.get("fromColumns") + or item.get("child_columns") + or item.get("childColumns") + or item.get("from") + ) + target_reference = _text( + item.get("referenced_table") + or item.get("referencedTable") + or item.get("to_table") + or item.get("toTable") + or item.get("parent_table") + or item.get("parentTable") + ) + target_columns = _string_tuple( + item.get("referenced_columns") + or item.get("referencedColumns") + or item.get("to_columns") + or item.get("toColumns") + or item.get("parent_columns") + or item.get("parentColumns") + or item.get("to") + ) + + if "FOREIGN" in normalized_type or target_reference: + relationships.append( + ObservedRelationship( + relationship_type=normalized_type, + source_columns=source_columns, + target_reference=target_reference, + target_columns=target_columns, + constraint_name=_text(item.get("name")), + metadata=_metadata_entries( + item, + excluded=_FLAT_CONSTRAINT_KEYS, + ), + ) + ) + else: + constraints.append( + ObservedConstraint( + constraint_type=normalized_type, + name=_text(item.get("name")), + columns=source_columns, + expression=_text( + item.get("expression") or item.get("check_expression") + ), + metadata=_metadata_entries( + item, + excluded=_FLAT_CONSTRAINT_KEYS | {"expression", "check_expression"}, + ), + ) + ) + + constraints.sort( + key=lambda item: ( + item.constraint_type, + (item.name or "").casefold(), + tuple(value.casefold() for value in item.columns), + ) + ) + relationships.sort( + key=lambda item: ( + item.relationship_type, + (item.constraint_name or "").casefold(), + (item.target_reference or "").casefold(), + tuple(value.casefold() for value in item.source_columns), + ) + ) + return tuple(constraints), tuple(relationships) + + +_FLAT_CONSTRAINT_KEYS = { + "constraint_type", + "constraintType", + "type", + "kind", + "name", + "columns", + "column_names", + "columnNames", + "from_columns", + "fromColumns", + "child_columns", + "childColumns", + "from", + "referenced_table", + "referencedTable", + "to_table", + "toTable", + "parent_table", + "parentTable", + "referenced_columns", + "referencedColumns", + "to_columns", + "toColumns", + "parent_columns", + "parentColumns", + "to", +} + + +def _constraint_items(metadata: Mapping[str, Any]) -> tuple[Mapping[str, Any], ...]: + items: list[Mapping[str, Any]] = [] + for key in ( + "table_constraints", + "tableConstraints", + "constraints", + "foreign_keys", + "foreignKeys", + ): + value = metadata.get(key) + if isinstance(value, list): + items.extend(item for item in value if isinstance(item, Mapping)) + return tuple(items) + + +def _runtime_evidence(metadata: Mapping[str, Any]) -> tuple[RuntimeEvidence, ...]: + raw = metadata.get("view_dependencies") or metadata.get("viewDependencies") + dependency_container = _mapping(raw) + if dependency_container is None: + return () + dependencies = dependency_container.get("dependencies") + if not isinstance(dependencies, list): + return () + + evidence: list[RuntimeEvidence] = [] + for dependency in dependencies: + if not isinstance(dependency, Mapping): + continue + table = _mapping(dependency.get("table")) + function = _mapping(dependency.get("function")) + if table is not None: + reference = _text( + table.get("table_full_name") or table.get("tableFullName") + ) + if reference: + evidence.append( + RuntimeEvidence( + kind="view_table_dependency", + reference=reference, + metadata=_metadata_entries( + table, + excluded={"table_full_name", "tableFullName"}, + ), + ) + ) + elif function is not None: + reference = _text( + function.get("function_full_name") + or function.get("functionFullName") + ) + if reference: + evidence.append( + RuntimeEvidence( + kind="view_function_dependency", + reference=reference, + metadata=_metadata_entries( + function, + excluded={"function_full_name", "functionFullName"}, + ), + ) + ) + + evidence.sort(key=lambda item: (item.kind, item.reference.casefold())) + return tuple(evidence) + + +def _tags(value: Any) -> tuple[RuntimeTag, ...]: + tags: list[RuntimeTag] = [] + if isinstance(value, Mapping): + for key, item_value in value.items(): + tags.append(RuntimeTag(key=str(key), value=str(item_value))) + elif isinstance(value, list): + for item in value: + if not isinstance(item, Mapping): + continue + key = _text(item.get("key") or item.get("tag_name") or item.get("tagName")) + item_value = _text( + item.get("value") or item.get("tag_value") or item.get("tagValue") + ) + if key and item_value is not None: + tags.append(RuntimeTag(key=key, value=item_value)) + tags.sort(key=lambda item: (item.key.casefold(), item.value)) + return tuple(tags) + + +def _metadata_entries( + mapping: Mapping[str, Any], *, excluded: set[str] +) -> tuple[RuntimeMetadata, ...]: + entries = [ + RuntimeMetadata(key=str(key), value_json=_canonical_json(value)) + for key, value in mapping.items() + if key not in excluded + ] + entries.sort(key=lambda item: item.key.casefold()) + return tuple(entries) + + +def _canonical_json(value: Any) -> str: + return json.dumps( + value, + ensure_ascii=False, + separators=(",", ":"), + sort_keys=True, + ) + + +def _mapping(value: Any) -> Mapping[str, Any] | None: + return value if isinstance(value, Mapping) else None + + +def _text(value: Any) -> str | None: + if value is None: + return None + text = str(value).strip() + return text or None + + +def _string_tuple(value: Any) -> tuple[str, ...]: + if value is None: + return () + if isinstance(value, str): + return tuple(part.strip() for part in value.split(",") if part.strip()) + if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): + result = tuple(text for item in value if (text := _text(item)) is not None) + return result + return () + + +def _bool_or_none(value: Any) -> bool | None: + if isinstance(value, bool): + return value + return None + + +def _int_or_none(value: Any) -> int | None: + if isinstance(value, bool): + return None + return value if isinstance(value, int) else None From 28f7e634bbd5790cfab83221fc439a62416e6560 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:05:37 +1000 Subject: [PATCH 03/28] feat(runtime): expose runtime observation boundary --- semapact/runtime/__init__.py | 37 ++++++++++++++++++++++++++++++++++++ 1 file changed, 37 insertions(+) create mode 100644 semapact/runtime/__init__.py diff --git a/semapact/runtime/__init__.py b/semapact/runtime/__init__.py new file mode 100644 index 0000000..166b747 --- /dev/null +++ b/semapact/runtime/__init__.py @@ -0,0 +1,37 @@ +"""Runtime observation boundary for SemaPact read-side state.""" + +from semapact.runtime.models import ( + ObservedAsset, + ObservedAssetIdentity, + ObservedConstraint, + ObservedContractState, + ObservedProperty, + ObservedPropertyIdentity, + ObservedRelationship, + RuntimeEvidence, + RuntimeMetadata, + RuntimeTag, + serialize_observed_state, +) +from semapact.runtime.unity import ( + UNITY_CATALOG_PLATFORM, + map_unity_table_metadata, + observe_unity_table, +) + +__all__ = [ + "UNITY_CATALOG_PLATFORM", + "ObservedAsset", + "ObservedAssetIdentity", + "ObservedConstraint", + "ObservedContractState", + "ObservedProperty", + "ObservedPropertyIdentity", + "ObservedRelationship", + "RuntimeEvidence", + "RuntimeMetadata", + "RuntimeTag", + "map_unity_table_metadata", + "observe_unity_table", + "serialize_observed_state", +] From 8895f16fd2cc086be2f897a8d3e53909870ba715 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:05:48 +1000 Subject: [PATCH 04/28] test(runtime): add representative Unity observation fixture --- .../runtime/unity/orders_table_info.json | 89 +++++++++++++++++++ 1 file changed, 89 insertions(+) create mode 100644 tests/fixtures/runtime/unity/orders_table_info.json diff --git a/tests/fixtures/runtime/unity/orders_table_info.json b/tests/fixtures/runtime/unity/orders_table_info.json new file mode 100644 index 0000000..54103cc --- /dev/null +++ b/tests/fixtures/runtime/unity/orders_table_info.json @@ -0,0 +1,89 @@ +{ + "name": "orders", + "catalog_name": "main", + "schema_name": "silver", + "full_name": "main.silver.orders", + "table_id": "5d3dc5ba-9070-4f24-b21f-4a153fe2a31d", + "table_type": "MANAGED", + "data_source_format": "DELTA", + "owner": "data-platform@example.com", + "comment": "Current production orders table", + "storage_location": "abfss://silver@example.dfs.core.windows.net/orders", + "created_at": 1788057600000, + "updated_at": 1788144000000, + "properties": { + "delta.enableChangeDataFeed": "true", + "delta.minReaderVersion": "2" + }, + "tags": { + "domain": "sales", + "tier": "silver" + }, + "columns": [ + { + "name": "customer_id", + "type_name": "STRING", + "type_text": "string", + "type_json": "{\"type\":\"string\"}", + "position": 1, + "nullable": false, + "comment": "Customer identifier", + "tags": { + "classification": "identifier" + } + }, + { + "name": "note", + "type_name": "STRING", + "type_text": "string", + "position": 2, + "nullable": true, + "comment": "Optional order note" + }, + { + "name": "order_id", + "type_name": "BIGINT", + "type_text": "bigint", + "type_json": "{\"type\":\"long\"}", + "position": 0, + "nullable": false, + "comment": "Order identifier" + } + ], + "table_constraints": [ + { + "foreign_key_constraint": { + "name": "fk_orders_customer", + "child_columns": ["customer_id"], + "parent_table": "main.silver.customers", + "parent_columns": ["customer_id"], + "rely": true + } + }, + { + "primary_key_constraint": { + "name": "pk_orders", + "child_columns": ["order_id"], + "rely": true + } + } + ], + "view_dependencies": { + "dependencies": [ + { + "table": { + "table_full_name": "main.bronze.orders_raw" + } + }, + { + "function": { + "function_full_name": "main.shared.normalize_order" + } + } + ] + }, + "runtime_extension": { + "source": "fixture", + "revision": 7 + } +} From 6029e38e801627e2fa3ed499c3b87fb6eb7c75d7 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:06:09 +1000 Subject: [PATCH 05/28] test(runtime): protect Unity observation boundary --- tests/test_runtime_unity_observation.py | 188 ++++++++++++++++++++++++ 1 file changed, 188 insertions(+) create mode 100644 tests/test_runtime_unity_observation.py diff --git a/tests/test_runtime_unity_observation.py b/tests/test_runtime_unity_observation.py new file mode 100644 index 0000000..80da8bf --- /dev/null +++ b/tests/test_runtime_unity_observation.py @@ -0,0 +1,188 @@ +from __future__ import annotations + +import copy +import inspect +import json +from datetime import datetime, timezone +from pathlib import Path + +import pytest +from pydantic import ValidationError + +from semapact.runtime import unity as unity_runtime +from semapact.runtime.models import ObservedContractState, serialize_observed_state +from semapact.runtime.unity import map_unity_table_metadata, observe_unity_table + +FIXTURE = Path(__file__).parent / "fixtures" / "runtime" / "unity" / "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")) + + +def _metadata(entries) -> dict[str, object]: # noqa: ANN001 + return {item.key: json.loads(item.value_json) for item in entries} + + +def test_unity_metadata_maps_to_observed_state_without_odcs_projection() -> None: + payload = _payload() + original = copy.deepcopy(payload) + + state = map_unity_table_metadata( + payload, + source_identifier="https://adb.example/", + table_fqn="main.silver.orders", + captured_at=CAPTURED_AT, + ) + + assert isinstance(state, ObservedContractState) + assert state.platform == "databricks-unity-catalog" + assert state.source_identifier == "https://adb.example" + assert state.captured_at == CAPTURED_AT + assert state.fingerprint is None + assert payload == original + + assert len(state.assets) == 1 + asset = state.assets[0] + assert asset.identity.canonical_key == ( + "databricks-unity-catalog", + "main", + "silver", + "orders", + ) + assert asset.asset_type == "MANAGED" + assert asset.data_source_format == "DELTA" + assert asset.owner == "data-platform@example.com" + assert [(tag.key, tag.value) for tag in asset.tags] == [ + ("domain", "sales"), + ("tier", "silver"), + ] + + assert [item.identity.property for item in asset.properties] == [ + "order_id", + "customer_id", + "note", + ] + order_id = asset.properties[0] + assert order_id.identity.canonical_key == ( + "databricks-unity-catalog", + "main", + "silver", + "orders", + "order_id", + ) + assert order_id.physical_type == "bigint" + assert order_id.nullable is False + assert order_id.required is True + assert _metadata(order_id.metadata)["type_json"] == '{"type":"long"}' + + assert len(asset.constraints) == 1 + assert asset.constraints[0].constraint_type == "PRIMARY_KEY" + assert asset.constraints[0].name == "pk_orders" + assert asset.constraints[0].columns == ("order_id",) + assert _metadata(asset.constraints[0].metadata)["rely"] is True + + assert len(asset.relationships) == 1 + relationship = asset.relationships[0] + assert relationship.relationship_type == "FOREIGN_KEY" + assert relationship.source_columns == ("customer_id",) + assert relationship.target_reference == "main.silver.customers" + assert relationship.target_columns == ("customer_id",) + assert relationship.constraint_name == "fk_orders_customer" + + runtime_metadata = _metadata(asset.metadata) + assert runtime_metadata["table_id"] == "5d3dc5ba-9070-4f24-b21f-4a153fe2a31d" + assert runtime_metadata["storage_location"].endswith("/orders") + assert runtime_metadata["properties"]["delta.enableChangeDataFeed"] == "true" + assert runtime_metadata["runtime_extension"] == { + "revision": 7, + "source": "fixture", + } + + assert [(item.kind, item.reference) for item in state.evidence] == [ + ("view_function_dependency", "main.shared.normalize_order"), + ("view_table_dependency", "main.bronze.orders_raw"), + ] + + +def test_observed_state_serialization_is_stable_when_upstream_order_changes() -> None: + payload = _payload() + reordered = dict(reversed(list(payload.items()))) + reordered["columns"] = list(reversed(payload["columns"])) # type: ignore[index] + reordered["table_constraints"] = list( + reversed(payload["table_constraints"]) # type: ignore[index] + ) + dependencies = payload["view_dependencies"] # type: ignore[index] + assert isinstance(dependencies, dict) + reordered["view_dependencies"] = { + "dependencies": list(reversed(dependencies["dependencies"])) + } + + left = map_unity_table_metadata( + payload, + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + right = map_unity_table_metadata( + 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_unity_table_uses_runtime_fetcher_and_requires_no_live_workspace() -> None: + calls: list[tuple[str, str, str]] = [] + + def fake_fetcher( + workspace_url: str, token: str, table_fqn: str + ) -> dict[str, object]: + calls.append((workspace_url, token, table_fqn)) + return _payload() + + state = observe_unity_table( + table_fqn="main.silver.orders", + workspace_url="https://adb.example/", + token="test-token", + captured_at=CAPTURED_AT, + fetcher=fake_fetcher, + ) + + assert calls == [ + ("https://adb.example/", "test-token", "main.silver.orders") + ] + assert state.assets[0].identity.asset == "orders" + assert state.fingerprint is None + + +def test_observation_models_are_immutable() -> None: + state = map_unity_table_metadata( + _payload(), + source_identifier="https://adb.example", + captured_at=CAPTURED_AT, + ) + + with pytest.raises(ValidationError): + state.platform = "other-platform" # type: ignore[misc] + + +def test_runtime_observer_has_no_contract_or_governance_layer_dependency() -> None: + source = inspect.getsource(unity_runtime) + + assert "open_data_contract_standard" not in source + assert "datacontract" 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_unity_table_metadata( + _payload(), + source_identifier="https://adb.example", + captured_at=datetime(2026, 8, 30, 3, 0), + ) From 9136a299de4ca0ff226cbf65a4a2e3f7f2585fa9 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:06:46 +1000 Subject: [PATCH 06/28] docs(architecture): distinguish governed and observed state --- .agents/skills/semapact-system/SKILL.md | 28 +++++++++++++++++-------- 1 file changed, 19 insertions(+), 9 deletions(-) diff --git a/.agents/skills/semapact-system/SKILL.md b/.agents/skills/semapact-system/SKILL.md index 73af944..940f330 100644 --- a/.agents/skills/semapact-system/SKILL.md +++ b/.agents/skills/semapact-system/SKILL.md @@ -18,31 +18,41 @@ 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 (Delta Tables, Spark DDL, Unity Catalog) 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 runtime 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. + - Runtime platform metadata must not become governed truth merely because it was observed. -### C. Lifecycle Governance Layer +### C. Runtime Observation Model +- **Role:** Represents point-in-time external runtime state for assurance and reconciliation workflows. +- **Rules:** + - `ObservedContractState` is a read-side model, not an alternative canonical contract format. + - Runtime identity is platform-local and must remain distinct from ODCS contract identity. + - Observation must not invoke lifecycle merge, governance evaluation, release mutation, or platform writeback. + - 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 runtime 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. From 3302c1b8f050de9a466482df2c5c54794df2e1a0 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:07:35 +1000 Subject: [PATCH 07/28] docs(architecture): define runtime observation boundary --- ARCHITECTURE.md | 63 ++++++++++++++++++++++++++++++++++++++++--------- 1 file changed, 52 insertions(+), 11 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 0b367d8..674996d 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -97,7 +97,31 @@ Important boundary: - `GovernanceService` creates `ChangeContext` - lifecycle/governance lower layers consume that context and must not regenerate it -### 5. Exporters +### 5. Runtime Observation + +Location: + +- `semapact/runtime/` + +Responsibilities: + +- represent point-in-time external runtime state independently from governed contracts +- observe platform-local assets without creating or mutating ODCS contracts +- preserve runtime-only metadata and evidence for later assurance/reconciliation + +Key modules: + +- `semapact/runtime/models.py` +- `semapact/runtime/unity.py` + +Important boundary: + +- observed runtime state is read-side evidence, not governed truth +- runtime asset identity is platform-local and distinct from ODCS contract identity +- 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 +139,7 @@ Key modules: - `semapact/quality/ge_exporter.py` - `semapact/exporters/sql_exporter.py` -### 6. Orchestration +### 7. Orchestration Location: @@ -130,7 +154,7 @@ Key module: - `semapact/orchestrator/pipeline.py` -### 7. Interfaces +### 8. Interfaces Location: @@ -148,14 +172,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. + +Runtime observations are intentionally different. `ObservedContractState` 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 +ObservedContractState += observed runtime 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 runtime state does not implicitly perform that conversion. ## Root Contract Governance @@ -367,10 +405,11 @@ 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 +- runtime observation is separate read-side evidence and cannot become governed truth implicitly - 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 +435,8 @@ Precedence: ## Known Next Steps +- separate Unity Catalog discovery, observation, and contract import application workflows +- add stable observed-state fingerprints and reconciliation semantics - 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 From 8800fa9f50b8c85189fb3b99ab0bbd31fe82c03e Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:08:58 +1000 Subject: [PATCH 08/28] refactor(runtime): keep Unity observation mapper narrow --- semapact/runtime/unity.py | 428 +++++++++++--------------------------- 1 file changed, 126 insertions(+), 302 deletions(-) diff --git a/semapact/runtime/unity.py b/semapact/runtime/unity.py index 65aa95a..043aa52 100644 --- a/semapact/runtime/unity.py +++ b/semapact/runtime/unity.py @@ -26,6 +26,32 @@ UNITY_CATALOG_PLATFORM = "databricks-unity-catalog" UnityMetadataFetcher = Callable[[str, str, str], Mapping[str, Any]] +_ASSET_FIELDS = { + "catalog_name", + "schema_name", + "name", + "full_name", + "table_type", + "data_source_format", + "owner", + "comment", + "columns", + "table_constraints", + "tags", + "view_dependencies", +} +_COLUMN_FIELDS = { + "name", + "logical_type", + "type_text", + "type_name", + "nullable", + "required", + "position", + "comment", + "tags", +} + def observe_unity_table( *, @@ -35,22 +61,21 @@ def observe_unity_table( captured_at: datetime | None = None, fetcher: UnityMetadataFetcher | None = None, ) -> ObservedContractState: - """Observe one Unity Catalog table without creating or mutating an ODCS contract.""" + """Observe one Unity Catalog table without creating or mutating ODCS state.""" + if not table_fqn: + raise ValueError("table_fqn is required for Unity Catalog observation") if not workspace_url: raise ValueError("workspace_url is required for Unity Catalog observation") if not token: raise ValueError("token is required for Unity Catalog observation") - if not table_fqn: - raise ValueError("table_fqn is required for Unity Catalog observation") observed_at = captured_at or datetime.now(timezone.utc) - if observed_at.tzinfo is None or observed_at.utcoffset() is None: - raise ValueError("captured_at must be timezone-aware") + _require_aware_datetime(observed_at) metadata_fetcher = fetcher or fetch_unity_table_metadata - payload = metadata_fetcher(workspace_url, token, table_fqn) + metadata = metadata_fetcher(workspace_url, token, table_fqn) return map_unity_table_metadata( - payload, + metadata, source_identifier=workspace_url.rstrip("/"), table_fqn=table_fqn, captured_at=observed_at, @@ -62,20 +87,14 @@ def fetch_unity_table_metadata( token: str, table_fqn: str, ) -> Mapping[str, Any]: - """Fetch the Unity Catalog TableInfo payload for one exact table.""" - if not workspace_url or not token: - raise ValueError("workspace_url and token are required for Unity observation") - + """Fetch one Unity Catalog ``TableInfo`` object.""" endpoint = ( f"{workspace_url.rstrip('/')}/api/2.1/unity-catalog/tables/" f"{quote(table_fqn, safe='')}" ) request = Request( endpoint, - headers={ - "Authorization": f"Bearer {token}", - "Accept": "application/json", - }, + headers={"Authorization": f"Bearer {token}", "Accept": "application/json"}, method="GET", ) @@ -104,86 +123,49 @@ def map_unity_table_metadata( captured_at: datetime, table_fqn: str | None = None, ) -> ObservedContractState: - """Map a Unity TableInfo-style payload into immutable observed runtime state.""" - if captured_at.tzinfo is None or captured_at.utcoffset() is None: - raise ValueError("captured_at must be timezone-aware") - + """Map a Unity Catalog ``TableInfo`` payload to observed runtime state.""" + _require_aware_datetime(captured_at) identity = _asset_identity(metadata, table_fqn=table_fqn) - properties = _properties(metadata, identity=identity) constraints, relationships = _constraints_and_relationships(metadata) - evidence = _runtime_evidence(metadata) asset = ObservedAsset( identity=identity, - asset_type=_text(metadata.get("table_type") or metadata.get("tableType")), - data_source_format=_text( - metadata.get("data_source_format") or metadata.get("dataSourceFormat") - ), + asset_type=_text(metadata.get("table_type")), + data_source_format=_text(metadata.get("data_source_format")), owner=_text(metadata.get("owner")), comment=_text(metadata.get("comment")), - properties=properties, + properties=_properties(metadata.get("columns"), identity=identity), constraints=constraints, relationships=relationships, tags=_tags(metadata.get("tags")), - metadata=_metadata_entries( - metadata, - excluded={ - "catalog_name", - "catalogName", - "schema_name", - "schemaName", - "name", - "full_name", - "fullName", - "table_type", - "tableType", - "data_source_format", - "dataSourceFormat", - "owner", - "comment", - "columns", - "table_constraints", - "tableConstraints", - "constraints", - "foreign_keys", - "foreignKeys", - "tags", - "view_dependencies", - "viewDependencies", - }, - ), + metadata=_metadata_entries(metadata, excluded=_ASSET_FIELDS), ) - return ObservedContractState( platform=UNITY_CATALOG_PLATFORM, source_identifier=source_identifier.rstrip("/"), assets=(asset,), captured_at=captured_at, fingerprint=None, - evidence=evidence, + evidence=_runtime_evidence(metadata.get("view_dependencies")), ) def _asset_identity( metadata: Mapping[str, Any], *, table_fqn: str | None ) -> ObservedAssetIdentity: - catalog = _text(metadata.get("catalog_name") or metadata.get("catalogName")) - schema = _text(metadata.get("schema_name") or metadata.get("schemaName")) + 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 metadata.get("fullName") or table_fqn - ) + full_name = _text(metadata.get("full_name") or table_fqn) fallback = _split_table_fqn(full_name) if full_name else None - if fallback is not 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( - "Unity metadata must identify catalog, schema, and table name" - ) + raise ValueError("Unity metadata must identify catalog, schema, and table name") return ObservedAssetIdentity( platform=UNITY_CATALOG_PLATFORM, @@ -194,285 +176,153 @@ def _asset_identity( def _split_table_fqn(value: str) -> tuple[str, str, str] | None: - parts = [part.strip() for part in value.split(".")] + parts = tuple(part.strip() for part in value.split(".")) if len(parts) != 3 or not all(parts): return None - return parts[0], parts[1], parts[2] + return parts def _properties( - metadata: Mapping[str, Any], *, identity: ObservedAssetIdentity + value: Any, *, identity: ObservedAssetIdentity ) -> tuple[ObservedProperty, ...]: - raw_columns = metadata.get("columns") - if not isinstance(raw_columns, list): + if not isinstance(value, list): return () properties: list[ObservedProperty] = [] - for item in raw_columns: + for item in value: if not isinstance(item, Mapping): continue name = _text(item.get("name")) if not name: continue - nullable = _bool_or_none(item.get("nullable")) - explicit_required = _bool_or_none(item.get("required")) - required = explicit_required if explicit_required is not None else ( - None if nullable is None else not nullable + nullable = item.get("nullable") if isinstance(item.get("nullable"), bool) else None + explicit_required = ( + item.get("required") if isinstance(item.get("required"), bool) else None ) + required = ( + explicit_required + if explicit_required is not None + else None if nullable is None else not nullable + ) + position = item.get("position") + if isinstance(position, bool) or not isinstance(position, int): + position = None properties.append( ObservedProperty( identity=ObservedPropertyIdentity(asset=identity, property=name), - logical_type=_text( - item.get("logical_type") or item.get("logicalType") - ), - physical_type=_text( - item.get("type_text") - or item.get("typeText") - or item.get("type_name") - or item.get("typeName") - or item.get("type") - ), + logical_type=_text(item.get("logical_type")), + physical_type=_text(item.get("type_text") or item.get("type_name")), nullable=nullable, required=required, - position=_int_or_none(item.get("position")), + position=position, comment=_text(item.get("comment")), tags=_tags(item.get("tags")), - metadata=_metadata_entries( - item, - excluded={ - "name", - "logical_type", - "logicalType", - "type_text", - "typeText", - "type_name", - "typeName", - "type", - "nullable", - "required", - "position", - "comment", - "tags", - }, - ), + metadata=_metadata_entries(item, excluded=_COLUMN_FIELDS), ) ) - properties.sort(key=_property_sort_key) + properties.sort( + key=lambda item: ( + item.position is None, + item.position if item.position is not None else 0, + item.identity.property.casefold(), + ) + ) return tuple(properties) -def _property_sort_key(item: ObservedProperty) -> tuple[int, int, str]: - if item.position is None: - return (1, 0, item.identity.property.casefold()) - return (0, item.position, item.identity.property.casefold()) - - def _constraints_and_relationships( metadata: Mapping[str, Any], ) -> tuple[tuple[ObservedConstraint, ...], tuple[ObservedRelationship, ...]]: + value = metadata.get("table_constraints") + if not isinstance(value, list): + return (), () + constraints: list[ObservedConstraint] = [] relationships: list[ObservedRelationship] = [] + for item in value: + if not isinstance(item, Mapping): + continue - for item in _constraint_items(metadata): - primary = _mapping(item.get("primary_key_constraint") or item.get("primaryKeyConstraint")) - foreign = _mapping(item.get("foreign_key_constraint") or item.get("foreignKeyConstraint")) - - if primary is not None: + primary = item.get("primary_key_constraint") + if isinstance(primary, Mapping): constraints.append( ObservedConstraint( constraint_type="PRIMARY_KEY", name=_text(primary.get("name")), - columns=_string_tuple( - primary.get("child_columns") or primary.get("childColumns") - ), + columns=_string_tuple(primary.get("child_columns")), metadata=_metadata_entries( primary, - excluded={"name", "child_columns", "childColumns"}, + excluded={"name", "child_columns"}, ), ) ) continue - if foreign is not None: + foreign = item.get("foreign_key_constraint") + if isinstance(foreign, Mapping): relationships.append( ObservedRelationship( relationship_type="FOREIGN_KEY", - source_columns=_string_tuple( - foreign.get("child_columns") or foreign.get("childColumns") - ), - target_reference=_text( - foreign.get("parent_table") or foreign.get("parentTable") - ), - target_columns=_string_tuple( - foreign.get("parent_columns") or foreign.get("parentColumns") - ), + source_columns=_string_tuple(foreign.get("child_columns")), + target_reference=_text(foreign.get("parent_table")), + target_columns=_string_tuple(foreign.get("parent_columns")), constraint_name=_text(foreign.get("name")), metadata=_metadata_entries( foreign, excluded={ "name", "child_columns", - "childColumns", "parent_table", - "parentTable", "parent_columns", - "parentColumns", }, ), ) ) continue - constraint_type = _text( - item.get("constraint_type") - or item.get("constraintType") - or item.get("type") - or item.get("kind") - ) - normalized_type = (constraint_type or "UNKNOWN").upper().replace(" ", "_") - source_columns = _string_tuple( - item.get("columns") - or item.get("column_names") - or item.get("columnNames") - or item.get("from_columns") - or item.get("fromColumns") - or item.get("child_columns") - or item.get("childColumns") - or item.get("from") - ) - target_reference = _text( - item.get("referenced_table") - or item.get("referencedTable") - or item.get("to_table") - or item.get("toTable") - or item.get("parent_table") - or item.get("parentTable") - ) - target_columns = _string_tuple( - item.get("referenced_columns") - or item.get("referencedColumns") - or item.get("to_columns") - or item.get("toColumns") - or item.get("parent_columns") - or item.get("parentColumns") - or item.get("to") - ) - - if "FOREIGN" in normalized_type or target_reference: - relationships.append( - ObservedRelationship( - relationship_type=normalized_type, - source_columns=source_columns, - target_reference=target_reference, - target_columns=target_columns, - constraint_name=_text(item.get("name")), - metadata=_metadata_entries( - item, - excluded=_FLAT_CONSTRAINT_KEYS, - ), - ) - ) - else: - constraints.append( - ObservedConstraint( - constraint_type=normalized_type, - name=_text(item.get("name")), - columns=source_columns, - expression=_text( - item.get("expression") or item.get("check_expression") - ), - metadata=_metadata_entries( - item, - excluded=_FLAT_CONSTRAINT_KEYS | {"expression", "check_expression"}, - ), - ) + constraints.append( + ObservedConstraint( + constraint_type="UNKNOWN", + metadata=(RuntimeMetadata(key="raw", value_json=_canonical_json(item)),), ) + ) constraints.sort( key=lambda item: ( item.constraint_type, (item.name or "").casefold(), - tuple(value.casefold() for value in item.columns), + tuple(column.casefold() for column in item.columns), + tuple(entry.value_json for entry in item.metadata), ) ) relationships.sort( key=lambda item: ( - item.relationship_type, (item.constraint_name or "").casefold(), (item.target_reference or "").casefold(), - tuple(value.casefold() for value in item.source_columns), + tuple(column.casefold() for column in item.source_columns), ) ) return tuple(constraints), tuple(relationships) -_FLAT_CONSTRAINT_KEYS = { - "constraint_type", - "constraintType", - "type", - "kind", - "name", - "columns", - "column_names", - "columnNames", - "from_columns", - "fromColumns", - "child_columns", - "childColumns", - "from", - "referenced_table", - "referencedTable", - "to_table", - "toTable", - "parent_table", - "parentTable", - "referenced_columns", - "referencedColumns", - "to_columns", - "toColumns", - "parent_columns", - "parentColumns", - "to", -} - - -def _constraint_items(metadata: Mapping[str, Any]) -> tuple[Mapping[str, Any], ...]: - items: list[Mapping[str, Any]] = [] - for key in ( - "table_constraints", - "tableConstraints", - "constraints", - "foreign_keys", - "foreignKeys", - ): - value = metadata.get(key) - if isinstance(value, list): - items.extend(item for item in value if isinstance(item, Mapping)) - return tuple(items) - - -def _runtime_evidence(metadata: Mapping[str, Any]) -> tuple[RuntimeEvidence, ...]: - raw = metadata.get("view_dependencies") or metadata.get("viewDependencies") - dependency_container = _mapping(raw) - if dependency_container is None: +def _runtime_evidence(value: Any) -> tuple[RuntimeEvidence, ...]: + if not isinstance(value, Mapping): return () - dependencies = dependency_container.get("dependencies") + dependencies = value.get("dependencies") if not isinstance(dependencies, list): return () evidence: list[RuntimeEvidence] = [] - for dependency in dependencies: - if not isinstance(dependency, Mapping): + for item in dependencies: + if not isinstance(item, Mapping): continue - table = _mapping(dependency.get("table")) - function = _mapping(dependency.get("function")) - if table is not None: - reference = _text( - table.get("table_full_name") or table.get("tableFullName") - ) + table = item.get("table") + function = item.get("function") + if isinstance(table, Mapping): + reference = _text(table.get("table_full_name")) if reference: evidence.append( RuntimeEvidence( @@ -480,15 +330,12 @@ def _runtime_evidence(metadata: Mapping[str, Any]) -> tuple[RuntimeEvidence, ... reference=reference, metadata=_metadata_entries( table, - excluded={"table_full_name", "tableFullName"}, + excluded={"table_full_name"}, ), ) ) - elif function is not None: - reference = _text( - function.get("function_full_name") - or function.get("functionFullName") - ) + elif isinstance(function, Mapping): + reference = _text(function.get("function_full_name")) if reference: evidence.append( RuntimeEvidence( @@ -496,7 +343,7 @@ def _runtime_evidence(metadata: Mapping[str, Any]) -> tuple[RuntimeEvidence, ... reference=reference, metadata=_metadata_entries( function, - excluded={"function_full_name", "functionFullName"}, + excluded={"function_full_name"}, ), ) ) @@ -506,22 +353,13 @@ def _runtime_evidence(metadata: Mapping[str, Any]) -> tuple[RuntimeEvidence, ... def _tags(value: Any) -> tuple[RuntimeTag, ...]: - tags: list[RuntimeTag] = [] - if isinstance(value, Mapping): - for key, item_value in value.items(): - tags.append(RuntimeTag(key=str(key), value=str(item_value))) - elif isinstance(value, list): - for item in value: - if not isinstance(item, Mapping): - continue - key = _text(item.get("key") or item.get("tag_name") or item.get("tagName")) - item_value = _text( - item.get("value") or item.get("tag_value") or item.get("tagValue") - ) - if key and item_value is not None: - tags.append(RuntimeTag(key=key, value=item_value)) - tags.sort(key=lambda item: (item.key.casefold(), item.value)) - return tuple(tags) + if not isinstance(value, Mapping): + return () + tags = tuple( + RuntimeTag(key=str(key), value=str(tag_value)) + for key, tag_value in sorted(value.items(), key=lambda pair: str(pair[0]).casefold()) + ) + return tags def _metadata_entries( @@ -545,8 +383,12 @@ def _canonical_json(value: Any) -> str: ) -def _mapping(value: Any) -> Mapping[str, Any] | None: - return value if isinstance(value, Mapping) else None +def _string_tuple(value: Any) -> tuple[str, ...]: + if isinstance(value, str): + return (value,) + if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): + return tuple(text for item in value if (text := _text(item)) is not None) + return () def _text(value: Any) -> str | None: @@ -556,24 +398,6 @@ def _text(value: Any) -> str | None: return text or None -def _string_tuple(value: Any) -> tuple[str, ...]: - if value is None: - return () - if isinstance(value, str): - return tuple(part.strip() for part in value.split(",") if part.strip()) - if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): - result = tuple(text for item in value if (text := _text(item)) is not None) - return result - return () - - -def _bool_or_none(value: Any) -> bool | None: - if isinstance(value, bool): - return value - return None - - -def _int_or_none(value: Any) -> int | None: - if isinstance(value, bool): - return None - return value if isinstance(value, int) else 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") From 3b4e63c2fd928db46b27eebb787dd8bc090a8f4e Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:49:33 +1000 Subject: [PATCH 09/28] refactor(observation): add platform-neutral observed state model --- semapact/observation/models.py | 94 ++++++++++++++++++++++++++++++++++ 1 file changed, 94 insertions(+) create mode 100644 semapact/observation/models.py diff --git a/semapact/observation/models.py b/semapact/observation/models.py new file mode 100644 index 0000000..2bbab4f --- /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 runtime 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, + ) From 1d0e71e18da9ecddd7ecf401a4e96c9210873f96 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:49:41 +1000 Subject: [PATCH 10/28] refactor(observation): expose platform-neutral read model --- semapact/observation/__init__.py | 19 +++++++++++++++++++ 1 file changed, 19 insertions(+) create mode 100644 semapact/observation/__init__.py 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", +] From a72319c39527755e02fd5b28d1ae8fbe6c54f9d3 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:50:09 +1000 Subject: [PATCH 11/28] refactor(observation): map Databricks SDK TableInfo directly --- semapact/observation/databricks.py | 183 +++++++++++++++++++++++++++++ 1 file changed, 183 insertions(+) create mode 100644 semapact/observation/databricks.py diff --git a/semapact/observation/databricks.py b/semapact/observation/databricks.py new file mode 100644 index 0000000..993aaa5 --- /dev/null +++ b/semapact/observation/databricks.py @@ -0,0 +1,183 @@ +"""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. +""" + +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 official ``databricks.sdk.WorkspaceClient`` + in production. It 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") From 4ee9ddf13967ac9a73abfac777065d275cc734a5 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:50:27 +1000 Subject: [PATCH 12/28] test(observation): add minimal Databricks TableInfo fixture --- .../databricks/orders_table_info.json | 30 +++++++++++++++++++ 1 file changed, 30 insertions(+) create mode 100644 tests/fixtures/observation/databricks/orders_table_info.json 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 + } + ] +} From 32308725301a1404e65311f28fd389b856c855dc Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:50:48 +1000 Subject: [PATCH 13/28] test(observation): cover platform-neutral Databricks mapping --- tests/test_observation_databricks.py | 180 +++++++++++++++++++++++++++ 1 file changed, 180 insertions(+) create mode 100644 tests/test_observation_databricks.py diff --git a/tests/test_observation_databricks.py b/tests/test_observation_databricks.py new file mode 100644 index 0000000..908e1ad --- /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(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), + ) From 56dbd8871510968f923eb2f1f248e7eb10fe7e6b Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:51:15 +1000 Subject: [PATCH 14/28] build(databricks): reuse datacontract-cli platform extra --- pyproject.toml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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", From 5aaf11c509208821f25506ffb0560afc907a357a Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:51:46 +1000 Subject: [PATCH 15/28] docs(architecture): define platform-neutral observation boundary --- .agents/skills/semapact-system/SKILL.md | 19 +++++++++++-------- 1 file changed, 11 insertions(+), 8 deletions(-) diff --git a/.agents/skills/semapact-system/SKILL.md b/.agents/skills/semapact-system/SKILL.md index 940f330..7da2e56 100644 --- a/.agents/skills/semapact-system/SKILL.md +++ b/.agents/skills/semapact-system/SKILL.md @@ -18,31 +18,34 @@ 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 when a contract import is explicitly requested. +- **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 runtime observation; observing a platform must not implicitly create or mutate an ODCS contract. + - Contract import is distinct from platform observation; observing a platform must not implicitly create or mutate an ODCS contract. ### 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 for governed contract state. - - Runtime platform metadata must not become governed truth merely because it was observed. + - External platform state must not become governed truth merely because it was observed. -### C. Runtime Observation Model -- **Role:** Represents point-in-time external runtime state for assurance and reconciliation workflows. +### C. Platform Observation Model +- **Role:** Represents point-in-time external platform state for assurance and reconciliation workflows. - **Rules:** - - `ObservedContractState` is a read-side model, not an alternative canonical contract format. - - Runtime identity is platform-local and must remain distinct from ODCS contract identity. + - `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, ingestion, and runtime observation layers. + - It must remain fully decoupled from the UI, ingestion, and platform observation layers. ### E. Export Layer - **Role:** Converts contracts to downstream assets (Great Expectations suites, Spark DDL, Graph cypher). From c08a172dc35e3ca49332287d9335417e88178242 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:52:34 +1000 Subject: [PATCH 16/28] docs(architecture): rename runtime module to platform observation --- ARCHITECTURE.md | 39 ++++++++++++++++++++++++--------------- 1 file changed, 24 insertions(+), 15 deletions(-) diff --git a/ARCHITECTURE.md b/ARCHITECTURE.md index 674996d..9ba0bdf 100644 --- a/ARCHITECTURE.md +++ b/ARCHITECTURE.md @@ -97,27 +97,32 @@ Important boundary: - `GovernanceService` creates `ChangeContext` - lifecycle/governance lower layers consume that context and must not regenerate it -### 5. Runtime Observation +### 5. Platform Observation Location: -- `semapact/runtime/` +- `semapact/observation/` Responsibilities: -- represent point-in-time external runtime state independently from governed contracts -- observe platform-local assets without creating or mutating ODCS contracts -- preserve runtime-only metadata and evidence for later assurance/reconciliation +- 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/runtime/models.py` -- `semapact/runtime/unity.py` +- `semapact/observation/models.py` +- `semapact/observation/databricks.py` Important boundary: -- observed runtime state is read-side evidence, not governed truth -- runtime asset identity is platform-local and distinct from ODCS contract identity +- `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 @@ -181,7 +186,7 @@ SemaPact assumes: 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. -Runtime observations are intentionally different. `ObservedContractState` 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. +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: @@ -189,11 +194,11 @@ Conceptually: Approved ODCS Contract = desired governed state -ObservedContractState -= observed runtime state +ObservedPlatformState += observed external platform state ``` -Converting external metadata into a new ODCS contract is an explicit import workflow. Observing runtime state does not implicitly perform that conversion. +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 @@ -407,7 +412,9 @@ Precedence: - main governed contract is canonical and immutable from presentation paths - ODCS is the canonical model for governed desired contract state -- runtime observation is separate read-side evidence and cannot become governed truth implicitly +- 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 - datacontract-cli is reused where possible instead of reimplemented @@ -435,8 +442,10 @@ Precedence: ## Known Next Steps -- separate Unity Catalog discovery, observation, and contract import application workflows +- 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 governed contract helper logic toward ODCS model-driven behavior From c5b802e5e77f8d2d2fefc2241d61d18c0a09f8e3 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:52:41 +1000 Subject: [PATCH 17/28] refactor(observation): remove runtime package alias --- semapact/runtime/__init__.py | 37 ------------------------------------ 1 file changed, 37 deletions(-) delete mode 100644 semapact/runtime/__init__.py diff --git a/semapact/runtime/__init__.py b/semapact/runtime/__init__.py deleted file mode 100644 index 166b747..0000000 --- a/semapact/runtime/__init__.py +++ /dev/null @@ -1,37 +0,0 @@ -"""Runtime observation boundary for SemaPact read-side state.""" - -from semapact.runtime.models import ( - ObservedAsset, - ObservedAssetIdentity, - ObservedConstraint, - ObservedContractState, - ObservedProperty, - ObservedPropertyIdentity, - ObservedRelationship, - RuntimeEvidence, - RuntimeMetadata, - RuntimeTag, - serialize_observed_state, -) -from semapact.runtime.unity import ( - UNITY_CATALOG_PLATFORM, - map_unity_table_metadata, - observe_unity_table, -) - -__all__ = [ - "UNITY_CATALOG_PLATFORM", - "ObservedAsset", - "ObservedAssetIdentity", - "ObservedConstraint", - "ObservedContractState", - "ObservedProperty", - "ObservedPropertyIdentity", - "ObservedRelationship", - "RuntimeEvidence", - "RuntimeMetadata", - "RuntimeTag", - "map_unity_table_metadata", - "observe_unity_table", - "serialize_observed_state", -] From dc1137c2084580b334152ac5f701cf33683e268d Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:52:47 +1000 Subject: [PATCH 18/28] refactor(observation): remove runtime-specific observed model --- semapact/runtime/models.py | 146 ------------------------------------- 1 file changed, 146 deletions(-) delete mode 100644 semapact/runtime/models.py diff --git a/semapact/runtime/models.py b/semapact/runtime/models.py deleted file mode 100644 index 1e77e49..0000000 --- a/semapact/runtime/models.py +++ /dev/null @@ -1,146 +0,0 @@ -"""Read-side domain models for observed runtime state. - -Observed runtime state is deliberately separate from the governed ODCS contract model. -It describes what a platform reports now; it is never governed truth by itself. -""" - -from __future__ import annotations - -import json -from datetime import datetime - -from pydantic import BaseModel, ConfigDict - - -class RuntimeModel(BaseModel): - """Shared immutable base for runtime observation models.""" - - model_config = ConfigDict(frozen=True, extra="forbid") - - -class ObservedAssetIdentity(RuntimeModel): - """Platform-local identity for one observed runtime asset.""" - - platform: str - catalog: str - schema: str - asset: str - - @property - def canonical_key(self) -> tuple[str, str, str, str]: - """Return the case-normalized runtime identity key.""" - return ( - self.platform.casefold(), - self.catalog.casefold(), - self.schema.casefold(), - self.asset.casefold(), - ) - - -class ObservedPropertyIdentity(RuntimeModel): - """Identity for a property within an observed runtime asset.""" - - asset: ObservedAssetIdentity - property: str - - @property - def canonical_key(self) -> tuple[str, str, str, str, str]: - """Return the case-normalized runtime property identity key.""" - return (*self.asset.canonical_key, self.property.casefold()) - - -class RuntimeMetadata(RuntimeModel): - """Canonical preservation of platform metadata not modeled directly.""" - - key: str - value_json: str - - -class RuntimeTag(RuntimeModel): - """Observed platform tag.""" - - key: str - value: str - - -class ObservedProperty(RuntimeModel): - """Observed runtime property/column state.""" - - identity: ObservedPropertyIdentity - logical_type: str | None = None - physical_type: str | None = None - nullable: bool | None = None - required: bool | None = None - position: int | None = None - comment: str | None = None - tags: tuple[RuntimeTag, ...] = () - metadata: tuple[RuntimeMetadata, ...] = () - - -class ObservedConstraint(RuntimeModel): - """Observed non-relational table constraint.""" - - constraint_type: str - name: str | None = None - columns: tuple[str, ...] = () - expression: str | None = None - metadata: tuple[RuntimeMetadata, ...] = () - - -class ObservedRelationship(RuntimeModel): - """Observed relationship from one runtime asset to another.""" - - relationship_type: str - source_columns: tuple[str, ...] = () - target_reference: str | None = None - target_columns: tuple[str, ...] = () - constraint_name: str | None = None - metadata: tuple[RuntimeMetadata, ...] = () - - -class RuntimeEvidence(RuntimeModel): - """Reference to runtime evidence reported by the source platform.""" - - kind: str - reference: str - metadata: tuple[RuntimeMetadata, ...] = () - - -class ObservedAsset(RuntimeModel): - """Observed state for one runtime asset.""" - - identity: ObservedAssetIdentity - asset_type: str | None = None - data_source_format: str | None = None - owner: str | None = None - comment: str | None = None - properties: tuple[ObservedProperty, ...] = () - constraints: tuple[ObservedConstraint, ...] = () - relationships: tuple[ObservedRelationship, ...] = () - tags: tuple[RuntimeTag, ...] = () - metadata: tuple[RuntimeMetadata, ...] = () - - -class ObservedContractState(RuntimeModel): - """Point-in-time runtime observation independent from governed ODCS state. - - ``fingerprint`` is intentionally optional. M1 fingerprint normalization and - hashing are a separate capability; this model only reserves the contract. - """ - - platform: str - source_identifier: str - assets: tuple[ObservedAsset, ...] - captured_at: datetime - fingerprint: str | None = None - evidence: tuple[RuntimeEvidence, ...] = () - - -def serialize_observed_state(state: ObservedContractState) -> str: - """Serialize observed state deterministically for machine-readable use.""" - return json.dumps( - state.model_dump(mode="json"), - ensure_ascii=False, - separators=(",", ":"), - sort_keys=True, - ) From 1fd72204fd02e5057e01e030014a25a6e4d5ac84 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:52:55 +1000 Subject: [PATCH 19/28] refactor(observation): remove handwritten Unity REST observer --- semapact/runtime/unity.py | 403 -------------------------------------- 1 file changed, 403 deletions(-) delete mode 100644 semapact/runtime/unity.py diff --git a/semapact/runtime/unity.py b/semapact/runtime/unity.py deleted file mode 100644 index 043aa52..0000000 --- a/semapact/runtime/unity.py +++ /dev/null @@ -1,403 +0,0 @@ -"""Unity Catalog read-side observation without ODCS contract generation.""" - -from __future__ import annotations - -import json -from datetime import datetime, timezone -from typing import Any, Callable, Mapping, Sequence -from urllib.error import HTTPError, URLError -from urllib.parse import quote -from urllib.request import Request, urlopen - -from semapact.exceptions import StorageError -from semapact.runtime.models import ( - ObservedAsset, - ObservedAssetIdentity, - ObservedConstraint, - ObservedContractState, - ObservedProperty, - ObservedPropertyIdentity, - ObservedRelationship, - RuntimeEvidence, - RuntimeMetadata, - RuntimeTag, -) - -UNITY_CATALOG_PLATFORM = "databricks-unity-catalog" -UnityMetadataFetcher = Callable[[str, str, str], Mapping[str, Any]] - -_ASSET_FIELDS = { - "catalog_name", - "schema_name", - "name", - "full_name", - "table_type", - "data_source_format", - "owner", - "comment", - "columns", - "table_constraints", - "tags", - "view_dependencies", -} -_COLUMN_FIELDS = { - "name", - "logical_type", - "type_text", - "type_name", - "nullable", - "required", - "position", - "comment", - "tags", -} - - -def observe_unity_table( - *, - table_fqn: str, - workspace_url: str, - token: str, - captured_at: datetime | None = None, - fetcher: UnityMetadataFetcher | None = None, -) -> ObservedContractState: - """Observe one Unity Catalog table without creating or mutating ODCS state.""" - if not table_fqn: - raise ValueError("table_fqn is required for Unity Catalog observation") - if not workspace_url: - raise ValueError("workspace_url is required for Unity Catalog observation") - if not token: - raise ValueError("token is required for Unity Catalog observation") - - observed_at = captured_at or datetime.now(timezone.utc) - _require_aware_datetime(observed_at) - - metadata_fetcher = fetcher or fetch_unity_table_metadata - metadata = metadata_fetcher(workspace_url, token, table_fqn) - return map_unity_table_metadata( - metadata, - source_identifier=workspace_url.rstrip("/"), - table_fqn=table_fqn, - captured_at=observed_at, - ) - - -def fetch_unity_table_metadata( - workspace_url: str, - token: str, - table_fqn: str, -) -> Mapping[str, Any]: - """Fetch one Unity Catalog ``TableInfo`` object.""" - endpoint = ( - f"{workspace_url.rstrip('/')}/api/2.1/unity-catalog/tables/" - f"{quote(table_fqn, safe='')}" - ) - request = Request( - endpoint, - headers={"Authorization": f"Bearer {token}", "Accept": "application/json"}, - method="GET", - ) - - try: - with urlopen(request, timeout=8) as response: - payload = response.read().decode("utf-8") - except HTTPError as exc: - raise StorageError( - f"Unity table observation request failed: HTTP {exc.code}" - ) from exc - except URLError as exc: - raise StorageError( - f"Unity table observation request failed: {exc.reason}" - ) from exc - - parsed = json.loads(payload) - if not isinstance(parsed, dict): - raise StorageError("Unity table observation response is not a JSON object") - return parsed - - -def map_unity_table_metadata( - metadata: Mapping[str, Any], - *, - source_identifier: str, - captured_at: datetime, - table_fqn: str | None = None, -) -> ObservedContractState: - """Map a Unity Catalog ``TableInfo`` payload to observed runtime state.""" - _require_aware_datetime(captured_at) - identity = _asset_identity(metadata, table_fqn=table_fqn) - constraints, relationships = _constraints_and_relationships(metadata) - - asset = ObservedAsset( - identity=identity, - asset_type=_text(metadata.get("table_type")), - data_source_format=_text(metadata.get("data_source_format")), - owner=_text(metadata.get("owner")), - comment=_text(metadata.get("comment")), - properties=_properties(metadata.get("columns"), identity=identity), - constraints=constraints, - relationships=relationships, - tags=_tags(metadata.get("tags")), - metadata=_metadata_entries(metadata, excluded=_ASSET_FIELDS), - ) - return ObservedContractState( - platform=UNITY_CATALOG_PLATFORM, - source_identifier=source_identifier.rstrip("/"), - assets=(asset,), - captured_at=captured_at, - fingerprint=None, - evidence=_runtime_evidence(metadata.get("view_dependencies")), - ) - - -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("Unity metadata must identify catalog, schema, and table name") - - return ObservedAssetIdentity( - platform=UNITY_CATALOG_PLATFORM, - catalog=catalog, - schema=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 () - - properties: list[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 - explicit_required = ( - item.get("required") if isinstance(item.get("required"), bool) else None - ) - required = ( - explicit_required - if explicit_required is not None - else None if nullable is None else not nullable - ) - position = item.get("position") - if isinstance(position, bool) or not isinstance(position, int): - position = None - - properties.append( - ObservedProperty( - identity=ObservedPropertyIdentity(asset=identity, property=name), - logical_type=_text(item.get("logical_type")), - physical_type=_text(item.get("type_text") or item.get("type_name")), - nullable=nullable, - required=required, - position=position, - comment=_text(item.get("comment")), - tags=_tags(item.get("tags")), - metadata=_metadata_entries(item, excluded=_COLUMN_FIELDS), - ) - ) - - properties.sort( - key=lambda item: ( - item.position is None, - item.position if item.position is not None else 0, - item.identity.property.casefold(), - ) - ) - return tuple(properties) - - -def _constraints_and_relationships( - metadata: Mapping[str, Any], -) -> tuple[tuple[ObservedConstraint, ...], tuple[ObservedRelationship, ...]]: - value = metadata.get("table_constraints") - if not isinstance(value, list): - return (), () - - constraints: list[ObservedConstraint] = [] - relationships: list[ObservedRelationship] = [] - for item in value: - if not isinstance(item, Mapping): - continue - - primary = item.get("primary_key_constraint") - if isinstance(primary, Mapping): - constraints.append( - ObservedConstraint( - constraint_type="PRIMARY_KEY", - name=_text(primary.get("name")), - columns=_string_tuple(primary.get("child_columns")), - metadata=_metadata_entries( - primary, - excluded={"name", "child_columns"}, - ), - ) - ) - continue - - foreign = item.get("foreign_key_constraint") - if isinstance(foreign, Mapping): - relationships.append( - ObservedRelationship( - relationship_type="FOREIGN_KEY", - source_columns=_string_tuple(foreign.get("child_columns")), - target_reference=_text(foreign.get("parent_table")), - target_columns=_string_tuple(foreign.get("parent_columns")), - constraint_name=_text(foreign.get("name")), - metadata=_metadata_entries( - foreign, - excluded={ - "name", - "child_columns", - "parent_table", - "parent_columns", - }, - ), - ) - ) - continue - - constraints.append( - ObservedConstraint( - constraint_type="UNKNOWN", - metadata=(RuntimeMetadata(key="raw", value_json=_canonical_json(item)),), - ) - ) - - constraints.sort( - key=lambda item: ( - item.constraint_type, - (item.name or "").casefold(), - tuple(column.casefold() for column in item.columns), - tuple(entry.value_json for entry in item.metadata), - ) - ) - relationships.sort( - key=lambda item: ( - (item.constraint_name or "").casefold(), - (item.target_reference or "").casefold(), - tuple(column.casefold() for column in item.source_columns), - ) - ) - return tuple(constraints), tuple(relationships) - - -def _runtime_evidence(value: Any) -> tuple[RuntimeEvidence, ...]: - if not isinstance(value, Mapping): - return () - dependencies = value.get("dependencies") - if not isinstance(dependencies, list): - return () - - evidence: list[RuntimeEvidence] = [] - for item in dependencies: - if not isinstance(item, Mapping): - continue - table = item.get("table") - function = item.get("function") - if isinstance(table, Mapping): - reference = _text(table.get("table_full_name")) - if reference: - evidence.append( - RuntimeEvidence( - kind="view_table_dependency", - reference=reference, - metadata=_metadata_entries( - table, - excluded={"table_full_name"}, - ), - ) - ) - elif isinstance(function, Mapping): - reference = _text(function.get("function_full_name")) - if reference: - evidence.append( - RuntimeEvidence( - kind="view_function_dependency", - reference=reference, - metadata=_metadata_entries( - function, - excluded={"function_full_name"}, - ), - ) - ) - - evidence.sort(key=lambda item: (item.kind, item.reference.casefold())) - return tuple(evidence) - - -def _tags(value: Any) -> tuple[RuntimeTag, ...]: - if not isinstance(value, Mapping): - return () - tags = tuple( - RuntimeTag(key=str(key), value=str(tag_value)) - for key, tag_value in sorted(value.items(), key=lambda pair: str(pair[0]).casefold()) - ) - return tags - - -def _metadata_entries( - mapping: Mapping[str, Any], *, excluded: set[str] -) -> tuple[RuntimeMetadata, ...]: - entries = [ - RuntimeMetadata(key=str(key), value_json=_canonical_json(value)) - for key, value in mapping.items() - if key not in excluded - ] - entries.sort(key=lambda item: item.key.casefold()) - return tuple(entries) - - -def _canonical_json(value: Any) -> str: - return json.dumps( - value, - ensure_ascii=False, - separators=(",", ":"), - sort_keys=True, - ) - - -def _string_tuple(value: Any) -> tuple[str, ...]: - if isinstance(value, str): - return (value,) - if isinstance(value, Sequence) and not isinstance(value, (str, bytes, bytearray)): - return tuple(text for item in value if (text := _text(item)) is not None) - return () - - -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") From fe97f1a94c2205ed52842bde6f64371ffbea329a Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:53:01 +1000 Subject: [PATCH 20/28] test(observation): remove runtime-specific Unity tests --- tests/test_runtime_unity_observation.py | 188 ------------------------ 1 file changed, 188 deletions(-) delete mode 100644 tests/test_runtime_unity_observation.py diff --git a/tests/test_runtime_unity_observation.py b/tests/test_runtime_unity_observation.py deleted file mode 100644 index 80da8bf..0000000 --- a/tests/test_runtime_unity_observation.py +++ /dev/null @@ -1,188 +0,0 @@ -from __future__ import annotations - -import copy -import inspect -import json -from datetime import datetime, timezone -from pathlib import Path - -import pytest -from pydantic import ValidationError - -from semapact.runtime import unity as unity_runtime -from semapact.runtime.models import ObservedContractState, serialize_observed_state -from semapact.runtime.unity import map_unity_table_metadata, observe_unity_table - -FIXTURE = Path(__file__).parent / "fixtures" / "runtime" / "unity" / "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")) - - -def _metadata(entries) -> dict[str, object]: # noqa: ANN001 - return {item.key: json.loads(item.value_json) for item in entries} - - -def test_unity_metadata_maps_to_observed_state_without_odcs_projection() -> None: - payload = _payload() - original = copy.deepcopy(payload) - - state = map_unity_table_metadata( - payload, - source_identifier="https://adb.example/", - table_fqn="main.silver.orders", - captured_at=CAPTURED_AT, - ) - - assert isinstance(state, ObservedContractState) - assert state.platform == "databricks-unity-catalog" - assert state.source_identifier == "https://adb.example" - assert state.captured_at == CAPTURED_AT - assert state.fingerprint is None - assert payload == original - - assert len(state.assets) == 1 - asset = state.assets[0] - assert asset.identity.canonical_key == ( - "databricks-unity-catalog", - "main", - "silver", - "orders", - ) - assert asset.asset_type == "MANAGED" - assert asset.data_source_format == "DELTA" - assert asset.owner == "data-platform@example.com" - assert [(tag.key, tag.value) for tag in asset.tags] == [ - ("domain", "sales"), - ("tier", "silver"), - ] - - assert [item.identity.property for item in asset.properties] == [ - "order_id", - "customer_id", - "note", - ] - order_id = asset.properties[0] - assert order_id.identity.canonical_key == ( - "databricks-unity-catalog", - "main", - "silver", - "orders", - "order_id", - ) - assert order_id.physical_type == "bigint" - assert order_id.nullable is False - assert order_id.required is True - assert _metadata(order_id.metadata)["type_json"] == '{"type":"long"}' - - assert len(asset.constraints) == 1 - assert asset.constraints[0].constraint_type == "PRIMARY_KEY" - assert asset.constraints[0].name == "pk_orders" - assert asset.constraints[0].columns == ("order_id",) - assert _metadata(asset.constraints[0].metadata)["rely"] is True - - assert len(asset.relationships) == 1 - relationship = asset.relationships[0] - assert relationship.relationship_type == "FOREIGN_KEY" - assert relationship.source_columns == ("customer_id",) - assert relationship.target_reference == "main.silver.customers" - assert relationship.target_columns == ("customer_id",) - assert relationship.constraint_name == "fk_orders_customer" - - runtime_metadata = _metadata(asset.metadata) - assert runtime_metadata["table_id"] == "5d3dc5ba-9070-4f24-b21f-4a153fe2a31d" - assert runtime_metadata["storage_location"].endswith("/orders") - assert runtime_metadata["properties"]["delta.enableChangeDataFeed"] == "true" - assert runtime_metadata["runtime_extension"] == { - "revision": 7, - "source": "fixture", - } - - assert [(item.kind, item.reference) for item in state.evidence] == [ - ("view_function_dependency", "main.shared.normalize_order"), - ("view_table_dependency", "main.bronze.orders_raw"), - ] - - -def test_observed_state_serialization_is_stable_when_upstream_order_changes() -> None: - payload = _payload() - reordered = dict(reversed(list(payload.items()))) - reordered["columns"] = list(reversed(payload["columns"])) # type: ignore[index] - reordered["table_constraints"] = list( - reversed(payload["table_constraints"]) # type: ignore[index] - ) - dependencies = payload["view_dependencies"] # type: ignore[index] - assert isinstance(dependencies, dict) - reordered["view_dependencies"] = { - "dependencies": list(reversed(dependencies["dependencies"])) - } - - left = map_unity_table_metadata( - payload, - source_identifier="https://adb.example", - captured_at=CAPTURED_AT, - ) - right = map_unity_table_metadata( - 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_unity_table_uses_runtime_fetcher_and_requires_no_live_workspace() -> None: - calls: list[tuple[str, str, str]] = [] - - def fake_fetcher( - workspace_url: str, token: str, table_fqn: str - ) -> dict[str, object]: - calls.append((workspace_url, token, table_fqn)) - return _payload() - - state = observe_unity_table( - table_fqn="main.silver.orders", - workspace_url="https://adb.example/", - token="test-token", - captured_at=CAPTURED_AT, - fetcher=fake_fetcher, - ) - - assert calls == [ - ("https://adb.example/", "test-token", "main.silver.orders") - ] - assert state.assets[0].identity.asset == "orders" - assert state.fingerprint is None - - -def test_observation_models_are_immutable() -> None: - state = map_unity_table_metadata( - _payload(), - source_identifier="https://adb.example", - captured_at=CAPTURED_AT, - ) - - with pytest.raises(ValidationError): - state.platform = "other-platform" # type: ignore[misc] - - -def test_runtime_observer_has_no_contract_or_governance_layer_dependency() -> None: - source = inspect.getsource(unity_runtime) - - assert "open_data_contract_standard" not in source - assert "datacontract" 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_unity_table_metadata( - _payload(), - source_identifier="https://adb.example", - captured_at=datetime(2026, 8, 30, 3, 0), - ) From 5c367f3a381177ab086752e04a350b9b0bcc791e Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:53:06 +1000 Subject: [PATCH 21/28] test(observation): remove runtime-specific Unity fixture --- .../runtime/unity/orders_table_info.json | 89 ------------------- 1 file changed, 89 deletions(-) delete mode 100644 tests/fixtures/runtime/unity/orders_table_info.json diff --git a/tests/fixtures/runtime/unity/orders_table_info.json b/tests/fixtures/runtime/unity/orders_table_info.json deleted file mode 100644 index 54103cc..0000000 --- a/tests/fixtures/runtime/unity/orders_table_info.json +++ /dev/null @@ -1,89 +0,0 @@ -{ - "name": "orders", - "catalog_name": "main", - "schema_name": "silver", - "full_name": "main.silver.orders", - "table_id": "5d3dc5ba-9070-4f24-b21f-4a153fe2a31d", - "table_type": "MANAGED", - "data_source_format": "DELTA", - "owner": "data-platform@example.com", - "comment": "Current production orders table", - "storage_location": "abfss://silver@example.dfs.core.windows.net/orders", - "created_at": 1788057600000, - "updated_at": 1788144000000, - "properties": { - "delta.enableChangeDataFeed": "true", - "delta.minReaderVersion": "2" - }, - "tags": { - "domain": "sales", - "tier": "silver" - }, - "columns": [ - { - "name": "customer_id", - "type_name": "STRING", - "type_text": "string", - "type_json": "{\"type\":\"string\"}", - "position": 1, - "nullable": false, - "comment": "Customer identifier", - "tags": { - "classification": "identifier" - } - }, - { - "name": "note", - "type_name": "STRING", - "type_text": "string", - "position": 2, - "nullable": true, - "comment": "Optional order note" - }, - { - "name": "order_id", - "type_name": "BIGINT", - "type_text": "bigint", - "type_json": "{\"type\":\"long\"}", - "position": 0, - "nullable": false, - "comment": "Order identifier" - } - ], - "table_constraints": [ - { - "foreign_key_constraint": { - "name": "fk_orders_customer", - "child_columns": ["customer_id"], - "parent_table": "main.silver.customers", - "parent_columns": ["customer_id"], - "rely": true - } - }, - { - "primary_key_constraint": { - "name": "pk_orders", - "child_columns": ["order_id"], - "rely": true - } - } - ], - "view_dependencies": { - "dependencies": [ - { - "table": { - "table_full_name": "main.bronze.orders_raw" - } - }, - { - "function": { - "function_full_name": "main.shared.normalize_order" - } - } - ] - }, - "runtime_extension": { - "source": "fixture", - "revision": 7 - } -} From 2d4ac8698e22ba7aafeadb241e382bece071148c Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:57:00 +1000 Subject: [PATCH 22/28] test(observation): avoid deprecated instance model_fields access --- tests/test_observation_databricks.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_observation_databricks.py b/tests/test_observation_databricks.py index 908e1ad..a4070cb 100644 --- a/tests/test_observation_databricks.py +++ b/tests/test_observation_databricks.py @@ -96,7 +96,7 @@ def test_observation_model_identity_does_not_encode_databricks_namespace_names() captured_at=CAPTURED_AT, ) - identity_fields = set(state.assets[0].identity.model_fields) + 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 From 2405cc4047050b11b25f2946d48ea67100e35d89 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:57:16 +1000 Subject: [PATCH 23/28] docs(observation): use observed-state fingerprint terminology --- semapact/observation/models.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/semapact/observation/models.py b/semapact/observation/models.py index 2bbab4f..683571b 100644 --- a/semapact/observation/models.py +++ b/semapact/observation/models.py @@ -73,8 +73,8 @@ 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 runtime fingerprint capability; this model only - reserves the field used by that later capability. + belong to the dedicated observed-state fingerprint capability; this model + only reserves the field used by that later capability. """ platform: str From 75908809f7650ef2a46d73185249789078977d87 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 13:59:32 +1000 Subject: [PATCH 24/28] ci(databricks): smoke test optional platform extra --- .github/workflows/ci.yml | 31 +++++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index f7af123..67b8d07 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 boundary + run: >- + .venv/bin/python -c + "from databricks.sdk import WorkspaceClient; from databricks.sdk.service.catalog import TableInfo; print(WorkspaceClient.__name__, TableInfo.__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 From 6cbdc73b1f7e5c011652e65618e159efe1d0f196 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 14:01:49 +1000 Subject: [PATCH 25/28] refactor(package): lazy-load public API exports --- semapact/__init__.py | 153 +++++++++++++++++++++++++++---------------- 1 file changed, 96 insertions(+), 57 deletions(-) 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__}) From dadd04b1d627814e498659b9ddeaadf1baefbf4d Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 14:02:03 +1000 Subject: [PATCH 26/28] refactor(importers): avoid eager Delta optional dependency --- semapact/importers/__init__.py | 30 ++++++++++++++++++++++-------- 1 file changed, 22 insertions(+), 8 deletions(-) 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") From 1c82867f6e25cb9f1aca5149766ab6ed728195a5 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 14:02:30 +1000 Subject: [PATCH 27/28] ci(databricks): verify existing Unity import remains loadable --- .github/workflows/ci.yml | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 67b8d07..30abc7c 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -66,10 +66,10 @@ jobs: - name: Install SemaPact Databricks extra from project metadata run: uv pip install --python .venv/bin/python -e ".[databricks]" pytest - - name: Verify Databricks SDK boundary + - 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; print(WorkspaceClient.__name__, TableInfo.__name__)" + "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 From 255282e786a5146b10da0f7572913e31c51aa7f7 Mon Sep 17 00:00:00 2001 From: Elliot Sun Date: Sun, 30 Aug 2026 14:44:17 +1000 Subject: [PATCH 28/28] docs(observation): clarify Databricks auth boundary --- semapact/observation/databricks.py | 11 +++++++++-- 1 file changed, 9 insertions(+), 2 deletions(-) diff --git a/semapact/observation/databricks.py b/semapact/observation/databricks.py index 993aaa5..5c959b7 100644 --- a/semapact/observation/databricks.py +++ b/semapact/observation/databricks.py @@ -3,6 +3,11 @@ 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 @@ -46,8 +51,10 @@ def observe_databricks_table( ) -> ObservedPlatformState: """Observe one Databricks table without generating or mutating an ODCS contract. - ``client`` is expected to be an official ``databricks.sdk.WorkspaceClient`` - in production. It is injectable so core observation tests require no live + ``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: