b473aac6aa
Three things were making the panel feel slow, all of them waiting rather than working. The trailing quiet window is waited out once per query, so it set the pace of the whole sweep: at 50 ms that was 400 ms of a 590 ms poll spent listening to silence. A reply streams at the baud rate — ~1 ms between bytes, no measurable gap between its lines — and the deadline restarts on every line read, so 20 ms outlasts the gap it exists for twenty times over. A tail that still arrives late is caught by _discard_input(), which is what actually protects the next query. Replayed against the rig transcript, a sweep goes from 590 ms to 356 ms. The poll interval was 1 s on top of that, so a value could be 1.6 s stale. At 0.3 s the panel comes round about every 0.65 s. And a poll held the port for its whole sweep, so a button pressed during one waited for all eight queries. _work_pending() on the worker base lets a poll drop what is left as soon as the operator queues something: a click now waits ~135 ms for the register read in progress instead of the full sweep, and the rest is picked up next time round. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
222 lines
6.3 KiB
Python
222 lines
6.3 KiB
Python
"""gui.qt_workers: the shared queue/poll worker base."""
|
||
import threading
|
||
import time
|
||
|
||
import pytest
|
||
from PyQt6.QtCore import QThread
|
||
from PyQt6.QtWidgets import QApplication
|
||
|
||
from gui.qt_workers import PollingQueueWorker, QueueWorker
|
||
|
||
|
||
@pytest.fixture(scope="module")
|
||
def qapp():
|
||
yield QApplication.instance() or QApplication([])
|
||
|
||
|
||
class _Recorder(QueueWorker):
|
||
def __init__(self):
|
||
super().__init__()
|
||
self.seen = []
|
||
self.stopped = threading.Event()
|
||
self._handlers.update({
|
||
"note": self._note,
|
||
"boom": self._boom,
|
||
})
|
||
|
||
def _note(self, value):
|
||
self.seen.append(value)
|
||
|
||
def _boom(self):
|
||
raise RuntimeError("handler failed")
|
||
|
||
def _on_stop(self):
|
||
self.stopped.set()
|
||
|
||
|
||
def _run_until(worker, predicate, timeout=5.0):
|
||
"""Run the worker loop on a plain thread until predicate() is true.
|
||
|
||
Pumps the Qt event loop while waiting: signals emitted from the worker
|
||
thread are delivered as queued events on this (main) thread.
|
||
"""
|
||
app = QApplication.instance()
|
||
t = threading.Thread(target=worker.run, daemon=True)
|
||
t.start()
|
||
deadline = time.monotonic() + timeout
|
||
while not predicate() and time.monotonic() < deadline:
|
||
app.processEvents()
|
||
time.sleep(0.01)
|
||
app.processEvents()
|
||
return t
|
||
|
||
|
||
def test_commands_dispatch_in_order(qapp):
|
||
w = _Recorder()
|
||
for i in range(5):
|
||
w._enqueue("note", value=i)
|
||
t = _run_until(w, lambda: len(w.seen) == 5)
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
assert w.seen == [0, 1, 2, 3, 4]
|
||
assert w.stopped.is_set()
|
||
|
||
|
||
def test_handler_exception_is_reported_not_fatal(qapp):
|
||
w = _Recorder()
|
||
errors = []
|
||
w.error_occurred.connect(errors.append)
|
||
w._enqueue("boom")
|
||
w._enqueue("note", value="after")
|
||
t = _run_until(w, lambda: w.seen == ["after"])
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
assert w.seen == ["after"], "loop died on a failing handler"
|
||
assert errors and "handler failed" in errors[0]
|
||
|
||
|
||
def test_unknown_command_reported(qapp):
|
||
w = _Recorder()
|
||
errors = []
|
||
w.error_occurred.connect(errors.append)
|
||
w._enqueue("nope")
|
||
t = _run_until(w, lambda: bool(errors))
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
assert errors and "Unknown command" in errors[0]
|
||
|
||
|
||
def test_idle_worker_does_not_spin(qapp):
|
||
"""The loop must block on the queue, not poll it on a timeout."""
|
||
w = _Recorder()
|
||
t = threading.Thread(target=w.run, daemon=True)
|
||
t.start()
|
||
time.sleep(0.3) # idle
|
||
cpu_before = time.process_time()
|
||
time.sleep(0.5) # still idle
|
||
cpu_used = time.process_time() - cpu_before
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
# A 10–20 Hz timeout-poll loop burns measurable CPU here; blocking uses ~0.
|
||
assert cpu_used < 0.05, f"idle worker used {cpu_used:.3f}s CPU"
|
||
|
||
|
||
class _Poller(PollingQueueWorker):
|
||
def __init__(self):
|
||
super().__init__(poll_interval_s=0.02)
|
||
self.polls = 0
|
||
self.in_flight = 0
|
||
self.overlaps = 0
|
||
self.is_connected = True
|
||
|
||
def _poll_once(self):
|
||
self.in_flight += 1
|
||
if self.in_flight > 1:
|
||
self.overlaps += 1
|
||
time.sleep(0.05) # deliberately slower than the poll interval
|
||
self.polls += 1
|
||
self.in_flight -= 1
|
||
|
||
|
||
def test_polling_never_overlaps_or_backs_up(qapp):
|
||
"""A device slower than the interval must not accumulate stale polls."""
|
||
w = _Poller()
|
||
t = threading.Thread(target=w.run, daemon=True)
|
||
t.start()
|
||
w.start_polling()
|
||
time.sleep(0.6)
|
||
w.stop_polling()
|
||
time.sleep(0.15)
|
||
queued = w._cmd_q.qsize()
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
|
||
assert w.polls >= 3, "polling did not run"
|
||
assert w.overlaps == 0, "polls overlapped"
|
||
# Self-rescheduling means at most one poll is ever pending.
|
||
assert queued <= 1, f"{queued} stale polls queued up"
|
||
|
||
|
||
class _YieldingPoller(PollingQueueWorker):
|
||
"""A poll made of several reads that gives up as soon as work arrives."""
|
||
|
||
def __init__(self):
|
||
super().__init__(poll_interval_s=0.02)
|
||
self.reads = 0
|
||
self.handled = []
|
||
self.is_connected = True
|
||
self._handlers["click"] = self._click
|
||
|
||
def _click(self, value):
|
||
self.handled.append(value)
|
||
|
||
def _poll_once(self):
|
||
for _ in range(6):
|
||
if self._work_pending():
|
||
return
|
||
time.sleep(0.02)
|
||
self.reads += 1
|
||
|
||
|
||
def test_a_queued_command_interrupts_a_poll(qapp):
|
||
"""A button pressed mid-poll should not wait out the whole sweep.
|
||
|
||
The Helios sweep is eight serial queries; before this, a command queued
|
||
behind one waited for every last read to finish.
|
||
"""
|
||
w = _YieldingPoller()
|
||
t = threading.Thread(target=w.run, daemon=True)
|
||
t.start()
|
||
w.start_polling()
|
||
time.sleep(0.03) # a poll is now in progress
|
||
|
||
pressed = time.monotonic()
|
||
w._enqueue("click", value="set current")
|
||
deadline = pressed + 2.0
|
||
while not w.handled and time.monotonic() < deadline:
|
||
time.sleep(0.002)
|
||
waited = time.monotonic() - pressed
|
||
|
||
w.stop_polling()
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
|
||
assert w.handled == ["set current"]
|
||
# A full sweep is 6 x 20 ms; the command must not have waited for it.
|
||
assert waited < 0.08, f"command waited {waited * 1000:.0f} ms for the poll"
|
||
|
||
|
||
def test_stop_polling_halts_the_cycle(qapp):
|
||
w = _Poller()
|
||
t = threading.Thread(target=w.run, daemon=True)
|
||
t.start()
|
||
w.start_polling()
|
||
time.sleep(0.2)
|
||
w.stop_polling()
|
||
time.sleep(0.2)
|
||
settled = w.polls
|
||
time.sleep(0.2)
|
||
w.stop_worker()
|
||
t.join(timeout=5)
|
||
assert w.polls == settled, "polling continued after stop_polling()"
|
||
|
||
|
||
def test_worker_runs_on_its_qthread(qapp):
|
||
"""Sanity check the intended usage: run() executes on the QThread."""
|
||
w = _Recorder()
|
||
thread = QThread()
|
||
w.moveToThread(thread)
|
||
thread.started.connect(w.run)
|
||
ids = []
|
||
w._handlers["note"] = lambda value: ids.append(threading.get_ident())
|
||
thread.start()
|
||
w._enqueue("note", value=None)
|
||
deadline = time.monotonic() + 5
|
||
while not ids and time.monotonic() < deadline:
|
||
qapp.processEvents()
|
||
time.sleep(0.01)
|
||
w.stop_worker()
|
||
thread.quit()
|
||
assert thread.wait(5000)
|
||
assert ids and ids[0] != threading.get_ident()
|