diff --git a/datajunction-clients/javascript/package-lock.json b/datajunction-clients/javascript/package-lock.json index e1b404b499..6a26df6c1c 100644 --- a/datajunction-clients/javascript/package-lock.json +++ b/datajunction-clients/javascript/package-lock.json @@ -1,12 +1,12 @@ { "name": "datajunction", - "version": "0.0.241", + "version": "0.0.244", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "datajunction", - "version": "0.0.241", + "version": "0.0.244", "license": "MIT", "dependencies": { "@babel/core": "^7.22.5", diff --git a/datajunction-clients/javascript/package.json b/datajunction-clients/javascript/package.json index e347e27f78..daa4e37e9b 100644 --- a/datajunction-clients/javascript/package.json +++ b/datajunction-clients/javascript/package.json @@ -1,6 +1,6 @@ { "name": "datajunction", - "version": "0.0.241", + "version": "0.0.244", "description": "A Javascript client for interacting with a DataJunction server", "module": "src/index.js", "scripts": { diff --git a/datajunction-clients/python/datajunction/__about__.py b/datajunction-clients/python/datajunction/__about__.py index ae2e65baa3..727425f40a 100644 --- a/datajunction-clients/python/datajunction/__about__.py +++ b/datajunction-clients/python/datajunction/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.241" +__version__ = "0.0.244" diff --git a/datajunction-clients/python/datajunction/models.py b/datajunction-clients/python/datajunction/models.py index 25c0c2d4a3..8eacc38194 100644 --- a/datajunction-clients/python/datajunction/models.py +++ b/datajunction-clients/python/datajunction/models.py @@ -406,6 +406,9 @@ class DeploymentResult: deploy_type: str = "" change_tier: Literal["none", "minor", "major"] | None = None semantic_fingerprint: SemanticFingerprintValue | None = None + # True when the node was re-deployed only to retry a pre-existing failure. + # None on responses from servers that predate the field. + revalidation_only: bool | None = None @classmethod def from_dict(cls, d: dict) -> DeploymentResult: @@ -419,6 +422,7 @@ def from_dict(cls, d: dict) -> DeploymentResult: deploy_type=d.get("deploy_type", ""), change_tier=d.get("change_tier"), semantic_fingerprint=_parse_semantic_fingerprint(fingerprint), + revalidation_only=d.get("revalidation_only"), ) diff --git a/datajunction-clients/python/tests/test_deploy.py b/datajunction-clients/python/tests/test_deploy.py index e74bbd3bf3..f3b6bec1e3 100644 --- a/datajunction-clients/python/tests/test_deploy.py +++ b/datajunction-clients/python/tests/test_deploy.py @@ -1418,6 +1418,7 @@ def test_get_impact_calls_api(self, tmp_path, monkeypatch): "version": 1, "digest": "a" * 64, }, + "revalidation_only": True, }, ], "downstream_impacts": [ @@ -1448,6 +1449,7 @@ def test_get_impact_calls_api(self, tmp_path, monkeypatch): parsed = DeploymentInfo.from_dict(result) assert parsed.results[0].deploy_type == "node" assert parsed.results[0].change_tier == "none" + assert parsed.results[0].revalidation_only is True assert parsed.results[0].semantic_fingerprint == SemanticFingerprint( digest="a" * 64, ) diff --git a/datajunction-query/djqs/__about__.py b/datajunction-query/djqs/__about__.py index ae2e65baa3..727425f40a 100644 --- a/datajunction-query/djqs/__about__.py +++ b/datajunction-query/djqs/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.241" +__version__ = "0.0.244" diff --git a/datajunction-reflection/datajunction_reflection/__about__.py b/datajunction-reflection/datajunction_reflection/__about__.py index ae2e65baa3..727425f40a 100644 --- a/datajunction-reflection/datajunction_reflection/__about__.py +++ b/datajunction-reflection/datajunction_reflection/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.241" +__version__ = "0.0.244" diff --git a/datajunction-server/datajunction_server/__about__.py b/datajunction-server/datajunction_server/__about__.py index ae2e65baa3..727425f40a 100644 --- a/datajunction-server/datajunction_server/__about__.py +++ b/datajunction-server/datajunction_server/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.241" +__version__ = "0.0.244" diff --git a/datajunction-server/datajunction_server/api/namespaces.py b/datajunction-server/datajunction_server/api/namespaces.py index 1eeee1eb7f..558eeac468 100644 --- a/datajunction-server/datajunction_server/api/namespaces.py +++ b/datajunction-server/datajunction_server/api/namespaces.py @@ -25,6 +25,9 @@ AccessDenialMode, get_access_checker, ) +from datajunction_server.internal.deployment.fingerprints import ( + build_current_fingerprint_graph, +) from datajunction_server.internal.git.github_service import ( GitHubService, GitHubServiceError, @@ -64,6 +67,8 @@ NamespaceSourcesResponse, ) from datajunction_server.models.namespace import ( + NamespaceFingerprintNode, + NamespaceFingerprintResponse, NamespaceProvisionRequest, NamespaceProvisionResponse, NamespaceWriteStatus, @@ -73,6 +78,10 @@ NodeMinimumDetail, ) from datajunction_server.models.node_type import NodeType +from datajunction_server.models.semantic_fingerprint import ( + LATEST_SEMANTIC_FINGERPRINT_VERSION, + SUPPORTED_SEMANTIC_FINGERPRINT_VERSIONS, +) from datajunction_server.service_clients import QueryServiceClient from datajunction_server.utils import ( get_current_user, @@ -289,6 +298,57 @@ async def list_nodes_in_namespace( return [node for node in nodes if node.name in approved_names] +@router.get( + "/namespaces/{namespace}/semantic-fingerprints", + response_model=NamespaceFingerprintResponse, +) +async def get_namespace_semantic_fingerprints( + namespace: str, + version: int = Query(default=LATEST_SEMANTIC_FINGERPRINT_VERSION), + *, + session: AsyncSession = Depends(get_session), + access_checker: AccessChecker = Depends(get_access_checker), +) -> NamespaceFingerprintResponse: + """Return current semantic fingerprints for a namespace and its descendants.""" + if version not in SUPPORTED_SEMANTIC_FINGERPRINT_VERSIONS: + raise DJInvalidInputException( + f"Unsupported semantic fingerprint version: {version}", + ) + + access_checker.add_namespace(namespace, ResourceAction.READ) + await access_checker.check(on_denied=AccessDenialMode.RAISE) + + nodes = await NodeNamespace.list_all_nodes( + session, + namespace, + options=Node.export_load_options(), + ) + access_checker.add_nodes(nodes, ResourceAction.READ) + await access_checker.check(on_denied=AccessDenialMode.RAISE) + + ordered_nodes = sorted(nodes, key=lambda node: node.name) + specs = [await node.to_spec(session) for node in ordered_nodes] + graph = await build_current_fingerprint_graph( + session, + specs, + version=version, + ) + fingerprints = graph.fingerprints(spec.rendered_name for spec in specs) + return NamespaceFingerprintResponse( + namespace=namespace, + version=version, + nodes=[ + NamespaceFingerprintNode( + name=node.name, + node_type=node.type, + owners=sorted(owner.username for owner in node.owners), + semantic_fingerprint=fingerprints[node.name], + ) + for node in ordered_nodes + ], + ) + + @router.delete("/namespaces/{namespace}/", status_code=HTTPStatus.OK) async def deactivate_a_namespace( namespace: str, diff --git a/datajunction-server/datajunction_server/database/namespace.py b/datajunction-server/datajunction_server/database/namespace.py index 9a9836f636..819a9e5da7 100644 --- a/datajunction-server/datajunction_server/database/namespace.py +++ b/datajunction-server/datajunction_server/database/namespace.py @@ -259,7 +259,7 @@ async def list_all_nodes( select(Node) .where( or_( - Node.namespace.like(f"{namespace}.%"), + Node.namespace.startswith(f"{namespace}.", autoescape=True), Node.namespace == namespace, ), ) diff --git a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index 5af5acf35d..ba75b47f4e 100644 --- a/datajunction-server/datajunction_server/internal/deployment/fingerprints.py +++ b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py @@ -563,3 +563,25 @@ async def build_deployment_fingerprints( submitted_names | additional_target_names, ) return current, proposed_hashes + + +async def build_current_fingerprint_graph( + session: AsyncSession, + specs: Iterable[NodeSpec], + *, + version: int = LATEST_SEMANTIC_FINGERPRINT_VERSION, +) -> SemanticFingerprintGraph: + """Build the current semantic graph for a node set.""" + specs_by_name = {spec.rendered_name: spec for spec in specs} + parent_cache: ParentCandidateCache = {} + external = await _load_external_specs( + session, + specs_by_name.values(), + ignored_parse_errors=set(), + parent_cache=parent_cache, + ) + return SemanticFingerprintGraph( + {**external, **specs_by_name}, + parent_cache=parent_cache, + version=version, + ) diff --git a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py index a1626821fa..711d985c4b 100644 --- a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py +++ b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py @@ -386,6 +386,8 @@ def __init__( self._cubes_bumped_by_upstream: dict[str, list[str]] = {} # Node name -> the tier its change earned, for those cubes to inherit. self._change_tiers: dict[str, ChangeTier] = {} + # Unchanged nodes re-deployed only to retry a pre-existing failure. + self._revalidation_only: set[str] = set() self._current_semantic_fingerprints: FingerprintMap = {} self._proposed_semantic_fingerprints: FingerprintMap = {} @@ -3999,6 +4001,7 @@ async def _create_cubes_from_validation( semantic_fingerprint=self._proposed_semantic_fingerprints.get( cube_spec.rendered_name, ), + revalidation_only=cube_spec.rendered_name in self._revalidation_only, ) deployment_results.append(deployment_result) @@ -4415,6 +4418,9 @@ def filter_nodes_to_deploy( version. So `force` and the INVALID re-deploy below can re-process a node without that implying anything about what changed. + Those re-deployed only for that retry are recorded in `_revalidation_only` + and marked on their `DeploymentResult`. + A cube whose own spec is unchanged is still processed when something upstream of it is changing, matching what `_propagate_update_downstream` does for a `PATCH`: the cube names the same metrics and dimensions, but @@ -4424,6 +4430,7 @@ def filter_nodes_to_deploy( to_create: list[NodeSpec] = [] to_update: list[NodeSpec] = [] to_skip: list[NodeSpec] = [] + revalidation_only: set[str] = set() force = self.deployment_spec.force for node_spec in self.deployment_spec.nodes: existing_spec = existing_nodes_map.get(node_spec.rendered_name) @@ -4461,9 +4468,11 @@ def filter_nodes_to_deploy( and existing_node.current.status == NodeStatus.INVALID ): to_update.append(node_spec) + revalidation_only.add(node_spec.rendered_name) else: to_skip.append(node_spec) + self._revalidation_only = revalidation_only changed_names = {spec.rendered_name for spec in to_create + to_update} self._cubes_bumped_by_upstream = self._cubes_below_changed_nodes( to_skip, @@ -5089,6 +5098,7 @@ async def _process_valid_node_deploy( semantic_fingerprint=self._proposed_semantic_fingerprints.get( result.spec.rendered_name, ), + revalidation_only=result.spec.rendered_name in self._revalidation_only, ) return deployment_result, new_node, new_revision diff --git a/datajunction-server/datajunction_server/models/deployment.py b/datajunction-server/datajunction_server/models/deployment.py index 4921e81447..bfc3545192 100644 --- a/datajunction-server/datajunction_server/models/deployment.py +++ b/datajunction-server/datajunction_server/models/deployment.py @@ -1783,6 +1783,9 @@ class Type(str, Enum): changed_fields: list[str] = Field(default_factory=list) change_tier: ChangeTierName | None = None semantic_fingerprint: SemanticFingerprintValue | None = None + # True when the node was re-deployed only to retry a pre-existing failure. + # Failure reasons are not compared. Nullable for older persisted rows. + revalidation_only: bool | None = None class DeploymentInfo(BaseModel): diff --git a/datajunction-server/datajunction_server/models/namespace.py b/datajunction-server/datajunction_server/models/namespace.py index 3dfc9d9483..db9792baae 100644 --- a/datajunction-server/datajunction_server/models/namespace.py +++ b/datajunction-server/datajunction_server/models/namespace.py @@ -2,6 +2,8 @@ from pydantic import BaseModel, Field +from datajunction_server.models.node_type import NodeType +from datajunction_server.models.semantic_fingerprint import SemanticFingerprintValue from datajunction_server.typing import StrEnum @@ -73,3 +75,20 @@ class NamespaceProvisionResponse(BaseModel): namespace: str owner_role: str deployer_role: str | None + + +class NamespaceFingerprintNode(BaseModel): + """A node and its current semantic fingerprint.""" + + name: str + node_type: NodeType + owners: list[str] = Field(default_factory=list) + semantic_fingerprint: SemanticFingerprintValue + + +class NamespaceFingerprintResponse(BaseModel): + """Current semantic fingerprints for a namespace tree.""" + + namespace: str + version: int + nodes: list[NamespaceFingerprintNode] = Field(default_factory=list) diff --git a/datajunction-server/tests/api/deployments_test.py b/datajunction-server/tests/api/deployments_test.py index ce802256bc..ca250163d8 100644 --- a/datajunction-server/tests/api/deployments_test.py +++ b/datajunction-server/tests/api/deployments_test.py @@ -1331,7 +1331,16 @@ def deployment_payload(deployment_spec: DeploymentSpec) -> dict: return deployment_spec.model_dump() -async def deploy_and_wait(client, deployment_spec: DeploymentSpec): +# Additive `DeploymentResult` fields, stripped rather than asserted below. +ADDITIVE_RESULT_FIELDS = ( + "change_tier", + "semantic_fingerprint", + "revalidation_only", +) + + +async def deploy_and_poll(client, deployment_spec: DeploymentSpec): + """Deploy and wait, keeping every field the API returned.""" response = await client.post( "/deployments", json=deployment_payload(deployment_spec), @@ -1345,11 +1354,15 @@ async def deploy_and_wait(client, deployment_spec: DeploymentSpec): await asyncio.sleep(1) response = await client.get(f"/deployments/{deployment_uuid}") data = response.json() + return data + + +async def deploy_and_wait(client, deployment_spec: DeploymentSpec): + data = await deploy_and_poll(client, deployment_spec) for result in data.get("results", []): - assert "change_tier" in result - assert "semantic_fingerprint" in result - result.pop("change_tier", None) - result.pop("semantic_fingerprint", None) + for additive in ADDITIVE_RESULT_FIELDS: + assert additive in result + result.pop(additive) return data @@ -1435,6 +1448,69 @@ async def test_deploy_failed_on_non_existent_upstream_deps( assert f"{namespace}.default.us_state" in link_result["name"] assert link_result["status"] == "failed" + @pytest.mark.asyncio + async def test_revalidation_only_marks_pre_existing_failures(self, client): + """ + An unchanged node re-deployed only to retry a pre-existing failure is + marked `revalidation_only`, so a caller can tell the failures this + deployment caused from the ones it inherited -- and still sees the + recovery when a later deploy fixes the node's upstream. + """ + namespace = "revalidation_only" + transform = TransformSpec( + name="${prefix}default.repair_totals", + query="SELECT repair_order_id, price FROM ${prefix}default.repairs", + owners=["dj"], + ) + source = SourceSpec( + name="${prefix}default.repairs", + table="repairs", + catalog="default", + schema_="roads", + columns=[ + ColumnSpec(name="repair_order_id", type="int"), + ColumnSpec(name="price", type="float"), + ], + owners=["dj"], + ) + + # The transform is created broken: its source is not in the deployment. + data = await deploy_and_poll( + client, + DeploymentSpec(namespace=namespace, nodes=[transform]), + ) + created = next(r for r in data["results"] if r["deploy_type"] == "node") + assert created["status"] == "invalid" + assert created["operation"] == "create" + assert created["revalidation_only"] is False + + # Re-deploying the same spec retries the node, which is still broken. + data = await deploy_and_poll( + client, + DeploymentSpec(namespace=namespace, nodes=[transform]), + ) + retried = next(r for r in data["results"] if r["deploy_type"] == "node") + assert retried["status"] == "invalid" + assert retried["operation"] == "update" + assert retried["revalidation_only"] is True + + # Adding the missing source fixes the node, which reports as a success. + data = await deploy_and_poll( + client, + DeploymentSpec(namespace=namespace, nodes=[transform, source]), + ) + recovered = next( + r + for r in data["results"] + if r["name"] == f"{namespace}.default.repair_totals" + ) + assert recovered["status"] == "success" + assert recovered["revalidation_only"] is True + added_source = next( + r for r in data["results"] if r["name"] == f"{namespace}.default.repairs" + ) + assert added_source["revalidation_only"] is False + @pytest.mark.asyncio async def test_deploy_failed_on_non_existent_link_deps( self, diff --git a/datajunction-server/tests/api/git_test.py b/datajunction-server/tests/api/git_test.py index 2f6450417d..02a2b07a6f 100644 --- a/datajunction-server/tests/api/git_test.py +++ b/datajunction-server/tests/api/git_test.py @@ -4410,11 +4410,20 @@ async def test_branch_creation_copies_nodes( result["semantic_fingerprint"]["version"] == 1 for result in data["deployment_results"] ) + assert all( + result["revalidation_only"] is False + for result in data["deployment_results"] + ) deployment_results = [ { key: value for key, value in result.items() - if key not in {"change_tier", "semantic_fingerprint"} + if key + not in { + "change_tier", + "semantic_fingerprint", + "revalidation_only", + } } for result in data["deployment_results"] ] diff --git a/datajunction-server/tests/api/namespace_fingerprints_test.py b/datajunction-server/tests/api/namespace_fingerprints_test.py new file mode 100644 index 0000000000..f783f12fe0 --- /dev/null +++ b/datajunction-server/tests/api/namespace_fingerprints_test.py @@ -0,0 +1,359 @@ +"""Tests for namespace semantic fingerprint snapshots.""" + +import asyncio +from unittest import mock + +import pytest + +from datajunction_server.database.node import Node +from datajunction_server.internal.access.authorization import AuthorizationService +from datajunction_server.models import access +from datajunction_server.models.deployment import ( + ColumnSpec, + CubeSpec, + DeploymentSpec, + MetricSpec, + SourceSpec, + TransformSpec, +) +from tests.authz import VALIDATOR_AUTH_SERVICE + + +@pytest.fixture(autouse=True, scope="module") +def patch_effective_writer_concurrency(): + from datajunction_server.internal.deployment.deployment import settings + + with mock.patch.object( + settings.__class__, + "effective_writer_concurrency", + new_callable=mock.PropertyMock, + return_value=1, + ): + yield + + +def lunch_kitchen(grill_table: str = "grill_orders") -> DeploymentSpec: + return DeploymentSpec( + namespace="lunch.kitchen", + nodes=[ + SourceSpec( + name="grill_orders", + catalog="default", + schema_="food", + table=grill_table, + columns=[ColumnSpec(name="id", type="int")], + owners=["burger_chef"], + ), + TransformSpec( + name="burger_orders", + query="SELECT id FROM ${prefix}grill_orders", + owners=["burger_chef"], + ), + MetricSpec( + name="burger_count", + query="SELECT COUNT(*) FROM ${prefix}burger_orders", + owners=["burger_chef"], + ), + CubeSpec( + name="picnic_basket", + metrics=["${prefix}burger_count"], + owners=["picnic_planner"], + ), + SourceSpec( + name="milkshake_orders", + catalog="default", + schema_="food", + table="milkshake_orders", + columns=[ColumnSpec(name="id", type="int")], + owners=["dessert_chef"], + ), + MetricSpec( + name="milkshake_count", + query="SELECT COUNT(*) FROM ${prefix}milkshake_orders", + owners=["dessert_chef"], + ), + ], + ) + + +def taco_truck() -> DeploymentSpec: + return DeploymentSpec( + namespace="lunch.taco_truck", + nodes=[ + TransformSpec( + name="taco_orders", + query="SELECT id FROM lunch.kitchen.grill_orders", + owners=["taco_chef"], + ), + MetricSpec( + name="taco_count", + query="SELECT COUNT(*) FROM ${prefix}taco_orders", + owners=["taco_chef"], + ), + ], + ) + + +async def deploy(client, *specs: DeploymentSpec) -> None: + for spec in specs: + response = await client.post( + "/deployments", + json=spec.model_dump(by_alias=True), + ) + assert response.status_code == 200 + deployment_id = response.json()["uuid"] + for _ in range(300): + result = await client.get(f"/deployments/{deployment_id}") + if result.json()["status"] in ("success", "failed"): + assert result.json()["status"] == "success" + break + await asyncio.sleep(0.1) + else: + pytest.fail(f"Deployment {deployment_id} did not finish") + + +async def snapshot(client, namespace: str) -> dict: + response = await client.get( + f"/namespaces/{namespace}/semantic-fingerprints?version=1", + ) + assert response.status_code == 200 + return response.json() + + +def fingerprint_map(*snapshots: dict) -> dict[str, dict | str]: + return { + node["name"]: node["semantic_fingerprint"] + for namespace_snapshot in snapshots + for node in namespace_snapshot["nodes"] + } + + +@pytest.mark.asyncio +async def test_git_backed_lunch_fingerprints_match_deployment_impact( + client_with_roads, +): + initial_kitchen = lunch_kitchen() + initial_taco_truck = taco_truck() + await deploy(client_with_roads, initial_kitchen, initial_taco_truck) + + root_git = await client_with_roads.patch( + "/namespaces/lunch/git", + json={ + "github_repo_path": "food-festival/lunch-menu", + "default_branch": "main", + }, + ) + assert root_git.status_code == 200 + kitchen_git = await client_with_roads.patch( + "/namespaces/lunch.kitchen/git", + json={"git_branch": "kitchen", "parent_namespace": "lunch"}, + ) + assert kitchen_git.status_code == 200 + taco_truck_git = await client_with_roads.patch( + "/namespaces/lunch.taco_truck/git", + json={"git_branch": "taco-truck", "parent_namespace": "lunch"}, + ) + assert taco_truck_git.status_code == 200 + + before_lunch = await snapshot(client_with_roads, "lunch") + assert before_lunch["namespace"] == "lunch" + assert before_lunch["version"] == 1 + assert [node["name"] for node in before_lunch["nodes"]] == sorted( + node.rendered_name + for spec in (initial_kitchen, initial_taco_truck) + for node in spec.nodes + ) + assert {node["name"]: node["owners"] for node in before_lunch["nodes"]} == { + "lunch.kitchen.burger_count": ["burger_chef"], + "lunch.kitchen.burger_orders": ["burger_chef"], + "lunch.kitchen.grill_orders": ["burger_chef"], + "lunch.kitchen.milkshake_count": ["dessert_chef"], + "lunch.kitchen.milkshake_orders": ["dessert_chef"], + "lunch.kitchen.picnic_basket": ["picnic_planner"], + "lunch.taco_truck.taco_count": ["taco_chef"], + "lunch.taco_truck.taco_orders": ["taco_chef"], + } + assert fingerprint_map(before_lunch) == { + "lunch.kitchen.burger_count": { + "version": 1, + "digest": "51876c5b1c1dafdcc1830730b9e5f5c41267daf51f94a3cec564d84b5cc0fbf7", + }, + "lunch.kitchen.burger_orders": { + "version": 1, + "digest": "0169918baf5a92f4545d6476214bdb2aa2d2975174c5e8e75c49785482f51ce1", + }, + "lunch.kitchen.grill_orders": { + "version": 1, + "digest": "fa4e66f1924973b8cca8aeb59752e2fe90066bdfc45894d9a31cac5c7b0c4d2c", + }, + "lunch.kitchen.milkshake_count": { + "version": 1, + "digest": "0438abad1a7bf147dc708a33ea20dac22b0d036647cea263d531f5a3a3e15a65", + }, + "lunch.kitchen.milkshake_orders": { + "version": 1, + "digest": "c613ea9bcd613edabebd0b7e09840e641652aed1a119817e4a524da88ec80f57", + }, + "lunch.kitchen.picnic_basket": { + "version": 1, + "digest": "2e506807af2296c51a63cca66e8ecb443385d482173161fe7b625cf6a7467aa8", + }, + "lunch.taco_truck.taco_count": { + "version": 1, + "digest": "9d92b7b9637148f386fdceb31f84c862ebcacf47a3098d93c9ccf8c57467c0d2", + }, + "lunch.taco_truck.taco_orders": { + "version": 1, + "digest": "0169918baf5a92f4545d6476214bdb2aa2d2975174c5e8e75c49785482f51ce1", + }, + } + + proposed_kitchen = lunch_kitchen(grill_table="grill_orders_v2") + impact_response = await client_with_roads.post( + "/deployments/impact", + json=proposed_kitchen.model_dump(by_alias=True), + ) + assert impact_response.status_code == 200 + impact = impact_response.json() + predicted = { + result["name"]: result["semantic_fingerprint"] + for result in impact["results"] + if result["deploy_type"] == "node" + } + predicted.update( + { + downstream["name"]: downstream["semantic_fingerprint"] + for downstream in impact["downstream_impacts"] + if downstream["semantic_fingerprint"] is not None + }, + ) + + await deploy(client_with_roads, proposed_kitchen) + after_lunch = await snapshot(client_with_roads, "lunch") + before = fingerprint_map(before_lunch) + after = fingerprint_map(after_lunch) + changed = {name for name in before if before[name] != after[name]} + + assert changed == { + "lunch.kitchen.burger_count", + "lunch.kitchen.burger_orders", + "lunch.kitchen.grill_orders", + "lunch.kitchen.picnic_basket", + "lunch.taco_truck.taco_count", + "lunch.taco_truck.taco_orders", + } + assert all(predicted[name] == after[name] for name in changed) + assert ( + before["lunch.kitchen.milkshake_count"] + == after["lunch.kitchen.milkshake_count"] + ) + assert ( + before["lunch.kitchen.milkshake_orders"] + == after["lunch.kitchen.milkshake_orders"] + ) + + +@pytest.mark.asyncio +async def test_namespace_snapshot_uses_unknown_for_legacy_invalid_sql( + client_with_roads, + session, +): + spec = DeploymentSpec( + namespace="lunch.mystery", + nodes=[ + SourceSpec( + name="grill", + catalog="default", + schema_="food", + table="grill", + columns=[ColumnSpec(name="id", type="int")], + ), + TransformSpec( + name="mystery_meat", + query="SELECT id FROM ${prefix}grill", + ), + TransformSpec( + name="customer", + query="SELECT id FROM ${prefix}mystery_meat", + ), + ], + ) + await deploy(client_with_roads, spec) + + mystery_meat = await Node.get_by_name(session, "lunch.mystery.mystery_meat") + assert mystery_meat is not None + mystery_meat.current.query = "SELECT (" + await session.commit() + + data = await snapshot(client_with_roads, "lunch.mystery") + fingerprints = fingerprint_map(data) + assert fingerprints["lunch.mystery.mystery_meat"] == "unknown" + assert fingerprints["lunch.mystery.customer"] == "unknown" + + +@pytest.mark.asyncio +async def test_namespace_snapshot_rejects_unsupported_version(client_with_roads): + response = await client_with_roads.get( + "/namespaces/default/semantic-fingerprints?version=2", + ) + + assert response.status_code == 422 + assert response.json()["message"] == "Unsupported semantic fingerprint version: 2" + + +@pytest.mark.asyncio +@pytest.mark.parametrize("restrict_reader", [False, True]) +async def test_namespace_snapshot_keeps_underscore_boundary_literal( + client_with_roads, + mocker, + restrict_reader, +): + namespace = "fingerprints.taco_truck" + specs = [ + DeploymentSpec( + namespace=target, + nodes=[ + SourceSpec( + name="orders", + catalog="default", + schema_="food", + table="orders", + columns=[ColumnSpec(name="id", type="int")], + ), + ], + ) + for target in ( + namespace, + f"{namespace}.child", + "fingerprints.tacoxtruck.child", + ) + ] + await deploy(client_with_roads, *specs) + + if restrict_reader: + + class SnapshotTreeAuthorizationService(AuthorizationService): + """Allow the requested tree while denying its similarly named sibling.""" + + name = "snapshot_tree" + + def authorize(self, auth_context, requests): + return [ + access.AccessDecision( + request=request, + approved=request.access_object.name == namespace + or request.access_object.name.startswith(f"{namespace}."), + ) + for request in requests + ] + + mocker.patch( + VALIDATOR_AUTH_SERVICE, + return_value=SnapshotTreeAuthorizationService(), + ) + + data = await snapshot(client_with_roads, namespace) + assert [node["name"] for node in data["nodes"]] == [ + f"{namespace}.child.orders", + f"{namespace}.orders", + ] diff --git a/datajunction-server/tests/internal/deployment/orchestration_test.py b/datajunction-server/tests/internal/deployment/orchestration_test.py index f960aa1a39..aa167a92a9 100644 --- a/datajunction-server/tests/internal/deployment/orchestration_test.py +++ b/datajunction-server/tests/internal/deployment/orchestration_test.py @@ -841,6 +841,35 @@ def test_filter_nodes_to_deploy_without_force( assert len(to_skip) == len(sample_deployment_spec.nodes) assert to_delete == [] + def test_filter_nodes_records_revalidation_only( + self, + orchestrator, + sample_deployment_spec, + ): + """An unchanged node that is stuck INVALID is queued for revalidation.""" + existing_specs = { + node.rendered_name: node for node in sample_deployment_spec.nodes + } + stuck = sample_deployment_spec.nodes[1] + orchestrator.registry.add_nodes( + { + stuck.rendered_name: SimpleNamespace( + current=SimpleNamespace( + status=NodeStatus.INVALID, + parents=[], + ), + ), + }, + ) + to_deploy, to_skip, _ = orchestrator.filter_nodes_to_deploy(existing_specs) + assert to_deploy == [stuck] + assert orchestrator._revalidation_only == {stuck.rendered_name} + assert [spec.rendered_name for spec in to_skip] == [ + spec.rendered_name + for spec in sample_deployment_spec.nodes + if spec is not stuck + ] + def test_filter_nodes_uses_normalized_query_and_source_columns(self): incoming = [ TransformSpec(name="transform", query=" SELECT id\nFROM source "), diff --git a/datajunction-ui/package.json b/datajunction-ui/package.json index 38e378239d..877896311d 100644 --- a/datajunction-ui/package.json +++ b/datajunction-ui/package.json @@ -1,6 +1,6 @@ { "name": "datajunction-ui", - "version": "0.0.241", + "version": "0.0.244", "description": "DataJunction UI", "type": "module", "module": "src/index.tsx", diff --git a/docs/content/0.1.0/docs/developers/semantic-fingerprint-versioning.md b/docs/content/0.1.0/docs/developers/semantic-fingerprint-versioning.md new file mode 100644 index 0000000000..3698f2f75a --- /dev/null +++ b/docs/content/0.1.0/docs/developers/semantic-fingerprint-versioning.md @@ -0,0 +1,155 @@ +--- +weight: 85 +title: "Semantic fingerprint compatibility" +description: "Compatibility rules and migration procedures for semantic fingerprint consumers" +draft: false +--- + +Semantic fingerprints are versioned identifiers for a node definition and its +semantic ancestors. Consumers may store them in Git, use changes to select +reviewers, or compare a deployment preview with committed server state. Stable +fingerprints therefore require a compatibility contract that extends beyond the +final hash function. + +For a given version, a fingerprint must remain stable when the complete DJ +metadata graph remains stable. The version covers all behavior that contributes +to the result: + +- selected node fields and their defaults; +- field normalization and canonical JSON encoding; +- structural SQL serialization; +- semantic parent discovery and resolution; +- cycle handling and parent hash composition; +- propagation of the `unknown` sentinel; and +- the digest algorithm and its domain separators. + +After consumers adopt a version, any change to this behavior requires a new +fingerprint version. A bug fix that changes existing digests also requires a new +version. + +## What can change a fingerprint + +The stability condition is: + +```text +same fingerprint version + same complete DJ metadata graph = same fingerprint +``` + +A repository can retain identical files while one of its fingerprints changes. +The complete graph can change through: + +- a Git-backed change in another namespace that supplies an ancestor; +- an API or UI update outside the repository deployment; +- source registration or schema reflection; +- resolution of a previously missing parent; or +- adoption of a new fingerprint version. + +Node owners, tags, display names, node descriptions, and other fields outside a +version's semantic field projection do not change its digest. Metadata attached +to a projected field, such as a source column, may contribute to the digest. + +## External tables + +DJ fingerprints metadata rather than the live contents of a warehouse. When an +external table has a DJ source node, its catalog, schema, table, and reflected +columns contribute to the source fingerprint. A physical schema change affects +the source and its descendants after DJ refreshes or redeploys that source +metadata. + +A physical change that has not reached DJ metadata is invisible to the +fingerprint. Data changes and properties absent from the source specification +are also invisible. A query parent with no corresponding DJ node makes the +fingerprint `unknown`, and that value propagates to its descendants. + +Source fingerprints represent the full reflected column set. Fingerprints do +not currently track column-level lineage, so a reflected change to an unused +source column can change every descendant fingerprint. + +## Normal server upgrade + +Consumers should request an explicit fingerprint version. They should not use +the server's latest-version default as a rollout control. + +Before a normal server upgrade: + +1. Evaluate a representative committed metadata snapshot with the current and + candidate server. +2. Request the fingerprint version pinned by each consumer. +3. Compare every returned value, including `unknown`. +4. Block the release if a pinned version changes unexpectedly. +5. Deploy the server without changing consumer baselines when the comparison is + clean. + +A normal server upgrade requires no fingerprint metadata update. + +## Introducing a fingerprint version + +New fingerprint behavior should be introduced through a staged migration: + +1. Add the new version while continuing to serve every version used by a + consumer. +2. Deploy the server with both the old and new versions available. +3. Keep consumers pinned to the old version. +4. Evaluate both versions over the production metadata graph and review the + changed-node and owner counts. +5. Generate a dedicated metadata migration pull request containing the new + version pin and its complete fingerprint baseline. +6. Mark that pull request as an algorithm migration so review automation does + not notify every owner whose hash changed only because of the migration. +7. Switch the consumer to the new version when the migration pull request + merges. +8. Retain the old server version until all consumers have migrated and the + rollback window has closed. + +The server deployment makes the new algorithm available. The migration pull +request chooses when a consumer adopts it. Fingerprints are computed from DJ +metadata, so this process does not require rewriting node metadata in the DJ +database. + +## Review automation safeguards + +Automation that turns fingerprint changes into review requests should enforce +the following controls: + +- Record the fingerprint version in every generated baseline. +- Reject a server response whose version differs from the requested version. +- Require an explicit migration mode for a version change. +- Stop before opening or updating a pull request when an unexpected change + exceeds a configured node or owner threshold. +- Report direct deployment changes separately from inherited Merkle changes. +- Treat transitions to or from `unknown` as a separate warning category. +- Cap automatic reviewer fanout and require manual approval above the cap. +- Suppress per-node owner notifications for an approved algorithm migration. + +These controls keep a server defect or accidental same-version change from +creating a repository-wide review event. + +## Rollback + +The previous fingerprint version must remain available during the migration +window. If the new version has a defect: + +1. Revert the consumer's version pin and generated baseline. +2. Resume requests for the previous version. +3. Fix the algorithm under another new version if the corrected output differs + from a version already adopted by consumers. + +No DJ node deployment is needed when rolling back only the fingerprint version. + +## Release checklist + +For a normal server release: + +- [ ] Consumers request explicit versions. +- [ ] Full-graph golden fixtures pass. +- [ ] The candidate server matches the current server for pinned versions. +- [ ] No generated metadata migration is required. + +For an algorithm migration: + +- [ ] The new behavior uses a new version number. +- [ ] The server serves both old and new versions. +- [ ] Production graph differences and owner fanout have been reviewed. +- [ ] The migration pull request records the new version and complete baseline. +- [ ] Migration mode prevents repository-wide owner notifications. +- [ ] The previous version remains available for rollback.