fix(voice): wake-word pause/stop no longer hangs on a wedged microphone stream
_Capture.read polls read_available against the stop event before every InputStream.read, close() aborts before closing, and _halt_thread joins outside the detector lock and keeps the handle while the thread is alive. Fixes #117096
This commit is contained in:
@@ -619,6 +619,56 @@ def test_detection_callback_can_pause_and_close_stream(monkeypatch, tmp_path):
|
||||
assert ww.stop_listening(owner=owner) is True
|
||||
|
||||
|
||||
def test_wedged_stream_halts_without_blocking_read_and_aborts_before_close(monkeypatch):
|
||||
"""A PortAudio device that never delivers samples (#117096) must not wedge pause().
|
||||
|
||||
``read(n)`` on such a device blocks forever; the detector must never call it
|
||||
while ``read_available`` is short, must return from pause() promptly once the
|
||||
stop event is set, and must ``abort()`` the stream before ``close()``.
|
||||
"""
|
||||
class _WedgedStream(_FakeStream):
|
||||
read_available = 0
|
||||
|
||||
def __init__(self, **kw):
|
||||
super().__init__(**kw)
|
||||
self.calls = []
|
||||
self.read_calls = 0
|
||||
|
||||
def read(self, n):
|
||||
self.read_calls += 1
|
||||
time.sleep(30) # a real wedged ALSA/PipeWire read never returns
|
||||
return [0] * n, False
|
||||
|
||||
def abort(self):
|
||||
self.calls.append("abort")
|
||||
|
||||
def stop(self):
|
||||
self.calls.append("stop")
|
||||
|
||||
def close(self):
|
||||
self.calls.append("close")
|
||||
self.closed = True
|
||||
|
||||
streams = []
|
||||
|
||||
def _stream(**kw):
|
||||
streams.append(_WedgedStream(**kw))
|
||||
return streams[-1]
|
||||
|
||||
monkeypatch.setattr(ww, "_import_audio", lambda: (types.SimpleNamespace(InputStream=_stream), None))
|
||||
det = ww.WakeWordDetector(_FakeEngine(fire=False), on_wake=lambda: None)
|
||||
det.start()
|
||||
assert det.running is True
|
||||
time.sleep(0.2)
|
||||
t0 = time.monotonic()
|
||||
det.pause()
|
||||
assert time.monotonic() - t0 < 1.5, "pause() must not wait out the join timeout"
|
||||
assert det.running is False
|
||||
stream = streams[0]
|
||||
assert stream.read_calls == 0, "read() must not be entered while read_available < frame_length"
|
||||
assert stream.calls[:2] == ["abort", "close"]
|
||||
|
||||
|
||||
def test_startup_failure_releases_owner_and_machine_lock(monkeypatch, tmp_path):
|
||||
class _BrokenSoundDevice:
|
||||
@staticmethod
|
||||
|
||||
@@ -31,6 +31,7 @@ SAMPLE_RATE = 16000 # 16 kHz mono int16 — Whisper-native and what every engin
|
||||
# several frames while the caller is still reacting.
|
||||
_FIRE_COOLDOWN_SECONDS = 2.0
|
||||
_START_TIMEOUT_SECONDS = 5.0
|
||||
_READ_POLL_SECONDS = 0.05 # slice between read_available polls; bounds halt latency
|
||||
|
||||
# Ambient-speech rejection: N consecutive over-threshold frames before firing
|
||||
# (a stray phoneme spikes one frame; a real phrase holds several).
|
||||
@@ -430,19 +431,39 @@ class _Capture:
|
||||
rate: int = SAMPLE_RATE
|
||||
frame_length: int = 1280 # samples per read at ``rate``
|
||||
|
||||
def read(self):
|
||||
"""One raw block; None when no client frame arrived within 250 ms. Stream errors propagate."""
|
||||
def read(self, stop: Optional[threading.Event] = None):
|
||||
"""One raw block; None when nothing arrived within ~250 ms (client) or ``stop`` was
|
||||
set while waiting (local). Stream errors propagate.
|
||||
|
||||
A PortAudio ``read(n)`` blocks until ``n`` samples exist and, on a wedged ALSA/
|
||||
PipeWire device, never returns — so the halting thread's ``join`` timed out and
|
||||
``close()`` raced the still-pending read. Poll ``read_available`` in short slices
|
||||
against ``stop`` and only call ``read`` once the block is guaranteed to be there.
|
||||
"""
|
||||
if self.stream is not None:
|
||||
available = getattr(self.stream, "read_available", None)
|
||||
if stop is not None and available is not None:
|
||||
while self.stream.read_available < self.frame_length:
|
||||
if stop.wait(_READ_POLL_SECONDS):
|
||||
return None
|
||||
return self.stream.read(self.frame_length)[0]
|
||||
with suppress(Exception):
|
||||
return self.queue.get(timeout=0.25)
|
||||
return None
|
||||
|
||||
def close(self) -> None:
|
||||
"""``abort()`` first: it discards pending buffers and unblocks any in-flight read,
|
||||
which ``stop()`` (drains, waits) cannot do on a dead device."""
|
||||
if self.stream is None:
|
||||
return
|
||||
abort = getattr(self.stream, "abort", None)
|
||||
with suppress(Exception):
|
||||
if self.stream is not None:
|
||||
if abort is not None:
|
||||
abort()
|
||||
else:
|
||||
self.stream.stop()
|
||||
self.stream.close()
|
||||
with suppress(Exception):
|
||||
self.stream.close()
|
||||
|
||||
|
||||
class WakeWordDetector:
|
||||
@@ -533,9 +554,14 @@ class WakeWordDetector:
|
||||
with self._lock:
|
||||
self._stop.set()
|
||||
t = self._thread
|
||||
if t is not None and t is not threading.current_thread():
|
||||
t.join(timeout=2.0)
|
||||
if self._thread is t:
|
||||
# Join OUTSIDE the lock: a reader wedged in PortAudio would otherwise pin the lock
|
||||
# for the whole timeout and stall every start()/pause() caller behind it.
|
||||
if t is not None and t is not threading.current_thread():
|
||||
t.join(timeout=2.0)
|
||||
with self._lock:
|
||||
# Keep the handle while the thread is still alive (join timed out) so
|
||||
# ``running`` stays truthful and the next start() does not double-arm.
|
||||
if self._thread is t and (t is None or not t.is_alive()):
|
||||
self._thread = None
|
||||
|
||||
def _dispatch_wake(self) -> None:
|
||||
@@ -633,7 +659,7 @@ class WakeWordDetector:
|
||||
try:
|
||||
while not self._stop.is_set():
|
||||
try:
|
||||
data = cap.read()
|
||||
data = cap.read(self._stop)
|
||||
except Exception as e:
|
||||
logger.warning("wake word: stream read error: %s", e)
|
||||
failed = not self._stop.is_set()
|
||||
|
||||
Reference in New Issue
Block a user