Skip to content
2 changes: 2 additions & 0 deletions doc/changes/changelog.md

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 10 additions & 0 deletions doc/changes/changes_0.2.1.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,10 @@
# 0.2.1 - 2026-09-30

## Summary

Fix of threading block at the process exit. Fixes `setup()` logic of reconfiguration (semi-private API).

## Bugs

- #10: Threading blocks at python exit
- Setup logic was wrong in one of tests.
1 change: 1 addition & 0 deletions exasol/telemetry/client/config.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ def was_enabled() -> bool:
def disable_config():
"""
Call disables telemetry entirely for all subsequent calls.
Telemetry still might be re-enabled by call to setup() with the opposite disable flag.
"""
conf = Config(enabled=False, endpoint=DEFAULT_ENDPOINT)
store(conf)
28 changes: 25 additions & 3 deletions exasol/telemetry/client/setup.py
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import atexit
import os
import typing as tt
from urllib.parse import urlparse
Expand Down Expand Up @@ -52,15 +53,32 @@ def setup_verbose_if_needed():
verbose.setup_logging()


def setup_exit_handlers():
"""
Trigger shutdown on exit of application.
Function have to be called only once - during the setup.
"""
atexit.register(shutdown)
Comment thread
ahsimb marked this conversation as resolved.


def drop_exit_handlers():
"""
Removes previously configured exit handlers
"""
atexit.unregister(shutdown)


def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None) -> bool:
"""
Telemetry client setup function.
Telemetry client setup function (should not be called from a client code).

Explicitly given arguments have the highest priority.
If they are not given, we check the environment variables (EXASOL_TELEMETRY_XXX),
if no environment value, we use defaults (DEFAULT_XXX).

If setup() was called before, we return the enabled status and do not reconfigure.
If setup() was called before, we reconfigure only when previous disable
option was different than a new one. On reconfiguration, we flush the buffers (if any).
In all other cases we return the enabled status without reconfiguration.

:param endpoint: Telemetry endpoint to send data. If not given,
default endpoint is used.
Expand All @@ -71,7 +89,9 @@ def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None)
False if it was disabled
"""
if config.was_setup():
Comment thread
ahsimb marked this conversation as resolved.
return config.was_enabled()
if disable is None or disable != config.was_enabled():
Comment thread
ahsimb marked this conversation as resolved.
Comment thread
ahsimb marked this conversation as resolved.
return config.was_enabled()
shutdown()
Comment thread
ahsimb marked this conversation as resolved.
val_endpoint = get_value(endpoint, config.ENV_ENDPOINT, config.DEFAULT_ENDPOINT)
Comment thread
ahsimb marked this conversation as resolved.

# Checking the presence of CI=true env variable
Expand Down Expand Up @@ -99,6 +119,7 @@ def setup(endpoint: tt.Optional[str] = None, disable: tt.Optional[bool] = None)
config.store(conf)
if enabled:
setup_verbose_if_needed()
setup_exit_handlers()
Comment thread
ahsimb marked this conversation as resolved.
worker.start_worker()
verbose.log("Setup is done, enabled=%s", conf.enabled)
return conf.enabled
Expand All @@ -115,6 +136,7 @@ def shutdown(flush_buffers: bool = True):
return
verbose.log("Shutdown")
worker.stop_worker(flush_buffers)
drop_exit_handlers()


def disable():
Expand Down
29 changes: 22 additions & 7 deletions exasol/telemetry/client/worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,9 @@
# requests' timeout value
SEND_TIMEOUT_SECONDS = 30

# how long to wait before thread exits
THREAD_EXIT_TIMEOUT_SECONDS = 1

# how long in seconds to wait before the first batch send
DATA_SEND_FIRST_INTERVAL_SECONDS = 0.5

Expand Down Expand Up @@ -257,25 +260,34 @@ def start_worker() -> bool:
if not config.was_enabled():
return False
_queue = queue.Queue(maxsize=MAX_QUEUE_CAPACITY)
_worker = threading.Thread(target=worker_proc, args=(_queue,))
_worker = threading.Thread(target=worker_proc, args=(_queue,), daemon=True)
_worker.start()
return True


