/
/
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 await session._log_output(
528 _stdout_of(
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 assert session._data_dir_busy is True
534
535
536async def test_a_lost_pairing_is_caught_the_moment_the_daemon_reports_it(
537 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
538) -> None:
539 """A daemon advertising for pairing fails the session at once, not on a timeout."""
540 session = _make_session(tmp_path)
541 session._logged_in = None
542 unload_with_error = MagicMock()
543 monkeypatch.setattr(session.backend.provider, "unload_with_error", unload_with_error)
544 await session._log_output(
545 _stdout_of('waiting for login - connect to "X" from your Spotify app')
546 )
547 assert session._unpaired is True
548 # the buffer gives up on the audio long before the startup budget runs out, so
549 # the session has to fail while an item is still waiting on it
550 assert session._error is not None
551 with pytest.raises(LoginFailed) as err:
552 session._raise_startup_error("did not connect and log in", TRACK_A)
553 assert err.value.translation_key == "soloist_pairing_required"
554 unload_with_error.assert_called_once()
555
556
557async def test_a_lost_pairing_fails_the_item_without_waiting_for_the_endpoint(
558 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
559) -> None:
560 """The item fails on the lost pairing, not on the endpoint that is no longer coming."""
561 session = _make_session(tmp_path)
562 session._logged_in = None
563 monkeypatch.setattr(session.backend.provider, "unload_with_error", MagicMock())
564 await session._log_output(
565 _stdout_of('waiting for login - connect to "X" from your Spotify app')
566 )
567 # the endpoint never appears, so the wait for it must not swallow the failure:
568 # sitting it out would outlast the queue's own patience for the audio
569 with pytest.raises(LoginFailed) as err:
570 async with asyncio.timeout(5):
571 await session._play(TRACK_A, 0, asyncio.Event())
572 assert err.value.translation_key == "soloist_pairing_required"
573
574
575async def test_a_daemon_still_restoring_its_session_is_left_alone(tmp_path: Path) -> None:
576 """The engine advertises for pairing while restoring too; the stored session decides."""
577 session = _make_session(tmp_path)
578 data_dir = session.backend._data_dir
579 (data_dir / "settings" / "Users" / "spotify-user-user").mkdir(parents=True)
580 await session._log_output(
581 _stdout_of('waiting for login - connect to "X" from your Spotify app')
582 )
583 assert session._unpaired is False
584 assert session._error is None
585
586
587async def test_a_pairing_that_never_logs_in_routes_through_setup(
588 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
589) -> None:
590 """A session that cannot log in sends the user back to setup, not a per-track error."""
591 session = _make_session(tmp_path)
592 unload_with_error = MagicMock()
593 monkeypatch.setattr(session.backend.provider, "unload_with_error", unload_with_error)
594 await session._handle_event(_auth_event(logged_in=False))
595 with pytest.raises(LoginFailed) as err:
596 session._raise_startup_error("timed out waiting for playback to start", TRACK_A)
597 assert err.value.translation_key == "soloist_pairing_required"
598 # ... and the provider is taken out of service, so the user is asked to redo setup
599 unload_with_error.assert_called_once()
600
601
602async def test_a_login_that_never_happened_is_not_confused_with_another_failure(
603 tmp_path: Path,
604) -> None:
605 """An unrelated failure keeps its own message even before any login was reported."""
606 session = _make_session(tmp_path)
607 session._fail("the capture sink was lost mid-stream")
608 with pytest.raises(AudioError, match="capture sink was lost"):
609 session._raise_startup_error("exited before playback started", TRACK_A)
610
611
612def test_a_seeked_item_only_expects_what_is_left_of_it(tmp_path: Path) -> None:
613 """A seeked item delivers the remainder, so its targets are based on that."""
614 session = _make_session(tmp_path)
615 item = _ItemAudio(TRACK_A, session)
616 item.duration_ms = 200_000
617 full = 200_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
618 assert item._duration_bytes() == full
619 item.seek_target_ms = 150_000
620 remainder = 50_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
621 assert item._duration_bytes() == remainder
622 # so the tail drain has a target it can actually reach
623 item.start_tail_drain()
624 item.write(b"\x01" * remainder)
625 assert item.tail_complete is True
626 # and the padding silence after it is refused
627 item.write(b"\x00" * 4096)
628 assert item.buffered == remainder
629
630
631def test_the_lead_trim_never_exceeds_its_budget() -> None:
632 """Silence beyond the budget is content, including where audio starts mid-chunk."""
633 budget = int(_MAX_LEAD_TRIM_S * _BYTES_PER_SECOND)
634 # already at the budget, with a chunk whose silence runs well past it
635 chunk = b"\x00" * 4096 + b"\x01" * 64
636 trimmed, skipped = _trim_lead_silence(chunk, budget - _FRAME_BYTES)
637 assert skipped == _FRAME_BYTES
638 assert len(trimmed) == len(chunk) - _FRAME_BYTES
639
640
641async def test_a_dying_log_reader_fails_the_session(tmp_path: Path) -> None:
642 """Nothing else drains the daemon's stdout, so a dead reader must not go unnoticed."""
643 session = _make_session(tmp_path)
644
645 async def _boom() -> None:
646 raise RuntimeError("reader blew up")
647
648 session._log_task = asyncio.create_task(_boom())
649 session._log_task.add_done_callback(session._task_done)
650 await asyncio.sleep(0)
651 await _wait_for(lambda: not session.usable)
652 assert session._error is not None
653 assert "reader blew up" in session._error
654
655
656async def test_feeding_never_replaces_a_channel_already_in_use(tmp_path: Path) -> None:
657 """If the engine reaches the fed item first, its live channel must survive."""
658 session = _make_session(tmp_path, queue_id="player1")
659 streamdetails = MagicMock()
660 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
661 queues = _queues_of(session)
662 queues.get.return_value = MagicMock(current_index=0)
663 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
664 queues.get_next_item.return_value = _queue_item(TRACK_B)
665
666 async def _engine_gets_there_first(_uri: str, **_kwargs: Any) -> None:
667 # the events task advances to the fed item while the command is in flight
668 await session._observe_current(TRACK_B, 200_000)
669
670 _client_of(session).add_to_queue.side_effect = _engine_gets_there_first
671 await session.feed_after(streamdetails, TRACK_A)
672 live = session.current
673 assert live is not None
674 assert live.uri == TRACK_B
675 # the channel the reader is writing to is the one a stream will be handed
676 assert session._items[TRACK_B] is live
677 assert session.item_for(TRACK_B) is live
678 # and it is not queued as pending, because it already started
679 assert session.has_pending is False
680
681
682async def test_seeking_the_playing_item_restarts_the_session(
683 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
684) -> None:
685 """
686 A seek re-opens the item that is playing, which is a restart of the session.
687
688 A realtime source has not captured anything past the play position, so any
689 forward seek lands outside the buffer and comes back here.
690 """
691 backend = _make_backend(tmp_path)
692 backend._server = MagicMock()
693 backend._binary = Path("/nonexistent/soloist")
694 session = _SoloistSession(backend, "player1")
695 backend._session = session
696 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
697 item.started.set()
698 # its own stream is still attached when the seek re-opens it
699 item.claim()
700 stopped = AsyncMock()
701 monkeypatch.setattr(session, "stop", stopped)
702 _install_fake_binary_manager(monkeypatch)
703 monkeypatch.setattr(
704 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
705 )
706 with pytest.raises(AudioError, match="spawn"):
707 await backend._acquire(TRACK_A, 90, "player1")
708 stopped.assert_awaited_once()
709
710
711@pytest.mark.parametrize(
712 ("requested", "other_queue"),
713 [
714 # another player, whatever it asks for - including the very track this
715 # session is in the middle of delivering
716 pytest.param(TRACK_B, "player2", id="other_player"),
717 pytest.param(TRACK_A, "player2", id="other_player_same_track"),
718 # an early fetch across a boundary this session does not drive, such as a
719 # podcast episode or audiobook chapter
720 pytest.param(TRACK_B, "player1", id="unstitched_boundary"),
721 ],
722)
723async def test_a_session_in_use_is_never_cut_short(
724 tmp_path: Path, requested: str, other_queue: str
725) -> None:
726 """
727 An item the session cannot serve must not stop one it is still delivering.
728
729 Reported as capacity, so a speculative prepare gives up softly.
730 """
731 backend = _make_backend(tmp_path)
732 backend._server = MagicMock()
733 backend._binary = Path("/nonexistent/soloist")
734 session = _SoloistSession(backend, "player1")
735 backend._session = session
736 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
737 item.started.set()
738 item.claim()
739 # the session really is playing TRACK_A, so a same-track request from another
740 # player cannot be mistaken for a seek
741 session._current = item
742 with pytest.raises(ProviderStreamLimitError) as err:
743 await backend._acquire(requested, 0, other_queue)
744 # a stream-limit error so the item is not marked unplayable, but the message
745 # is about the session, not the provider's source-stream budget
746 assert err.value.limit == 1
747 assert err.value.translation_key == "soloist_session_busy"
748 # the session that was playing is untouched
749 assert backend._session is session
750 assert session.usable is True
751
752
753async def test_a_session_nobody_reads_is_replaced_for_another_item(
754 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
755) -> None:
756 """Once the other item has been released, the same request gets the session."""
757 backend = _make_backend(tmp_path)
758 backend._server = MagicMock()
759 backend._binary = Path("/nonexistent/soloist")
760 session = _SoloistSession(backend, "player1")
761 backend._session = session
762 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
763 item.started.set()
764 item.claim()
765 item.close()
766 item.release()
767 _install_fake_binary_manager(monkeypatch)
768 monkeypatch.setattr(
769 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
770 )
771 monkeypatch.setattr(session, "stop", AsyncMock())
772 with pytest.raises(AudioError, match="spawn"):
773 await backend._acquire(TRACK_B, 0, "player1")
774
775
776async def test_a_replacement_waits_for_the_old_daemon_to_be_gone(
777 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
778) -> None:
779 """The engine refuses to start while another daemon still holds its data dir."""
780 backend = _make_backend(tmp_path)
781 backend._server = MagicMock()
782 backend._binary = Path("/nonexistent/soloist")
783 session = _SoloistSession(backend, "player1")
784 backend._session = session
785 order: list[str] = []
786
787 async def _slow_stop() -> None:
788 order.append("stop-start")
789 await asyncio.sleep(0.05)
790 order.append("stop-done")
791
792 monkeypatch.setattr(session, "stop", _slow_stop)
793 _install_fake_binary_manager(monkeypatch)
794
795 async def _spawn(_self: Any, _uri: str, _seek: int) -> None:
796 order.append("spawn")
797 raise AudioError("spawn")
798
799 monkeypatch.setattr(soloist_backend._SoloistSession, "start", _spawn)
800 # the session failed, so its teardown is under way when the next item arrives
801 discard = asyncio.create_task(backend.discard_session(session))
802 await asyncio.sleep(0)
803 with pytest.raises(AudioError, match="spawn"):
804 await backend._acquire(TRACK_B, 0, "player1")
805 await discard
806 assert order == ["stop-start", "stop-done", "spawn"]
807
808
809async def test_an_idle_session_is_taken_over(
810 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
811) -> None:
812 """A session nobody is reading is replaced instead of blocking another player."""
813 backend = _make_backend(tmp_path)
814 backend._server = MagicMock()
815 backend._binary = Path("/nonexistent/soloist")
816 session = _SoloistSession(backend, "player1")
817 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
818 backend._session = session
819 stopped = AsyncMock()
820 monkeypatch.setattr(session, "stop", stopped)
821 _install_fake_binary_manager(monkeypatch)
822 # the replacement spawn is out of scope here; only the takeover decision is
823 monkeypatch.setattr(
824 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
825 )
826 with pytest.raises(AudioError, match="spawn"):
827 await backend._acquire(TRACK_B, 0, "player2")
828 stopped.assert_awaited_once()
829
830
831def test_a_dead_session_task_fails_the_session(tmp_path: Path) -> None:
832 """A session task that dies of an unexpected error takes the session with it."""
833 session = _make_session(tmp_path)
834 task: Any = MagicMock()
835 task.cancelled.return_value = False
836 task.exception.return_value = RuntimeError("reader blew up")
837 session._task_done(task)
838 assert session.usable is False
839 assert session._error is not None
840 assert "reader blew up" in session._error
841
842
843def test_a_cancelled_session_task_is_not_a_failure(tmp_path: Path) -> None:
844 """Teardown cancels the session's tasks; that must not be reported as an error."""
845 session = _make_session(tmp_path)
846 task: Any = MagicMock()
847 task.cancelled.return_value = True
848 session._task_done(task)
849 assert session.usable is True
850
851
852async def test_failed_sink_control_fails_the_session(tmp_path: Path) -> None:
853 """A failed suspend/resume fails the session instead of leaking stall silence."""
854 session = _make_session(tmp_path)
855 session._demand_started = True
856 session._sink_running = True
857 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
858 session._pending.append(TRACK_B)
859 _sink_of(session).suspend.side_effect = RuntimeError("pactl failed")
860 await session._handle_event(_playback_event("buffering"))
861 assert session._error is not None
862 assert "capture sink control failed" in session._error
863
864
865async def test_app_pause_is_fought_with_a_resume(tmp_path: Path) -> None:
866 """A pause from the Spotify app is undone: this session has no user-facing pause."""
867 session = _make_session(tmp_path)
868 session._demand_started = True
869 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
870 session._pending.append(TRACK_B)
871 await session._handle_event(_playback_event("paused"))
872 _client_of(session).resume.assert_awaited_once()
873
874
875async def test_an_app_pause_is_only_undone_so_many_times(tmp_path: Path) -> None:
876 """Someone who keeps pausing means it: the session gives up instead of fighting on."""
877 session = _make_session(tmp_path)
878 session._demand_started = True
879 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
880 session._pending.append(TRACK_B)
881 for _ in range(_MAX_APP_PAUSE_RESUMES):
882 await session._handle_event(_playback_event("playing"))
883 await session._handle_event(_playback_event("paused"))
884 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
885 assert session._error is None
886
887 await session._handle_event(_playback_event("playing"))
888 await session._handle_event(_playback_event("paused"))
889 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
890 assert session.usable is False
891 assert session._app_control is SoloistAppControl.PAUSED
892
893
894async def test_one_pause_reported_twice_counts_once(tmp_path: Path) -> None:
895 """A repeated snapshot of the same pause is not a new pause."""
896 session = _make_session(tmp_path)
897 session._demand_started = True
898 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
899 session._pending.append(TRACK_B)
900 for _ in range(_MAX_APP_PAUSE_RESUMES + 2):
901 await session._handle_event(_playback_event("paused"))
902 assert session.usable is True
903
904
905async def test_the_pause_budget_resets_on_the_next_item(tmp_path: Path) -> None:
906 """Each item gets its own budget; pausing one track does not spend the next one's."""
907 session = _make_session(tmp_path)
908 session._demand_started = True
909 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
910 session._pending.append(TRACK_B)
911 for _ in range(_MAX_APP_PAUSE_RESUMES):
912 await session._handle_event(_playback_event("playing"))
913 await session._handle_event(_playback_event("paused"))
914 await session._observe_current(TRACK_B, 200_000)
915 assert session._app_pauses == 0
916
917
918async def test_a_pause_is_not_undone_once_the_device_is_gone(tmp_path: Path) -> None:
919 """A bare resume on a device Spotify no longer routes to would play to nobody."""
920 session = _make_session(tmp_path)
921 session._demand_started = True
922 session._was_active = False
923 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
924 session._pending.append(TRACK_B)
925 await session._handle_event(_playback_event("paused"))
926 _client_of(session).resume.assert_not_awaited()
927
928
929async def test_losing_the_active_device_ends_the_session(tmp_path: Path) -> None:
930 """Playback moved to another device from the Spotify app: this session is over."""
931 session = _make_session(tmp_path)
932 await session._handle_event(_device_event(is_active=False))
933 assert session.usable is False
934 assert session._app_control is SoloistAppControl.TOOK_OVER
935
936
937async def test_a_takeover_reported_on_the_auth_state_ends_the_session(tmp_path: Path) -> None:
938 """The active-device state also rides on auth_state, and counts the same there."""
939 session = _make_session(tmp_path)
940 await session._handle_event(_auth_event(logged_in=True, is_active=False))
941 assert session.usable is False
942
943
944async def test_an_inactive_device_before_activation_is_not_a_takeover(tmp_path: Path) -> None:
945 """A fresh daemon is inactive until the session claims it; that is not a takeover."""
946 session = _make_session(tmp_path)
947 session._was_active = False
948 await session._handle_event(_device_event(is_active=False))
949 await session._handle_event(_auth_event(logged_in=True, is_active=False))
950 assert session.usable is True
951
952 # nor does a respawned daemon reporting the session Spotify still has for
953 # the account: only the status _play claimed is followed
954 await session._handle_event(_device_event(is_active=True))
955 await session._handle_event(_device_event(is_active=False))
956 assert session.usable is True
957
958
959async def test_a_reconnect_snapshot_keeps_an_active_session_alive(tmp_path: Path) -> None:
960 """The events connection re-snapshots after a drop; that is not a device change."""
961 session = _make_session(tmp_path)
962 await session._handle_event(_auth_event(logged_in=True, is_active=True))
963 await session._handle_event(_device_event(is_active=True))
964 assert session.usable is True
965
966
967async def test_the_playback_snapshots_active_flag_is_ignored(tmp_path: Path) -> None:
968 """It is optional and rides on deltas, so only the dedicated reports are followed."""
969 session = _make_session(tmp_path)
970 session._demand_started = True
971 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
972 await session._handle_event(
973 SoloistEvent(
974 type="playback_changed",
975 data=SoloistPlaybackState(status="playing", is_active=False),
976 raw={},
977 )
978 )
979 assert session.usable is True
980
981
982async def test_backpressure_does_not_spend_the_pause_budget(tmp_path: Path) -> None:
983 """A sink suspended to cap the cushion is our doing, not the user pausing."""
984 session = _make_session(tmp_path)
985 session._demand_started = True
986 session._engine_playing = True
987 session._backpressured = True
988 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
989 session._pending.append(TRACK_B)
990 await session._handle_event(_playback_event("paused"))
991 _client_of(session).resume.assert_not_awaited()
992 assert session._app_pauses == 0
993
994
995async def test_a_lost_login_is_not_reported_as_a_takeover(tmp_path: Path) -> None:
996 """Losing the login wins over the inactive device it brings with it."""
997 session = _make_session(tmp_path)
998 await session._handle_event(_auth_event(logged_in=False, is_active=False))
999 assert session.usable is False
1000 assert session._app_control is None
1001
1002
1003async def test_a_track_started_from_the_app_ends_the_session(tmp_path: Path) -> None:
1004 """The engine pulled off an item part-way through is the app playing something else."""
1005 session = _make_session(tmp_path)
1006 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1007 item.duration_ms = 200_000
1008 await session._observe_current(TRACK_A, 200_000)
1009 item.observe_position(20_000)
1010
1011 await session._observe_current("spotify:track:theirs", 180_000)
1012 assert session.usable is False
1013 assert session._app_control is SoloistAppControl.TOOK_OVER
1014 assert session.current is item
1015
1016
1017async def test_a_track_played_earlier_started_from_the_app_ends_the_session(
1018 tmp_path: Path,
1019) -> None:
1020 """A known uri is no exemption: only the item fed behind this one is where we sent it."""
1021 session = _make_session(tmp_path)
1022 played = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1023 played.spent = True
1024 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1025 item.duration_ms = 200_000
1026 await session._observe_current(TRACK_A, 200_000)
1027 item.observe_position(20_000)
1028
1029 await session._observe_current(TRACK_B, 180_000)
1030 assert session.usable is False
1031 assert session._app_control is SoloistAppControl.TOOK_OVER
1032
1033
1034async def test_skipping_from_the_app_to_the_fed_item_is_followed(tmp_path: Path) -> None:
1035 """The queue moves to that same track, so following the engine keeps the two in step."""
1036 session = _make_session(tmp_path)
1037 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1038 item.duration_ms = 200_000
1039 await session._observe_current(TRACK_A, 200_000)
1040 item.observe_position(20_000)
1041 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1042 session._pending.append(TRACK_B)
1043
1044 await session._observe_current(TRACK_B, 180_000)
1045 assert session.usable is True
1046 assert session.current is fed
1047 assert session.item_for(TRACK_B) is fed
1048
1049
1050async def test_a_takeover_snapshot_stops_pinning_volume_and_options(tmp_path: Path) -> None:
1051 """Once the app has the session, the rest of its snapshot must not reach the daemon."""
1052 session = _make_session(tmp_path)
1053 session._demand_started = True
1054 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1055 item.duration_ms = 200_000
1056 await session._observe_current(TRACK_A, 200_000)
1057 item.observe_position(20_000)
1058
1059 await session._handle_event(
1060 SoloistEvent(
1061 type="playback_changed",
1062 data=SoloistPlaybackState(
1063 status="playing",
1064 item=SoloistEntity(uri="spotify:track:theirs", entity_type="track"),
1065 volume=40,
1066 options=SoloistPlaybackOptions(shuffle=True, repeat="context"),
1067 ),
1068 raw={},
1069 )
1070 )
1071 assert session.usable is False
1072 _client_of(session).set_volume.assert_not_awaited()
1073 _client_of(session).set_shuffle.assert_not_awaited()
1074
1075
1076async def test_the_engine_moving_on_at_a_track_end_is_not_a_takeover(tmp_path: Path) -> None:
1077 """An unasked-for item the engine reaches at a boundary is its own autoplay."""
1078 session = _make_session(tmp_path)
1079 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1080 item.duration_ms = 200_000
1081 await session._observe_current(TRACK_A, 200_000)
1082 item.observe_position(200_000)
1083
1084 await session._observe_current("spotify:track:autoplay", 180_000)
1085 assert session.usable is True
1086 assert session.item_for("spotify:track:autoplay") is None
1087
1088
1089async def test_an_ended_item_says_what_the_app_did(tmp_path: Path) -> None:
1090 """The item's stream fails with the takeover, not a generic session error."""
1091 session = _make_session(tmp_path)
1092 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1093 await session._handle_event(_device_event(is_active=False))
1094 with pytest.raises(SoloistAppControlError) as err:
1095 await session.validate_item(item)
1096 assert err.value.translation_key == SoloistAppControl.TOOK_OVER.value
1097 assert isinstance(err.value, ProviderStreamLimitError)
1098
1099
1100async def test_a_session_being_torn_down_does_not_hold_off_the_next_one(tmp_path: Path) -> None:
1101 """Teardown pauses the daemon; that must not read as the user pausing."""
1102 session = _make_session(tmp_path)
1103 session._demand_started = True
1104 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1105 session._pending.append(TRACK_B)
1106 session._stopped = True
1107 for _ in range(_MAX_APP_PAUSE_RESUMES + 1):
1108 await session._handle_event(_playback_event("playing"))
1109 await session._handle_event(_playback_event("paused"))
1110 await session._handle_event(_device_event(is_active=False))
1111 session.backend._raise_if_app_controlled()
1112 _client_of(session).resume.assert_not_awaited()
1113
1114
1115async def test_no_session_is_started_while_the_app_holds_the_last_one(tmp_path: Path) -> None:
1116 """A replacement would claim the Connect device straight back off the user."""
1117 backend = _make_backend(tmp_path)
1118 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1119 with pytest.raises(SoloistAppControlError):
1120 await backend._acquire(TRACK_A, 0, "player1")
1121
1122
1123async def test_the_hold_on_a_new_session_expires(tmp_path: Path) -> None:
1124 """Coming back to Music Assistant later plays again without any fuss."""
1125 backend = _make_backend(tmp_path)
1126 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1127 backend._app_control_until = time.monotonic() - 1
1128 backend._raise_if_app_controlled()
1129 assert backend._held_by_app() is None
1130
1131
1132async def test_an_audiobook_gives_up_on_capacity_instead_of_burning_chapters(
1133 tmp_path: Path,
1134) -> None:
1135 """Skipping ahead would cost the audiobook its availability and the caller its retry."""
1136 provider = _make_provider(tmp_path)
1137 calls: list[str] = []
1138
1139 async def _refuse(uri: str, *_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
1140 calls.append(uri)
1141 for _ in (): # never yields; only makes this an async generator
1142 yield b""
1143 raise SoloistAppControlError(provider, SoloistAppControl.TOOK_OVER)
1144
1145 provider.backend = MagicMock(stream_spotify_uri=_refuse)
1146 streamdetails = MagicMock(
1147 media_type=MediaType.AUDIOBOOK,
1148 data={"chapters": [TRACK_A, TRACK_B, "spotify:track:ccc"], "chapters_data": []},
1149 )
1150
1151 with pytest.raises(SoloistAppControlError):
1152 async for _ in provider.get_audio_stream(streamdetails):
1153 pass
1154 # the first chapter's refusal ends it: no chapter is skipped over
1155 assert calls == [TRACK_A]
1156
1157
1158def test_the_playback_device_is_named_apart_from_the_connect_one() -> None:
1159 """Two identically named devices in the Spotify app is what causes the takeovers."""
1160 assert SOLOIST_DEVICE_NAME != DEFAULT_PUBLISH_NAME
1161
1162
1163async def test_app_volume_change_is_pinned_back_to_unity(tmp_path: Path) -> None:
1164 """An off-unity volume set from the Spotify app is pinned back to 100."""
1165 session = _make_session(tmp_path)
1166 await session._handle_event(
1167 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=40), raw={})
1168 )
1169 _client_of(session).set_volume.assert_awaited_once_with(100)
1170 _client_of(session).set_volume.reset_mock()
1171 await session._handle_event(
1172 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=100), raw={})
1173 )
1174 _client_of(session).set_volume.assert_not_awaited()
1175
1176
1177async def test_track_change_signals_the_queue_when_it_matches_the_next_item(
1178 tmp_path: Path,
1179) -> None:
1180 """Reaching a fed item tells the queue to start filling that item's buffer."""
1181 session = _make_session(tmp_path, queue_id="player1")
1182 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1183 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1184 queues = _queues_of(session)
1185 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1186 await session._handle_event(
1187 SoloistEvent(
1188 type="track_changed",
1189 data=SoloistTrackChanged(item=SoloistEntity(uri=TRACK_B, entity_type="track")),
1190 raw={},
1191 )
1192 )
1193 queues.prepare_next_audio_buffer.assert_called_once_with("player1")
1194
1195
1196async def test_track_change_to_another_item_signals_nothing(tmp_path: Path) -> None:
1197 """An item the queue is not asking for next must not trigger a prebuffer."""
1198 session = _make_session(tmp_path, queue_id="player1")
1199 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1200 queues = _queues_of(session)
1201 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1202 await session._handle_event(
1203 SoloistEvent(
1204 type="track_changed",
1205 data=SoloistTrackChanged(
1206 item=SoloistEntity(uri="spotify:track:surprise", entity_type="track")
1207 ),
1208 raw={},
1209 )
1210 )
1211 queues.prepare_next_audio_buffer.assert_not_called()
1212
1213
1214async def test_the_follower_of_the_streamed_item_is_fed(tmp_path: Path) -> None:
1215 """The item after the one being streamed is handed to the engine."""
1216 session = _make_session(tmp_path, queue_id="player1")
1217 streamdetails = MagicMock()
1218 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1219 follower = _queue_item(TRACK_B)
1220 queues = _queues_of(session)
1221 queues.get.return_value = MagicMock(current_index=3)
1222 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 3 else None
1223 queues.get_next_item.return_value = follower
1224 await session.feed_after(streamdetails, TRACK_A)
1225 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1226 assert TRACK_B in session._items
1227 assert session.has_pending is True
1228
1229
1230async def test_an_item_the_queue_resolved_elsewhere_is_not_fed(tmp_path: Path) -> None:
1231 """A track the queue will stream from another provider must not be queued here."""
1232 session = _make_session(tmp_path, queue_id="player1")
1233 streamdetails = MagicMock()
1234 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1235 # same track, but the queue already picked a different provider for it
1236 follower = _queue_item(TRACK_B, streamdetails=MagicMock(provider="tidal--x"))
1237 queues = _queues_of(session)
1238 queues.get.return_value = MagicMock(current_index=0)
1239 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1240 queues.get_next_item.return_value = follower
1241 await session.feed_after(streamdetails, TRACK_A)
1242 _client_of(session).add_to_queue.assert_not_awaited()
1243
1244
1245async def test_skipping_to_the_fed_item_keeps_the_session(tmp_path: Path) -> None:
1246 """A next-track lands on the item already fed, so the engine jumps instead of respawning."""
1247 backend = _make_backend(tmp_path)
1248 backend._server = MagicMock()
1249 backend._binary = Path("/nonexistent/soloist")
1250 session = _SoloistSession(backend, "player1")
1251 session._client = AsyncMock()
1252 session._logged_in = True
1253 backend._session = session
1254 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1255 playing.started.set()
1256 # fed one ahead and not reached yet, which is where a next-track goes
1257 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1258 session._pending.append(TRACK_B)
1259
1260 async def _engine_gets_there(**_kwargs: Any) -> None:
1261 await session._observe_current(TRACK_B, 200_000)
1262
1263 _client_of(session).skip_next.side_effect = _engine_gets_there
1264 got_session, got_item = await backend._acquire(TRACK_B, 0, "player1")
1265 # the same session, no respawn, and the item that was already queued
1266 assert got_session is session
1267 assert got_item is fed
1268 assert backend._session is session
1269 _client_of(session).skip_next.assert_awaited_once()
1270
1271
1272async def test_a_skip_drops_what_arrives_while_the_command_is_in_flight(
1273 tmp_path: Path,
1274) -> None:
1275 """
1276 Audio captured between the skip command and the engine's answer is dropped.
1277
1278 Only covers the marker's own window; what the pipeline still holds when the
1279 answer arrives is measured at the cut instead.
1280 """
1281 session = _make_session(tmp_path)
1282 leaving = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1283 leaving.started.set()
1284 target = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1285 session._pending.append(TRACK_B)
1286 captured: list[bytes] = []
1287
1288 async def _engine_gets_there(**_kwargs: Any) -> None:
1289 # the pipeline still holds the old track while the command is in flight
1290 session._write_if_wanted(b"\x01" * 32)
1291 await session._observe_current(TRACK_B, 200_000)
1292 # from here on the audio really is the new item's
1293 session._write_if_wanted(b"\x02" * 32)
1294
1295 _client_of(session).skip_next.side_effect = _engine_gets_there
1296 await session.skip_to(target)
1297 captured.extend(target._chunks)
1298 assert b"".join(captured) == b"\x02" * 32
1299 assert session._discard_until is None
1300
1301
1302async def test_a_skip_the_engine_never_reaches_fails(
1303 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1304) -> None:
1305 """A skip that does not land is an error, not a wait for the track to end."""
1306 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
1307 session = _make_session(tmp_path)
1308 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1309 session._pending.append(TRACK_B)
1310 with pytest.raises(AudioError, match="did not reach"):
1311 await session.skip_to(fed)
1312
1313
1314async def test_a_fed_item_the_engine_has_not_reached_is_not_served(tmp_path: Path) -> None:
1315 """Skipping to an already-fed item must not hand over a channel that fills later."""
1316 session = _make_session(tmp_path)
1317 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1318 # fed, but the engine is still on the previous track
1319 assert session.item_for(TRACK_B) is None
1320 # once the engine gets there it is servable
1321 await session._observe_current(TRACK_B, 200_000)
1322 assert session.item_for(TRACK_B) is fed
1323
1324
1325def test_the_shaper_only_emits_whole_frames() -> None:
1326 """A read that ends mid-frame must never split a frame across two items."""
1327 shaper = soloist_backend._CaptureShaper()
1328 # the session's first bytes are infrastructure silence, and are dropped
1329 assert shaper.shape(b"\x00" * 4096) == b""
1330 # a mis-aligned read emits whole frames and carries the remainder
1331 first = shaper.shape(b"\x01" * (_FRAME_BYTES + 3))
1332 assert len(first) == _FRAME_BYTES
1333 # which is then completed by the next read, losing nothing
1334 second = shaper.shape(b"\x02" * (_FRAME_BYTES - 3))
1335 assert len(second) == _FRAME_BYTES
1336 assert second[:3] == b"\x01" * 3
1337 # an aligned read passes straight through
1338 assert shaper.shape(b"\x03" * _FRAME_BYTES) == b"\x03" * _FRAME_BYTES
1339
1340
1341def test_the_shaper_trims_lead_silence_only_once() -> None:
1342 """Silence after the audio has started is content, not pre-roll."""
1343 shaper = soloist_backend._CaptureShaper()
1344 assert shaper.shape(b"\x01" * _FRAME_BYTES) == b"\x01" * _FRAME_BYTES
1345 silence = b"\x00" * _FRAME_BYTES
1346 assert shaper.shape(silence) == silence
1347
1348
1349async def test_only_whole_sample_frames_are_handed_over(
1350 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1351) -> None:
1352 """A read that ends mid-frame must not split a frame across two items."""
1353 session = _make_session(tmp_path)
1354 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1355 item.started.set()
1356 item.claim()
1357 session._demand_started = True
1358 session._sink_running = True
1359 # two reads that are each mis-aligned but whole together
1360 reads = [b"\x01" * (_FRAME_BYTES + 3), b"\x02" * (_FRAME_BYTES - 3), b""]
1361 reader = MagicMock()
1362
1363 async def _read(_size: int) -> bytes:
1364 return reads.pop(0) if reads else b""
1365
1366 reader.read = _read
1367 session._reader = reader
1368 monkeypatch.setattr(soloist_backend, "_PACE_RATE", 1000.0)
1369 await session._read_capture()
1370 # every write was frame-aligned, and no byte was lost
1371 assert item.buffered % _FRAME_BYTES == 0
1372 assert item.buffered == _FRAME_BYTES * 2
1373
1374
1375async def test_an_already_known_item_is_not_fed_twice(tmp_path: Path) -> None:
1376 """An item the session already plays or was fed is not queued again."""
1377 session = _make_session(tmp_path, queue_id="player1")
1378 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1379 streamdetails = MagicMock()
1380 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
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 = _queue_item(TRACK_B)
1385 await session.feed_after(streamdetails, TRACK_A)
1386 _client_of(session).add_to_queue.assert_not_awaited()
1387
1388
1389async def test_only_tracks_are_fed_ahead(tmp_path: Path) -> None:
1390 """A podcast episode or audiobook chapter is played on its own, never stitched."""
1391 session = _make_session(tmp_path, queue_id="player1")
1392 await session.feed_after(MagicMock(), "spotify:episode:xyz")
1393 _client_of(session).add_to_queue.assert_not_awaited()
1394
1395
1396async def test_a_non_spotify_follower_is_not_fed(tmp_path: Path) -> None:
1397 """The run simply ends where the queue leaves this provider."""
1398 session = _make_session(tmp_path, queue_id="player1")
1399 streamdetails = MagicMock()
1400 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1401 follower = MagicMock(
1402 media_item=MagicMock(media_type=MediaType.TRACK, provider="tidal--x"), streamdetails=None
1403 )
1404 follower.media_item.provider_mappings = []
1405 queues = _queues_of(session)
1406 queues.get.return_value = MagicMock(current_index=0)
1407 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1408 queues.get_next_item.return_value = follower
1409 await session.feed_after(streamdetails, TRACK_A)
1410 _client_of(session).add_to_queue.assert_not_awaited()
1411
1412
1413async def test_a_library_item_is_fed_through_its_spotify_mapping(tmp_path: Path) -> None:
1414 """A library track is fed with the item id this provider instance knows it by."""
1415 session = _make_session(tmp_path, queue_id="player1")
1416 streamdetails = MagicMock()
1417 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1418 follower = MagicMock(
1419 media_item=MagicMock(media_type=MediaType.TRACK, provider="library", item_id="42"),
1420 streamdetails=None,
1421 )
1422 follower.media_item.provider_mappings = [
1423 MagicMock(provider_instance="other--y", item_id="wrong"),
1424 MagicMock(provider_instance="spotify--test", item_id="bbb"),
1425 ]
1426 queues = _queues_of(session)
1427 queues.get.return_value = MagicMock(current_index=0)
1428 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1429 queues.get_next_item.return_value = follower
1430 await session.feed_after(streamdetails, TRACK_A)
1431 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1432
1433
1434@pytest.mark.parametrize(
1435 ("provider_option", "player_setting", "expected"),
1436 [
1437 (True, "enabled", True),
1438 # the player's own switch decides first: off means nobody normalizes,
1439 # not that the job passes to Spotify
1440 (True, "disabled", False),
1441 (False, "enabled", False),
1442 (False, "disabled", False),
1443 ],
1444)
1445def test_who_normalizes_needs_both_switches(
1446 tmp_path: Path,
1447 monkeypatch: pytest.MonkeyPatch,
1448 provider_option: bool,
1449 player_setting: str,
1450 expected: bool,
1451) -> None:
1452 """The engine normalizes only when the provider option and the player agree."""
1453 session = _make_session(tmp_path, queue_id="player1")
1454 monkeypatch.setattr(
1455 type(session.backend.provider),
1456 "spotify_normalization_configured",
1457 property(lambda _self: provider_option),
1458 )
1459 cast("MagicMock", session.mass.config).get_effective_player_queue_config_value = MagicMock(
1460 return_value=player_setting
1461 )
1462 assert session._engine_normalization_enabled() is expected
1463
1464
1465def test_a_running_session_answers_for_what_the_engine_is_doing(tmp_path: Path) -> None:
1466 """
1467 The engine reads its settings at startup, so a later toggle must not split them.
1468
1469 Otherwise the streams core would start normalizing on top of audio the engine
1470 is still normalizing, or stop while it no longer is.
1471 """
1472 backend = _make_backend(tmp_path)
1473 provider = backend.provider
1474 streamdetails = _streamdetails_for(queue_id="player1")
1475 # nothing playing yet: the configuration is all there is to go on
1476 before_any_session = backend.session_normalizes(streamdetails)
1477 session = _SoloistSession(backend, "player1")
1478 session.engine_normalizes = True
1479 backend._session = session
1480 while_playing = backend.session_normalizes(streamdetails)
1481 # ... and a session that has been torn down no longer speaks for the engine
1482 session._stopped = True
1483 after_teardown = backend.session_normalizes(streamdetails)
1484 assert before_any_session is None
1485 assert while_playing is True
1486 assert after_teardown is None
1487 assert (
1488 provider.delivers_normalized_audio(streamdetails)
1489 is provider.spotify_normalization_configured
1490 )
1491
1492
1493async def test_short_delivery_is_rejected_as_incomplete(tmp_path: Path) -> None:
1494 """PCM that stops well short of the item's duration is rejected."""
1495 session = _make_session(tmp_path)
1496 item = _ItemAudio(TRACK_A, session)
1497 item.playing_seen = True
1498 item.duration_ms = 200_000
1499 item.last_position_ms = 100_000
1500 with pytest.raises(AudioError, match="incomplete"):
1501 await session.validate_item(item)
1502
1503
1504async def test_missing_position_is_rejected_as_incomplete(tmp_path: Path) -> None:
1505 """Without any position report there is no evidence the item played out."""
1506 session = _make_session(tmp_path)
1507 item = _ItemAudio(TRACK_A, session)
1508 item.playing_seen = True
1509 item.duration_ms = 200_000
1510 with pytest.raises(AudioError, match="incomplete"):
1511 await session.validate_item(item)
1512
1513
1514async def test_short_item_cannot_pass_at_position_zero(tmp_path: Path) -> None:
1515 """The tolerance never spans a whole item, so a short item cannot pass unplayed."""
1516 session = _make_session(tmp_path)
1517 item = _ItemAudio(TRACK_A, session)
1518 item.playing_seen = True
1519 item.duration_ms = 8_000
1520 item.last_position_ms = 0
1521 with pytest.raises(AudioError, match="incomplete"):
1522 await session.validate_item(item)
1523
1524
1525async def test_an_item_that_never_played_is_rejected(tmp_path: Path) -> None:
1526 """An item the engine never reported playing is a failure, whatever was delivered."""
1527 session = _make_session(tmp_path)
1528 item = _ItemAudio(TRACK_A, session)
1529 item.duration_ms = 200_000
1530 item.last_position_ms = 200_000
1531 with pytest.raises(AudioError, match="never started playing"):
1532 await session.validate_item(item)
1533
1534
1535async def test_a_duration_less_item_is_not_judged(tmp_path: Path) -> None:
1536 """Without a duration there is nothing to judge completeness against."""
1537 session = _make_session(tmp_path)
1538 item = _ItemAudio(TRACK_A, session)
1539 item.playing_seen = True
1540 await session.validate_item(item)
1541
1542
1543def test_an_unread_session_expires(tmp_path: Path) -> None:
1544 """A session no item stream reads from is ended so its daemon does not linger."""
1545 session = _make_session(tmp_path)
1546 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1547 session._expire_idle()
1548 assert session._idle_since is not None
1549 assert session.usable is True
1550 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S - 1
1551 session._expire_idle()
1552 assert session.usable is False
1553
1554
1555def test_a_session_being_read_never_expires(tmp_path: Path) -> None:
1556 """An item stream reading the session keeps it alive indefinitely."""
1557 session = _make_session(tmp_path)
1558 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1559 item.claim()
1560 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S * 10
1561 session._expire_idle()
1562 assert session.usable is True
1563
1564
1565def test_pre_roll_silence_is_dropped_a_whole_frame_at_a_time() -> None:
1566 """
1567 Trimming pre-roll must leave the audio on the session's frame grid.
1568
1569 A FIFO read is not always a whole number of frames, and dropping a partial
1570 one would shift every sample that follows for the rest of the session.
1571 """
1572 shaper = _CaptureShaper()
1573 # pre-roll that ends mid-frame: the real audio starts at byte 1024
1574 assert shaper.shape(b"\x00" * 1021) == b""
1575 audio = bytes(range(1, 9)) * 4
1576 shaped = shaper.shape(b"\x00" * 3 + audio)
1577 assert shaped == audio
1578 assert shaper._lead_skipped % _FRAME_BYTES == 0
1579
1580
1581async def test_a_refused_skip_does_not_leave_the_audio_discarded(tmp_path: Path) -> None:
1582 """
1583 A skip that never landed must not keep the session dropping its audio.
1584
1585 The marker silences everything the session captures, so a command that
1586 failed has to clear it on the way out.
1587 """
1588 session = _make_session(tmp_path)
1589 client = cast("MagicMock", session._client)
1590 client.skip_next = AsyncMock(side_effect=TimeoutError)
1591 item = _ItemAudio(TRACK_B, session)
1592
1593 with pytest.raises(AudioError, match="would not skip"):
1594 await session.skip_to(item)
1595
1596 assert session._discard_until is None
1597
1598
1599async def test_a_daemon_that_will_not_die_is_reported_and_released(tmp_path: Path) -> None:
1600 """A close that could not terminate the daemon still finishes the teardown."""
1601 session = _make_session(tmp_path)
1602 proc = cast("MagicMock", session._proc)
1603 proc.close = AsyncMock()
1604 # AsyncProcess.close() gives up after a handful of kill attempts
1605 proc.returncode = None
1606 with patch.object(session.logger, "warning") as warning:
1607 await session.stop()
1608 assert warning.called
1609 assert session._teardown_done is True
1610 assert session._proc is None
1611
1612
1613async def test_a_cancelled_teardown_still_closes_the_daemon(tmp_path: Path) -> None:
1614 """
1615 A cancelled teardown must leave the retry something to close.
1616
1617 Dropping the references first is how a daemon survives to hold the data
1618 directory, which every later session is then refused for.
1619 """
1620 session = _make_session(tmp_path)
1621 proc = cast("MagicMock", session._proc)
1622 sink = cast("AsyncMock", session._sink)
1623
1624 async def _never_returns() -> None:
1625 await asyncio.Event().wait()
1626
1627 proc.close = _never_returns
1628 task = asyncio.create_task(session.stop())
1629 await asyncio.sleep(0.01)
1630 task.cancel()
1631 with suppress(asyncio.CancelledError):
1632 await task
1633 # the teardown did not finish, so nothing was dropped and it can be redone
1634 unfinished = session._teardown_done
1635 kept_proc = session._proc
1636 kept_sink = session._sink
1637 proc.close = AsyncMock()
1638 proc.returncode = 0
1639 await session.stop()
1640 assert unfinished is False
1641 assert kept_proc is proc
1642 assert kept_sink is sink
1643 assert session._teardown_done is True
1644 assert session._proc is None
1645 assert session._sink is None
1646 proc.close.assert_awaited()
1647 sink.unload.assert_awaited()
1648
1649
1650def test_a_failed_session_is_torn_down(tmp_path: Path) -> None:
1651 """A session that fails is discarded, so its daemon does not keep playing to nobody."""
1652 session = _make_session(tmp_path)
1653 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1654 item.claim()
1655 session._fail("audio stalled")
1656 assert session.usable is False
1657 # every waiting item is released and the teardown is scheduled
1658 assert item._closed is True
1659 # a startup wait must not sit out its timeout on a session that already failed
1660 assert item.started.is_set() is True
1661 discard = cast("MagicMock", session.mass.create_task)
1662 discard.assert_called_once_with(session.backend.discard_session, session)
1663 # a second failure does not queue a second teardown
1664 session._fail("and again")
1665 assert session._error == "audio stalled"
1666 assert discard.call_count == 1
1667
1668
1669async def test_an_item_the_engine_skipped_past_fails_instead_of_hanging(
1670 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1671) -> None:
1672 """A claimed channel the engine never reaches gives up rather than blocking forever."""
1673 monkeypatch.setattr(soloist_backend, "_READ_SLICE_S", 0.01)
1674 monkeypatch.setattr(soloist_backend, "_STALL_TIMEOUT_S", 0.05)
1675 session = _make_session(tmp_path)
1676 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1677 item.claim()
1678 # the engine is playing something else, so nothing is ever written here
1679 session._items["spotify:track:other"] = session._current = _ItemAudio(
1680 "spotify:track:other", session
1681 )
1682 with pytest.raises(AudioError, match="no audio"):
1683 async for _ in item.read():
1684 pass
1685
1686
1687async def test_adopt_paired_session_copies_into_the_canonical_dir(
1688 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1689) -> None:
1690 """A session paired by the setup flow is adopted into the per-instance data dir."""
1691 storage = tmp_path / "storage"
1692 pending = storage / "spotify" / "pairing" / "flow1"
1693 pending.mkdir(parents=True)
1694 (pending / "session.bin").write_bytes(b"session")
1695 prov = _make_provider(tmp_path, {CONF_SOLOIST_SESSION_DIR: "spotify/pairing/flow1"})
1696 update_setup_data = MagicMock()
1697 monkeypatch.setattr(prov, "_update_setup_data", update_setup_data)
1698 backend = SoloistBackend(prov)
1699 await backend._adopt_paired_session()
1700 canonical = storage / "spotify" / "spotify--test" / SOLOIST_DATA_DIR_NAME
1701 assert (canonical / "session.bin").read_bytes() == b"session"
1702 # a copy, not a move: the flow-private source must survive a failed
1703 # provider load so the setup flow can retry (the flow removes it at its end)
1704 assert (pending / "session.bin").exists()
1705 update_setup_data.assert_called_once_with(CONF_SOLOIST_SESSION_DIR, None)
1706
1707
1708def test_the_engine_is_told_not_to_normalize(tmp_path: Path) -> None:
1709 """MA normalizes this audio itself, so the engine's own normalization is switched off."""
1710 backend = _make_backend(tmp_path)
1711 prefs = backend._data_dir / "settings" / "Users" / "alice-user" / "prefs"
1712 prefs.parent.mkdir(parents=True)
1713 prefs.write_text("some.engine.key=1\n", encoding="utf-8")
1714 backend._prepare_data_dir(normalize=False)
1715 content = prefs.read_text(encoding="utf-8").splitlines()
1716 assert "some.engine.key=1" in content
1717 assert "audio.normalize_v2=false" in content
1718 # MA mixes the queue's crossfade itself, so the engine's own is always off
1719 assert "audio.crossfade_v2=false" in content
1720 # the ceiling is stated rather than left to the engine's own default
1721 assert "audio.play_bitrate_enumeration=5" in content
1722 assert "audio.play_bitrate_non_metered_enumeration=5" in content
1723 assert "audio.play_bitrate_non_metered_migrated=true" in content
1724
1725
1726def test_disabling_crossfade_writes_the_boolean(tmp_path: Path) -> None:
1727 """Crossfade off is written explicitly, so a stale 'on' cannot survive."""
1728 backend = _make_backend(tmp_path)
1729 prefs = backend._data_dir / "settings" / "prefs"
1730 prefs.parent.mkdir(parents=True)
1731 prefs.write_text("audio.crossfade_v2=true\naudio.crossfade.time_v2=8000\n", encoding="utf-8")
1732 backend._prepare_data_dir(normalize=False)
1733 content = prefs.read_text(encoding="utf-8").splitlines()
1734 assert "audio.crossfade_v2=false" in content
1735 assert not any(line.startswith("audio.crossfade.time_v2") for line in content)
1736
1737
1738async def test_setup_requires_an_api_key(tmp_path: Path) -> None:
1739 """Without a stored API key the user must be sent back through the setup flow."""
1740 backend = _make_backend(tmp_path)
1741 with pytest.raises(LoginFailed) as err:
1742 await backend.setup()
1743 assert err.value.translation_key == "soloist_pairing_required"
1744
1745
1746async def test_setup_requires_a_paired_session(
1747 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1748) -> None:
1749 """An API key without a paired session also routes back to the setup flow."""
1750 backend = _make_backend(tmp_path, {CONF_SOLOIST_API_KEY: "k" * 20, CONF_SOLOIST_CONSENT: True})
1751 _install_fake_binary_manager(monkeypatch)
1752 with pytest.raises(LoginFailed) as err:
1753 await backend.setup()
1754 assert err.value.translation_key == "soloist_pairing_required"
1755
1756
1757async def test_streaming_without_setup_is_refused(tmp_path: Path) -> None:
1758 """A backend whose setup never ran refuses to stream instead of half-starting."""
1759 backend = _make_backend(tmp_path)
1760 with pytest.raises(AudioError, match="not started"):
1761 async for _ in backend.stream_spotify_uri(TRACK_A):
1762 pass
1763
1764
1765def test_session_present_detection(tmp_path: Path) -> None:
1766 """Only the engine's per-account state counts as paired."""
1767 data_dir = tmp_path / "soloist-data"
1768 assert soloist_session_present(data_dir) is False
1769 data_dir.mkdir()
1770 (data_dir / WS_ADDR_FILE).write_text("127.0.0.1", encoding="utf-8")
1771 (data_dir / WS_PORT_FILE).write_text("1234", encoding="utf-8")
1772 assert soloist_session_present(data_dir) is False
1773 # everything a spawn leaves behind outlives the pairing it ran on: the engine
1774 # keeps its identity, lock, cache and crash handler in the data dir even
1775 # though it is given a cache dir of its own, and Music Assistant writes the
1776 # prefs there before every spawn
1777 (data_dir / "settings").mkdir()
1778 (data_dir / "settings" / "prefs").write_text("audio.normalize_v2=false\n", encoding="utf-8")
1779 (data_dir / ".device_id").write_text("6b6c2a07", encoding="utf-8")
1780 (data_dir / ".lock").write_bytes(b"")
1781 (data_dir / "cache" / "Users" / "spotify-user-user").mkdir(parents=True)
1782 (data_dir / "crashpad").mkdir()
1783 assert soloist_session_present(data_dir) is False
1784 (data_dir / "settings" / "Users" / "spotify-user-user").mkdir(parents=True)
1785 assert soloist_session_present(data_dir) is True
1786
1787
1788async def test_a_skip_drops_the_audio_still_in_flight(tmp_path: Path) -> None:
1789 """The item jumped to opens with its own audio, not the tail of the one left behind."""
1790 session = _make_session(tmp_path)
1791 left_behind = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1792 session._current = left_behind
1793 left_behind.started.set()
1794 left_behind.claim()
1795 jumped_to = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1796 session._pending.append(TRACK_B)
1797 session._discard_until = TRACK_B
1798 with _capture_holding(session, fifo_bytes=2 * _FRAME_BYTES, reader_bytes=2 * _FRAME_BYTES):
1799 await session._observe_current(TRACK_B, 200_000)
1800 assert session._stale_budget == 4 * _FRAME_BYTES
1801 session._write_if_wanted(b"s" * (4 * _FRAME_BYTES))
1802 session._write_if_wanted(b"n" * (2 * _FRAME_BYTES))
1803 jumped_to.claim()
1804 jumped_to.close()
1805 assert b"".join([chunk async for chunk in jumped_to.read()]) == b"n" * (2 * _FRAME_BYTES)
1806
1807
1808async def test_a_skip_drops_the_stale_audio_across_reads(tmp_path: Path) -> None:
1809 """A budget larger than one read keeps dropping, and resumes on a frame boundary."""
1810 session = _make_session(tmp_path)
1811 session._stale_budget = 3 * _FRAME_BYTES
1812 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1813 item.claim()
1814 session._write_if_wanted(b"s" * (2 * _FRAME_BYTES))
1815 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1816 item.close()
1817 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1818
1819
1820async def test_the_marker_spends_an_earlier_jumps_budget(tmp_path: Path) -> None:
1821 """What the marker drops still counts against a budget left from an earlier jump."""
1822 session = _make_session(tmp_path)
1823 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1824 session._stale_budget = 4 * _FRAME_BYTES
1825 session._discard_until = TRACK_B
1826 session._write_if_wanted(b"s" * (3 * _FRAME_BYTES))
1827 assert session._stale_budget == _FRAME_BYTES
1828 # a refused command leaves only what is genuinely still in flight to drop
1829 session._discard_until = None
1830 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1831 item = session._current
1832 item.claim()
1833 item.close()
1834 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1835
1836
1837async def test_a_natural_cut_keeps_the_audio_in_flight(tmp_path: Path) -> None:
1838 """Nothing is dropped without a jump: what is in flight is the continuation."""
1839 session = _make_session(tmp_path)
1840 playing = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1841 session._current = playing
1842 playing.started.set()
1843 playing.claim()
1844 with _capture_holding(session, fifo_bytes=4 * _FRAME_BYTES, reader_bytes=4 * _FRAME_BYTES):
1845 await session._observe_current(TRACK_B, 200_000)
1846 assert session._stale_budget == 0
1847
1848
1849def test_stale_bytes_spans_both_buffers_in_whole_frames(tmp_path: Path) -> None:
1850 """The in-flight measure covers the FIFO and the reader, and never splits a frame."""
1851 session = _make_session(tmp_path)
1852 with _capture_holding(session, fifo_bytes=3 * _FRAME_BYTES + 3, reader_bytes=2 * _FRAME_BYTES):
1853 assert session._stale_bytes() == 5 * _FRAME_BYTES
1854
1855
1856def test_stale_bytes_falls_back_when_the_reader_cannot_be_sized(tmp_path: Path) -> None:
1857 """Losing the reader's internal view drops extra rather than leaving audio behind."""
1858 session = _make_session(tmp_path)
1859 with _capture_holding(session, fifo_bytes=0, reader_bytes=None):
1860 assert session._stale_bytes() == 6 * _READ_CHUNK_SIZE
1861
1862
1863async def test_a_channel_abandoned_at_the_cut_stops_holding_the_cushion(
1864 tmp_path: Path,
1865) -> None:
1866 """A skip closes the channel first and only then unwinds its stream."""
1867 session = _make_session(tmp_path)
1868 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1869 item.started.set()
1870 item.claim()
1871 item.write(b"x" * 4096)
1872 # the cut lands while the abandoned stream is still unwinding
1873 await session._observe_current(TRACK_B, 200_000)
1874 assert session._retained_bytes() == 4096
1875 item.release()
1876 assert session._retained_bytes() == 0
1877
1878
1879def test_an_abandoned_channel_stops_holding_the_cushion(tmp_path: Path) -> None:
1880 """A channel skipped away from frees its buffer instead of gating the sink for good."""
1881 session = _make_session(tmp_path)
1882 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1883 item.claim()
1884 item.write(b"x" * 4096)
1885 assert session._retained_bytes() == 4096
1886 # the stream is gone, then the cut closes the channel
1887 item.release()
1888 item.close()
1889 assert session._retained_bytes() == 0
1890
1891
1892async def test_a_channel_still_being_read_keeps_its_tail(tmp_path: Path) -> None:
1893 """Closing the playing item at a cut must not discard what its stream is still owed."""
1894 session = _make_session(tmp_path)
1895 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1896 item.claim()
1897 item.write(b"tail" * 4)
1898 item.close()
1899 assert item.buffered == 16
1900 assert b"".join([chunk async for chunk in item.read()]) == b"tail" * 4
1901
1902
1903@contextmanager
1904def _capture_holding(
1905 session: _SoloistSession, *, fifo_bytes: int, reader_bytes: int | None
1906) -> Iterator[None]:
1907 """
1908 Give the session a real capture FIFO and a reader holding the given amounts.
1909
1910 A real pipe is used so the byte count comes from the same ioctl the backend
1911 relies on. Pass ``reader_bytes=None`` for a reader whose buffer cannot be read.
1912 """
1913 read_fd, write_fd = os.pipe()
1914 try:
1915 if fifo_bytes:
1916 os.write(write_fd, bytes(fifo_bytes))
1917 pipe = MagicMock()
1918 pipe.fileno.return_value = read_fd
1919 transport = MagicMock()
1920 transport.get_extra_info.return_value = pipe
1921 session._transport = transport
1922 reader = MagicMock(spec=[]) if reader_bytes is None else MagicMock()
1923 if reader_bytes is not None:
1924 reader._buffer = bytearray(reader_bytes)
1925 session._reader = reader
1926 yield
1927 finally:
1928 session._transport = None
1929 session._reader = None
1930 os.close(read_fd)
1931 os.close(write_fd)
1932
1933
1934def _stdout_of(*lines: str) -> MagicMock:
1935 """Return a process mock whose stdout yields the given daemon log lines."""
1936
1937 async def _iter_stdout() -> AsyncGenerator[str]:
1938 for line in lines:
1939 yield line
1940
1941 proc = MagicMock()
1942 proc.iter_stdout = _iter_stdout
1943 return proc
1944
1945
1946def _make_provider(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SpotifyProvider:
1947 """Return a SpotifyProvider (bypassing __init__) with the given setup_data."""
1948 prov = object.__new__(SpotifyProvider)
1949 config = MagicMock(instance_id="spotify--test")
1950 config.get_value = MagicMock(return_value=None)
1951 config.values = {}
1952 prov.config = config
1953 prov.manifest = MagicMock(domain="spotify")
1954 prov.logger = MagicMock()
1955 prov.available = True
1956 mass = MagicMock()
1957 mass.storage_path = str(tmp_path / "storage")
1958 mass.cache_path = str(tmp_path / "cache")
1959 # get_setup_value reads the live setup_data blob from the store
1960 mass.config.get = MagicMock(return_value=setup_data or {})
1961 mass.config.get_raw_provider_config_value = MagicMock(return_value=None)
1962 # the store keeps values encrypted; decrypt is an identity map for the test
1963 mass.config.decrypt_string = MagicMock(side_effect=lambda value: value)
1964 prov.mass = mass
1965 return prov
1966
1967
1968def _make_backend(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SoloistBackend:
1969 """Return a SoloistBackend on a mocked provider."""
1970 return SoloistBackend(_make_provider(tmp_path, setup_data))
1971
1972
1973def _make_session(tmp_path: Path, queue_id: str | None = "player1") -> _SoloistSession:
1974 """Return a session with its process/sink/client replaced by mocks."""
1975 session = _SoloistSession(_make_backend(tmp_path), queue_id)
1976 session._sink = AsyncMock()
1977 session._client = AsyncMock()
1978 session._proc = MagicMock(returncode=None)
1979 # a session under test is past the engine's login and has claimed the
1980 # Connect device, unless a test says otherwise
1981 session._logged_in = True
1982 session._was_active = True
1983 return session
1984
1985
1986def _streamdetails_for(
1987 *,
1988 queue_id: str | None = "player1",
1989 uri: str = TRACK_A,
1990 media_type: MediaType = MediaType.TRACK,
1991) -> StreamDetails:
1992 """Return stream details for a Spotify item served by the test instance."""
1993 return StreamDetails(
1994 provider="spotify--test",
1995 item_id=uri.rsplit(":", 1)[1],
1996 audio_format=AudioFormat(content_type=ContentType.PCM_S16LE),
1997 media_type=media_type,
1998 queue_id=queue_id,
1999 )
2000
2001
2002def _make_item(tmp_path: Path, uri: str) -> _ItemAudio:
2003 """Return a bare item channel on a mocked session."""
2004 return _ItemAudio(uri, _make_session(tmp_path))
2005
2006
2007def _queue_item(uri: str, streamdetails: Any = None) -> MagicMock:
2008 """Return a queue item stand-in for a Spotify track on the test instance."""
2009 item_id = uri.rsplit(":", 1)[1]
2010 media_item = MagicMock(media_type=MediaType.TRACK, provider="spotify--test", item_id=item_id)
2011 media_item.provider_mappings = []
2012 return MagicMock(
2013 media_item=media_item,
2014 queue_item_id=f"qi-{item_id}",
2015 streamdetails=streamdetails,
2016 )
2017
2018
2019async def _wait_for(predicate: Callable[[], bool], timeout: float = 2.0) -> None:
2020 """Wait until the predicate holds, so a background task can get there."""
2021 loop = asyncio.get_running_loop()
2022 deadline = loop.time() + timeout
2023 while loop.time() < deadline:
2024 if predicate():
2025 return
2026 await asyncio.sleep(0.01)
2027 raise AssertionError("condition not met within timeout")
2028
2029
2030def _client_of(session: _SoloistSession) -> AsyncMock:
2031 """Return the session's mocked WebSocket client."""
2032 return cast("AsyncMock", session._client)
2033
2034
2035def _sink_of(session: _SoloistSession) -> AsyncMock:
2036 """Return the session's mocked capture sink."""
2037 return cast("AsyncMock", session._sink)
2038
2039
2040def _queues_of(session: _SoloistSession) -> MagicMock:
2041 """Return the mocked player_queues controller the session consults."""
2042 return cast("MagicMock", session.mass.player_queues)
2043
2044
2045def _auth_event(*, logged_in: bool, is_active: bool = True) -> SoloistEvent:
2046 """Return an auth_state event with the given login and active-device state."""
2047 return SoloistEvent(
2048 type="auth_state",
2049 data=SoloistAuthState(logged_in=logged_in, is_active=is_active),
2050 raw={},
2051 )
2052
2053
2054def _device_event(*, is_active: bool) -> SoloistEvent:
2055 """Return a device_changed event with the given active-device state."""
2056 return SoloistEvent(
2057 type="device_changed", data=SoloistDeviceChanged(is_active=is_active), raw={}
2058 )
2059
2060
2061def _playback_event(status: str, position_ms: int = 0) -> SoloistEvent:
2062 """Return a playback_state event for the current item with the given status."""
2063 return SoloistEvent(
2064 type="playback_state",
2065 data=SoloistPlaybackState(
2066 status=status,
2067 item=SoloistEntity(uri=TRACK_A, entity_type="track"),
2068 position=SoloistPosition(position_ms=position_ms, timestamp_ms=0),
2069 ),
2070 raw={},
2071 )
2072
2073
2074def _install_fake_binary_manager(monkeypatch: pytest.MonkeyPatch) -> None:
2075 """Replace the shared binary manager so no download or exec is attempted."""
2076 manager = MagicMock()
2077 manager.ensure_fresh = AsyncMock(return_value=Path("/nonexistent/soloist"))
2078 monkeypatch.setattr(soloist_backend, "SoloistBinaryManager", MagicMock(return_value=manager))
2079