diff --git a/pchealthstream2py/pchealth.py b/pchealthstream2py/pchealth.py index e96ab55..14dfa83 100644 --- a/pchealthstream2py/pchealth.py +++ b/pchealthstream2py/pchealth.py @@ -280,6 +280,11 @@ def _stop_previous_worker(self): stop, leaving two threads appending to the same queue. """ worker = self._worker + if worker is None and threading.Thread.is_alive(self): + # Started the legacy way (`start()`, so running in `self`) and then + # `open()`ed: that run must stop too, or it would miss the stop flag + # `open()` clears below and keep feeding the queue next to the new one. + worker = self if worker is None or not worker.is_alive(): return diff --git a/pchealthstream2py/tests/simple_test.py b/pchealthstream2py/tests/simple_test.py index 59ea77b..6b311f9 100644 --- a/pchealthstream2py/tests/simple_test.py +++ b/pchealthstream2py/tests/simple_test.py @@ -1,6 +1,7 @@ """Simple tests""" from pchealthstream2py.pchealth import StatusInfoReader +import threading import time from pprint import pprint @@ -43,7 +44,8 @@ def test_stop_does_not_shadow_thread_internals(): stuck on True after the worker had finished. """ reader = StatusInfoReader(read_interval_ms=50) - assert callable(reader._stop) + # (`Thread._stop` exists up to Python 3.12; 3.13 removed it.) + assert not isinstance(getattr(reader, '_stop', None), threading.Event) reader.open() try: @@ -127,3 +129,20 @@ def test_thread_api_follows_the_current_worker_after_reopen(): reader.join(timeout=10) assert not reader.is_alive() assert not reader._worker.is_alive() + + +def test_open_after_a_legacy_start_leaves_a_single_producer(): + """A reader started with `start()` (running in `self`), closed, then `open()`ed + must not keep the legacy run going: `open()` clears the stop flag, so a run it + didn't wait for would never see the stop and would keep feeding the queue.""" + reader = StatusInfoReader(read_interval_ms=50) + reader.start() + time.sleep(0.1) + reader.close() + reader.open() + try: + assert not threading.Thread.is_alive(reader) # the legacy run has stopped + assert reader.is_alive() # and the new worker is reading + finally: + reader.close() + reader.join(timeout=10)