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