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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 12 additions & 0 deletions semapact/observation/__init__.py
Original file line number Diff line number Diff line change
@@ -1,5 +1,12 @@
"""Platform-neutral observation domain for SemaPact read-side state."""

from semapact.observation.fingerprint import (
OBSERVED_STATE_FINGERPRINT_ALGORITHM,
OBSERVED_STATE_FINGERPRINT_VERSION,
canonical_observed_state_payload,
fingerprint_observed_state,
with_observed_state_fingerprint,
)
from semapact.observation.models import (
ObservedAsset,
ObservedAssetIdentity,
Expand All @@ -10,10 +17,15 @@
)

__all__ = [
"OBSERVED_STATE_FINGERPRINT_ALGORITHM",
"OBSERVED_STATE_FINGERPRINT_VERSION",
"ObservedAsset",
"ObservedAssetIdentity",
"ObservedPlatformState",
"ObservedProperty",
"ObservedPropertyIdentity",
"canonical_observed_state_payload",
"fingerprint_observed_state",
"serialize_observed_state",
"with_observed_state_fingerprint",
]
6 changes: 4 additions & 2 deletions semapact/observation/databricks.py
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
from datetime import datetime, timezone
from typing import TYPE_CHECKING, Any, Mapping, Protocol

from semapact.observation.fingerprint import with_observed_state_fingerprint
from semapact.observation.models import (
ObservedAsset,
ObservedAssetIdentity,
Expand Down Expand Up @@ -81,7 +82,7 @@ def map_databricks_table_info(
captured_at: datetime,
table_fqn: str | None = None,
) -> ObservedPlatformState:
"""Project an SDK ``TableInfo`` into the platform-neutral observation model."""
"""Project an SDK ``TableInfo`` into fingerprinted platform-neutral state."""
_require_aware_datetime(captured_at)
if not source_identifier:
raise ValueError("source_identifier is required for Databricks observation")
Expand All @@ -97,13 +98,14 @@ def map_databricks_table_info(
properties=_properties(metadata.get("columns"), identity=identity),
)

return ObservedPlatformState(
state = ObservedPlatformState(
platform=DATABRICKS_PLATFORM,
source_identifier=source_identifier.rstrip("/"),
assets=(asset,),
captured_at=captured_at,
fingerprint=None,
)
return with_observed_state_fingerprint(state)


def _asset_identity(
Expand Down
88 changes: 88 additions & 0 deletions semapact/observation/fingerprint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,88 @@
"""Stable semantic fingerprints for platform-neutral observed state.

A fingerprint identifies the semantic content of an ``ObservedPlatformState``.
It deliberately excludes observation-envelope fields such as ``captured_at``
and ``source_identifier`` so repeated captures of the same platform state have
the same fingerprint.

Version ``obs-v1`` hashes the minimal observation model introduced by M1:
platform, asset identity/type, and property identity/physical type/nullability.
Provider-specific semantic normalization remains the responsibility of the
provider adapter before it constructs the platform-neutral observation model.
"""

from __future__ import annotations

import hashlib
import json
from typing import Any

from semapact.observation.models import (
ObservedAsset,
ObservedPlatformState,
ObservedProperty,
)

OBSERVED_STATE_FINGERPRINT_VERSION = "obs-v1"
OBSERVED_STATE_FINGERPRINT_ALGORITHM = "sha256"


def canonical_observed_state_payload(state: ObservedPlatformState) -> dict[str, object]:
"""Return the versioned semantic payload used by the v1 fingerprint."""
assets = [_canonical_asset(asset) for asset in state.assets]
assets.sort(key=_canonical_json)

return {
"fingerprint_version": OBSERVED_STATE_FINGERPRINT_VERSION,
"platform": state.platform.casefold(),
"assets": assets,
}


def fingerprint_observed_state(state: ObservedPlatformState) -> str:
"""Return the deterministic fingerprint for observed semantic state."""
canonical = _canonical_json(canonical_observed_state_payload(state)).encode("utf-8")
digest = hashlib.sha256(canonical).hexdigest()
return (
f"{OBSERVED_STATE_FINGERPRINT_VERSION}:"
f"{OBSERVED_STATE_FINGERPRINT_ALGORITHM}:{digest}"
)


def with_observed_state_fingerprint(state: ObservedPlatformState) -> ObservedPlatformState:
"""Return an immutable copy of ``state`` with its canonical fingerprint set."""
return state.model_copy(update={"fingerprint": fingerprint_observed_state(state)})


def _canonical_asset(asset: ObservedAsset) -> dict[str, object]:
properties = [_canonical_property(prop) for prop in asset.properties]
properties.sort(key=_canonical_json)

return {
"identity": list(asset.identity.canonical_key),
"asset_type": _normalize_optional_text(asset.asset_type),
"properties": properties,
}


def _canonical_property(prop: ObservedProperty) -> dict[str, object]:
return {
"identity": list(prop.identity.canonical_key),
"physical_type": _normalize_optional_text(prop.physical_type),
"nullable": prop.nullable,
}


def _normalize_optional_text(value: str | None) -> str | None:
if value is None:
return None
return value.strip()


def _canonical_json(value: Any) -> str:
return json.dumps(
value,
ensure_ascii=False,
separators=(",", ":"),
sort_keys=True,
)
7 changes: 4 additions & 3 deletions semapact/observation/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -72,9 +72,10 @@ class ObservedAsset(ObservationModel):
class ObservedPlatformState(ObservationModel):
"""Point-in-time platform observation independent from governed ODCS state.

``fingerprint`` is intentionally optional. Stable normalization and hashing
belong to the dedicated observed-state fingerprint capability; this model
only reserves the field used by that later capability.
``fingerprint`` is the optional versioned semantic content hash for this
observed state. Provider adapters may populate it through the shared
fingerprint capability; manually constructed observations may leave it
unset until fingerprinting is requested.
"""

platform: str
Expand Down
8 changes: 7 additions & 1 deletion tests/test_observation_databricks.py
Original file line number Diff line number Diff line change
Expand Up @@ -63,7 +63,10 @@ def test_databricks_table_info_maps_to_platform_neutral_observation() -> None:
assert state.platform == "databricks"
assert state.source_identifier == "https://adb.example"
assert state.captured_at == CAPTURED_AT
assert state.fingerprint is None
assert state.fingerprint == (
"obs-v1:sha256:"
"cc72fb5e20b673784b0847f6f903c05178b88c3d88ffb25f830da0a373457287"
)

assert len(state.assets) == 1
asset = state.assets[0]
Expand Down Expand Up @@ -119,6 +122,7 @@ def test_serialization_is_stable_when_databricks_column_order_changes() -> None:
)

assert left == right
assert left.fingerprint == right.fingerprint
assert serialize_observed_state(left) == serialize_observed_state(right)


Expand All @@ -134,6 +138,7 @@ def test_observe_databricks_table_uses_workspace_client_tables_get() -> None:

assert client.tables.calls == ["main.silver.orders"]
assert state.assets[0].identity.asset == "orders"
assert state.fingerprint is not None


def test_mapper_accepts_official_databricks_sdk_table_info_when_extra_is_installed() -> None:
Expand All @@ -148,6 +153,7 @@ def test_mapper_accepts_official_databricks_sdk_table_info_when_extra_is_install

assert state.assets[0].identity.namespace == ("main", "silver")
assert state.assets[0].properties[0].identity.property == "order_id"
assert state.fingerprint is not None


def test_observation_models_are_immutable() -> None:
Expand Down
177 changes: 177 additions & 0 deletions tests/test_observation_fingerprint.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
from __future__ import annotations

from datetime import datetime, timezone

from semapact.observation.fingerprint import (
canonical_observed_state_payload,
fingerprint_observed_state,
with_observed_state_fingerprint,
)
from semapact.observation.models import (
ObservedAsset,
ObservedAssetIdentity,
ObservedPlatformState,
ObservedProperty,
ObservedPropertyIdentity,
)

CAPTURED_AT = datetime(2026, 8, 30, 3, 0, tzinfo=timezone.utc)


def _orders_asset(
*,
platform: str = "databricks",
namespace: tuple[str, ...] = ("main", "silver"),
asset_name: str = "orders",
property_order: tuple[str, ...] = ("order_id", "customer_id", "note"),
physical_types: dict[str, str] | None = None,
nullable: dict[str, bool | None] | None = None,
asset_type: str = "MANAGED",
) -> ObservedAsset:
identity = ObservedAssetIdentity(
platform=platform,
namespace=namespace,
asset=asset_name,
)
types = physical_types or {
"order_id": "bigint",
"customer_id": "string",
"note": "string",
}
nullability = nullable or {
"order_id": False,
"customer_id": False,
"note": True,
}
return ObservedAsset(
identity=identity,
asset_type=asset_type,
properties=tuple(
ObservedProperty(
identity=ObservedPropertyIdentity(asset=identity, property=name),
physical_type=types[name],
nullable=nullability[name],
)
for name in property_order
),
)


def _state(
*,
assets: tuple[ObservedAsset, ...] | None = None,
platform: str = "databricks",
source_identifier: str = "https://adb.example",
captured_at: datetime = CAPTURED_AT,
fingerprint: str | None = None,
) -> ObservedPlatformState:
return ObservedPlatformState(
platform=platform,
source_identifier=source_identifier,
assets=assets or (_orders_asset(),),
captured_at=captured_at,
fingerprint=fingerprint,
)


def test_observed_state_fingerprint_has_stable_versioned_golden_value() -> None:
state = _state()

assert fingerprint_observed_state(state) == (
"obs-v1:sha256:"
"cc72fb5e20b673784b0847f6f903c05178b88c3d88ffb25f830da0a373457287"
)
assert canonical_observed_state_payload(state)["fingerprint_version"] == "obs-v1"


def test_observation_envelope_fields_do_not_change_semantic_fingerprint() -> None:
baseline = _state()
later = _state(
source_identifier="workspace-alias-b",
captured_at=datetime(2026, 8, 31, 4, 30, tzinfo=timezone.utc),
fingerprint="previous-value",
)

assert fingerprint_observed_state(baseline) == fingerprint_observed_state(later)


def test_asset_property_order_and_identity_case_do_not_change_fingerprint() -> None:
first = _orders_asset()
second_identity = ObservedAssetIdentity(
platform="databricks",
namespace=("main", "silver"),
asset="customers",
)
second = ObservedAsset(identity=second_identity, asset_type="MANAGED")

reordered = _orders_asset(
platform="DATABRICKS",
namespace=("MAIN", "SILVER"),
asset_name="ORDERS",
property_order=("note", "customer_id", "order_id"),
)

left = _state(assets=(first, second))
right = _state(platform="DATABRICKS", assets=(second, reordered))

assert fingerprint_observed_state(left) == fingerprint_observed_state(right)


def test_governance_relevant_observed_changes_change_fingerprint() -> None:
baseline = fingerprint_observed_state(_state())

changed_asset_type = _state(assets=(_orders_asset(asset_type="EXTERNAL"),))
changed_property_type = _state(
assets=(
_orders_asset(
physical_types={
"order_id": "string",
"customer_id": "string",
"note": "string",
}
),
)
)
changed_nullability = _state(
assets=(
_orders_asset(
nullable={
"order_id": True,
"customer_id": False,
"note": True,
}
),
)
)
changed_property_identity = _state(
assets=(
_orders_asset(
property_order=("order_id", "customer_id", "comment"),
physical_types={
"order_id": "bigint",
"customer_id": "string",
"comment": "string",
},
nullable={
"order_id": False,
"customer_id": False,
"comment": True,
},
),
)
)

assert fingerprint_observed_state(changed_asset_type) != baseline
assert fingerprint_observed_state(changed_property_type) != baseline
assert fingerprint_observed_state(changed_nullability) != baseline
assert fingerprint_observed_state(changed_property_identity) != baseline


def test_with_observed_state_fingerprint_returns_immutable_copy() -> None:
state = _state()

fingerprinted = with_observed_state_fingerprint(state)

assert state.fingerprint is None
assert fingerprinted is not state
assert fingerprinted.fingerprint == fingerprint_observed_state(state)
Loading