/
/
1"""
2Stream feeding for the Player Queues controller.
3
4Handles handing the next queue item to the player and preparing its audio: enqueuing the upcoming
5item on the player, preloading its stream details, warming the next track's AudioBuffer ahead of
6playback, and cleaning up stale buffers. Owns no per-queue state; it is mixed into the controller
7and reads/mutates the controller's `PlayerQueueData` records.
8"""
9
10from __future__ import annotations
11
12import asyncio
13from typing import TYPE_CHECKING
14
15from music_assistant_models.enums import (
16 MediaType,
17 PlaybackState,
18)
19from music_assistant_models.errors import (
20 AudioError,
21 MediaNotFoundError,
22 QueueEmpty,
23)
24
25from music_assistant.constants import (
26 VERBOSE_LOG_LEVEL,
27)
28from music_assistant.controllers.player_queues.base import _PlayerQueuesBase
29from music_assistant.controllers.streams.constants import STREAM_SLOT_WAIT_TIMEOUT
30
31if TYPE_CHECKING:
32 from music_assistant_models.queue_item import QueueItem
33
34
35class StreamFeederMixin(_PlayerQueuesBase):
36 """Feed the player's stream: enqueue the next item, preload/prepare its audio, clean up."""
37
38 def prepare_next_audio_buffer(self, queue_id: str) -> None:
39 """
40 Prepare the AudioBuffer for the next track in the queue.
41
42 Called ~30-60 seconds before the current track ends to ensure
43 the buffer is warm when the next track starts playing.
44 """
45 queue = self.get(queue_id)
46 if not queue or not queue.next_item:
47 return
48 next_item = queue.next_item
49 # AudioSource items are realtime/live and bypass the AudioBuffer
50 if next_item.media_type == MediaType.AUDIO_SOURCE:
51 return
52 # guard against race condition where queue.next_item still points to the
53 # currently playing track because the player state hasn't been updated yet
54 if queue.current_item and next_item.queue_item_id == queue.current_item.queue_item_id:
55 return
56 # check if buffer already exists and is valid
57 if (
58 next_item.streamdetails
59 and next_item.streamdetails.buffer
60 and next_item.streamdetails.buffer.is_valid()
61 ):
62 return
63
64 async def _do_prepare() -> None:
65 try:
66 # fetch streamdetails if not yet available
67 if not next_item.streamdetails:
68 next_item.streamdetails = await self.mass.streams.audio.get_stream_details(
69 queue_item=next_item
70 )
71 self.logger.debug(
72 "Preparing audio buffer for next track %s on queue %s",
73 next_item.name,
74 queue.display_name,
75 )
76 await self.mass.streams.audio.get_audio_buffer(
77 next_item,
78 reason="prepare_next",
79 capacity_wait_timeout=STREAM_SLOT_WAIT_TIMEOUT,
80 )
81 except (AudioError, MediaNotFoundError) as err:
82 self.logger.debug("Failed to prepare next audio buffer: %s", err)
83 except asyncio.CancelledError:
84 # a replacement prepare aborted this one: release the half-filled source
85 # so its slot is not pinned until the inactivity sweep
86 if (sd := next_item.streamdetails) and (buf := sd.buffer) and buf.is_buffering:
87 await asyncio.shield(buf.clear())
88 raise
89
90 self.mass.create_task(
91 _do_prepare,
92 task_id=f"prepare_next_audio_buffer_{queue_id}",
93 abort_existing=True,
94 )
95
96 def _enqueue_next_item(self, queue_id: str, next_item: QueueItem | None) -> None:
97 """Enqueue the next item on the player."""
98 if not next_item:
99 # no next item, nothing to do...
100 return
101
102 queue_data = self._queue_data[queue_id]
103 queue = queue_data.queue
104 session_id = queue_data.session_id
105 if queue.flow_mode:
106 # ignore this for flow mode
107 return
108
109 async def _enqueue_next_item_on_player(next_item: QueueItem) -> None:
110 # Player state updates can lag behind queue loading, so wait before validating.
111 async with self.mass.players.wait_for_player_update(
112 queue_id,
113 attribute_name="playback_state",
114 attribute_value=PlaybackState.PLAYING,
115 ):
116 pass
117
118 player = self.mass.players.get_player(queue_id)
119 if (
120 player is None
121 or player.state.playback_state != PlaybackState.PLAYING
122 or player.state.active_source not in (queue.queue_id, None)
123 or queue_data.session_id != session_id
124 or queue.flow_mode
125 ):
126 # nothing re-attempts this handover, so a skip here means the player runs out
127 # of audio when the current track ends - leave a trace of why it was skipped
128 self.logger.debug(
129 "Not enqueuing next track %s on queue %s "
130 "(state: %s, source: %s, same session: %s, flow mode: %s)",
131 next_item.name,
132 queue.display_name,
133 player.state.playback_state if player else "player unavailable",
134 player.state.active_source if player else None,
135 queue_data.session_id == session_id,
136 queue.flow_mode,
137 )
138 return
139
140 current_item = queue.current_item
141 if current_item is None:
142 return
143 current_next = self.get_next_item(queue_id, current_item.queue_item_id)
144 if current_next is None or current_next.queue_item_id != next_item.queue_item_id:
145 return
146
147 await self.mass.players.enqueue_next_media(
148 player_id=queue_id,
149 media=await self.player_media_from_queue_item(next_item),
150 )
151 if queue_data.next_item_id_enqueued != next_item.queue_item_id:
152 queue_data.next_item_id_enqueued = next_item.queue_item_id
153 self.logger.debug(
154 "Enqueued next track %s on queue %s",
155 next_item.name,
156 self._queue_data[queue_id].queue.display_name,
157 )
158
159 task_id = f"enqueue_next_item_{queue_id}"
160 self.mass.call_later(1, _enqueue_next_item_on_player, next_item, task_id=task_id)
161
162 def _preload_next_item(self, queue_id: str, item_id_in_buffer: str) -> None:
163 """
164 Preload the streamdetails for the next item in the queue/buffer.
165
166 This basically ensures the item is playable and fetches the stream details.
167 If an error occurs, the item will be skipped and the next item will be loaded.
168 """
169 queue = self._queue_data[queue_id].queue
170
171 async def _preload_streamdetails(item_id_in_buffer: str) -> None:
172 try:
173 # wait for the item that was loaded in the buffer is the actually playing item
174 # this prevents a race condition when we preload the next item too soon
175 # while the player is actually preloading the previously enqueued item.
176 current_item = queue.current_item
177 if current_item is None:
178 return # guard
179 retries = max(120, int(current_item.duration or 0) + 10)
180 for _ in range(retries):
181 # the queue can drain to empty while we sleep (e.g. all remaining
182 # items skipped as unplayable); stop waiting once it has no current item
183 current_item = queue.current_item
184 if current_item is None:
185 return
186 if current_item.queue_item_id == item_id_in_buffer:
187 break
188 await asyncio.sleep(1)
189 if next_item := await self.load_next_queue_item(queue_id, item_id_in_buffer):
190 self.logger.debug(
191 "Preloaded next item %s for queue %s",
192 next_item.name,
193 queue.display_name,
194 )
195 # enqueue the next item on the player
196 self._enqueue_next_item(queue_id, next_item)
197
198 except QueueEmpty:
199 return
200
201 if not (current_item := self.get_item(queue_id, item_id_in_buffer)):
202 # this should not happen, but guard anyways
203 return
204 if current_item.media_type == MediaType.RADIO or not current_item.duration:
205 # radio items or no duration, nothing to do
206 return
207
208 task_id = f"preload_next_item_{queue_id}"
209 self.mass.create_task(
210 _preload_streamdetails,
211 item_id_in_buffer,
212 task_id=task_id,
213 abort_existing=True,
214 )
215
216 async def _cleanup_stale_queue_buffers(self, queue_id: str, current_index: int) -> None:
217 """
218 Clean up audio buffers for queue items that are no longer needed.
219
220 This clears buffers for items at index <= current_index - 2, keeping only:
221 - The previous track (current_index - 1)
222 - The current track (current_index)
223 - The next track (current_index + 1, handled by preloading)
224
225 :param queue_id: The queue ID to clean up buffers for.
226 :param current_index: The current playing index in the queue.
227 """
228 if current_index < 2:
229 return # Nothing to clean up yet
230
231 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
232 cleanup_threshold = current_index - 2
233 buffers_cleared = 0
234
235 for idx, item in enumerate(queue_items):
236 if idx > cleanup_threshold:
237 break # No need to check further
238 if item.streamdetails and item.streamdetails.buffer:
239 self.logger.log(
240 VERBOSE_LOG_LEVEL,
241 "Clearing stale audio buffer for queue item %s (index %d) in queue %s",
242 item.name,
243 idx,
244 queue_id,
245 )
246 await item.streamdetails.buffer.clear()
247 item.streamdetails.buffer = None
248 buffers_cleared += 1
249
250 if buffers_cleared > 0:
251 self.logger.debug(
252 "Cleared %d stale audio buffer(s) for queue %s (items before index %d)",
253 buffers_cleared,
254 queue_id,
255 cleanup_threshold + 1,
256 )
257
258 async def _cleanup_queue_audio_data(self, queue_id: str) -> None:
259 """
260 Clean up all audio-related data for a queue when it is stopped or cleared.
261
262 This clears:
263 - All audio buffers attached to queue item streamdetails
264 - Any pending crossfade data for the queue
265
266 :param queue_id: The queue ID to clean up.
267 """
268 self.mass.streams.audio.clear_crossfade_data(queue_id)
269
270 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
271 buffers_cleared = 0
272
273 for item in queue_items:
274 if item.streamdetails and item.streamdetails.buffer:
275 await item.streamdetails.buffer.clear()
276 item.streamdetails.buffer = None
277 buffers_cleared += 1
278
279 if buffers_cleared > 0:
280 self.logger.debug(
281 "Cleared %d audio buffer(s) for stopped/cleared queue %s",
282 buffers_cleared,
283 queue_id,
284 )
285