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