/
/
/
1"""Tests that starting an item releases the source slots its own queue no longer needs."""
2
3from __future__ import annotations
4
5from typing import cast
6from unittest.mock import AsyncMock, MagicMock
7
8from music_assistant_models.enums import ContentType, MediaType, StreamType
9from music_assistant_models.media_items import AudioFormat
10from music_assistant_models.queue_item import QueueItem
11from music_assistant_models.streamdetails import StreamDetails
12
13from music_assistant.controllers.player_queues import PlayerQueuesController
14from music_assistant.controllers.player_queues.state import PlayerQueueData
15from music_assistant.controllers.streams.audio_buffer import AudioBuffer
16from music_assistant.models.music_provider import MusicProvider
17
18QUEUE_ID = "q1"
19LIMITED = "limited--1"
20UNLIMITED = "local--1"
21
22
23def _item(queue_id: str, item_id: str, provider: str | None, buffering: bool = True) -> QueueItem:
24 """
25 Build a queue item, optionally with a buffer attached to its stream details.
26
27 :param queue_id: The queue the item belongs to.
28 :param item_id: The queue item id.
29 :param provider: Provider instance owning the source, or None for an item without details.
30 :param buffering: Whether the attached buffer is still filling from its source.
31 """
32 queue_item = QueueItem(queue_id=queue_id, queue_item_id=item_id, name=item_id, duration=180)
33 if provider is None:
34 return queue_item
35 audio_buffer = MagicMock(spec=AudioBuffer)
36 audio_buffer.is_buffering = buffering
37 audio_buffer.clear = AsyncMock()
38 queue_item.streamdetails = StreamDetails(
39 provider=provider,
40 item_id=item_id,
41 audio_format=AudioFormat(content_type=ContentType.MP3),
42 media_type=MediaType.TRACK,
43 stream_type=StreamType.HTTP,
44 path=f"http://test.invalid/{item_id}.mp3",
45 )
46 queue_item.streamdetails.buffer = audio_buffer
47 return queue_item
48
49
50def _controller(*queues: tuple[str, list[QueueItem]]) -> PlayerQueuesController:
51 """Build a bare controller holding the given queues and their items."""
52 ctrl = PlayerQueuesController.__new__(PlayerQueuesController)
53 ctrl.logger = MagicMock()
54 ctrl._queue_data = {
55 queue_id: PlayerQueueData(queue=MagicMock(), items=items, session_id=queue_id)
56 for queue_id, items in queues
57 }
58 limited = MagicMock(spec=MusicProvider)
59 limited.name = "Limited"
60 limited.max_concurrent_streams = 2
61 # pinned rather than left to MagicMock truthiness: this decides whether a prewarm survives
62 limited.has_available_stream_slot = True
63 unlimited = MagicMock(spec=MusicProvider)
64 unlimited.name = "Local"
65 unlimited.max_concurrent_streams = None
66 unlimited.has_available_stream_slot = True
67 providers = {LIMITED: limited, UNLIMITED: unlimited}
68 ctrl.mass = MagicMock()
69 ctrl.mass.get_provider.side_effect = lambda instance, **_kwargs: providers.get(instance)
70 return ctrl
71
72
73def _limited_double(ctrl: PlayerQueuesController) -> MagicMock:
74 """Return the slot-limited provider double the controller resolves LIMITED to."""
75 return cast("MagicMock", ctrl.mass.get_provider(LIMITED))
76
77
78async def test_starting_an_item_aborts_the_other_filling_sources_of_its_queue() -> None:
79 """A still-filling source of a preceding item in the same queue hands its slot over."""
80 target = _item(QUEUE_ID, "target", LIMITED)
81 filling = _item(QUEUE_ID, "filling", LIMITED)
82 ctrl = _controller((QUEUE_ID, [filling, target]))
83
84 await ctrl._abort_superseded_source_buffers(target)
85
86 # the aborted buffer stays attached so the flow stream can see the abort
87 assert filling.streamdetails is not None
88 assert filling.streamdetails.buffer is not None
89 filling.streamdetails.buffer.clear.assert_awaited_once()
90 assert target.streamdetails is not None
91 assert target.streamdetails.buffer is not None
92 assert target.streamdetails.buffer.clear.await_count == 0
93
94
95async def test_the_started_items_successor_keeps_its_prewarm() -> None:
96 """A resume or seek must not cost the crossfade prewarm of the upcoming track."""
97 target = _item(QUEUE_ID, "target", LIMITED)
98 upcoming = _item(QUEUE_ID, "upcoming", LIMITED)
99 stale = _item(QUEUE_ID, "stale", LIMITED)
100 ctrl = _controller((QUEUE_ID, [target, upcoming, stale]))
101
102 await ctrl._abort_superseded_source_buffers(target)
103
104 assert upcoming.streamdetails is not None
105 assert upcoming.streamdetails.buffer is not None
106 assert upcoming.streamdetails.buffer.clear.await_count == 0
107 assert stale.streamdetails is not None
108 assert stale.streamdetails.buffer is not None
109 stale.streamdetails.buffer.clear.assert_awaited_once()
110
111
112async def test_a_saturated_provider_takes_back_the_successors_prewarm() -> None:
113 """A prewarm may not sit on the last slot the item being started needs."""
114 target = _item(QUEUE_ID, "target", LIMITED)
115 upcoming = _item(QUEUE_ID, "upcoming", LIMITED)
116 ctrl = _controller((QUEUE_ID, [target, upcoming]))
117 _limited_double(ctrl).has_available_stream_slot = False
118
119 await ctrl._abort_superseded_source_buffers(target)
120
121 assert upcoming.streamdetails is not None
122 assert upcoming.streamdetails.buffer is not None
123 upcoming.streamdetails.buffer.clear.assert_awaited_once()
124 assert target.streamdetails is not None
125 assert target.streamdetails.buffer is not None
126 assert target.streamdetails.buffer.clear.await_count == 0
127
128
129async def test_freeing_a_slot_first_lets_the_successor_keep_its_prewarm() -> None:
130 """The prewarm is only given up when aborting the stale sources did not free a slot."""
131 stale = _item(QUEUE_ID, "stale", LIMITED)
132 target = _item(QUEUE_ID, "target", LIMITED)
133 upcoming = _item(QUEUE_ID, "upcoming", LIMITED)
134 ctrl = _controller((QUEUE_ID, [stale, target, upcoming]))
135 provider = _limited_double(ctrl)
136 provider.has_available_stream_slot = False
137
138 async def _release_slot() -> None:
139 provider.has_available_stream_slot = True
140
141 assert stale.streamdetails is not None
142 assert stale.streamdetails.buffer is not None
143 stale.streamdetails.buffer.clear.side_effect = _release_slot
144
145 await ctrl._abort_superseded_source_buffers(target)
146
147 # the successor is only re-checked after the stale aborts, so it survives
148 stale.streamdetails.buffer.clear.assert_awaited_once()
149 assert upcoming.streamdetails is not None
150 assert upcoming.streamdetails.buffer is not None
151 assert upcoming.streamdetails.buffer.clear.await_count == 0
152
153
154async def test_completed_and_unlimited_sources_are_left_alone() -> None:
155 """Only a source that still holds a capped provider slot is worth aborting."""
156 target = _item(QUEUE_ID, "target", LIMITED)
157 completed = _item(QUEUE_ID, "completed", LIMITED, buffering=False)
158 unlimited = _item(QUEUE_ID, "unlimited", UNLIMITED)
159 ctrl = _controller((QUEUE_ID, [completed, unlimited, target]))
160
161 await ctrl._abort_superseded_source_buffers(target)
162
163 for item in (completed, unlimited):
164 assert item.streamdetails is not None
165 assert item.streamdetails.buffer is not None
166 assert item.streamdetails.buffer.clear.await_count == 0
167
168
169async def test_other_queues_keep_their_sources() -> None:
170 """Playback on one queue never takes a source away from another queue."""
171 target = _item(QUEUE_ID, "target", LIMITED)
172 other = _item("q2", "other", LIMITED)
173 ctrl = _controller((QUEUE_ID, [target]), ("q2", [other]))
174
175 await ctrl._abort_superseded_source_buffers(target)
176
177 assert other.streamdetails is not None
178 assert other.streamdetails.buffer is not None
179 assert other.streamdetails.buffer.clear.await_count == 0
180