diff --git a/src/snowflake/snowpark/mock/_telemetry.py b/src/snowflake/snowpark/mock/_telemetry.py index 2c090d0617..c4b33d7264 100644 --- a/src/snowflake/snowpark/mock/_telemetry.py +++ b/src/snowflake/snowpark/mock/_telemetry.py @@ -5,6 +5,7 @@ import json import logging import os +import queue as _queue import threading import uuid from datetime import datetime @@ -13,7 +14,6 @@ from typing import Optional from snowflake.connector.secret_detector import SecretDetector -from snowflake.connector.telemetry_oob import TelemetryService from snowflake.snowpark._internal.utils import ( get_os_name, get_python_version, @@ -83,17 +83,29 @@ class LocalTestTelemetryEventType(Enum): SESSION_CONNECTION = "session" -class LocalTestOOBTelemetryService(TelemetryService): +class LocalTestOOBTelemetryService: PROD = "https://client-telemetry.snowflakecomputing.com/enqueue" + _instance: "LocalTestOOBTelemetryService | None" = None + _instance_lock: threading.Lock = threading.Lock() + + @classmethod + def get_instance(cls) -> "LocalTestOOBTelemetryService": + if cls._instance is None: + with cls._instance_lock: + if cls._instance is None: + cls._instance = cls() + return cls._instance + def __init__(self) -> None: - super().__init__() self._is_internal_usage = bool( os.getenv("SNOWPARK_LOCAL_TESTING_INTERNAL_TELEMETRY", False) ) self._deployment_url = self.PROD - self._enable = True + self._enabled = True self._lock = threading.RLock() + self.queue: _queue.Queue = _queue.Queue() + self.batch_size: int = 100 def _upload_payload(self, payload) -> None: if not REQUESTS_AVAILABLE: @@ -189,6 +201,9 @@ def export_queue_to_string(self): _, masked_text, _ = SecretDetector.mask_secrets(payload) return masked_text + def close(self) -> None: + self.flush() + def log_session_creation(self, connection_uuid: Optional[str] = None): try: telemetry_data = generate_base_oob_telemetry_data_dict(