diff --git a/doc/changes/changelog.md b/doc/changes/changelog.md index fd4a51d..e1874a0 100644 --- a/doc/changes/changelog.md +++ b/doc/changes/changelog.md @@ -1,6 +1,7 @@ # Changes * [unreleased](unreleased.md) +* [0.2.0](changes_0.2.0.md) * [0.1.6](changes_0.1.6.md) ```{toctree} @@ -8,5 +9,6 @@ hidden: --- unreleased +changes_0.2.0 changes_0.1.6 ``` diff --git a/doc/changes/changes_0.2.0.md b/doc/changes/changes_0.2.0.md new file mode 100644 index 0000000..fed1d9a --- /dev/null +++ b/doc/changes/changes_0.2.0.md @@ -0,0 +1,10 @@ +# 0.2.0 - 2026-09-15 + +## Summary + +Major refactoring of client API. Support for global disable of telemetry and verbose mode. + +## Refactoring + +- #4: Support v0.2 of the protocol + API refactor +- #8: Provide py.typed marker diff --git a/doc/client-python.rst b/doc/client-python.rst index 604f240..cd9039c 100644 --- a/doc/client-python.rst +++ b/doc/client-python.rst @@ -10,54 +10,36 @@ public repo yet) ``pip install exasol-telemetry-client`` Usage ----- -Once installed, package ``exasel.telemetry.client`` provides three -methods and one exception: - ``setup()``: configures the library, has to -be called once in the beginning - ``track(feature_name)``: tracks -feature as used (string) - ``shutdown()``: should be called at the end -of the program. If not called, some tracked features could be lost. - -``TelemetryError``: exception could be thrown during ``setup()`` call, -if environment variables are wrong. +Once installed, package ``exasol.telemetry.client`` provides three +methods: -Function ``was_setup()`` could be used to check the ``setup()`` was called -before (possibly in another library). +- ``track(product_name, product_version, feature_name)``: tracks feature as used (string) +- ``shutdown()``: should be called at the end of the program. If not called, some tracked features could be lost. +- ``disable()``: disables telemetry entirely for the whole process till the termination of the process. It is useful for cases when some software wants to disable telemetry even when some of its libraries are using it. + +Explicit initialization of the library is not needed --- it will be set up on the first ``track()`` call. Example of minimalistic program: .. code:: python - import logging from exasol.telemetry.client import * if __name__ == "__main__": try: - try: - if not was_setup(): - setup() - except TelemetryError as e: - logging.warning("Telemetry disabled due to error: %s", str(e)) - - track("feature1") - track("feature2") + track("hello-world", "0.1", "started") + core_of_the_program() finally: shutdown() -Exceptions ----------- - -Exception ``TelemetryError`` could be thrown from ``setup()`` and -``shutdown()`` in case of errors. Call of ``track()`` never raises exceptions, in case of errors -tracked feature is ignored. - Environment variables --------------------- To change the telemetry configuration, you can set the following -environment variables. Those values also could be changed via -``setup()`` arguments. +environment variables. - ``EXASOL_TELEMETRY_DISABLE`` - any value disables the telemetry data collection and sending - ``EXASOL_TELEMETRY_ENDPOINT`` - redefines telemetry endpoint url. - -In addition, if environment variable ``CI=true`` (which is the case during Github CI workflows run) -the telemetry is disabled unless explicitly enabled with ``setup()`` arguments. +- ``EXASOL_TELEMETRY_VERBOSE`` -- enables logging messages from the library. Could be used to make sure integration was done properly. +- ``CI=true`` -- disables telemetry to prevent tracking during CI. diff --git a/doc/index.rst b/doc/index.rst index 7edb639..37ee185 100644 --- a/doc/index.rst +++ b/doc/index.rst @@ -30,6 +30,7 @@ Documentation of telemetry user_guide client-python + protocol developer_guide api faq diff --git a/doc/protocol.rst b/doc/protocol.rst new file mode 100644 index 0000000..fd5efdc --- /dev/null +++ b/doc/protocol.rst @@ -0,0 +1,52 @@ +Telemetry Protocol Specification +================================ + +Exasol telemetry uses simplistic protocol sending events happened in the the software. +Every event has a timestamp attached to be used for server-side analytics. +All the data on the server is immediately aggregated and anonymized and no personal information +is transferred or stored. + +The data is transferred in json format and at the moment there are two versions of the protocol. + +Version 0.1 +----------- +.. code:: json + + { + "version": "0.1", + "timestamp": 1787036195, + "features": { + "EMCP.started": [1787036195] + } + } + +Transferred data has the following fields: + +- ``version``: string specifying the protocol version +- ``timestamp``: UTC timestamp of the transmission attempt +- ``features``: dictionary with pairs ``feature-name`` and vector of timestamps when the event happened. + +Recording of both event timestamp and transmission timestamp allows to check the clock discrepancies on the client +side and filter out outliers. + +Version 0.2 +----------- + +This is an extension of version 0.1, sample data is below :: + + { + "version": "0.2", + "category": "EMCP", + "productVersion": "0.22", + "timestamp": 1787036195, + "features": { + "started": [1787036195] + } + } + +In this version we have two new top-level fields added: + +- ``category``: name of the product +- ``productVersion``: version of the product + +The name of the product is no longer prepended to the features, which makes the data more compact. \ No newline at end of file diff --git a/exasol/telemetry/client/__init__.py b/exasol/telemetry/client/__init__.py index d3ecff8..695439b 100644 --- a/exasol/telemetry/client/__init__.py +++ b/exasol/telemetry/client/__init__.py @@ -2,19 +2,15 @@ Telemetry client library for python. Public API is three methods: -- setup: initializes the library based on explicit arguments or environment variables -- shutdown: cleans up the resources and sends the data still in buffers - track: remembers the feature name in the buffer (will be sent in background) - -On error we raise exception TelemetryError with error message. +- disable: disables the telemetry +- shutdown: cleans up the resources and sends the data still in buffers """ -from exasol.telemetry.client.config import was_setup -from exasol.telemetry.client.errors import TelemetryError from exasol.telemetry.client.setup import ( - setup, + disable, shutdown, ) from exasol.telemetry.client.worker import track -__all__ = ["was_setup", "setup", "shutdown", "track", "TelemetryError"] +__all__ = ["track", "disable", "shutdown"] diff --git a/exasol/telemetry/client/config.py b/exasol/telemetry/client/config.py index 76327f2..dc18595 100644 --- a/exasol/telemetry/client/config.py +++ b/exasol/telemetry/client/config.py @@ -6,6 +6,8 @@ ENV_DISABLE = "EXASOL_TELEMETRY_DISABLE" # Endpoint (has to be valid http/https URL) ENV_ENDPOINT = "EXASOL_TELEMETRY_ENDPOINT" +# Enable console logging of telemetry events +ENV_VERBOSE = "EXASOL_TELEMETRY_VERBOSE" # GitHub CI sets this to true during execution ENV_CI = "CI" @@ -49,3 +51,11 @@ def was_enabled() -> bool: """ conf = get() return conf is not None and conf.enabled + + +def disable_config(): + """ + Call disables telemetry entirely for all subsequent calls. + """ + conf = Config(enabled=False, endpoint=DEFAULT_ENDPOINT) + store(conf) diff --git a/exasol/telemetry/client/errors.py b/exasol/telemetry/client/errors.py deleted file mode 100644 index 91c1bb9..0000000 --- a/exasol/telemetry/client/errors.py +++ /dev/null @@ -1,7 +0,0 @@ -class TelemetryError(Exception): - """ - Telemetry exception, thrown from telemetry client methods. - """ - - def __init__(self, message): - super().__init__(message) diff --git a/exasol/telemetry/client/protocol.py b/exasol/telemetry/client/protocol.py index ba94fc0..f91c845 100644 --- a/exasol/telemetry/client/protocol.py +++ b/exasol/telemetry/client/protocol.py @@ -7,10 +7,12 @@ import typing as tt # Version of the protocol -VERSION = "0.1" +VERSION = "0.2" -# Name of the feature +# type aliases Feature = str +ProductName = str +ProductVersion = str # Timestamp of the measurement Timestamp = tt.Union[int, float] @@ -28,6 +30,12 @@ class Message: # Version of the protocol this message corresponds to. version: str + # Name of the product ('category' in v0.2 protocol) + product_name: ProductName + + # Version of the product + product_version: ProductVersion + # Current unit timestamp when the message was created # (used by the server to get the age of the individual reports) timestamp: Timestamp @@ -38,20 +46,31 @@ class Message: def to_json(self) -> dict: return { "version": self.version, + "category": self.product_name, + "productVersion": self.product_version, "timestamp": self.timestamp, "features": self.features, } @classmethod - def from_features(cls, features: Features) -> "Message": + def from_features( + cls, + product_name: ProductName, + product_version: ProductVersion, + features: Features, + ) -> "Message": """ Construct the message object from collected features. We're not deep copy of features, just store the reference of it. :param features: collection of features + :param product_name: name of the product + :param product_version: version of the product :return: Message created """ return Message( version=VERSION, + product_name=product_name, + product_version=product_version, timestamp=get_current_ts(), features=features, ) diff --git a/exasol/telemetry/client/py.typed b/exasol/telemetry/client/py.typed new file mode 100644 index 0000000..e69de29 diff --git a/exasol/telemetry/client/setup.py b/exasol/telemetry/client/setup.py index 3a34de7..e5ee75d 100644 --- a/exasol/telemetry/client/setup.py +++ b/exasol/telemetry/client/setup.py @@ -4,9 +4,9 @@ from exasol.telemetry.client import ( config, + verbose, worker, ) -from exasol.telemetry.client.errors import TelemetryError def get_value( @@ -42,10 +42,19 @@ def is_valid_endpoint_url(url: str) -> bool: return res.scheme in ("http", "https") and len(res.netloc) > 0 +def setup_verbose_if_needed(): + """ + Function enables verbose mode for telemetry prefix if env variable is set. + """ + # if env variable is not set, do nothing + if get_value(None, config.ENV_VERBOSE, None) is None: + return + verbose.setup_logging() + + def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None) -> bool: """ - Telemetry client setup function. Has to be called before any other - calls to the client. + Telemetry client setup function. Explicitly given arguments have the highest priority. If they are not given, we check the environment variables (EXASOL_TELEMETRY_XXX), @@ -58,7 +67,6 @@ def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None) :param disable: If True, disable telemetry communication and data accumulation. - :raises TelemetryError: if error has happened :returns True if telemetry is active according to configuration, False if it was disabled """ @@ -85,11 +93,14 @@ def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None) enabled = not val_disabled if enabled and not is_valid_endpoint_url(val_endpoint): - raise TelemetryError("Endpoint is invalid: " + val_endpoint) + enabled = False conf = config.Config(endpoint=val_endpoint, enabled=enabled) config.store(conf) - worker.start_worker() + if enabled: + setup_verbose_if_needed() + worker.start_worker() + verbose.log("Setup is done, enabled=%s", conf.enabled) return conf.enabled @@ -101,6 +112,14 @@ def shutdown(flush_buffers: bool = True): so some values might be lost. """ if not config.was_setup(): - raise TelemetryError("Telemetry was not initialized") + return + verbose.log("Shutdown") worker.stop_worker(flush_buffers) - config.store(None) + + +def disable(): + """ + Shuts down workers and disables the telemetry globally. + """ + config.disable_config() + worker.stop_worker(flush_buffers=False) diff --git a/exasol/telemetry/client/verbose.py b/exasol/telemetry/client/verbose.py new file mode 100644 index 0000000..1f9d906 --- /dev/null +++ b/exasol/telemetry/client/verbose.py @@ -0,0 +1,25 @@ +import logging +import typing as tt + +LOGGER = "exasol.telemetry.client" +LEVEL = logging.DEBUG + +logger: tt.Optional[logging.Logger] = None + + +def setup_logging(): + """ + Enable logging for our package. + """ + global logger + # prevent double-initialization + if logger is not None: + return + logger = logging.getLogger(LOGGER) + logger.setLevel(LEVEL) + + +def log(msg: str, *args, **kwargs): + global logger + if logger is not None: + logger.log(LEVEL, msg, *args, **kwargs) diff --git a/exasol/telemetry/client/worker.py b/exasol/telemetry/client/worker.py index cb11f4c..5277689 100644 --- a/exasol/telemetry/client/worker.py +++ b/exasol/telemetry/client/worker.py @@ -2,7 +2,6 @@ import dataclasses import enum import json -import logging import queue import threading import typing as tt @@ -12,6 +11,7 @@ from exasol.telemetry.client import ( config, protocol, + verbose, ) MAX_QUEUE_CAPACITY = 10 @@ -28,9 +28,9 @@ # for how long we keep features in buffers before removing them MAX_DATA_KEEP_SECONDS = 60 * 60 -log = logging.getLogger("worker") _worker: tt.Optional[threading.Thread] = None _queue: tt.Optional[queue.Queue] = None +_setup_lock: threading.Lock = threading.Lock() class WorkerCommand(enum.Enum): @@ -42,53 +42,160 @@ class WorkerCommand(enum.Enum): @dataclasses.dataclass(frozen=True) class WorkerMessage: command: WorkerCommand - feature: tt.Optional[protocol.Feature] - timestamp: tt.Optional[protocol.Timestamp] + product_name: tt.Optional[protocol.ProductName] = None + product_version: tt.Optional[protocol.ProductVersion] = None + feature: tt.Optional[protocol.Feature] = None + timestamp: tt.Optional[protocol.Timestamp] = None @classmethod - def make_track(cls, feature: protocol.Feature) -> "WorkerMessage": + def make_track( + cls, + product_name: protocol.ProductName, + product_version: protocol.ProductVersion, + feature: protocol.Feature, + ) -> "WorkerMessage": return WorkerMessage( command=WorkerCommand.Track, + product_name=product_name, + product_version=product_version, feature=feature, timestamp=protocol.get_current_ts(), ) @classmethod def make_send_buffers(cls) -> "WorkerMessage": - return WorkerMessage( - command=WorkerCommand.SendBuffers, feature=None, timestamp=None - ) + return WorkerMessage(command=WorkerCommand.SendBuffers) @classmethod def make_terminate(cls) -> "WorkerMessage": - return WorkerMessage( - command=WorkerCommand.Terminate, feature=None, timestamp=None + return WorkerMessage(command=WorkerCommand.Terminate) + + +class WorkerDeadlineQueue: + """ + Queue with deadline - moment in the future when we + want to stop waiting for a message to arrive. + """ + + def __init__(self, msg_queue: queue.Queue): + self._queue = msg_queue + self._deadline_ts: tt.Optional[protocol.Timestamp] = None + + def set_deadline(self, seconds: float): + self._deadline_ts = protocol.get_current_ts() + seconds + + def deadline_expired(self) -> bool: + return ( + self._deadline_ts is None or protocol.get_current_ts() > self._deadline_ts ) + def seconds_to_deadline(self, now: tt.Optional[protocol.Timestamp] = None) -> float: + """ + Get amount of seconds until deadline. If expired, return 0 + :param now: optional current time (used for testing) + :return: count of seconds + """ + if self._deadline_ts is None: + return 0.0 + if now is None: + now = protocol.get_current_ts() + dt = self._deadline_ts - now + return max(dt, 0.0) + + def get_msg(self) -> tt.Optional[WorkerMessage]: + try: + msg = self._queue.get(timeout=self.seconds_to_deadline()) + return msg if isinstance(msg, WorkerMessage) else None + except queue.Empty: + return None + + +DataPoolKey = tt.Tuple[protocol.ProductName, protocol.ProductVersion] -def send_features(features: protocol.Features) -> bool: + +class DataPool: + def __init__(self): + self._pool: tt.Dict[DataPoolKey, protocol.Features] = {} + + def is_empty(self) -> bool: + return not bool(self._pool) + + def send(self) -> bool: + to_clear: tt.List[DataPoolKey] = [] + try: + for key, features in self._pool.items(): + product, version = key + if send_features(product, version, features): + to_clear.append(key) + else: + # stop on first error + break + finally: + for key in to_clear: + self._pool.pop(key) + return self.is_empty() + + def clear_expired(self, now: protocol.Timestamp): + to_clear: tt.List[DataPoolKey] = [] + for key, features in self._pool.items(): + clear_expired_features(features, now) + if not features: + to_clear.append(key) + for key in to_clear: + self._pool.pop(key) + + def append( + self, + product_name: tt.Optional[protocol.ProductName], + product_version: tt.Optional[protocol.ProductVersion], + feature: tt.Optional[protocol.Feature], + timestamp: tt.Optional[protocol.Timestamp], + ): + # should never happen, but to make linter happy :shrug + if ( + product_name is None + or product_version is None + or feature is None + or timestamp is None + ): + return + key = (product_name, product_version) + features = self._pool.get(key) + if features is None: + features = collections.defaultdict(list) + self._pool[key] = features + features[feature].append(timestamp) + + +def send_features( + product_name: protocol.ProductName, + product_version: protocol.ProductVersion, + features: protocol.Features, +) -> bool: """ Internal method to send the accumulated data to endpoint. + :param product_name: name of the product + :param product_version: version of the product :param features: data to be sent :return: True if data was sent successfully, - False if something happened and we need to keep data for some time. + False if something happened, and we need to keep data for some time. """ if not features: return True conf = config.get() if conf is None: return True - message = protocol.Message.from_features(features) + message = protocol.Message.from_features(product_name, product_version, features) try: url = conf.endpoint data = json.dumps(message.to_json()) res = requests.post(url, data, timeout=SEND_TIMEOUT_SECONDS) if res.status_code != 200: - log.debug("Feature send error: %s", str(res)) + verbose.log("Features send error: %s", str(res)) return False return True except requests.exceptions.RequestException as e: - log.debug("Features send error: %s", str(e)) + verbose.log("Send exception: %s", str(e)) return False @@ -109,52 +216,30 @@ def clear_expired_features(features: protocol.Features, now: protocol.Timestamp) features.pop(key) -def get_msg_timeout( - next_sent_ts: protocol.Timestamp, now: tt.Optional[protocol.Timestamp] = None -) -> protocol.Timestamp: - """ - Get amount of seconds to wait for new message. - :param next_sent_ts: when next data push is planned - :param now: optional argument to redefine the current timestamp. Used for testing - :return: Amount of seconds to sleep before the next send attempt. - """ - if now is None: - now = protocol.get_current_ts() - dt = next_sent_ts - now - # if interval has expired, return 0 - return max(dt, 0) - - def worker_proc(msg_queue: queue.Queue): """ Worker procedure - consumes the queue, periodically sends the accumulated data. :param msg_queue: queue to consume """ - data: protocol.Features = collections.defaultdict(list) - next_sent_ts = protocol.get_current_ts() + DATA_SEND_FIRST_INTERVAL_SECONDS + data_pool = DataPool() + deadline_queue = WorkerDeadlineQueue(msg_queue) + deadline_queue.set_deadline(DATA_SEND_FIRST_INTERVAL_SECONDS) while True: - try: - msg = msg_queue.get(timeout=get_msg_timeout(next_sent_ts)) - except queue.Empty: - msg = None - now = protocol.get_current_ts() - - # if interval has expired, or we have an explicit send request, send the data - if now > next_sent_ts or ( - msg is not None and msg.command == WorkerCommand.SendBuffers - ): - if send_features(data): - data.clear() - else: - clear_expired_features(data, now) - next_sent_ts = now + DATA_SEND_INTERVAL_SECONDS - - if not isinstance(msg, WorkerMessage): + msg = deadline_queue.get_msg() + verbose.log("Message: %s", str(msg)) + if msg is None or msg.command == WorkerCommand.SendBuffers: + # deadline has expired or we've asked to flush buffers + if not data_pool.send(): + data_pool.clear_expired(protocol.get_current_ts()) + if deadline_queue.deadline_expired(): + deadline_queue.set_deadline(DATA_SEND_INTERVAL_SECONDS) + if msg is None: continue if msg.command == WorkerCommand.Track: - if msg.feature is not None and msg.timestamp is not None: - data[msg.feature].append(msg.timestamp) + data_pool.append( + msg.product_name, msg.product_version, msg.feature, msg.timestamp + ) elif msg.command == WorkerCommand.Terminate: return @@ -169,6 +254,8 @@ def start_worker() -> bool: if _worker is not None: return False + if not config.was_enabled(): + return False _queue = queue.Queue(maxsize=MAX_QUEUE_CAPACITY) _worker = threading.Thread(target=worker_proc, args=(_queue,)) _worker.start() @@ -193,15 +280,32 @@ def stop_worker(flush_buffers: bool): _queue = None -def track(feature: protocol.Feature): +def _do_setup(): + from .setup import setup + + global _setup_lock + + with _setup_lock: + setup() + + +def track( + product_name: protocol.ProductName, + product_version: protocol.ProductVersion, + feature: protocol.Feature, +): """ - Track feature usage. Library has to be initialized with `setup()` call. + Track feature usage. + :param product_name: product name + :param product_version: product version :param feature: string feature to track. """ + if not config.was_setup(): + _do_setup() if not config.was_enabled(): return global _queue if _queue is not None: if _queue.not_full: - _queue.put(WorkerMessage.make_track(feature)) + _queue.put(WorkerMessage.make_track(product_name, product_version, feature)) diff --git a/test/integration/test_client.py b/test/integration/test_client.py index 1ecc9c0..2821ce2 100644 --- a/test/integration/test_client.py +++ b/test/integration/test_client.py @@ -1,15 +1,11 @@ import pytest from exasol.telemetry.client import * -from exasol.telemetry.client import worker - -ENDPOINT = "" @pytest.mark.skip() def test_client(): - assert setup(ENDPOINT, disable=False) try: - assert worker.send_features({"test_feat": [1]}) + track("test", "0.1", "test-feature") finally: shutdown(flush_buffers=True) diff --git a/test/unit/client/conftest.py b/test/unit/client/conftest.py index 9923c22..6fae9bb 100644 --- a/test/unit/client/conftest.py +++ b/test/unit/client/conftest.py @@ -1,8 +1,8 @@ import pytest from exasol.telemetry.client import ( - TelemetryError, config, + verbose, ) from exasol.telemetry.client.setup import shutdown @@ -10,19 +10,32 @@ @pytest.fixture def telemetry_reset(): """ - Call `shutdown()` after the test. + Resets the telemetry into initial state """ yield - try: - shutdown() - except TelemetryError: - pass + shutdown() + config.store(None) -@pytest.fixture() +@pytest.fixture def telemetry_unset_ci(monkeypatch): """ Temporary remove CI env variable if present """ monkeypatch.delenv(config.ENV_CI, raising=False) + + +@pytest.fixture +def telemetry_unset_disable(monkeypatch): + """ + Temporary remove EXASOL_TELEMETRY_DISABLE env variable if present + """ + monkeypatch.delenv(config.ENV_DISABLE, raising=False) + + +@pytest.fixture +def telemetry_verbose(monkeypatch): + monkeypatch.setenv(config.ENV_VERBOSE, "t") yield + verbose.logger = None + monkeypatch.delenv(config.ENV_VERBOSE) diff --git a/test/unit/client/test_setup.py b/test/unit/client/test_setup.py index 9518a10..3085c4c 100644 --- a/test/unit/client/test_setup.py +++ b/test/unit/client/test_setup.py @@ -1,18 +1,19 @@ import pytest -from exasol.telemetry.client import * -from exasol.telemetry.client import config +from exasol.telemetry.client import ( + config, + verbose, +) +from exasol.telemetry.client.config import was_setup from exasol.telemetry.client.setup import ( get_value, is_valid_endpoint_url, + setup, + setup_verbose_if_needed, + shutdown, ) -def test_error_when_not_initialized(): - with pytest.raises(TelemetryError, match="not initialized"): - shutdown() - - @pytest.mark.parametrize( "url, expected", [ @@ -64,13 +65,17 @@ def test_setup_explicit_disabled(telemetry_reset): assert config.get().endpoint.startswith("https") -def test_setup_wrong_endpoint(telemetry_reset, telemetry_unset_ci): - with pytest.raises(TelemetryError, match="Endpoint is invalid"): - setup("ftp://test.com") - assert not config.was_setup() +def test_setup_wrong_endpoint( + telemetry_reset, telemetry_unset_ci, telemetry_unset_disable +): + assert not setup("ftp://test.com") + assert config.was_setup() + assert not config.was_enabled() -def test_setup_env_disabled(monkeypatch, telemetry_reset): +def test_setup_env_disabled( + monkeypatch, telemetry_reset, telemetry_unset_ci, telemetry_unset_disable +): monkeypatch.setenv(config.ENV_DISABLE, "1") assert not setup("http://endpoint") assert config.was_setup() @@ -79,14 +84,16 @@ def test_setup_env_disabled(monkeypatch, telemetry_reset): assert not setup(disable=False) -def test_setup_env_enabled(monkeypatch, telemetry_reset, telemetry_unset_ci): +def test_setup_env_enabled( + monkeypatch, telemetry_reset, telemetry_unset_ci, telemetry_unset_disable +): monkeypatch.setenv(config.ENV_ENDPOINT, "http://test") assert setup() assert config.was_enabled() assert config.get().endpoint == "http://test" -def test_setup_defaults(telemetry_reset, telemetry_unset_ci): +def test_setup_defaults(telemetry_reset, telemetry_unset_ci, telemetry_unset_disable): assert setup() == (not config.DEFAULT_DISABLED) assert config.get().endpoint == config.DEFAULT_ENDPOINT @@ -97,11 +104,22 @@ def test_setup_ci_true(monkeypatch, telemetry_reset): assert not config.was_enabled() -def test_setup_ci_true_explicit(monkeypatch, telemetry_reset): +def test_setup_ci_true_explicit(monkeypatch, telemetry_reset, telemetry_unset_disable): monkeypatch.setenv(config.ENV_CI, "true") assert setup(disable=False) -def test_setup_ci_false(monkeypatch, telemetry_reset): +def test_setup_ci_false(monkeypatch, telemetry_reset, telemetry_unset_disable): monkeypatch.setenv(config.ENV_CI, "t") assert setup() + + +def test_shutdown_not_setup(telemetry_reset): + shutdown() + assert not was_setup() + + +def test_verbose_mode(monkeypatch, telemetry_reset, telemetry_verbose): + assert verbose.logger is None + setup_verbose_if_needed() + assert verbose.logger is not None diff --git a/test/unit/client/test_verbose.py b/test/unit/client/test_verbose.py new file mode 100644 index 0000000..b3cd456 --- /dev/null +++ b/test/unit/client/test_verbose.py @@ -0,0 +1,24 @@ +import logging + +from exasol.telemetry.client import verbose + + +def test_no_show_unconfigured(caplog, telemetry_reset): + logging.log(verbose.LEVEL, "Test") + assert "Test" not in caplog.text + caplog.clear() + + verbose.log("Test") + assert "Test" not in caplog.text + + +def test_show_configured(caplog, telemetry_reset): + verbose.setup_logging() + verbose.log("Test") + assert "Test" in caplog.text + + # second setup changes nothing + verbose.setup_logging() + caplog.clear() + verbose.log("Test2") + assert "Test2" in caplog.text diff --git a/test/unit/client/test_worker.py b/test/unit/client/test_worker.py index 675b875..3ad5d77 100644 --- a/test/unit/client/test_worker.py +++ b/test/unit/client/test_worker.py @@ -6,8 +6,10 @@ from exasol.telemetry.client import ( config, protocol, + verbose, worker, ) +from exasol.telemetry.client.setup import setup def test_stop_worker_doing_nothing_without_worker(): @@ -15,10 +17,10 @@ def test_stop_worker_doing_nothing_without_worker(): def test_track_not_init(telemetry_reset): - worker.track("feature") + worker.track("test-product", "0.1", "feature") assert not setup(disable=True) - worker.track("feature") + worker.track("test-product", "0.1", "feature") def test_clear_expired_features(): @@ -41,49 +43,55 @@ def test_clear_expired_features(): @mock.patch("requests.post") -def test_track(post_mock: mock.MagicMock, telemetry_reset, telemetry_unset_ci): +def test_track( + post_mock: mock.MagicMock, + telemetry_reset, + telemetry_unset_ci, + telemetry_unset_disable, +): post_mock.return_value = mock.MagicMock(status_code=200) - assert setup() - worker.track("feature1") - worker.track("feature2") + worker.track("test", "0.1", "feature1") + worker.track("test", "0.1", "feature2") shutdown(flush_buffers=True) post_mock.assert_called_once() -def test_get_msg_timeout(): - now = protocol.get_current_ts() - default_interval = worker.DATA_SEND_INTERVAL_SECONDS - assert worker.get_msg_timeout(now, now) == 0 - assert worker.get_msg_timeout(now + default_interval, now) == default_interval - assert ( - worker.get_msg_timeout(now + default_interval - 10, now) - == default_interval - 10 - ) - assert worker.get_msg_timeout(now - default_interval * 2, now) == 0 +def test_deadline_queue_deadline(): + q = worker.WorkerDeadlineQueue(queue.Queue()) + assert q.deadline_expired() + assert q.seconds_to_deadline() == 0.0 + q.set_deadline(0.2) + assert not q.deadline_expired() + assert q.seconds_to_deadline() > 0.0 + time.sleep(1) + assert q.deadline_expired() def test_send_features_not_conf(): assert config.get() is None - assert worker.send_features({}) - assert worker.send_features({"f": [1]}) + assert worker.send_features("prod", "ver", {}) + assert worker.send_features("prod", "ver", {"f": [1]}) def test_send_features_wrong_endpoint(telemetry_reset, caplog): - caplog.set_level("DEBUG") + verbose.setup_logging() + caplog.set_level(verbose.LEVEL) assert setup(endpoint="http://non-existent-domain.weird", disable=False) - assert not worker.send_features({"f": [1]}) - assert "Features send error" in caplog.text + assert not worker.send_features("prod", "ver", {"f": [1]}) + assert "Send exception" in caplog.text assert "Name or service not known" in caplog.text # Make sure that first buffer is sent quickly after the initialization @mock.patch("requests.post", return_value=mock.MagicMock(status_code=200)) def test_worker_proc_sent_quick( - mock_post: mock.MagicMock, telemetry_reset, telemetry_unset_ci + mock_post: mock.MagicMock, + telemetry_reset, + telemetry_unset_ci, + telemetry_unset_disable, ): - assert setup() - track("test") + track("product", "ver", "test") time.sleep(worker.DATA_SEND_FIRST_INTERVAL_SECONDS * 2) shutdown(flush_buffers=False) mock_post.assert_called_once() @@ -92,12 +100,15 @@ def test_worker_proc_sent_quick( # Make sure that features are not sent if not enabled @mock.patch("requests.post") def test_worker_proc_not_sent_when_disabled( - mock_post: mock.MagicMock, telemetry_reset, telemetry_unset_ci + mock_post: mock.MagicMock, + telemetry_reset, + telemetry_unset_ci, + telemetry_unset_disable, ): assert not setup(disable=True) assert config.was_setup() assert not config.was_enabled() - track("test") + track("test", "0.1", "test-feature") shutdown(flush_buffers=True) mock_post.assert_not_called() @@ -106,7 +117,7 @@ def test_worker_proc_not_sent_when_disabled( @mock.patch("exasol.telemetry.client.worker.send_features") def test_worker_proc_no_send(mock_send_features: mock.MagicMock): msg_queue = queue.Queue() - msg_queue.put(worker.WorkerMessage.make_track("f1")) + msg_queue.put(worker.WorkerMessage.make_track("prod", "ver", "f1")) msg_queue.put(worker.WorkerMessage.make_terminate()) worker.worker_proc(msg_queue) @@ -116,7 +127,7 @@ def test_worker_proc_no_send(mock_send_features: mock.MagicMock): @mock.patch("exasol.telemetry.client.worker.send_features", return_value=True) def test_worker_proc_send_success(mock_send_features: mock.MagicMock): msg_queue = queue.Queue() - msg_queue.put(worker.WorkerMessage.make_track("f1")) + msg_queue.put(worker.WorkerMessage.make_track("prod", "ver", "f1")) msg_queue.put(worker.WorkerMessage.make_send_buffers()) msg_queue.put(worker.WorkerMessage.make_terminate()) @@ -132,9 +143,8 @@ def test_worker_proc_send_fail( mock_send_features: mock.MagicMock, mock_clear_expired_features: mock.MagicMock ): msg_queue = queue.Queue() - msg_queue.put(worker.WorkerMessage.make_track("f1")) + msg_queue.put(worker.WorkerMessage.make_track("prod", "ver", "f1")) msg_queue.put(None) - msg_queue.put(worker.WorkerMessage.make_send_buffers()) msg_queue.put(worker.WorkerMessage.make_terminate()) worker.worker_proc(msg_queue)