/
/
1"""
2Unit tests for the output-device recovery in the Sendspin -> Local Audio bridge.
3
4Device calls are answered by mock streams -- run inline, or handed back on a
5future the test controls -- so no test ever opens a sound card and none of them
6depend on the host's wall clock. Cover six things, across both writers
7(PulseAudio and sounddevice) wherever the two paths differ:
8
9* mid-stream recovery: a write that fails closes the dead stream, reopens the
10 device and carries on writing to the fresh one, with a reopened PA sink both
11 re-arming the hardware-volume re-apply that the idle->RUNNING transition would
12 otherwise let PulseAudio's own device-restore undo, and re-pinned to unity
13 where this player's volume is applied in software instead;
14* the reopen budget: a device that lasted beyond the guard window earns another
15 reopen, while one that fails again inside it is not churned for the rest of the
16 stream, and every writer starts with a budget of its own;
17* giving up: a device that cannot be (re)opened -- the guard refusing, the reopen
18 raising, the very first open failing, or the writer dying on something
19 unclassified -- stops the feed and takes the bridge out of its Sendspin
20 session, so the group stops holding the player on PLAYING for audio nobody can
21 hear, and a device that opens but will not start is handed back rather than
22 leaked;
23* not giving up: a stream that simply ended, was stopped, or had its writer
24 cancelled keeps its place in the session;
25* leaving the session itself: the client is quiesced, never unregistered, an
26 absent client is a silent no-op and a refusal is reported rather than raised;
27* writer supersession: a writer replaced by a newer one can neither silence its
28 replacement nor be the reason a teardown scheduled for itself stops it.
29"""
30
31import asyncio
32import importlib
33import time
34from collections.abc import Callable, Coroutine
35from contextlib import AbstractContextManager, suppress
36from types import ModuleType
37from typing import Any, cast
38from unittest.mock import AsyncMock, MagicMock, patch
39
40import pytest
41
42from music_assistant.providers.local_audio.sendspin_bridge import (
43 _DEVICE_REOPEN_GUARD_SECONDS,
44 SendspinLocalAudioBridge,
45)
46from music_assistant.providers.sendspin.bridge_role import (
47 BRIDGE_BIT_DEPTH,
48 BRIDGE_CHANNELS,
49 BRIDGE_SAMPLE_RATE,
50)
51from tests.common import collect_loop_errors
52
53_MODULE = "music_assistant.providers.local_audio.sendspin_bridge"
54
55# An arbitrary reading of the bridge clock. Chunks are queued at exactly this
56# instant, so they are due now: neither slept on nor dropped as late.
57NOW_US = 4_200_000_000
58
59
60def _import_sounddevice() -> ModuleType | None:
61 """Return the sounddevice module, or None where the PortAudio library is missing."""
62 try:
63 return importlib.import_module("sounddevice")
64 except OSError:
65 # only the missing shared library is a reason to skip; anything else is
66 # a real problem and silently skipping it would report green on nothing
67 return None
68
69
70# The sounddevice writer imports the module for real, and importing it fails
71# outright on a host that has the Python package but not the PortAudio shared
72# library it binds to. That writer is then untestable rather than broken.
73_SOUNDDEVICE = _import_sounddevice()
74needs_portaudio = pytest.mark.skipif(
75 _SOUNDDEVICE is None, reason="the sounddevice writer needs the PortAudio shared library"
76)
77
78BACKENDS = [
79 pytest.param("pulse", id="pulse"),
80 pytest.param("sounddevice", id="sounddevice", marks=needs_portaudio),
81]
82
83
84def _run_inline(_executor: object, func: Callable[..., Any], *args: Any) -> asyncio.Future[Any]:
85 """
86 Stand in for ``loop.run_in_executor``, running the blocking call inline.
87
88 The outcome is handed back on an already-completed future, so a device that
89 raises on write really propagates into the writer awaiting it.
90 """
91 future: asyncio.Future[Any] = asyncio.get_running_loop().create_future()
92 try:
93 future.set_result(func(*args))
94 except Exception as err:
95 future.set_exception(err)
96 return future
97
98
99def _schedule(target: Any, *_args: Any, eager_start: bool = True, **_kwargs: Any) -> Any:
100 """
101 Stand in for ``mass.create_task``, running real coroutines as real tasks.
102
103 A target that a test patched away is only recorded, so what the bridge
104 scheduled stays assertable without a coroutine ever being built.
105 """
106 if not asyncio.iscoroutine(target):
107 return MagicMock()
108 return asyncio.Task(target, loop=asyncio.get_running_loop(), eager_start=eager_start)
109
110
111class _MonotonicClock:
112 """A monotonic clock the test moves by hand, so the reopen guard window is exact."""
113
114 def __init__(self) -> None:
115 """Start from the current real reading, leaving the event loop's own timers sane."""
116 self._now = time.monotonic()
117
118 def __call__(self) -> float:
119 """Return the current reading."""
120 return self._now
121
122 def advance(self, seconds: float) -> None:
123 """Move the clock forward by the given number of seconds."""
124 self._now += seconds
125
126
127def _make_bridge(
128 backend: str = "pulse", volume_controller: MagicMock | None = None
129) -> SendspinLocalAudioBridge:
130 """
131 Build a streaming bridge on mocked infrastructure, bound to no real device.
132
133 :param backend: The audio backend the writer dispatches on.
134 :param volume_controller: A PA volume controller driving the sink's hardware
135 volume, or None for the software-volume path.
136 """
137 mass = MagicMock()
138 mass.loop.run_in_executor = MagicMock(side_effect=_run_inline)
139 mass.create_task = MagicMock(side_effect=_schedule)
140 provider = MagicMock()
141 provider.mass = mass
142 device_info = {
143 "name": "alsa_output.pci-0000_00_1f.3.analog-stereo",
144 "description": "Study Speakers",
145 "pa_sink_name": "alsa_output.pci-0000_00_1f.3.analog-stereo",
146 "index": 2,
147 "sample_rate": BRIDGE_SAMPLE_RATE,
148 "bit_depth": BRIDGE_BIT_DEPTH,
149 "max_output_channels": BRIDGE_CHANNELS,
150 }
151 bridge = SendspinLocalAudioBridge(
152 provider, device_info, MagicMock(), backend=backend, volume_controller=volume_controller
153 )
154 # Unity gain, so a chunk reaches the device as the exact bytes it was queued with.
155 bridge._volume_level = 100
156 bridge._is_streaming = True
157 return bridge
158
159
160def _make_device_stream(write_effect: Any = None) -> MagicMock:
161 """
162 Build a mock output stream standing in for an open PA sink or sound device.
163
164 :param write_effect: What the blocking write does â an exception it raises,
165 or a sequence of per-call outcomes.
166 """
167 stream = MagicMock()
168 stream.write = MagicMock(side_effect=write_effect)
169 return stream
170
171
172def _patch_device_open(
173 bridge: SendspinLocalAudioBridge, *streams: Any
174) -> AbstractContextManager[Any]:
175 """
176 Patch the bridge's device-open call to hand out the given streams in order.
177
178 An exception among them is raised by that open instead. Patching the call
179 rather than the backend library keeps the tests running on any platform.
180 """
181 if bridge.backend == "pulse":
182 return patch.object(bridge, "_get_pa_stream", AsyncMock(side_effect=list(streams)))
183 return patch.object(bridge, "_get_sounddevice_stream", AsyncMock(side_effect=list(streams)))
184
185
186def _device_error(bridge: SendspinLocalAudioBridge) -> Exception:
187 """Return the failure the bridge's backend raises when its device goes away."""
188 if bridge.backend == "pulse":
189 return OSError("sink vanished")
190 assert _SOUNDDEVICE is not None
191 error: Exception = _SOUNDDEVICE.PortAudioError("device disconnected")
192 return error
193
194
195def _pcm_chunk(marker: int) -> bytes:
196 """Build a short, distinguishable 16-bit stereo PCM payload."""
197 return marker.to_bytes(2, "little") * BRIDGE_CHANNELS * 8
198
199
200def _queue_chunks(bridge: SendspinLocalAudioBridge, *chunks: bytes) -> None:
201 """Queue chunks that are all due right now, followed by the writer's stop sentinel."""
202 for chunk in chunks:
203 bridge._write_queue.put_nowait((NOW_US, chunk))
204 bridge._write_queue.put_nowait(None)
205
206
207def _install_writer(
208 bridge: SendspinLocalAudioBridge, coro: Coroutine[Any, Any, None]
209) -> asyncio.Task[None]:
210 """Start the given coroutine as the writer task the bridge owns."""
211 task = asyncio.create_task(coro)
212 bridge._writer_task = task
213 return task
214
215
216async def _run_writer(bridge: SendspinLocalAudioBridge) -> None:
217 """
218 Run the audio writer to completion as the writer the bridge owns.
219
220 The bridge clock is pinned to the instant the queued chunks are due, so the
221 writer neither sleeps on them nor drops them as late.
222 """
223 with patch(f"{_MODULE}._now_us", return_value=NOW_US):
224 await _install_writer(bridge, bridge._audio_writer())
225
226
227async def _never_returns() -> None:
228 """Stand in for a writer still feeding its stream."""
229 await asyncio.Event().wait()
230
231
232async def _settle() -> None:
233 """Let every task the bridge scheduled reach its own next suspension."""
234 # a handful of turns: a teardown can be several tasks deep before it settles
235 for _ in range(5):
236 await asyncio.sleep(0)
237
238
239def _scheduled_targets(bridge: SendspinLocalAudioBridge) -> list[Any]:
240 """Return everything the bridge handed to ``mass.create_task``."""
241 return [call.args[0] for call in cast("MagicMock", bridge.mass).create_task.call_args_list]
242
243
244# --- Mid-stream recovery: a dead device is replaced and the audio carries on ---
245
246
247async def test_a_failed_pa_write_reopens_the_sink_and_keeps_writing() -> None:
248 """
249 A PA sink that dies mid-stream is replaced and the audio keeps flowing.
250
251 The chunk that hit the dead sink goes with it, but every chunk after it
252 reaches the fresh stream instead of being dropped for the rest of the track.
253 """
254 bridge = _make_bridge()
255 dead = _make_device_stream(write_effect=OSError("connection terminated"))
256 fresh = _make_device_stream()
257 first, second = _pcm_chunk(1), _pcm_chunk(2)
258 _queue_chunks(bridge, first, second)
259
260 with (
261 _patch_device_open(bridge, dead, fresh),
262 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
263 ):
264 await _run_writer(bridge)
265
266 dead.write.assert_called_once_with(first)
267 fresh.write.assert_called_once_with(second)
268 # libpulse has no reconnect, so the failed connection is released outright
269 dead.close.assert_called_once_with()
270 leave.assert_not_called()
271
272
273@needs_portaudio
274async def test_a_failed_sounddevice_write_reopens_the_device_and_keeps_writing() -> None:
275 """
276 A sound device that drops out mid-stream is reopened and the audio keeps flowing.
277
278 The chunk that hit the dead device goes with it, but every chunk after it
279 reaches the fresh stream instead of being dropped for the rest of the track.
280 """
281 bridge = _make_bridge(backend="sounddevice")
282 dead = _make_device_stream(write_effect=_device_error(bridge))
283 fresh = _make_device_stream()
284 first, second = _pcm_chunk(1), _pcm_chunk(2)
285 _queue_chunks(bridge, first, second)
286
287 with (
288 _patch_device_open(bridge, dead, fresh),
289 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
290 ):
291 await _run_writer(bridge)
292
293 dead.write.assert_called_once_with(first)
294 fresh.write.assert_called_once_with(second)
295 dead.close.assert_called_once_with()
296 # a device that stopped responding can never finish playing its buffer out
297 dead.stop.assert_not_called()
298 leave.assert_not_called()
299
300
301async def test_a_reopened_pa_sink_re_arms_the_hardware_volume_reapply() -> None:
302 """
303 A reopened sink has the cached hardware volume applied to it a second time.
304
305 The fresh sink runs through idle->RUNNING again on its own first write, which
306 is exactly when PulseAudio's device-restore can put back a stale level.
307 """
308 bridge = _make_bridge(volume_controller=MagicMock())
309 dead = _make_device_stream(write_effect=[None, OSError("sink vanished")])
310 fresh = _make_device_stream()
311 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2), _pcm_chunk(3))
312
313 with (
314 _patch_device_open(bridge, dead, fresh),
315 patch.object(bridge, "_schedule_hardware_volume_reapply", MagicMock()) as reapply,
316 ):
317 await _run_writer(bridge)
318
319 assert dead.write.call_count == 2
320 fresh.write.assert_called_once_with(_pcm_chunk(3))
321 assert reapply.call_count == 2
322
323
324async def test_a_reopened_pa_sink_is_re_pinned_to_unity_without_hardware_volume() -> None:
325 """
326 A sink reopened on the software-volume path is pinned back to unity.
327
328 Opening a stream can bleed its own volume into the sink's hardware volume,
329 which on this path has to stay at unity or it quietly attenuates everything
330 feeding through it.
331 """
332 bridge = _make_bridge()
333 dead = _make_device_stream(write_effect=OSError("sink vanished"))
334 fresh = _make_device_stream()
335 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2))
336
337 with (
338 _patch_device_open(bridge, dead, fresh),
339 patch.object(bridge, "_reset_sink_volume", AsyncMock()) as reset_volume,
340 ):
341 await _run_writer(bridge)
342
343 fresh.write.assert_called_once_with(_pcm_chunk(2))
344 reset_volume.assert_awaited_once_with()
345
346
347# --- The reopen budget ---------------------------------------------------------
348
349
350async def test_a_device_that_lasted_beyond_the_guard_window_is_reopened_again() -> None:
351 """
352 A sink that played on for a while before failing again earns another reopen.
353
354 Two unrelated dropouts within one stream are a device having a bad day, not
355 one that is unusable, so the bridge keeps it playing.
356 """
357 bridge = _make_bridge()
358 clock = _MonotonicClock()
359
360 def _fail_after_the_guard_window(_data: bytes) -> None:
361 clock.advance(_DEVICE_REOPEN_GUARD_SECONDS + 1)
362 raise OSError("sink vanished again")
363
364 first_dead = _make_device_stream(write_effect=OSError("sink vanished"))
365 second_dead = _make_device_stream(write_effect=_fail_after_the_guard_window)
366 fresh = _make_device_stream()
367 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2), _pcm_chunk(3))
368
369 with (
370 patch(f"{_MODULE}.time.monotonic", clock),
371 _patch_device_open(bridge, first_dead, second_dead, fresh),
372 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
373 ):
374 await _run_writer(bridge)
375
376 fresh.write.assert_called_once_with(_pcm_chunk(3))
377 leave.assert_not_called()
378
379
380async def test_a_device_failing_again_inside_the_guard_window_leaves_the_session() -> None:
381 """
382 A sink that will not stay open is taken out of the Sendspin session.
383
384 Churning it for the rest of the stream would leave the group holding the
385 player on PLAYING while it renders nothing at all.
386 """
387 bridge = _make_bridge()
388 clock = _MonotonicClock()
389 first_dead = _make_device_stream(write_effect=OSError("sink vanished"))
390 second_dead = _make_device_stream(write_effect=OSError("sink vanished again"))
391 spare = _make_device_stream()
392 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2), _pcm_chunk(3))
393
394 with (
395 patch(f"{_MODULE}.time.monotonic", clock),
396 _patch_device_open(bridge, first_dead, second_dead, spare) as open_device,
397 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
398 ):
399 await _run_writer(bridge)
400
401 # the refused reopen never reached the device
402 assert open_device.call_count == 2
403 assert bridge._is_streaming is False
404 leave.assert_called_once_with()
405 # scheduled, not merely constructed: an unscheduled coroutine never leaves
406 assert leave.return_value in _scheduled_targets(bridge)
407
408
409async def test_each_writer_starts_with_a_fresh_reopen_budget() -> None:
410 """
411 A device failure on the previous stream cannot condemn the device on this one.
412
413 The guard measures how long a reopened device lasted, so a stamp carried over
414 from an earlier stream would give up on the very first failure of a device
415 that has been playing fine ever since.
416 """
417 bridge = _make_bridge()
418 clock = _MonotonicClock()
419 # a stamp left behind by a previous stream, well inside the guard window
420 bridge._last_device_reopen = clock()
421 dead = _make_device_stream(write_effect=OSError("sink vanished"))
422 fresh = _make_device_stream()
423 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2))
424
425 with (
426 patch(f"{_MODULE}.time.monotonic", clock),
427 _patch_device_open(bridge, dead, fresh),
428 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
429 ):
430 await _run_writer(bridge)
431
432 fresh.write.assert_called_once_with(_pcm_chunk(2))
433 leave.assert_not_called()
434
435
436# --- Giving up on a device that cannot be (re)opened ---------------------------
437
438
439@pytest.mark.parametrize("backend", BACKENDS)
440async def test_a_device_that_cannot_be_reopened_leaves_the_session(backend: str) -> None:
441 """A reopen the device itself refuses ends the stream and leaves the session."""
442 bridge = _make_bridge(backend=backend)
443 dead = _make_device_stream(write_effect=_device_error(bridge))
444 _queue_chunks(bridge, _pcm_chunk(1), _pcm_chunk(2))
445
446 with (
447 _patch_device_open(bridge, dead, _device_error(bridge)),
448 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
449 ):
450 await _run_writer(bridge)
451
452 assert bridge._is_streaming is False
453 leave.assert_called_once_with()
454 assert leave.return_value in _scheduled_targets(bridge)
455
456
457@pytest.mark.parametrize("backend", BACKENDS)
458async def test_a_device_that_never_opens_leaves_the_session(backend: str) -> None:
459 """
460 A device already gone when the stream starts leaves the session too.
461
462 Nothing restarts the writer within a stream, so this player has no audio
463 path at all for the whole of it.
464 """
465 bridge = _make_bridge(backend=backend)
466 _queue_chunks(bridge, _pcm_chunk(1))
467
468 with (
469 _patch_device_open(bridge, _device_error(bridge)),
470 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
471 ):
472 await _run_writer(bridge)
473
474 assert bridge._is_streaming is False
475 leave.assert_called_once_with()
476 assert leave.return_value in _scheduled_targets(bridge)
477
478
479@pytest.mark.parametrize("backend", BACKENDS)
480async def test_an_unexpected_writer_failure_leaves_the_session(backend: str) -> None:
481 """
482 A writer dying on something nobody anticipated takes the player with it.
483
484 The device is this player's only audio path, so an unclassified failure is
485 still silence, and the group must stop reporting it as playback. The device
486 is handed back on the way out, whatever the writer died on.
487 """
488 bridge = _make_bridge(backend=backend)
489 stream = _make_device_stream(write_effect=ValueError("unclassified device failure"))
490 _queue_chunks(bridge, _pcm_chunk(1))
491
492 with (
493 _patch_device_open(bridge, stream),
494 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
495 ):
496 await _run_writer(bridge)
497
498 assert bridge._is_streaming is False
499 stream.close.assert_called_once_with()
500 leave.assert_called_once_with()
501 assert leave.return_value in _scheduled_targets(bridge)
502
503
504async def test_giving_up_really_quiesces_the_sendspin_client() -> None:
505 """
506 A writer that gives up takes the client out of its session for real.
507
508 Scheduling the leave and performing it are separate steps, so this drives
509 the whole path end to end rather than either half of it.
510 """
511 bridge = _make_bridge()
512 client = MagicMock()
513 client.quiesce_to_solo_stopped = AsyncMock()
514 bridge._sendspin_client = client
515 _queue_chunks(bridge, _pcm_chunk(1))
516
517 with _patch_device_open(bridge, OSError("sink vanished")):
518 await _run_writer(bridge)
519 await _settle() # the leave is deferred off the writer
520
521 client.quiesce_to_solo_stopped.assert_awaited_once_with()
522
523
524@pytest.mark.parametrize("backend", BACKENDS)
525async def test_a_device_that_never_finishes_opening_leaves_the_session(backend: str) -> None:
526 """
527 A device that hangs on open is given up on rather than waited for.
528
529 Waiting holds the writer with nothing consuming its queue, so the chunks
530 pile up behind a player that goes on reporting playback it cannot produce.
531 """
532 bridge = _make_bridge(backend=backend)
533 opening: asyncio.Future[Any] = asyncio.get_running_loop().create_future()
534 cast("MagicMock", bridge.mass).loop.run_in_executor = MagicMock(return_value=opening)
535 _queue_chunks(bridge, _pcm_chunk(1))
536
537 with (
538 patch(f"{_MODULE}._DEVICE_OPEN_TIMEOUT_SECONDS", 0.01),
539 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
540 ):
541 await _run_writer(bridge)
542
543 assert bridge._is_streaming is False
544 leave.assert_called_once_with()
545
546
547async def test_a_device_failing_after_the_open_timeout_logs_no_loop_error() -> None:
548 """A device open that fails once the writer gave up on it is not reported to the loop."""
549 bridge = _make_bridge()
550 opening: asyncio.Future[Any] = asyncio.get_running_loop().create_future()
551 cast("MagicMock", bridge.mass).loop.run_in_executor = MagicMock(return_value=opening)
552 _queue_chunks(bridge, _pcm_chunk(1))
553
554 with (
555 collect_loop_errors() as reported,
556 patch(f"{_MODULE}._DEVICE_OPEN_TIMEOUT_SECONDS", 0.01),
557 patch.object(bridge, "_leave_sendspin_session", MagicMock()),
558 ):
559 await _run_writer(bridge)
560 # fail the open only once the writer has given up waiting for it, so the failure
561 # reliably lands after the caller is gone
562 opening.set_exception(OSError("device is already in use"))
563 await _settle()
564
565 assert reported == []
566
567
568async def test_abandoning_a_stream_defers_the_leave_off_the_writer() -> None:
569 """
570 Giving up stops the feed and hands the leave to a task of its own.
571
572 Leaving ends the Sendspin stream, which unwinds straight back into the
573 bridge's own teardown â that must not run inside the writer on its way out.
574 """
575 bridge = _make_bridge()
576 # only the registered writer may give up, which is what every caller is
577 bridge._writer_task = cast("asyncio.Task[None]", asyncio.current_task())
578
579 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
580 bridge._abandon_streaming()
581
582 assert bridge._is_streaming is False
583 leave.assert_called_once_with()
584 create_task = cast("MagicMock", bridge.mass).create_task
585 assert create_task.call_args.args[0] is leave.return_value
586 assert create_task.call_args.kwargs["eager_start"] is False
587
588
589async def test_a_writer_a_stream_start_installed_can_still_give_up() -> None:
590 """
591 A writer started by a real stream start is registered before it runs.
592
593 Everything it touches on the bridge is keyed on being the registered
594 writer, so a device failing at the very first open has to still take the
595 player out of the session rather than leave it reporting playback.
596 """
597 bridge = _make_bridge()
598
599 with (
600 _patch_device_open(bridge, OSError("sink vanished")),
601 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
602 ):
603 bridge._on_bridge_stream_start()
604 await _settle()
605
606 assert bridge._is_streaming is False
607 leave.assert_called_once_with()
608
609
610# --- Handing the device back -------------------------------------------------
611
612
613@needs_portaudio
614async def test_a_stream_that_will_not_start_is_handed_back_to_portaudio() -> None:
615 """
616 A device that opens but refuses to start is released rather than leaked.
617
618 PortAudio has already handed the device out by then, and a handle nobody
619 gives back keeps every later attempt on that device failing.
620 """
621 bridge = _make_bridge(backend="sounddevice")
622 stream = _make_device_stream()
623 stream.start = MagicMock(side_effect=_device_error(bridge))
624
625 with (
626 patch("sounddevice.RawOutputStream", MagicMock(return_value=stream)),
627 pytest.raises(Exception, match="device disconnected"),
628 ):
629 bridge._create_sounddevice_stream()
630
631 stream.close.assert_called_once_with()
632
633
634@needs_portaudio
635async def test_a_cancelled_open_releases_the_device_it_still_gets() -> None:
636 """
637 A device that finishes opening after its writer is gone is closed, not orphaned.
638
639 PortAudio hands the device out whether or not anyone is still waiting, and a
640 handle nobody gives back makes every later open of that device fail.
641 """
642 bridge = _make_bridge(backend="sounddevice")
643 stream = _make_device_stream()
644 opening: asyncio.Future[Any] = asyncio.get_running_loop().create_future()
645
646 def _defer_the_open(
647 executor: object, func: Callable[..., Any], *args: Any
648 ) -> asyncio.Future[Any]:
649 # the open takes no arguments; anything else is the release closing up
650 return _run_inline(executor, func, *args) if args else opening
651
652 cast("MagicMock", bridge.mass).loop.run_in_executor = MagicMock(side_effect=_defer_the_open)
653 task = asyncio.create_task(bridge._get_sounddevice_stream())
654 await _settle()
655 task.cancel()
656 with suppress(asyncio.CancelledError):
657 await task
658
659 # the device only becomes available once nobody is waiting for it any more
660 opening.set_result(stream)
661 await _settle()
662
663 stream.close.assert_called_once_with()
664
665
666async def test_a_timed_out_open_releases_the_device_when_it_finally_arrives() -> None:
667 """
668 A device that opens long after the wait gave up is closed, not left behind.
669
670 The open runs to completion in its own thread whatever the bridge does, and
671 a sink nobody hands back keeps every later open of it failing.
672 """
673 bridge = _make_bridge()
674 stream = _make_device_stream()
675 opening: asyncio.Future[Any] = asyncio.get_running_loop().create_future()
676
677 def _defer_the_open(
678 executor: object, func: Callable[..., Any], *args: Any
679 ) -> asyncio.Future[Any]:
680 # the open takes no arguments; anything else is the release closing up
681 return _run_inline(executor, func, *args) if args else opening
682
683 cast("MagicMock", bridge.mass).loop.run_in_executor = MagicMock(side_effect=_defer_the_open)
684
685 with (
686 patch(f"{_MODULE}._DEVICE_OPEN_TIMEOUT_SECONDS", 0.01),
687 pytest.raises(TimeoutError),
688 ):
689 await bridge._get_pa_stream()
690
691 opening.set_result(stream)
692 await _settle()
693
694 stream.close.assert_called_once_with()
695
696
697# --- Not giving up: an ordinary end of stream keeps the bridge in its session ---
698
699
700@pytest.mark.parametrize("backend", BACKENDS)
701async def test_a_stream_that_simply_ended_keeps_its_place_in_the_session(backend: str) -> None:
702 """
703 The stop sentinel is the normal end of a stream, not a device failure.
704
705 Leaving on it would drop the player out of its group at the end of every
706 single track.
707 """
708 bridge = _make_bridge(backend=backend)
709 stream = _make_device_stream()
710 _queue_chunks(bridge, _pcm_chunk(1))
711
712 with (
713 _patch_device_open(bridge, stream),
714 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
715 ):
716 await _run_writer(bridge)
717
718 stream.write.assert_called_once_with(_pcm_chunk(1))
719 assert bridge._is_streaming is False
720 leave.assert_not_called()
721
722
723async def test_a_writer_whose_stream_already_stopped_keeps_its_place_in_the_session() -> None:
724 """
725 A writer that finds streaming already switched off simply stops feeding.
726
727 That is a teardown racing the queue, not a device that failed, so there is
728 nothing for Sendspin to hear about.
729 """
730 bridge = _make_bridge()
731 stream = _make_device_stream()
732 _queue_chunks(bridge, _pcm_chunk(1))
733 bridge._is_streaming = False
734
735 with (
736 _patch_device_open(bridge, stream),
737 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
738 ):
739 await _run_writer(bridge)
740
741 stream.write.assert_not_called()
742 leave.assert_not_called()
743
744
745@needs_portaudio
746async def test_a_stopped_stream_keeps_its_place_in_the_session() -> None:
747 """A writer taken down by the teardown path has nothing to report to Sendspin."""
748 bridge = _make_bridge(backend="sounddevice")
749 stream = _make_device_stream()
750
751 with (
752 _patch_device_open(bridge, stream),
753 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
754 ):
755 task = _install_writer(bridge, bridge._audio_writer())
756 await asyncio.sleep(0) # let the writer park on the empty queue
757 await bridge._stop_streaming()
758
759 assert task.done()
760 assert not task.cancelled()
761 assert bridge._writer_task is None
762 assert bridge._is_streaming is False
763 stream.close.assert_called_once_with()
764 leave.assert_not_called()
765
766
767async def test_a_cancelled_writer_keeps_its_place_in_the_session() -> None:
768 """
769 A writer cancelled out from under a healthy sink is not a device failure.
770
771 Cancellation is how the PA writer is torn down between streams, so leaving
772 on it would un-group the player on every track change.
773 """
774 bridge = _make_bridge()
775 stream = _make_device_stream()
776
777 with (
778 _patch_device_open(bridge, stream),
779 patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave,
780 ):
781 task = _install_writer(bridge, bridge._audio_writer())
782 await asyncio.sleep(0) # let the writer park on the empty queue
783 task.cancel()
784 with suppress(asyncio.CancelledError):
785 await task
786
787 assert bridge._is_streaming is False
788 leave.assert_not_called()
789
790
791# --- Leaving the Sendspin session ----------------------------------------------
792
793
794async def test_leaving_the_session_quiesces_the_client_without_unregistering() -> None:
795 """The bridge leaves a shared group (or stops a solo one) but stays registered."""
796 bridge = _make_bridge()
797 client = MagicMock()
798 client.quiesce_to_solo_stopped = AsyncMock()
799 bridge._sendspin_client = client
800
801 await bridge._leave_sendspin_session()
802
803 client.quiesce_to_solo_stopped.assert_awaited_once_with()
804 # staying registered is what keeps the player around to be grouped again
805 cast("MagicMock", bridge.sendspin_server).remove_client.assert_not_called()
806
807
808async def test_leaving_the_session_without_a_client_is_a_silent_noop() -> None:
809 """
810 Giving up before registration completed has no session to leave.
811
812 The call has to return without touching anything: swallowing an error from
813 an absent client would look identical from the outside, so the absence of a
814 complaint is what distinguishes the two.
815 """
816 bridge = _make_bridge()
817 assert bridge._sendspin_client is None
818
819 await bridge._leave_sendspin_session()
820
821 cast("MagicMock", bridge.logger).warning.assert_not_called()
822
823
824async def test_a_refused_quiesce_is_reported_and_contained() -> None:
825 """
826 A session that will not release the bridge is reported, never raised.
827
828 The caller is a writer already on its way out, so an exception here would be
829 swallowed by the task machinery and the silence would go unexplained.
830 """
831 bridge = _make_bridge()
832 client = MagicMock()
833 client.quiesce_to_solo_stopped = AsyncMock(side_effect=RuntimeError("no active session"))
834 bridge._sendspin_client = client
835
836 await bridge._leave_sendspin_session()
837
838 cast("MagicMock", bridge.logger).warning.assert_called_once()
839
840
841# --- Writer supersession: a replaced writer speaks for a stream that is gone ---
842
843
844async def test_a_superseded_writer_cannot_give_up_on_its_replacement() -> None:
845 """
846 A writer a newer stream replaced gives up on nothing.
847
848 Leaving the session would stop the group for a stream that is playing
849 perfectly well, and the flag it would clear now belongs to that stream.
850 """
851 bridge = _make_bridge()
852 bridge._writer_task = MagicMock() # a newer stream's writer
853
854 with patch.object(bridge, "_leave_sendspin_session", MagicMock()) as leave:
855 bridge._abandon_streaming()
856
857 assert bridge._is_streaming is True
858 leave.assert_not_called()
859
860
861@pytest.mark.parametrize("backend", BACKENDS)
862async def test_a_superseded_writer_does_not_silence_its_replacement(backend: str) -> None:
863 """
864 A writer left behind by a newer stream keeps its hands off the shared state.
865
866 The streaming flag now belongs to its replacement, which reads it on every
867 chunk; clearing it would silence a stream that is playing perfectly well.
868 Its own stream is still its own to close.
869 """
870 bridge = _make_bridge(backend=backend)
871 stream = _make_device_stream()
872 newer_writer = MagicMock()
873 bridge._writer_task = newer_writer
874 bridge._write_queue.put_nowait(None)
875
876 with _patch_device_open(bridge, stream):
877 await bridge._audio_writer()
878
879 assert bridge._is_streaming is True
880 assert bridge._writer_task is newer_writer
881 stream.close.assert_called_once_with()
882
883
884async def test_a_stale_teardown_leaves_the_newer_writer_running() -> None:
885 """
886 A teardown scheduled for a replaced writer must not stop the one that replaced it.
887
888 The stream end that scheduled it belongs to the previous stream; acting on it
889 would stop the audio the new stream has only just started.
890 """
891 bridge = _make_bridge()
892 newer_writer = _install_writer(bridge, _never_returns())
893
894 await bridge._stop_streaming_locked(MagicMock())
895
896 assert not newer_writer.done()
897 assert not newer_writer.cancelled()
898 assert bridge._writer_task is newer_writer
899 assert bridge._is_streaming is True
900
901 newer_writer.cancel()
902 with suppress(asyncio.CancelledError):
903 await newer_writer
904
905
906async def test_a_stream_end_tears_down_the_writer_that_was_running() -> None:
907 """
908 The teardown a stream end schedules is bound to the writer running at the time.
909
910 That binding is what lets a later teardown recognise itself as stale, so it
911 has to name a writer â a teardown naming none matches nothing and quietly
912 stops tearing anything down at all.
913 """
914 bridge = _make_bridge()
915 writer = _install_writer(bridge, _never_returns())
916 await _settle()
917
918 bridge._on_bridge_stream_end()
919 await _settle()
920
921 assert writer.cancelled()
922 assert bridge._writer_task is None
923 assert bridge._is_streaming is False
924
925
926async def test_a_stream_end_behind_a_running_teardown_leaves_the_newer_writer_alone() -> None:
927 """
928 A stream end whose teardown turns out to be stale leaves the streaming flag alone.
929
930 Every chunk is checked against that flag, so clearing it on behalf of a
931 stream that has already finished would silence the one that replaced it.
932 """
933 bridge = _make_bridge()
934 old_writer = _install_writer(bridge, _never_returns())
935 await _settle()
936
937 await bridge._lock.acquire() # a teardown is already in progress
938 bridge._on_bridge_stream_end()
939 await _settle()
940 # a new stream installs its writer before the queued teardown gets the lock
941 newer_writer = _install_writer(bridge, _never_returns())
942 bridge._lock.release()
943 await _settle()
944
945 assert bridge._is_streaming is True
946 assert bridge._writer_task is newer_writer
947 assert not newer_writer.cancelled()
948
949 for task in (old_writer, newer_writer):
950 task.cancel()
951 with suppress(asyncio.CancelledError):
952 await task
953
954
955async def test_a_timed_out_stop_cancels_the_writer_it_captured() -> None:
956 """
957 A writer that will not stop on its own is cancelled â and only that writer.
958
959 The cancel comes after a wait, by which time a new stream may already have
960 installed its own writer; cancelling that one would silence a stream that
961 has only just started.
962 """
963 bridge = _make_bridge(backend="sounddevice")
964 newer_writer = MagicMock()
965 stubborn = _install_writer(bridge, _never_returns())
966 await _settle()
967
968 async def _install_the_replacement() -> None:
969 await asyncio.sleep(0)
970 bridge._writer_task = newer_writer
971
972 with patch(f"{_MODULE}._WRITER_STOP_TIMEOUT_SECONDS", 0.01):
973 replacing = asyncio.create_task(_install_the_replacement())
974 await bridge._stop_streaming()
975 await replacing
976
977 assert stubborn.cancelled()
978 assert bridge._writer_task is newer_writer
979 newer_writer.cancel.assert_not_called()
980
981
982async def test_a_teardown_does_not_disown_a_writer_installed_while_it_waits() -> None:
983 """
984 A writer installed while the teardown waits for its own keeps its registration.
985
986 Stopping a writer takes an await, and a new stream can start within it. The
987 registration then belongs to the replacement, and clearing it would leave
988 that writer running unowned â invisible to the next teardown, and splitting
989 the following stream's chunks with the writer started alongside it.
990 """
991 bridge = _make_bridge()
992 newer_writer = MagicMock()
993
994 async def _hand_over_when_stopped() -> None:
995 try:
996 await asyncio.Event().wait()
997 except asyncio.CancelledError:
998 bridge._writer_task = newer_writer
999 raise
1000
1001 old_writer = _install_writer(bridge, _hand_over_when_stopped())
1002 await asyncio.sleep(0) # let the writer park before it is torn down
1003
1004 await bridge._stop_streaming()
1005
1006 assert old_writer.cancelled()
1007 assert bridge._writer_task is newer_writer
1008