/
/
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_prepare_next_uses_the_speculative_capacity_budget() -> None:
140 """Warming the next track never waits longer for capacity than a speculative attempt may."""
141 controller, next_item, mass = _controller_with_next_item()
142 mass.streams.audio.get_audio_buffer = AsyncMock()
143
144 controller.prepare_next_audio_buffer("queue-1")
145 await mass.create_task.call_args.args[0]()
146
147 mass.streams.audio.get_audio_buffer.assert_awaited_once_with(
148 next_item,
149 reason="prepare_next",
150 capacity_wait_timeout=STREAM_SLOT_WAIT_TIMEOUT,
151 allow_provider_match=False,
152 )
153 assert mass.create_task.call_args.kwargs == {
154 "task_id": "prepare_next_audio_buffer_queue-1",
155 "abort_existing": True,
156 }
157
158
159@pytest.mark.parametrize(
160 ("is_buffering", "expect_cleared"),
161 [(True, True), (False, False)],
162 ids=["still_filling", "completed"],
163)
164async def test_an_aborted_prepare_releases_its_half_filled_source(
165 is_buffering: bool, expect_cleared: bool
166) -> None:
167 """Aborting a prewarm must free its slot instead of pinning it until the inactivity sweep."""
168 controller, next_item, mass = _controller_with_next_item()
169 buffer = MagicMock()
170 buffer.is_buffering = is_buffering
171 buffer.clear = AsyncMock()
172 started = asyncio.Event()
173
174 async def _hang(*_args: object, **_kwargs: object) -> None:
175 # the producer is already running and owns a source slot at this point
176 next_item.streamdetails.buffer = buffer
177 started.set()
178 await asyncio.Event().wait()
179
180 mass.streams.audio.get_audio_buffer = _hang
181
182 controller.prepare_next_audio_buffer("queue-1")
183 prepare_task = asyncio.create_task(mass.create_task.call_args.args[0]())
184 await started.wait()
185 prepare_task.cancel()
186 with pytest.raises(asyncio.CancelledError):
187 await prepare_task
188
189 assert buffer.clear.await_count == (1 if expect_cleared else 0)
190
191
192async def test_prepare_next_gives_up_softly_on_a_capacity_failure() -> None:
193 """A speculative source-capacity miss leaves the next item playable."""
194 controller, next_item, mass = _controller_with_next_item()
195 provider = MagicMock(spec=MusicProvider)
196 provider.max_concurrent_streams = 1
197 provider.name = "Limited"
198 provider.instance_id = "limited--1"
199 mass.streams.audio.get_audio_buffer = AsyncMock(
200 side_effect=ProviderStreamLimitError(provider, STREAM_SLOT_WAIT_TIMEOUT)
201 )
202
203 controller.prepare_next_audio_buffer("queue-1")
204 await mass.create_task.call_args.args[0]()
205
206 assert next_item.available
207