/
/
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.audio_buffer import AudioBuffer
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 AudioBuffer.get_buffer(
77 self.mass,
78 next_item.streamdetails,
79 reason="prepare_next",
80 wait_ready=True,
81 )
82 except (AudioError, MediaNotFoundError) as err:
83 self.logger.debug("Failed to prepare next audio buffer: %s", err)
84
85 self.mass.create_task(_do_prepare)
86
87 def _enqueue_next_item(self, queue_id: str, next_item: QueueItem | None) -> None:
88 """Enqueue the next item on the player."""
89 if not next_item:
90 # no next item, nothing to do...
91 return
92
93 queue_data = self._queue_data[queue_id]
94 queue = queue_data.queue
95 session_id = queue_data.session_id
96 if queue.flow_mode:
97 # ignore this for flow mode
98 return
99
100 async def _enqueue_next_item_on_player(next_item: QueueItem) -> None:
101 # Player state updates can lag behind queue loading, so wait before validating.
102 async with self.mass.players.wait_for_player_update(
103 queue_id,
104 attribute_name="playback_state",
105 attribute_value=PlaybackState.PLAYING,
106 ):
107 pass
108
109 player = self.mass.players.get_player(queue_id)
110 if (
111 player is None
112 or player.state.playback_state != PlaybackState.PLAYING
113 or player.state.active_source not in (queue.queue_id, None)
114 or queue_data.session_id != session_id
115 or queue.flow_mode
116 ):
117 # nothing re-attempts this handover, so a skip here means the player runs out
118 # of audio when the current track ends - leave a trace of why it was skipped
119 self.logger.debug(
120 "Not enqueuing next track %s on queue %s "
121 "(state: %s, source: %s, same session: %s, flow mode: %s)",
122 next_item.name,
123 queue.display_name,
124 player.state.playback_state if player else "player unavailable",
125 player.state.active_source if player else None,
126 queue_data.session_id == session_id,
127 queue.flow_mode,
128 )
129 return
130
131 current_item = queue.current_item
132 if current_item is None:
133 return
134 current_next = self.get_next_item(queue_id, current_item.queue_item_id)
135 if current_next is None or current_next.queue_item_id != next_item.queue_item_id:
136 return
137
138 await self.mass.players.enqueue_next_media(
139 player_id=queue_id,
140 media=await self.player_media_from_queue_item(next_item),
141 )
142 if queue_data.next_item_id_enqueued != next_item.queue_item_id:
143 queue_data.next_item_id_enqueued = next_item.queue_item_id
144 self.logger.debug(
145 "Enqueued next track %s on queue %s",
146 next_item.name,
147 self._queue_data[queue_id].queue.display_name,
148 )
149
150 task_id = f"enqueue_next_item_{queue_id}"
151 self.mass.call_later(1, _enqueue_next_item_on_player, next_item, task_id=task_id)
152
153 def _preload_next_item(self, queue_id: str, item_id_in_buffer: str) -> None:
154 """
155 Preload the streamdetails for the next item in the queue/buffer.
156
157 This basically ensures the item is playable and fetches the stream details.
158 If an error occurs, the item will be skipped and the next item will be loaded.
159 """
160 queue = self._queue_data[queue_id].queue
161
162 async def _preload_streamdetails(item_id_in_buffer: str) -> None:
163 try:
164 # wait for the item that was loaded in the buffer is the actually playing item
165 # this prevents a race condition when we preload the next item too soon
166 # while the player is actually preloading the previously enqueued item.
167 current_item = queue.current_item
168 if current_item is None:
169 return # guard
170 retries = max(120, int(current_item.duration or 0) + 10)
171 for _ in range(retries):
172 # the queue can drain to empty while we sleep (e.g. all remaining
173 # items skipped as unplayable); stop waiting once it has no current item
174 current_item = queue.current_item
175 if current_item is None:
176 return
177 if current_item.queue_item_id == item_id_in_buffer:
178 break
179 await asyncio.sleep(1)
180 if next_item := await self.load_next_queue_item(queue_id, item_id_in_buffer):
181 self.logger.debug(
182 "Preloaded next item %s for queue %s",
183 next_item.name,
184 queue.display_name,
185 )
186 # enqueue the next item on the player
187 self._enqueue_next_item(queue_id, next_item)
188
189 except QueueEmpty:
190 return
191
192 if not (current_item := self.get_item(queue_id, item_id_in_buffer)):
193 # this should not happen, but guard anyways
194 return
195 if current_item.media_type == MediaType.RADIO or not current_item.duration:
196 # radio items or no duration, nothing to do
197 return
198
199 task_id = f"preload_next_item_{queue_id}"
200 self.mass.create_task(
201 _preload_streamdetails,
202 item_id_in_buffer,
203 task_id=task_id,
204 abort_existing=True,
205 )
206
207 async def _cleanup_stale_queue_buffers(self, queue_id: str, current_index: int) -> None:
208 """
209 Clean up audio buffers for queue items that are no longer needed.
210
211 This clears buffers for items at index <= current_index - 2, keeping only:
212 - The previous track (current_index - 1)
213 - The current track (current_index)
214 - The next track (current_index + 1, handled by preloading)
215
216 :param queue_id: The queue ID to clean up buffers for.
217 :param current_index: The current playing index in the queue.
218 """
219 if current_index < 2:
220 return # Nothing to clean up yet
221
222 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
223 cleanup_threshold = current_index - 2
224 buffers_cleared = 0
225
226 for idx, item in enumerate(queue_items):
227 if idx > cleanup_threshold:
228 break # No need to check further
229 if item.streamdetails and item.streamdetails.buffer:
230 self.logger.log(
231 VERBOSE_LOG_LEVEL,
232 "Clearing stale audio buffer for queue item %s (index %d) in queue %s",
233 item.name,
234 idx,
235 queue_id,
236 )
237 await item.streamdetails.buffer.clear()
238 item.streamdetails.buffer = None
239 buffers_cleared += 1
240
241 if buffers_cleared > 0:
242 self.logger.debug(
243 "Cleared %d stale audio buffer(s) for queue %s (items before index %d)",
244 buffers_cleared,
245 queue_id,
246 cleanup_threshold + 1,
247 )
248
249 async def _cleanup_queue_audio_data(self, queue_id: str) -> None:
250 """
251 Clean up all audio-related data for a queue when it is stopped or cleared.
252
253 This clears:
254 - All audio buffers attached to queue item streamdetails
255 - Any pending crossfade data for the queue
256
257 :param queue_id: The queue ID to clean up.
258 """
259 self.mass.streams.audio.clear_crossfade_data(queue_id)
260
261 queue_items = queue_data.items if (queue_data := self._queue_data.get(queue_id)) else []
262 buffers_cleared = 0
263
264 for item in queue_items:
265 if item.streamdetails and item.streamdetails.buffer:
266 await item.streamdetails.buffer.clear()
267 item.streamdetails.buffer = None
268 buffers_cleared += 1
269
270 if buffers_cleared > 0:
271 self.logger.debug(
272 "Cleared %d audio buffer(s) for stopped/cleared queue %s",
273 buffers_cleared,
274 queue_id,
275 )
276