From dbf09883e15efdaca311f180d5ce87734763e6f1 Mon Sep 17 00:00:00 2001 From: Sai Shree Pradhan Date: Wed, 29 Jul 2026 15:41:33 +0530 Subject: [PATCH 1/2] add connection log --- dbt/adapters/databricks/connections.py | 11 +++ dbt/adapters/databricks/handle.py | 4 + dbt/adapters/databricks/telemetry/__init__.py | 27 +++++++ dbt/adapters/databricks/telemetry/builder.py | 27 +++++++ dbt/adapters/databricks/telemetry/client.py | 73 +++++++++++++++++++ dbt/adapters/databricks/telemetry/config.py | 13 ++++ dbt/adapters/databricks/telemetry/encoder.py | 46 ++++++++++++ dbt/adapters/databricks/telemetry/hooks.py | 45 ++++++++++++ dbt/adapters/databricks/telemetry/models.py | 14 ++++ 9 files changed, 260 insertions(+) create mode 100644 dbt/adapters/databricks/telemetry/__init__.py create mode 100644 dbt/adapters/databricks/telemetry/builder.py create mode 100644 dbt/adapters/databricks/telemetry/client.py create mode 100644 dbt/adapters/databricks/telemetry/config.py create mode 100644 dbt/adapters/databricks/telemetry/encoder.py create mode 100644 dbt/adapters/databricks/telemetry/hooks.py create mode 100644 dbt/adapters/databricks/telemetry/models.py diff --git a/dbt/adapters/databricks/connections.py b/dbt/adapters/databricks/connections.py index b6825edd6..fb9e0145f 100644 --- a/dbt/adapters/databricks/connections.py +++ b/dbt/adapters/databricks/connections.py @@ -55,6 +55,7 @@ from dbt.adapters.databricks.logging import logger from dbt.adapters.databricks.python_models.run_tracking import PythonRunTracker from dbt.adapters.databricks.spog.decision import check_spog_preconditions +from dbt.adapters.databricks.telemetry import hooks as telemetry_hooks from dbt.adapters.databricks.utils import QueryTagsUtils, is_cluster_http_path, redact_credentials if TYPE_CHECKING: @@ -518,6 +519,16 @@ def connect() -> DatabricksHandle: databricks_connection.capabilities = cls._get_capabilities_for_http_path( databricks_connection.http_path ) + + telemetry_hooks.on_connection_open( + creds, + cls.credentials_manager, + session_id=conn.session_id, + http_path=databricks_connection.http_path, + is_cluster=is_cluster_http_path( + databricks_connection.http_path, creds.cluster_id + ), + ) return conn else: raise DbtDatabaseError("Failed to create connection") diff --git a/dbt/adapters/databricks/handle.py b/dbt/adapters/databricks/handle.py index 3da799c9c..c26446afe 100644 --- a/dbt/adapters/databricks/handle.py +++ b/dbt/adapters/databricks/handle.py @@ -394,6 +394,10 @@ def prepare_connection_arguments( connection_parameters = creds.connection_parameters.copy() # type: ignore[union-attr] + # dbt telemetry opt-in lives in connection_parameters but is consumed by + # the adapter's own collector; drop it so it never reaches dbsql.connect. + connection_parameters.pop("enable_dbt_telemetry", None) + http_headers: list[tuple[str, str]] = list( creds.get_all_http_headers(connection_parameters.pop("http_headers", {})).items() ) diff --git a/dbt/adapters/databricks/telemetry/__init__.py b/dbt/adapters/databricks/telemetry/__init__.py new file mode 100644 index 000000000..51c050c72 --- /dev/null +++ b/dbt/adapters/databricks/telemetry/__init__.py @@ -0,0 +1,27 @@ +from typing import Optional + +from dbt.adapters.databricks.logging import logger +from dbt.adapters.databricks.telemetry import builder, client, encoder +from dbt.adapters.databricks.telemetry.client import HeaderFactory +from dbt.adapters.databricks.telemetry.config import is_enabled + +__all__ = ["is_enabled", "send_connection_log"] + + +def send_connection_log( + host: Optional[str], + header_factory: Optional[HeaderFactory] = None, + workspace_id: Optional[int] = None, + session_id: Optional[str] = None, + http_path: Optional[str] = None, + is_cluster: Optional[bool] = None, +) -> None: + """Encode a single connection log and POST it.""" + try: + log = builder.build_connection_log( + session_id=session_id, http_path=http_path, is_cluster=is_cluster + ) + body = encoder.encode_request(log, workspace_id=workspace_id) + client.send(host, body, header_factory=header_factory, workspace_id=workspace_id) + except Exception as e: # pragma: no cover - defensive + logger.debug(f"dbt telemetry: send_connection_log failed (ignored): {e}") diff --git a/dbt/adapters/databricks/telemetry/builder.py b/dbt/adapters/databricks/telemetry/builder.py new file mode 100644 index 000000000..b5ad42653 --- /dev/null +++ b/dbt/adapters/databricks/telemetry/builder.py @@ -0,0 +1,27 @@ +"""Build the initial connection telemetry log from connection facts.""" + +from importlib.metadata import version as _pkg_version +from typing import Optional + +from databricks.sql import __version__ as dbsql_version +from dbt.adapters.databricks.__version__ import version as __version__ +from dbt.adapters.databricks.telemetry.models import ConnectionLog + + +def build_connection_log( + session_id: Optional[str] = None, + http_path: Optional[str] = None, + is_cluster: Optional[bool] = None, +) -> ConnectionLog: + """Assemble the connection log; only non-sensitive metadata.""" + compute_type = None + if is_cluster is not None: + compute_type = "cluster" if is_cluster else "sql_warehouse" + return ConnectionLog( + dbt_databricks_version=__version__, + dbt_core_version=_pkg_version("dbt-core"), + databricks_sql_connector_version=dbsql_version, + session_id=session_id, + http_path=http_path, + compute_type=compute_type, + ) diff --git a/dbt/adapters/databricks/telemetry/client.py b/dbt/adapters/databricks/telemetry/client.py new file mode 100644 index 000000000..8c9ab0958 --- /dev/null +++ b/dbt/adapters/databricks/telemetry/client.py @@ -0,0 +1,73 @@ +import json +from typing import Any, Callable, Optional + +import requests + +from dbt.adapters.databricks.logging import logger + +TELEMETRY_AUTHENTICATED_PATH = "/telemetry-ext" +TELEMETRY_UNAUTHENTICATED_PATH = "/telemetry-unauth" + +_TIMEOUT_SECONDS = 10 + +HeaderFactory = Callable[[], dict[str, str]] + + +def _normalize_host(host: str) -> str: + host = host.rstrip("/") + if not host.startswith(("http://", "https://")): + host = f"https://{host}" + return host + + +def send( + host: Optional[str], + body: dict[str, Any], + header_factory: Optional[HeaderFactory] = None, + workspace_id: Optional[int] = None, +) -> bool: + """POST a TelemetryRequest body.""" + if not host: + logger.debug("dbt telemetry: no host available; skipping send") + return False + + try: + headers = {"Accept": "application/json", "Content-Type": "application/json"} + path = TELEMETRY_AUTHENTICATED_PATH + if header_factory is not None: + try: + headers.update(header_factory()) + except Exception as e: + logger.debug(f"dbt telemetry: failed to build auth headers: {e}") + return False + else: + path = TELEMETRY_UNAUTHENTICATED_PATH + + if workspace_id is not None: + headers["x-databricks-org-id"] = str(workspace_id) + + url = _normalize_host(host) + path + + logger.debug(f"dbt telemetry: log = {json.dumps(body)}") + logger.debug(f"dbt telemetry: endpoint = {url}") + + response = requests.post(url, json=body, headers=headers, timeout=_TIMEOUT_SECONDS) + + body_preview = response.text if len(response.text) <= 500 else response.text[:500] + "…" + logger.debug(f"dbt telemetry: response = [{response.status_code}] {body_preview}") + + if response.status_code // 100 != 2: + logger.debug(f"dbt telemetry: not accepted (status {response.status_code})") + return False + try: + ack = response.json() + except ValueError: + logger.debug("dbt telemetry: not accepted (non-JSON response)") + return False + accepted = ack.get("numProtoSuccess", 0) >= 1 and not ack.get("errors") + if not accepted: + logger.debug(f"dbt telemetry: not accepted (ack: {ack})") + return accepted + except Exception as e: + logger.debug(f"dbt telemetry: send failed (ignored): {e}") + return False diff --git a/dbt/adapters/databricks/telemetry/config.py b/dbt/adapters/databricks/telemetry/config.py new file mode 100644 index 000000000..29e45266a --- /dev/null +++ b/dbt/adapters/databricks/telemetry/config.py @@ -0,0 +1,13 @@ +from typing import Optional + +from dbt.adapters.databricks.credentials import DatabricksCredentials + +ENABLE_FLAG = "enable_dbt_telemetry" + + +def is_enabled(credentials: Optional[DatabricksCredentials]) -> bool: + """True only when the user explicitly opted in on this target.""" + if credentials is None: + return False + params = credentials.connection_parameters or {} + return bool(params.get(ENABLE_FLAG, False)) diff --git a/dbt/adapters/databricks/telemetry/encoder.py b/dbt/adapters/databricks/telemetry/encoder.py new file mode 100644 index 000000000..9c630dc44 --- /dev/null +++ b/dbt/adapters/databricks/telemetry/encoder.py @@ -0,0 +1,46 @@ +import json +import time +import uuid +from typing import Any, Optional + +from dbt.adapters.databricks.telemetry.models import ConnectionLog + +DRIVER_NAME = "dbt-databricks" +CLIENT_APP_NAME = "dbt" + + +def encode_event(log: ConnectionLog, workspace_id: Optional[int] = None) -> str: + """Encode one ConnectionLog as a TelemetryFrontendLog JSON string.""" + sql_driver_log: dict[str, Any] = { + "session_id": log.session_id, + "system_configuration": { + "driver_name": DRIVER_NAME, + "driver_version": log.dbt_databricks_version, + "client_app_name": CLIENT_APP_NAME, + }, + "driver_connection_params": { + "http_path": log.http_path, + }, + } + + frontend_log: dict[str, Any] = { + "frontend_log_event_id": str(uuid.uuid4()), + "context": { + "client_context": { + "timestamp_millis": int(time.time() * 1000), + "user_agent": f"{DRIVER_NAME}/{log.dbt_databricks_version}", + } + }, + "entry": {"sql_driver_log": sql_driver_log}, + "workspace_id": workspace_id, + } + return json.dumps(frontend_log) + + +def encode_request(log: ConnectionLog, workspace_id: Optional[int] = None) -> dict[str, Any]: + """Build the TelemetryRequest body wrapping a single connection log.""" + return { + "uploadTime": int(time.time() * 1000), + "items": [], + "protoLogs": [encode_event(log, workspace_id)], + } diff --git a/dbt/adapters/databricks/telemetry/hooks.py b/dbt/adapters/databricks/telemetry/hooks.py new file mode 100644 index 000000000..5e9f0abeb --- /dev/null +++ b/dbt/adapters/databricks/telemetry/hooks.py @@ -0,0 +1,45 @@ +"""Thin hooks called from the connection manager. + +On the first successful connection we build one initial connection log and POST +it immediately, using the connection's host + auth. +""" + +from typing import Optional + +from dbt.adapters.databricks import telemetry +from dbt.adapters.databricks.credentials import ( + DatabricksCredentialManager, + DatabricksCredentials, +) +from dbt.adapters.databricks.logging import logger + +_sent = False + + +def on_connection_open( + credentials: Optional[DatabricksCredentials], + credentials_manager: Optional[DatabricksCredentialManager], + session_id: Optional[str] = None, + http_path: Optional[str] = None, + is_cluster: Optional[bool] = None, +) -> None: + """Send a single initial connection telemetry log on the first opted-in connection.""" + global _sent + if _sent: + return + if not telemetry.is_enabled(credentials): + return + if credentials_manager is None: + return + try: + _sent = True + telemetry.send_connection_log( + host=credentials_manager.host, + header_factory=credentials_manager.header_factory, + workspace_id=getattr(credentials_manager, "workspace_id", None), + session_id=session_id, + http_path=http_path, + is_cluster=is_cluster, + ) + except Exception as e: + logger.debug(f"dbt telemetry: on_connection_open failed (ignored): {e}") diff --git a/dbt/adapters/databricks/telemetry/models.py b/dbt/adapters/databricks/telemetry/models.py new file mode 100644 index 000000000..fb271476e --- /dev/null +++ b/dbt/adapters/databricks/telemetry/models.py @@ -0,0 +1,14 @@ +"""Event model for the dbt-databricks initial connection telemetry log.""" + +from dataclasses import dataclass +from typing import Optional + + +@dataclass +class ConnectionLog: + dbt_databricks_version: str + dbt_core_version: str + databricks_sql_connector_version: str + session_id: Optional[str] = None + http_path: Optional[str] = None + compute_type: Optional[str] = None # "cluster" | "sql_warehouse" From 6c6da671ca9118a6f0af54e03119d65fb77eb9f1 Mon Sep 17 00:00:00 2001 From: Sai Shree Pradhan Date: Wed, 29 Jul 2026 16:20:09 +0530 Subject: [PATCH 2/2] modified comments --- dbt/adapters/databricks/handle.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/dbt/adapters/databricks/handle.py b/dbt/adapters/databricks/handle.py index c26446afe..356fc932b 100644 --- a/dbt/adapters/databricks/handle.py +++ b/dbt/adapters/databricks/handle.py @@ -394,8 +394,6 @@ def prepare_connection_arguments( connection_parameters = creds.connection_parameters.copy() # type: ignore[union-attr] - # dbt telemetry opt-in lives in connection_parameters but is consumed by - # the adapter's own collector; drop it so it never reaches dbsql.connect. connection_parameters.pop("enable_dbt_telemetry", None) http_headers: list[tuple[str, str]] = list(