def stop_worker(flush_buffers: bool):
"""
Gracefully stops the worker process.

In case of connectivity issues, flush of buffers might not happen.
:param flush_buffers: if True, we'll try to send the buffers (if any),
otherwise, we'll just shut down the worker process.
"""
global _worker, _queue

if _worker is None or _queue is None:
return
if flush_buffers:
_queue.put(WorkerMessage.make_send_buffers())
_queue.put(WorkerMessage.make_terminate())
_worker.join()
try:
if flush_buffers:
_queue.put_nowait(WorkerMessage.make_send_buffers())
_queue.put_nowait(WorkerMessage.make_terminate())
_worker.join(timeout=THREAD_EXIT_TIMEOUT_SECONDS)
except queue.Full:
# rare situation - if the thread is blocked on send and
# we have lots of messages in the queue, we can have no capacity
# in the queue. In such cases, we just don't stop the thread,
# which is fine as thread is daemon.
pass
_worker = None
_queue = None

Expand Down Expand Up @@ -307,5 +319,8 @@ def track(

global _queue
if _queue is not None:
if _queue.not_full:
_queue.put(WorkerMessage.make_track(product_name, product_version, feature))
try:
msg = WorkerMessage.make_track(product_name, product_version, feature)
_queue.put_nowait(msg)
except queue.Full:
pass
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
[project]
name = "exasol-telemetry-client"
version = "0.2.0"
version = "0.2.1"
description = ""
authors = [{ name = "Exasol AG", email = "opensource@exasol.com" }]
requires-python = ">=3.10,<4.0"
Expand Down
5 changes: 3 additions & 2 deletions test/unit/client/test_setup.py
Original file line number Diff line number Diff line change
Expand Up @@ -80,8 +80,9 @@ def test_setup_env_disabled(
assert not setup("http://endpoint")
assert config.was_setup()
assert not config.was_enabled()
# double-call to setup skips reconfiguration and returns the enable status
assert not setup(disable=False)
# double-call to setup skips reconfiguration if disable is the same
assert not setup(disable=True)
Comment thread
ahsimb marked this conversation as resolved.
assert setup(disable=False)


def test_setup_env_enabled(
Expand Down
31 changes: 31 additions & 0 deletions test/unit/client/test_worker.py
Original file line number Diff line number Diff line change
Expand Up @@ -97,6 +97,24 @@ def test_worker_proc_sent_quick(
mock_post.assert_called_once()


# Make sure that buffers are sent after reconfiguration
@mock.patch("requests.post", return_value=mock.MagicMock(status_code=200))
def test_worker_proc_sent_after_disable(
mock_post: mock.MagicMock,
telemetry_reset,
telemetry_unset_ci,
telemetry_unset_disable,
):
# enable and send track feature
assert setup(disable=False)
track("product", "ver", "test")
# after reconfiguration we should be disabled and flushed the buffers
assert not setup(disable=True)
mock_post.assert_called_once()
assert config.was_setup()
assert not config.was_enabled()


# Make sure that features are not sent if not enabled
@mock.patch("requests.post")
def test_worker_proc_not_sent_when_disabled(
Expand All @@ -114,6 +132,19 @@ def test_worker_proc_not_sent_when_disabled(
mock_post.assert_not_called()


@mock.patch("requests.post", side_effect=lambda *w, **kw: time.sleep(1000))
def test_shutdown_exits_on_blocked_worker(
mock_post: mock.MagicMock,
telemetry_reset,
telemetry_unset_ci,
telemetry_unset_disable,
):
assert setup(disable=False)
track("test", "0.1", "test-feature")
shutdown(flush_buffers=True)
mock_post.assert_called_once()


@mock.patch("exasol.telemetry.client.worker.send_features")
def test_worker_proc_no_send(mock_send_features: mock.MagicMock):
msg_queue = queue.Queue()
Expand Down