Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions pchealthstream2py/pchealth.py
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
21 changes: 20 additions & 1 deletion pchealthstream2py/tests/simple_test.py
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
"""Simple tests"""

from pchealthstream2py.pchealth import StatusInfoReader
import threading
import time
from pprint import pprint

Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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)
Loading