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