/
/
1"""
2Unit tests for the Spotify Soloist playback backend.
3
4The backend runs one continuous soloist session, feeds it one track ahead and
5splits the captured PCM into per-item streams. These tests lock down the pure
6logic around that: lead-silence trimming, the item channels and where the
7session cuts between them, event handling, the crossfade handed to the engine,
8feeding the follower, completeness validation, paired-session adoption and
9setup. No real process or PulseAudio is involved.
10"""
11
12from __future__ import annotations
13
14import asyncio
15import os
16import time
17from collections.abc import AsyncGenerator, Callable, Iterator
18from contextlib import contextmanager, suppress
19from pathlib import Path
20from typing import Any, cast
21from unittest.mock import AsyncMock, MagicMock, patch
22
23import pytest
24from music_assistant_models.enums import MediaType
25from music_assistant_models.errors import AudioError, LoginFailed
26
27from music_assistant.helpers.pulse_capture import CAPTURE_SAMPLE_RATE
28from music_assistant.models.music_provider import ProviderStreamLimitError
29from music_assistant.providers.spotify.backends import soloist as soloist_backend
30from music_assistant.providers.spotify.backends.soloist import (
31 _BYTES_PER_SECOND,
32 _FRAME_BYTES,
33 _IDLE_TIMEOUT_S,
34 _MAX_APP_PAUSE_RESUMES,
35 _MAX_LEAD_TRIM_S,
36 _READ_CHUNK_SIZE,
37 SoloistAppControl,
38 SoloistAppControlError,
39 SoloistBackend,
40 _CaptureShaper,
41 _ItemAudio,
42 _SoloistSession,
43 _trim_lead_silence,
44)
45from music_assistant.providers.spotify.constants import (
46 CONF_SOLOIST_API_KEY,
47 CONF_SOLOIST_CONSENT,
48 CONF_SOLOIST_SESSION_DIR,
49 SOLOIST_DATA_DIR_NAME,
50 SOLOIST_DEVICE_NAME,
51)
52from music_assistant.providers.spotify.helpers import soloist_session_present
53from music_assistant.providers.spotify.provider import SpotifyProvider
54from music_assistant.providers.spotify_connect.provider import DEFAULT_PUBLISH_NAME
55from music_assistant.providers.spotify_connect.soloist.runtime import (
56 WS_ADDR_FILE,
57 WS_PORT_FILE,
58 SoloistAuthState,
59 SoloistDeviceChanged,
60 SoloistEntity,
61 SoloistError,
62 SoloistEvent,
63 SoloistOptionsChanged,
64 SoloistPlaybackOptions,
65 SoloistPlaybackState,
66 SoloistPosition,
67 SoloistTrackChanged,
68 SoloistVolumeChanged,
69)
70
71TRACK_A = "spotify:track:aaa"
72TRACK_B = "spotify:track:bbb"
73
74
75def test_trim_drops_an_all_zero_chunk_within_the_bound() -> None:
76 """A pure-silence chunk inside the trim budget is dropped entirely."""
77 chunk = b"\x00" * 1024
78 trimmed, skipped = _trim_lead_silence(chunk, 0)
79 assert trimmed == b""
80 assert skipped == 1024
81
82
83def test_trim_keeps_frame_alignment_when_audio_starts_mid_chunk() -> None:
84 """Audio starting mid-chunk is cut on a sample-frame boundary."""
85 # audio starts one byte into the third frame: the trim must keep that frame whole
86 chunk = b"\x00" * (_FRAME_BYTES * 2 + 1) + b"\x01" * 64
87 trimmed, skipped = _trim_lead_silence(chunk, 0)
88 assert skipped == _FRAME_BYTES * 2
89 assert len(trimmed) % _FRAME_BYTES == 1 # the partial frame's remainder is preserved
90 assert trimmed.endswith(b"\x01" * 64)
91
92
93def test_trim_passes_silence_through_once_the_bound_is_exceeded() -> None:
94 """Beyond the trim budget, silence is genuine content and is delivered."""
95 chunk = b"\x00" * 1024
96 trimmed, skipped = _trim_lead_silence(chunk, int(_MAX_LEAD_TRIM_S * _BYTES_PER_SECOND))
97 assert trimmed == chunk
98 assert skipped == 0
99
100
101def test_seek_is_confirmed_only_within_tolerance(tmp_path: Path) -> None:
102 """A position report confirms a seek only once it reaches the tolerance window."""
103 item = _make_item(tmp_path, TRACK_A)
104 item.seek_target_ms = 60_000
105 item.observe_position(50_000)
106 assert not item.seek_confirmed.is_set()
107 item.observe_position(58_500)
108 assert item.seek_confirmed.is_set()
109
110
111def test_small_seek_target_is_not_confirmed_by_a_pre_seek_zero_report(tmp_path: Path) -> None:
112 """A position-0 report before the seek lands cannot confirm a small target."""
113 item = _make_item(tmp_path, TRACK_A)
114 item.seek_target_ms = 1_500
115 item.observe_position(0)
116 assert not item.seek_confirmed.is_set()
117 item.observe_position(1_500)
118 assert item.seek_confirmed.is_set()
119
120
121def test_position_never_regresses_and_stops_at_the_cut(tmp_path: Path) -> None:
122 """The furthest position is kept, and reports after the cut belong to the next item."""
123 item = _make_item(tmp_path, TRACK_A)
124 item.observe_position(120_000)
125 # the engine's stop/idle snapshot at the end of an item reports position 0
126 item.observe_position(0)
127 assert item.last_position_ms == 120_000
128 item.close()
129 item.observe_position(5_000)
130 assert item.last_position_ms == 120_000
131
132
133async def test_item_stream_ends_where_the_session_moves_on(tmp_path: Path) -> None:
134 """An item's audio ends at the track change, and the next item's begins there."""
135 session = _make_session(tmp_path)
136 item_a = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
137 session._current = item_a
138 item_a.started.set()
139 item_a.claim()
140 item_a.write(b"a" * 16)
141 await session._observe_current(TRACK_B, 200_000)
142 item_a.write(b"late" * 4) # written after the cut: goes nowhere
143 chunks = [chunk async for chunk in item_a.read()]
144 assert b"".join(chunks) == b"a" * 16
145 # the next item exists, carries the duration and now receives the audio
146 item_b = session._items[TRACK_B]
147 assert session.current is item_b
148 assert item_b.duration_ms == 200_000
149
150
151async def test_the_engines_restored_state_does_not_cut_a_pending_item(
152 tmp_path: Path,
153) -> None:
154 """A daemon reports the item it restored before playing ours; that is not a boundary."""
155 session = _make_session(tmp_path)
156 requested = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
157 session._current = requested
158 requested.claim()
159 # the engine announces the state it came up with, which is someone else's item
160 await session._observe_current("spotify:track:restored", 152_000)
161 # closing our item here would end its stream before it delivered anything
162 assert requested._closed is False
163 assert requested.started.is_set() is False
164 # ... and the restored item is never offered as an item's audio
165 assert session.item_for("spotify:track:restored") is None
166 # then ours starts for real, and picks up from there
167 await session._observe_current(TRACK_A, 200_000)
168 assert session.current is requested
169 assert requested.started.is_set() is True
170 requested.write(b"\x01" * 32)
171 requested.close()
172 assert b"".join([chunk async for chunk in requested.read()]) == b"\x01" * 32
173
174
175async def test_leaving_the_engines_restored_item_is_not_a_takeover(tmp_path: Path) -> None:
176 """The restored item is part-way through a track, and we are about to leave it."""
177 session = _make_session(tmp_path)
178 # as _play leaves it: the channel exists, its stream is not reading it yet
179 requested = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
180 session._current = requested
181 await session._observe_current("spotify:track:restored", 152_000)
182 restored = session.current
183 assert restored is not None
184 restored.observe_position(20_000)
185
186 # our own play() lands and the engine leaves the restored item for ours
187 await session._observe_current(TRACK_A, 200_000)
188 assert session.usable is True
189 assert session.current is requested
190
191
192async def test_audio_read_before_the_stream_opens_is_kept(tmp_path: Path) -> None:
193 """Audio captured before an item's stream opens is buffered, not dropped."""
194 session = _make_session(tmp_path)
195 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
196 session._current = item
197 item.write(b"head" * 8)
198 item.claim()
199 item.close()
200 chunks = [chunk async for chunk in item.read()]
201 assert b"".join(chunks) == b"head" * 8
202
203
204async def test_a_channel_is_only_ever_served_once(tmp_path: Path) -> None:
205 """A consumed channel cannot be replayed, so the item needs a fresh session."""
206 session = _make_session(tmp_path)
207 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
208 item.started.set()
209 assert session.item_for(TRACK_A) is item
210 item.claim()
211 item.close()
212 item.release()
213 # this is what a queue holding the same track twice, or repeat wrapping back
214 # to the top, asks for: it must not be handed a drained channel
215 assert session.item_for(TRACK_A) is None
216
217
218async def test_an_abandoned_channel_cannot_be_continued(tmp_path: Path) -> None:
219 """A stream abandoned mid-item cannot resume where it left off either."""
220 session = _make_session(tmp_path)
221 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
222 item.started.set()
223 item.claim()
224 item.release()
225 assert session.item_for(TRACK_A) is None
226
227
228async def test_a_stuck_item_fails_instead_of_streaming_forever(tmp_path: Path) -> None:
229 """An item that runs far past its duration without a track change fails."""
230 session = _make_session(tmp_path)
231 item = _ItemAudio(TRACK_A, session)
232 item.duration_ms = 1_000
233 item.claim()
234 limit = item._overrun_limit()
235 assert limit is not None
236 item.write(b"\x01" * (limit + _FRAME_BYTES))
237 with pytest.raises(AudioError, match="never moved on"):
238 async for _ in item.read():
239 pass
240
241
242async def test_the_first_logged_out_snapshot_is_not_a_lost_pairing(tmp_path: Path) -> None:
243 """A daemon reports logged_in=False until it has restored its session."""
244 session = _make_session(tmp_path)
245 session._logged_in = None
246 session._was_active = False
247 await session._handle_event(_auth_event(logged_in=False, is_active=False))
248 # failing here would break every playback on a perfectly good pairing
249 assert session.usable is True
250 await session._handle_event(_auth_event(logged_in=True, is_active=False))
251 assert session.usable is True
252
253
254async def test_losing_an_established_login_fails_the_session(tmp_path: Path) -> None:
255 """A login that goes away mid-session is real, and ends the session."""
256 session = _make_session(tmp_path)
257 await session._handle_event(_auth_event(logged_in=True))
258 await session._handle_event(_auth_event(logged_in=False))
259 assert session.usable is False
260 assert session._error == "the session was logged out"
261
262
263async def test_buffering_gates_the_sink_once_demand_started(tmp_path: Path) -> None:
264 """Once PCM demand started, playing runs the sink and buffering suspends it again."""
265 session = _make_session(tmp_path)
266 session._demand_started = True
267 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
268 session._pending.append(TRACK_B)
269 sink = _sink_of(session)
270 # the sink is created suspended, so there is nothing to suspend yet
271 await session._handle_event(_playback_event("buffering"))
272 sink.suspend.assert_not_awaited()
273 await session._handle_event(_playback_event("playing"))
274 sink.resume.assert_awaited_once()
275 assert session._current is not None
276 assert session._current.playing_seen is True
277 # the engine stalling on a rebuffer keeps that silence out of the PCM
278 await session._handle_event(_playback_event("buffering"))
279 sink.suspend.assert_awaited_once()
280
281
282async def test_sink_is_not_gated_before_demand_started(tmp_path: Path) -> None:
283 """Buffering/playing before PCM demand leave the (still suspended) sink alone."""
284 session = _make_session(tmp_path)
285 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
286 sink = _sink_of(session)
287 await session._handle_event(_playback_event("buffering"))
288 await session._handle_event(_playback_event("playing"))
289 sink.suspend.assert_not_awaited()
290 sink.resume.assert_not_awaited()
291 # the status is recorded either way, so session start can decide when to resume
292 assert session._current.status == "playing"
293
294
295@pytest.mark.parametrize("end_status", ["stopped", "idle", "paused"])
296async def test_the_last_item_is_drained_rather_than_cut(tmp_path: Path, end_status: str) -> None:
297 """However the engine reports the end of a run, the last item drains and closes."""
298 session = _make_session(tmp_path)
299 session._demand_started = True
300 session._sink_running = True
301 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
302 item.duration_ms = 1_000
303 item.last_position_ms = 1_000
304 one_second = 1_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
305 sink = _sink_of(session)
306 await session._handle_event(_playback_event(end_status, position_ms=1_000))
307 # the sink stays open for now, so audio still in the FIFO can arrive...
308 sink.suspend.assert_not_awaited()
309 assert item.draining is True
310 assert item._closed is False
311 # ... but only that item's own audio is taken, never the padding silence the
312 # sink keeps rendering afterwards
313 item.write(b"\x01" * one_second)
314 item.write(b"\x00" * 4096)
315 assert item.buffered == one_second
316 await _wait_for(lambda: item._closed)
317 sink.suspend.assert_awaited_once()
318
319
320async def test_an_app_pause_midway_through_the_last_item_is_not_the_end(
321 tmp_path: Path,
322) -> None:
323 """Pausing in the Spotify app halfway through the last track must not truncate it."""
324 session = _make_session(tmp_path)
325 session._demand_started = True
326 session._sink_running = True
327 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
328 item.duration_ms = 200_000
329 item.last_position_ms = 90_000
330 await session._handle_event(_playback_event("paused", position_ms=90_000))
331 assert item.draining is False
332 assert item._closed is False
333 # treated as interference instead: the sink is gated and playback resumed
334 _sink_of(session).suspend.assert_awaited_once()
335 _client_of(session).resume.assert_awaited_once()
336
337
338async def test_a_resumed_item_cancels_its_tail_drain(tmp_path: Path) -> None:
339 """An armed drain is undone when the engine turns out to have been rebuffering."""
340 session = _make_session(tmp_path)
341 session._demand_started = True
342 session._sink_running = True
343 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
344 item.duration_ms = 200_000
345 item.last_position_ms = 199_000
346 await session._handle_event(_playback_event("stopped", position_ms=199_000))
347 armed = item.draining
348 await session._handle_event(_playback_event("playing", position_ms=199_500))
349 assert armed is True
350 assert item.draining is False
351 assert item._closed is False
352 assert item.drain_task is None
353
354
355async def test_the_cushion_is_capped_by_suspending_the_sink(tmp_path: Path) -> None:
356 """Undelivered audio is handed back as backpressure rather than piling up."""
357 session = _make_session(tmp_path)
358 session._demand_started = True
359 session._sink_running = True
360 session._engine_playing = True
361 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
362 item.claim()
363 sink = _sink_of(session)
364 await session._apply_sink_state()
365 sink.suspend.assert_not_awaited()
366 # the engine has run this far ahead of what the player has taken
367 item.write(b"\x01" * int((soloist_backend._MAX_RETAINED_S + 1) * _BYTES_PER_SECOND))
368 await session._apply_sink_state()
369 sink.suspend.assert_awaited_once()
370 assert session._backpressured is True
371 # and it comes back once the player has drained enough of it
372 item._buffered = int(soloist_backend._RESUME_RETAINED_S * _BYTES_PER_SECOND) - 1
373 await session._apply_sink_state()
374 sink.resume.assert_awaited_once()
375 assert session._backpressured is False
376
377
378async def test_a_pause_with_more_queued_suspends_the_sink(tmp_path: Path) -> None:
379 """A pause while another item is queued behind is ordinary interference, not the end."""
380 session = _make_session(tmp_path)
381 session._demand_started = True
382 session._sink_running = True
383 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
384 session._pending.append(TRACK_B)
385 await session._handle_event(_playback_event("paused"))
386 _sink_of(session).suspend.assert_awaited_once()
387
388
389async def test_nothing_is_sent_before_the_websocket_is_up(
390 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
391) -> None:
392 """Commands travel over the events socket: a published endpoint is not enough."""
393 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
394 session = _make_session(tmp_path)
395 client = _client_of(session)
396 # the endpoint file exists, but the events task has not connected yet
397 client.connected = False
398 endpoint_published = asyncio.Event()
399 endpoint_published.set()
400 with pytest.raises(AudioError, match="did not connect and log in"):
401 await session._play(TRACK_A, 0, endpoint_published)
402 client.activate.assert_not_awaited()
403 client.play.assert_not_awaited()
404
405
406async def test_nothing_is_sent_before_the_engine_has_logged_in(
407 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
408) -> None:
409 """The engine drops commands sent before it has restored its session."""
410 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
411 session = _make_session(tmp_path)
412 client = _client_of(session)
413 client.connected = True
414 # connected, but the engine has not announced its login yet
415 session._logged_in = None
416 endpoint_published = asyncio.Event()
417 endpoint_published.set()
418 with pytest.raises(AudioError, match="did not connect and log in"):
419 await session._play(TRACK_A, 0, endpoint_published)
420 client.activate.assert_not_awaited()
421 client.play.assert_not_awaited()
422
423
424async def test_startup_activates_before_it_plays(
425 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
426) -> None:
427 """A fresh daemon has to become the active device before it is told to play."""
428 session = _make_session(tmp_path)
429 client = _client_of(session)
430 client.connected = True
431 monkeypatch.setattr(session, "_await_item_ready", AsyncMock())
432 endpoint_published = asyncio.Event()
433 endpoint_published.set()
434 item = await session._play(TRACK_A, 0, endpoint_published)
435 assert item.uri == TRACK_A
436 client.activate.assert_awaited_once_with(await_result=True)
437 client.play.assert_awaited_once_with(TRACK_A)
438
439
440async def test_a_takeover_between_activate_and_play_stops_the_start(tmp_path: Path) -> None:
441 """Playing here would claim the device straight back off wherever the user moved to."""
442 session = _make_session(tmp_path)
443 session._was_active = False
444 client = _client_of(session)
445
446 async def _take_over(*_args: Any, **_kwargs: Any) -> None:
447 session._observe_active_device(is_active=False)
448
449 client.set_repeat_track.side_effect = _take_over
450 ready = asyncio.Event()
451 ready.set()
452 with pytest.raises(SoloistAppControlError):
453 await session._play(TRACK_A, 0, ready)
454 client.play.assert_not_awaited()
455
456
457async def test_a_refused_start_command_reports_soloist(tmp_path: Path) -> None:
458 """A dropped start command surfaces as a Soloist error, not a raw client one."""
459 session = _make_session(tmp_path)
460 client = _client_of(session)
461 client.connected = True
462 client.activate.side_effect = SoloistError("websocket is not connected")
463 endpoint_published = asyncio.Event()
464 endpoint_published.set()
465 with pytest.raises(AudioError, match="Spotify Soloist would not start"):
466 await session._play(TRACK_A, 0, endpoint_published)
467
468
469async def test_the_engine_is_told_not_to_shuffle_or_repeat(
470 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
471) -> None:
472 """MA owns the order, and a repeating engine would never reach the item fed behind."""
473 session = _make_session(tmp_path)
474 client = _client_of(session)
475 client.connected = True
476 monkeypatch.setattr(session, "_await_item_ready", AsyncMock())
477 endpoint_published = asyncio.Event()
478 endpoint_published.set()
479 await session._play(TRACK_A, 0, endpoint_published)
480 client.set_shuffle.assert_awaited_once_with(False)
481 client.set_repeat_context.assert_awaited_once_with(False)
482 client.set_repeat_track.assert_awaited_once_with(False)
483
484
485async def test_repeat_turned_on_from_the_app_is_pinned_back_off(tmp_path: Path) -> None:
486 """Repeat enabled in the Spotify app is undone before it can loop the item."""
487 session = _make_session(tmp_path)
488 await session._handle_event(
489 SoloistEvent(
490 type="options_changed",
491 data=SoloistOptionsChanged(
492 options=SoloistPlaybackOptions(shuffle=True, repeat="track")
493 ),
494 raw={},
495 )
496 )
497 client = _client_of(session)
498 client.set_shuffle.assert_awaited_once_with(False)
499 client.set_repeat_track.assert_awaited_once_with(False)
500 client.set_repeat_context.assert_awaited_once_with(False)
501 # options that are already off are left alone
502 client.set_shuffle.reset_mock()
503 await session._handle_event(
504 SoloistEvent(
505 type="options_changed",
506 data=SoloistOptionsChanged(options=SoloistPlaybackOptions()),
507 raw={},
508 )
509 )
510 client.set_shuffle.assert_not_awaited()
511
512
513async def test_a_busy_data_directory_is_reported_as_such(tmp_path: Path) -> None:
514 """A daemon left over from an earlier run is named, not reported as a generic failure."""
515 session = _make_session(tmp_path)
516 # the daemon's own parting complaint, which is all it gives (it exits with 1)
517 session._data_dir_busy = True
518 with pytest.raises(AudioError, match="Another Spotify Soloist session is still running"):
519 session._raise_startup_error("exited before playback started", TRACK_A)
520
521
522async def test_the_busy_marker_is_picked_up_from_the_daemon_output(tmp_path: Path) -> None:
523 """The marker is read off the daemon's stdout, with the API key still redacted."""
524 session = _make_session(tmp_path)
525 proc = MagicMock()
526 lines = [
527 'Error: another session is running for data directory "/data/x/soloist-data".',
528 "Stop the running session before starting soloist again.",
529 ]
530
531 async def _iter_stdout() -> AsyncGenerator[str]:
532 for line in lines:
533 yield line
534
535 proc.iter_stdout = _iter_stdout
536 await session._log_output(proc)
537 assert session._data_dir_busy is True
538
539
540async def test_a_pairing_that_never_logs_in_routes_through_setup(
541 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
542) -> None:
543 """A session that cannot log in sends the user back to setup, not a per-track error."""
544 session = _make_session(tmp_path)
545 unload_with_error = MagicMock()
546 monkeypatch.setattr(session.backend.provider, "unload_with_error", unload_with_error)
547 await session._handle_event(_auth_event(logged_in=False))
548 with pytest.raises(LoginFailed) as err:
549 session._raise_startup_error("timed out waiting for playback to start", TRACK_A)
550 assert err.value.translation_key == "soloist_pairing_required"
551 # ... and the provider is taken out of service, so the user is asked to redo setup
552 unload_with_error.assert_called_once()
553
554
555async def test_a_login_that_never_happened_is_not_confused_with_another_failure(
556 tmp_path: Path,
557) -> None:
558 """An unrelated failure keeps its own message even before any login was reported."""
559 session = _make_session(tmp_path)
560 session._fail("the capture sink was lost mid-stream")
561 with pytest.raises(AudioError, match="capture sink was lost"):
562 session._raise_startup_error("exited before playback started", TRACK_A)
563
564
565def test_a_seeked_item_only_expects_what_is_left_of_it(tmp_path: Path) -> None:
566 """A seeked item delivers the remainder, so its targets are based on that."""
567 session = _make_session(tmp_path)
568 item = _ItemAudio(TRACK_A, session)
569 item.duration_ms = 200_000
570 full = 200_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
571 assert item._duration_bytes() == full
572 item.seek_target_ms = 150_000
573 remainder = 50_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
574 assert item._duration_bytes() == remainder
575 # so the tail drain has a target it can actually reach
576 item.start_tail_drain()
577 item.write(b"\x01" * remainder)
578 assert item.tail_complete is True
579 # and the padding silence after it is refused
580 item.write(b"\x00" * 4096)
581 assert item.buffered == remainder
582
583
584def test_the_lead_trim_never_exceeds_its_budget() -> None:
585 """Silence beyond the budget is content, including where audio starts mid-chunk."""
586 budget = int(_MAX_LEAD_TRIM_S * _BYTES_PER_SECOND)
587 # already at the budget, with a chunk whose silence runs well past it
588 chunk = b"\x00" * 4096 + b"\x01" * 64
589 trimmed, skipped = _trim_lead_silence(chunk, budget - _FRAME_BYTES)
590 assert skipped == _FRAME_BYTES
591 assert len(trimmed) == len(chunk) - _FRAME_BYTES
592
593
594async def test_a_dying_log_reader_fails_the_session(tmp_path: Path) -> None:
595 """Nothing else drains the daemon's stdout, so a dead reader must not go unnoticed."""
596 session = _make_session(tmp_path)
597
598 async def _boom() -> None:
599 raise RuntimeError("reader blew up")
600
601 session._log_task = asyncio.create_task(_boom())
602 session._log_task.add_done_callback(session._task_done)
603 await asyncio.sleep(0)
604 await _wait_for(lambda: not session.usable)
605 assert session._error is not None
606 assert "reader blew up" in session._error
607
608
609async def test_feeding_never_replaces_a_channel_already_in_use(tmp_path: Path) -> None:
610 """If the engine reaches the fed item first, its live channel must survive."""
611 session = _make_session(tmp_path, queue_id="player1")
612 streamdetails = MagicMock()
613 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
614 queues = _queues_of(session)
615 queues.get.return_value = MagicMock(current_index=0)
616 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
617 queues.get_next_item.return_value = _queue_item(TRACK_B)
618
619 async def _engine_gets_there_first(_uri: str, **_kwargs: Any) -> None:
620 # the events task advances to the fed item while the command is in flight
621 await session._observe_current(TRACK_B, 200_000)
622
623 _client_of(session).add_to_queue.side_effect = _engine_gets_there_first
624 await session.feed_after(streamdetails, TRACK_A)
625 live = session.current
626 assert live is not None
627 assert live.uri == TRACK_B
628 # the channel the reader is writing to is the one a stream will be handed
629 assert session._items[TRACK_B] is live
630 assert session.item_for(TRACK_B) is live
631 # and it is not queued as pending, because it already started
632 assert session.has_pending is False
633
634
635async def test_seeking_the_playing_item_restarts_the_session(
636 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
637) -> None:
638 """
639 A seek re-opens the item that is playing, which is a restart of the session.
640
641 A realtime source has not captured anything past the play position, so any
642 forward seek lands outside the buffer and comes back here.
643 """
644 backend = _make_backend(tmp_path)
645 backend._server = MagicMock()
646 backend._binary = Path("/nonexistent/soloist")
647 session = _SoloistSession(backend, "player1")
648 backend._session = session
649 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
650 item.started.set()
651 # its own stream is still attached when the seek re-opens it
652 item.claim()
653 stopped = AsyncMock()
654 monkeypatch.setattr(session, "stop", stopped)
655 _install_fake_binary_manager(monkeypatch)
656 monkeypatch.setattr(
657 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
658 )
659 with pytest.raises(AudioError, match="spawn"):
660 await backend._acquire(TRACK_A, 90, "player1")
661 stopped.assert_awaited_once()
662
663
664@pytest.mark.parametrize(
665 ("requested", "other_queue"),
666 [
667 # another player, whatever it asks for - including the very track this
668 # session is in the middle of delivering
669 pytest.param(TRACK_B, "player2", id="other_player"),
670 pytest.param(TRACK_A, "player2", id="other_player_same_track"),
671 # an early fetch across a boundary this session does not drive, such as a
672 # podcast episode or audiobook chapter
673 pytest.param(TRACK_B, "player1", id="unstitched_boundary"),
674 ],
675)
676async def test_a_session_in_use_is_never_cut_short(
677 tmp_path: Path, requested: str, other_queue: str
678) -> None:
679 """
680 An item the session cannot serve must not stop one it is still delivering.
681
682 Reported as capacity, so a speculative prepare gives up softly.
683 """
684 backend = _make_backend(tmp_path)
685 backend._server = MagicMock()
686 backend._binary = Path("/nonexistent/soloist")
687 session = _SoloistSession(backend, "player1")
688 backend._session = session
689 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
690 item.started.set()
691 item.claim()
692 # the session really is playing TRACK_A, so a same-track request from another
693 # player cannot be mistaken for a seek
694 session._current = item
695 with pytest.raises(ProviderStreamLimitError) as err:
696 await backend._acquire(requested, 0, other_queue)
697 # a stream-limit error so the item is not marked unplayable, but the message
698 # is about the session, not the provider's source-stream budget
699 assert err.value.limit == 1
700 assert err.value.translation_key == "soloist_session_busy"
701 # the session that was playing is untouched
702 assert backend._session is session
703 assert session.usable is True
704
705
706async def test_a_session_nobody_reads_is_replaced_for_another_item(
707 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
708) -> None:
709 """Once the other item has been released, the same request gets the session."""
710 backend = _make_backend(tmp_path)
711 backend._server = MagicMock()
712 backend._binary = Path("/nonexistent/soloist")
713 session = _SoloistSession(backend, "player1")
714 backend._session = session
715 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
716 item.started.set()
717 item.claim()
718 item.close()
719 item.release()
720 _install_fake_binary_manager(monkeypatch)
721 monkeypatch.setattr(
722 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
723 )
724 monkeypatch.setattr(session, "stop", AsyncMock())
725 with pytest.raises(AudioError, match="spawn"):
726 await backend._acquire(TRACK_B, 0, "player1")
727
728
729async def test_a_replacement_waits_for_the_old_daemon_to_be_gone(
730 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
731) -> None:
732 """The engine refuses to start while another daemon still holds its data dir."""
733 backend = _make_backend(tmp_path)
734 backend._server = MagicMock()
735 backend._binary = Path("/nonexistent/soloist")
736 session = _SoloistSession(backend, "player1")
737 backend._session = session
738 order: list[str] = []
739
740 async def _slow_stop() -> None:
741 order.append("stop-start")
742 await asyncio.sleep(0.05)
743 order.append("stop-done")
744
745 monkeypatch.setattr(session, "stop", _slow_stop)
746 _install_fake_binary_manager(monkeypatch)
747
748 async def _spawn(_self: Any, _uri: str, _seek: int) -> None:
749 order.append("spawn")
750 raise AudioError("spawn")
751
752 monkeypatch.setattr(soloist_backend._SoloistSession, "start", _spawn)
753 # the session failed, so its teardown is under way when the next item arrives
754 discard = asyncio.create_task(backend.discard_session(session))
755 await asyncio.sleep(0)
756 with pytest.raises(AudioError, match="spawn"):
757 await backend._acquire(TRACK_B, 0, "player1")
758 await discard
759 assert order == ["stop-start", "stop-done", "spawn"]
760
761
762async def test_an_idle_session_is_taken_over(
763 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
764) -> None:
765 """A session nobody is reading is replaced instead of blocking another player."""
766 backend = _make_backend(tmp_path)
767 backend._server = MagicMock()
768 backend._binary = Path("/nonexistent/soloist")
769 session = _SoloistSession(backend, "player1")
770 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
771 backend._session = session
772 stopped = AsyncMock()
773 monkeypatch.setattr(session, "stop", stopped)
774 _install_fake_binary_manager(monkeypatch)
775 # the replacement spawn is out of scope here; only the takeover decision is
776 monkeypatch.setattr(
777 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
778 )
779 with pytest.raises(AudioError, match="spawn"):
780 await backend._acquire(TRACK_B, 0, "player2")
781 stopped.assert_awaited_once()
782
783
784def test_a_dead_session_task_fails_the_session(tmp_path: Path) -> None:
785 """A session task that dies of an unexpected error takes the session with it."""
786 session = _make_session(tmp_path)
787 task: Any = MagicMock()
788 task.cancelled.return_value = False
789 task.exception.return_value = RuntimeError("reader blew up")
790 session._task_done(task)
791 assert session.usable is False
792 assert session._error is not None
793 assert "reader blew up" in session._error
794
795
796def test_a_cancelled_session_task_is_not_a_failure(tmp_path: Path) -> None:
797 """Teardown cancels the session's tasks; that must not be reported as an error."""
798 session = _make_session(tmp_path)
799 task: Any = MagicMock()
800 task.cancelled.return_value = True
801 session._task_done(task)
802 assert session.usable is True
803
804
805async def test_failed_sink_control_fails_the_session(tmp_path: Path) -> None:
806 """A failed suspend/resume fails the session instead of leaking stall silence."""
807 session = _make_session(tmp_path)
808 session._demand_started = True
809 session._sink_running = True
810 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
811 session._pending.append(TRACK_B)
812 _sink_of(session).suspend.side_effect = RuntimeError("pactl failed")
813 await session._handle_event(_playback_event("buffering"))
814 assert session._error is not None
815 assert "capture sink control failed" in session._error
816
817
818async def test_app_pause_is_fought_with_a_resume(tmp_path: Path) -> None:
819 """A pause from the Spotify app is undone: this session has no user-facing pause."""
820 session = _make_session(tmp_path)
821 session._demand_started = True
822 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
823 session._pending.append(TRACK_B)
824 await session._handle_event(_playback_event("paused"))
825 _client_of(session).resume.assert_awaited_once()
826
827
828async def test_an_app_pause_is_only_undone_so_many_times(tmp_path: Path) -> None:
829 """Someone who keeps pausing means it: the session gives up instead of fighting on."""
830 session = _make_session(tmp_path)
831 session._demand_started = True
832 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
833 session._pending.append(TRACK_B)
834 for _ in range(_MAX_APP_PAUSE_RESUMES):
835 await session._handle_event(_playback_event("playing"))
836 await session._handle_event(_playback_event("paused"))
837 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
838 assert session._error is None
839
840 await session._handle_event(_playback_event("playing"))
841 await session._handle_event(_playback_event("paused"))
842 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
843 assert session.usable is False
844 assert session._app_control is SoloistAppControl.PAUSED
845
846
847async def test_one_pause_reported_twice_counts_once(tmp_path: Path) -> None:
848 """A repeated snapshot of the same pause is not a new pause."""
849 session = _make_session(tmp_path)
850 session._demand_started = True
851 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
852 session._pending.append(TRACK_B)
853 for _ in range(_MAX_APP_PAUSE_RESUMES + 2):
854 await session._handle_event(_playback_event("paused"))
855 assert session.usable is True
856
857
858async def test_the_pause_budget_resets_on_the_next_item(tmp_path: Path) -> None:
859 """Each item gets its own budget; pausing one track does not spend the next one's."""
860 session = _make_session(tmp_path)
861 session._demand_started = True
862 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
863 session._pending.append(TRACK_B)
864 for _ in range(_MAX_APP_PAUSE_RESUMES):
865 await session._handle_event(_playback_event("playing"))
866 await session._handle_event(_playback_event("paused"))
867 await session._observe_current(TRACK_B, 200_000)
868 assert session._app_pauses == 0
869
870
871async def test_a_pause_is_not_undone_once_the_device_is_gone(tmp_path: Path) -> None:
872 """A bare resume on a device Spotify no longer routes to would play to nobody."""
873 session = _make_session(tmp_path)
874 session._demand_started = True
875 session._was_active = False
876 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
877 session._pending.append(TRACK_B)
878 await session._handle_event(_playback_event("paused"))
879 _client_of(session).resume.assert_not_awaited()
880
881
882async def test_losing_the_active_device_ends_the_session(tmp_path: Path) -> None:
883 """Playback moved to another device from the Spotify app: this session is over."""
884 session = _make_session(tmp_path)
885 await session._handle_event(_device_event(is_active=False))
886 assert session.usable is False
887 assert session._app_control is SoloistAppControl.TOOK_OVER
888
889
890async def test_a_takeover_reported_on_the_auth_state_ends_the_session(tmp_path: Path) -> None:
891 """The active-device state also rides on auth_state, and counts the same there."""
892 session = _make_session(tmp_path)
893 await session._handle_event(_auth_event(logged_in=True, is_active=False))
894 assert session.usable is False
895
896
897async def test_an_inactive_device_before_activation_is_not_a_takeover(tmp_path: Path) -> None:
898 """A fresh daemon is inactive until the session claims it; that is not a takeover."""
899 session = _make_session(tmp_path)
900 session._was_active = False
901 await session._handle_event(_device_event(is_active=False))
902 await session._handle_event(_auth_event(logged_in=True, is_active=False))
903 assert session.usable is True
904
905 # nor does a respawned daemon reporting the session Spotify still has for
906 # the account: only the status _play claimed is followed
907 await session._handle_event(_device_event(is_active=True))
908 await session._handle_event(_device_event(is_active=False))
909 assert session.usable is True
910
911
912async def test_a_reconnect_snapshot_keeps_an_active_session_alive(tmp_path: Path) -> None:
913 """The events connection re-snapshots after a drop; that is not a device change."""
914 session = _make_session(tmp_path)
915 await session._handle_event(_auth_event(logged_in=True, is_active=True))
916 await session._handle_event(_device_event(is_active=True))
917 assert session.usable is True
918
919
920async def test_the_playback_snapshots_active_flag_is_ignored(tmp_path: Path) -> None:
921 """It is optional and rides on deltas, so only the dedicated reports are followed."""
922 session = _make_session(tmp_path)
923 session._demand_started = True
924 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
925 await session._handle_event(
926 SoloistEvent(
927 type="playback_changed",
928 data=SoloistPlaybackState(status="playing", is_active=False),
929 raw={},
930 )
931 )
932 assert session.usable is True
933
934
935async def test_backpressure_does_not_spend_the_pause_budget(tmp_path: Path) -> None:
936 """A sink suspended to cap the cushion is our doing, not the user pausing."""
937 session = _make_session(tmp_path)
938 session._demand_started = True
939 session._engine_playing = True
940 session._backpressured = True
941 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
942 session._pending.append(TRACK_B)
943 await session._handle_event(_playback_event("paused"))
944 _client_of(session).resume.assert_not_awaited()
945 assert session._app_pauses == 0
946
947
948async def test_a_lost_login_is_not_reported_as_a_takeover(tmp_path: Path) -> None:
949 """Losing the login wins over the inactive device it brings with it."""
950 session = _make_session(tmp_path)
951 await session._handle_event(_auth_event(logged_in=False, is_active=False))
952 assert session.usable is False
953 assert session._app_control is None
954
955
956async def test_a_track_started_from_the_app_ends_the_session(tmp_path: Path) -> None:
957 """The engine pulled off an item part-way through is the app playing something else."""
958 session = _make_session(tmp_path)
959 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
960 item.duration_ms = 200_000
961 await session._observe_current(TRACK_A, 200_000)
962 item.observe_position(20_000)
963
964 await session._observe_current("spotify:track:theirs", 180_000)
965 assert session.usable is False
966 assert session._app_control is SoloistAppControl.TOOK_OVER
967 assert session.current is item
968
969
970async def test_a_track_played_earlier_started_from_the_app_ends_the_session(
971 tmp_path: Path,
972) -> None:
973 """A known uri is no exemption: only the item fed behind this one is where we sent it."""
974 session = _make_session(tmp_path)
975 played = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
976 played.spent = True
977 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
978 item.duration_ms = 200_000
979 await session._observe_current(TRACK_A, 200_000)
980 item.observe_position(20_000)
981
982 await session._observe_current(TRACK_B, 180_000)
983 assert session.usable is False
984 assert session._app_control is SoloistAppControl.TOOK_OVER
985
986
987async def test_skipping_from_the_app_to_the_fed_item_is_followed(tmp_path: Path) -> None:
988 """The queue moves to that same track, so following the engine keeps the two in step."""
989 session = _make_session(tmp_path)
990 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
991 item.duration_ms = 200_000
992 await session._observe_current(TRACK_A, 200_000)
993 item.observe_position(20_000)
994 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
995 session._pending.append(TRACK_B)
996
997 await session._observe_current(TRACK_B, 180_000)
998 assert session.usable is True
999 assert session.current is fed
1000 assert session.item_for(TRACK_B) is fed
1001
1002
1003async def test_a_takeover_snapshot_stops_pinning_volume_and_options(tmp_path: Path) -> None:
1004 """Once the app has the session, the rest of its snapshot must not reach the daemon."""
1005 session = _make_session(tmp_path)
1006 session._demand_started = True
1007 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1008 item.duration_ms = 200_000
1009 await session._observe_current(TRACK_A, 200_000)
1010 item.observe_position(20_000)
1011
1012 await session._handle_event(
1013 SoloistEvent(
1014 type="playback_changed",
1015 data=SoloistPlaybackState(
1016 status="playing",
1017 item=SoloistEntity(uri="spotify:track:theirs", entity_type="track"),
1018 volume=40,
1019 options=SoloistPlaybackOptions(shuffle=True, repeat="context"),
1020 ),
1021 raw={},
1022 )
1023 )
1024 assert session.usable is False
1025 _client_of(session).set_volume.assert_not_awaited()
1026 _client_of(session).set_shuffle.assert_not_awaited()
1027
1028
1029async def test_the_engine_moving_on_at_a_track_end_is_not_a_takeover(tmp_path: Path) -> None:
1030 """An unasked-for item the engine reaches at a boundary is its own autoplay."""
1031 session = _make_session(tmp_path)
1032 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1033 item.duration_ms = 200_000
1034 await session._observe_current(TRACK_A, 200_000)
1035 item.observe_position(200_000)
1036
1037 await session._observe_current("spotify:track:autoplay", 180_000)
1038 assert session.usable is True
1039 assert session.item_for("spotify:track:autoplay") is None
1040
1041
1042async def test_a_long_crossfade_boundary_is_not_a_takeover(tmp_path: Path) -> None:
1043 """With crossfade the engine moves on a crossfade short of the duration."""
1044 session = _make_session(tmp_path)
1045 session.crossfade_ms = 15_000
1046 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1047 await session._observe_current(TRACK_A, 200_000)
1048 # the last position reported before the engine crossfades into the next track
1049 item.observe_position(200_000 - session.crossfade_ms)
1050
1051 await session._observe_current("spotify:track:autoplay", 180_000)
1052 assert session.usable is True
1053
1054
1055async def test_an_ended_item_says_what_the_app_did(tmp_path: Path) -> None:
1056 """The item's stream fails with the takeover, not a generic session error."""
1057 session = _make_session(tmp_path)
1058 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1059 await session._handle_event(_device_event(is_active=False))
1060 with pytest.raises(SoloistAppControlError) as err:
1061 await session.validate_item(item)
1062 assert err.value.translation_key == SoloistAppControl.TOOK_OVER.value
1063 assert isinstance(err.value, ProviderStreamLimitError)
1064
1065
1066async def test_a_session_being_torn_down_does_not_hold_off_the_next_one(tmp_path: Path) -> None:
1067 """Teardown pauses the daemon; that must not read as the user pausing."""
1068 session = _make_session(tmp_path)
1069 session._demand_started = True
1070 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1071 session._pending.append(TRACK_B)
1072 session._stopped = True
1073 for _ in range(_MAX_APP_PAUSE_RESUMES + 1):
1074 await session._handle_event(_playback_event("playing"))
1075 await session._handle_event(_playback_event("paused"))
1076 await session._handle_event(_device_event(is_active=False))
1077 session.backend._raise_if_app_controlled()
1078 _client_of(session).resume.assert_not_awaited()
1079
1080
1081async def test_no_session_is_started_while_the_app_holds_the_last_one(tmp_path: Path) -> None:
1082 """A replacement would claim the Connect device straight back off the user."""
1083 backend = _make_backend(tmp_path)
1084 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1085 with pytest.raises(SoloistAppControlError):
1086 await backend._acquire(TRACK_A, 0, "player1")
1087
1088
1089async def test_the_hold_on_a_new_session_expires(tmp_path: Path) -> None:
1090 """Coming back to Music Assistant later plays again without any fuss."""
1091 backend = _make_backend(tmp_path)
1092 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1093 backend._app_control_until = time.monotonic() - 1
1094 backend._raise_if_app_controlled()
1095 assert backend._held_by_app() is None
1096
1097
1098async def test_an_audiobook_gives_up_on_capacity_instead_of_burning_chapters(
1099 tmp_path: Path,
1100) -> None:
1101 """Skipping ahead would cost the audiobook its availability and the caller its retry."""
1102 provider = _make_provider(tmp_path)
1103 calls: list[str] = []
1104
1105 async def _refuse(uri: str, *_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
1106 calls.append(uri)
1107 for _ in (): # never yields; only makes this an async generator
1108 yield b""
1109 raise SoloistAppControlError(provider, SoloistAppControl.TOOK_OVER)
1110
1111 provider.backend = MagicMock(stream_spotify_uri=_refuse)
1112 streamdetails = MagicMock(
1113 media_type=MediaType.AUDIOBOOK,
1114 data={"chapters": [TRACK_A, TRACK_B, "spotify:track:ccc"], "chapters_data": []},
1115 )
1116
1117 with pytest.raises(SoloistAppControlError):
1118 async for _ in provider.get_audio_stream(streamdetails):
1119 pass
1120 # the first chapter's refusal ends it: no chapter is skipped over
1121 assert calls == [TRACK_A]
1122
1123
1124def test_the_playback_device_is_named_apart_from_the_connect_one() -> None:
1125 """Two identically named devices in the Spotify app is what causes the takeovers."""
1126 assert SOLOIST_DEVICE_NAME != DEFAULT_PUBLISH_NAME
1127
1128
1129async def test_app_volume_change_is_pinned_back_to_unity(tmp_path: Path) -> None:
1130 """An off-unity volume set from the Spotify app is pinned back to 100."""
1131 session = _make_session(tmp_path)
1132 await session._handle_event(
1133 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=40), raw={})
1134 )
1135 _client_of(session).set_volume.assert_awaited_once_with(100)
1136 _client_of(session).set_volume.reset_mock()
1137 await session._handle_event(
1138 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=100), raw={})
1139 )
1140 _client_of(session).set_volume.assert_not_awaited()
1141
1142
1143async def test_track_change_signals_the_queue_when_it_matches_the_next_item(
1144 tmp_path: Path,
1145) -> None:
1146 """Reaching a fed item tells the queue to start filling that item's buffer."""
1147 session = _make_session(tmp_path, queue_id="player1")
1148 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1149 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1150 queues = _queues_of(session)
1151 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1152 await session._handle_event(
1153 SoloistEvent(
1154 type="track_changed",
1155 data=SoloistTrackChanged(item=SoloistEntity(uri=TRACK_B, entity_type="track")),
1156 raw={},
1157 )
1158 )
1159 queues.prepare_next_audio_buffer.assert_called_once_with("player1")
1160
1161
1162async def test_track_change_to_another_item_signals_nothing(tmp_path: Path) -> None:
1163 """An item the queue is not asking for next must not trigger a prebuffer."""
1164 session = _make_session(tmp_path, queue_id="player1")
1165 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1166 queues = _queues_of(session)
1167 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1168 await session._handle_event(
1169 SoloistEvent(
1170 type="track_changed",
1171 data=SoloistTrackChanged(
1172 item=SoloistEntity(uri="spotify:track:surprise", entity_type="track")
1173 ),
1174 raw={},
1175 )
1176 )
1177 queues.prepare_next_audio_buffer.assert_not_called()
1178
1179
1180async def test_the_follower_of_the_streamed_item_is_fed(tmp_path: Path) -> None:
1181 """The item after the one being streamed is handed to the engine."""
1182 session = _make_session(tmp_path, queue_id="player1")
1183 streamdetails = MagicMock()
1184 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1185 follower = _queue_item(TRACK_B)
1186 queues = _queues_of(session)
1187 queues.get.return_value = MagicMock(current_index=3)
1188 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 3 else None
1189 queues.get_next_item.return_value = follower
1190 await session.feed_after(streamdetails, TRACK_A)
1191 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1192 assert TRACK_B in session._items
1193 assert session.has_pending is True
1194
1195
1196async def test_an_item_the_queue_resolved_elsewhere_is_not_fed(tmp_path: Path) -> None:
1197 """A track the queue will stream from another provider must not be queued here."""
1198 session = _make_session(tmp_path, queue_id="player1")
1199 streamdetails = MagicMock()
1200 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1201 # same track, but the queue already picked a different provider for it
1202 follower = _queue_item(TRACK_B, streamdetails=MagicMock(provider="tidal--x"))
1203 queues = _queues_of(session)
1204 queues.get.return_value = MagicMock(current_index=0)
1205 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1206 queues.get_next_item.return_value = follower
1207 await session.feed_after(streamdetails, TRACK_A)
1208 _client_of(session).add_to_queue.assert_not_awaited()
1209
1210
1211async def test_skipping_to_the_fed_item_keeps_the_session(tmp_path: Path) -> None:
1212 """A next-track lands on the item already fed, so the engine jumps instead of respawning."""
1213 backend = _make_backend(tmp_path)
1214 backend._server = MagicMock()
1215 backend._binary = Path("/nonexistent/soloist")
1216 session = _SoloistSession(backend, "player1")
1217 session._client = AsyncMock()
1218 session._logged_in = True
1219 backend._session = session
1220 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1221 playing.started.set()
1222 # fed one ahead and not reached yet, which is where a next-track goes
1223 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1224 session._pending.append(TRACK_B)
1225
1226 async def _engine_gets_there(**_kwargs: Any) -> None:
1227 await session._observe_current(TRACK_B, 200_000)
1228
1229 _client_of(session).skip_next.side_effect = _engine_gets_there
1230 got_session, got_item = await backend._acquire(TRACK_B, 0, "player1")
1231 # the same session, no respawn, and the item that was already queued
1232 assert got_session is session
1233 assert got_item is fed
1234 assert backend._session is session
1235 _client_of(session).skip_next.assert_awaited_once()
1236
1237
1238async def test_a_skip_drops_what_arrives_while_the_command_is_in_flight(
1239 tmp_path: Path,
1240) -> None:
1241 """
1242 Audio captured between the skip command and the engine's answer is dropped.
1243
1244 Only covers the marker's own window; what the pipeline still holds when the
1245 answer arrives is measured at the cut instead.
1246 """
1247 session = _make_session(tmp_path)
1248 leaving = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1249 leaving.started.set()
1250 target = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1251 session._pending.append(TRACK_B)
1252 captured: list[bytes] = []
1253
1254 async def _engine_gets_there(**_kwargs: Any) -> None:
1255 # the pipeline still holds the old track while the command is in flight
1256 session._write_if_wanted(b"\x01" * 32)
1257 await session._observe_current(TRACK_B, 200_000)
1258 # from here on the audio really is the new item's
1259 session._write_if_wanted(b"\x02" * 32)
1260
1261 _client_of(session).skip_next.side_effect = _engine_gets_there
1262 await session.skip_to(target)
1263 captured.extend(target._chunks)
1264 assert b"".join(captured) == b"\x02" * 32
1265 assert session._discard_until is None
1266
1267
1268async def test_a_skip_the_engine_never_reaches_fails(
1269 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1270) -> None:
1271 """A skip that does not land is an error, not a wait for the track to end."""
1272 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
1273 session = _make_session(tmp_path)
1274 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1275 session._pending.append(TRACK_B)
1276 with pytest.raises(AudioError, match="did not reach"):
1277 await session.skip_to(fed)
1278
1279
1280async def test_a_fed_item_the_engine_has_not_reached_is_not_served(tmp_path: Path) -> None:
1281 """Skipping to an already-fed item must not hand over a channel that fills later."""
1282 session = _make_session(tmp_path)
1283 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1284 # fed, but the engine is still on the previous track
1285 assert session.item_for(TRACK_B) is None
1286 # once the engine gets there it is servable
1287 await session._observe_current(TRACK_B, 200_000)
1288 assert session.item_for(TRACK_B) is fed
1289
1290
1291def test_the_shaper_only_emits_whole_frames() -> None:
1292 """A read that ends mid-frame must never split a frame across two items."""
1293 shaper = soloist_backend._CaptureShaper()
1294 # the session's first bytes are infrastructure silence, and are dropped
1295 assert shaper.shape(b"\x00" * 4096) == b""
1296 # a mis-aligned read emits whole frames and carries the remainder
1297 first = shaper.shape(b"\x01" * (_FRAME_BYTES + 3))
1298 assert len(first) == _FRAME_BYTES
1299 # which is then completed by the next read, losing nothing
1300 second = shaper.shape(b"\x02" * (_FRAME_BYTES - 3))
1301 assert len(second) == _FRAME_BYTES
1302 assert second[:3] == b"\x01" * 3
1303 # an aligned read passes straight through
1304 assert shaper.shape(b"\x03" * _FRAME_BYTES) == b"\x03" * _FRAME_BYTES
1305
1306
1307def test_the_shaper_trims_lead_silence_only_once() -> None:
1308 """Silence after the audio has started is content, not pre-roll."""
1309 shaper = soloist_backend._CaptureShaper()
1310 assert shaper.shape(b"\x01" * _FRAME_BYTES) == b"\x01" * _FRAME_BYTES
1311 silence = b"\x00" * _FRAME_BYTES
1312 assert shaper.shape(silence) == silence
1313
1314
1315async def test_only_whole_sample_frames_are_handed_over(
1316 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1317) -> None:
1318 """A read that ends mid-frame must not split a frame across two items."""
1319 session = _make_session(tmp_path)
1320 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1321 item.started.set()
1322 item.claim()
1323 session._demand_started = True
1324 session._sink_running = True
1325 # two reads that are each mis-aligned but whole together
1326 reads = [b"\x01" * (_FRAME_BYTES + 3), b"\x02" * (_FRAME_BYTES - 3), b""]
1327 reader = MagicMock()
1328
1329 async def _read(_size: int) -> bytes:
1330 return reads.pop(0) if reads else b""
1331
1332 reader.read = _read
1333 session._reader = reader
1334 monkeypatch.setattr(soloist_backend, "_PACE_RATE", 1000.0)
1335 await session._read_capture()
1336 # every write was frame-aligned, and no byte was lost
1337 assert item.buffered % _FRAME_BYTES == 0
1338 assert item.buffered == _FRAME_BYTES * 2
1339
1340
1341async def test_an_already_known_item_is_not_fed_twice(tmp_path: Path) -> None:
1342 """An item the session already plays or was fed is not queued again."""
1343 session = _make_session(tmp_path, queue_id="player1")
1344 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1345 streamdetails = MagicMock()
1346 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1347 queues = _queues_of(session)
1348 queues.get.return_value = MagicMock(current_index=0)
1349 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1350 queues.get_next_item.return_value = _queue_item(TRACK_B)
1351 await session.feed_after(streamdetails, TRACK_A)
1352 _client_of(session).add_to_queue.assert_not_awaited()
1353
1354
1355async def test_only_tracks_are_fed_ahead(tmp_path: Path) -> None:
1356 """A podcast episode or audiobook chapter is played on its own, never stitched."""
1357 session = _make_session(tmp_path, queue_id="player1")
1358 await session.feed_after(MagicMock(), "spotify:episode:xyz")
1359 _client_of(session).add_to_queue.assert_not_awaited()
1360
1361
1362async def test_a_non_spotify_follower_is_not_fed(tmp_path: Path) -> None:
1363 """The run simply ends where the queue leaves this provider."""
1364 session = _make_session(tmp_path, queue_id="player1")
1365 streamdetails = MagicMock()
1366 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1367 follower = MagicMock(
1368 media_item=MagicMock(media_type=MediaType.TRACK, provider="tidal--x"), streamdetails=None
1369 )
1370 follower.media_item.provider_mappings = []
1371 queues = _queues_of(session)
1372 queues.get.return_value = MagicMock(current_index=0)
1373 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1374 queues.get_next_item.return_value = follower
1375 await session.feed_after(streamdetails, TRACK_A)
1376 _client_of(session).add_to_queue.assert_not_awaited()
1377
1378
1379async def test_a_library_item_is_fed_through_its_spotify_mapping(tmp_path: Path) -> None:
1380 """A library track is fed with the item id this provider instance knows it by."""
1381 session = _make_session(tmp_path, queue_id="player1")
1382 streamdetails = MagicMock()
1383 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1384 follower = MagicMock(
1385 media_item=MagicMock(media_type=MediaType.TRACK, provider="library", item_id="42"),
1386 streamdetails=None,
1387 )
1388 follower.media_item.provider_mappings = [
1389 MagicMock(provider_instance="other--y", item_id="wrong"),
1390 MagicMock(provider_instance="spotify--test", item_id="bbb"),
1391 ]
1392 queues = _queues_of(session)
1393 queues.get.return_value = MagicMock(current_index=0)
1394 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1395 queues.get_next_item.return_value = follower
1396 await session.feed_after(streamdetails, TRACK_A)
1397 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1398
1399
1400@pytest.mark.parametrize(
1401 ("provider_option", "player_setting", "expected"),
1402 [
1403 (True, "enabled", True),
1404 # the player's own switch decides first: off means nobody normalizes,
1405 # not that the job passes to Spotify
1406 (True, "disabled", False),
1407 (False, "enabled", False),
1408 (False, "disabled", False),
1409 ],
1410)
1411def test_who_normalizes_needs_both_switches(
1412 tmp_path: Path,
1413 monkeypatch: pytest.MonkeyPatch,
1414 provider_option: bool,
1415 player_setting: str,
1416 expected: bool,
1417) -> None:
1418 """The engine normalizes only when the provider option and the player agree."""
1419 session = _make_session(tmp_path, queue_id="player1")
1420 monkeypatch.setattr(
1421 type(session.backend.provider),
1422 "spotify_normalization_configured",
1423 property(lambda _self: provider_option),
1424 )
1425 cast("MagicMock", session.mass.config).get_effective_player_queue_config_value = MagicMock(
1426 return_value=player_setting
1427 )
1428 assert session._engine_normalization_enabled() is expected
1429
1430
1431def test_a_running_session_answers_for_what_the_engine_is_doing(tmp_path: Path) -> None:
1432 """
1433 The engine reads its settings at startup, so a later toggle must not split them.
1434
1435 Otherwise the streams core would start normalizing on top of audio the engine
1436 is still normalizing, or stop while it no longer is.
1437 """
1438 backend = _make_backend(tmp_path)
1439 provider = backend.provider
1440 # nothing playing yet: the configuration is all there is to go on
1441 before_any_session = backend.session_normalizes
1442 session = _SoloistSession(backend, "player1")
1443 session.engine_normalizes = True
1444 backend._session = session
1445 while_playing = backend.session_normalizes
1446 # ... and a session that has been torn down no longer speaks for the engine
1447 session._stopped = True
1448 after_teardown = backend.session_normalizes
1449 assert before_any_session is None
1450 assert while_playing is True
1451 assert after_teardown is None
1452 assert provider.delivers_normalized_audio is provider.spotify_normalization_configured
1453
1454
1455def test_crossfade_comes_from_the_queue_preference(tmp_path: Path) -> None:
1456 """The queue's crossfade setting is handed to the engine, in milliseconds."""
1457 session = _make_session(tmp_path, queue_id="player1")
1458 _queues_of(session).get.return_value = MagicMock(queue_id="player1", crossfade_enabled=True)
1459 cast("MagicMock", session.mass.config).get_raw_core_config_value = MagicMock(return_value=6)
1460 assert session._queue_crossfade_ms() == 6000
1461
1462
1463def test_crossfade_off_is_zero(tmp_path: Path) -> None:
1464 """A queue with crossfade disabled gets an explicit zero (which clears the pref)."""
1465 session = _make_session(tmp_path, queue_id="player1")
1466 _queues_of(session).get.return_value = MagicMock(crossfade_enabled=False)
1467 assert session._queue_crossfade_ms() == 0
1468
1469
1470def test_no_queue_means_no_crossfade(tmp_path: Path) -> None:
1471 """Without a queue to read the preference from, the engine gets no crossfade."""
1472 session = _make_session(tmp_path, queue_id=None)
1473 assert session._queue_crossfade_ms() == 0
1474
1475
1476async def test_short_delivery_is_rejected_as_incomplete(tmp_path: Path) -> None:
1477 """PCM that stops well short of the item's duration is rejected."""
1478 session = _make_session(tmp_path)
1479 item = _ItemAudio(TRACK_A, session)
1480 item.playing_seen = True
1481 item.duration_ms = 200_000
1482 item.last_position_ms = 100_000
1483 with pytest.raises(AudioError, match="incomplete"):
1484 await session.validate_item(item)
1485
1486
1487async def test_a_crossfade_shortfall_is_tolerated(tmp_path: Path) -> None:
1488 """With crossfade the engine reports the item short by design; that is not a failure."""
1489 session = _make_session(tmp_path)
1490 session.crossfade_ms = 12_000
1491 item = _ItemAudio(TRACK_A, session)
1492 item.playing_seen = True
1493 item.duration_ms = 200_000
1494 # 12s of crossfade plus the ordinary tolerance
1495 item.last_position_ms = 200_000 - 21_000
1496 await session.validate_item(item)
1497
1498
1499async def test_missing_position_is_rejected_as_incomplete(tmp_path: Path) -> None:
1500 """Without any position report there is no evidence the item played out."""
1501 session = _make_session(tmp_path)
1502 item = _ItemAudio(TRACK_A, session)
1503 item.playing_seen = True
1504 item.duration_ms = 200_000
1505 with pytest.raises(AudioError, match="incomplete"):
1506 await session.validate_item(item)
1507
1508
1509async def test_short_item_cannot_pass_at_position_zero(tmp_path: Path) -> None:
1510 """The tolerance never spans a whole item, so a short item cannot pass unplayed."""
1511 session = _make_session(tmp_path)
1512 item = _ItemAudio(TRACK_A, session)
1513 item.playing_seen = True
1514 item.duration_ms = 8_000
1515 item.last_position_ms = 0
1516 with pytest.raises(AudioError, match="incomplete"):
1517 await session.validate_item(item)
1518
1519
1520async def test_an_item_that_never_played_is_rejected(tmp_path: Path) -> None:
1521 """An item the engine never reported playing is a failure, whatever was delivered."""
1522 session = _make_session(tmp_path)
1523 item = _ItemAudio(TRACK_A, session)
1524 item.duration_ms = 200_000
1525 item.last_position_ms = 200_000
1526 with pytest.raises(AudioError, match="never started playing"):
1527 await session.validate_item(item)
1528
1529
1530async def test_a_duration_less_item_is_not_judged(tmp_path: Path) -> None:
1531 """Without a duration there is nothing to judge completeness against."""
1532 session = _make_session(tmp_path)
1533 item = _ItemAudio(TRACK_A, session)
1534 item.playing_seen = True
1535 await session.validate_item(item)
1536
1537
1538def test_an_unread_session_expires(tmp_path: Path) -> None:
1539 """A session no item stream reads from is ended so its daemon does not linger."""
1540 session = _make_session(tmp_path)
1541 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1542 session._expire_idle()
1543 assert session._idle_since is not None
1544 assert session.usable is True
1545 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S - 1
1546 session._expire_idle()
1547 assert session.usable is False
1548
1549
1550def test_a_session_being_read_never_expires(tmp_path: Path) -> None:
1551 """An item stream reading the session keeps it alive indefinitely."""
1552 session = _make_session(tmp_path)
1553 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1554 item.claim()
1555 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S * 10
1556 session._expire_idle()
1557 assert session.usable is True
1558
1559
1560def test_pre_roll_silence_is_dropped_a_whole_frame_at_a_time() -> None:
1561 """
1562 Trimming pre-roll must leave the audio on the session's frame grid.
1563
1564 A FIFO read is not always a whole number of frames, and dropping a partial
1565 one would shift every sample that follows for the rest of the session.
1566 """
1567 shaper = _CaptureShaper()
1568 # pre-roll that ends mid-frame: the real audio starts at byte 1024
1569 assert shaper.shape(b"\x00" * 1021) == b""
1570 audio = bytes(range(1, 9)) * 4
1571 shaped = shaper.shape(b"\x00" * 3 + audio)
1572 assert shaped == audio
1573 assert shaper._lead_skipped % _FRAME_BYTES == 0
1574
1575
1576async def test_a_refused_skip_does_not_leave_the_audio_discarded(tmp_path: Path) -> None:
1577 """
1578 A skip that never landed must not keep the session dropping its audio.
1579
1580 The marker silences everything the session captures, so a command that
1581 failed has to clear it on the way out.
1582 """
1583 session = _make_session(tmp_path)
1584 client = cast("MagicMock", session._client)
1585 client.skip_next = AsyncMock(side_effect=TimeoutError)
1586 item = _ItemAudio(TRACK_B, session)
1587
1588 with pytest.raises(AudioError, match="would not skip"):
1589 await session.skip_to(item)
1590
1591 assert session._discard_until is None
1592
1593
1594async def test_a_daemon_that_will_not_die_is_reported_and_released(tmp_path: Path) -> None:
1595 """A close that could not terminate the daemon still finishes the teardown."""
1596 session = _make_session(tmp_path)
1597 proc = cast("MagicMock", session._proc)
1598 proc.close = AsyncMock()
1599 # AsyncProcess.close() gives up after a handful of kill attempts
1600 proc.returncode = None
1601 with patch.object(session.logger, "warning") as warning:
1602 await session.stop()
1603 assert warning.called
1604 assert session._teardown_done is True
1605 assert session._proc is None
1606
1607
1608async def test_a_cancelled_teardown_still_closes_the_daemon(tmp_path: Path) -> None:
1609 """
1610 A cancelled teardown must leave the retry something to close.
1611
1612 Dropping the references first is how a daemon survives to hold the data
1613 directory, which every later session is then refused for.
1614 """
1615 session = _make_session(tmp_path)
1616 proc = cast("MagicMock", session._proc)
1617 sink = cast("AsyncMock", session._sink)
1618
1619 async def _never_returns() -> None:
1620 await asyncio.Event().wait()
1621
1622 proc.close = _never_returns
1623 task = asyncio.create_task(session.stop())
1624 await asyncio.sleep(0.01)
1625 task.cancel()
1626 with suppress(asyncio.CancelledError):
1627 await task
1628 # the teardown did not finish, so nothing was dropped and it can be redone
1629 unfinished = session._teardown_done
1630 kept_proc = session._proc
1631 kept_sink = session._sink
1632 proc.close = AsyncMock()
1633 proc.returncode = 0
1634 await session.stop()
1635 assert unfinished is False
1636 assert kept_proc is proc
1637 assert kept_sink is sink
1638 assert session._teardown_done is True
1639 assert session._proc is None
1640 assert session._sink is None
1641 proc.close.assert_awaited()
1642 sink.unload.assert_awaited()
1643
1644
1645def test_a_failed_session_is_torn_down(tmp_path: Path) -> None:
1646 """A session that fails is discarded, so its daemon does not keep playing to nobody."""
1647 session = _make_session(tmp_path)
1648 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1649 item.claim()
1650 session._fail("audio stalled")
1651 assert session.usable is False
1652 # every waiting item is released and the teardown is scheduled
1653 assert item._closed is True
1654 # a startup wait must not sit out its timeout on a session that already failed
1655 assert item.started.is_set() is True
1656 discard = cast("MagicMock", session.mass.create_task)
1657 discard.assert_called_once_with(session.backend.discard_session, session)
1658 # a second failure does not queue a second teardown
1659 session._fail("and again")
1660 assert session._error == "audio stalled"
1661 assert discard.call_count == 1
1662
1663
1664async def test_an_item_the_engine_skipped_past_fails_instead_of_hanging(
1665 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1666) -> None:
1667 """A claimed channel the engine never reaches gives up rather than blocking forever."""
1668 monkeypatch.setattr(soloist_backend, "_READ_SLICE_S", 0.01)
1669 monkeypatch.setattr(soloist_backend, "_STALL_TIMEOUT_S", 0.05)
1670 session = _make_session(tmp_path)
1671 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1672 item.claim()
1673 # the engine is playing something else, so nothing is ever written here
1674 session._items["spotify:track:other"] = session._current = _ItemAudio(
1675 "spotify:track:other", session
1676 )
1677 with pytest.raises(AudioError, match="no audio"):
1678 async for _ in item.read():
1679 pass
1680
1681
1682async def test_adopt_paired_session_copies_into_the_canonical_dir(
1683 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1684) -> None:
1685 """A session paired by the setup flow is adopted into the per-instance data dir."""
1686 storage = tmp_path / "storage"
1687 pending = storage / "spotify" / "pairing" / "flow1"
1688 pending.mkdir(parents=True)
1689 (pending / "session.bin").write_bytes(b"session")
1690 prov = _make_provider(tmp_path, {CONF_SOLOIST_SESSION_DIR: "spotify/pairing/flow1"})
1691 update_setup_data = MagicMock()
1692 monkeypatch.setattr(prov, "_update_setup_data", update_setup_data)
1693 backend = SoloistBackend(prov)
1694 await backend._adopt_paired_session()
1695 canonical = storage / "spotify" / "spotify--test" / SOLOIST_DATA_DIR_NAME
1696 assert (canonical / "session.bin").read_bytes() == b"session"
1697 # a copy, not a move: the flow-private source must survive a failed
1698 # provider load so the setup flow can retry (the flow removes it at its end)
1699 assert (pending / "session.bin").exists()
1700 update_setup_data.assert_called_once_with(CONF_SOLOIST_SESSION_DIR, None)
1701
1702
1703def test_the_engine_is_told_not_to_normalize(tmp_path: Path) -> None:
1704 """MA normalizes this audio itself, so the engine's own normalization is switched off."""
1705 backend = _make_backend(tmp_path)
1706 prefs = backend._data_dir / "settings" / "Users" / "alice-user" / "prefs"
1707 prefs.parent.mkdir(parents=True)
1708 prefs.write_text("some.engine.key=1\n", encoding="utf-8")
1709 backend._prepare_data_dir(8000, normalize=False)
1710 content = prefs.read_text(encoding="utf-8").splitlines()
1711 assert "some.engine.key=1" in content
1712 assert "audio.normalize_v2=false" in content
1713 assert "audio.crossfade_v2=true" in content
1714 assert "audio.crossfade.time_v2=8000" in content
1715 # the ceiling is stated rather than left to the engine's own default
1716 assert "audio.play_bitrate_enumeration=5" in content
1717 assert "audio.play_bitrate_non_metered_enumeration=5" in content
1718 assert "audio.play_bitrate_non_metered_migrated=true" in content
1719
1720
1721def test_disabling_crossfade_writes_the_boolean(tmp_path: Path) -> None:
1722 """Crossfade off is written explicitly, so a stale 'on' cannot survive."""
1723 backend = _make_backend(tmp_path)
1724 prefs = backend._data_dir / "settings" / "prefs"
1725 prefs.parent.mkdir(parents=True)
1726 prefs.write_text("audio.crossfade_v2=true\naudio.crossfade.time_v2=8000\n", encoding="utf-8")
1727 backend._prepare_data_dir(0, normalize=False)
1728 content = prefs.read_text(encoding="utf-8").splitlines()
1729 assert "audio.crossfade_v2=false" in content
1730 assert not any(line.startswith("audio.crossfade.time_v2") for line in content)
1731
1732
1733async def test_setup_requires_an_api_key(tmp_path: Path) -> None:
1734 """Without a stored API key the user must be sent back through the setup flow."""
1735 backend = _make_backend(tmp_path)
1736 with pytest.raises(LoginFailed) as err:
1737 await backend.setup()
1738 assert err.value.translation_key == "soloist_pairing_required"
1739
1740
1741async def test_setup_requires_a_paired_session(
1742 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1743) -> None:
1744 """An API key without a paired session also routes back to the setup flow."""
1745 backend = _make_backend(tmp_path, {CONF_SOLOIST_API_KEY: "k" * 20, CONF_SOLOIST_CONSENT: True})
1746 _install_fake_binary_manager(monkeypatch)
1747 with pytest.raises(LoginFailed) as err:
1748 await backend.setup()
1749 assert err.value.translation_key == "soloist_pairing_required"
1750
1751
1752async def test_streaming_without_setup_is_refused(tmp_path: Path) -> None:
1753 """A backend whose setup never ran refuses to stream instead of half-starting."""
1754 backend = _make_backend(tmp_path)
1755 with pytest.raises(AudioError, match="not started"):
1756 async for _ in backend.stream_spotify_uri(TRACK_A):
1757 pass
1758
1759
1760def test_session_present_detection(tmp_path: Path) -> None:
1761 """Only a dir holding something besides the WS endpoint files counts as paired."""
1762 data_dir = tmp_path / "soloist-data"
1763 assert soloist_session_present(data_dir) is False
1764 data_dir.mkdir()
1765 (data_dir / WS_ADDR_FILE).write_text("127.0.0.1", encoding="utf-8")
1766 (data_dir / WS_PORT_FILE).write_text("1234", encoding="utf-8")
1767 assert soloist_session_present(data_dir) is False
1768 (data_dir / "session.bin").write_bytes(b"x")
1769 assert soloist_session_present(data_dir) is True
1770
1771
1772async def test_a_skip_drops_the_audio_still_in_flight(tmp_path: Path) -> None:
1773 """The item jumped to opens with its own audio, not the tail of the one left behind."""
1774 session = _make_session(tmp_path)
1775 left_behind = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1776 session._current = left_behind
1777 left_behind.started.set()
1778 left_behind.claim()
1779 jumped_to = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1780 session._pending.append(TRACK_B)
1781 session._discard_until = TRACK_B
1782 with _capture_holding(session, fifo_bytes=2 * _FRAME_BYTES, reader_bytes=2 * _FRAME_BYTES):
1783 await session._observe_current(TRACK_B, 200_000)
1784 assert session._stale_budget == 4 * _FRAME_BYTES
1785 session._write_if_wanted(b"s" * (4 * _FRAME_BYTES))
1786 session._write_if_wanted(b"n" * (2 * _FRAME_BYTES))
1787 jumped_to.claim()
1788 jumped_to.close()
1789 assert b"".join([chunk async for chunk in jumped_to.read()]) == b"n" * (2 * _FRAME_BYTES)
1790
1791
1792async def test_a_skip_drops_the_stale_audio_across_reads(tmp_path: Path) -> None:
1793 """A budget larger than one read keeps dropping, and resumes on a frame boundary."""
1794 session = _make_session(tmp_path)
1795 session._stale_budget = 3 * _FRAME_BYTES
1796 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1797 item.claim()
1798 session._write_if_wanted(b"s" * (2 * _FRAME_BYTES))
1799 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1800 item.close()
1801 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1802
1803
1804async def test_the_marker_spends_an_earlier_jumps_budget(tmp_path: Path) -> None:
1805 """What the marker drops still counts against a budget left from an earlier jump."""
1806 session = _make_session(tmp_path)
1807 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1808 session._stale_budget = 4 * _FRAME_BYTES
1809 session._discard_until = TRACK_B
1810 session._write_if_wanted(b"s" * (3 * _FRAME_BYTES))
1811 assert session._stale_budget == _FRAME_BYTES
1812 # a refused command leaves only what is genuinely still in flight to drop
1813 session._discard_until = None
1814 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1815 item = session._current
1816 item.claim()
1817 item.close()
1818 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1819
1820
1821async def test_a_natural_cut_keeps_the_audio_in_flight(tmp_path: Path) -> None:
1822 """Nothing is dropped without a jump: what is in flight is the continuation."""
1823 session = _make_session(tmp_path)
1824 playing = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1825 session._current = playing
1826 playing.started.set()
1827 playing.claim()
1828 with _capture_holding(session, fifo_bytes=4 * _FRAME_BYTES, reader_bytes=4 * _FRAME_BYTES):
1829 await session._observe_current(TRACK_B, 200_000)
1830 assert session._stale_budget == 0
1831
1832
1833def test_stale_bytes_spans_both_buffers_in_whole_frames(tmp_path: Path) -> None:
1834 """The in-flight measure covers the FIFO and the reader, and never splits a frame."""
1835 session = _make_session(tmp_path)
1836 with _capture_holding(session, fifo_bytes=3 * _FRAME_BYTES + 3, reader_bytes=2 * _FRAME_BYTES):
1837 assert session._stale_bytes() == 5 * _FRAME_BYTES
1838
1839
1840def test_stale_bytes_falls_back_when_the_reader_cannot_be_sized(tmp_path: Path) -> None:
1841 """Losing the reader's internal view drops extra rather than leaving audio behind."""
1842 session = _make_session(tmp_path)
1843 with _capture_holding(session, fifo_bytes=0, reader_bytes=None):
1844 assert session._stale_bytes() == 6 * _READ_CHUNK_SIZE
1845
1846
1847async def test_a_channel_abandoned_at_the_cut_stops_holding_the_cushion(
1848 tmp_path: Path,
1849) -> None:
1850 """A skip closes the channel first and only then unwinds its stream."""
1851 session = _make_session(tmp_path)
1852 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1853 item.started.set()
1854 item.claim()
1855 item.write(b"x" * 4096)
1856 # the cut lands while the abandoned stream is still unwinding
1857 await session._observe_current(TRACK_B, 200_000)
1858 assert session._retained_bytes() == 4096
1859 item.release()
1860 assert session._retained_bytes() == 0
1861
1862
1863def test_an_abandoned_channel_stops_holding_the_cushion(tmp_path: Path) -> None:
1864 """A channel skipped away from frees its buffer instead of gating the sink for good."""
1865 session = _make_session(tmp_path)
1866 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1867 item.claim()
1868 item.write(b"x" * 4096)
1869 assert session._retained_bytes() == 4096
1870 # the stream is gone, then the cut closes the channel
1871 item.release()
1872 item.close()
1873 assert session._retained_bytes() == 0
1874
1875
1876async def test_a_channel_still_being_read_keeps_its_tail(tmp_path: Path) -> None:
1877 """Closing the playing item at a cut must not discard what its stream is still owed."""
1878 session = _make_session(tmp_path)
1879 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1880 item.claim()
1881 item.write(b"tail" * 4)
1882 item.close()
1883 assert item.buffered == 16
1884 assert b"".join([chunk async for chunk in item.read()]) == b"tail" * 4
1885
1886
1887@contextmanager
1888def _capture_holding(
1889 session: _SoloistSession, *, fifo_bytes: int, reader_bytes: int | None
1890) -> Iterator[None]:
1891 """
1892 Give the session a real capture FIFO and a reader holding the given amounts.
1893
1894 A real pipe is used so the byte count comes from the same ioctl the backend
1895 relies on. Pass ``reader_bytes=None`` for a reader whose buffer cannot be read.
1896 """
1897 read_fd, write_fd = os.pipe()
1898 try:
1899 if fifo_bytes:
1900 os.write(write_fd, bytes(fifo_bytes))
1901 pipe = MagicMock()
1902 pipe.fileno.return_value = read_fd
1903 transport = MagicMock()
1904 transport.get_extra_info.return_value = pipe
1905 session._transport = transport
1906 reader = MagicMock(spec=[]) if reader_bytes is None else MagicMock()
1907 if reader_bytes is not None:
1908 reader._buffer = bytearray(reader_bytes)
1909 session._reader = reader
1910 yield
1911 finally:
1912 session._transport = None
1913 session._reader = None
1914 os.close(read_fd)
1915 os.close(write_fd)
1916
1917
1918def _make_provider(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SpotifyProvider:
1919 """Return a SpotifyProvider (bypassing __init__) with the given setup_data."""
1920 prov = object.__new__(SpotifyProvider)
1921 config = MagicMock(instance_id="spotify--test")
1922 config.get_value = MagicMock(return_value=None)
1923 config.values = {}
1924 prov.config = config
1925 prov.manifest = MagicMock(domain="spotify")
1926 prov.logger = MagicMock()
1927 prov.available = True
1928 mass = MagicMock()
1929 mass.storage_path = str(tmp_path / "storage")
1930 mass.cache_path = str(tmp_path / "cache")
1931 # get_setup_value reads the live setup_data blob from the store
1932 mass.config.get = MagicMock(return_value=setup_data or {})
1933 mass.config.get_raw_provider_config_value = MagicMock(return_value=None)
1934 # the store keeps values encrypted; decrypt is an identity map for the test
1935 mass.config.decrypt_string = MagicMock(side_effect=lambda value: value)
1936 prov.mass = mass
1937 return prov
1938
1939
1940def _make_backend(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SoloistBackend:
1941 """Return a SoloistBackend on a mocked provider."""
1942 return SoloistBackend(_make_provider(tmp_path, setup_data))
1943
1944
1945def _make_session(tmp_path: Path, queue_id: str | None = "player1") -> _SoloistSession:
1946 """Return a session with its process/sink/client replaced by mocks."""
1947 session = _SoloistSession(_make_backend(tmp_path), queue_id)
1948 session._sink = AsyncMock()
1949 session._client = AsyncMock()
1950 session._proc = MagicMock(returncode=None)
1951 # a session under test is past the engine's login and has claimed the
1952 # Connect device, unless a test says otherwise
1953 session._logged_in = True
1954 session._was_active = True
1955 return session
1956
1957
1958def _make_item(tmp_path: Path, uri: str) -> _ItemAudio:
1959 """Return a bare item channel on a mocked session."""
1960 return _ItemAudio(uri, _make_session(tmp_path))
1961
1962
1963def _queue_item(uri: str, streamdetails: Any = None) -> MagicMock:
1964 """Return a queue item stand-in for a Spotify track on the test instance."""
1965 item_id = uri.rsplit(":", 1)[1]
1966 media_item = MagicMock(media_type=MediaType.TRACK, provider="spotify--test", item_id=item_id)
1967 media_item.provider_mappings = []
1968 return MagicMock(
1969 media_item=media_item,
1970 queue_item_id=f"qi-{item_id}",
1971 streamdetails=streamdetails,
1972 )
1973
1974
1975async def _wait_for(predicate: Callable[[], bool], timeout: float = 2.0) -> None:
1976 """Wait until the predicate holds, so a background task can get there."""
1977 loop = asyncio.get_running_loop()
1978 deadline = loop.time() + timeout
1979 while loop.time() < deadline:
1980 if predicate():
1981 return
1982 await asyncio.sleep(0.01)
1983 raise AssertionError("condition not met within timeout")
1984
1985
1986def _client_of(session: _SoloistSession) -> AsyncMock:
1987 """Return the session's mocked WebSocket client."""
1988 return cast("AsyncMock", session._client)
1989
1990
1991def _sink_of(session: _SoloistSession) -> AsyncMock:
1992 """Return the session's mocked capture sink."""
1993 return cast("AsyncMock", session._sink)
1994
1995
1996def _queues_of(session: _SoloistSession) -> MagicMock:
1997 """Return the mocked player_queues controller the session consults."""
1998 return cast("MagicMock", session.mass.player_queues)
1999
2000
2001def _auth_event(*, logged_in: bool, is_active: bool = True) -> SoloistEvent:
2002 """Return an auth_state event with the given login and active-device state."""
2003 return SoloistEvent(
2004 type="auth_state",
2005 data=SoloistAuthState(logged_in=logged_in, is_active=is_active),
2006 raw={},
2007 )
2008
2009
2010def _device_event(*, is_active: bool) -> SoloistEvent:
2011 """Return a device_changed event with the given active-device state."""
2012 return SoloistEvent(
2013 type="device_changed", data=SoloistDeviceChanged(is_active=is_active), raw={}
2014 )
2015
2016
2017def _playback_event(status: str, position_ms: int = 0) -> SoloistEvent:
2018 """Return a playback_state event for the current item with the given status."""
2019 return SoloistEvent(
2020 type="playback_state",
2021 data=SoloistPlaybackState(
2022 status=status,
2023 item=SoloistEntity(uri=TRACK_A, entity_type="track"),
2024 position=SoloistPosition(position_ms=position_ms, timestamp_ms=0),
2025 ),
2026 raw={},
2027 )
2028
2029
2030def _install_fake_binary_manager(monkeypatch: pytest.MonkeyPatch) -> None:
2031 """Replace the shared binary manager so no download or exec is attempted."""
2032 manager = MagicMock()
2033 manager.ensure_fresh = AsyncMock(return_value=Path("/nonexistent/soloist"))
2034 monkeypatch.setattr(soloist_backend, "SoloistBinaryManager", MagicMock(return_value=manager))
2035