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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions dbt/adapters/databricks/connections.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down Expand Up @@ -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")
Expand Down
2 changes: 2 additions & 0 deletions dbt/adapters/databricks/handle.py
Original file line number Diff line number Diff line change
Expand Up @@ -394,6 +394,8 @@ def prepare_connection_arguments(

connection_parameters = creds.connection_parameters.copy() # type: ignore[union-attr]

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()
)
Expand Down
27 changes: 27 additions & 0 deletions dbt/adapters/databricks/telemetry/__init__.py
Original file line number Diff line number Diff line change
@@ -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}")
27 changes: 27 additions & 0 deletions dbt/adapters/databricks/telemetry/builder.py
Original file line number Diff line number Diff line change
@@ -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,
)
73 changes: 73 additions & 0 deletions dbt/adapters/databricks/telemetry/client.py
Original file line number Diff line number Diff line change
@@ -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
13 changes: 13 additions & 0 deletions dbt/adapters/databricks/telemetry/config.py
Original file line number Diff line number Diff line change
@@ -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))
46 changes: 46 additions & 0 deletions dbt/adapters/databricks/telemetry/encoder.py
Original file line number Diff line number Diff line change
@@ -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)],
}
45 changes: 45 additions & 0 deletions dbt/adapters/databricks/telemetry/hooks.py
Original file line number Diff line number Diff line change
@@ -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}")
14 changes: 14 additions & 0 deletions dbt/adapters/databricks/telemetry/models.py
Original file line number Diff line number Diff line change
@@ -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"