/
/
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 _ITEM_OVERRUN_S,
37 _MAX_APP_PAUSE_RESUMES,
38 _MAX_LEAD_TRIM_S,
39 _READ_CHUNK_SIZE,
40 SoloistAppControl,
41 SoloistAppControlError,
42 SoloistBackend,
43 _CaptureShaper,
44 _ItemAudio,
45 _SoloistSession,
46 _trim_lead_silence,
47)
48from music_assistant.providers.spotify.constants import (
49 CONF_SOLOIST_API_KEY,
50 CONF_SOLOIST_CONSENT,
51 CONF_SOLOIST_SESSION_DIR,
52 SOLOIST_DATA_DIR_NAME,
53 SOLOIST_DEVICE_NAME,
54)
55from music_assistant.providers.spotify.helpers import soloist_session_present
56from music_assistant.providers.spotify.provider import SpotifyProvider
57from music_assistant.providers.spotify_connect.provider import DEFAULT_PUBLISH_NAME
58from music_assistant.providers.spotify_connect.soloist.runtime import (
59 WS_ADDR_FILE,
60 WS_PORT_FILE,
61 SoloistAuthState,
62 SoloistDeviceChanged,
63 SoloistEntity,
64 SoloistError,
65 SoloistEvent,
66 SoloistOptionsChanged,
67 SoloistPlaybackOptions,
68 SoloistPlaybackState,
69 SoloistPosition,
70 SoloistTrackChanged,
71 SoloistVolumeChanged,
72)
73
74TRACK_A = "spotify:track:aaa"
75TRACK_B = "spotify:track:bbb"
76
77
78def test_trim_drops_an_all_zero_chunk_within_the_bound() -> None:
79 """A pure-silence chunk inside the trim budget is dropped entirely."""
80 chunk = b"\x00" * 1024
81 trimmed, skipped = _trim_lead_silence(chunk, 0)
82 assert trimmed == b""
83 assert skipped == 1024
84
85
86def test_trim_keeps_frame_alignment_when_audio_starts_mid_chunk() -> None:
87 """Audio starting mid-chunk is cut on a sample-frame boundary."""
88 # audio starts one byte into the third frame: the trim must keep that frame whole
89 chunk = b"\x00" * (_FRAME_BYTES * 2 + 1) + b"\x01" * 64
90 trimmed, skipped = _trim_lead_silence(chunk, 0)
91 assert skipped == _FRAME_BYTES * 2
92 assert len(trimmed) % _FRAME_BYTES == 1 # the partial frame's remainder is preserved
93 assert trimmed.endswith(b"\x01" * 64)
94
95
96def test_trim_passes_silence_through_once_the_bound_is_exceeded() -> None:
97 """Beyond the trim budget, silence is genuine content and is delivered."""
98 chunk = b"\x00" * 1024
99 trimmed, skipped = _trim_lead_silence(chunk, int(_MAX_LEAD_TRIM_S * _BYTES_PER_SECOND))
100 assert trimmed == chunk
101 assert skipped == 0
102
103
104def test_seek_is_confirmed_only_within_tolerance(tmp_path: Path) -> None:
105 """A position report confirms a seek only once it reaches the tolerance window."""
106 item = _make_item(tmp_path, TRACK_A)
107 item.arm_seek(60_000)
108 item.observe_position(50_000)
109 assert not item.seek_confirmed.is_set()
110 item.observe_position(58_500)
111 assert item.seek_confirmed.is_set()
112 assert item.started_at_ms == 58_500
113
114
115def test_small_seek_target_is_not_confirmed_by_a_pre_seek_zero_report(tmp_path: Path) -> None:
116 """A position-0 report before the seek lands cannot confirm a small target."""
117 item = _make_item(tmp_path, TRACK_A)
118 item.arm_seek(1_500)
119 item.observe_position(0)
120 assert not item.seek_confirmed.is_set()
121 item.observe_position(1_500)
122 assert item.seek_confirmed.is_set()
123
124
125def test_a_small_seek_is_confirmed_without_a_report_of_exactly_zero(tmp_path: Path) -> None:
126 """A target inside the tolerance window has no room below it to be anchored on."""
127 item = _make_item(tmp_path, TRACK_A)
128 # the engine restored this item a second in, so the seek is short enough
129 # that no report can fall below its tolerance window
130 item.observe_position(1_200)
131 item.arm_seek(2_000)
132 item.observe_position(400)
133 item.observe_position(2_000)
134 assert item.seek_confirmed.is_set()
135
136
137def test_the_restored_position_of_the_same_item_cannot_confirm_a_seek(tmp_path: Path) -> None:
138 """The state a fresh session restores does not pass for the seek landing."""
139 item = _make_item(tmp_path, TRACK_A)
140 item.duration_ms = 176_000
141 # the engine restores the account's last state: this very item, sitting at
142 # the position the seek is aiming for
143 item.observe_position(117_000)
144 item.arm_seek(117_000)
145 item.observe_position(117_000)
146 assert not item.seek_confirmed.is_set()
147 # only once the engine has reloaded the track does its seek count
148 item.observe_position(0)
149 item.observe_position(117_000)
150 assert item.seek_confirmed.is_set()
151
152
153def test_a_backward_seek_is_confirmed_below_where_the_engine_was(tmp_path: Path) -> None:
154 """Seeking back into an item confirms on the target, not on where it came from."""
155 item = _make_item(tmp_path, TRACK_A)
156 # the engine restored this item well past the point being seeked back to
157 item.observe_position(117_000)
158 item.arm_seek(30_000)
159 item.observe_position(0)
160 item.observe_position(30_000)
161 assert item.seek_confirmed.is_set()
162 assert item.started_at_ms == 30_000
163
164
165def test_position_never_regresses_and_stops_at_the_cut(tmp_path: Path) -> None:
166 """The furthest position is kept, and reports after the cut belong to the next item."""
167 item = _make_item(tmp_path, TRACK_A)
168 item.observe_position(120_000)
169 # the engine's stop/idle snapshot at the end of an item reports position 0
170 item.observe_position(0)
171 assert item.last_position_ms == 120_000
172 item.close()
173 item.observe_position(5_000)
174 assert item.last_position_ms == 120_000
175
176
177async def test_item_stream_ends_where_the_session_moves_on(tmp_path: Path) -> None:
178 """An item's audio ends at the track change, and the next item's begins there."""
179 session = _make_session(tmp_path)
180 item_a = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
181 session._current = item_a
182 item_a.started.set()
183 item_a.claim()
184 item_a.write(b"a" * 16)
185 await session._observe_current(TRACK_B, 200_000)
186 item_a.write(b"late" * 4) # written after the cut: goes nowhere
187 chunks = [chunk async for chunk in item_a.read()]
188 assert b"".join(chunks) == b"a" * 16
189 # the next item exists, carries the duration and now receives the audio
190 item_b = session._items[TRACK_B]
191 assert session.current is item_b
192 assert item_b.duration_ms == 200_000
193
194
195async def test_the_engines_restored_state_does_not_cut_a_pending_item(
196 tmp_path: Path,
197) -> None:
198 """A daemon reports the item it restored before playing ours; that is not a boundary."""
199 session = _make_session(tmp_path)
200 requested = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
201 session._current = requested
202 requested.claim()
203 # the engine announces the state it came up with, which is someone else's item
204 await session._observe_current("spotify:track:restored", 152_000)
205 # closing our item here would end its stream before it delivered anything
206 assert requested._closed is False
207 assert requested.started.is_set() is False
208 # ... and the restored item is never offered as an item's audio
209 assert session.item_for("spotify:track:restored") is None
210 # then ours starts for real, and picks up from there
211 await session._observe_current(TRACK_A, 200_000)
212 assert session.current is requested
213 assert requested.started.is_set() is True
214 requested.write(b"\x01" * 32)
215 requested.close()
216 assert b"".join([chunk async for chunk in requested.read()]) == b"\x01" * 32
217
218
219async def test_leaving_the_engines_restored_item_is_not_a_takeover(tmp_path: Path) -> None:
220 """The restored item is part-way through a track, and we are about to leave it."""
221 session = _make_session(tmp_path)
222 # as _play leaves it: the channel exists, its stream is not reading it yet
223 requested = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
224 session._current = requested
225 await session._observe_current("spotify:track:restored", 152_000)
226 restored = session.current
227 assert restored is not None
228 restored.observe_position(20_000)
229
230 # our own play() lands and the engine leaves the restored item for ours
231 await session._observe_current(TRACK_A, 200_000)
232 assert session.usable is True
233 assert session.current is requested
234
235
236async def test_audio_read_before_the_stream_opens_is_kept(tmp_path: Path) -> None:
237 """Audio captured before an item's stream opens is buffered, not dropped."""
238 session = _make_session(tmp_path)
239 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
240 session._current = item
241 item.write(b"head" * 8)
242 item.claim()
243 item.close()
244 chunks = [chunk async for chunk in item.read()]
245 assert b"".join(chunks) == b"head" * 8
246
247
248async def test_a_channel_is_only_ever_served_once(tmp_path: Path) -> None:
249 """A consumed channel cannot be replayed, so the item needs a fresh session."""
250 session = _make_session(tmp_path)
251 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
252 item.started.set()
253 assert session.item_for(TRACK_A) is item
254 item.claim()
255 item.close()
256 item.release()
257 # this is what a queue holding the same track twice, or repeat wrapping back
258 # to the top, asks for: it must not be handed a drained channel
259 assert session.item_for(TRACK_A) is None
260
261
262async def test_an_abandoned_channel_cannot_be_continued(tmp_path: Path) -> None:
263 """A stream abandoned mid-item cannot resume where it left off either."""
264 session = _make_session(tmp_path)
265 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
266 item.started.set()
267 item.claim()
268 item.release()
269 assert session.item_for(TRACK_A) is None
270
271
272async def test_a_stuck_item_fails_instead_of_streaming_forever(tmp_path: Path) -> None:
273 """An item that runs far past its duration without a track change fails."""
274 session = _make_session(tmp_path)
275 item = _ItemAudio(TRACK_A, session)
276 item.duration_ms = 1_000
277 item.claim()
278 limit = item._overrun_limit()
279 assert limit is not None
280 item.write(b"\x01" * (limit + _FRAME_BYTES))
281 with pytest.raises(AudioError, match="never moved on"):
282 async for _ in item.read():
283 pass
284
285
286async def test_the_first_logged_out_snapshot_is_not_a_lost_pairing(tmp_path: Path) -> None:
287 """A daemon reports logged_in=False until it has restored its session."""
288 session = _make_session(tmp_path)
289 session._logged_in = None
290 session._was_active = False
291 await session._handle_event(_auth_event(logged_in=False, is_active=False))
292 # failing here would break every playback on a perfectly good pairing
293 assert session.usable is True
294 await session._handle_event(_auth_event(logged_in=True, is_active=False))
295 assert session.usable is True
296
297
298async def test_losing_an_established_login_fails_the_session(tmp_path: Path) -> None:
299 """A login that goes away mid-session is real, and ends the session."""
300 session = _make_session(tmp_path)
301 await session._handle_event(_auth_event(logged_in=True))
302 await session._handle_event(_auth_event(logged_in=False))
303 assert session.usable is False
304 assert session._error == "the session was logged out"
305
306
307async def test_buffering_gates_the_sink_once_demand_started(tmp_path: Path) -> None:
308 """Once PCM demand started, playing runs the sink and buffering suspends it again."""
309 session = _make_session(tmp_path)
310 session._demand_started = True
311 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
312 session._pending.append(TRACK_B)
313 sink = _sink_of(session)
314 # the sink is created suspended, so there is nothing to suspend yet
315 await session._handle_event(_playback_event("buffering"))
316 sink.suspend.assert_not_awaited()
317 await session._handle_event(_playback_event("playing"))
318 sink.resume.assert_awaited_once()
319 assert session._current is not None
320 assert session._current.playing_seen is True
321 # the engine stalling on a rebuffer keeps that silence out of the PCM
322 await session._handle_event(_playback_event("buffering"))
323 sink.suspend.assert_awaited_once()
324
325
326async def test_sink_is_not_gated_before_demand_started(tmp_path: Path) -> None:
327 """Buffering/playing before PCM demand leave the (still suspended) sink alone."""
328 session = _make_session(tmp_path)
329 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
330 sink = _sink_of(session)
331 await session._handle_event(_playback_event("buffering"))
332 await session._handle_event(_playback_event("playing"))
333 sink.suspend.assert_not_awaited()
334 sink.resume.assert_not_awaited()
335 # the status is recorded either way, so session start can decide when to resume
336 assert session._current.status == "playing"
337
338
339@pytest.mark.parametrize("end_status", ["stopped", "idle", "paused"])
340async def test_the_last_item_is_drained_rather_than_cut(tmp_path: Path, end_status: str) -> None:
341 """However the engine reports the end of a run, the last item drains and closes."""
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 = 1_000
347 item.last_position_ms = 1_000
348 one_second = 1_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
349 sink = _sink_of(session)
350 await session._handle_event(_playback_event(end_status, position_ms=1_000))
351 # the sink stays open for now, so audio still in the FIFO can arrive...
352 sink.suspend.assert_not_awaited()
353 assert item.draining is True
354 assert item._closed is False
355 # ... but only that item's own audio is taken, never the padding silence the
356 # sink keeps rendering afterwards
357 item.write(b"\x01" * one_second)
358 item.write(b"\x00" * 4096)
359 assert item.buffered == one_second
360 await _wait_for(lambda: item._closed)
361 sink.suspend.assert_awaited_once()
362
363
364async def test_an_app_pause_midway_through_the_last_item_is_not_the_end(
365 tmp_path: Path,
366) -> None:
367 """Pausing in the Spotify app halfway through the last track must not truncate it."""
368 session = _make_session(tmp_path)
369 session._demand_started = True
370 session._sink_running = True
371 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
372 item.duration_ms = 200_000
373 item.last_position_ms = 90_000
374 await session._handle_event(_playback_event("paused", position_ms=90_000))
375 assert item.draining is False
376 assert item._closed is False
377 # treated as interference instead: the sink is gated and playback resumed
378 _sink_of(session).suspend.assert_awaited_once()
379 _client_of(session).resume.assert_awaited_once()
380
381
382async def test_a_resumed_item_cancels_its_tail_drain(tmp_path: Path) -> None:
383 """An armed drain is undone when the engine turns out to have been rebuffering."""
384 session = _make_session(tmp_path)
385 session._demand_started = True
386 session._sink_running = True
387 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
388 item.duration_ms = 200_000
389 item.last_position_ms = 199_000
390 await session._handle_event(_playback_event("stopped", position_ms=199_000))
391 armed = item.draining
392 await session._handle_event(_playback_event("playing", position_ms=199_500))
393 assert armed is True
394 assert item.draining is False
395 assert item._closed is False
396 assert item.drain_task is None
397
398
399async def test_the_cushion_is_capped_by_suspending_the_sink(tmp_path: Path) -> None:
400 """Undelivered audio is handed back as backpressure rather than piling up."""
401 session = _make_session(tmp_path)
402 session._demand_started = True
403 session._sink_running = True
404 session._engine_playing = True
405 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
406 item.claim()
407 sink = _sink_of(session)
408 await session._apply_sink_state()
409 sink.suspend.assert_not_awaited()
410 # the engine has run this far ahead of what the player has taken
411 item.write(b"\x01" * int((soloist_backend._MAX_RETAINED_S + 1) * _BYTES_PER_SECOND))
412 await session._apply_sink_state()
413 sink.suspend.assert_awaited_once()
414 assert session._backpressured is True
415 # and it comes back once the player has drained enough of it
416 item._buffered = int(soloist_backend._RESUME_RETAINED_S * _BYTES_PER_SECOND) - 1
417 await session._apply_sink_state()
418 sink.resume.assert_awaited_once()
419 assert session._backpressured is False
420
421
422async def test_a_pause_with_more_queued_suspends_the_sink(tmp_path: Path) -> None:
423 """A pause while another item is queued behind is ordinary interference, not the end."""
424 session = _make_session(tmp_path)
425 session._demand_started = True
426 session._sink_running = True
427 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
428 session._pending.append(TRACK_B)
429 await session._handle_event(_playback_event("paused"))
430 _sink_of(session).suspend.assert_awaited_once()
431
432
433async def test_nothing_is_sent_before_the_websocket_is_up(
434 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
435) -> None:
436 """Commands travel over the events socket: a published endpoint is not enough."""
437 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
438 session = _make_session(tmp_path)
439 client = _client_of(session)
440 # the endpoint file exists, but the events task has not connected yet
441 client.connected = False
442 endpoint_published = asyncio.Event()
443 endpoint_published.set()
444 with pytest.raises(AudioError, match="did not connect and log in"):
445 await session._play(TRACK_A, 0, endpoint_published)
446 client.activate.assert_not_awaited()
447 client.play.assert_not_awaited()
448
449
450async def test_nothing_is_sent_before_the_engine_has_logged_in(
451 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
452) -> None:
453 """The engine drops commands sent before it has restored its session."""
454 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
455 session = _make_session(tmp_path)
456 client = _client_of(session)
457 client.connected = True
458 # connected, but the engine has not announced its login yet
459 session._logged_in = None
460 endpoint_published = asyncio.Event()
461 endpoint_published.set()
462 with pytest.raises(AudioError, match="did not connect and log in"):
463 await session._play(TRACK_A, 0, endpoint_published)
464 client.activate.assert_not_awaited()
465 client.play.assert_not_awaited()
466
467
468async def test_startup_activates_before_it_plays(
469 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
470) -> None:
471 """A fresh daemon has to become the active device before it is told to play."""
472 session = _make_session(tmp_path)
473 client = _client_of(session)
474 client.connected = True
475 monkeypatch.setattr(session, "_await_item_ready", AsyncMock())
476 endpoint_published = asyncio.Event()
477 endpoint_published.set()
478 item = await session._play(TRACK_A, 0, endpoint_published)
479 assert item.uri == TRACK_A
480 client.activate.assert_awaited_once_with(await_result=True)
481 client.play.assert_awaited_once_with(TRACK_A)
482
483
484async def test_a_takeover_between_activate_and_play_stops_the_start(tmp_path: Path) -> None:
485 """Playing here would claim the device straight back off wherever the user moved to."""
486 session = _make_session(tmp_path)
487 session._was_active = False
488 client = _client_of(session)
489
490 async def _take_over(*_args: Any, **_kwargs: Any) -> None:
491 session._observe_active_device(is_active=False)
492
493 client.set_repeat_track.side_effect = _take_over
494 ready = asyncio.Event()
495 ready.set()
496 with pytest.raises(SoloistAppControlError):
497 await session._play(TRACK_A, 0, ready)
498 client.play.assert_not_awaited()
499
500
501async def test_a_refused_start_command_reports_soloist(tmp_path: Path) -> None:
502 """A dropped start command surfaces as a Soloist error, not a raw client one."""
503 session = _make_session(tmp_path)
504 client = _client_of(session)
505 client.connected = True
506 client.activate.side_effect = SoloistError("websocket is not connected")
507 endpoint_published = asyncio.Event()
508 endpoint_published.set()
509 with pytest.raises(AudioError, match="Spotify Soloist would not start"):
510 await session._play(TRACK_A, 0, endpoint_published)
511
512
513async def test_the_engine_is_told_not_to_shuffle_or_repeat(
514 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
515) -> None:
516 """MA owns the order, and a repeating engine would never reach the item fed behind."""
517 session = _make_session(tmp_path)
518 client = _client_of(session)
519 client.connected = True
520 monkeypatch.setattr(session, "_await_item_ready", AsyncMock())
521 endpoint_published = asyncio.Event()
522 endpoint_published.set()
523 await session._play(TRACK_A, 0, endpoint_published)
524 client.set_shuffle.assert_awaited_once_with(False)
525 client.set_repeat_context.assert_awaited_once_with(False)
526 client.set_repeat_track.assert_awaited_once_with(False)
527
528
529async def test_repeat_turned_on_from_the_app_is_pinned_back_off(tmp_path: Path) -> None:
530 """Repeat enabled in the Spotify app is undone before it can loop the item."""
531 session = _make_session(tmp_path)
532 await session._handle_event(
533 SoloistEvent(
534 type="options_changed",
535 data=SoloistOptionsChanged(
536 options=SoloistPlaybackOptions(shuffle=True, repeat="track")
537 ),
538 raw={},
539 )
540 )
541 client = _client_of(session)
542 client.set_shuffle.assert_awaited_once_with(False)
543 client.set_repeat_track.assert_awaited_once_with(False)
544 client.set_repeat_context.assert_awaited_once_with(False)
545 # options that are already off are left alone
546 client.set_shuffle.reset_mock()
547 await session._handle_event(
548 SoloistEvent(
549 type="options_changed",
550 data=SoloistOptionsChanged(options=SoloistPlaybackOptions()),
551 raw={},
552 )
553 )
554 client.set_shuffle.assert_not_awaited()
555
556
557async def test_a_busy_data_directory_is_reported_as_such(tmp_path: Path) -> None:
558 """A daemon left over from an earlier run is named, not reported as a generic failure."""
559 session = _make_session(tmp_path)
560 # the daemon's own parting complaint, which is all it gives (it exits with 1)
561 session._data_dir_busy = True
562 with pytest.raises(AudioError, match="Another Spotify Soloist session is still running"):
563 session._raise_startup_error("exited before playback started", TRACK_A)
564
565
566async def test_the_busy_marker_is_picked_up_from_the_daemon_output(tmp_path: Path) -> None:
567 """The marker is read off the daemon's stdout, with the API key still redacted."""
568 session = _make_session(tmp_path)
569 await session._log_output(
570 _stdout_of(
571 'Error: another session is running for data directory "/data/x/soloist-data".',
572 "Stop the running session before starting soloist again.",
573 )
574 )
575 assert session._data_dir_busy is True
576
577
578async def test_a_lost_pairing_is_caught_the_moment_the_daemon_reports_it(
579 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
580) -> None:
581 """A daemon advertising for pairing fails the session at once, not on a timeout."""
582 session = _make_session(tmp_path)
583 session._logged_in = None
584 unload_with_error = MagicMock()
585 monkeypatch.setattr(session.backend.provider, "unload_with_error", unload_with_error)
586 await session._log_output(
587 _stdout_of('waiting for login - connect to "X" from your Spotify app')
588 )
589 assert session._unpaired is True
590 # the buffer gives up on the audio long before the startup budget runs out, so
591 # the session has to fail while an item is still waiting on it
592 assert session._error is not None
593 with pytest.raises(LoginFailed) as err:
594 session._raise_startup_error("did not connect and log in", TRACK_A)
595 assert err.value.translation_key == "soloist_pairing_required"
596 unload_with_error.assert_called_once()
597
598
599async def test_a_lost_pairing_fails_the_item_without_waiting_for_the_endpoint(
600 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
601) -> None:
602 """The item fails on the lost pairing, not on the endpoint that is no longer coming."""
603 session = _make_session(tmp_path)
604 session._logged_in = None
605 monkeypatch.setattr(session.backend.provider, "unload_with_error", MagicMock())
606 await session._log_output(
607 _stdout_of('waiting for login - connect to "X" from your Spotify app')
608 )
609 # the endpoint never appears, so the wait for it must not swallow the failure:
610 # sitting it out would outlast the queue's own patience for the audio
611 with pytest.raises(LoginFailed) as err:
612 async with asyncio.timeout(5):
613 await session._play(TRACK_A, 0, asyncio.Event())
614 assert err.value.translation_key == "soloist_pairing_required"
615
616
617async def test_a_daemon_still_restoring_its_session_is_left_alone(tmp_path: Path) -> None:
618 """The engine advertises for pairing while restoring too; the stored session decides."""
619 session = _make_session(tmp_path)
620 data_dir = session.backend._data_dir
621 (data_dir / "settings" / "Users" / "spotify-user-user").mkdir(parents=True)
622 await session._log_output(
623 _stdout_of('waiting for login - connect to "X" from your Spotify app')
624 )
625 assert session._unpaired is False
626 assert session._error is None
627
628
629async def test_a_pairing_that_never_logs_in_routes_through_setup(
630 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
631) -> None:
632 """A session that cannot log in sends the user back to setup, not a per-track error."""
633 session = _make_session(tmp_path)
634 unload_with_error = MagicMock()
635 monkeypatch.setattr(session.backend.provider, "unload_with_error", unload_with_error)
636 await session._handle_event(_auth_event(logged_in=False))
637 with pytest.raises(LoginFailed) as err:
638 session._raise_startup_error("timed out waiting for playback to start", TRACK_A)
639 assert err.value.translation_key == "soloist_pairing_required"
640 # ... and the provider is taken out of service, so the user is asked to redo setup
641 unload_with_error.assert_called_once()
642
643
644async def test_a_login_that_never_happened_is_not_confused_with_another_failure(
645 tmp_path: Path,
646) -> None:
647 """An unrelated failure keeps its own message even before any login was reported."""
648 session = _make_session(tmp_path)
649 session._fail("the capture sink was lost mid-stream")
650 with pytest.raises(AudioError, match="capture sink was lost"):
651 session._raise_startup_error("exited before playback started", TRACK_A)
652
653
654def test_a_seeked_item_only_expects_what_is_left_of_it(tmp_path: Path) -> None:
655 """A seeked item delivers the remainder, so its targets are based on that."""
656 session = _make_session(tmp_path)
657 item = _ItemAudio(TRACK_A, session)
658 item.duration_ms = 200_000
659 full = 200_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
660 assert item._duration_bytes() == full
661 item.seek_target_ms = 150_000
662 remainder = 50_000 * CAPTURE_SAMPLE_RATE // 1000 * _FRAME_BYTES
663 assert item._duration_bytes() == remainder
664 # so the tail drain has a target it can actually reach
665 item.start_tail_drain()
666 item.write(b"\x01" * remainder)
667 assert item.tail_complete is True
668 # and the padding silence after it is refused
669 item.write(b"\x00" * 4096)
670 assert item.buffered == remainder
671
672
673@pytest.mark.parametrize(
674 ("target_ms", "reports"),
675 [
676 # seeking the restored item to where the engine already was
677 (117_000, (117_000, 0, 117_000)),
678 # seeking back into it, where every report lands below where it was
679 (30_000, (117_000, 0, 30_000)),
680 ],
681)
682async def test_the_seek_retries_until_the_engine_reports_the_target(
683 tmp_path: Path, target_ms: int, reports: tuple[int, ...]
684) -> None:
685 """A seek dropped while the track loads is re-sent until a report confirms it."""
686 session = _make_session(tmp_path)
687 item = _ItemAudio(TRACK_A, session)
688 client = cast("Any", session._client)
689 # the engine restored this item part-way in, before the seek goes out
690 item.observe_position(117_000)
691
692 async def _report_positions() -> None:
693 for position_ms in reports:
694 await asyncio.sleep(0)
695 item.observe_position(position_ms)
696
697 with (
698 patch.object(soloist_backend, "_SEEK_RETRY_INTERVAL_S", 0.01),
699 # bounded so a regression fails fast instead of sitting out the real budget
700 patch.object(soloist_backend, "_SEEK_CONFIRM_TIMEOUT_S", 1.0),
701 ):
702 reporter = asyncio.create_task(_report_positions())
703 await session._cold_seek(client, item, target_ms)
704 await reporter
705 assert item.seek_confirmed.is_set()
706 assert item.started_at_ms == target_ms
707 assert client.seek.await_count >= 1
708
709
710async def test_a_seek_that_only_ever_sees_the_restored_position_fails(tmp_path: Path) -> None:
711 """A seek nothing confirms fails loudly rather than streaming from elsewhere."""
712 session = _make_session(tmp_path)
713 item = _ItemAudio(TRACK_A, session)
714 # the engine sits at the restored position and never reloads the track
715 item.observe_position(117_000)
716 with (
717 patch.object(soloist_backend, "_SEEK_RETRY_INTERVAL_S", 0.01),
718 patch.object(soloist_backend, "_SEEK_CONFIRM_TIMEOUT_S", 0.05),
719 pytest.raises(AudioError, match="did not confirm seeking"),
720 ):
721 await session._cold_seek(cast("Any", session._client), item, 117_000)
722
723
724async def test_a_seek_the_engine_ignored_does_not_cut_the_item_short(tmp_path: Path) -> None:
725 """An item the engine plays from its start is bounded by its full duration."""
726 session = _make_session(tmp_path)
727 item = _ItemAudio(TRACK_A, session)
728 item.duration_ms = 176_000
729 item.arm_seek(117_000)
730 item.claim()
731 # the engine never made the seek and is playing the item from its start, so
732 # the audio it delivers runs well past what the seeked remainder would allow
733 item.observe_position(80_000)
734 item.write(b"\x01" * (89 * _BYTES_PER_SECOND))
735 item.close()
736 delivered = 0
737 async for chunk in item.read():
738 delivered += len(chunk)
739 assert delivered == 89 * _BYTES_PER_SECOND
740
741
742async def test_a_seek_that_landed_still_bounds_the_item_at_its_remainder(tmp_path: Path) -> None:
743 """An item the engine really seeked into stays bounded by what is left of it."""
744 session = _make_session(tmp_path)
745 item = _ItemAudio(TRACK_A, session)
746 item.duration_ms = 176_000
747 item.arm_seek(117_000)
748 item.observe_position(0)
749 item.observe_position(117_000)
750 assert item.seek_confirmed.is_set()
751 assert item._overrun_limit() == 59 * _BYTES_PER_SECOND + int(
752 _ITEM_OVERRUN_S * _BYTES_PER_SECOND
753 )
754 # later reports do not move the latch, which would shrink the bound
755 item.observe_position(150_000)
756 assert item.started_at_ms == 117_000
757 # so it still fails once it runs that far past the seek point
758 item.claim()
759 item.write(b"\x01" * (89 * _BYTES_PER_SECOND))
760 with pytest.raises(AudioError, match="never moved on"):
761 async for _ in item.read():
762 pass
763
764
765def test_the_lead_trim_never_exceeds_its_budget() -> None:
766 """Silence beyond the budget is content, including where audio starts mid-chunk."""
767 budget = int(_MAX_LEAD_TRIM_S * _BYTES_PER_SECOND)
768 # already at the budget, with a chunk whose silence runs well past it
769 chunk = b"\x00" * 4096 + b"\x01" * 64
770 trimmed, skipped = _trim_lead_silence(chunk, budget - _FRAME_BYTES)
771 assert skipped == _FRAME_BYTES
772 assert len(trimmed) == len(chunk) - _FRAME_BYTES
773
774
775async def test_a_dying_log_reader_fails_the_session(tmp_path: Path) -> None:
776 """Nothing else drains the daemon's stdout, so a dead reader must not go unnoticed."""
777 session = _make_session(tmp_path)
778
779 async def _boom() -> None:
780 raise RuntimeError("reader blew up")
781
782 session._log_task = asyncio.create_task(_boom())
783 session._log_task.add_done_callback(session._task_done)
784 await asyncio.sleep(0)
785 await _wait_for(lambda: not session.usable)
786 assert session._error is not None
787 assert "reader blew up" in session._error
788
789
790async def test_feeding_never_replaces_a_channel_already_in_use(tmp_path: Path) -> None:
791 """If the engine reaches the fed item first, its live channel must survive."""
792 session = _make_session(tmp_path, queue_id="player1")
793 streamdetails = MagicMock()
794 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
795 queues = _queues_of(session)
796 queues.get.return_value = MagicMock(current_index=0)
797 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
798 queues.get_next_item.return_value = _queue_item(TRACK_B)
799
800 async def _engine_gets_there_first(_uri: str, **_kwargs: Any) -> None:
801 # the events task advances to the fed item while the command is in flight
802 await session._observe_current(TRACK_B, 200_000)
803
804 _client_of(session).add_to_queue.side_effect = _engine_gets_there_first
805 await session.feed_after(streamdetails, TRACK_A)
806 live = session.current
807 assert live is not None
808 assert live.uri == TRACK_B
809 # the channel the reader is writing to is the one a stream will be handed
810 assert session._items[TRACK_B] is live
811 assert session.item_for(TRACK_B) is live
812 # and it is not queued as pending, because it already started
813 assert session.has_pending is False
814
815
816async def test_a_seek_the_session_cannot_take_restarts_it(
817 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
818) -> None:
819 """
820 A seek the running session cannot serve falls back to restarting it.
821
822 A realtime source has not captured anything past the play position, so any
823 forward seek lands outside the buffer and comes back here; the session is
824 seeked in place when it can be, and replaced when it cannot.
825 """
826 backend = _make_backend(tmp_path)
827 backend._server = MagicMock()
828 backend._binary = Path("/nonexistent/soloist")
829 session = _SoloistSession(backend, "player1")
830 backend._session = session
831 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
832 item.started.set()
833 # its own stream is still attached when the seek re-opens it
834 item.claim()
835 stopped = AsyncMock()
836 monkeypatch.setattr(session, "stop", stopped)
837 _install_fake_binary_manager(monkeypatch)
838 monkeypatch.setattr(
839 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
840 )
841 with pytest.raises(AudioError, match="spawn"):
842 await backend._acquire(TRACK_A, 90, "player1")
843 stopped.assert_awaited_once()
844
845
846@pytest.mark.parametrize(
847 ("requested", "other_queue"),
848 [
849 # another player, whatever it asks for - including the very track this
850 # session is in the middle of delivering
851 pytest.param(TRACK_B, "player2", id="other_player"),
852 pytest.param(TRACK_A, "player2", id="other_player_same_track"),
853 # an early fetch across a boundary this session does not drive, such as a
854 # podcast episode or audiobook chapter
855 pytest.param(TRACK_B, "player1", id="unstitched_boundary"),
856 ],
857)
858async def test_a_session_in_use_is_never_cut_short(
859 tmp_path: Path, requested: str, other_queue: str
860) -> None:
861 """
862 An item the session cannot serve must not stop one it is still delivering.
863
864 Reported as capacity, so a speculative prepare gives up softly.
865 """
866 backend = _make_backend(tmp_path)
867 backend._server = MagicMock()
868 backend._binary = Path("/nonexistent/soloist")
869 session = _SoloistSession(backend, "player1")
870 backend._session = session
871 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
872 item.started.set()
873 item.claim()
874 # the session really is playing TRACK_A, so a same-track request from another
875 # player cannot be mistaken for a seek
876 session._current = item
877 with pytest.raises(ProviderStreamLimitError) as err:
878 await backend._acquire(requested, 0, other_queue)
879 # a stream-limit error so the item is not marked unplayable, but the message
880 # is about the session, not the provider's source-stream budget
881 assert err.value.limit == 1
882 assert err.value.translation_key == "soloist_session_busy"
883 # the session that was playing is untouched
884 assert backend._session is session
885 assert session.usable is True
886
887
888async def test_a_session_nobody_reads_is_replaced_for_another_item(
889 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
890) -> None:
891 """Once the other item has been released, the same request gets the session."""
892 backend = _make_backend(tmp_path)
893 backend._server = MagicMock()
894 backend._binary = Path("/nonexistent/soloist")
895 session = _SoloistSession(backend, "player1")
896 backend._session = session
897 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
898 item.started.set()
899 item.claim()
900 item.close()
901 item.release()
902 _install_fake_binary_manager(monkeypatch)
903 monkeypatch.setattr(
904 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
905 )
906 monkeypatch.setattr(session, "stop", AsyncMock())
907 with pytest.raises(AudioError, match="spawn"):
908 await backend._acquire(TRACK_B, 0, "player1")
909
910
911async def test_a_replacement_waits_for_the_old_daemon_to_be_gone(
912 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
913) -> None:
914 """The engine refuses to start while another daemon still holds its data dir."""
915 backend = _make_backend(tmp_path)
916 backend._server = MagicMock()
917 backend._binary = Path("/nonexistent/soloist")
918 session = _SoloistSession(backend, "player1")
919 backend._session = session
920 order: list[str] = []
921
922 async def _slow_stop() -> None:
923 order.append("stop-start")
924 await asyncio.sleep(0.05)
925 order.append("stop-done")
926
927 monkeypatch.setattr(session, "stop", _slow_stop)
928 _install_fake_binary_manager(monkeypatch)
929
930 async def _spawn(_self: Any, _uri: str, _seek: int) -> None:
931 order.append("spawn")
932 raise AudioError("spawn")
933
934 monkeypatch.setattr(soloist_backend._SoloistSession, "start", _spawn)
935 # the session failed, so its teardown is under way when the next item arrives
936 discard = asyncio.create_task(backend.discard_session(session))
937 await asyncio.sleep(0)
938 with pytest.raises(AudioError, match="spawn"):
939 await backend._acquire(TRACK_B, 0, "player1")
940 await discard
941 assert order == ["stop-start", "stop-done", "spawn"]
942
943
944async def test_an_idle_session_is_taken_over(
945 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
946) -> None:
947 """A session nobody is reading is replaced instead of blocking another player."""
948 backend = _make_backend(tmp_path)
949 backend._server = MagicMock()
950 backend._binary = Path("/nonexistent/soloist")
951 session = _SoloistSession(backend, "player1")
952 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
953 backend._session = session
954 stopped = AsyncMock()
955 monkeypatch.setattr(session, "stop", stopped)
956 _install_fake_binary_manager(monkeypatch)
957 # the replacement spawn is out of scope here; only the takeover decision is
958 monkeypatch.setattr(
959 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
960 )
961 with pytest.raises(AudioError, match="spawn"):
962 await backend._acquire(TRACK_B, 0, "player2")
963 stopped.assert_awaited_once()
964
965
966def test_a_dead_session_task_fails_the_session(tmp_path: Path) -> None:
967 """A session task that dies of an unexpected error takes the session with it."""
968 session = _make_session(tmp_path)
969 task: Any = MagicMock()
970 task.cancelled.return_value = False
971 task.exception.return_value = RuntimeError("reader blew up")
972 session._task_done(task)
973 assert session.usable is False
974 assert session._error is not None
975 assert "reader blew up" in session._error
976
977
978def test_a_cancelled_session_task_is_not_a_failure(tmp_path: Path) -> None:
979 """Teardown cancels the session's tasks; that must not be reported as an error."""
980 session = _make_session(tmp_path)
981 task: Any = MagicMock()
982 task.cancelled.return_value = True
983 session._task_done(task)
984 assert session.usable is True
985
986
987async def test_failed_sink_control_fails_the_session(tmp_path: Path) -> None:
988 """A failed suspend/resume fails the session instead of leaking stall silence."""
989 session = _make_session(tmp_path)
990 session._demand_started = True
991 session._sink_running = True
992 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
993 session._pending.append(TRACK_B)
994 _sink_of(session).suspend.side_effect = RuntimeError("pactl failed")
995 await session._handle_event(_playback_event("buffering"))
996 assert session._error is not None
997 assert "capture sink control failed" in session._error
998
999
1000async def test_app_pause_is_fought_with_a_resume(tmp_path: Path) -> None:
1001 """A pause from the Spotify app is undone: this session has no user-facing pause."""
1002 session = _make_session(tmp_path)
1003 session._demand_started = True
1004 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1005 session._pending.append(TRACK_B)
1006 await session._handle_event(_playback_event("paused"))
1007 _client_of(session).resume.assert_awaited_once()
1008
1009
1010async def test_an_app_pause_is_only_undone_so_many_times(tmp_path: Path) -> None:
1011 """Someone who keeps pausing means it: the session gives up instead of fighting on."""
1012 session = _make_session(tmp_path)
1013 session._demand_started = True
1014 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1015 session._pending.append(TRACK_B)
1016 for _ in range(_MAX_APP_PAUSE_RESUMES):
1017 await session._handle_event(_playback_event("playing"))
1018 await session._handle_event(_playback_event("paused"))
1019 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
1020 assert session._error is None
1021
1022 await session._handle_event(_playback_event("playing"))
1023 await session._handle_event(_playback_event("paused"))
1024 assert _client_of(session).resume.await_count == _MAX_APP_PAUSE_RESUMES
1025 assert session.usable is False
1026 assert session._app_control is SoloistAppControl.PAUSED
1027
1028
1029async def test_one_pause_reported_twice_counts_once(tmp_path: Path) -> None:
1030 """A repeated snapshot of the same pause is not a new pause."""
1031 session = _make_session(tmp_path)
1032 session._demand_started = True
1033 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1034 session._pending.append(TRACK_B)
1035 for _ in range(_MAX_APP_PAUSE_RESUMES + 2):
1036 await session._handle_event(_playback_event("paused"))
1037 assert session.usable is True
1038
1039
1040async def test_the_pause_budget_resets_on_the_next_item(tmp_path: Path) -> None:
1041 """Each item gets its own budget; pausing one track does not spend the next one's."""
1042 session = _make_session(tmp_path)
1043 session._demand_started = True
1044 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1045 session._pending.append(TRACK_B)
1046 for _ in range(_MAX_APP_PAUSE_RESUMES):
1047 await session._handle_event(_playback_event("playing"))
1048 await session._handle_event(_playback_event("paused"))
1049 await session._observe_current(TRACK_B, 200_000)
1050 assert session._app_pauses == 0
1051
1052
1053async def test_a_pause_is_not_undone_once_the_device_is_gone(tmp_path: Path) -> None:
1054 """A bare resume on a device Spotify no longer routes to would play to nobody."""
1055 session = _make_session(tmp_path)
1056 session._demand_started = True
1057 session._was_active = False
1058 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1059 session._pending.append(TRACK_B)
1060 await session._handle_event(_playback_event("paused"))
1061 _client_of(session).resume.assert_not_awaited()
1062
1063
1064async def test_losing_the_active_device_ends_the_session(tmp_path: Path) -> None:
1065 """Playback moved to another device from the Spotify app: this session is over."""
1066 session = _make_session(tmp_path)
1067 await session._handle_event(_device_event(is_active=False))
1068 assert session.usable is False
1069 assert session._app_control is SoloistAppControl.TOOK_OVER
1070
1071
1072async def test_a_takeover_reported_on_the_auth_state_ends_the_session(tmp_path: Path) -> None:
1073 """The active-device state also rides on auth_state, and counts the same there."""
1074 session = _make_session(tmp_path)
1075 await session._handle_event(_auth_event(logged_in=True, is_active=False))
1076 assert session.usable is False
1077
1078
1079async def test_an_inactive_device_before_activation_is_not_a_takeover(tmp_path: Path) -> None:
1080 """A fresh daemon is inactive until the session claims it; that is not a takeover."""
1081 session = _make_session(tmp_path)
1082 session._was_active = False
1083 await session._handle_event(_device_event(is_active=False))
1084 await session._handle_event(_auth_event(logged_in=True, is_active=False))
1085 assert session.usable is True
1086
1087 # nor does a respawned daemon reporting the session Spotify still has for
1088 # the account: only the status _play claimed is followed
1089 await session._handle_event(_device_event(is_active=True))
1090 await session._handle_event(_device_event(is_active=False))
1091 assert session.usable is True
1092
1093
1094async def test_a_reconnect_snapshot_keeps_an_active_session_alive(tmp_path: Path) -> None:
1095 """The events connection re-snapshots after a drop; that is not a device change."""
1096 session = _make_session(tmp_path)
1097 await session._handle_event(_auth_event(logged_in=True, is_active=True))
1098 await session._handle_event(_device_event(is_active=True))
1099 assert session.usable is True
1100
1101
1102async def test_the_playback_snapshots_active_flag_is_ignored(tmp_path: Path) -> None:
1103 """It is optional and rides on deltas, so only the dedicated reports are followed."""
1104 session = _make_session(tmp_path)
1105 session._demand_started = True
1106 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1107 await session._handle_event(
1108 SoloistEvent(
1109 type="playback_changed",
1110 data=SoloistPlaybackState(status="playing", is_active=False),
1111 raw={},
1112 )
1113 )
1114 assert session.usable is True
1115
1116
1117async def test_backpressure_does_not_spend_the_pause_budget(tmp_path: Path) -> None:
1118 """A sink suspended to cap the cushion is our doing, not the user pausing."""
1119 session = _make_session(tmp_path)
1120 session._demand_started = True
1121 session._engine_playing = True
1122 session._backpressured = True
1123 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1124 session._pending.append(TRACK_B)
1125 await session._handle_event(_playback_event("paused"))
1126 _client_of(session).resume.assert_not_awaited()
1127 assert session._app_pauses == 0
1128
1129
1130async def test_a_lost_login_is_not_reported_as_a_takeover(tmp_path: Path) -> None:
1131 """Losing the login wins over the inactive device it brings with it."""
1132 session = _make_session(tmp_path)
1133 await session._handle_event(_auth_event(logged_in=False, is_active=False))
1134 assert session.usable is False
1135 assert session._app_control is None
1136
1137
1138async def test_a_track_started_from_the_app_ends_the_session(tmp_path: Path) -> None:
1139 """The engine pulled off an item part-way through is the app playing something else."""
1140 session = _make_session(tmp_path)
1141 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1142 item.duration_ms = 200_000
1143 await session._observe_current(TRACK_A, 200_000)
1144 item.observe_position(20_000)
1145
1146 await session._observe_current("spotify:track:theirs", 180_000)
1147 assert session.usable is False
1148 assert session._app_control is SoloistAppControl.TOOK_OVER
1149 assert session.current is item
1150
1151
1152async def test_a_track_played_earlier_started_from_the_app_ends_the_session(
1153 tmp_path: Path,
1154) -> None:
1155 """A known uri is no exemption: only the item fed behind this one is where we sent it."""
1156 session = _make_session(tmp_path)
1157 played = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1158 played.spent = True
1159 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1160 item.duration_ms = 200_000
1161 await session._observe_current(TRACK_A, 200_000)
1162 item.observe_position(20_000)
1163
1164 await session._observe_current(TRACK_B, 180_000)
1165 assert session.usable is False
1166 assert session._app_control is SoloistAppControl.TOOK_OVER
1167
1168
1169async def test_skipping_from_the_app_to_the_fed_item_is_followed(tmp_path: Path) -> None:
1170 """The queue moves to that same track, so following the engine keeps the two in step."""
1171 session = _make_session(tmp_path)
1172 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1173 item.duration_ms = 200_000
1174 await session._observe_current(TRACK_A, 200_000)
1175 item.observe_position(20_000)
1176 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1177 session._pending.append(TRACK_B)
1178
1179 await session._observe_current(TRACK_B, 180_000)
1180 assert session.usable is True
1181 assert session.current is fed
1182 assert session.item_for(TRACK_B) is fed
1183
1184
1185async def test_a_takeover_snapshot_stops_pinning_volume_and_options(tmp_path: Path) -> None:
1186 """Once the app has the session, the rest of its snapshot must not reach the daemon."""
1187 session = _make_session(tmp_path)
1188 session._demand_started = True
1189 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1190 item.duration_ms = 200_000
1191 await session._observe_current(TRACK_A, 200_000)
1192 item.observe_position(20_000)
1193
1194 await session._handle_event(
1195 SoloistEvent(
1196 type="playback_changed",
1197 data=SoloistPlaybackState(
1198 status="playing",
1199 item=SoloistEntity(uri="spotify:track:theirs", entity_type="track"),
1200 volume=40,
1201 options=SoloistPlaybackOptions(shuffle=True, repeat="context"),
1202 ),
1203 raw={},
1204 )
1205 )
1206 assert session.usable is False
1207 _client_of(session).set_volume.assert_not_awaited()
1208 _client_of(session).set_shuffle.assert_not_awaited()
1209
1210
1211async def test_the_engine_moving_on_at_a_track_end_is_not_a_takeover(tmp_path: Path) -> None:
1212 """An unasked-for item the engine reaches at a boundary is its own autoplay."""
1213 session = _make_session(tmp_path)
1214 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1215 item.duration_ms = 200_000
1216 await session._observe_current(TRACK_A, 200_000)
1217 item.observe_position(200_000)
1218
1219 await session._observe_current("spotify:track:autoplay", 180_000)
1220 assert session.usable is True
1221 assert session.item_for("spotify:track:autoplay") is None
1222
1223
1224async def test_an_ended_item_says_what_the_app_did(tmp_path: Path) -> None:
1225 """The item's stream fails with the takeover, not a generic session error."""
1226 session = _make_session(tmp_path)
1227 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1228 await session._handle_event(_device_event(is_active=False))
1229 with pytest.raises(SoloistAppControlError) as err:
1230 await session.validate_item(item)
1231 assert err.value.translation_key == SoloistAppControl.TOOK_OVER.value
1232 assert isinstance(err.value, ProviderStreamLimitError)
1233
1234
1235async def test_a_session_being_torn_down_does_not_hold_off_the_next_one(tmp_path: Path) -> None:
1236 """Teardown pauses the daemon; that must not read as the user pausing."""
1237 session = _make_session(tmp_path)
1238 session._demand_started = True
1239 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1240 session._pending.append(TRACK_B)
1241 session._stopped = True
1242 for _ in range(_MAX_APP_PAUSE_RESUMES + 1):
1243 await session._handle_event(_playback_event("playing"))
1244 await session._handle_event(_playback_event("paused"))
1245 await session._handle_event(_device_event(is_active=False))
1246 session.backend._raise_if_app_controlled()
1247 _client_of(session).resume.assert_not_awaited()
1248
1249
1250async def test_no_session_is_started_while_the_app_holds_the_last_one(tmp_path: Path) -> None:
1251 """A replacement would claim the Connect device straight back off the user."""
1252 backend = _make_backend(tmp_path)
1253 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1254 with pytest.raises(SoloistAppControlError):
1255 await backend._acquire(TRACK_A, 0, "player1")
1256
1257
1258async def test_the_hold_on_a_new_session_expires(tmp_path: Path) -> None:
1259 """Coming back to Music Assistant later plays again without any fuss."""
1260 backend = _make_backend(tmp_path)
1261 backend._note_app_control(SoloistAppControl.TOOK_OVER)
1262 backend._app_control_until = time.monotonic() - 1
1263 backend._raise_if_app_controlled()
1264 assert backend._held_by_app() is None
1265
1266
1267async def test_an_audiobook_gives_up_on_capacity_instead_of_burning_chapters(
1268 tmp_path: Path,
1269) -> None:
1270 """Skipping ahead would cost the audiobook its availability and the caller its retry."""
1271 provider = _make_provider(tmp_path)
1272 calls: list[str] = []
1273
1274 async def _refuse(uri: str, *_args: Any, **_kwargs: Any) -> AsyncGenerator[bytes]:
1275 calls.append(uri)
1276 for _ in (): # never yields; only makes this an async generator
1277 yield b""
1278 raise SoloistAppControlError(provider, SoloistAppControl.TOOK_OVER)
1279
1280 provider.backend = MagicMock(stream_spotify_uri=_refuse)
1281 streamdetails = MagicMock(
1282 media_type=MediaType.AUDIOBOOK,
1283 data={"chapters": [TRACK_A, TRACK_B, "spotify:track:ccc"], "chapters_data": []},
1284 )
1285
1286 with pytest.raises(SoloistAppControlError):
1287 async for _ in provider.get_audio_stream(streamdetails):
1288 pass
1289 # the first chapter's refusal ends it: no chapter is skipped over
1290 assert calls == [TRACK_A]
1291
1292
1293def test_the_playback_device_is_named_apart_from_the_connect_one() -> None:
1294 """Two identically named devices in the Spotify app is what causes the takeovers."""
1295 assert SOLOIST_DEVICE_NAME != DEFAULT_PUBLISH_NAME
1296
1297
1298async def test_app_volume_change_is_pinned_back_to_unity(tmp_path: Path) -> None:
1299 """An off-unity volume set from the Spotify app is pinned back to 100."""
1300 session = _make_session(tmp_path)
1301 await session._handle_event(
1302 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=40), raw={})
1303 )
1304 _client_of(session).set_volume.assert_awaited_once_with(100)
1305 _client_of(session).set_volume.reset_mock()
1306 await session._handle_event(
1307 SoloistEvent(type="volume_changed", data=SoloistVolumeChanged(volume=100), raw={})
1308 )
1309 _client_of(session).set_volume.assert_not_awaited()
1310
1311
1312async def test_track_change_signals_the_queue_when_it_matches_the_next_item(
1313 tmp_path: Path,
1314) -> None:
1315 """Reaching a fed item tells the queue to start filling that item's buffer."""
1316 session = _make_session(tmp_path, queue_id="player1")
1317 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1318 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1319 queues = _queues_of(session)
1320 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1321 await session._handle_event(
1322 SoloistEvent(
1323 type="track_changed",
1324 data=SoloistTrackChanged(item=SoloistEntity(uri=TRACK_B, entity_type="track")),
1325 raw={},
1326 )
1327 )
1328 queues.prepare_next_audio_buffer.assert_called_once_with("player1")
1329
1330
1331async def test_track_change_to_another_item_signals_nothing(tmp_path: Path) -> None:
1332 """An item the queue is not asking for next must not trigger a prebuffer."""
1333 session = _make_session(tmp_path, queue_id="player1")
1334 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1335 queues = _queues_of(session)
1336 queues.get.return_value = MagicMock(next_item=_queue_item(TRACK_B), current_index=0)
1337 await session._handle_event(
1338 SoloistEvent(
1339 type="track_changed",
1340 data=SoloistTrackChanged(
1341 item=SoloistEntity(uri="spotify:track:surprise", entity_type="track")
1342 ),
1343 raw={},
1344 )
1345 )
1346 queues.prepare_next_audio_buffer.assert_not_called()
1347
1348
1349async def test_the_follower_of_the_streamed_item_is_fed(tmp_path: Path) -> None:
1350 """The item after the one being streamed is handed to the engine."""
1351 session = _make_session(tmp_path, queue_id="player1")
1352 streamdetails = MagicMock()
1353 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1354 follower = _queue_item(TRACK_B)
1355 queues = _queues_of(session)
1356 queues.get.return_value = MagicMock(current_index=3)
1357 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 3 else None
1358 queues.get_next_item.return_value = follower
1359 await session.feed_after(streamdetails, TRACK_A)
1360 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1361 assert TRACK_B in session._items
1362 assert session.has_pending is True
1363
1364
1365async def test_an_item_the_queue_resolved_elsewhere_is_not_fed(tmp_path: Path) -> None:
1366 """A track the queue will stream from another provider must not be queued here."""
1367 session = _make_session(tmp_path, queue_id="player1")
1368 streamdetails = MagicMock()
1369 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1370 # same track, but the queue already picked a different provider for it
1371 follower = _queue_item(TRACK_B, streamdetails=MagicMock(provider="tidal--x"))
1372 queues = _queues_of(session)
1373 queues.get.return_value = MagicMock(current_index=0)
1374 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1375 queues.get_next_item.return_value = follower
1376 await session.feed_after(streamdetails, TRACK_A)
1377 _client_of(session).add_to_queue.assert_not_awaited()
1378
1379
1380async def test_skipping_to_the_fed_item_keeps_the_session(tmp_path: Path) -> None:
1381 """A next-track lands on the item already fed, so the engine jumps instead of respawning."""
1382 backend = _make_backend(tmp_path)
1383 backend._server = MagicMock()
1384 backend._binary = Path("/nonexistent/soloist")
1385 session = _SoloistSession(backend, "player1")
1386 session._client = AsyncMock()
1387 session._logged_in = True
1388 backend._session = session
1389 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1390 playing.started.set()
1391 # fed one ahead and not reached yet, which is where a next-track goes
1392 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1393 session._pending.append(TRACK_B)
1394
1395 async def _engine_gets_there(**_kwargs: Any) -> None:
1396 await session._observe_current(TRACK_B, 200_000)
1397
1398 _client_of(session).skip_next.side_effect = _engine_gets_there
1399 got_session, got_item = await backend._acquire(TRACK_B, 0, "player1")
1400 # the same session, no respawn, and the item that was already queued
1401 assert got_session is session
1402 assert got_item is fed
1403 assert backend._session is session
1404 _client_of(session).skip_next.assert_awaited_once()
1405
1406
1407async def test_a_skip_drops_what_arrives_while_the_command_is_in_flight(
1408 tmp_path: Path,
1409) -> None:
1410 """
1411 Audio captured between the skip command and the engine's answer is dropped.
1412
1413 Only covers the marker's own window; what the pipeline still holds when the
1414 answer arrives is measured at the cut instead.
1415 """
1416 session = _make_session(tmp_path)
1417 leaving = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1418 leaving.started.set()
1419 target = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1420 session._pending.append(TRACK_B)
1421 captured: list[bytes] = []
1422
1423 async def _engine_gets_there(**_kwargs: Any) -> None:
1424 # the pipeline still holds the old track while the command is in flight
1425 session._write_if_wanted(b"\x01" * 32)
1426 await session._observe_current(TRACK_B, 200_000)
1427 # from here on the audio really is the new item's
1428 session._write_if_wanted(b"\x02" * 32)
1429
1430 _client_of(session).skip_next.side_effect = _engine_gets_there
1431 await session.skip_to(target)
1432 captured.extend(target._chunks)
1433 assert b"".join(captured) == b"\x02" * 32
1434 assert session._discard_until is None
1435
1436
1437async def test_a_skip_the_engine_never_reaches_fails(
1438 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1439) -> None:
1440 """A skip that does not land is an error, not a wait for the track to end."""
1441 monkeypatch.setattr(soloist_backend, "_STARTUP_TIMEOUT_S", 0.05)
1442 session = _make_session(tmp_path)
1443 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1444 session._pending.append(TRACK_B)
1445 with pytest.raises(AudioError, match="did not reach"):
1446 await session.skip_to(fed)
1447
1448
1449async def test_a_fed_item_the_engine_has_not_reached_is_not_served(tmp_path: Path) -> None:
1450 """Skipping to an already-fed item must not hand over a channel that fills later."""
1451 session = _make_session(tmp_path)
1452 fed = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1453 # fed, but the engine is still on the previous track
1454 assert session.item_for(TRACK_B) is None
1455 # once the engine gets there it is servable
1456 await session._observe_current(TRACK_B, 200_000)
1457 assert session.item_for(TRACK_B) is fed
1458
1459
1460def test_the_shaper_only_emits_whole_frames() -> None:
1461 """A read that ends mid-frame must never split a frame across two items."""
1462 shaper = soloist_backend._CaptureShaper()
1463 # the session's first bytes are infrastructure silence, and are dropped
1464 assert shaper.shape(b"\x00" * 4096) == b""
1465 # a mis-aligned read emits whole frames and carries the remainder
1466 first = shaper.shape(b"\x01" * (_FRAME_BYTES + 3))
1467 assert len(first) == _FRAME_BYTES
1468 # which is then completed by the next read, losing nothing
1469 second = shaper.shape(b"\x02" * (_FRAME_BYTES - 3))
1470 assert len(second) == _FRAME_BYTES
1471 assert second[:3] == b"\x01" * 3
1472 # an aligned read passes straight through
1473 assert shaper.shape(b"\x03" * _FRAME_BYTES) == b"\x03" * _FRAME_BYTES
1474
1475
1476def test_the_shaper_trims_lead_silence_only_once() -> None:
1477 """Silence after the audio has started is content, not pre-roll."""
1478 shaper = soloist_backend._CaptureShaper()
1479 assert shaper.shape(b"\x01" * _FRAME_BYTES) == b"\x01" * _FRAME_BYTES
1480 silence = b"\x00" * _FRAME_BYTES
1481 assert shaper.shape(silence) == silence
1482
1483
1484async def test_only_whole_sample_frames_are_handed_over(
1485 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1486) -> None:
1487 """A read that ends mid-frame must not split a frame across two items."""
1488 session = _make_session(tmp_path)
1489 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1490 item.started.set()
1491 item.claim()
1492 session._demand_started = True
1493 session._sink_running = True
1494 # two reads that are each mis-aligned but whole together
1495 reads = [b"\x01" * (_FRAME_BYTES + 3), b"\x02" * (_FRAME_BYTES - 3), b""]
1496 reader = MagicMock()
1497
1498 async def _read(_size: int) -> bytes:
1499 return reads.pop(0) if reads else b""
1500
1501 reader.read = _read
1502 session._reader = reader
1503 monkeypatch.setattr(soloist_backend, "_PACE_RATE", 1000.0)
1504 await session._read_capture()
1505 # every write was frame-aligned, and no byte was lost
1506 assert item.buffered % _FRAME_BYTES == 0
1507 assert item.buffered == _FRAME_BYTES * 2
1508
1509
1510async def test_an_already_known_item_is_not_fed_twice(tmp_path: Path) -> None:
1511 """An item the session already plays or was fed is not queued again."""
1512 session = _make_session(tmp_path, queue_id="player1")
1513 session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1514 streamdetails = MagicMock()
1515 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1516 queues = _queues_of(session)
1517 queues.get.return_value = MagicMock(current_index=0)
1518 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1519 queues.get_next_item.return_value = _queue_item(TRACK_B)
1520 await session.feed_after(streamdetails, TRACK_A)
1521 _client_of(session).add_to_queue.assert_not_awaited()
1522
1523
1524async def test_only_tracks_are_fed_ahead(tmp_path: Path) -> None:
1525 """A podcast episode or audiobook chapter is played on its own, never stitched."""
1526 session = _make_session(tmp_path, queue_id="player1")
1527 await session.feed_after(MagicMock(), "spotify:episode:xyz")
1528 _client_of(session).add_to_queue.assert_not_awaited()
1529
1530
1531async def test_a_non_spotify_follower_is_not_fed(tmp_path: Path) -> None:
1532 """The run simply ends where the queue leaves this provider."""
1533 session = _make_session(tmp_path, queue_id="player1")
1534 streamdetails = MagicMock()
1535 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1536 follower = MagicMock(
1537 media_item=MagicMock(media_type=MediaType.TRACK, provider="tidal--x"), streamdetails=None
1538 )
1539 follower.media_item.provider_mappings = []
1540 queues = _queues_of(session)
1541 queues.get.return_value = MagicMock(current_index=0)
1542 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1543 queues.get_next_item.return_value = follower
1544 await session.feed_after(streamdetails, TRACK_A)
1545 _client_of(session).add_to_queue.assert_not_awaited()
1546
1547
1548async def test_a_library_item_is_fed_through_its_spotify_mapping(tmp_path: Path) -> None:
1549 """A library track is fed with the item id this provider instance knows it by."""
1550 session = _make_session(tmp_path, queue_id="player1")
1551 streamdetails = MagicMock()
1552 playing = _queue_item(TRACK_A, streamdetails=streamdetails)
1553 follower = MagicMock(
1554 media_item=MagicMock(media_type=MediaType.TRACK, provider="library", item_id="42"),
1555 streamdetails=None,
1556 )
1557 follower.media_item.provider_mappings = [
1558 MagicMock(provider_instance="other--y", item_id="wrong"),
1559 MagicMock(provider_instance="spotify--test", item_id="bbb"),
1560 ]
1561 queues = _queues_of(session)
1562 queues.get.return_value = MagicMock(current_index=0)
1563 queues.get_item.side_effect = lambda _queue_id, index: playing if index == 0 else None
1564 queues.get_next_item.return_value = follower
1565 await session.feed_after(streamdetails, TRACK_A)
1566 _client_of(session).add_to_queue.assert_awaited_once_with(TRACK_B)
1567
1568
1569@pytest.mark.parametrize(
1570 ("provider_option", "player_setting", "expected"),
1571 [
1572 (True, "enabled", True),
1573 # the player's own switch decides first: off means nobody normalizes,
1574 # not that the job passes to Spotify
1575 (True, "disabled", False),
1576 (False, "enabled", False),
1577 (False, "disabled", False),
1578 ],
1579)
1580def test_who_normalizes_needs_both_switches(
1581 tmp_path: Path,
1582 monkeypatch: pytest.MonkeyPatch,
1583 provider_option: bool,
1584 player_setting: str,
1585 expected: bool,
1586) -> None:
1587 """The engine normalizes only when the provider option and the player agree."""
1588 session = _make_session(tmp_path, queue_id="player1")
1589 monkeypatch.setattr(
1590 type(session.backend.provider),
1591 "spotify_normalization_configured",
1592 property(lambda _self: provider_option),
1593 )
1594 cast("MagicMock", session.mass.config).get_effective_player_queue_config_value = MagicMock(
1595 return_value=player_setting
1596 )
1597 assert session._engine_normalization_enabled() is expected
1598
1599
1600def test_a_running_session_answers_for_what_the_engine_is_doing(tmp_path: Path) -> None:
1601 """
1602 The engine reads its settings at startup, so a later toggle must not split them.
1603
1604 Otherwise the streams core would start normalizing on top of audio the engine
1605 is still normalizing, or stop while it no longer is.
1606 """
1607 backend = _make_backend(tmp_path)
1608 provider = backend.provider
1609 streamdetails = _streamdetails_for(queue_id="player1")
1610 # nothing playing yet: the configuration is all there is to go on
1611 before_any_session = backend.session_normalizes(streamdetails)
1612 session = _SoloistSession(backend, "player1")
1613 session.engine_normalizes = True
1614 backend._session = session
1615 while_playing = backend.session_normalizes(streamdetails)
1616 # ... and a session that has been torn down no longer speaks for the engine
1617 session._stopped = True
1618 after_teardown = backend.session_normalizes(streamdetails)
1619 assert before_any_session is None
1620 assert while_playing is True
1621 assert after_teardown is None
1622 assert (
1623 provider.delivers_normalized_audio(streamdetails)
1624 is provider.spotify_normalization_configured
1625 )
1626
1627
1628async def test_short_delivery_is_rejected_as_incomplete(tmp_path: Path) -> None:
1629 """PCM that stops well short of the item's duration is rejected."""
1630 session = _make_session(tmp_path)
1631 item = _ItemAudio(TRACK_A, session)
1632 item.playing_seen = True
1633 item.duration_ms = 200_000
1634 item.last_position_ms = 100_000
1635 with pytest.raises(AudioError, match="incomplete"):
1636 await session.validate_item(item)
1637
1638
1639async def test_missing_position_is_rejected_as_incomplete(tmp_path: Path) -> None:
1640 """Without any position report there is no evidence the item played out."""
1641 session = _make_session(tmp_path)
1642 item = _ItemAudio(TRACK_A, session)
1643 item.playing_seen = True
1644 item.duration_ms = 200_000
1645 with pytest.raises(AudioError, match="incomplete"):
1646 await session.validate_item(item)
1647
1648
1649async def test_short_item_cannot_pass_at_position_zero(tmp_path: Path) -> None:
1650 """The tolerance never spans a whole item, so a short item cannot pass unplayed."""
1651 session = _make_session(tmp_path)
1652 item = _ItemAudio(TRACK_A, session)
1653 item.playing_seen = True
1654 item.duration_ms = 8_000
1655 item.last_position_ms = 0
1656 with pytest.raises(AudioError, match="incomplete"):
1657 await session.validate_item(item)
1658
1659
1660async def test_an_item_that_never_played_is_rejected(tmp_path: Path) -> None:
1661 """An item the engine never reported playing is a failure, whatever was delivered."""
1662 session = _make_session(tmp_path)
1663 item = _ItemAudio(TRACK_A, session)
1664 item.duration_ms = 200_000
1665 item.last_position_ms = 200_000
1666 with pytest.raises(AudioError, match="never started playing"):
1667 await session.validate_item(item)
1668
1669
1670async def test_a_duration_less_item_is_not_judged(tmp_path: Path) -> None:
1671 """Without a duration there is nothing to judge completeness against."""
1672 session = _make_session(tmp_path)
1673 item = _ItemAudio(TRACK_A, session)
1674 item.playing_seen = True
1675 await session.validate_item(item)
1676
1677
1678def test_an_unread_session_expires(tmp_path: Path) -> None:
1679 """A session no item stream reads from is ended so its daemon does not linger."""
1680 session = _make_session(tmp_path)
1681 session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1682 session._expire_idle()
1683 assert session._idle_since is not None
1684 assert session.usable is True
1685 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S - 1
1686 session._expire_idle()
1687 assert session.usable is False
1688
1689
1690def test_a_session_being_read_never_expires(tmp_path: Path) -> None:
1691 """An item stream reading the session keeps it alive indefinitely."""
1692 session = _make_session(tmp_path)
1693 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1694 item.claim()
1695 session._idle_since = time.monotonic() - _IDLE_TIMEOUT_S * 10
1696 session._expire_idle()
1697 assert session.usable is True
1698
1699
1700def test_pre_roll_silence_is_dropped_a_whole_frame_at_a_time() -> None:
1701 """
1702 Trimming pre-roll must leave the audio on the session's frame grid.
1703
1704 A FIFO read is not always a whole number of frames, and dropping a partial
1705 one would shift every sample that follows for the rest of the session.
1706 """
1707 shaper = _CaptureShaper()
1708 # pre-roll that ends mid-frame: the real audio starts at byte 1024
1709 assert shaper.shape(b"\x00" * 1021) == b""
1710 audio = bytes(range(1, 9)) * 4
1711 shaped = shaper.shape(b"\x00" * 3 + audio)
1712 assert shaped == audio
1713 assert shaper._lead_skipped % _FRAME_BYTES == 0
1714
1715
1716async def test_a_refused_skip_does_not_leave_the_audio_discarded(tmp_path: Path) -> None:
1717 """
1718 A skip that never landed must not keep the session dropping its audio.
1719
1720 The marker silences everything the session captures, so a command that
1721 failed has to clear it on the way out.
1722 """
1723 session = _make_session(tmp_path)
1724 client = cast("MagicMock", session._client)
1725 client.skip_next = AsyncMock(side_effect=TimeoutError)
1726 item = _ItemAudio(TRACK_B, session)
1727
1728 with pytest.raises(AudioError, match="would not skip"):
1729 await session.skip_to(item)
1730
1731 assert session._discard_until is None
1732
1733
1734async def test_a_daemon_that_will_not_die_is_reported_and_released(tmp_path: Path) -> None:
1735 """A close that could not terminate the daemon still finishes the teardown."""
1736 session = _make_session(tmp_path)
1737 proc = cast("MagicMock", session._proc)
1738 proc.close = AsyncMock()
1739 # AsyncProcess.close() gives up after a handful of kill attempts
1740 proc.returncode = None
1741 with patch.object(session.logger, "warning") as warning:
1742 await session.stop()
1743 assert warning.called
1744 assert session._teardown_done is True
1745 assert session._proc is None
1746
1747
1748async def test_a_cancelled_teardown_still_closes_the_daemon(tmp_path: Path) -> None:
1749 """
1750 A cancelled teardown must leave the retry something to close.
1751
1752 Dropping the references first is how a daemon survives to hold the data
1753 directory, which every later session is then refused for.
1754 """
1755 session = _make_session(tmp_path)
1756 proc = cast("MagicMock", session._proc)
1757 sink = cast("AsyncMock", session._sink)
1758
1759 async def _never_returns() -> None:
1760 await asyncio.Event().wait()
1761
1762 proc.close = _never_returns
1763 task = asyncio.create_task(session.stop())
1764 await asyncio.sleep(0.01)
1765 task.cancel()
1766 with suppress(asyncio.CancelledError):
1767 await task
1768 # the teardown did not finish, so nothing was dropped and it can be redone
1769 unfinished = session._teardown_done
1770 kept_proc = session._proc
1771 kept_sink = session._sink
1772 proc.close = AsyncMock()
1773 proc.returncode = 0
1774 await session.stop()
1775 assert unfinished is False
1776 assert kept_proc is proc
1777 assert kept_sink is sink
1778 assert session._teardown_done is True
1779 assert session._proc is None
1780 assert session._sink is None
1781 proc.close.assert_awaited()
1782 sink.unload.assert_awaited()
1783
1784
1785def test_a_failed_session_is_torn_down(tmp_path: Path) -> None:
1786 """A session that fails is discarded, so its daemon does not keep playing to nobody."""
1787 session = _make_session(tmp_path)
1788 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1789 item.claim()
1790 session._fail("audio stalled")
1791 assert session.usable is False
1792 # every waiting item is released and the teardown is scheduled
1793 assert item._closed is True
1794 # a startup wait must not sit out its timeout on a session that already failed
1795 assert item.started.is_set() is True
1796 discard = cast("MagicMock", session.mass.create_task)
1797 discard.assert_called_once_with(session.backend.discard_session, session)
1798 # a second failure does not queue a second teardown
1799 session._fail("and again")
1800 assert session._error == "audio stalled"
1801 assert discard.call_count == 1
1802
1803
1804async def test_an_item_the_engine_skipped_past_fails_instead_of_hanging(
1805 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1806) -> None:
1807 """A claimed channel the engine never reaches gives up rather than blocking forever."""
1808 monkeypatch.setattr(soloist_backend, "_READ_SLICE_S", 0.01)
1809 monkeypatch.setattr(soloist_backend, "_STALL_TIMEOUT_S", 0.05)
1810 session = _make_session(tmp_path)
1811 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1812 item.claim()
1813 # the engine is playing something else, so nothing is ever written here
1814 session._items["spotify:track:other"] = session._current = _ItemAudio(
1815 "spotify:track:other", session
1816 )
1817 with pytest.raises(AudioError, match="no audio"):
1818 async for _ in item.read():
1819 pass
1820
1821
1822async def test_adopt_paired_session_copies_into_the_canonical_dir(
1823 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1824) -> None:
1825 """A session paired by the setup flow is adopted into the per-instance data dir."""
1826 storage = tmp_path / "storage"
1827 pending = storage / "spotify" / "pairing" / "flow1"
1828 pending.mkdir(parents=True)
1829 (pending / "session.bin").write_bytes(b"session")
1830 prov = _make_provider(tmp_path, {CONF_SOLOIST_SESSION_DIR: "spotify/pairing/flow1"})
1831 update_setup_data = MagicMock()
1832 monkeypatch.setattr(prov, "_update_setup_data", update_setup_data)
1833 backend = SoloistBackend(prov)
1834 await backend._adopt_paired_session()
1835 canonical = storage / "spotify" / "spotify--test" / SOLOIST_DATA_DIR_NAME
1836 assert (canonical / "session.bin").read_bytes() == b"session"
1837 # a copy, not a move: the flow-private source must survive a failed
1838 # provider load so the setup flow can retry (the flow removes it at its end)
1839 assert (pending / "session.bin").exists()
1840 update_setup_data.assert_called_once_with(CONF_SOLOIST_SESSION_DIR, None)
1841
1842
1843def test_the_engine_is_told_not_to_normalize(tmp_path: Path) -> None:
1844 """MA normalizes this audio itself, so the engine's own normalization is switched off."""
1845 backend = _make_backend(tmp_path)
1846 prefs = backend._data_dir / "settings" / "Users" / "alice-user" / "prefs"
1847 prefs.parent.mkdir(parents=True)
1848 prefs.write_text("some.engine.key=1\n", encoding="utf-8")
1849 backend._prepare_data_dir(normalize=False)
1850 content = prefs.read_text(encoding="utf-8").splitlines()
1851 assert "some.engine.key=1" in content
1852 assert "audio.normalize_v2=false" in content
1853 # MA mixes the queue's crossfade itself, so the engine's own is always off
1854 assert "audio.crossfade_v2=false" in content
1855 # the ceiling is stated rather than left to the engine's own default
1856 assert "audio.play_bitrate_enumeration=5" in content
1857 assert "audio.play_bitrate_non_metered_enumeration=5" in content
1858 assert "audio.play_bitrate_non_metered_migrated=true" in content
1859
1860
1861def test_disabling_crossfade_writes_the_boolean(tmp_path: Path) -> None:
1862 """Crossfade off is written explicitly, so a stale 'on' cannot survive."""
1863 backend = _make_backend(tmp_path)
1864 prefs = backend._data_dir / "settings" / "prefs"
1865 prefs.parent.mkdir(parents=True)
1866 prefs.write_text("audio.crossfade_v2=true\naudio.crossfade.time_v2=8000\n", encoding="utf-8")
1867 backend._prepare_data_dir(normalize=False)
1868 content = prefs.read_text(encoding="utf-8").splitlines()
1869 assert "audio.crossfade_v2=false" in content
1870 assert not any(line.startswith("audio.crossfade.time_v2") for line in content)
1871
1872
1873async def test_setup_requires_an_api_key(tmp_path: Path) -> None:
1874 """Without a stored API key the user must be sent back through the setup flow."""
1875 backend = _make_backend(tmp_path)
1876 with pytest.raises(LoginFailed) as err:
1877 await backend.setup()
1878 assert err.value.translation_key == "soloist_pairing_required"
1879
1880
1881async def test_setup_requires_a_paired_session(
1882 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
1883) -> None:
1884 """An API key without a paired session also routes back to the setup flow."""
1885 backend = _make_backend(tmp_path, {CONF_SOLOIST_API_KEY: "k" * 20, CONF_SOLOIST_CONSENT: True})
1886 _install_fake_binary_manager(monkeypatch)
1887 with pytest.raises(LoginFailed) as err:
1888 await backend.setup()
1889 assert err.value.translation_key == "soloist_pairing_required"
1890
1891
1892async def test_streaming_without_setup_is_refused(tmp_path: Path) -> None:
1893 """A backend whose setup never ran refuses to stream instead of half-starting."""
1894 backend = _make_backend(tmp_path)
1895 with pytest.raises(AudioError, match="not started"):
1896 async for _ in backend.stream_spotify_uri(TRACK_A):
1897 pass
1898
1899
1900def test_session_present_detection(tmp_path: Path) -> None:
1901 """Only the engine's per-account state counts as paired."""
1902 data_dir = tmp_path / "soloist-data"
1903 assert soloist_session_present(data_dir) is False
1904 data_dir.mkdir()
1905 (data_dir / WS_ADDR_FILE).write_text("127.0.0.1", encoding="utf-8")
1906 (data_dir / WS_PORT_FILE).write_text("1234", encoding="utf-8")
1907 assert soloist_session_present(data_dir) is False
1908 # everything a spawn leaves behind outlives the pairing it ran on: the engine
1909 # keeps its identity, lock, cache and crash handler in the data dir even
1910 # though it is given a cache dir of its own, and Music Assistant writes the
1911 # prefs there before every spawn
1912 (data_dir / "settings").mkdir()
1913 (data_dir / "settings" / "prefs").write_text("audio.normalize_v2=false\n", encoding="utf-8")
1914 (data_dir / ".device_id").write_text("6b6c2a07", encoding="utf-8")
1915 (data_dir / ".lock").write_bytes(b"")
1916 (data_dir / "cache" / "Users" / "spotify-user-user").mkdir(parents=True)
1917 (data_dir / "crashpad").mkdir()
1918 assert soloist_session_present(data_dir) is False
1919 (data_dir / "settings" / "Users" / "spotify-user-user").mkdir(parents=True)
1920 assert soloist_session_present(data_dir) is True
1921
1922
1923async def test_a_skip_drops_the_audio_still_in_flight(tmp_path: Path) -> None:
1924 """The item jumped to opens with its own audio, not the tail of the one left behind."""
1925 session = _make_session(tmp_path)
1926 left_behind = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1927 session._current = left_behind
1928 left_behind.started.set()
1929 left_behind.claim()
1930 jumped_to = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
1931 session._pending.append(TRACK_B)
1932 session._discard_until = TRACK_B
1933 with _capture_holding(session, fifo_bytes=2 * _FRAME_BYTES, reader_bytes=2 * _FRAME_BYTES):
1934 await session._observe_current(TRACK_B, 200_000)
1935 assert session._stale_budget == 4 * _FRAME_BYTES
1936 session._write_if_wanted(b"s" * (4 * _FRAME_BYTES))
1937 session._write_if_wanted(b"n" * (2 * _FRAME_BYTES))
1938 jumped_to.claim()
1939 jumped_to.close()
1940 assert b"".join([chunk async for chunk in jumped_to.read()]) == b"n" * (2 * _FRAME_BYTES)
1941
1942
1943async def test_a_skip_drops_the_stale_audio_across_reads(tmp_path: Path) -> None:
1944 """A budget larger than one read keeps dropping, and resumes on a frame boundary."""
1945 session = _make_session(tmp_path)
1946 session._stale_budget = 3 * _FRAME_BYTES
1947 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1948 item.claim()
1949 session._write_if_wanted(b"s" * (2 * _FRAME_BYTES))
1950 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1951 item.close()
1952 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1953
1954
1955async def test_the_marker_spends_an_earlier_jumps_budget(tmp_path: Path) -> None:
1956 """What the marker drops still counts against a budget left from an earlier jump."""
1957 session = _make_session(tmp_path)
1958 session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
1959 session._stale_budget = 4 * _FRAME_BYTES
1960 session._discard_until = TRACK_B
1961 session._write_if_wanted(b"s" * (3 * _FRAME_BYTES))
1962 assert session._stale_budget == _FRAME_BYTES
1963 # a refused command leaves only what is genuinely still in flight to drop
1964 session._discard_until = None
1965 session._write_if_wanted(b"s" * _FRAME_BYTES + b"n" * _FRAME_BYTES)
1966 item = session._current
1967 item.claim()
1968 item.close()
1969 assert b"".join([chunk async for chunk in item.read()]) == b"n" * _FRAME_BYTES
1970
1971
1972async def test_a_natural_cut_keeps_the_audio_in_flight(tmp_path: Path) -> None:
1973 """Nothing is dropped without a jump: what is in flight is the continuation."""
1974 session = _make_session(tmp_path)
1975 playing = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
1976 session._current = playing
1977 playing.started.set()
1978 playing.claim()
1979 with _capture_holding(session, fifo_bytes=4 * _FRAME_BYTES, reader_bytes=4 * _FRAME_BYTES):
1980 await session._observe_current(TRACK_B, 200_000)
1981 assert session._stale_budget == 0
1982
1983
1984def test_stale_bytes_spans_both_buffers_in_whole_frames(tmp_path: Path) -> None:
1985 """The in-flight measure covers the FIFO and the reader, and never splits a frame."""
1986 session = _make_session(tmp_path)
1987 with _capture_holding(session, fifo_bytes=3 * _FRAME_BYTES + 3, reader_bytes=2 * _FRAME_BYTES):
1988 assert session._stale_bytes() == 5 * _FRAME_BYTES
1989
1990
1991def test_stale_bytes_falls_back_when_the_reader_cannot_be_sized(tmp_path: Path) -> None:
1992 """Losing the reader's internal view drops extra rather than leaving audio behind."""
1993 session = _make_session(tmp_path)
1994 with _capture_holding(session, fifo_bytes=0, reader_bytes=None):
1995 assert session._stale_bytes() == 6 * _READ_CHUNK_SIZE
1996
1997
1998async def test_a_channel_abandoned_at_the_cut_stops_holding_the_cushion(
1999 tmp_path: Path,
2000) -> None:
2001 """A skip closes the channel first and only then unwinds its stream."""
2002 session = _make_session(tmp_path)
2003 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2004 item.started.set()
2005 item.claim()
2006 item.write(b"x" * 4096)
2007 # the cut lands while the abandoned stream is still unwinding
2008 await session._observe_current(TRACK_B, 200_000)
2009 assert session._retained_bytes() == 4096
2010 item.release()
2011 assert session._retained_bytes() == 0
2012
2013
2014def test_an_abandoned_channel_stops_holding_the_cushion(tmp_path: Path) -> None:
2015 """A channel skipped away from frees its buffer instead of gating the sink for good."""
2016 session = _make_session(tmp_path)
2017 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
2018 item.claim()
2019 item.write(b"x" * 4096)
2020 assert session._retained_bytes() == 4096
2021 # the stream is gone, then the cut closes the channel
2022 item.release()
2023 item.close()
2024 assert session._retained_bytes() == 0
2025
2026
2027async def test_a_channel_still_being_read_keeps_its_tail(tmp_path: Path) -> None:
2028 """Closing the playing item at a cut must not discard what its stream is still owed."""
2029 session = _make_session(tmp_path)
2030 item = session._items[TRACK_A] = _ItemAudio(TRACK_A, session)
2031 item.claim()
2032 item.write(b"tail" * 4)
2033 item.close()
2034 assert item.buffered == 16
2035 assert b"".join([chunk async for chunk in item.read()]) == b"tail" * 4
2036
2037
2038@contextmanager
2039def _capture_holding(
2040 session: _SoloistSession, *, fifo_bytes: int, reader_bytes: int | None
2041) -> Iterator[None]:
2042 """
2043 Give the session a real capture FIFO and a reader holding the given amounts.
2044
2045 A real pipe is used so the byte count comes from the same ioctl the backend
2046 relies on. Pass ``reader_bytes=None`` for a reader whose buffer cannot be read.
2047 """
2048 read_fd, write_fd = os.pipe()
2049 try:
2050 if fifo_bytes:
2051 os.write(write_fd, bytes(fifo_bytes))
2052 pipe = MagicMock()
2053 pipe.fileno.return_value = read_fd
2054 transport = MagicMock()
2055 transport.get_extra_info.return_value = pipe
2056 session._transport = transport
2057 reader = MagicMock(spec=[]) if reader_bytes is None else MagicMock()
2058 if reader_bytes is not None:
2059 reader._buffer = bytearray(reader_bytes)
2060 session._reader = reader
2061 yield
2062 finally:
2063 session._transport = None
2064 session._reader = None
2065 os.close(read_fd)
2066 os.close(write_fd)
2067
2068
2069def _stdout_of(*lines: str) -> MagicMock:
2070 """Return a process mock whose stdout yields the given daemon log lines."""
2071
2072 async def _iter_stdout() -> AsyncGenerator[str]:
2073 for line in lines:
2074 yield line
2075
2076 proc = MagicMock()
2077 proc.iter_stdout = _iter_stdout
2078 return proc
2079
2080
2081def _make_provider(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SpotifyProvider:
2082 """Return a SpotifyProvider (bypassing __init__) with the given setup_data."""
2083 prov = object.__new__(SpotifyProvider)
2084 config = MagicMock(instance_id="spotify--test")
2085 config.get_value = MagicMock(return_value=None)
2086 config.values = {}
2087 prov.config = config
2088 prov.manifest = MagicMock(domain="spotify")
2089 prov.logger = MagicMock()
2090 prov.available = True
2091 mass = MagicMock()
2092 mass.storage_path = str(tmp_path / "storage")
2093 mass.cache_path = str(tmp_path / "cache")
2094 # get_setup_value reads the live setup_data blob from the store
2095 mass.config.get = MagicMock(return_value=setup_data or {})
2096 mass.config.get_raw_provider_config_value = MagicMock(return_value=None)
2097 # the store keeps values encrypted; decrypt is an identity map for the test
2098 mass.config.decrypt_string = MagicMock(side_effect=lambda value: value)
2099 prov.mass = mass
2100 return prov
2101
2102
2103def _make_backend(tmp_path: Path, setup_data: dict[str, Any] | None = None) -> SoloistBackend:
2104 """Return a SoloistBackend on a mocked provider."""
2105 return SoloistBackend(_make_provider(tmp_path, setup_data))
2106
2107
2108def _make_session(tmp_path: Path, queue_id: str | None = "player1") -> _SoloistSession:
2109 """Return a session with its process/sink/client replaced by mocks."""
2110 session = _SoloistSession(_make_backend(tmp_path), queue_id)
2111 session._sink = AsyncMock()
2112 session._client = AsyncMock()
2113 session._proc = MagicMock(returncode=None)
2114 # a session under test is past the engine's login and has claimed the
2115 # Connect device, unless a test says otherwise
2116 session._logged_in = True
2117 session._was_active = True
2118 return session
2119
2120
2121def _streamdetails_for(
2122 *,
2123 queue_id: str | None = "player1",
2124 uri: str = TRACK_A,
2125 media_type: MediaType = MediaType.TRACK,
2126) -> StreamDetails:
2127 """Return stream details for a Spotify item served by the test instance."""
2128 return StreamDetails(
2129 provider="spotify--test",
2130 item_id=uri.rsplit(":", 1)[1],
2131 audio_format=AudioFormat(content_type=ContentType.PCM_S16LE),
2132 media_type=media_type,
2133 queue_id=queue_id,
2134 )
2135
2136
2137def _make_item(tmp_path: Path, uri: str) -> _ItemAudio:
2138 """Return a bare item channel on a mocked session."""
2139 return _ItemAudio(uri, _make_session(tmp_path))
2140
2141
2142def _queue_item(uri: str, streamdetails: Any = None) -> MagicMock:
2143 """Return a queue item stand-in for a Spotify track on the test instance."""
2144 item_id = uri.rsplit(":", 1)[1]
2145 media_item = MagicMock(media_type=MediaType.TRACK, provider="spotify--test", item_id=item_id)
2146 media_item.provider_mappings = []
2147 return MagicMock(
2148 media_item=media_item,
2149 queue_item_id=f"qi-{item_id}",
2150 streamdetails=streamdetails,
2151 )
2152
2153
2154async def _wait_for(predicate: Callable[[], bool], timeout: float = 2.0) -> None:
2155 """Wait until the predicate holds, so a background task can get there."""
2156 loop = asyncio.get_running_loop()
2157 deadline = loop.time() + timeout
2158 while loop.time() < deadline:
2159 if predicate():
2160 return
2161 await asyncio.sleep(0.01)
2162 raise AssertionError("condition not met within timeout")
2163
2164
2165def _client_of(session: _SoloistSession) -> AsyncMock:
2166 """Return the session's mocked WebSocket client."""
2167 return cast("AsyncMock", session._client)
2168
2169
2170def _current_of(session: _SoloistSession) -> _ItemAudio:
2171 """Return the channel the session is playing, which the caller knows exists."""
2172 item = session._current
2173 assert item is not None
2174 return item
2175
2176
2177def _sink_of(session: _SoloistSession) -> AsyncMock:
2178 """Return the session's mocked capture sink."""
2179 return cast("AsyncMock", session._sink)
2180
2181
2182def _queues_of(session: _SoloistSession) -> MagicMock:
2183 """Return the mocked player_queues controller the session consults."""
2184 return cast("MagicMock", session.mass.player_queues)
2185
2186
2187def _auth_event(*, logged_in: bool, is_active: bool = True) -> SoloistEvent:
2188 """Return an auth_state event with the given login and active-device state."""
2189 return SoloistEvent(
2190 type="auth_state",
2191 data=SoloistAuthState(logged_in=logged_in, is_active=is_active),
2192 raw={},
2193 )
2194
2195
2196def _device_event(*, is_active: bool) -> SoloistEvent:
2197 """Return a device_changed event with the given active-device state."""
2198 return SoloistEvent(
2199 type="device_changed", data=SoloistDeviceChanged(is_active=is_active), raw={}
2200 )
2201
2202
2203def _playback_event(status: str, position_ms: int = 0) -> SoloistEvent:
2204 """Return a playback_state event for the current item with the given status."""
2205 return SoloistEvent(
2206 type="playback_state",
2207 data=SoloistPlaybackState(
2208 status=status,
2209 item=SoloistEntity(uri=TRACK_A, entity_type="track"),
2210 position=SoloistPosition(position_ms=position_ms, timestamp_ms=0),
2211 ),
2212 raw={},
2213 )
2214
2215
2216def _install_fake_binary_manager(monkeypatch: pytest.MonkeyPatch) -> None:
2217 """Replace the shared binary manager so no download or exec is attempted."""
2218 manager = MagicMock()
2219 manager.ensure_fresh = AsyncMock(return_value=Path("/nonexistent/soloist"))
2220 monkeypatch.setattr(soloist_backend, "SoloistBinaryManager", MagicMock(return_value=manager))
2221
2222
2223async def test_seeking_the_playing_item_keeps_the_session(tmp_path: Path) -> None:
2224 """The engine is moved where it stands rather than the session being respawned."""
2225 session = _make_session(tmp_path)
2226 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2227 playing.started.set()
2228 playing.claim()
2229 playing.duration_ms = 260_000
2230 playing.playing_seen = True
2231 playing.observe_position(30_000)
2232 # the pre-seek audio nobody may hear again
2233 playing.write(b"\x01" * 32)
2234
2235 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2236 _current_of(session).observe_position(position_ms)
2237
2238 _client_of(session).seek.side_effect = _engine_seeks
2239 item = await session.seek_current(TRACK_A, 120_000)
2240
2241 assert item is not playing
2242 assert session.current is item
2243 assert item.claimed
2244 # what the track is stays with it; where it was does not
2245 assert item.duration_ms == 260_000
2246 assert item.playing_seen
2247 assert item.started_at_ms == 120_000
2248 # the outgoing channel is closed, which is what ends the stream reading it
2249 assert playing._closed
2250 _client_of(session).seek.assert_awaited_once_with(120_000, await_result=True)
2251
2252
2253async def test_a_seek_of_the_playing_item_is_sent_only_once(tmp_path: Path) -> None:
2254 """
2255 A landed seek is never repeated: a repeat restarts the item, audibly.
2256
2257 The engine answers late on purpose, so a re-send loop around the wait would
2258 have fired several times over before the confirmation arrives.
2259 """
2260 session = _make_session(tmp_path)
2261 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2262 playing.started.set()
2263 playing.observe_position(30_000)
2264 pending: list[asyncio.Task[None]] = []
2265
2266 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2267 async def _confirm_late() -> None:
2268 await asyncio.sleep(0.05)
2269 _current_of(session).observe_position(position_ms)
2270
2271 pending.append(asyncio.create_task(_confirm_late()))
2272
2273 _client_of(session).seek.side_effect = _engine_seeks
2274 with patch.object(soloist_backend, "_SEEK_RETRY_INTERVAL_S", 0.01):
2275 await session.seek_current(TRACK_A, 120_000)
2276 await asyncio.gather(*pending)
2277 assert _client_of(session).seek.await_count == 1
2278
2279
2280async def test_seeking_back_is_not_confirmed_by_the_position_seeked_away_from(
2281 tmp_path: Path,
2282) -> None:
2283 """A report still describing the pre-seek position cannot land a backward seek."""
2284 session = _make_session(tmp_path)
2285 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2286 playing.started.set()
2287 playing.duration_ms = 260_000
2288 playing.observe_position(200_000)
2289
2290 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2291 item = _current_of(session)
2292 # a report from before the seek is still in flight; it sits above the
2293 # target's tolerance window and must not pass for the landing
2294 item.observe_position(200_000)
2295 assert not item.seek_confirmed.is_set()
2296 item.observe_position(position_ms)
2297 item.observe_position(position_ms + 2)
2298
2299 _client_of(session).seek.side_effect = _engine_seeks
2300 item = await session.seek_current(TRACK_A, 60_000)
2301 assert item.started_at_ms == 60_002
2302 # and the pre-seek report is not left standing in for progress this item
2303 # never made, which at_own_end and the completeness check would believe
2304 assert item.last_position_ms == 60_002
2305
2306
2307async def test_audio_in_flight_across_an_in_place_seek_is_dropped(tmp_path: Path) -> None:
2308 """Only audio from past the seek reaches the fresh channel."""
2309 session = _make_session(tmp_path)
2310 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2311 playing.started.set()
2312 playing.observe_position(30_000)
2313
2314 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2315 # still rendering the position being left behind
2316 session._write_if_wanted(b"\x01" * 32)
2317 _current_of(session).observe_position(position_ms)
2318
2319 _client_of(session).seek.side_effect = _engine_seeks
2320 with _capture_holding(session, fifo_bytes=2 * _FRAME_BYTES, reader_bytes=_FRAME_BYTES):
2321 item = await session.seek_current(TRACK_A, 120_000)
2322 # nothing rendered while the engine was being moved reached the channel
2323 assert not item._chunks
2324 # and what the pipeline still held at the confirmation is dropped after it
2325 assert session._stale_budget == 3 * _FRAME_BYTES
2326 session._write_if_wanted(b"\x02" * (3 * _FRAME_BYTES))
2327 session._write_if_wanted(b"\x03" * 16)
2328 assert b"".join(item._chunks) == b"\x03" * 16
2329
2330
2331async def test_the_sink_is_suspended_while_a_seek_is_in_flight(tmp_path: Path) -> None:
2332 """No pre-seek audio enters the capture while the engine is being moved."""
2333 session = _make_session(tmp_path)
2334 session._demand_started = True
2335 session._engine_playing = True
2336 session._sink_running = True
2337 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2338 playing.started.set()
2339 playing.observe_position(30_000)
2340 suspended_during_seek = False
2341
2342 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2343 nonlocal suspended_during_seek
2344 suspended_during_seek = not session._sink_running
2345 _current_of(session).observe_position(position_ms)
2346
2347 _client_of(session).seek.side_effect = _engine_seeks
2348 await session.seek_current(TRACK_A, 120_000)
2349 assert suspended_during_seek
2350 _sink_of(session).suspend.assert_awaited()
2351
2352
2353async def test_a_seek_the_engine_never_confirms_fails_the_item(tmp_path: Path) -> None:
2354 """An unconfirmed seek is reported rather than served from the wrong position."""
2355 session = _make_session(tmp_path)
2356 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2357 playing.started.set()
2358 playing.observe_position(30_000)
2359 with (
2360 patch.object(soloist_backend, "_SEEK_CONFIRM_TIMEOUT_S", 0.01),
2361 pytest.raises(AudioError, match="did not confirm"),
2362 ):
2363 await session.seek_current(TRACK_A, 120_000)
2364
2365
2366async def test_a_refused_seek_command_reports_soloist(tmp_path: Path) -> None:
2367 """A rejected seek names the engine, so the caller can fall back."""
2368 session = _make_session(tmp_path)
2369 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2370 playing.started.set()
2371 _client_of(session).seek.side_effect = SoloistError("nope")
2372 with pytest.raises(AudioError, match="would not seek"):
2373 await session.seek_current(TRACK_A, 120_000)
2374
2375
2376async def test_seeking_an_item_the_engine_is_not_on_is_refused(tmp_path: Path) -> None:
2377 """Only the item the session is actually playing can be seeked in place."""
2378 session = _make_session(tmp_path)
2379 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2380 playing.started.set()
2381 with pytest.raises(AudioError, match="is not playing"):
2382 await session.seek_current(TRACK_B, 120_000)
2383
2384
2385async def test_a_seek_of_the_playing_item_is_served_by_the_running_session(
2386 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
2387) -> None:
2388 """The session is seeked where it stands instead of being replaced."""
2389 backend = _make_backend(tmp_path)
2390 backend._server = MagicMock()
2391 backend._binary = Path("/nonexistent/soloist")
2392 session = _make_session(tmp_path)
2393 backend._session = session
2394 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2395 item.started.set()
2396 # its own stream is still attached when the seek re-opens it
2397 item.claim()
2398 item.observe_position(30_000)
2399 stopped = AsyncMock()
2400 monkeypatch.setattr(session, "stop", stopped)
2401
2402 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2403 _current_of(session).observe_position(position_ms)
2404
2405 _client_of(session).seek.side_effect = _engine_seeks
2406 got_session, got_item = await backend._acquire(TRACK_A, 90, "player1")
2407
2408 assert got_session is session
2409 assert got_item is not item
2410 assert got_item.claimed
2411 stopped.assert_not_awaited()
2412 _client_of(session).seek.assert_awaited_once_with(90_000, await_result=True)
2413
2414
2415async def test_a_cancelled_seek_does_not_leave_the_session_wedged(tmp_path: Path) -> None:
2416 """
2417 A superseded seek ends the session instead of holding it claimed for good.
2418
2419 A second seek cancels the stream the first one is being made for, and the
2420 channel it had already claimed would otherwise keep the session in use:
2421 unable to expire, and refusing every later item as busy.
2422 """
2423 session = _make_session(tmp_path)
2424 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2425 playing.started.set()
2426 playing.observe_position(30_000)
2427 seeking = asyncio.create_task(session.seek_current(TRACK_A, 120_000))
2428 # let it get as far as waiting for the engine to confirm
2429 while not _client_of(session).seek.await_count:
2430 await asyncio.sleep(0)
2431 seeking.cancel()
2432 with suppress(asyncio.CancelledError):
2433 await seeking
2434 assert not session.usable
2435 assert not session._seeking
2436
2437
2438async def test_a_seek_cancelled_before_the_channel_is_swapped_keeps_the_session(
2439 tmp_path: Path,
2440) -> None:
2441 """Nothing has been given up yet while the sink is still being held."""
2442 session = _make_session(tmp_path)
2443 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2444 playing.started.set()
2445 held = asyncio.Event()
2446
2447 async def _slow_suspend(**_kwargs: Any) -> None:
2448 held.set()
2449 await asyncio.sleep(60)
2450
2451 with patch.object(session, "_apply_sink_state", _slow_suspend):
2452 seeking = asyncio.create_task(session.seek_current(TRACK_A, 120_000))
2453 await held.wait()
2454 seeking.cancel()
2455 with suppress(asyncio.CancelledError):
2456 await seeking
2457 # the session is untouched and, crucially, not left dropping every chunk
2458 assert session.usable
2459 assert not session._seeking
2460 assert session.current is playing
2461
2462
2463async def test_a_seek_is_refused_once_the_engine_has_moved_on(tmp_path: Path) -> None:
2464 """The item seeked must still be the one the engine is on when the sink settles."""
2465 session = _make_session(tmp_path)
2466 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2467 playing.started.set()
2468 follower = session._items[TRACK_B] = _ItemAudio(TRACK_B, session)
2469
2470 async def _boundary_lands(**_kwargs: Any) -> None:
2471 session._current = follower
2472
2473 with (
2474 patch.object(session, "_apply_sink_state", _boundary_lands),
2475 pytest.raises(AudioError, match="moved on from"),
2476 ):
2477 await session.seek_current(TRACK_A, 120_000)
2478 # the follower the engine actually reached keeps its own channel
2479 assert session.current is follower
2480 assert session._items[TRACK_A] is playing
2481
2482
2483async def test_a_seek_that_fails_part_way_restarts_the_session(
2484 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
2485) -> None:
2486 """A seek the engine refuses after the channel was swapped still gets its audio."""
2487 backend = _make_backend(tmp_path)
2488 backend._server = MagicMock()
2489 backend._binary = Path("/nonexistent/soloist")
2490 session = _make_session(tmp_path)
2491 backend._session = session
2492 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2493 item.started.set()
2494 item.claim()
2495 _client_of(session).seek.side_effect = SoloistError("refused")
2496 stopped = AsyncMock()
2497 monkeypatch.setattr(session, "stop", stopped)
2498 _install_fake_binary_manager(monkeypatch)
2499 monkeypatch.setattr(
2500 soloist_backend._SoloistSession, "start", AsyncMock(side_effect=AudioError("spawn"))
2501 )
2502 with pytest.raises(AudioError, match="spawn"):
2503 await backend._acquire(TRACK_A, 90, "player1")
2504 stopped.assert_awaited_once()
2505
2506
2507async def test_a_seek_refused_because_the_app_took_over_does_not_respawn(
2508 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
2509) -> None:
2510 """A replacement would claim the Connect device back off the Spotify app."""
2511 backend = _make_backend(tmp_path)
2512 backend._server = MagicMock()
2513 backend._binary = Path("/nonexistent/soloist")
2514 session = _make_session(tmp_path)
2515 backend._session = session
2516 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2517 item.started.set()
2518 item.claim()
2519 item.observe_position(30_000)
2520 started = AsyncMock()
2521 monkeypatch.setattr(soloist_backend._SoloistSession, "start", started)
2522
2523 async def _app_takes_over(_position_ms: int, **_kwargs: Any) -> None:
2524 # the user moved playback elsewhere from their Spotify app while the
2525 # seek was in flight; the wait is released by the session ending
2526 session._end_on_app_control(SoloistAppControl.TOOK_OVER)
2527
2528 _client_of(session).seek.side_effect = _app_takes_over
2529 with pytest.raises(SoloistAppControlError):
2530 await backend._acquire(TRACK_A, 120, "player1")
2531 started.assert_not_awaited()
2532
2533
2534async def test_a_seek_is_abandoned_when_the_app_took_over_during_the_suspend(
2535 tmp_path: Path,
2536) -> None:
2537 """Nothing is seeked on a session the Spotify app has already taken over."""
2538 session = _make_session(tmp_path)
2539 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2540 playing.started.set()
2541
2542 async def _app_takes_over(**_kwargs: Any) -> None:
2543 session._end_on_app_control(SoloistAppControl.TOOK_OVER)
2544
2545 with (
2546 patch.object(session, "_apply_sink_state", _app_takes_over),
2547 pytest.raises(SoloistAppControlError),
2548 ):
2549 await session.seek_current(TRACK_A, 120_000)
2550 # the engine was never asked to move
2551 _client_of(session).seek.assert_not_awaited()
2552
2553
2554async def test_a_cancelled_sink_transition_is_re_issued(tmp_path: Path) -> None:
2555 """
2556 A suspend that may or may not have landed is never taken as done.
2557
2558 A sink that did suspend would otherwise still read as running, and the
2559 resume that should follow would be skipped as a no-op: silence for good.
2560 """
2561 session = _make_session(tmp_path)
2562 session._demand_started = True
2563 session._engine_playing = True
2564 session._sink_running = True
2565 session._seeking = True
2566
2567 async def _cancelled_suspend() -> None:
2568 raise asyncio.CancelledError
2569
2570 _sink_of(session).suspend.side_effect = _cancelled_suspend
2571 with suppress(asyncio.CancelledError):
2572 await session._apply_sink_state()
2573 # the engine plays on and the sink is wanted running again
2574 session._seeking = False
2575 await session._apply_sink_state()
2576 _sink_of(session).resume.assert_awaited_once()
2577
2578
2579async def test_a_seek_cancelled_at_the_final_resume_does_not_leave_the_session_usable(
2580 tmp_path: Path,
2581) -> None:
2582 """The channel is claimed by then, so an abandoned seek must still end the session."""
2583 session = _make_session(tmp_path)
2584 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2585 playing.started.set()
2586 playing.observe_position(30_000)
2587 calls = 0
2588 real_apply = session._apply_sink_state
2589
2590 async def _cancel_on_the_way_out(**kwargs: Any) -> None:
2591 nonlocal calls
2592 calls += 1
2593 if calls > 1:
2594 raise asyncio.CancelledError
2595 await real_apply(**kwargs)
2596
2597 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2598 _current_of(session).observe_position(position_ms)
2599
2600 _client_of(session).seek.side_effect = _engine_seeks
2601 with (
2602 patch.object(session, "_apply_sink_state", _cancel_on_the_way_out),
2603 suppress(asyncio.CancelledError),
2604 ):
2605 await session.seek_current(TRACK_A, 120_000)
2606 assert not session.usable
2607
2608
2609async def test_a_cold_seek_does_not_read_a_failed_session_as_landed(tmp_path: Path) -> None:
2610 """The wake-up a fatal failure gives every channel is not a confirmed seek."""
2611 session = _make_session(tmp_path)
2612 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2613
2614 async def _engine_dies(_position_ms: int, **_kwargs: Any) -> None:
2615 session._fail("the session exited")
2616
2617 _client_of(session).seek.side_effect = _engine_dies
2618 with pytest.raises(AudioError, match="the session exited"):
2619 await session._cold_seek(_client_of(session), item, 60_000)
2620
2621
2622async def test_seeking_the_playing_item_back_to_its_start_keeps_the_session(
2623 tmp_path: Path, monkeypatch: pytest.MonkeyPatch
2624) -> None:
2625 """
2626 A seek to zero is a seek: reachable once an earlier one moved the buffer.
2627
2628 The buffer only hands a position to the provider when it cannot serve it
2629 itself, so seeking back before an earlier seek's target arrives here with a
2630 target of zero.
2631 """
2632 backend = _make_backend(tmp_path)
2633 backend._server = MagicMock()
2634 backend._binary = Path("/nonexistent/soloist")
2635 session = _make_session(tmp_path)
2636 backend._session = session
2637 item = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2638 item.started.set()
2639 item.claim()
2640 item.observe_position(200_000)
2641 stopped = AsyncMock()
2642 monkeypatch.setattr(session, "stop", stopped)
2643
2644 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2645 current = _current_of(session)
2646 current.observe_position(position_ms)
2647 current.observe_position(position_ms + 2)
2648
2649 _client_of(session).seek.side_effect = _engine_seeks
2650 got_session, got_item = await backend._acquire(TRACK_A, 0, "player1")
2651
2652 assert got_session is session
2653 assert got_item is not item
2654 stopped.assert_not_awaited()
2655 _client_of(session).seek.assert_awaited_once_with(0, await_result=True)
2656 assert got_item.started_at_ms == 2
2657
2658
2659async def test_a_short_forward_seek_still_confirms(tmp_path: Path) -> None:
2660 """
2661 A seek only a little past where the engine is must not wait itself out.
2662
2663 Reachable because the engine runs ahead of what has been delivered - up to
2664 the retained cushion - so a target the buffer will not serve can still be
2665 inside the tolerance window of the engine's own position. Demanding the
2666 engine drop below that mark would never be satisfied by a seek forward.
2667 """
2668 session = _make_session(tmp_path)
2669 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2670 playing.started.set()
2671 playing.duration_ms = 260_000
2672 # the engine is at 49s while only ~30s has been delivered
2673 playing.observe_position(49_000)
2674
2675 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2676 _current_of(session).observe_position(position_ms)
2677
2678 _client_of(session).seek.side_effect = _engine_seeks
2679 item = await session.seek_current(TRACK_A, 50_500)
2680 assert item.seek_confirmed.is_set()
2681 assert item.started_at_ms == 50_500
2682
2683
2684async def test_a_seek_does_not_arm_the_last_items_drain(tmp_path: Path) -> None:
2685 """
2686 The channel opened for a seek has no position yet, which is not an ended item.
2687
2688 Holding the sink can make the engine report a state that is not playing, and
2689 the run's last item would then be drained out from under the seek.
2690 """
2691 session = _make_session(tmp_path)
2692 session._demand_started = True
2693 playing = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2694 playing.started.set()
2695 playing.duration_ms = 260_000
2696 playing.observe_position(30_000)
2697 armed: list[bool] = []
2698
2699 async def _engine_seeks(position_ms: int, **_kwargs: Any) -> None:
2700 # nothing to judge the fresh channel by yet, and no follower queued
2701 await session._handle_playback_state(
2702 SoloistPlaybackState(status="buffering", item=None, position=None)
2703 )
2704 armed.append(_current_of(session).draining)
2705 _current_of(session).observe_position(position_ms)
2706
2707 _client_of(session).seek.side_effect = _engine_seeks
2708 item = await session.seek_current(TRACK_A, 120_000)
2709 assert armed == [False]
2710 assert not item.draining
2711
2712
2713async def test_the_item_a_finished_run_stopped_on_is_not_seeked_in_place(
2714 tmp_path: Path,
2715) -> None:
2716 """
2717 A matching uri is not proof the engine is still playing it.
2718
2719 The channel stays current through the idle grace after the run ended, and a
2720 seek would wait out its confirmation on an engine that has stopped.
2721 """
2722 session = _make_session(tmp_path)
2723 ended = session._items[TRACK_A] = session._current = _ItemAudio(TRACK_A, session)
2724 ended.started.set()
2725 ended.close()
2726 with pytest.raises(AudioError, match="is not playing"):
2727 await session.seek_current(TRACK_A, 0)
2728 _client_of(session).seek.assert_not_awaited()
2729