diff --git a/doc/changes/changelog.md b/doc/changes/changelog.md index e1874a0..1b26859 100644 --- a/doc/changes/changelog.md +++ b/doc/changes/changelog.md @@ -1,6 +1,7 @@ # Changes * [unreleased](unreleased.md) +* [0.2.1](changes_0.2.1.md) * [0.2.0](changes_0.2.0.md) * [0.1.6](changes_0.1.6.md) @@ -9,6 +10,7 @@ hidden: --- unreleased +changes_0.2.1 changes_0.2.0 changes_0.1.6 ``` diff --git a/doc/changes/changes_0.2.1.md b/doc/changes/changes_0.2.1.md new file mode 100644 index 0000000..c25c25f --- /dev/null +++ b/doc/changes/changes_0.2.1.md @@ -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. diff --git a/exasol/telemetry/client/config.py b/exasol/telemetry/client/config.py index dc18595..0f3598f 100644 --- a/exasol/telemetry/client/config.py +++ b/exasol/telemetry/client/config.py @@ -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) diff --git a/exasol/telemetry/client/setup.py b/exasol/telemetry/client/setup.py index e5ee75d..0f757e8 100644 --- a/exasol/telemetry/client/setup.py +++ b/exasol/telemetry/client/setup.py @@ -1,3 +1,4 @@ +import atexit import os import typing as tt from urllib.parse import urlparse @@ -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) + + +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. @@ -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(): - return config.was_enabled() + if disable is None or disable != config.was_enabled(): + return config.was_enabled() + shutdown() val_endpoint = get_value(endpoint, config.ENV_ENDPOINT, config.DEFAULT_ENDPOINT) # Checking the presence of CI=true env variable @@ -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() worker.start_worker() verbose.log("Setup is done, enabled=%s", conf.enabled) return conf.enabled @@ -115,6 +136,7 @@ def shutdown(flush_buffers: bool = True): return verbose.log("Shutdown") worker.stop_worker(flush_buffers) + drop_exit_handlers() def disable(): diff --git a/exasol/telemetry/client/worker.py b/exasol/telemetry/client/worker.py index 5277689..e39f67e 100644 --- a/exasol/telemetry/client/worker.py +++ b/exasol/telemetry/client/worker.py @@ -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 @@ -257,7 +260,7 @@ 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 @@ -265,6 +268,8 @@ def start_worker() -> bool: 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. """ @@ -272,10 +277,17 @@ def stop_worker(flush_buffers: bool): 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 @@ -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 diff --git a/pyproject.toml b/pyproject.toml index 621daf2..21fef52 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -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" diff --git a/test/unit/client/test_setup.py b/test/unit/client/test_setup.py index 3085c4c..9b3e4fe 100644 --- a/test/unit/client/test_setup.py +++ b/test/unit/client/test_setup.py @@ -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) + assert setup(disable=False) def test_setup_env_enabled( diff --git a/test/unit/client/test_worker.py b/test/unit/client/test_worker.py index 3ad5d77..abd34b0 100644 --- a/test/unit/client/test_worker.py +++ b/test/unit/client/test_worker.py @@ -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( @@ -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()