Conversation
ahsimb
left a comment
There was a problem hiding this comment.
The core fix for #10 looks right: making the worker a daemon thread means Python no longer waits for it at shutdown, and the atexit handler still sends the remaining data and stops the worker cleanly. Previously the interpreter joined the non-daemon worker before running atexit handlers, so the handler that would have stopped the worker never ran.
My main concern is the new re-configuration logic in setup() (second call with a different disable value). It isn't needed to fix #10, and it brings in several problems (see inline comments). Since setup() isn't part of the public API (__init__ exports only track, disable and shutdown) and the internal caller _do_setup() passes no arguments, the simplest option might be to drop that change. If it's needed, please fix the inline points and add tests for it.
One more thing, not in this diff but related to blocking: in worker.track() (worker.py ~line 310), if _queue.not_full: tests a threading.Condition object, which is always truthy, so the check never skips. If the worker is stuck in a slow requests.post and more than MAX_QUEUE_CAPACITY (10) messages are tracked, the next track() blocks the application thread until the send times out. Suggest _queue.put_nowait(...) with except queue.Full: pass, maybe as a separate issue.
| Trigger shutdown on exit of application. | ||
| Function have to be called only once - during the setup. | ||
| """ | ||
| atexit.register(shutdown) |
There was a problem hiding this comment.
At exit this runs shutdown() with flush_buffers=True, and stop_worker() then calls _worker.join() with no timeout. If data is still buffered (likely for short-lived CLIs, since the first send happens after 0.5s) and the endpoint is unreachable with packets dropped rather than refused, exit is delayed by up to SEND_TIMEOUT_SECONDS for connect plus the same again for read. DNS lookup isn't covered by that timeout. Much better than hanging forever, but still noticeable. Now that the thread is a daemon, a bounded wait at exit (e.g. join(timeout=...)) would be safe.
There was a problem hiding this comment.
Thanks, the bounded join() works well now that the thread is a daemon. One remaining path can still block before it: stop_worker() uses a blocking _queue.put() for SendBuffers and Terminate. If the worker is stuck in requests.post and the queue is full (easy to reach, since the _queue.not_full check in track() is always true), exit still waits for the send to time out. put(..., timeout=...) or put_nowait with except queue.Full would close that. Also, the shutdown() docstring says it "might introduce delay"; it could mention that buffered data may now be dropped after the timeout.
There was a problem hiding this comment.
Ah, yes, missed this part. Added
There was a problem hiding this comment.
Thanks, the put_nowait fix in track() is good, but it covers a different spot. stop_worker() still uses blocking puts:
if flush_buffers:
_queue.put(WorkerMessage.make_send_buffers())
_queue.put(WorkerMessage.make_terminate())
_worker.join(timeout=THREAD_EXIT_TIMEOUT_SECONDS)The queue can still fill up to MAX_QUEUE_CAPACITY while the worker is stuck in requests.post (offline network, the app tracks 10+ features). Then these put() calls wait until the send times out, before the bounded join() is reached, so exit is still delayed by ~30s or more. test_shutdown_exits_on_blocked_worker tracks only one feature, so it doesn't hit this. Using put(..., timeout=...) or put_nowait with except queue.Full here, as in track(), would fix it.
Small one: the shutdown() docstring could mention that, with the join timeout, buffered data may be dropped if sending takes too long.
There was a problem hiding this comment.
Thanks, the Terminate part is fixed now, but the SendBuffers put right before it is still blocking:
if flush_buffers:
_queue.put(WorkerMessage.make_send_buffers()) # still blocks on a full queue
try:
_queue.put_nowait(WorkerMessage.make_terminate())
_worker.join(timeout=THREAD_EXIT_TIMEOUT_SECONDS)
except queue.Full:
passThis is exactly the exit path: the atexit handler calls shutdown(), and its default is flush_buffers=True. So the scenario still hangs: offline network → worker stuck in requests.post → app tracks 10+ features (queue full) → at exit, stop_worker() waits on that put() until the send times out (~30s or more). Only shutdown(flush_buffers=False) avoids it.
Suggested fix: move the SendBuffers put into the same try and make it non-blocking:
if _worker is None or _queue is None:
return
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 = NoneIf the queue is full, both messages are skipped and the (stuck, daemon) worker is abandoned. If only Terminate doesn't fit, the worker will still flush once it recovers. In all cases stop_worker() returns within THREAD_EXIT_TIMEOUT_SECONDS.
Suggested test (test_worker.py). test_shutdown_exits_on_blocked_worker tracks only one feature, so the queue never fills. This one first blocks the worker in a send, then fills the queue:
import threading
def test_shutdown_exits_on_blocked_worker_with_full_queue(
telemetry_reset,
telemetry_unset_ci,
telemetry_unset_disable,
):
in_send = threading.Event()
release = threading.Event()
def blocked_post(*args, **kwargs):
in_send.set()
release.wait()
return mock.MagicMock(status_code=200)
# safety net: with the old blocking put(), the test fails instead of hanging
unblock = threading.Timer(5, release.set)
unblock.start()
with mock.patch("requests.post", side_effect=blocked_post):
assert setup(disable=False)
track("test", "0.1", "test-feature")
# first send happens after DATA_SEND_FIRST_INTERVAL_SECONDS
assert in_send.wait(timeout=5)
# worker is blocked in send -> fill the queue
for i in range(worker.MAX_QUEUE_CAPACITY + 1):
track("test", "0.1", f"feature-{i}")
old_worker, old_queue = worker._worker, worker._queue
start = time.monotonic()
shutdown(flush_buffers=True)
elapsed = time.monotonic() - start
# clean up the abandoned daemon thread
release.set()
old_queue.put(worker.WorkerMessage.make_terminate())
old_worker.join(timeout=5)
unblock.cancel()
assert elapsed < worker.THREAD_EXIT_TIMEOUT_SECONDS + 1- Waiting on
in_sendmakes sure the worker is really blocked before the queue is filled, so the test doesn't depend on timing. - With the current code,
shutdown()would block forever; thethreading.Timerreleases the send after 5s, so the timing assert fails instead of hanging the test run. - The cleanup at the end stops the abandoned worker, so no thread leaks into other tests.
I haven't run this locally, so small adjustments may be needed.
There was a problem hiding this comment.
Ok, thanks! I fixed the put.
In terms of testing, I think this should be implemented in integration tests, unit tests are leaking too much information from internal implementation. Created a task for them: #12
Overall observation is correct, setup is not exposed at the moment. The fix was needed, as one of the prior tests were broken because of the wrong logic. But at the same time, we might need a way to explicitly configure the telemetry in the future, so I think setup could be seen not as "private", but rather as "protected" API. Will add the fix description to the release notes to track the fix. |
|



Closes #10