From e279df1b1a7a8664f5ed61952b6150a00a61fe0c Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Tue, 1 Sep 2026 18:48:26 -0700 Subject: [PATCH 01/13] Add semantic fingerprints to deployment impact Return server-owned fingerprints and change tiers for deployment results and downstream impacts, including client parsing and unknown-state propagation. --- .../python/datajunction/models.py | 49 ++- .../python/tests/test_deploy.py | 103 ++++- .../internal/deployment/orchestrator.py | 211 +++++++--- .../datajunction_server/internal/impact.py | 94 ++++- .../datajunction_server/models/deployment.py | 16 + .../datajunction_server/models/impact.py | 2 + .../tests/api/deployment_impact_test.py | 384 +++++++++++++++++- .../internal/deployment/orchestration_test.py | 108 ++++- .../deployment/test_dimension_reachability.py | 47 +-- .../tests/models/deployment_test.py | 24 +- 10 files changed, 925 insertions(+), 113 deletions(-) diff --git a/datajunction-clients/python/datajunction/models.py b/datajunction-clients/python/datajunction/models.py index bfa0cf9c6..132f28c60 100644 --- a/datajunction-clients/python/datajunction/models.py +++ b/datajunction-clients/python/datajunction/models.py @@ -3,8 +3,9 @@ from __future__ import annotations import enum +import re from dataclasses import dataclass, field -from typing import TYPE_CHECKING, Any +from typing import TYPE_CHECKING, Any, Literal, TypeAlias from datajunction._base import SerializableMixin @@ -361,6 +362,42 @@ def from_dict( # --------------------------------------------------------------------------- +@dataclass +class SemanticFingerprint: + """A semantic node digest returned by the server.""" + + digest: str + version: Literal[1] = 1 + + def __post_init__(self) -> None: + if self.version != 1: + raise ValueError( + f"Unsupported semantic fingerprint version: {self.version}", + ) + if re.fullmatch(r"[0-9a-f]{64}", self.digest) is None: + raise ValueError( + "Semantic fingerprint digest must be 64 lowercase hexadecimal characters", + ) + + @classmethod + def from_dict(cls, d: dict) -> SemanticFingerprint: + return cls( + version=d.get("version", 1), + digest=d.get("digest", ""), + ) + + +SemanticFingerprintValue: TypeAlias = SemanticFingerprint | Literal["unknown"] + + +def _parse_semantic_fingerprint(value: Any) -> SemanticFingerprintValue | None: + if value is None or value == "unknown": + return value + if not isinstance(value, dict): + raise ValueError("Semantic fingerprint must be an object or 'unknown'") + return SemanticFingerprint.from_dict(value) + + @dataclass class DeploymentResult: """A single node-level result within a deployment.""" @@ -370,15 +407,22 @@ class DeploymentResult: status: str message: str = "" changed_fields: list[str] = field(default_factory=list) + deploy_type: str = "" + change_tier: Literal["none", "minor", "major"] | None = None + semantic_fingerprint: SemanticFingerprintValue | None = None @classmethod def from_dict(cls, d: dict) -> DeploymentResult: + fingerprint = d.get("semantic_fingerprint") return cls( name=d.get("name", ""), operation=d.get("operation", ""), status=d.get("status", ""), message=d.get("message", ""), changed_fields=d.get("changed_fields") or [], + deploy_type=d.get("deploy_type", ""), + change_tier=d.get("change_tier"), + semantic_fingerprint=_parse_semantic_fingerprint(fingerprint), ) @@ -392,9 +436,11 @@ class DownstreamImpact: caused_by: list[str] = field(default_factory=list) depth: int = 0 impact_type: str = "" + semantic_fingerprint: SemanticFingerprintValue | None = None @classmethod def from_dict(cls, d: dict) -> DownstreamImpact: + fingerprint = d.get("semantic_fingerprint") return cls( name=d.get("name", ""), node_type=d.get("node_type", ""), @@ -402,6 +448,7 @@ def from_dict(cls, d: dict) -> DownstreamImpact: caused_by=d.get("caused_by") or [], depth=d.get("depth", 0), impact_type=d.get("impact_type", ""), + semantic_fingerprint=_parse_semantic_fingerprint(fingerprint), ) diff --git a/datajunction-clients/python/tests/test_deploy.py b/datajunction-clients/python/tests/test_deploy.py index dc2c33822..a67455c59 100644 --- a/datajunction-clients/python/tests/test_deploy.py +++ b/datajunction-clients/python/tests/test_deploy.py @@ -1,21 +1,22 @@ import importlib.metadata import io -from pathlib import Path import time import zipfile +from pathlib import Path from unittest import mock -import pytest from unittest.mock import MagicMock, patch + +import pytest +import yaml from datajunction.deployment import DeploymentService from datajunction.exceptions import DJClientException, DJDeploymentFailure -from datajunction.models import DeploymentInfo +from datajunction.models import DeploymentInfo, SemanticFingerprint from datajunction.rendering import ( _render_error_bullets, _strip_summary_lines, print_deployment_header, print_results, ) -import yaml from rich.console import Console @@ -1337,24 +1338,41 @@ def test_get_impact_calls_api(self, tmp_path, monkeypatch): monkeypatch.delenv("DJ_DEPLOY_REPO", raising=False) # Create mock client - mock_client = MagicMock() - mock_client.get_deployment_impact.return_value = { + response = { "uuid": "dry_run", "namespace": "test.ns", "status": "success", "results": [ { "name": "test.ns.my_node", + "deploy_type": "node", "operation": "noop", "status": "success", "message": "", + "change_tier": "none", + "semantic_fingerprint": { + "version": 1, + "digest": "a" * 64, + }, + }, + ], + "downstream_impacts": [ + { + "name": "external.metric", + "node_type": "metric", + "predicted_status": "valid", + "semantic_fingerprint": { + "version": 1, + "digest": "b" * 64, + }, }, ], - "downstream_impacts": [], } + mock_client = MagicMock() + mock_client.get_deployment_impact.return_value = response svc = DeploymentService(mock_client) - result = svc.get_impact(tmp_path) + result = svc.get_impact(tmp_path, display=False) # Verify the API was called mock_client.get_deployment_impact.assert_called_once() @@ -1362,9 +1380,72 @@ def test_get_impact_calls_api(self, tmp_path, monkeypatch): assert call_args["namespace"] == "test.ns" assert "nodes" in call_args - # Verify the result is returned - assert result["namespace"] == "test.ns" - assert result["uuid"] == "dry_run" + assert result is response + parsed = DeploymentInfo.from_dict(result) + assert parsed.results[0].deploy_type == "node" + assert parsed.results[0].change_tier == "none" + assert parsed.results[0].semantic_fingerprint == SemanticFingerprint( + digest="a" * 64, + ) + assert parsed.downstream_impacts[0].semantic_fingerprint == SemanticFingerprint( + digest="b" * 64, + ) + + def test_deployment_info_parses_unknown_fingerprints(self): + parsed = DeploymentInfo.from_dict( + { + "uuid": "dry_run", + "namespace": "test.ns", + "status": "success", + "results": [{"semantic_fingerprint": "unknown"}], + "downstream_impacts": [{"semantic_fingerprint": "unknown"}], + }, + ) + + assert parsed.results[0].semantic_fingerprint == "unknown" + assert parsed.downstream_impacts[0].semantic_fingerprint == "unknown" + + with pytest.raises( + ValueError, + match="Semantic fingerprint must be an object or 'unknown'", + ): + DeploymentInfo.from_dict( + { + "results": [{"semantic_fingerprint": "invalid"}], + }, + ) + + def test_deployment_info_parses_older_impact_response(self): + parsed = DeploymentInfo.from_dict( + { + "uuid": "dry_run", + "namespace": "test.ns", + "status": "success", + "results": [ + { + "name": "test.ns.my_node", + "operation": "noop", + "status": "skipped", + }, + ], + }, + ) + assert parsed.results[0].deploy_type == "" + assert parsed.results[0].change_tier is None + assert parsed.results[0].semantic_fingerprint is None + assert parsed.downstream_impacts == [] + + @pytest.mark.parametrize("digest", ["a" * 63, "A" * 64, "g" * 64]) + def test_semantic_fingerprint_rejects_invalid_digest(self, digest): + with pytest.raises(ValueError, match="64 lowercase hexadecimal"): + SemanticFingerprint(digest=digest) + + def test_semantic_fingerprint_rejects_unknown_version(self): + with pytest.raises( + ValueError, + match="Unsupported semantic fingerprint version", + ): + SemanticFingerprint(version=2, digest="a" * 64) def test_get_impact_with_namespace_override(self, tmp_path, monkeypatch): """get_impact should respect namespace override.""" diff --git a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py index 4039ba5bd..67f5b98d0 100644 --- a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py +++ b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py @@ -51,10 +51,14 @@ from datajunction_server.internal.deployment.dimension_reachability import ( DimensionReachability, ) +from datajunction_server.internal.deployment.fingerprints import ( + FingerprintMap, + build_deployment_fingerprints, +) from datajunction_server.internal.deployment.utils import ( DeploymentContext, classify_parents, - extract_dimension_refs_from_filters as _extract_dimension_refs_from_filters, + extract_dimension_refs_from_filters, extract_node_graph, topological_levels, ) @@ -104,6 +108,7 @@ SourceSpec, TagSpec, bump_version, + change_tier_name, declared_materialization_blocks, eq_or_fallback, render_prefixes, @@ -317,12 +322,21 @@ class DeploymentPlan: node_graph: dict[str, list[str]] external_deps: set[str] to_delete_namespaces: list[str] = field(default_factory=list) + delete_references: dict[str, list[str]] = field(default_factory=dict) def is_empty(self) -> bool: return ( not self.to_deploy and not self.to_delete and not self.to_delete_namespaces ) + @property + def deletable_specs(self) -> list[NodeSpec]: + return [ + spec + for spec in self.to_delete + if spec.rendered_name not in self.delete_references + ] + @property def linked_dimension_nodes(self) -> set[str]: return { @@ -368,6 +382,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] = {} + self._current_semantic_fingerprints: FingerprintMap = {} + self._proposed_semantic_fingerprints: FingerprintMap = {} @property def _history_user(self) -> str: @@ -510,6 +526,21 @@ async def _plan_and_execute(self) -> DeploymentExecuteResult: # SAVEPOINT, so setup-phase writes roll back too. await self._authorize_deployment_plan(deployment_plan) + with self._timer.phase("build semantic fingerprints"): + ( + self._current_semantic_fingerprints, + self._proposed_semantic_fingerprints, + ) = await build_deployment_fingerprints( + self.session, + deployment_plan.existing_specs, + self.deployment_spec.nodes, + [], + additional_target_names={ + spec.rendered_name for spec in deployment_plan.to_delete + }, + ) + self._apply_semantic_fingerprints() + if deployment_plan.is_empty() and not self.deployment_spec.hierarchies: # Pre-aggregations still need reconciling on an otherwise-empty # deploy: to register specs, or to delete external pre-aggs that were @@ -538,6 +569,7 @@ async def _plan_and_execute(self) -> DeploymentExecuteResult: ) downstream = await self._execute_deployment_plan(deployment_plan) + self._apply_semantic_fingerprints() return DeploymentExecuteResult( results=self.deployed_results, downstream_impacts=downstream, @@ -848,6 +880,39 @@ async def _handle_no_changes(self) -> list[DeploymentResult]: await self._update_deployment_status() return self.deployed_results + def _apply_semantic_fingerprints(self) -> None: + for result in self.deployed_results: + if result.deploy_type != DeploymentResult.Type.NODE: + continue + fingerprints = ( + self._current_semantic_fingerprints + if result.operation == DeploymentResult.Operation.DELETE + else self._proposed_semantic_fingerprints + ) + result.semantic_fingerprint = fingerprints.get(result.name) + + async def _apply_downstream_semantic_fingerprints( + self, + plan: DeploymentPlan, + downstream: list, + ) -> None: + target_names = {impact.name for impact in downstream} | { + spec.rendered_name for spec in plan.to_delete + } + current, proposed = await build_deployment_fingerprints( + self.session, + plan.existing_specs, + self.deployment_spec.nodes, + plan.deletable_specs, + additional_target_names=target_names, + ) + self._current_semantic_fingerprints = current + self._proposed_semantic_fingerprints = proposed + for impact in downstream: + proposed_fingerprint = proposed.get(impact.name) + if current.get(impact.name) != proposed_fingerprint: + impact.semantic_fingerprint = proposed_fingerprint + async def _find_namespaces_to_create(self) -> set[str]: """ Identify all namespaces that need to be created based on nodes in the deployment. @@ -1593,6 +1658,7 @@ async def _create_deployment_plan( else DeploymentResult.Status.SKIPPED, operation=DeploymentResult.Operation.NOOP, message="Unchanged, still INVALID" if is_invalid else "Unchanged", + change_tier=change_tier_name(ChangeTier.NONE), ), ) @@ -1603,6 +1669,7 @@ async def _create_deployment_plan( with self._timer.phase(" plan: extract node graph") as p: node_graph = extract_node_graph( [node for node in to_deploy if not isinstance(node, CubeSpec)], + tolerate_parse_errors=self.dry_run, ) p.append(f"{len(node_graph)} nodes in graph") with self._timer.phase(" plan: check external deps") as p: @@ -1724,6 +1791,7 @@ async def _create_deployment_plan( node_graph = extract_node_graph( [node for node in to_deploy if not isinstance(node, CubeSpec)], + tolerate_parse_errors=self.dry_run, ) # Re-check external deps (should be empty now) @@ -1831,23 +1899,29 @@ async def _execute_deployment_plan(self, plan: DeploymentPlan) -> list: if r.deploy_type == DeploymentResult.Type.LINK and r.status != DeploymentResult.Status.SKIPPED } + plan.delete_references = await self._validate_node_deletion(plan.to_delete) with timer.phase("propagate impact") as p: downstream = await propagate_impact( session=self.session, namespace=self.deployment_spec.namespace, changed_node_names=changed_names, deleted_node_names=frozenset( - spec.rendered_name for spec in plan.to_delete + spec.rendered_name for spec in plan.deletable_specs ), changed_link_node_names=changed_link_names, ) p.append(f"{len(downstream)} downstream") + with timer.phase("build downstream semantic fingerprints"): + await self._apply_downstream_semantic_fingerprints(plan, downstream) # Hard-delete after impact propagation (cascade-deletes # NodeRelationship rows that were needed for the BFS above). if plan.to_delete: with timer.phase("delete nodes") as p: - delete_results = await self._delete_nodes(plan.to_delete) + delete_results = await self._delete_nodes( + plan.to_delete, + references=plan.delete_references, + ) p.append(f"{len(delete_results)} deleted") self.deployed_results.extend(delete_results) await self._update_deployment_status() @@ -3243,7 +3317,7 @@ async def _bulk_validate_cubes( if cube.rendered_filters: all_dim_node_names |= { node_name - for node_name, _ in _extract_dimension_refs_from_filters( + for node_name, _ in extract_dimension_refs_from_filters( cube.rendered_filters, ) } @@ -3522,7 +3596,7 @@ def _validate_single_cube( # reachable only under a role can deploy green and flip on revalidation, # the same gap just closed above for cube dimensions. if cube_spec.rendered_filters and cube_parent_rev_ids: - filter_refs = _extract_dimension_refs_from_filters( + filter_refs = extract_dimension_refs_from_filters( cube_spec.rendered_filters, ) filter_dim_nodes = {node_name for node_name, _ in filter_refs} @@ -3850,6 +3924,12 @@ async def _create_cubes_from_validation( + ("\n".join([""] + changelog)) + invalid_note, changed_fields=changed_fields, + change_tier=change_tier_name( + change_tier if existing else ChangeTier.MAJOR, + ), + semantic_fingerprint=self._proposed_semantic_fingerprints.get( + cube_spec.rendered_name, + ), ) deployment_results.append(deployment_result) @@ -3991,7 +4071,13 @@ async def _validate_node_deletion( references: dict[str, list[str]] = {} # Query just IDs and names of nodes being deleted (more efficient than loading full objects) - stmt = select(Node.id, Node.name).where(Node.name.in_(list(nodes_to_delete))) + # Hold the target rows until commit so no concurrent transaction can add + # a new foreign-key reference after this validation snapshot. + stmt = ( + select(Node.id, Node.name) + .where(Node.name.in_(list(nodes_to_delete))) + .with_for_update() + ) result = await self.session.execute(stmt) id_to_name = {node_id: node_name for node_id, node_name in result} deleted_node_ids = set(id_to_name.keys()) @@ -4075,11 +4161,16 @@ async def _validate_node_deletion( return references - async def _delete_nodes(self, to_delete: list[NodeSpec]) -> list[DeploymentResult]: + async def _delete_nodes( + self, + to_delete: list[NodeSpec], + references: dict[str, list[str]] | None = None, + ) -> list[DeploymentResult]: logger.info("Starting deletion of %d nodes", len(to_delete)) # Check which nodes have references that would prevent deletion - references = await self._validate_node_deletion(to_delete) + if references is None: + references = await self._validate_node_deletion(to_delete) # Bulk-delete every deletable node via the shared ``hard_delete_nodes`` # helper (the same machinery ``hard_delete_namespace`` uses): it @@ -4127,6 +4218,7 @@ async def _delete_nodes(self, to_delete: list[NodeSpec]) -> list[DeploymentResul results = [] for node_spec in to_delete: node_name = node_spec.rendered_name + semantic_fingerprint = self._current_semantic_fingerprints.get(node_name) if node_name in references: # Node has references - skip deletion and return FAILED result referencing_nodes = references[node_name] @@ -4142,6 +4234,8 @@ async def _delete_nodes(self, to_delete: list[NodeSpec]) -> list[DeploymentResul status=DeploymentResult.Status.FAILED, operation=DeploymentResult.Operation.DELETE, message=error_msg, + change_tier=change_tier_name(ChangeTier.MAJOR), + semantic_fingerprint=semantic_fingerprint, ), ) elif node_name in deleted_names: @@ -4152,6 +4246,8 @@ async def _delete_nodes(self, to_delete: list[NodeSpec]) -> list[DeploymentResul status=DeploymentResult.Status.SUCCESS, operation=DeploymentResult.Operation.DELETE, message=f"Node {node_name} has been removed.", + change_tier=change_tier_name(ChangeTier.MAJOR), + semantic_fingerprint=semantic_fingerprint, ), ) else: @@ -4163,6 +4259,8 @@ async def _delete_nodes(self, to_delete: list[NodeSpec]) -> list[DeploymentResul status=DeploymentResult.Status.FAILED, operation=DeploymentResult.Operation.DELETE, message=f"Node {node_name} not found.", + change_tier=change_tier_name(ChangeTier.MAJOR), + semantic_fingerprint=semantic_fingerprint, ), ) @@ -4263,9 +4361,28 @@ def filter_nodes_to_deploy( existing_spec = existing_nodes_map.get(node_spec.rendered_name) if not existing_spec: to_create.append(node_spec) - elif force or node_spec != existing_spec: - to_update.append(node_spec) else: + if force: + to_update.append(node_spec) + continue + resolved_columns = None + proposed_columns = None + if isinstance(existing_spec, SourceSpec) and isinstance( + node_spec, + SourceSpec, + ): + if not existing_spec.columns: + resolved_columns = node_spec.columns + if not node_spec.columns: + proposed_columns = existing_spec.columns + changed_fields, reordered_fields = existing_spec.semantic_diff( + node_spec, + resolved_columns=resolved_columns, + other_resolved_columns=proposed_columns, + ) + if changed_fields or reordered_fields: + to_update.append(node_spec) + continue # Re-deploy unchanged nodes that are stuck in INVALID state so # they get revalidated (e.g. after an upstream fix). existing_node = self.registry.nodes.get(node_spec.rendered_name) @@ -4897,6 +5014,12 @@ async def _process_valid_node_deploy( + ("\n".join([""] + changelog)) + invalid_note, changed_fields=changed_fields, + change_tier=change_tier_name( + change_tier if existing else ChangeTier.MAJOR, + ), + semantic_fingerprint=self._proposed_semantic_fingerprints.get( + result.spec.rendered_name, + ), ) return deployment_result, new_node, new_revision @@ -4941,59 +5064,42 @@ async def _generate_changelog( f"└─ Set properties for {sum(changed_count)} columns", ) - # Track changes to other node fields + # Classify changes from the same normalized values used by fingerprints. existing_node_spec = await existing.to_spec(self.session) - changed_fields = existing_node_spec.diff(result.spec) if existing else [] - - # Check if query changed (diff() ignores it, but we want to surface it) - if hasattr( - existing_node_spec, - "rendered_query", - ) and hasattr( # pragma: no branch + existing_columns: list[ColumnSpec] | None = None + proposed_columns: list[ColumnSpec] | None = result.inferred_columns + if isinstance(existing_node_spec, SourceSpec) and isinstance( result.spec, - "rendered_query", + SourceSpec, ): - old_query = existing_node_spec.rendered_query - new_query = result.spec.rendered_spec().rendered_query - if old_query != new_query: - changed_fields = ["query"] + changed_fields + if not existing_node_spec.columns: + existing_columns = proposed_columns + if not proposed_columns: + proposed_columns = existing_node_spec.columns + changed_fields, reordered_fields = existing_node_spec.semantic_diff( + result.spec, + resolved_columns=existing_columns, + other_resolved_columns=proposed_columns, + ) - # Check if column metadata changed (diff() ignores columns) + # Keep detailed column notes for the human-readable changelog. from datajunction_server.models.deployment import LinkableNodeSpec as LNS - if isinstance(result.spec, LNS) and isinstance(existing_node_spec, LNS): + if ( + "columns" in changed_fields + and isinstance(result.spec, LNS) + and isinstance(existing_node_spec, LNS) + ): col_change_notes = _diff_column_metadata( result.spec.rendered_spec().columns, existing_node_spec.columns, ) - if col_change_notes: - changed_fields = changed_fields + ["columns"] - for note in col_change_notes: - changelog.append(f"└─ {note}") - - # A cube's columns are derived from its metrics and dimensions; the only - # user-authored thing on them is partition config, which is exactly what - # CubeSpec.__eq__ compares. So a partition-only edit reaches the update path - # and has to be visible here too, or it would earn no version at all. - if isinstance(result.spec, CubeSpec) and isinstance( - existing_node_spec, - CubeSpec, - ): - from datajunction_server.semantic_fingerprints.normalization import ( - normalize_cube_columns, - ) - - if normalize_cube_columns( - result.spec.matched_rendered_columns, - ) != normalize_cube_columns(existing_node_spec.matched_rendered_columns): - changed_fields = changed_fields + ["columns"] + for note in col_change_notes: + changelog.append(f"└─ {note}") if changed_fields: changelog.append("└─ Updated " + ", ".join(changed_fields)) - # Fields whose contents are unchanged but whose ordering moved. diff() - # compares list fields as sets and so cannot see these on its own. - reordered_fields = existing_node_spec.order_diff(result.spec) if reordered_fields: changelog.append("└─ Reordered " + ", ".join(reordered_fields)) @@ -5274,13 +5380,18 @@ async def _create_node_revision( ) new_revision.schema_ = schema new_revision.table = table + source_columns = source_spec.columns + if not source_columns and new_node.current: + source_columns = [ + column.to_spec() for column in new_node.current.columns + ] new_revision.columns = [ self._create_column_from_spec( col, pk_columns, order=col.order if col.order is not None else idx, ) - for idx, col in enumerate(result.spec.columns) + for idx, col in enumerate(source_columns or []) ] if result.spec.node_type == NodeType.METRIC: diff --git a/datajunction-server/datajunction_server/internal/impact.py b/datajunction-server/datajunction_server/internal/impact.py index 389fba1e6..ff93eb7ef 100644 --- a/datajunction-server/datajunction_server/internal/impact.py +++ b/datajunction-server/datajunction_server/internal/impact.py @@ -16,10 +16,17 @@ from sqlalchemy import select from sqlalchemy.ext.asyncio import AsyncSession -from sqlalchemy.orm import joinedload, selectinload +from sqlalchemy.orm import aliased, joinedload, selectinload from sqlalchemy.sql.operators import is_ -from datajunction_server.database.node import Node, NodeRelationship, NodeRevision +from datajunction_server.database.column import Column as DBColumn +from datajunction_server.database.dimensionlink import DimensionLink +from datajunction_server.database.node import ( + BoundDimensionsRelationship, + Node, + NodeRelationship, + NodeRevision, +) from datajunction_server.database.user import User from datajunction_server.instrumentation.provider import get_metrics_provider from datajunction_server.internal.deployment.dimension_reachability import ( @@ -191,11 +198,37 @@ async def _build_propagation_context( # --------------------------------------------------------------------------- +async def _cube_filter_children_by_parent( + session: AsyncSession, +) -> dict[str, set[int]]: + from datajunction_server.internal.deployment.utils import ( + extract_dimension_refs_from_filters, + ) + + rows = ( + await session.execute( + select(Node.id, NodeRevision.cube_filters) + .join( + NodeRevision, + (NodeRevision.node_id == Node.id) + & (NodeRevision.version == Node.current_version), + ) + .where(Node.type == NodeType.CUBE) + .where(is_(Node.deactivated_at, None)), + ) + ).all() + children_by_parent: dict[str, set[int]] = defaultdict(set) + for node_id, filters in rows: + for parent_name, _ in extract_dimension_refs_from_filters(filters or []): + children_by_parent[parent_name].add(node_id) + return children_by_parent + + async def _propagate_via_parent_graph( session: AsyncSession, ctx: PropagationContext, ) -> list[DownstreamImpact]: - """BFS through NodeRelationship to find all downstream nodes. + """BFS through persisted semantic relationships to find downstream nodes. Returns impacts without mutating DB state — Phase 3 determines the actual impact type via revalidation. @@ -207,19 +240,68 @@ async def _propagate_via_parent_graph( visited_node_ids = set(frontier_ids) results: list[DownstreamImpact] = [] depth = 1 + cube_filter_children = await _cube_filter_children_by_parent(session) while frontier_ids: - rows = ( + relationship_rows = ( await session.execute( select(NodeRevision.node_id, NodeRelationship.parent_id) .join(NodeRelationship, NodeRelationship.child_id == NodeRevision.id) - .where(NodeRelationship.parent_id.in_(frontier_ids)), + .join(Node, Node.id == NodeRevision.node_id) + .where(NodeRelationship.parent_id.in_(frontier_ids)) + .where(Node.current_version == NodeRevision.version), + ) + ).all() + dimension_link_rows = ( + await session.execute( + select(NodeRevision.node_id, DimensionLink.dimension_id) + .join( + DimensionLink, + DimensionLink.node_revision_id == NodeRevision.id, + ) + .join(Node, Node.id == NodeRevision.node_id) + .where(DimensionLink.dimension_id.in_(frontier_ids)) + .where(Node.current_version == NodeRevision.version), + ) + ).all() + + parent_revision = aliased(NodeRevision) + metric_revision = aliased(NodeRevision) + required_dimension_rows = ( + await session.execute( + select(metric_revision.node_id, parent_revision.node_id) + .select_from(BoundDimensionsRelationship) + .join( + DBColumn, + DBColumn.id == BoundDimensionsRelationship.bound_dimension_id, + ) + .join( + parent_revision, + parent_revision.id == DBColumn.node_revision_id, + ) + .join( + metric_revision, + metric_revision.id == BoundDimensionsRelationship.metric_id, + ) + .join(Node, Node.id == metric_revision.node_id) + .where(parent_revision.node_id.in_(frontier_ids)) + .where(Node.current_version == metric_revision.version), ) ).all() child_to_parents: dict[int, set[int]] = {} - for child_node_id, parent_id in rows: + for child_node_id, parent_id in [ + *relationship_rows, + *dimension_link_rows, + *required_dimension_rows, + ]: child_to_parents.setdefault(child_node_id, set()).add(parent_id) + for parent_id in frontier_ids: + parent_node = ctx.visited_nodes_by_id.get(parent_id) + if parent_node is None: # pragma: no cover + continue + for child_node_id in cube_filter_children.get(parent_node.name, set()): + child_to_parents.setdefault(child_node_id, set()).add(parent_id) unvisited = [nid for nid in child_to_parents if nid not in visited_node_ids] if not unvisited: diff --git a/datajunction-server/datajunction_server/models/deployment.py b/datajunction-server/datajunction_server/models/deployment.py index 9b1377db7..3297348f0 100644 --- a/datajunction-server/datajunction_server/models/deployment.py +++ b/datajunction-server/datajunction_server/models/deployment.py @@ -40,6 +40,7 @@ NodeType, ) from datajunction_server.models.partition import Granularity, PartitionType +from datajunction_server.models.semantic_fingerprint import SemanticFingerprintValue from datajunction_server.models.unit import ( Unit, legacy_unit_to_structured, @@ -67,6 +68,19 @@ class ChangeTier(IntEnum): MAJOR = 20 +ChangeTierName = Literal["none", "minor", "major"] + + +def change_tier_name(tier: ChangeTier) -> ChangeTierName: + """Return the stable API representation of a change tier.""" + names: dict[ChangeTier, ChangeTierName] = { + ChangeTier.NONE: "none", + ChangeTier.MINOR: "minor", + ChangeTier.MAJOR: "major", + } + return names[tier] + + def fold_change_tiers(tiers: Iterable[ChangeTier]) -> ChangeTier: """ Reduce the tiers of several individual changes to the tier of the change as a @@ -1766,6 +1780,8 @@ class Type(str, Enum): operation: Operation message: str = "" changed_fields: list[str] = Field(default_factory=list) + change_tier: ChangeTierName | None = None + semantic_fingerprint: SemanticFingerprintValue | None = None class DeploymentInfo(BaseModel): diff --git a/datajunction-server/datajunction_server/models/impact.py b/datajunction-server/datajunction_server/models/impact.py index c0b67cdf1..866764904 100644 --- a/datajunction-server/datajunction_server/models/impact.py +++ b/datajunction-server/datajunction_server/models/impact.py @@ -7,6 +7,7 @@ from pydantic import BaseModel, Field from datajunction_server.models.node import NodeStatus, NodeType +from datajunction_server.models.semantic_fingerprint import SemanticFingerprintValue class ImpactType(str, Enum): @@ -34,3 +35,4 @@ class DownstreamImpact(BaseModel): # Defaulted so that deployment rows persisted before this field existed can # still be rehydrated from JSON. owners: list[str] = Field(default_factory=list) + semantic_fingerprint: SemanticFingerprintValue | None = None diff --git a/datajunction-server/tests/api/deployment_impact_test.py b/datajunction-server/tests/api/deployment_impact_test.py index b5f88e8ab..068d41487 100644 --- a/datajunction-server/tests/api/deployment_impact_test.py +++ b/datajunction-server/tests/api/deployment_impact_test.py @@ -9,8 +9,12 @@ from datajunction_server.models.deployment import ( ColumnSpec, + CubeSpec, DeploymentInfo, DeploymentSpec, + DimensionJoinLinkSpec, + DimensionSpec, + MetricSpec, SourceSpec, TransformSpec, ) @@ -39,6 +43,79 @@ async def _wait_for_deployment(client, deployment_id: str, timeout: int = 30): return (await client.get(f"/deployments/{deployment_id}")).json() +def _source( + name: str, + *, + description: str | None = None, + dimension_links: list[DimensionJoinLinkSpec] | None = None, +) -> SourceSpec: + return SourceSpec( + name=name, + catalog="default", + schema_="test", + table=name, + columns=[ColumnSpec(name="id", type="int")], + description=description, + dimension_links=dimension_links or [], + ) + + +def _dimension_project(namespace: str) -> DeploymentSpec: + return DeploymentSpec( + namespace=namespace, + nodes=[ + _source("raw"), + DimensionSpec( + name="dimension", + query="SELECT id FROM ${prefix}raw", + primary_key=["id"], + ), + ], + ) + + +def _updated_dimension_project(project: DeploymentSpec) -> DeploymentSpec: + updated = project.model_copy(deep=True) + dimension = updated.nodes[1] + assert isinstance(dimension, DimensionSpec) + dimension.query = "SELECT id FROM ${prefix}raw WHERE id IS NOT NULL" + return updated + + +def _dimension_fingerprint(project: DeploymentSpec): + return project.nodes[1].semantic_fingerprint( + parent_fingerprints=[project.nodes[0].semantic_fingerprint()], + ) + + +async def _deploy(client, *specs: DeploymentSpec): + for spec in specs: + response = await client.post( + "/deployments", + json=spec.model_dump(by_alias=True), + ) + assert response.status_code == 200 + await _wait_for_deployment(client, response.json()["uuid"]) + + +async def _impact(client, spec: DeploymentSpec): + response = await client.post( + "/deployments/impact", + json=spec.model_dump(by_alias=True), + ) + assert response.status_code == 200 + return response.json() + + +async def _impact_nodes(client, spec): + data = await _impact(client, spec) + return { + result["name"]: result + for result in data["results"] + if result["deploy_type"] == "node" + } + + class TestDeploymentImpactEndpoint: """Tests for POST /deployments/impact (orchestrator dry-run).""" @@ -76,10 +153,14 @@ async def test_impact_create_new_nodes(self, client_with_roads): assert len(node_results) == 1 assert node_results[0].name == "impact_create_test.orders" assert node_results[0].operation == "create" + assert node_results[0].change_tier == "major" + assert ( + node_results[0].semantic_fingerprint == spec.nodes[0].semantic_fingerprint() + ) @pytest.mark.asyncio async def test_impact_detects_updates(self, client_with_roads): - """After deploying a node, a dry-run with a changed query shows UPDATE.""" + """A source change updates its own result and descendant fingerprint.""" initial_spec = DeploymentSpec( namespace="impact_update_test", nodes=[ @@ -109,13 +190,13 @@ async def test_impact_detects_updates(self, client_with_roads): nodes=[ TransformSpec( name="orders_summary", - query="SELECT 1 AS order_id, 'updated' AS status FROM ${prefix}raw", + query="SELECT 1 AS order_id FROM ${prefix}raw", ), SourceSpec( name="raw", catalog="default", schema_="test", - table="raw", + table="raw_v2", columns=[ColumnSpec(name="order_id", type="int")], ), ], @@ -135,8 +216,34 @@ async def test_impact_detects_updates(self, client_with_roads): ] skip_results = [r for r in data["results"] if r["operation"] == "noop"] assert len(update_results) == 1 - assert "orders_summary" in update_results[0]["name"] + assert update_results[0]["name"] == "impact_update_test.raw" assert len(skip_results) >= 1 + assert update_results[0]["change_tier"] == "major" + assert ( + update_results[0]["semantic_fingerprint"] + == updated_spec.nodes[1].semantic_fingerprint().model_dump() + ) + updated_fingerprint = updated_spec.nodes[0].semantic_fingerprint( + parent_fingerprints=[updated_spec.nodes[1].semantic_fingerprint()], + ) + initial_fingerprint = initial_spec.nodes[0].semantic_fingerprint( + parent_fingerprints=[initial_spec.nodes[1].semantic_fingerprint()], + ) + unchanged_nodes = [ + result for result in skip_results if result["deploy_type"] == "node" + ] + assert all(result["change_tier"] == "none" for result in unchanged_nodes) + transform_result = next( + result + for result in unchanged_nodes + if result["name"] == "impact_update_test.orders_summary" + ) + assert ( + transform_result["semantic_fingerprint"] == updated_fingerprint.model_dump() + ) + assert ( + transform_result["semantic_fingerprint"] != initial_fingerprint.model_dump() + ) @pytest.mark.asyncio async def test_impact_detects_deletions(self, client_with_roads): @@ -192,6 +299,59 @@ async def test_impact_detects_deletions(self, client_with_roads): delete_results = [r for r in data["results"] if r["operation"] == "delete"] assert len(delete_results) == 1 assert "to_delete" in delete_results[0]["name"] + assert delete_results[0]["change_tier"] == "major" + assert ( + delete_results[0]["semantic_fingerprint"] + == initial_spec.nodes[1].semantic_fingerprint().model_dump() + ) + + @pytest.mark.asyncio + async def test_impact_minor_full_noop_and_forced_revalidation( + self, + client_with_roads, + ): + initial = DeploymentSpec( + namespace="impact_tiers_test", + nodes=[ + _source("one", description="Before"), + _source("two"), + ], + ) + await _deploy(client_with_roads, initial) + + minor = initial.model_copy(deep=True) + minor.nodes[0].description = "After" + by_name = await _impact_nodes(client_with_roads, minor) + assert by_name["impact_tiers_test.one"]["change_tier"] == "minor" + assert ( + by_name["impact_tiers_test.one"]["semantic_fingerprint"] + == initial.nodes[0].semantic_fingerprint().model_dump() + ) + assert by_name["impact_tiers_test.two"]["change_tier"] == "none" + + equivalent = initial.model_copy(deep=True) + equivalent.nodes[1].columns = None + noop_nodes = await _impact_nodes(client_with_roads, equivalent) + assert set(noop_nodes) == { + "impact_tiers_test.one", + "impact_tiers_test.two", + } + assert all(result["change_tier"] == "none" for result in noop_nodes.values()) + assert all(result["semantic_fingerprint"] for result in noop_nodes.values()) + assert ( + noop_nodes["impact_tiers_test.two"]["semantic_fingerprint"] + == initial.nodes[1].semantic_fingerprint().model_dump() + ) + + forced = equivalent.model_copy(update={"force": True}) + forced_nodes = await _impact_nodes(client_with_roads, forced) + assert all(result["operation"] == "update" for result in forced_nodes.values()) + assert all(result["change_tier"] == "none" for result in forced_nodes.values()) + assert all(result["semantic_fingerprint"] for result in forced_nodes.values()) + assert ( + forced_nodes["impact_tiers_test.two"]["semantic_fingerprint"] + == initial.nodes[1].semantic_fingerprint().model_dump() + ) @pytest.mark.asyncio async def test_dry_run_does_not_mutate_db(self, client_with_roads): @@ -317,5 +477,221 @@ async def test_downstream_impacts_returned_for_invalid_parent( "caused_by": ["impact_downstream_test.base"], "is_external": False, "owners": ["dj"], + "semantic_fingerprint": modified_spec.nodes[1] + .semantic_fingerprint( + parent_fingerprints=[ + modified_spec.nodes[0].semantic_fingerprint(), + ], + ) + .model_dump(), }, ] + + @pytest.mark.asyncio + async def test_external_downstream_impacts_include_changed_fingerprint( + self, + client_with_roads, + ): + parent = DeploymentSpec( + namespace="impact_external_parent", + nodes=[_source("base")], + ) + consumer = DeploymentSpec( + namespace="impact_external_consumer", + nodes=[ + TransformSpec( + name="derived", + query="SELECT id FROM impact_external_parent.base", + ), + ], + ) + await _deploy(client_with_roads, parent, consumer) + + updated_parent = parent.model_copy(deep=True) + assert isinstance(updated_parent.nodes[0], SourceSpec) + updated_parent.nodes[0].table = "base_v2" + data = await _impact(client_with_roads, updated_parent) + + external = next( + impact + for impact in data["downstream_impacts"] + if impact["name"] == "impact_external_consumer.derived" + ) + expected = consumer.nodes[0].semantic_fingerprint( + parent_fingerprints=[ + updated_parent.nodes[0].semantic_fingerprint(), + ], + ) + assert external["is_external"] is True + assert external["semantic_fingerprint"] == expected.model_dump() + + metadata_only = parent.model_copy(deep=True) + metadata_only.nodes[0].description = "Updated description" + data = await _impact(client_with_roads, metadata_only) + external = next( + impact + for impact in data["downstream_impacts"] + if impact["name"] == "impact_external_consumer.derived" + ) + assert external["semantic_fingerprint"] is None + + @pytest.mark.asyncio + async def test_unparseable_node_and_dependents_return_unknown_fingerprints( + self, + client_with_roads, + ): + initial = DeploymentSpec( + namespace="impact_unknown", + nodes=[ + _source("base"), + TransformSpec( + name="broken", + query="SELECT id FROM ${prefix}base", + ), + TransformSpec( + name="dependent", + query="SELECT id FROM ${prefix}broken", + ), + ], + ) + await _deploy(client_with_roads, initial) + + proposed = DeploymentSpec( + namespace="impact_unknown", + nodes=[ + _source("base"), + TransformSpec(name="broken", query="SELECT ("), + TransformSpec( + name="dependent", + query="SELECT id FROM ${prefix}broken", + ), + ], + ) + data = await _impact(client_with_roads, proposed) + results = { + result["name"]: result + for result in data["results"] + if result["deploy_type"] == "node" + } + + assert results["impact_unknown.broken"]["semantic_fingerprint"] == "unknown" + assert results["impact_unknown.dependent"]["semantic_fingerprint"] == "unknown" + + @pytest.mark.asyncio + async def test_blocked_delete_keeps_current_fingerprint_state( + self, + client_with_roads, + ): + parent = DeploymentSpec( + namespace="impact_blocked_delete", + nodes=[_source("keep"), _source("base")], + ) + consumer = DeploymentSpec( + namespace="impact_blocked_consumer", + nodes=[ + TransformSpec( + name="derived", + query="SELECT id FROM impact_blocked_delete.base", + ), + ], + ) + await _deploy(client_with_roads, parent, consumer) + + without_parent = parent.model_copy( + deep=True, + update={"nodes": [parent.nodes[0]]}, + ) + data = await _impact(client_with_roads, without_parent) + + delete_result = next( + result + for result in data["results"] + if result["name"] == "impact_blocked_delete.base" + ) + assert delete_result["status"] == "failed" + assert ( + delete_result["semantic_fingerprint"] + == parent.nodes[1].semantic_fingerprint().model_dump() + ) + assert data["downstream_impacts"] == [] + + @pytest.mark.asyncio + async def test_dimension_link_descendants_are_discovered( + self, + client_with_roads, + ): + parent = _dimension_project("impact_link_parent") + consumer = DeploymentSpec( + namespace="impact_link_consumer", + nodes=[ + _source( + "fact", + dimension_links=[ + DimensionJoinLinkSpec( + dimension_node="impact_link_parent.dimension", + join_on=( + "impact_link_consumer.fact.id = " + "impact_link_parent.dimension.id" + ), + ), + ], + ), + ], + ) + await _deploy(client_with_roads, parent, consumer) + + updated_parent = _updated_dimension_project(parent) + data = await _impact(client_with_roads, updated_parent) + + external = next( + impact + for impact in data["downstream_impacts"] + if impact["name"] == "impact_link_consumer.fact" + ) + expected = consumer.nodes[0].semantic_fingerprint( + parent_fingerprints=[_dimension_fingerprint(updated_parent)], + ) + assert external["semantic_fingerprint"] == expected.model_dump() + + @pytest.mark.asyncio + async def test_cube_filter_descendants_are_discovered( + self, + client_with_roads, + ): + parent = _dimension_project("impact_filter_parent") + consumer = DeploymentSpec( + namespace="impact_filter_consumer", + nodes=[ + _source("fact"), + MetricSpec( + name="metric", + query="SELECT COUNT(*) FROM ${prefix}fact", + ), + CubeSpec( + name="cube", + metrics=["${prefix}metric"], + filters=["impact_filter_parent.dimension.id > 0"], + ), + ], + ) + await _deploy(client_with_roads, parent, consumer) + + updated_parent = _updated_dimension_project(parent) + data = await _impact(client_with_roads, updated_parent) + + external = next( + impact + for impact in data["downstream_impacts"] + if impact["name"] == "impact_filter_consumer.cube" + ) + fact_fingerprint = consumer.nodes[0].semantic_fingerprint() + metric_fingerprint = consumer.nodes[1].semantic_fingerprint( + parent_fingerprints=[fact_fingerprint], + ) + expected = consumer.nodes[2].semantic_fingerprint( + parent_fingerprints=[ + _dimension_fingerprint(updated_parent), + metric_fingerprint, + ], + ) + assert external["semantic_fingerprint"] == expected.model_dump() diff --git a/datajunction-server/tests/internal/deployment/orchestration_test.py b/datajunction-server/tests/internal/deployment/orchestration_test.py index a163ba5ab..bdf7cfb03 100644 --- a/datajunction-server/tests/internal/deployment/orchestration_test.py +++ b/datajunction-server/tests/internal/deployment/orchestration_test.py @@ -57,6 +57,7 @@ MetricSpec, PartitionSpec, PartitionType, + SemanticFingerprint, SourceSpec, TagSpec, TransformSpec, @@ -837,6 +838,47 @@ def test_filter_nodes_to_deploy_without_force( assert len(to_skip) == len(sample_deployment_spec.nodes) assert to_delete == [] + def test_filter_nodes_uses_normalized_query_and_source_columns(self): + incoming = [ + TransformSpec(name="transform", query=" SELECT id\nFROM source "), + SourceSpec( + name="source", + catalog="catalog", + schema_="schema", + table="table", + columns=None, + ), + ] + orchestrator = DeploymentOrchestrator( + deployment_spec=DeploymentSpec(namespace="test", nodes=incoming), + deployment_id="normalized-filter", + session=MagicMock(), + context=MagicMock(), + ) + existing = { + "test.transform": TransformSpec( + name="transform", + namespace="test", + query="SELECT id FROM source", + ), + "test.source": SourceSpec( + name="source", + namespace="test", + catalog="catalog", + schema_="schema", + table="table", + columns=[ColumnSpec(name="id", type="bigint")], + ), + } + + to_deploy, to_skip, _ = orchestrator.filter_nodes_to_deploy(existing) + assert to_deploy == [] + assert to_skip == incoming + + incoming[1].columns = [ColumnSpec(name="id", type="string")] + to_deploy, _, _ = orchestrator.filter_nodes_to_deploy(existing) + assert to_deploy == [incoming[1]] + def test_filter_nodes_to_deploy_with_force( self, session, @@ -897,6 +939,7 @@ async def test_execute_full_deployment_success(self, orchestrator): mock_plan.to_deploy = [] mock_plan.to_delete = [] mock_plan.to_delete_namespaces = [] + mock_plan.existing_specs = {} mock_create_plan.return_value = (mock_plan, []) # Execute @@ -924,6 +967,7 @@ async def test_execute_empty_deployment(self, orchestrator): mock_plan.to_deploy = [] mock_plan.to_delete = [] mock_plan.to_delete_namespaces = [] + mock_plan.existing_specs = {} mock_create_plan.return_value = (mock_plan, []) mock_handle_no_changes.return_value = [] @@ -1157,6 +1201,10 @@ async def test_create_cubes_from_validation_invalid_cubes(self, orchestrator): orchestrator._generate_changelog = AsyncMock( return_value=([], [], ChangeTier.NONE), ) + cube_fingerprint = invalid_results[0].spec.semantic_fingerprint() + orchestrator._proposed_semantic_fingerprints = { + invalid_results[0].spec.rendered_name: cube_fingerprint, + } with patch( "datajunction_server.internal.deployment.orchestrator.get_node_namespace", @@ -1174,6 +1222,8 @@ async def test_create_cubes_from_validation_invalid_cubes(self, orchestrator): assert len(revisions) == 1 assert len(results) == 1 assert results[0].status == "invalid" + assert results[0].change_tier == "major" + assert results[0].semantic_fingerprint == cube_fingerprint @pytest.mark.asyncio async def test_cube_column_partition_applied_from_spec( @@ -2647,9 +2697,13 @@ async def test_delete_nodes_bulk_deletes_existing_node( context=mock_deployment_context, dry_run=False, ) - # "default.hard_hat" is present in the pre-loaded roads example DB. - spec = Mock() - spec.rendered_name = "default.hard_hat" + # "default.hard_hat" is present in the pre-loaded roads example DB. Its + # unparseable legacy query must not prevent deletion. + spec = TransformSpec( + name="hard_hat", + namespace="default", + query="SELECT (", + ) # No external references block the delete. with patch.object(orch, "_validate_node_deletion", AsyncMock(return_value={})): results = await orch._delete_nodes([spec]) @@ -2658,6 +2712,8 @@ async def test_delete_nodes_bulk_deletes_existing_node( assert results[0].status == DeploymentResult.Status.SUCCESS assert results[0].operation == DeploymentResult.Operation.DELETE assert results[0].name == "default.hard_hat" + assert results[0].change_tier == "major" + assert results[0].semantic_fingerprint is None # The node row is gone. gone = ( @@ -2686,6 +2742,10 @@ async def test_delete_nodes_reports_referenced_and_missing( referenced.rendered_name = "default.referenced" absent = Mock() absent.rendered_name = "default.does_not_exist" + orch._current_semantic_fingerprints = { + "default.referenced": SemanticFingerprint(digest="b" * 64), + "default.does_not_exist": SemanticFingerprint(digest="c" * 64), + } with patch.object( orch, @@ -2699,6 +2759,8 @@ async def test_delete_nodes_reports_referenced_and_missing( assert "referenced by" in by_name["default.referenced"].message assert by_name["default.does_not_exist"].status == DeploymentResult.Status.FAILED assert "not found" in by_name["default.does_not_exist"].message + assert all(result.change_tier == "major" for result in results) + assert by_name["default.referenced"].semantic_fingerprint.digest == "b" * 64 class TestGenerateChangelog: @@ -2720,7 +2782,7 @@ async def test_no_changed_fields_with_dimension_links( transform_spec = TransformSpec( name="test_node", namespace="default", - query="SELECT id FROM default.source_table", + query=" SELECT id\nFROM default.source_table ", dimension_links=[dim_link], ) # existing_spec is identical — diff() will return [] @@ -2763,6 +2825,44 @@ async def test_no_changed_fields_with_dimension_links( assert changelog == ["└─ Updated dimension_links"] assert change_tier == ChangeTier.NONE + @pytest.mark.asyncio + async def test_source_column_type_change_is_major(self): + existing_spec = SourceSpec( + name="source", + namespace="test", + catalog="catalog", + schema_="schema", + table="table", + columns=[ColumnSpec(name="id", type="bigint")], + ) + proposed = existing_spec.model_copy(deep=True) + proposed.columns[0].type = "string" + existing = MagicMock() + existing.current.columns = [] + existing.to_spec = AsyncMock(return_value=existing_spec) + orchestrator = DeploymentOrchestrator( + deployment_spec=DeploymentSpec(namespace="test", nodes=[]), + deployment_id="source-changelog", + session=MagicMock(), + context=MagicMock(), + ) + orchestrator.registry.nodes["test.source"] = existing + result = NodeValidationResult( + spec=proposed, + status=NodeStatus.VALID, + inferred_columns=proposed.columns, + errors=[], + dependencies=[], + ) + + changelog, changed_fields, tier = await orchestrator._generate_changelog( + result, + ) + + assert changed_fields == ["columns"] + assert tier == ChangeTier.MAJOR + assert changelog[-1] == "└─ Updated columns" + @pytest.mark.asyncio async def test_cube_column_change_uses_role_qualified_identity( self, diff --git a/datajunction-server/tests/internal/deployment/test_dimension_reachability.py b/datajunction-server/tests/internal/deployment/test_dimension_reachability.py index 91994efbf..49067f321 100644 --- a/datajunction-server/tests/internal/deployment/test_dimension_reachability.py +++ b/datajunction-server/tests/internal/deployment/test_dimension_reachability.py @@ -10,6 +10,9 @@ DimensionReachability, find_reference_dimensions_batch, ) +from datajunction_server.internal.deployment.utils import ( + extract_dimension_refs_from_filters, +) class TestDimensionReachabilityInMemory: @@ -271,65 +274,37 @@ async def test_build_with_local_names_no_targets(self): class TestExtractDimensionRefsFromFilters: - """Tests for _extract_dimension_refs_from_filters.""" + """Tests for extract_dimension_refs_from_filters.""" def test_single_filter(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters( + result = extract_dimension_refs_from_filters( ["ns.hard_hat.state = 'CA'"], ) assert result == [("ns.hard_hat", "state")] def test_multiple_filters(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters( + result = extract_dimension_refs_from_filters( ["ns.hard_hat.state = 'CA'", "ns.date_dim.year > 2020"], ) assert sorted(result) == [("ns.date_dim", "year"), ("ns.hard_hat", "state")] def test_empty_filters(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - assert _extract_dimension_refs_from_filters([]) == [] + assert extract_dimension_refs_from_filters([]) == [] def test_unparseable_filter(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters(["not valid sql !!!"]) + result = extract_dimension_refs_from_filters(["not valid sql !!!"]) assert result == [] def test_filter_with_no_namespace(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters(["x > 5"]) + result = extract_dimension_refs_from_filters(["x > 5"]) assert result == [] def test_filter_with_single_namespace_segment(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters(["hard_hat.state = 'CA'"]) + result = extract_dimension_refs_from_filters(["hard_hat.state = 'CA'"]) assert result == [] def test_filter_with_multiple_refs_in_one_expression(self): - from datajunction_server.internal.deployment.orchestrator import ( - _extract_dimension_refs_from_filters, - ) - - result = _extract_dimension_refs_from_filters( + result = extract_dimension_refs_from_filters( ["ns.dim_a.col1 > 5 AND ns.dim_b.col2 = 'x'"], ) assert sorted(result) == [("ns.dim_a", "col1"), ("ns.dim_b", "col2")] diff --git a/datajunction-server/tests/models/deployment_test.py b/datajunction-server/tests/models/deployment_test.py index 91507b62d..6a40a9e96 100644 --- a/datajunction-server/tests/models/deployment_test.py +++ b/datajunction-server/tests/models/deployment_test.py @@ -659,13 +659,35 @@ def test_deployment_results_property_getter(): "status": "success", "operation": "create", "message": "Created", + "change_tier": "major", + "semantic_fingerprint": { + "digest": "a" * 64, + }, + }, + { + "name": "legacy_node", + "deploy_type": "node", + "status": "invalid", + "operation": "update", + "semantic_fingerprint": "unknown", + }, + { + "name": "test_node -> test_dimension", + "deploy_type": "link", + "status": "success", + "operation": "create", }, ], ) results = deployment.deployment_results - assert len(results) == 1 + assert len(results) == 3 assert results[0].name == "test_node" assert results[0].status == DeploymentResult.Status.SUCCESS + assert results[0].change_tier == "major" + assert results[0].semantic_fingerprint == SemanticFingerprint(digest="a" * 64) + assert results[1].semantic_fingerprint == "unknown" + assert results[2].change_tier is None + assert results[2].semantic_fingerprint is None def test_deployment_spec_preserves_explicit_preagg_namespace(): From dabaafb3e019efe43ea00f60a8dd23f906fc1369 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Tue, 1 Sep 2026 18:51:27 -0700 Subject: [PATCH 02/13] Align deletion result expectations --- datajunction-server/tests/internal/deployment_test.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/datajunction-server/tests/internal/deployment_test.py b/datajunction-server/tests/internal/deployment_test.py index 21febd10f..5c5af1b2a 100644 --- a/datajunction-server/tests/internal/deployment_test.py +++ b/datajunction-server/tests/internal/deployment_test.py @@ -516,6 +516,7 @@ async def test_delete_nodes_success( status=DeploymentResult.Status.SUCCESS, operation=DeploymentResult.Operation.DELETE, message="Node catalog.dim.categories has been removed.", + change_tier="major", ), ] assert await Node.get_by_name(session, categories.name) is None @@ -542,6 +543,7 @@ async def test_delete_nodes_missing( status=DeploymentResult.Status.FAILED, operation=DeploymentResult.Operation.DELETE, message="Node catalog.dim.categoriesbogus not found.", + change_tier="major", ), ] From 8a96d1f8cb988eb168fad3315bcc0974535c3605 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Wed, 2 Sep 2026 12:47:06 -0700 Subject: [PATCH 03/13] Use graph snapshots in fingerprint expectations Impact tests now derive expected hashes from complete graph snapshots, matching production parent resolution. --- .../tests/api/deployment_impact_test.py | 70 +++++++++---------- .../internal/deployment/orchestration_test.py | 10 ++- 2 files changed, 43 insertions(+), 37 deletions(-) diff --git a/datajunction-server/tests/api/deployment_impact_test.py b/datajunction-server/tests/api/deployment_impact_test.py index 068d41487..b26fdfde9 100644 --- a/datajunction-server/tests/api/deployment_impact_test.py +++ b/datajunction-server/tests/api/deployment_impact_test.py @@ -7,6 +7,9 @@ import pytest +from datajunction_server.internal.deployment.fingerprints import ( + SemanticFingerprintGraph, +) from datajunction_server.models.deployment import ( ColumnSpec, CubeSpec, @@ -82,9 +85,9 @@ def _updated_dimension_project(project: DeploymentSpec) -> DeploymentSpec: return updated -def _dimension_fingerprint(project: DeploymentSpec): - return project.nodes[1].semantic_fingerprint( - parent_fingerprints=[project.nodes[0].semantic_fingerprint()], +def _fingerprint_graph(*projects: DeploymentSpec) -> SemanticFingerprintGraph: + return SemanticFingerprintGraph( + {node.rendered_name: node for project in projects for node in project.nodes}, ) @@ -154,9 +157,10 @@ async def test_impact_create_new_nodes(self, client_with_roads): assert node_results[0].name == "impact_create_test.orders" assert node_results[0].operation == "create" assert node_results[0].change_tier == "major" - assert ( - node_results[0].semantic_fingerprint == spec.nodes[0].semantic_fingerprint() + expected = _fingerprint_graph(spec).fingerprint( + spec.nodes[0].rendered_name, ) + assert node_results[0].semantic_fingerprint == expected @pytest.mark.asyncio async def test_impact_detects_updates(self, client_with_roads): @@ -221,13 +225,15 @@ async def test_impact_detects_updates(self, client_with_roads): assert update_results[0]["change_tier"] == "major" assert ( update_results[0]["semantic_fingerprint"] - == updated_spec.nodes[1].semantic_fingerprint().model_dump() + == _fingerprint_graph(updated_spec) + .fingerprint(updated_spec.nodes[1].rendered_name) + .model_dump() ) - updated_fingerprint = updated_spec.nodes[0].semantic_fingerprint( - parent_fingerprints=[updated_spec.nodes[1].semantic_fingerprint()], + updated_fingerprint = _fingerprint_graph(updated_spec).fingerprint( + updated_spec.nodes[0].rendered_name, ) - initial_fingerprint = initial_spec.nodes[0].semantic_fingerprint( - parent_fingerprints=[initial_spec.nodes[1].semantic_fingerprint()], + initial_fingerprint = _fingerprint_graph(initial_spec).fingerprint( + initial_spec.nodes[0].rendered_name, ) unchanged_nodes = [ result for result in skip_results if result["deploy_type"] == "node" @@ -302,7 +308,9 @@ async def test_impact_detects_deletions(self, client_with_roads): assert delete_results[0]["change_tier"] == "major" assert ( delete_results[0]["semantic_fingerprint"] - == initial_spec.nodes[1].semantic_fingerprint().model_dump() + == _fingerprint_graph(initial_spec) + .fingerprint(initial_spec.nodes[1].rendered_name) + .model_dump() ) @pytest.mark.asyncio @@ -318,6 +326,7 @@ async def test_impact_minor_full_noop_and_forced_revalidation( ], ) await _deploy(client_with_roads, initial) + initial_graph = _fingerprint_graph(initial) minor = initial.model_copy(deep=True) minor.nodes[0].description = "After" @@ -325,7 +334,7 @@ async def test_impact_minor_full_noop_and_forced_revalidation( assert by_name["impact_tiers_test.one"]["change_tier"] == "minor" assert ( by_name["impact_tiers_test.one"]["semantic_fingerprint"] - == initial.nodes[0].semantic_fingerprint().model_dump() + == initial_graph.fingerprint(initial.nodes[0].rendered_name).model_dump() ) assert by_name["impact_tiers_test.two"]["change_tier"] == "none" @@ -340,7 +349,7 @@ async def test_impact_minor_full_noop_and_forced_revalidation( assert all(result["semantic_fingerprint"] for result in noop_nodes.values()) assert ( noop_nodes["impact_tiers_test.two"]["semantic_fingerprint"] - == initial.nodes[1].semantic_fingerprint().model_dump() + == initial_graph.fingerprint(initial.nodes[1].rendered_name).model_dump() ) forced = equivalent.model_copy(update={"force": True}) @@ -350,7 +359,7 @@ async def test_impact_minor_full_noop_and_forced_revalidation( assert all(result["semantic_fingerprint"] for result in forced_nodes.values()) assert ( forced_nodes["impact_tiers_test.two"]["semantic_fingerprint"] - == initial.nodes[1].semantic_fingerprint().model_dump() + == initial_graph.fingerprint(initial.nodes[1].rendered_name).model_dump() ) @pytest.mark.asyncio @@ -477,11 +486,9 @@ async def test_downstream_impacts_returned_for_invalid_parent( "caused_by": ["impact_downstream_test.base"], "is_external": False, "owners": ["dj"], - "semantic_fingerprint": modified_spec.nodes[1] - .semantic_fingerprint( - parent_fingerprints=[ - modified_spec.nodes[0].semantic_fingerprint(), - ], + "semantic_fingerprint": _fingerprint_graph(modified_spec) + .fingerprint( + modified_spec.nodes[1].rendered_name, ) .model_dump(), }, @@ -517,10 +524,8 @@ async def test_external_downstream_impacts_include_changed_fingerprint( for impact in data["downstream_impacts"] if impact["name"] == "impact_external_consumer.derived" ) - expected = consumer.nodes[0].semantic_fingerprint( - parent_fingerprints=[ - updated_parent.nodes[0].semantic_fingerprint(), - ], + expected = _fingerprint_graph(updated_parent, consumer).fingerprint( + consumer.nodes[0].rendered_name, ) assert external["is_external"] is True assert external["semantic_fingerprint"] == expected.model_dump() @@ -611,7 +616,9 @@ async def test_blocked_delete_keeps_current_fingerprint_state( assert delete_result["status"] == "failed" assert ( delete_result["semantic_fingerprint"] - == parent.nodes[1].semantic_fingerprint().model_dump() + == _fingerprint_graph(parent) + .fingerprint(parent.nodes[1].rendered_name) + .model_dump() ) assert data["downstream_impacts"] == [] @@ -648,8 +655,8 @@ async def test_dimension_link_descendants_are_discovered( for impact in data["downstream_impacts"] if impact["name"] == "impact_link_consumer.fact" ) - expected = consumer.nodes[0].semantic_fingerprint( - parent_fingerprints=[_dimension_fingerprint(updated_parent)], + expected = _fingerprint_graph(updated_parent, consumer).fingerprint( + consumer.nodes[0].rendered_name, ) assert external["semantic_fingerprint"] == expected.model_dump() @@ -684,14 +691,7 @@ async def test_cube_filter_descendants_are_discovered( for impact in data["downstream_impacts"] if impact["name"] == "impact_filter_consumer.cube" ) - fact_fingerprint = consumer.nodes[0].semantic_fingerprint() - metric_fingerprint = consumer.nodes[1].semantic_fingerprint( - parent_fingerprints=[fact_fingerprint], - ) - expected = consumer.nodes[2].semantic_fingerprint( - parent_fingerprints=[ - _dimension_fingerprint(updated_parent), - metric_fingerprint, - ], + expected = _fingerprint_graph(updated_parent, consumer).fingerprint( + consumer.nodes[2].rendered_name, ) assert external["semantic_fingerprint"] == expected.model_dump() diff --git a/datajunction-server/tests/internal/deployment/orchestration_test.py b/datajunction-server/tests/internal/deployment/orchestration_test.py index bdf7cfb03..8f3c725bb 100644 --- a/datajunction-server/tests/internal/deployment/orchestration_test.py +++ b/datajunction-server/tests/internal/deployment/orchestration_test.py @@ -24,6 +24,9 @@ from datajunction_server.database.tag import Tag from datajunction_server.database.user import OAuthProvider, User from datajunction_server.errors import DJError, DJInvalidDeploymentConfig, ErrorCode +from datajunction_server.internal.deployment.fingerprints import ( + SemanticFingerprintGraph, +) from datajunction_server.internal.deployment.orchestrator import ( DeploymentOrchestrator, DeploymentPlan, @@ -57,7 +60,6 @@ MetricSpec, PartitionSpec, PartitionType, - SemanticFingerprint, SourceSpec, TagSpec, TransformSpec, @@ -67,6 +69,7 @@ from datajunction_server.models.node import MetricUnit, NodeStatus from datajunction_server.models.node_type import NodeType from datajunction_server.models.partition import Granularity +from datajunction_server.models.semantic_fingerprint import SemanticFingerprint from datajunction_server.sql.parsing.types import StringType @@ -1201,7 +1204,10 @@ async def test_create_cubes_from_validation_invalid_cubes(self, orchestrator): orchestrator._generate_changelog = AsyncMock( return_value=([], [], ChangeTier.NONE), ) - cube_fingerprint = invalid_results[0].spec.semantic_fingerprint() + cube = invalid_results[0].spec + cube_fingerprint = SemanticFingerprintGraph( + {cube.rendered_name: cube}, + ).fingerprint(cube.rendered_name) orchestrator._proposed_semantic_fingerprints = { invalid_results[0].spec.rendered_name: cube_fingerprint, } From 52d4f758c952e85b23e57d6447c52062e3f8157e Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Wed, 2 Sep 2026 18:28:56 -0700 Subject: [PATCH 04/13] Keep dry-run graph extraction tolerant Recover per node during impact analysis so invalid SQL can produce unknown fingerprints without weakening normal deployment parsing. --- .../internal/deployment/orchestrator.py | 22 +++++++++++++++---- 1 file changed, 18 insertions(+), 4 deletions(-) diff --git a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py index 67f5b98d0..ca3749390 100644 --- a/datajunction-server/datajunction_server/internal/deployment/orchestrator.py +++ b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py @@ -136,6 +136,7 @@ structured_to_legacy_unit, ) from datajunction_server.sql.dag import get_metric_parents_map +from datajunction_server.sql.parsing.backends.exceptions import DJParseException from datajunction_server.typing import UTCDatetime from datajunction_server.utils import ( SEPARATOR, @@ -1606,6 +1607,21 @@ async def _build_copy_plan(self) -> DeploymentPlan: external_deps=external_dep_names, ) + def _extract_plan_node_graph( + self, + nodes: list[NodeSpec], + ) -> dict[str, list[str]]: + if not self.dry_run: + return extract_node_graph(nodes) + + node_graph = {} + for node in nodes: + try: + node_graph.update(extract_node_graph([node])) + except DJParseException: + node_graph[node.rendered_name] = [] + return node_graph + async def _create_deployment_plan( self, ) -> tuple[DeploymentPlan, list[DeploymentResult]]: @@ -1667,9 +1683,8 @@ async def _create_deployment_plan( external_deps: set[str] = set() if to_deploy or to_delete: with self._timer.phase(" plan: extract node graph") as p: - node_graph = extract_node_graph( + node_graph = self._extract_plan_node_graph( [node for node in to_deploy if not isinstance(node, CubeSpec)], - tolerate_parse_errors=self.dry_run, ) p.append(f"{len(node_graph)} nodes in graph") with self._timer.phase(" plan: check external deps") as p: @@ -1789,9 +1804,8 @@ async def _create_deployment_plan( existing_node.name, ) - node_graph = extract_node_graph( + node_graph = self._extract_plan_node_graph( [node for node in to_deploy if not isinstance(node, CubeSpec)], - tolerate_parse_errors=self.dry_run, ) # Re-check external deps (should be empty now) From 216e1f8e7088ce483257d3d3d3d8c275b8bb82ce Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Fri, 4 Sep 2026 14:19:41 -0700 Subject: [PATCH 05/13] Fix deployment impact regressions after rebase Canonicalize required-dimension identities during change detection so equivalent metric specs remain no-ops, and align deployment assertions with the additive fingerprint response fields. --- .../semantic_fingerprints/normalization.py | 30 ++++++++++++++++--- .../tests/api/deployments_test.py | 21 +++++++------ datajunction-server/tests/api/git_test.py | 18 ++++++++++- .../tests/models/deployment_test.py | 15 ++++++++++ 4 files changed, 70 insertions(+), 14 deletions(-) diff --git a/datajunction-server/datajunction_server/semantic_fingerprints/normalization.py b/datajunction-server/datajunction_server/semantic_fingerprints/normalization.py index 41b034f36..90c1202ef 100644 --- a/datajunction-server/datajunction_server/semantic_fingerprints/normalization.py +++ b/datajunction-server/datajunction_server/semantic_fingerprints/normalization.py @@ -211,6 +211,28 @@ def normalize_field( ) +def _normalize_comparison_field( + spec: NodeSpec, + field: str, + *, + resolved_columns: list[ColumnSpec] | None, + preserve_order: bool = False, +) -> Any: + """Normalize change detection without altering versioned fingerprints.""" + if field == "required_dimensions" and isinstance(spec, MetricSpec): + try: + value = spec.canonical_required_dimensions + except DJParseException: + value = spec.rendered_required_dimensions + return normalize_sequence(value, preserve_order=preserve_order) + return normalize_field( + spec, + field, + resolved_columns=resolved_columns, + preserve_order=preserve_order, + ) + + def semantic_diff( one: NodeSpec, two: NodeSpec, @@ -239,12 +261,12 @@ def semantic_diff( continue try: - left = normalize_field( + left = _normalize_comparison_field( rendered_one, field, resolved_columns=resolved_columns, ) - right = normalize_field( + right = _normalize_comparison_field( rendered_two, field, resolved_columns=other_resolved_columns, @@ -260,13 +282,13 @@ def semantic_diff( if type(rendered_two).field_order_change_tier(field) == ChangeTier.NONE: continue - left_ordered = normalize_field( + left_ordered = _normalize_comparison_field( rendered_one, field, resolved_columns=resolved_columns, preserve_order=True, ) - right_ordered = normalize_field( + right_ordered = _normalize_comparison_field( rendered_two, field, resolved_columns=other_resolved_columns, diff --git a/datajunction-server/tests/api/deployments_test.py b/datajunction-server/tests/api/deployments_test.py index daf50bda9..01adb426f 100644 --- a/datajunction-server/tests/api/deployments_test.py +++ b/datajunction-server/tests/api/deployments_test.py @@ -1345,6 +1345,9 @@ async def deploy_and_wait(client, deployment_spec: DeploymentSpec): await asyncio.sleep(1) response = await client.get(f"/deployments/{deployment_uuid}") data = response.json() + for result in data.get("results", []): + result.pop("change_tier", None) + result.pop("semantic_fingerprint", None) return data @@ -2216,8 +2219,8 @@ async def test_deploy_dimension_with_update( "name": f"{namespace}.default.hard_hat", "status": "success", "operation": "update", - "changed_fields": ["query", "columns"], - "message": "Updated dimension (v2.0)\n└─ Column removed: hard_hat_id, state\n└─ Updated query, columns", + "changed_fields": ["query"], + "message": "Updated dimension (v2.0)\n└─ Updated query", } update_us_state = next( res @@ -2276,12 +2279,12 @@ async def test_deploy_metric_with_update( ) assert metric_result == { "deploy_type": "node", - "message": "Updated metric (v2.0)\n└─ Updated query, display_name\n" + "message": "Updated metric (v2.0)\n└─ Updated query\n" "[invalid] Metric metric_update.default.avg_length_of_employment has an invalid " "query, should have an aggregate expression", "name": "metric_update.default.avg_length_of_employment", "operation": "update", - "changed_fields": ["query", "display_name"], + "changed_fields": ["query"], "status": "invalid", } @@ -2301,10 +2304,10 @@ async def test_deploy_metric_with_update( ) assert metric_result == { "deploy_type": "node", - "message": "Updated metric (v3.0)\n└─ Updated query, display_name", + "message": "Updated metric (v3.0)\n└─ Updated query", "name": "metric_update.default.avg_length_of_employment", "operation": "update", - "changed_fields": ["query", "display_name"], + "changed_fields": ["query"], "status": "success", } @@ -2968,7 +2971,7 @@ async def deploy(**overrides) -> None: await deploy(description="Hard hats, revised") response = await client.get(f"/nodes/{name}/") assert response.status_code == 200, response.json() - assert response.json()["version"] == "v2.0" + assert response.json()["version"] == "v1.1" assert response.json()["description"] == "Hard hats, revised" @pytest.mark.asyncio @@ -4423,10 +4426,10 @@ async def test_deploy_tags( ) assert data["results"][-1] == { "deploy_type": "node", - "message": "Updated dimension (v2.0)\n└─ Column removed: state_id, state_name, state_region, state_short\n└─ Updated tags, columns", + "message": "Updated dimension (v1.1)\n└─ Updated tags", "name": "node_update.default.us_state", "operation": "update", - "changed_fields": ["tags", "columns"], + "changed_fields": ["tags"], "status": "success", } node = await Node.get_by_name(session, f"{namespace}.default.us_state") diff --git a/datajunction-server/tests/api/git_test.py b/datajunction-server/tests/api/git_test.py index 58c722270..2f6450417 100644 --- a/datajunction-server/tests/api/git_test.py +++ b/datajunction-server/tests/api/git_test.py @@ -4402,7 +4402,23 @@ async def test_branch_creation_copies_nodes( assert response.status_code == HTTPStatus.CREATED data = response.json() assert data["branch"]["namespace"] == "copy_test.feature_copy" - assert data["deployment_results"] == [ + assert all( + result["change_tier"] == "major" + for result in data["deployment_results"] + ) + assert all( + result["semantic_fingerprint"]["version"] == 1 + for result in data["deployment_results"] + ) + deployment_results = [ + { + key: value + for key, value in result.items() + if key not in {"change_tier", "semantic_fingerprint"} + } + for result in data["deployment_results"] + ] + assert deployment_results == [ { "deploy_type": "node", "message": "Created source (v1.0)", diff --git a/datajunction-server/tests/models/deployment_test.py b/datajunction-server/tests/models/deployment_test.py index 6a40a9e96..11eb29f5f 100644 --- a/datajunction-server/tests/models/deployment_test.py +++ b/datajunction-server/tests/models/deployment_test.py @@ -1928,6 +1928,21 @@ def test_semantic_diff_and_fingerprint_share_change_rules(): assert MetricSpec.change_tier(changed, reordered) == ChangeTier.MINOR +def test_semantic_diff_canonicalizes_required_dimension_identity(): + bare = MetricSpec( + namespace="analytics", + name="order_count", + query="SELECT COUNT(*) FROM analytics.orders", + required_dimensions=["order_id"], + ) + qualified = bare.model_copy( + update={"required_dimensions": ["analytics.orders.order_id"]}, + ) + + assert bare.canonical_required_dimensions == qualified.canonical_required_dimensions + assert bare.semantic_diff(qualified) == ([], []) + + def test_semantic_diff_compares_unparseable_queries_as_raw_sql(): original = TransformSpec(name="node", query="SELECT (") same = TransformSpec(name="node", query="SELECT (") From a0169f617b0ea4f02624a2b97a00f0fe3b758150 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Tue, 1 Sep 2026 18:48:30 -0700 Subject: [PATCH 06/13] Add namespace semantic fingerprint snapshots Expose committed namespace-tree hashes so metadata files can bootstrap and reconcile outside deployment previews. --- .../datajunction_server/api/namespaces.py | 59 ++++ .../internal/deployment/fingerprints.py | 22 ++ .../datajunction_server/models/namespace.py | 19 ++ .../tests/api/namespace_fingerprints_test.py | 298 ++++++++++++++++++ 4 files changed, 398 insertions(+) create mode 100644 datajunction-server/tests/api/namespace_fingerprints_test.py diff --git a/datajunction-server/datajunction_server/api/namespaces.py b/datajunction-server/datajunction_server/api/namespaces.py index 1eeee1eb7..1b71262f4 100644 --- a/datajunction-server/datajunction_server/api/namespaces.py +++ b/datajunction-server/datajunction_server/api/namespaces.py @@ -33,6 +33,9 @@ fetch_existing_yaml_map, generate_namespace_yaml_files, ) +from datajunction_server.internal.deployment.fingerprints import ( + build_current_fingerprints, +) from datajunction_server.internal.namespaces import ( create_or_reactivate_namespace, detect_parent_cycle, @@ -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,56 @@ 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] + fingerprints = await build_current_fingerprints( + session, + specs, + version=version, + ) + 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/internal/deployment/fingerprints.py b/datajunction-server/datajunction_server/internal/deployment/fingerprints.py index 5af5acf35..ba75b47f4 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/models/namespace.py b/datajunction-server/datajunction_server/models/namespace.py index 3dfc9d948..db9792baa 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/namespace_fingerprints_test.py b/datajunction-server/tests/api/namespace_fingerprints_test.py new file mode 100644 index 000000000..d1b002d58 --- /dev/null +++ b/datajunction-server/tests/api/namespace_fingerprints_test.py @@ -0,0 +1,298 @@ +"""Tests for namespace semantic fingerprint snapshots.""" + +import asyncio +from unittest import mock + +import pytest + +from datajunction_server.database.node import Node +from datajunction_server.models.deployment import ( + ColumnSpec, + CubeSpec, + DeploymentSpec, + MetricSpec, + SourceSpec, + TransformSpec, +) + + +@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" From 72dff13ee8f65bff3281dc662ea4ada94754328c Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Wed, 2 Sep 2026 12:51:40 -0700 Subject: [PATCH 07/13] Use graph snapshots for namespace fingerprints The namespace endpoint now evaluates one current graph and selects its requested rows from the memoized result. --- .../datajunction_server/api/namespaces.py | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/datajunction-server/datajunction_server/api/namespaces.py b/datajunction-server/datajunction_server/api/namespaces.py index 1b71262f4..558eeac46 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, @@ -33,9 +36,6 @@ fetch_existing_yaml_map, generate_namespace_yaml_files, ) -from datajunction_server.internal.deployment.fingerprints import ( - build_current_fingerprints, -) from datajunction_server.internal.namespaces import ( create_or_reactivate_namespace, detect_parent_cycle, @@ -328,11 +328,12 @@ async def get_namespace_semantic_fingerprints( ordered_nodes = sorted(nodes, key=lambda node: node.name) specs = [await node.to_spec(session) for node in ordered_nodes] - fingerprints = await build_current_fingerprints( + 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, From 40f37e027f84a9d75ae7cf51e6875a7417b073d9 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Wed, 9 Sep 2026 17:06:14 -0700 Subject: [PATCH 08/13] Document semantic fingerprint version migrations --- .../semantic-fingerprint-versioning.md | 155 ++++++++++++++++++ 1 file changed, 155 insertions(+) create mode 100644 docs/content/0.1.0/docs/developers/semantic-fingerprint-versioning.md 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 000000000..3698f2f75 --- /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. From eb42c5dfcc84d4113d2bc52bc7b4118f4695bf53 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Fri, 11 Sep 2026 13:04:31 -0700 Subject: [PATCH 09/13] Bump version to v0.0.242 (#2538) Co-authored-by: Philip Weiss --- datajunction-clients/javascript/package-lock.json | 4 ++-- datajunction-clients/javascript/package.json | 2 +- datajunction-clients/python/datajunction/__about__.py | 2 +- datajunction-query/djqs/__about__.py | 2 +- datajunction-reflection/datajunction_reflection/__about__.py | 2 +- datajunction-server/datajunction_server/__about__.py | 2 +- datajunction-ui/package.json | 2 +- 7 files changed, 8 insertions(+), 8 deletions(-) diff --git a/datajunction-clients/javascript/package-lock.json b/datajunction-clients/javascript/package-lock.json index e1b404b49..6b7d3924d 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.242", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "datajunction", - "version": "0.0.241", + "version": "0.0.242", "license": "MIT", "dependencies": { "@babel/core": "^7.22.5", diff --git a/datajunction-clients/javascript/package.json b/datajunction-clients/javascript/package.json index e347e27f7..52413309a 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.242", "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 ae2e65baa..cf39e2930 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.242" diff --git a/datajunction-query/djqs/__about__.py b/datajunction-query/djqs/__about__.py index ae2e65baa..cf39e2930 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.242" diff --git a/datajunction-reflection/datajunction_reflection/__about__.py b/datajunction-reflection/datajunction_reflection/__about__.py index ae2e65baa..cf39e2930 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.242" diff --git a/datajunction-server/datajunction_server/__about__.py b/datajunction-server/datajunction_server/__about__.py index ae2e65baa..cf39e2930 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.242" diff --git a/datajunction-ui/package.json b/datajunction-ui/package.json index 38e378239..2b0f55da9 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.242", "description": "DataJunction UI", "type": "module", "module": "src/index.tsx", From 1ca2cfe394400df7d0d254724563621b43a04d1a Mon Sep 17 00:00:00 2001 From: Yian Date: Fri, 11 Sep 2026 13:51:33 -0700 Subject: [PATCH 10/13] Bump version to v0.0.243 (#2539) Co-authored-by: GitHub Actions Bot --- datajunction-clients/javascript/package-lock.json | 4 ++-- datajunction-clients/javascript/package.json | 2 +- datajunction-clients/python/datajunction/__about__.py | 2 +- datajunction-query/djqs/__about__.py | 2 +- datajunction-reflection/datajunction_reflection/__about__.py | 2 +- datajunction-server/datajunction_server/__about__.py | 2 +- datajunction-ui/package.json | 2 +- 7 files changed, 8 insertions(+), 8 deletions(-) diff --git a/datajunction-clients/javascript/package-lock.json b/datajunction-clients/javascript/package-lock.json index 6b7d3924d..41d31bc23 100644 --- a/datajunction-clients/javascript/package-lock.json +++ b/datajunction-clients/javascript/package-lock.json @@ -1,12 +1,12 @@ { "name": "datajunction", - "version": "0.0.242", + "version": "0.0.243", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "datajunction", - "version": "0.0.242", + "version": "0.0.243", "license": "MIT", "dependencies": { "@babel/core": "^7.22.5", diff --git a/datajunction-clients/javascript/package.json b/datajunction-clients/javascript/package.json index 52413309a..e93f0b3f7 100644 --- a/datajunction-clients/javascript/package.json +++ b/datajunction-clients/javascript/package.json @@ -1,6 +1,6 @@ { "name": "datajunction", - "version": "0.0.242", + "version": "0.0.243", "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 cf39e2930..7e6adf0cf 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.242" +__version__ = "0.0.243" diff --git a/datajunction-query/djqs/__about__.py b/datajunction-query/djqs/__about__.py index cf39e2930..7e6adf0cf 100644 --- a/datajunction-query/djqs/__about__.py +++ b/datajunction-query/djqs/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.242" +__version__ = "0.0.243" diff --git a/datajunction-reflection/datajunction_reflection/__about__.py b/datajunction-reflection/datajunction_reflection/__about__.py index cf39e2930..7e6adf0cf 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.242" +__version__ = "0.0.243" diff --git a/datajunction-server/datajunction_server/__about__.py b/datajunction-server/datajunction_server/__about__.py index cf39e2930..7e6adf0cf 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.242" +__version__ = "0.0.243" diff --git a/datajunction-ui/package.json b/datajunction-ui/package.json index 2b0f55da9..2865001d3 100644 --- a/datajunction-ui/package.json +++ b/datajunction-ui/package.json @@ -1,6 +1,6 @@ { "name": "datajunction-ui", - "version": "0.0.242", + "version": "0.0.243", "description": "DataJunction UI", "type": "module", "module": "src/index.tsx", From 4b39f985ec3959281b425d8a23688ccf8d27c1e6 Mon Sep 17 00:00:00 2001 From: Yian Date: Sun, 13 Sep 2026 10:08:01 -0700 Subject: [PATCH 11/13] Distinguish a pre-existing failure from a new one (#2540) * Distinguish a re-deployed pre-existing failure from a new one A node whose spec is unchanged but whose stored status is INVALID is deliberately promoted into the update set so that the deployment gets a chance to revalidate it. That is the right thing to do -- the node may well have been broken by something outside this deployment that has since been fixed -- but it made an already-broken, untouched node indistinguishable in the response from one the author just broke. A deployment then failed on nodes its change never touched. Those nodes are now marked `revalidation_only`, so a caller can tell the failures a deployment caused from the ones it inherited, and still sees the recovery when a later deploy fixes the node's upstream. The flag is not on its own enough to excuse a failure, and a later gate must not treat it that way: it does not compare failure reasons, so a node whose own spec never moved can fail for a new reason caused by an edit upstream of it. Attribution has to be "changed, or reachable from something changed", and the inherited breakage surfaces separately as a downstream impact. Additive and observational. Exit codes, gating, and overall pass/fail are untouched, and the field is nullable so deployment rows persisted before it existed still rehydrate. * Trim the comments --- .../python/datajunction/models.py | 4 + .../python/tests/test_deploy.py | 2 + .../internal/deployment/orchestrator.py | 10 +++ .../datajunction_server/models/deployment.py | 3 + .../tests/api/deployments_test.py | 86 +++++++++++++++++-- datajunction-server/tests/api/git_test.py | 11 ++- .../internal/deployment/orchestration_test.py | 29 +++++++ 7 files changed, 139 insertions(+), 6 deletions(-) diff --git a/datajunction-clients/python/datajunction/models.py b/datajunction-clients/python/datajunction/models.py index 25c0c2d4a..8eacc3819 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 e74bbd3bf..f3b6bec1e 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-server/datajunction_server/internal/deployment/orchestrator.py b/datajunction-server/datajunction_server/internal/deployment/orchestrator.py index a1626821f..711d985c4 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 4921e8144..bfc354519 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/tests/api/deployments_test.py b/datajunction-server/tests/api/deployments_test.py index ce802256b..ca250163d 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 2f6450417..02a2b07a6 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/internal/deployment/orchestration_test.py b/datajunction-server/tests/internal/deployment/orchestration_test.py index f960aa1a3..aa167a92a 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 "), From a7f632549d09df706e165203aab48d11e623b52b Mon Sep 17 00:00:00 2001 From: Yian Date: Sun, 13 Sep 2026 10:17:28 -0700 Subject: [PATCH 12/13] Bump version to v0.0.244 (#2541) Co-authored-by: GitHub Actions Bot --- datajunction-clients/javascript/package-lock.json | 4 ++-- datajunction-clients/javascript/package.json | 2 +- datajunction-clients/python/datajunction/__about__.py | 2 +- datajunction-query/djqs/__about__.py | 2 +- datajunction-reflection/datajunction_reflection/__about__.py | 2 +- datajunction-server/datajunction_server/__about__.py | 2 +- datajunction-ui/package.json | 2 +- 7 files changed, 8 insertions(+), 8 deletions(-) diff --git a/datajunction-clients/javascript/package-lock.json b/datajunction-clients/javascript/package-lock.json index 41d31bc23..6a26df6c1 100644 --- a/datajunction-clients/javascript/package-lock.json +++ b/datajunction-clients/javascript/package-lock.json @@ -1,12 +1,12 @@ { "name": "datajunction", - "version": "0.0.243", + "version": "0.0.244", "lockfileVersion": 2, "requires": true, "packages": { "": { "name": "datajunction", - "version": "0.0.243", + "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 e93f0b3f7..daa4e37e9 100644 --- a/datajunction-clients/javascript/package.json +++ b/datajunction-clients/javascript/package.json @@ -1,6 +1,6 @@ { "name": "datajunction", - "version": "0.0.243", + "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 7e6adf0cf..727425f40 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.243" +__version__ = "0.0.244" diff --git a/datajunction-query/djqs/__about__.py b/datajunction-query/djqs/__about__.py index 7e6adf0cf..727425f40 100644 --- a/datajunction-query/djqs/__about__.py +++ b/datajunction-query/djqs/__about__.py @@ -2,4 +2,4 @@ Version for Hatch """ -__version__ = "0.0.243" +__version__ = "0.0.244" diff --git a/datajunction-reflection/datajunction_reflection/__about__.py b/datajunction-reflection/datajunction_reflection/__about__.py index 7e6adf0cf..727425f40 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.243" +__version__ = "0.0.244" diff --git a/datajunction-server/datajunction_server/__about__.py b/datajunction-server/datajunction_server/__about__.py index 7e6adf0cf..727425f40 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.243" +__version__ = "0.0.244" diff --git a/datajunction-ui/package.json b/datajunction-ui/package.json index 2865001d3..877896311 100644 --- a/datajunction-ui/package.json +++ b/datajunction-ui/package.json @@ -1,6 +1,6 @@ { "name": "datajunction-ui", - "version": "0.0.243", + "version": "0.0.244", "description": "DataJunction UI", "type": "module", "module": "src/index.tsx", From 83bc23c92ba9edd2a7b29256768e5ae03ac12762 Mon Sep 17 00:00:00 2001 From: Philip Weiss Date: Mon, 14 Sep 2026 08:21:38 -0700 Subject: [PATCH 13/13] Match namespace snapshot descendants literally --- .../datajunction_server/database/namespace.py | 2 +- .../tests/api/namespace_fingerprints_test.py | 61 +++++++++++++++++++ 2 files changed, 62 insertions(+), 1 deletion(-) diff --git a/datajunction-server/datajunction_server/database/namespace.py b/datajunction-server/datajunction_server/database/namespace.py index 9a9836f63..819a9e5da 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/tests/api/namespace_fingerprints_test.py b/datajunction-server/tests/api/namespace_fingerprints_test.py index d1b002d58..f783f12fe 100644 --- a/datajunction-server/tests/api/namespace_fingerprints_test.py +++ b/datajunction-server/tests/api/namespace_fingerprints_test.py @@ -6,6 +6,8 @@ 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, @@ -14,6 +16,7 @@ SourceSpec, TransformSpec, ) +from tests.authz import VALIDATOR_AUTH_SERVICE @pytest.fixture(autouse=True, scope="module") @@ -296,3 +299,61 @@ async def test_namespace_snapshot_rejects_unsupported_version(client_with_roads) 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", + ]