/
/
/
1"""Regression tests for delayed next-track enqueueing in the player queue stream feeder."""
2
3from __future__ import annotations
4
5import asyncio
6from collections.abc import AsyncIterator
7from contextlib import asynccontextmanager
8from types import SimpleNamespace
9from typing import Any, cast
10from unittest.mock import AsyncMock, MagicMock
11
12import pytest
13from music_assistant_models.enums import MediaType, PlaybackState
14from music_assistant_models.player_queue import PlayerQueue
15from music_assistant_models.queue_item import QueueItem
16
17from music_assistant.controllers.player_queues import PlayerQueuesController
18from music_assistant.controllers.player_queues.state import PlayerQueueData
19from music_assistant.controllers.streams.constants import STREAM_SLOT_WAIT_TIMEOUT
20from music_assistant.models.music_provider import MusicProvider, ProviderStreamLimitError
21
22
23@pytest.mark.parametrize(
24 "index_in_buffer",
25 [0, 1],
26 ids=["aligned-buffer-index", "dynamic-queue-reindexed"],
27)
28async def test_enqueue_next_item_waits_for_playing_player_update(index_in_buffer: int) -> None:
29 """Enqueue the expected next item after an update, even if its buffered index is stale."""
30 controller = PlayerQueuesController.__new__(PlayerQueuesController)
31 controller.logger = MagicMock()
32
33 wait_entered = asyncio.Event()
34 release_wait = asyncio.Event()
35
36 @asynccontextmanager
37 async def wait_for_player_update(*_args: object, **_kwargs: object) -> AsyncIterator[None]:
38 wait_entered.set()
39 await release_wait.wait()
40 yield
41
42 player_state = SimpleNamespace(
43 playback_state=PlaybackState.IDLE,
44 active_source="q1",
45 )
46 player = SimpleNamespace(state=player_state)
47
48 mass = MagicMock()
49 mass.players = MagicMock()
50 mass.players.wait_for_player_update = MagicMock(side_effect=wait_for_player_update)
51 mass.players.get_player = MagicMock(return_value=player)
52 mass.players.enqueue_next_media = AsyncMock()
53 controller.mass = mass
54
55 current_item = _make_queue_item("q1", "nerin")
56 next_item = _make_queue_item("q1", "another-love")
57 future_item = _make_queue_item("q1", "future-track")
58 queue_items = [current_item, next_item, future_item]
59 queue = PlayerQueue(
60 queue_id="q1",
61 active=True,
62 display_name="Q1",
63 available=True,
64 items=len(queue_items),
65 state=PlaybackState.IDLE,
66 current_index=0,
67 index_in_buffer=index_in_buffer,
68 current_item=current_item,
69 )
70 controller._queue_data = {
71 "q1": PlayerQueueData(
72 queue=queue,
73 items=queue_items,
74 session_id="session-1",
75 )
76 }
77
78 controller._enqueue_next_item("q1", next_item)
79 enqueue_callback = mass.call_later.call_args.args[1]
80 enqueue_task = asyncio.create_task(enqueue_callback(next_item))
81 await asyncio.sleep(0)
82
83 assert wait_entered.is_set()
84 mass.players.enqueue_next_media.assert_not_awaited()
85
86 player_state.playback_state = PlaybackState.PLAYING
87 release_wait.set()
88 await enqueue_task
89
90 mass.players.wait_for_player_update.assert_called_once_with(
91 "q1",
92 attribute_name="playback_state",
93 attribute_value=PlaybackState.PLAYING,
94 )
95 mass.players.enqueue_next_media.assert_awaited_once()
96 assert mass.players.enqueue_next_media.await_args.kwargs["player_id"] == "q1"
97 assert (
98 mass.players.enqueue_next_media.await_args.kwargs["media"].queue_item_id
99 == next_item.queue_item_id
100 )
101 assert controller._queue_data["q1"].next_item_id_enqueued == next_item.queue_item_id
102
103
104def _make_queue_item(queue_id: str, item_id: str) -> QueueItem:
105 """Build a minimal playable queue item."""
106 return QueueItem(
107 queue_id=queue_id,
108 queue_item_id=item_id,
109 name=item_id,
110 duration=60,
111 )
112
113
114def _controller_with_next_item() -> tuple[PlayerQueuesController, SimpleNamespace, MagicMock]:
115 """Build a bare controller whose queue has an unprepared next item."""
116 controller = PlayerQueuesController.__new__(PlayerQueuesController)
117 controller.logger = MagicMock()
118 next_item = SimpleNamespace(
119 queue_item_id="next",
120 media_type=MediaType.TRACK,
121 streamdetails=SimpleNamespace(buffer=None),
122 name="Next",
123 available=True,
124 )
125 queue = SimpleNamespace(
126 current_item=SimpleNamespace(queue_item_id="current"),
127 next_item=next_item,
128 display_name="Queue",
129 )
130 controller.get = MagicMock(return_value=queue) # type: ignore[method-assign]
131 controller._queue_data = {
132 "queue-1": cast("Any", SimpleNamespace(queue=queue, session_id="session-1"))
133 }
134 mass = MagicMock()
135 controller.mass = mass
136 return controller, next_item, mass
137
138
139async def test_reusing_a_warm_buffer_claims_it_for_the_current_session() -> None:
140 """
141 A prewarm that is already warm still becomes this session's audio.
142
143 Without the claim the buffer keeps the session that filled it, and that session's stop
144 releases audio the current one is relying on.
145 """
146 controller, next_item, mass = _controller_with_next_item()
147 warm = MagicMock()
148 warm.is_valid.return_value = True
149 next_item.streamdetails = SimpleNamespace(buffer=warm, queue_session_id="session-0")
150
151 controller.prepare_next_audio_buffer("queue-1")
152
153 assert next_item.streamdetails.queue_session_id == "session-1"
154 mass.create_task.assert_not_called()
155
156
157async def test_prepare_next_uses_the_speculative_capacity_budget() -> None:
158 """Warming the next track never waits longer for capacity than a speculative attempt may."""
159 controller, next_item, mass = _controller_with_next_item()
160 mass.streams.audio.get_audio_buffer = AsyncMock()
161
162 controller.prepare_next_audio_buffer("queue-1")
163 await mass.create_task.call_args.args[0]()
164
165 mass.streams.audio.get_audio_buffer.assert_awaited_once_with(
166 next_item,
167 reason="prepare_next",
168 capacity_wait_timeout=STREAM_SLOT_WAIT_TIMEOUT,
169 allow_provider_match=False,
170 )
171 assert mass.create_task.call_args.kwargs == {
172 "task_id": "prepare_next_audio_buffer_queue-1",
173 "abort_existing": True,
174 }
175
176
177@pytest.mark.parametrize(
178 ("is_buffering", "expect_cleared"),
179 [(True, True), (False, False)],
180 ids=["still_filling", "completed"],
181)
182async def test_an_aborted_prepare_releases_its_half_filled_source(
183 is_buffering: bool, expect_cleared: bool
184) -> None:
185 """Aborting a prewarm must free its slot instead of pinning it until the inactivity sweep."""
186 controller, next_item, mass = _controller_with_next_item()
187 buffer = MagicMock()
188 buffer.is_buffering = is_buffering
189 buffer.clear = AsyncMock()
190 started = asyncio.Event()
191
192 async def _hang(*_args: object, **_kwargs: object) -> None:
193 # the producer is already running and owns a source slot at this point
194 next_item.streamdetails.buffer = buffer
195 started.set()
196 await asyncio.Event().wait()
197
198 mass.streams.audio.get_audio_buffer = _hang
199
200 controller.prepare_next_audio_buffer("queue-1")
201 prepare_task = asyncio.create_task(mass.create_task.call_args.args[0]())
202 await started.wait()
203 prepare_task.cancel()
204 with pytest.raises(asyncio.CancelledError):
205 await prepare_task
206
207 assert buffer.clear.await_count == (1 if expect_cleared else 0)
208
209
210async def test_prepare_next_gives_up_softly_on_a_capacity_failure() -> None:
211 """A speculative source-capacity miss leaves the next item playable."""
212 controller, next_item, mass = _controller_with_next_item()
213 provider = MagicMock(spec=MusicProvider)
214 provider.max_concurrent_streams = 1
215 provider.name = "Limited"
216 provider.instance_id = "limited--1"
217 mass.streams.audio.get_audio_buffer = AsyncMock(
218 side_effect=ProviderStreamLimitError(provider, STREAM_SLOT_WAIT_TIMEOUT)
219 )
220
221 controller.prepare_next_audio_buffer("queue-1")
222 await mass.create_task.call_args.args[0]()
223
224 assert next_item.available
225